ffi: Create separate timeline object, mirroring the Rust API

This commit is contained in:
Jonas Platte
2023-11-23 18:02:19 +01:00
committed by Jonas Platte
parent e761ad8f97
commit 959e90252b
3 changed files with 640 additions and 736 deletions
+21 -729
View File
@@ -1,64 +1,26 @@
use std::{convert::TryFrom, fmt::Write, fs, sync::Arc};
use std::{convert::TryFrom, sync::Arc};
use anyhow::{anyhow, Context, Result};
use futures_util::{pin_mut, StreamExt};
use matrix_sdk::{
attachment::{
AttachmentConfig, AttachmentInfo, BaseAudioInfo, BaseFileInfo, BaseImageInfo,
BaseThumbnailInfo, BaseVideoInfo, Thumbnail,
},
room::Room as SdkRoom,
ruma::{
api::client::{receipt::create_receipt::v3::ReceiptType, room::report_content},
events::{
location::{AssetType as RumaAssetType, LocationContent, ZoomLevel},
poll::unstable_start::{
UnstablePollAnswer, UnstablePollAnswers, UnstablePollStartContentBlock,
},
receipt::ReceiptThread,
relation::Annotation,
room::{
avatar::ImageInfo as RumaAvatarImageInfo,
message::{
ForwardThread, LocationMessageEventContent, MessageType,
RoomMessageEventContentWithoutRelation,
},
},
AnyMessageLikeEventContent,
},
EventId, UserId,
},
RoomMemberships, RoomState,
};
use matrix_sdk_ui::timeline::{BackPaginationStatus, RoomExt, Timeline};
use anyhow::{Context, Result};
use matrix_sdk::{room::Room as SdkRoom, RoomMemberships, RoomState};
use matrix_sdk_ui::timeline::RoomExt;
use mime::Mime;
use ruma::{
api::client::room::report_content,
assign,
events::{
poll::{
unstable_end::UnstablePollEndEventContent,
unstable_response::UnstablePollResponseEventContent,
unstable_start::NewUnstablePollStartEventContent,
},
room::MediaSource,
},
events::room::{avatar::ImageInfo as RumaAvatarImageInfo, MediaSource},
EventId, UserId,
};
use tokio::{
sync::{Mutex, RwLock},
task::{AbortHandle, JoinHandle},
};
use tracing::{error, info};
use uuid::Uuid;
use tokio::sync::RwLock;
use tracing::error;
use super::RUNTIME;
use crate::{
chunk_iterator::ChunkIterator,
client::ProgressWatcher,
error::{ClientError, MediaInfoError, RoomError},
room_info::RoomInfo,
room_member::{MessageLikeEventType, RoomMember, StateEventType},
ruma::{AssetType, AudioInfo, FileInfo, ImageInfo, PollKind, ThumbnailInfo, VideoInfo},
timeline::{EventTimelineItem, TimelineDiff, TimelineItem, TimelineListener},
ruma::ImageInfo,
timeline::{EventTimelineItem, Timeline},
utils::u64_to_uint,
TaskHandle,
};
@@ -177,31 +139,15 @@ impl Room {
}
}
pub fn retry_decryption(&self, session_ids: Vec<String>) {
let timeline = match &*RUNTIME.block_on(self.timeline.read()) {
Some(t) => Arc::clone(t),
None => {
error!("Timeline not set up, can't retry decryption");
return;
}
};
RUNTIME.spawn(async move {
timeline.retry_decryption(&session_ids).await;
});
}
pub async fn fetch_members(&self) -> Result<(), ClientError> {
let timeline = self
.timeline
.read()
.await
.clone()
.context("Timeline not set up, can't fetch members")?;
timeline.fetch_members().await;
Ok(())
pub async fn timeline(&self) -> Arc<Timeline> {
let mut write_guard = self.timeline.write().await;
if let Some(timeline) = &*write_guard {
timeline.clone()
} else {
let timeline = Timeline::new(self.inner.timeline().await);
*write_guard = Some(timeline.clone());
timeline
}
}
pub fn display_name(&self) -> Result<String, ClientError> {
@@ -247,37 +193,6 @@ impl Room {
})
}
pub async fn add_timeline_listener(
&self,
listener: Box<dyn TimelineListener>,
) -> RoomTimelineListenerResult {
let timeline = {
let mut write_guard = self.timeline.write().await;
if let Some(timeline) = &*write_guard {
timeline.clone()
} else {
let timeline = Arc::new(self.inner.timeline().await);
*write_guard = Some(timeline.clone());
timeline
}
};
let (timeline_items, timeline_stream) = timeline.subscribe_batched().await;
let timeline_stream = TaskHandle::new(RUNTIME.spawn(async move {
pin_mut!(timeline_stream);
while let Some(diffs) = timeline_stream.next().await {
listener
.on_update(diffs.into_iter().map(|d| Arc::new(TimelineDiff::new(d))).collect());
}
}));
RoomTimelineListenerResult {
items: timeline_items.into_iter().map(TimelineItem::from_arc).collect(),
items_stream: Arc::new(timeline_stream),
}
}
pub async fn room_info(&self) -> Result<RoomInfo, ClientError> {
let avatar_url = self.inner.avatar_url();
@@ -286,7 +201,7 @@ impl Room {
// First off, let's see if a `Timeline` exists…
if let Some(timeline) = self.timeline.read().await.clone() {
// If it contains a `latest_event`…
if let Some(timeline_last_event) = timeline.latest_event().await {
if let Some(timeline_last_event) = timeline.inner.latest_event().await {
// If it's a local echo…
if timeline_last_event.is_local_echo() {
return Ok(RoomInfo::new(
@@ -331,211 +246,6 @@ impl Room {
})))
}
pub fn subscribe_to_back_pagination_status(
&self,
listener: Box<dyn BackPaginationStatusListener>,
) -> Result<Arc<TaskHandle>, ClientError> {
let mut subscriber = match &*RUNTIME.block_on(self.timeline.read()) {
Some(t) => t.back_pagination_status(),
None => {
return Err(anyhow!(
"Timeline not set up, can't subscribe to back-pagination status"
)
.into());
}
};
Ok(Arc::new(TaskHandle::new(RUNTIME.spawn(async move {
// Send the current state even if it hasn't changed right away.
listener.on_update(subscriber.next_now());
while let Some(status) = subscriber.next().await {
listener.on_update(status);
}
}))))
}
/// Loads older messages into the timeline.
///
/// Raises an exception if there are no timeline listeners.
pub fn paginate_backwards(&self, opts: PaginationOptions) -> Result<(), ClientError> {
RUNTIME.block_on(async move {
let timeline: Arc<_> = self
.timeline
.read()
.await
.clone()
.context("No timeline listeners registered, can't paginate")?;
Ok(timeline.paginate_backwards(opts.into()).await?)
})
}
pub fn send_read_receipt(&self, event_id: String) -> Result<(), ClientError> {
let event_id = EventId::parse(event_id)?;
RUNTIME.block_on(async move {
self.timeline
.read()
.await
.clone()
.context("Timeline not set up, can't send read receipt")?
.send_single_receipt(ReceiptType::Read, ReceiptThread::Unthreaded, event_id)
.await?;
Ok(())
})
}
pub fn send(&self, msg: Arc<RoomMessageEventContentWithoutRelation>) {
let timeline = match &*RUNTIME.block_on(self.timeline.read()) {
Some(t) => Arc::clone(t),
None => {
error!("Timeline not set up, can't send message");
return;
}
};
RUNTIME.spawn(async move {
timeline.send((*msg).to_owned().with_relation(None).into()).await;
});
}
pub fn create_poll(
&self,
question: String,
answers: Vec<String>,
max_selections: u8,
poll_kind: PollKind,
) -> Result<(), ClientError> {
let timeline = match &*RUNTIME.block_on(self.timeline.read()) {
Some(t) => Arc::clone(t),
None => {
return Err(anyhow!("Timeline not set up, can't send the poll").into());
}
};
let poll_data = PollData { question, answers, max_selections, poll_kind };
let poll_start_event_content = NewUnstablePollStartEventContent::plain_text(
poll_data.fallback_text(),
poll_data.try_into()?,
);
let event_content =
AnyMessageLikeEventContent::UnstablePollStart(poll_start_event_content.into());
RUNTIME.spawn(async move {
timeline.send(event_content).await;
});
Ok(())
}
pub fn send_poll_response(
&self,
poll_start_id: String,
answers: Vec<String>,
) -> Result<(), ClientError> {
let timeline = match &*RUNTIME.block_on(self.timeline.read()) {
Some(t) => Arc::clone(t),
None => {
return Err(anyhow!("Timeline not set up, can't send the poll vote").into());
}
};
let poll_start_event_id =
EventId::parse(poll_start_id).context("Failed to parse EventId")?;
let poll_response_event_content =
UnstablePollResponseEventContent::new(answers, poll_start_event_id);
let event_content =
AnyMessageLikeEventContent::UnstablePollResponse(poll_response_event_content);
RUNTIME.spawn(async move {
timeline.send(event_content).await;
});
Ok(())
}
pub fn end_poll(&self, poll_start_id: String, text: String) -> Result<(), ClientError> {
let timeline = match &*RUNTIME.block_on(self.timeline.read()) {
Some(t) => Arc::clone(t),
None => {
return Err(anyhow!("Timeline not set up, can't end the poll").into());
}
};
let poll_start_event_id =
EventId::parse(poll_start_id).context("Failed to parse EventId")?;
let poll_end_event_content = UnstablePollEndEventContent::new(text, poll_start_event_id);
let event_content = AnyMessageLikeEventContent::UnstablePollEnd(poll_end_event_content);
RUNTIME.spawn(async move {
timeline.send(event_content).await;
});
Ok(())
}
pub fn send_reply(
&self,
msg: Arc<RoomMessageEventContentWithoutRelation>,
reply_item: Arc<EventTimelineItem>,
) -> Result<(), ClientError> {
let timeline = match &*RUNTIME.block_on(self.timeline.read()) {
Some(t) => Arc::clone(t),
None => return Err(anyhow!("Timeline not set up, can't send message").into()),
};
RUNTIME.block_on(async move {
timeline.send_reply((*msg).clone(), &reply_item.0, ForwardThread::Yes).await?;
anyhow::Ok(())
})?;
Ok(())
}
pub fn edit(
&self,
new_content: Arc<RoomMessageEventContentWithoutRelation>,
edit_item: Arc<EventTimelineItem>,
) -> Result<(), ClientError> {
let timeline = match &*RUNTIME.block_on(self.timeline.read()) {
Some(t) => Arc::clone(t),
None => return Err(anyhow!("Timeline not set up, can't send message").into()),
};
RUNTIME.block_on(async move {
timeline.edit((*new_content).clone().with_relation(None), &edit_item.0).await?;
anyhow::Ok(())
})?;
Ok(())
}
pub fn edit_poll(
&self,
question: String,
answers: Vec<String>,
max_selections: u8,
poll_kind: PollKind,
edit_item: Arc<EventTimelineItem>,
) -> Result<(), ClientError> {
let timeline = match &*RUNTIME.block_on(self.timeline.read()) {
Some(t) => Arc::clone(t),
None => return Err(anyhow!("Timeline not set up, can't edit poll").into()),
};
let poll_data = PollData { question, answers, max_selections, poll_kind };
RUNTIME.block_on(async move {
timeline
.edit_poll(poll_data.fallback_text(), poll_data.try_into()?, &edit_item.0)
.await?;
anyhow::Ok(())
})?;
Ok(())
}
/// Redacts an event from the room.
///
/// # Arguments
@@ -552,19 +262,6 @@ impl Room {
})
}
pub fn toggle_reaction(&self, event_id: String, key: String) -> Result<(), ClientError> {
let timeline = match &*RUNTIME.block_on(self.timeline.read()) {
Some(t) => Arc::clone(t),
None => return Err(anyhow!("Timeline not set up, can't send message").into()),
};
RUNTIME.block_on(async move {
let event_id = EventId::parse(event_id)?;
timeline.toggle_reaction(&Annotation::new(event_id, key)).await?;
Ok(())
})
}
pub fn active_members_count(&self) -> u64 {
self.inner.active_members_count()
}
@@ -713,262 +410,6 @@ impl Room {
})
}
pub fn fetch_details_for_event(&self, event_id: String) -> Result<(), ClientError> {
let timeline = RUNTIME
.block_on(self.timeline.read())
.as_ref()
.context("Timeline not set up, can't fetch event details")?
.clone();
RUNTIME.block_on(async move {
let event_id = <&EventId>::try_from(event_id.as_str())?;
timeline.fetch_details_for_event(event_id).await.context("Fetching event details")?;
Ok(())
})
}
pub fn send_image(
self: Arc<Self>,
url: String,
thumbnail_url: String,
image_info: ImageInfo,
progress_watcher: Option<Box<dyn ProgressWatcher>>,
) -> Arc<SendAttachmentJoinHandle> {
SendAttachmentJoinHandle::new(RUNTIME.spawn(async move {
let mime_str =
image_info.mimetype.as_ref().ok_or(RoomError::InvalidAttachmentMimeType)?;
let mime_type =
mime_str.parse::<Mime>().map_err(|_| RoomError::InvalidAttachmentMimeType)?;
let base_image_info = BaseImageInfo::try_from(&image_info)
.map_err(|_| RoomError::InvalidAttachmentData)?;
let attachment_info = AttachmentInfo::Image(base_image_info);
let attachment_config = match image_info.thumbnail_info {
Some(thumbnail_image_info) => {
let thumbnail =
self.build_thumbnail_info(thumbnail_url, thumbnail_image_info)?;
AttachmentConfig::with_thumbnail(thumbnail).info(attachment_info)
}
None => AttachmentConfig::new().info(attachment_info),
};
self.send_attachment(url, mime_type, attachment_config, progress_watcher).await
}))
}
pub fn send_video(
self: Arc<Self>,
url: String,
thumbnail_url: String,
video_info: VideoInfo,
progress_watcher: Option<Box<dyn ProgressWatcher>>,
) -> Arc<SendAttachmentJoinHandle> {
SendAttachmentJoinHandle::new(RUNTIME.spawn(async move {
let mime_str =
video_info.mimetype.as_ref().ok_or(RoomError::InvalidAttachmentMimeType)?;
let mime_type =
mime_str.parse::<Mime>().map_err(|_| RoomError::InvalidAttachmentMimeType)?;
let base_video_info: BaseVideoInfo = BaseVideoInfo::try_from(&video_info)
.map_err(|_| RoomError::InvalidAttachmentData)?;
let attachment_info = AttachmentInfo::Video(base_video_info);
let attachment_config = match video_info.thumbnail_info {
Some(thumbnail_image_info) => {
let thumbnail =
self.build_thumbnail_info(thumbnail_url, thumbnail_image_info)?;
AttachmentConfig::with_thumbnail(thumbnail).info(attachment_info)
}
None => AttachmentConfig::new().info(attachment_info),
};
self.send_attachment(url, mime_type, attachment_config, progress_watcher).await
}))
}
pub fn send_audio(
self: Arc<Self>,
url: String,
audio_info: AudioInfo,
progress_watcher: Option<Box<dyn ProgressWatcher>>,
) -> Arc<SendAttachmentJoinHandle> {
SendAttachmentJoinHandle::new(RUNTIME.spawn(async move {
let mime_str =
audio_info.mimetype.as_ref().ok_or(RoomError::InvalidAttachmentMimeType)?;
let mime_type =
mime_str.parse::<Mime>().map_err(|_| RoomError::InvalidAttachmentMimeType)?;
let base_audio_info: BaseAudioInfo = BaseAudioInfo::try_from(&audio_info)
.map_err(|_| RoomError::InvalidAttachmentData)?;
let attachment_info = AttachmentInfo::Audio(base_audio_info);
let attachment_config = AttachmentConfig::new().info(attachment_info);
self.send_attachment(url, mime_type, attachment_config, progress_watcher).await
}))
}
pub fn send_voice_message(
self: Arc<Self>,
url: String,
audio_info: AudioInfo,
waveform: Vec<u16>,
progress_watcher: Option<Box<dyn ProgressWatcher>>,
) -> Arc<SendAttachmentJoinHandle> {
SendAttachmentJoinHandle::new(RUNTIME.spawn(async move {
let mime_str =
audio_info.mimetype.as_ref().ok_or(RoomError::InvalidAttachmentMimeType)?;
let mime_type =
mime_str.parse::<Mime>().map_err(|_| RoomError::InvalidAttachmentMimeType)?;
let base_audio_info: BaseAudioInfo = BaseAudioInfo::try_from(&audio_info)
.map_err(|_| RoomError::InvalidAttachmentData)?;
let attachment_info =
AttachmentInfo::Voice { audio_info: base_audio_info, waveform: Some(waveform) };
let attachment_config = AttachmentConfig::new().info(attachment_info);
self.send_attachment(url, mime_type, attachment_config, progress_watcher).await
}))
}
pub fn send_file(
self: Arc<Self>,
url: String,
file_info: FileInfo,
progress_watcher: Option<Box<dyn ProgressWatcher>>,
) -> Arc<SendAttachmentJoinHandle> {
SendAttachmentJoinHandle::new(RUNTIME.spawn(async move {
let mime_str =
file_info.mimetype.as_ref().ok_or(RoomError::InvalidAttachmentMimeType)?;
let mime_type =
mime_str.parse::<Mime>().map_err(|_| RoomError::InvalidAttachmentMimeType)?;
let base_file_info: BaseFileInfo =
BaseFileInfo::try_from(&file_info).map_err(|_| RoomError::InvalidAttachmentData)?;
let attachment_info = AttachmentInfo::File(base_file_info);
let attachment_config = AttachmentConfig::new().info(attachment_info);
self.send_attachment(url, mime_type, attachment_config, progress_watcher).await
}))
}
pub fn retry_send(&self, txn_id: String) {
let timeline = match &*RUNTIME.block_on(self.timeline.read()) {
Some(t) => Arc::clone(t),
None => {
error!("Timeline not set up, can't retry sending message");
return;
}
};
RUNTIME.spawn(async move {
if let Err(e) = timeline.retry_send(txn_id.as_str().into()).await {
error!(txn_id, "Failed to retry sending: {e}");
}
});
}
pub fn send_location(
&self,
body: String,
geo_uri: String,
description: Option<String>,
zoom_level: Option<u8>,
asset_type: Option<AssetType>,
) {
let mut location_event_message_content =
LocationMessageEventContent::new(body, geo_uri.clone());
if let Some(asset_type) = asset_type {
location_event_message_content =
location_event_message_content.with_asset_type(RumaAssetType::from(asset_type));
}
let mut location_content = LocationContent::new(geo_uri);
location_content.description = description;
location_content.zoom_level = zoom_level.and_then(ZoomLevel::new);
location_event_message_content.location = Some(location_content);
let room_message_event_content = RoomMessageEventContentWithoutRelation::new(
MessageType::Location(location_event_message_content),
);
self.send(Arc::new(room_message_event_content))
}
pub fn cancel_send(&self, txn_id: String) {
let timeline = match &*RUNTIME.block_on(self.timeline.read()) {
Some(t) => Arc::clone(t),
None => {
error!("Timeline not set up, can't retry sending message");
return;
}
};
RUNTIME.spawn(async move {
if !timeline.cancel_send(txn_id.as_str().into()).await {
info!(txn_id, "Failed to discard local echo: Not found");
}
});
}
pub fn get_event_timeline_item_by_event_id(
&self,
event_id: String,
) -> Result<Arc<EventTimelineItem>, ClientError> {
RUNTIME.block_on(async move {
let timeline = self
.timeline
.read()
.await
.clone()
.context("Timeline not set up, can't get event ")?;
let event_id = EventId::parse(event_id)?;
let item = timeline
.item_by_event_id(&event_id)
.await
.context("Item with given event ID not found")?;
Ok(Arc::new(EventTimelineItem(item)))
})
}
pub fn get_timeline_event_content_by_event_id(
&self,
event_id: String,
) -> Result<Arc<RoomMessageEventContentWithoutRelation>, ClientError> {
RUNTIME.block_on(async move {
let timeline = self
.timeline
.read()
.await
.clone()
.context("Timeline not set up, can't get event content")?;
let event_id = EventId::parse(event_id)?;
let item = timeline
.item_by_event_id(&event_id)
.await
.context("Item with given event ID not found")?;
let msgtype = item
.content()
.as_message()
.context("Item with given event ID is not a message")?
.msgtype()
.to_owned();
Ok(Arc::new(RoomMessageEventContentWithoutRelation::new(msgtype)))
})
}
pub async fn can_user_redact(&self, user_id: String) -> Result<bool, ClientError> {
let user_id = UserId::parse(&user_id)?;
Ok(self.inner.can_user_redact(&user_id).await?)
@@ -1020,122 +461,11 @@ impl Room {
}
}
impl Room {
fn build_thumbnail_info(
&self,
thumbnail_url: String,
thumbnail_info: ThumbnailInfo,
) -> Result<Thumbnail, RoomError> {
let thumbnail_data =
fs::read(thumbnail_url).map_err(|_| RoomError::InvalidThumbnailData)?;
let base_thumbnail_info = BaseThumbnailInfo::try_from(&thumbnail_info)
.map_err(|_| RoomError::InvalidAttachmentData)?;
let mime_str =
thumbnail_info.mimetype.as_ref().ok_or(RoomError::InvalidAttachmentMimeType)?;
let mime_type =
mime_str.parse::<Mime>().map_err(|_| RoomError::InvalidAttachmentMimeType)?;
Ok(Thumbnail {
data: thumbnail_data,
content_type: mime_type,
info: Some(base_thumbnail_info),
})
}
async fn send_attachment(
&self,
url: String,
mime_type: Mime,
attachment_config: AttachmentConfig,
progress_watcher: Option<Box<dyn ProgressWatcher>>,
) -> Result<(), RoomError> {
let timeline = self.timeline.read().await.clone().ok_or(RoomError::TimelineUnavailable)?;
let request = timeline.send_attachment(url, mime_type, attachment_config);
if let Some(progress_watcher) = progress_watcher {
let mut subscriber = request.subscribe_to_send_progress();
RUNTIME.spawn(async move {
while let Some(progress) = subscriber.next().await {
progress_watcher.transmission_progress(progress.into());
}
});
}
request.await.map_err(|_| RoomError::FailedSendingAttachment)?;
Ok(())
}
}
#[derive(uniffi::Object)]
pub struct SendAttachmentJoinHandle {
join_hdl: Arc<Mutex<JoinHandle<Result<(), RoomError>>>>,
abort_hdl: AbortHandle,
}
impl SendAttachmentJoinHandle {
fn new(join_hdl: JoinHandle<Result<(), RoomError>>) -> Arc<Self> {
let abort_hdl = join_hdl.abort_handle();
let join_hdl = Arc::new(Mutex::new(join_hdl));
Arc::new(Self { join_hdl, abort_hdl })
}
}
#[uniffi::export(async_runtime = "tokio")]
impl SendAttachmentJoinHandle {
pub async fn join(&self) -> Result<(), RoomError> {
let join_hdl = self.join_hdl.clone();
RUNTIME.spawn(async move { (&mut *join_hdl.lock().await).await.unwrap() }).await.unwrap()
}
pub fn cancel(&self) {
self.abort_hdl.abort();
}
}
#[derive(uniffi::Record)]
pub struct RoomTimelineListenerResult {
pub items: Vec<Arc<TimelineItem>>,
pub items_stream: Arc<TaskHandle>,
}
#[uniffi::export(callback_interface)]
pub trait RoomInfoListener: Sync + Send {
fn call(&self, room_info: RoomInfo);
}
#[uniffi::export(callback_interface)]
pub trait BackPaginationStatusListener: Sync + Send {
fn on_update(&self, status: BackPaginationStatus);
}
#[derive(uniffi::Enum)]
pub enum PaginationOptions {
SimpleRequest { event_limit: u16, wait_for_token: bool },
UntilNumItems { event_limit: u16, items: u16, wait_for_token: bool },
}
impl From<PaginationOptions> for matrix_sdk_ui::timeline::PaginationOptions<'static> {
fn from(value: PaginationOptions) -> Self {
use matrix_sdk_ui::timeline::PaginationOptions as Opts;
let (wait_for_token, mut opts) = match value {
PaginationOptions::SimpleRequest { event_limit, wait_for_token } => {
(wait_for_token, Opts::simple_request(event_limit))
}
PaginationOptions::UntilNumItems { event_limit, items, wait_for_token } => {
(wait_for_token, Opts::until_num_items(event_limit, items))
}
};
if wait_for_token {
opts = opts.wait_for_token();
}
opts
}
}
#[derive(uniffi::Object)]
pub struct RoomMembersIterator {
chunk_iterator: ChunkIterator<matrix_sdk::room::RoomMember>,
@@ -1184,41 +514,3 @@ impl TryFrom<ImageInfo> for RumaAvatarImageInfo {
}))
}
}
struct PollData {
question: String,
answers: Vec<String>,
max_selections: u8,
poll_kind: PollKind,
}
impl PollData {
fn fallback_text(&self) -> String {
self.answers.iter().enumerate().fold(self.question.clone(), |mut acc, (index, answer)| {
write!(&mut acc, "\n{}. {answer}", index + 1).unwrap();
acc
})
}
}
impl TryFrom<PollData> for UnstablePollStartContentBlock {
type Error = ClientError;
fn try_from(value: PollData) -> Result<Self, Self::Error> {
let poll_answers_vec: Vec<UnstablePollAnswer> = value
.answers
.iter()
.map(|answer| UnstablePollAnswer::new(Uuid::new_v4().to_string(), answer))
.collect();
let poll_answers = UnstablePollAnswers::try_from(poll_answers_vec)
.context("Failed to create poll answers")?;
let mut poll_content_block =
UnstablePollStartContentBlock::new(value.question.clone(), poll_answers);
poll_content_block.kind = value.poll_kind.into();
poll_content_block.max_selections = value.max_selections.into();
Ok(poll_content_block)
}
}
+6 -3
View File
@@ -19,8 +19,11 @@ use matrix_sdk_ui::room_list_service::filters::{
use tokio::sync::RwLock;
use crate::{
error::ClientError, room::Room, room_info::RoomInfo, timeline::EventTimelineItem, TaskHandle,
RUNTIME,
error::ClientError,
room::Room,
room_info::RoomInfo,
timeline::{EventTimelineItem, Timeline},
TaskHandle, RUNTIME,
};
#[derive(Debug, thiserror::Error, uniffi::Error)]
@@ -445,7 +448,7 @@ impl RoomListItem {
async fn full_room(&self) -> Arc<Room> {
Arc::new(Room::with_timeline(
self.inner.inner_room().clone(),
Arc::new(RwLock::new(Some(self.inner.timeline().await))),
Arc::new(RwLock::new(Some(Timeline::from_arc(self.inner.timeline().await)))),
))
}
+613 -4
View File
@@ -12,24 +12,543 @@
// See the License for the specific language governing permissions and
// limitations under the License.
use std::{collections::HashMap, sync::Arc};
use std::{collections::HashMap, fmt::Write as _, fs, sync::Arc};
use anyhow::{Context, Result};
use as_variant::as_variant;
use eyeball_im::VectorDiff;
use matrix_sdk_ui::timeline::{EventItemOrigin, Profile, TimelineDetails};
use tracing::warn;
use futures_util::{pin_mut, StreamExt};
use matrix_sdk::attachment::{
AttachmentConfig, AttachmentInfo, BaseAudioInfo, BaseFileInfo, BaseImageInfo,
BaseThumbnailInfo, BaseVideoInfo, Thumbnail,
};
use matrix_sdk_ui::timeline::{BackPaginationStatus, EventItemOrigin, Profile, TimelineDetails};
use mime::Mime;
use ruma::{
api::client::receipt::create_receipt::v3::ReceiptType,
events::{
location::{AssetType as RumaAssetType, LocationContent, ZoomLevel},
poll::{
unstable_end::UnstablePollEndEventContent,
unstable_response::UnstablePollResponseEventContent,
unstable_start::{
NewUnstablePollStartEventContent, UnstablePollAnswer, UnstablePollAnswers,
UnstablePollStartContentBlock,
},
},
receipt::ReceiptThread,
relation::Annotation,
room::message::{
ForwardThread, LocationMessageEventContent, MessageType,
RoomMessageEventContentWithoutRelation,
},
AnyMessageLikeEventContent,
},
EventId,
};
use tokio::{
sync::Mutex,
task::{AbortHandle, JoinHandle},
};
use tracing::{error, info, warn};
use uuid::Uuid;
use crate::helpers::unwrap_or_clone_arc;
use crate::{
client::ProgressWatcher,
error::{ClientError, RoomError},
helpers::unwrap_or_clone_arc,
ruma::{AssetType, AudioInfo, FileInfo, ImageInfo, PollKind, ThumbnailInfo, VideoInfo},
task_handle::TaskHandle,
RUNTIME,
};
mod content;
pub use self::content::{Reaction, ReactionSenderData, TimelineItemContent};
#[derive(uniffi::Object)]
#[repr(transparent)]
pub struct Timeline {
pub(crate) inner: matrix_sdk_ui::timeline::Timeline,
}
impl Timeline {
pub(crate) fn new(inner: matrix_sdk_ui::timeline::Timeline) -> Arc<Self> {
Arc::new(Self { inner })
}
pub(crate) fn from_arc(inner: Arc<matrix_sdk_ui::timeline::Timeline>) -> Arc<Self> {
// SAFETY: repr(transparent) means transmuting the arc this way is allowed
unsafe { Arc::from_raw(Arc::into_raw(inner) as _) }
}
fn build_thumbnail_info(
&self,
thumbnail_url: String,
thumbnail_info: ThumbnailInfo,
) -> Result<Thumbnail, RoomError> {
let thumbnail_data =
fs::read(thumbnail_url).map_err(|_| RoomError::InvalidThumbnailData)?;
let base_thumbnail_info = BaseThumbnailInfo::try_from(&thumbnail_info)
.map_err(|_| RoomError::InvalidAttachmentData)?;
let mime_str =
thumbnail_info.mimetype.as_ref().ok_or(RoomError::InvalidAttachmentMimeType)?;
let mime_type =
mime_str.parse::<Mime>().map_err(|_| RoomError::InvalidAttachmentMimeType)?;
Ok(Thumbnail {
data: thumbnail_data,
content_type: mime_type,
info: Some(base_thumbnail_info),
})
}
async fn send_attachment(
&self,
url: String,
mime_type: Mime,
attachment_config: AttachmentConfig,
progress_watcher: Option<Box<dyn ProgressWatcher>>,
) -> Result<(), RoomError> {
let request = self.inner.send_attachment(url, mime_type, attachment_config);
if let Some(progress_watcher) = progress_watcher {
let mut subscriber = request.subscribe_to_send_progress();
RUNTIME.spawn(async move {
while let Some(progress) = subscriber.next().await {
progress_watcher.transmission_progress(progress.into());
}
});
}
request.await.map_err(|_| RoomError::FailedSendingAttachment)?;
Ok(())
}
}
#[uniffi::export(async_runtime = "tokio")]
impl Timeline {
pub async fn add_listener(
&self,
listener: Box<dyn TimelineListener>,
) -> RoomTimelineListenerResult {
let (timeline_items, timeline_stream) = self.inner.subscribe_batched().await;
let timeline_stream = TaskHandle::new(RUNTIME.spawn(async move {
pin_mut!(timeline_stream);
while let Some(diffs) = timeline_stream.next().await {
listener
.on_update(diffs.into_iter().map(|d| Arc::new(TimelineDiff::new(d))).collect());
}
}));
RoomTimelineListenerResult {
items: timeline_items.into_iter().map(TimelineItem::from_arc).collect(),
items_stream: Arc::new(timeline_stream),
}
}
pub fn retry_decryption(self: Arc<Self>, session_ids: Vec<String>) {
RUNTIME.spawn(async move {
self.inner.retry_decryption(&session_ids).await;
});
}
pub async fn fetch_members(&self) {
self.inner.fetch_members().await
}
pub fn subscribe_to_back_pagination_status(
&self,
listener: Box<dyn BackPaginationStatusListener>,
) -> Result<Arc<TaskHandle>, ClientError> {
let mut subscriber = self.inner.back_pagination_status();
Ok(Arc::new(TaskHandle::new(RUNTIME.spawn(async move {
// Send the current state even if it hasn't changed right away.
listener.on_update(subscriber.next_now());
while let Some(status) = subscriber.next().await {
listener.on_update(status);
}
}))))
}
/// Loads older messages into the timeline.
///
/// Raises an exception if there are no timeline listeners.
pub fn paginate_backwards(&self, opts: PaginationOptions) -> Result<(), ClientError> {
RUNTIME.block_on(async { Ok(self.inner.paginate_backwards(opts.into()).await?) })
}
pub fn send_read_receipt(&self, event_id: String) -> Result<(), ClientError> {
let event_id = EventId::parse(event_id)?;
RUNTIME.block_on(async {
self.inner
.send_single_receipt(ReceiptType::Read, ReceiptThread::Unthreaded, event_id)
.await?;
Ok(())
})
}
pub fn send(self: Arc<Self>, msg: Arc<RoomMessageEventContentWithoutRelation>) {
RUNTIME.spawn(async move {
self.inner.send((*msg).to_owned().with_relation(None).into()).await;
});
}
pub fn send_image(
self: Arc<Self>,
url: String,
thumbnail_url: String,
image_info: ImageInfo,
progress_watcher: Option<Box<dyn ProgressWatcher>>,
) -> Arc<SendAttachmentJoinHandle> {
SendAttachmentJoinHandle::new(RUNTIME.spawn(async move {
let mime_str =
image_info.mimetype.as_ref().ok_or(RoomError::InvalidAttachmentMimeType)?;
let mime_type =
mime_str.parse::<Mime>().map_err(|_| RoomError::InvalidAttachmentMimeType)?;
let base_image_info = BaseImageInfo::try_from(&image_info)
.map_err(|_| RoomError::InvalidAttachmentData)?;
let attachment_info = AttachmentInfo::Image(base_image_info);
let attachment_config = match image_info.thumbnail_info {
Some(thumbnail_image_info) => {
let thumbnail =
self.build_thumbnail_info(thumbnail_url, thumbnail_image_info)?;
AttachmentConfig::with_thumbnail(thumbnail).info(attachment_info)
}
None => AttachmentConfig::new().info(attachment_info),
};
self.send_attachment(url, mime_type, attachment_config, progress_watcher).await
}))
}
pub fn send_video(
self: Arc<Self>,
url: String,
thumbnail_url: String,
video_info: VideoInfo,
progress_watcher: Option<Box<dyn ProgressWatcher>>,
) -> Arc<SendAttachmentJoinHandle> {
SendAttachmentJoinHandle::new(RUNTIME.spawn(async move {
let mime_str =
video_info.mimetype.as_ref().ok_or(RoomError::InvalidAttachmentMimeType)?;
let mime_type =
mime_str.parse::<Mime>().map_err(|_| RoomError::InvalidAttachmentMimeType)?;
let base_video_info: BaseVideoInfo = BaseVideoInfo::try_from(&video_info)
.map_err(|_| RoomError::InvalidAttachmentData)?;
let attachment_info = AttachmentInfo::Video(base_video_info);
let attachment_config = match video_info.thumbnail_info {
Some(thumbnail_image_info) => {
let thumbnail =
self.build_thumbnail_info(thumbnail_url, thumbnail_image_info)?;
AttachmentConfig::with_thumbnail(thumbnail).info(attachment_info)
}
None => AttachmentConfig::new().info(attachment_info),
};
self.send_attachment(url, mime_type, attachment_config, progress_watcher).await
}))
}
pub fn send_audio(
self: Arc<Self>,
url: String,
audio_info: AudioInfo,
progress_watcher: Option<Box<dyn ProgressWatcher>>,
) -> Arc<SendAttachmentJoinHandle> {
SendAttachmentJoinHandle::new(RUNTIME.spawn(async move {
let mime_str =
audio_info.mimetype.as_ref().ok_or(RoomError::InvalidAttachmentMimeType)?;
let mime_type =
mime_str.parse::<Mime>().map_err(|_| RoomError::InvalidAttachmentMimeType)?;
let base_audio_info: BaseAudioInfo = BaseAudioInfo::try_from(&audio_info)
.map_err(|_| RoomError::InvalidAttachmentData)?;
let attachment_info = AttachmentInfo::Audio(base_audio_info);
let attachment_config = AttachmentConfig::new().info(attachment_info);
self.send_attachment(url, mime_type, attachment_config, progress_watcher).await
}))
}
pub fn send_voice_message(
self: Arc<Self>,
url: String,
audio_info: AudioInfo,
waveform: Vec<u16>,
progress_watcher: Option<Box<dyn ProgressWatcher>>,
) -> Arc<SendAttachmentJoinHandle> {
SendAttachmentJoinHandle::new(RUNTIME.spawn(async move {
let mime_str =
audio_info.mimetype.as_ref().ok_or(RoomError::InvalidAttachmentMimeType)?;
let mime_type =
mime_str.parse::<Mime>().map_err(|_| RoomError::InvalidAttachmentMimeType)?;
let base_audio_info: BaseAudioInfo = BaseAudioInfo::try_from(&audio_info)
.map_err(|_| RoomError::InvalidAttachmentData)?;
let attachment_info =
AttachmentInfo::Voice { audio_info: base_audio_info, waveform: Some(waveform) };
let attachment_config = AttachmentConfig::new().info(attachment_info);
self.send_attachment(url, mime_type, attachment_config, progress_watcher).await
}))
}
pub fn send_file(
self: Arc<Self>,
url: String,
file_info: FileInfo,
progress_watcher: Option<Box<dyn ProgressWatcher>>,
) -> Arc<SendAttachmentJoinHandle> {
SendAttachmentJoinHandle::new(RUNTIME.spawn(async move {
let mime_str =
file_info.mimetype.as_ref().ok_or(RoomError::InvalidAttachmentMimeType)?;
let mime_type =
mime_str.parse::<Mime>().map_err(|_| RoomError::InvalidAttachmentMimeType)?;
let base_file_info: BaseFileInfo =
BaseFileInfo::try_from(&file_info).map_err(|_| RoomError::InvalidAttachmentData)?;
let attachment_info = AttachmentInfo::File(base_file_info);
let attachment_config = AttachmentConfig::new().info(attachment_info);
self.send_attachment(url, mime_type, attachment_config, progress_watcher).await
}))
}
pub fn create_poll(
self: Arc<Self>,
question: String,
answers: Vec<String>,
max_selections: u8,
poll_kind: PollKind,
) -> Result<(), ClientError> {
let poll_data = PollData { question, answers, max_selections, poll_kind };
let poll_start_event_content = NewUnstablePollStartEventContent::plain_text(
poll_data.fallback_text(),
poll_data.try_into()?,
);
let event_content =
AnyMessageLikeEventContent::UnstablePollStart(poll_start_event_content.into());
RUNTIME.spawn(async move {
self.inner.send(event_content).await;
});
Ok(())
}
pub fn send_poll_response(
self: Arc<Self>,
poll_start_id: String,
answers: Vec<String>,
) -> Result<(), ClientError> {
let poll_start_event_id =
EventId::parse(poll_start_id).context("Failed to parse EventId")?;
let poll_response_event_content =
UnstablePollResponseEventContent::new(answers, poll_start_event_id);
let event_content =
AnyMessageLikeEventContent::UnstablePollResponse(poll_response_event_content);
RUNTIME.spawn(async move {
self.inner.send(event_content).await;
});
Ok(())
}
pub fn end_poll(
self: Arc<Self>,
poll_start_id: String,
text: String,
) -> Result<(), ClientError> {
let poll_start_event_id =
EventId::parse(poll_start_id).context("Failed to parse EventId")?;
let poll_end_event_content = UnstablePollEndEventContent::new(text, poll_start_event_id);
let event_content = AnyMessageLikeEventContent::UnstablePollEnd(poll_end_event_content);
RUNTIME.spawn(async move {
self.inner.send(event_content).await;
});
Ok(())
}
pub fn send_reply(
&self,
msg: Arc<RoomMessageEventContentWithoutRelation>,
reply_item: Arc<EventTimelineItem>,
) -> Result<(), ClientError> {
RUNTIME.block_on(async {
self.inner.send_reply((*msg).clone(), &reply_item.0, ForwardThread::Yes).await?;
anyhow::Ok(())
})?;
Ok(())
}
pub fn edit(
&self,
new_content: Arc<RoomMessageEventContentWithoutRelation>,
edit_item: Arc<EventTimelineItem>,
) -> Result<(), ClientError> {
RUNTIME.block_on(async {
self.inner.edit((*new_content).clone().with_relation(None), &edit_item.0).await?;
anyhow::Ok(())
})?;
Ok(())
}
pub async fn edit_poll(
&self,
question: String,
answers: Vec<String>,
max_selections: u8,
poll_kind: PollKind,
edit_item: Arc<EventTimelineItem>,
) -> Result<(), ClientError> {
let poll_data = PollData { question, answers, max_selections, poll_kind };
RUNTIME.block_on(async {
self.inner
.edit_poll(poll_data.fallback_text(), poll_data.try_into()?, &edit_item.0)
.await?;
anyhow::Ok(())
})?;
Ok(())
}
pub fn send_location(
self: Arc<Self>,
body: String,
geo_uri: String,
description: Option<String>,
zoom_level: Option<u8>,
asset_type: Option<AssetType>,
) {
let mut location_event_message_content =
LocationMessageEventContent::new(body, geo_uri.clone());
if let Some(asset_type) = asset_type {
location_event_message_content =
location_event_message_content.with_asset_type(RumaAssetType::from(asset_type));
}
let mut location_content = LocationContent::new(geo_uri);
location_content.description = description;
location_content.zoom_level = zoom_level.and_then(ZoomLevel::new);
location_event_message_content.location = Some(location_content);
let room_message_event_content = RoomMessageEventContentWithoutRelation::new(
MessageType::Location(location_event_message_content),
);
self.send(Arc::new(room_message_event_content))
}
pub fn toggle_reaction(&self, event_id: String, key: String) -> Result<(), ClientError> {
let event_id = EventId::parse(event_id)?;
RUNTIME.block_on(async {
self.inner.toggle_reaction(&Annotation::new(event_id, key)).await?;
Ok(())
})
}
pub fn fetch_details_for_event(&self, event_id: String) -> Result<(), ClientError> {
let event_id = <&EventId>::try_from(event_id.as_str())?;
RUNTIME.block_on(async {
self.inner.fetch_details_for_event(event_id).await.context("Fetching event details")?;
Ok(())
})
}
pub fn retry_send(self: Arc<Self>, txn_id: String) {
RUNTIME.spawn(async move {
if let Err(e) = self.inner.retry_send(txn_id.as_str().into()).await {
error!(txn_id, "Failed to retry sending: {e}");
}
});
}
pub fn cancel_send(self: Arc<Self>, txn_id: String) {
RUNTIME.spawn(async move {
if !self.inner.cancel_send(txn_id.as_str().into()).await {
info!(txn_id, "Failed to discard local echo: Not found");
}
});
}
pub fn get_event_timeline_item_by_event_id(
&self,
event_id: String,
) -> Result<Arc<EventTimelineItem>, ClientError> {
let event_id = EventId::parse(event_id)?;
RUNTIME.block_on(async {
let item = self
.inner
.item_by_event_id(&event_id)
.await
.context("Item with given event ID not found")?;
Ok(Arc::new(EventTimelineItem(item)))
})
}
pub fn get_timeline_event_content_by_event_id(
&self,
event_id: String,
) -> Result<Arc<RoomMessageEventContentWithoutRelation>, ClientError> {
let event_id = EventId::parse(event_id)?;
RUNTIME.block_on(async {
let item = self
.inner
.item_by_event_id(&event_id)
.await
.context("Item with given event ID not found")?;
let msgtype = item
.content()
.as_message()
.context("Item with given event ID is not a message")?
.msgtype()
.to_owned();
Ok(Arc::new(RoomMessageEventContentWithoutRelation::new(msgtype)))
})
}
}
#[derive(uniffi::Record)]
pub struct RoomTimelineListenerResult {
pub items: Vec<Arc<TimelineItem>>,
pub items_stream: Arc<TaskHandle>,
}
#[uniffi::export(callback_interface)]
pub trait TimelineListener: Sync + Send {
fn on_update(&self, diff: Vec<Arc<TimelineDiff>>);
}
#[uniffi::export(callback_interface)]
pub trait BackPaginationStatusListener: Sync + Send {
fn on_update(&self, status: BackPaginationStatus);
}
#[derive(Clone, uniffi::Object)]
pub enum TimelineDiff {
Append { values: Vec<Arc<TimelineItem>> },
@@ -353,6 +872,96 @@ impl From<&TimelineDetails<Profile>> for ProfileDetails {
}
}
struct PollData {
question: String,
answers: Vec<String>,
max_selections: u8,
poll_kind: PollKind,
}
impl PollData {
fn fallback_text(&self) -> String {
self.answers.iter().enumerate().fold(self.question.clone(), |mut acc, (index, answer)| {
write!(&mut acc, "\n{}. {answer}", index + 1).unwrap();
acc
})
}
}
impl TryFrom<PollData> for UnstablePollStartContentBlock {
type Error = ClientError;
fn try_from(value: PollData) -> Result<Self, Self::Error> {
let poll_answers_vec: Vec<UnstablePollAnswer> = value
.answers
.iter()
.map(|answer| UnstablePollAnswer::new(Uuid::new_v4().to_string(), answer))
.collect();
let poll_answers = UnstablePollAnswers::try_from(poll_answers_vec)
.context("Failed to create poll answers")?;
let mut poll_content_block =
UnstablePollStartContentBlock::new(value.question.clone(), poll_answers);
poll_content_block.kind = value.poll_kind.into();
poll_content_block.max_selections = value.max_selections.into();
Ok(poll_content_block)
}
}
#[derive(uniffi::Object)]
pub struct SendAttachmentJoinHandle {
join_hdl: Arc<Mutex<JoinHandle<Result<(), RoomError>>>>,
abort_hdl: AbortHandle,
}
impl SendAttachmentJoinHandle {
fn new(join_hdl: JoinHandle<Result<(), RoomError>>) -> Arc<Self> {
let abort_hdl = join_hdl.abort_handle();
let join_hdl = Arc::new(Mutex::new(join_hdl));
Arc::new(Self { join_hdl, abort_hdl })
}
}
#[uniffi::export(async_runtime = "tokio")]
impl SendAttachmentJoinHandle {
pub async fn join(&self) -> Result<(), RoomError> {
let join_hdl = self.join_hdl.clone();
RUNTIME.spawn(async move { (&mut *join_hdl.lock().await).await.unwrap() }).await.unwrap()
}
pub fn cancel(&self) {
self.abort_hdl.abort();
}
}
#[derive(uniffi::Enum)]
pub enum PaginationOptions {
SimpleRequest { event_limit: u16, wait_for_token: bool },
UntilNumItems { event_limit: u16, items: u16, wait_for_token: bool },
}
impl From<PaginationOptions> for matrix_sdk_ui::timeline::PaginationOptions<'static> {
fn from(value: PaginationOptions) -> Self {
use matrix_sdk_ui::timeline::PaginationOptions as Opts;
let (wait_for_token, mut opts) = match value {
PaginationOptions::SimpleRequest { event_limit, wait_for_token } => {
(wait_for_token, Opts::simple_request(event_limit))
}
PaginationOptions::UntilNumItems { event_limit, items, wait_for_token } => {
(wait_for_token, Opts::until_num_items(event_limit, items))
}
};
if wait_for_token {
opts = opts.wait_for_token();
}
opts
}
}
/// A [`TimelineItem`](super::TimelineItem) that doesn't correspond to an event.
#[derive(uniffi::Enum)]
pub enum VirtualTimelineItem {