mod.rs (19174B)
1 use std::collections::HashSet; 2 use std::sync::Arc; 3 4 use log::*; 5 use matrix_sdk::{ 6 api::r0::{ 7 filter::RoomEventFilter, 8 message::{ 9 get_message_events::Direction, get_message_events::Request as GetMessagesRequest, 10 }, 11 }, 12 events::{ 13 room::message::{FormattedBody, MessageEventContent, TextMessageEventContent}, 14 AnyMessageEvent, AnyRoomEvent, AnySyncMessageEvent, 15 }, 16 identifiers::RoomId, 17 js_int::uint, 18 locks::RwLock, 19 Client, Raw, Room, 20 }; 21 use pulldown_cmark::{html, Options, Parser}; 22 use serde::{Deserialize, Serialize}; 23 use wasm_bindgen_futures::spawn_local; 24 use yew::worker::*; 25 26 use crate::app::matrix::types::{get_media_download_url, get_video_media_download_url}; 27 use crate::errors::MatrixError; 28 use login::{login, SessionStore}; 29 30 pub mod login; 31 mod sync; 32 pub mod types; 33 34 #[derive(Default, Clone, Debug)] 35 pub struct MatrixClient { 36 pub(crate) homeserver: Option<String>, 37 pub(crate) username: Option<String>, 38 pub(crate) password: Option<String>, 39 } 40 41 #[derive(Clone, Debug)] 42 pub struct MatrixAgent { 43 link: AgentLink<MatrixAgent>, 44 matrix_state: MatrixClient, 45 matrix_client: Option<Client>, 46 // TODO make arc mutex :( 47 subscribers: HashSet<HandlerId>, 48 session: Option<SessionStore>, 49 } 50 51 #[derive(Serialize, Deserialize, Debug)] 52 pub enum Request { 53 SetHomeserver(String), 54 SetUsername(String), 55 SetPassword(String), 56 SetSession(SessionStore), 57 Login, 58 GetLoggedIn, 59 GetOldMessages((RoomId, Option<String>)), 60 StartSync, 61 GetJoinedRoom(RoomId), 62 SendMessage((RoomId, String)), 63 } 64 65 #[allow(clippy::large_enum_variant)] 66 #[derive(Serialize, Deserialize, Debug, Clone)] 67 pub enum Response { 68 Error(MatrixError), 69 LoggedIn(bool), 70 // TODO properly handle sync events 71 Sync((RoomId, Raw<AnySyncMessageEvent>)), 72 JoinedRoomSync(RoomId), 73 SyncPing, 74 OldMessages((RoomId, Vec<Raw<AnyMessageEvent>>)), 75 JoinedRoom((RoomId, Room)), 76 SaveSession(SessionStore), 77 } 78 79 #[derive(Debug, Clone)] 80 pub enum Msg { 81 OnSyncResponse(Response), 82 } 83 84 impl Agent for MatrixAgent { 85 type Reach = Public<Self>; 86 type Message = Msg; 87 type Input = Request; 88 type Output = Response; 89 90 fn create(link: AgentLink<Self>) -> Self { 91 MatrixAgent { 92 link, 93 matrix_state: Default::default(), 94 matrix_client: None, 95 subscribers: HashSet::new(), 96 session: Default::default(), 97 } 98 } 99 100 fn update(&mut self, msg: Self::Message) { 101 match msg { 102 Msg::OnSyncResponse(resp) => { 103 for sub in self.subscribers.iter() { 104 self.link.respond(*sub, resp.clone()); 105 } 106 } 107 } 108 } 109 110 fn connected(&mut self, id: HandlerId) { 111 self.subscribers.insert(id); 112 } 113 fn handle_input(&mut self, msg: Self::Input, _: HandlerId) { 114 match msg { 115 Request::SetSession(session) => { 116 self.session = Some(session); 117 } 118 119 Request::SetHomeserver(homeserver) => { 120 self.matrix_state.homeserver = Some(homeserver); 121 } 122 Request::SetUsername(username) => { 123 self.matrix_state.username = Some(username); 124 } 125 Request::SetPassword(password) => { 126 self.matrix_state.password = Some(password); 127 } 128 Request::Login => { 129 info!("Starting Login"); 130 let homeserver = self.matrix_state.homeserver.as_ref(); 131 let session = self.session.as_ref(); 132 let client = login(session, homeserver); 133 match client { 134 Ok(client) => { 135 if let Some(_session) = session { 136 for sub in self.subscribers.iter() { 137 let resp = Response::LoggedIn(true); 138 self.link.respond(*sub, resp); 139 } 140 } 141 self.matrix_client = Some(client.clone()); 142 let username = self.matrix_state.username.clone().unwrap(); 143 let password = self.matrix_state.password.clone().unwrap(); 144 let agent = self.clone(); 145 spawn_local(async move { 146 // FIXME gracefully handle login errors 147 // TODO make the String to &str conversion smarter if possible 148 let login_response = agent 149 .matrix_client 150 .as_ref() 151 .unwrap() 152 .login(&username, &password, None, Some("Daydream")) 153 .await; 154 match login_response { 155 Ok(login_response) => { 156 let session_store = SessionStore { 157 access_token: login_response.access_token, 158 user_id: login_response.user_id.to_string(), 159 device_id: login_response.device_id.into(), 160 homeserver_url: client.homeserver().to_string(), 161 }; 162 for sub in agent.subscribers.iter() { 163 let resp = Response::SaveSession(session_store.clone()); 164 agent.link.respond(*sub, resp); 165 let resp = Response::LoggedIn(true); 166 agent.link.respond(*sub, resp); 167 } 168 } 169 Err(e) => { 170 if let matrix_sdk::Error::Reqwest(e) = e { 171 match e.status() { 172 None => { 173 for sub in agent.subscribers.iter() { 174 let resp = Response::Error( 175 MatrixError::SDKError(e.to_string()), 176 ); 177 agent.link.respond(*sub, resp); 178 } 179 } 180 Some(v) => { 181 if v.is_server_error() { 182 for sub in agent.subscribers.iter() { 183 let resp = Response::Error( 184 MatrixError::LoginTimeout, 185 ); 186 agent.link.respond(*sub, resp); 187 } 188 } else { 189 for sub in agent.subscribers.iter() { 190 let resp = Response::Error( 191 MatrixError::SDKError(e.to_string()), 192 ); 193 agent.link.respond(*sub, resp); 194 } 195 } 196 } 197 } 198 } else { 199 for sub in agent.subscribers.iter() { 200 let resp = Response::Error(MatrixError::SDKError( 201 e.to_string(), 202 )); 203 agent.link.respond(*sub, resp); 204 } 205 } 206 } 207 } 208 }); 209 } 210 Err(e) => { 211 for sub in self.subscribers.iter() { 212 let resp = Response::Error(e.clone()); 213 self.link.respond(*sub, resp); 214 } 215 } 216 } 217 } 218 Request::GetLoggedIn => { 219 let homeserver = self.matrix_state.homeserver.as_ref(); 220 let session = self.session.clone(); 221 let client = login(session.as_ref(), homeserver); 222 match client { 223 Ok(client) => { 224 info!("Got client"); 225 self.matrix_client = Some(client); 226 info!("Client set"); 227 let agent = self.clone(); 228 spawn_local(async move { 229 let logged_in = agent.get_logged_in().await; 230 231 if !logged_in && session.is_some() { 232 error!("Not logged in but got session"); 233 } else { 234 for sub in agent.subscribers.iter() { 235 let resp = Response::LoggedIn(logged_in); 236 agent.link.respond(*sub, resp); 237 } 238 } 239 }); 240 } 241 Err(e) => { 242 error!("Got no client: {:?}", e); 243 for sub in self.subscribers.iter() { 244 let resp = Response::Error(e.clone()); 245 self.link.respond(*sub, resp); 246 } 247 } 248 } 249 } 250 Request::StartSync => { 251 // Always clone agent after having tried to login! 252 let agent = self.clone(); 253 spawn_local(async move { 254 agent.start_sync().await; 255 }); 256 } 257 Request::GetOldMessages((room_id, from)) => { 258 let agent = self.clone(); 259 spawn_local(async move { 260 let sync_token = match from { 261 Some(from) => from, 262 None => agent 263 .matrix_client 264 .as_ref() 265 .unwrap() 266 .sync_token() 267 .await 268 .unwrap(), 269 }; 270 let mut req = 271 GetMessagesRequest::new(&room_id, &sync_token, Direction::Backward); 272 let filter = RoomEventFilter { 273 types: Some(vec!["m.room.message".to_string()]), 274 ..Default::default() 275 }; 276 // TODO find better way than cloning 277 req.filter = Some(filter); 278 req.limit = uint!(30); 279 280 // TODO handle error gracefully 281 let messsages = agent 282 .matrix_client 283 .clone() 284 .unwrap() 285 .room_messages(req) 286 .await 287 .unwrap(); 288 // TODO save end point for future loading 289 290 let mut wrapped_messages: Vec<Raw<AnyMessageEvent>> = Vec::new(); 291 let chunk_iter: Vec<Raw<AnyRoomEvent>> = messsages.chunk; 292 let (oks, _): (Vec<_>, Vec<_>) = chunk_iter 293 .iter() 294 .map(|event| event.deserialize()) 295 .partition(Result::is_ok); 296 297 let deserialized_events: Vec<AnyRoomEvent> = 298 oks.into_iter().map(Result::unwrap).collect(); 299 300 for event in deserialized_events.into_iter().rev() { 301 // TODO deduplicate betweeen this and sync 302 if let AnyRoomEvent::Message(AnyMessageEvent::RoomMessage(mut event)) = 303 event 304 { 305 if let MessageEventContent::Image(mut image_event) = 306 event.clone().content 307 { 308 if let Some(image_event_url) = image_event.url { 309 let new_url = get_media_download_url( 310 agent.matrix_client.as_ref().unwrap().homeserver(), 311 &image_event_url, 312 ); 313 image_event.url = Some(new_url.to_string()); 314 } 315 if let Some(mut info) = image_event.info { 316 if let Some(thumbnail_url) = info.thumbnail_url.as_ref() { 317 let new_url = get_media_download_url( 318 agent.matrix_client.as_ref().unwrap().homeserver(), 319 thumbnail_url, 320 ); 321 info.thumbnail_url = Some(new_url.to_string()); 322 } 323 image_event.info = Some(info); 324 } 325 event.content = MessageEventContent::Image(image_event); 326 } 327 if let MessageEventContent::Video(mut video_event) = event.content { 328 if let Some(video_event_url) = video_event.url { 329 let new_url = get_video_media_download_url( 330 agent.matrix_client.as_ref().unwrap().homeserver(), 331 video_event_url, 332 ); 333 video_event.url = Some(new_url.to_string()); 334 } 335 if let Some(mut info) = video_event.info { 336 if let Some(thumbnail_url) = info.thumbnail_url { 337 let new_url = Some(get_media_download_url( 338 agent.matrix_client.as_ref().unwrap().homeserver(), 339 &thumbnail_url, 340 )) 341 .unwrap(); 342 info.thumbnail_url = Some(new_url.to_string()); 343 } 344 video_event.info = Some(info); 345 } 346 event.content = MessageEventContent::Video(video_event); 347 } 348 349 let serialized_event = 350 Raw::from(AnyMessageEvent::RoomMessage(event.clone())); 351 wrapped_messages.push(serialized_event); 352 } 353 } 354 355 for sub in agent.subscribers.iter() { 356 let resp = 357 Response::OldMessages((room_id.clone(), wrapped_messages.clone())); 358 agent.link.respond(*sub, resp); 359 } 360 }); 361 } 362 Request::GetJoinedRoom(room_id) => { 363 let agent = self.clone(); 364 spawn_local(async move { 365 let room: Arc<RwLock<Room>> = agent 366 .matrix_client 367 .unwrap() 368 .get_joined_room(&room_id) 369 .await 370 .unwrap(); 371 let read_clone = room.read().await; 372 let clean_room = (*read_clone).clone(); 373 for sub in agent.subscribers.iter() { 374 let resp = Response::JoinedRoom((room_id.clone(), clean_room.clone())); 375 agent.link.respond(*sub, resp); 376 } 377 }); 378 } 379 Request::SendMessage((room_id, raw_message)) => { 380 let client = self.matrix_client.clone().unwrap(); 381 spawn_local(async move { 382 let replacer = gh_emoji::Replacer::new(); 383 let message = replacer.replace_all(raw_message.as_str()); 384 385 let mut options = Options::empty(); 386 options.insert(Options::ENABLE_STRIKETHROUGH); 387 let parser = Parser::new_ext(message.as_ref(), options); 388 389 let mut formatted_message: String = 390 String::with_capacity(message.len() * 3 / 2); 391 html::push_html(&mut formatted_message, parser); 392 formatted_message = formatted_message.replace("<p>", "").replace("</p>", ""); 393 formatted_message.pop(); 394 395 let content = if formatted_message == message { 396 MessageEventContent::Text(TextMessageEventContent::plain(message)) 397 } else { 398 MessageEventContent::Text(TextMessageEventContent { 399 body: message.to_string(), 400 relates_to: None, 401 formatted: Some(FormattedBody::html(formatted_message)), 402 }) 403 }; 404 if let Err(e) = client.room_send(&room_id, content, None).await { 405 // TODO show error in UI or try again if possible 406 error!("Error sending message: {}", e); 407 } 408 }); 409 } 410 } 411 } 412 413 fn disconnected(&mut self, id: HandlerId) { 414 self.subscribers.remove(&id); 415 } 416 417 fn name_of_resource() -> &'static str { 418 "worker.js" 419 } 420 } 421 422 unsafe impl Send for MatrixAgent {} 423 424 unsafe impl std::marker::Sync for MatrixAgent {} 425 426 impl MatrixAgent { 427 async fn start_sync(&self) { 428 let sync = sync::Sync { 429 matrix_client: self.matrix_client.clone().unwrap(), 430 callback: self.link.callback(Msg::OnSyncResponse), 431 }; 432 sync.start_sync().await; 433 } 434 435 async fn get_logged_in(&self) -> bool { 436 if self.matrix_client.is_none() { 437 return false; 438 } 439 self.matrix_client.as_ref().unwrap().logged_in().await 440 } 441 }