diff --git a/bindings/matrix-sdk-ffi/src/room.rs b/bindings/matrix-sdk-ffi/src/room.rs index e250b619e..5c87e3fdf 100644 --- a/bindings/matrix-sdk-ffi/src/room.rs +++ b/bindings/matrix-sdk-ffi/src/room.rs @@ -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) { - 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 { + 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 { @@ -247,37 +193,6 @@ impl Room { }) } - pub async fn add_timeline_listener( - &self, - listener: Box, - ) -> 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 { 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, - ) -> Result, 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) { - 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, - 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, - ) -> 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, - reply_item: Arc, - ) -> 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, - edit_item: Arc, - ) -> 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, - max_selections: u8, - poll_kind: PollKind, - edit_item: Arc, - ) -> 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, - url: String, - thumbnail_url: String, - image_info: ImageInfo, - progress_watcher: Option>, - ) -> Arc { - SendAttachmentJoinHandle::new(RUNTIME.spawn(async move { - let mime_str = - image_info.mimetype.as_ref().ok_or(RoomError::InvalidAttachmentMimeType)?; - let mime_type = - mime_str.parse::().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, - url: String, - thumbnail_url: String, - video_info: VideoInfo, - progress_watcher: Option>, - ) -> Arc { - SendAttachmentJoinHandle::new(RUNTIME.spawn(async move { - let mime_str = - video_info.mimetype.as_ref().ok_or(RoomError::InvalidAttachmentMimeType)?; - let mime_type = - mime_str.parse::().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, - url: String, - audio_info: AudioInfo, - progress_watcher: Option>, - ) -> Arc { - SendAttachmentJoinHandle::new(RUNTIME.spawn(async move { - let mime_str = - audio_info.mimetype.as_ref().ok_or(RoomError::InvalidAttachmentMimeType)?; - let mime_type = - mime_str.parse::().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, - url: String, - audio_info: AudioInfo, - waveform: Vec, - progress_watcher: Option>, - ) -> Arc { - SendAttachmentJoinHandle::new(RUNTIME.spawn(async move { - let mime_str = - audio_info.mimetype.as_ref().ok_or(RoomError::InvalidAttachmentMimeType)?; - let mime_type = - mime_str.parse::().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, - url: String, - file_info: FileInfo, - progress_watcher: Option>, - ) -> Arc { - SendAttachmentJoinHandle::new(RUNTIME.spawn(async move { - let mime_str = - file_info.mimetype.as_ref().ok_or(RoomError::InvalidAttachmentMimeType)?; - let mime_type = - mime_str.parse::().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, - zoom_level: Option, - asset_type: Option, - ) { - 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, 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, 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 { 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 { - 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::().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>, - ) -> 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>>>, - abort_hdl: AbortHandle, -} - -impl SendAttachmentJoinHandle { - fn new(join_hdl: JoinHandle>) -> Arc { - 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>, - pub items_stream: Arc, -} - #[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 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, @@ -1184,41 +514,3 @@ impl TryFrom for RumaAvatarImageInfo { })) } } - -struct PollData { - question: String, - answers: Vec, - 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 for UnstablePollStartContentBlock { - type Error = ClientError; - - fn try_from(value: PollData) -> Result { - let poll_answers_vec: Vec = 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) - } -} diff --git a/bindings/matrix-sdk-ffi/src/room_list.rs b/bindings/matrix-sdk-ffi/src/room_list.rs index 5b615a975..372dd47a6 100644 --- a/bindings/matrix-sdk-ffi/src/room_list.rs +++ b/bindings/matrix-sdk-ffi/src/room_list.rs @@ -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 { 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)))), )) } diff --git a/bindings/matrix-sdk-ffi/src/timeline/mod.rs b/bindings/matrix-sdk-ffi/src/timeline/mod.rs index 989aa6228..b4bde590e 100644 --- a/bindings/matrix-sdk-ffi/src/timeline/mod.rs +++ b/bindings/matrix-sdk-ffi/src/timeline/mod.rs @@ -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 { + Arc::new(Self { inner }) + } + + pub(crate) fn from_arc(inner: Arc) -> Arc { + // 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 { + 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::().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>, + ) -> 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, + ) -> 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, session_ids: Vec) { + 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, + ) -> Result, 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, msg: Arc) { + RUNTIME.spawn(async move { + self.inner.send((*msg).to_owned().with_relation(None).into()).await; + }); + } + + pub fn send_image( + self: Arc, + url: String, + thumbnail_url: String, + image_info: ImageInfo, + progress_watcher: Option>, + ) -> Arc { + SendAttachmentJoinHandle::new(RUNTIME.spawn(async move { + let mime_str = + image_info.mimetype.as_ref().ok_or(RoomError::InvalidAttachmentMimeType)?; + let mime_type = + mime_str.parse::().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, + url: String, + thumbnail_url: String, + video_info: VideoInfo, + progress_watcher: Option>, + ) -> Arc { + SendAttachmentJoinHandle::new(RUNTIME.spawn(async move { + let mime_str = + video_info.mimetype.as_ref().ok_or(RoomError::InvalidAttachmentMimeType)?; + let mime_type = + mime_str.parse::().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, + url: String, + audio_info: AudioInfo, + progress_watcher: Option>, + ) -> Arc { + SendAttachmentJoinHandle::new(RUNTIME.spawn(async move { + let mime_str = + audio_info.mimetype.as_ref().ok_or(RoomError::InvalidAttachmentMimeType)?; + let mime_type = + mime_str.parse::().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, + url: String, + audio_info: AudioInfo, + waveform: Vec, + progress_watcher: Option>, + ) -> Arc { + SendAttachmentJoinHandle::new(RUNTIME.spawn(async move { + let mime_str = + audio_info.mimetype.as_ref().ok_or(RoomError::InvalidAttachmentMimeType)?; + let mime_type = + mime_str.parse::().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, + url: String, + file_info: FileInfo, + progress_watcher: Option>, + ) -> Arc { + SendAttachmentJoinHandle::new(RUNTIME.spawn(async move { + let mime_str = + file_info.mimetype.as_ref().ok_or(RoomError::InvalidAttachmentMimeType)?; + let mime_type = + mime_str.parse::().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, + question: String, + answers: Vec, + 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, + poll_start_id: String, + answers: Vec, + ) -> 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, + 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, + reply_item: Arc, + ) -> 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, + edit_item: Arc, + ) -> 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, + max_selections: u8, + poll_kind: PollKind, + edit_item: Arc, + ) -> 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, + body: String, + geo_uri: String, + description: Option, + zoom_level: Option, + asset_type: Option, + ) { + 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, 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, 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, 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, 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>, + pub items_stream: Arc, +} + #[uniffi::export(callback_interface)] pub trait TimelineListener: Sync + Send { fn on_update(&self, diff: Vec>); } +#[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> }, @@ -353,6 +872,96 @@ impl From<&TimelineDetails> for ProfileDetails { } } +struct PollData { + question: String, + answers: Vec, + 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 for UnstablePollStartContentBlock { + type Error = ClientError; + + fn try_from(value: PollData) -> Result { + let poll_answers_vec: Vec = 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>>>, + abort_hdl: AbortHandle, +} + +impl SendAttachmentJoinHandle { + fn new(join_hdl: JoinHandle>) -> Arc { + 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 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 {