main.rs (19173B)
1 use std::convert::TryFrom; 2 use std::fs::{File, OpenOptions}; 3 use std::io::Write; 4 use std::path::PathBuf; 5 use std::{env, fs, process::exit}; 6 7 use bytes::buf::Buf; 8 use ipfs_api::{IpfsClient, TryFromUri}; 9 use matrix_sdk::{ 10 self, 11 events::collections::all::RoomEvent, 12 events::room::{ 13 ThumbnailInfo, 14 ImageInfo, 15 member::MemberEventContent, 16 message::{MessageEvent, MessageEventContent, NoticeMessageEventContent, RelatesTo}, 17 }, 18 events::stripped::StrippedRoomMember, 19 identifiers::{RoomId, UserId}, 20 Client, ClientConfig, EventEmitter, Session as SDKSession, SyncRoom, SyncSettings, 21 }; 22 use tracing::{info, warn, Level}; 23 use tracing_subscriber::FmtSubscriber; 24 use url::Url; 25 26 use crate::config::Config; 27 use crate::utils::{get_media_download_url, Session}; 28 29 mod config; 30 mod get_room_event; 31 mod utils; 32 33 struct CommandBot { 34 /// This clone of the `Client` will send requests to the server, 35 /// while the other keeps us in sync with the server using `sync_forever`. 36 client: Client, 37 ipfs_client: IpfsClient, 38 config: Config, 39 } 40 41 impl CommandBot { 42 pub fn new(client: Client, config: Config) -> Self { 43 let ipfs_client = IpfsClient::from_str(&config.ipfs_api).unwrap(); 44 Self { 45 client, 46 ipfs_client, 47 config, 48 } 49 } 50 51 fn get_temp_file(&self, filename: String) -> PathBuf { 52 let tmp_dir = env::temp_dir(); 53 tmp_dir.join(filename) 54 } 55 fn save_file(&self, filename: String, body: &[u8]) { 56 println!("file to download: '{}'", filename); 57 println!("will be located under: '{:?}'", &filename); 58 let filename = self.get_temp_file(filename); 59 let mut dest = { File::create(&filename).unwrap() }; 60 61 dest.write_all(body).unwrap(); 62 } 63 64 fn remove_file(&self, filename: String) { 65 let filename = self.get_temp_file(filename); 66 fs::remove_file(filename).unwrap(); 67 } 68 69 async fn send_link( 70 &self, 71 room_id: &RoomId, 72 filename: String, 73 hash: String, 74 related_event_original: Option<RelatesTo>, 75 ) { 76 let content = MessageEventContent::Notice(NoticeMessageEventContent { 77 body: format!( 78 "{}/ipfs/{}?filename={}", 79 self.config.ipfs_gateway, hash, filename 80 ), 81 format: None, 82 formatted_body: None, 83 relates_to: related_event_original, 84 }); 85 86 self.client 87 // send our message to the room we found the "!party" command in 88 // the last parameter is an optional Uuid which we don't care about. 89 .room_send(room_id, content, None) 90 .await 91 .unwrap(); 92 } 93 94 async fn handle_media(&self, mxc_url: String, raw_filename: String) -> String { 95 let download_url = get_media_download_url(mxc_url); 96 97 let response = reqwest::get(&download_url).await.unwrap(); 98 99 let content = response.bytes().await.unwrap(); 100 101 self.save_file(raw_filename.clone(), content.bytes()); 102 103 let filename = self.get_temp_file(raw_filename.clone()); 104 let ipfs_resp = self.ipfs_client.add_path(&filename).await.unwrap(); 105 self.remove_file(raw_filename); 106 107 let hash = ipfs_resp.first().unwrap().hash.clone(); 108 self.ipfs_client.pin_add(&hash, true).await.unwrap(); 109 110 hash 111 } 112 } 113 114 #[matrix_sdk_common_macros::async_trait] 115 impl EventEmitter for CommandBot { 116 async fn on_stripped_state_member( 117 &self, 118 room: SyncRoom, 119 _: &StrippedRoomMember, 120 _: Option<MemberEventContent>, 121 ) { 122 if let SyncRoom::Invited(room) = room { 123 let room_id = room.read().await.room_id.clone(); 124 self.client.join_room_by_id(&room_id).await.unwrap(); 125 } 126 } 127 async fn on_room_message(&self, room: SyncRoom, event: &MessageEvent) { 128 if let SyncRoom::Joined(room) = room { 129 if let MessageEventContent::Text(text_event) = event.clone().content { 130 let test = serde_json::from_str::<ImageInfo>(r#"{"mimetype":"image/jpeg", "w":4998, "h":3333,"size":6467842}"#); 131 println!("{:?}", test); 132 let msg_body = text_event.body.clone(); 133 134 // TODO fix e2ee relates_to with something like https://github.com/matrix-org/matrix-rust-sdk/blob/master/matrix_sdk_base/src/client.rs#L93 inside of receive_joined_timeline_event 135 136 if msg_body.contains("!ipfs") && text_event.relates_to.is_some() { 137 println!("new !ipfs message"); 138 let related_event_original = text_event.relates_to.clone(); 139 140 // we clone here to hold the lock for as little time as possible. 141 let room_id = room.read().await.room_id.clone(); 142 let mut related_events: Vec<MessageEvent> = room 143 .read() 144 .await 145 .messages 146 .iter() 147 .filter(|x| { 148 (**x).event_id 149 == related_event_original 150 .as_ref() 151 .unwrap() 152 .in_reply_to 153 .event_id 154 }) 155 .map(|x| (**x).clone()) 156 .collect(); 157 if related_events.is_empty() { 158 // Fetch missing event 159 let resp = self 160 .client 161 .send(get_room_event::Request { 162 room_id: room_id.clone(), 163 event_id: related_event_original 164 .clone() 165 .unwrap() 166 .in_reply_to 167 .event_id, 168 }) 169 .await; 170 171 println!("event: {:?}", resp); 172 173 match resp { 174 Ok(mut resp) => { 175 println!("{:?}", resp.event.deserialize()); 176 let (event, _updated) = self 177 .client 178 .base_client 179 .receive_joined_timeline_event(&room_id, &mut resp.event) 180 .await 181 .unwrap(); 182 match event { 183 Some(event) => { 184 if let Ok(RoomEvent::RoomMessage(msg_event)) = 185 event.deserialize() 186 { 187 related_events.push(msg_event); 188 } 189 } 190 None => { 191 if let Ok(RoomEvent::RoomMessage(msg_event)) = 192 resp.event.deserialize() 193 { 194 related_events.push(msg_event); 195 } 196 } 197 } 198 } 199 Err(e) => { 200 println!("error: {:?}", e); 201 } 202 } 203 } 204 if !related_events.is_empty() { 205 let related_event = related_events.first(); 206 207 if let Some(related_event) = related_event { 208 // TODO handle media content 209 info!("got related_event"); 210 211 match related_event.clone().content { 212 MessageEventContent::Image(image_event) => { 213 info!("handling image event"); 214 215 // Saving image 216 let filename = image_event.body.clone(); 217 let hash = match image_event.url { 218 None => { 219 self.handle_media( 220 image_event.file.unwrap().url, 221 filename.clone(), 222 ) 223 .await 224 } 225 Some(url) => self.handle_media(url, filename.clone()).await, 226 }; 227 228 // Sending link 229 self.send_link( 230 &room_id, 231 filename.clone(), 232 hash, 233 related_event_original.clone(), 234 ) 235 .await; 236 237 info!("image event message sent"); 238 } 239 MessageEventContent::Video(video_event) => { 240 info!("handling video event"); 241 242 // Saving video 243 let filename = video_event.body.clone(); 244 let hash = match video_event.url { 245 None => { 246 self.handle_media( 247 video_event.file.unwrap().url, 248 filename.clone(), 249 ) 250 .await 251 } 252 Some(url) => self.handle_media(url, filename.clone()).await, 253 }; 254 255 // Sending link 256 self.send_link( 257 &room_id, 258 filename.clone(), 259 hash, 260 related_event_original.clone(), 261 ) 262 .await; 263 264 info!("video event message sent"); 265 } 266 MessageEventContent::File(file_event) => { 267 info!("handling file event"); 268 269 // Saving file 270 let filename = file_event.body.clone(); 271 let hash = match file_event.url { 272 None => { 273 self.handle_media( 274 file_event.file.unwrap().url, 275 filename.clone(), 276 ) 277 .await 278 } 279 Some(url) => self.handle_media(url, filename.clone()).await, 280 }; 281 282 // Sending link 283 self.send_link( 284 &room_id, 285 filename.clone(), 286 hash, 287 related_event_original.clone(), 288 ) 289 .await; 290 291 info!("file event message sent"); 292 } 293 MessageEventContent::Audio(audio_event) => { 294 info!("handling audio event"); 295 296 // Saving audio 297 let filename = audio_event.body.clone(); 298 let hash = match audio_event.url { 299 None => { 300 self.handle_media( 301 audio_event.file.unwrap().url, 302 filename.clone(), 303 ) 304 .await 305 } 306 Some(url) => self.handle_media(url, filename.clone()).await, 307 }; 308 309 // Sending link 310 self.send_link( 311 &room_id, 312 filename.clone(), 313 hash, 314 related_event_original.clone(), 315 ) 316 .await; 317 318 info!("audio event message sent"); 319 } 320 _ => { 321 info!("sending fallback response"); 322 323 let content = MessageEventContent::Notice(NoticeMessageEventContent { 324 body: "Only Image, Video, File and Audio events are supported!".to_string(), 325 format: None, 326 formatted_body: None, 327 relates_to: related_event_original.clone(), 328 }); 329 330 self.client 331 // send our message to the room we found the "!party" command in 332 // the last parameter is an optional Uuid which we don't care about. 333 .room_send(&room_id, content, None) 334 .await 335 .unwrap(); 336 337 info!("fallback response message sent"); 338 } 339 } 340 } 341 } else { 342 let content = MessageEventContent::Notice(NoticeMessageEventContent { 343 body: "Unable to find related event!".to_string(), 344 format: None, 345 formatted_body: None, 346 relates_to: related_event_original.clone(), 347 }); 348 349 self.client 350 // send our message to the room we found the "!party" command in 351 // the last parameter is an optional Uuid which we don't care about. 352 .room_send(&room_id, content, None) 353 .await 354 .unwrap(); 355 356 warn!("Unable to find related_event"); 357 } 358 } 359 } 360 } 361 } 362 } 363 364 async fn login_and_sync( 365 homeserver_url: String, 366 username: String, 367 password: String, 368 ) -> Result<(), matrix_sdk::Error> { 369 // the location for `JsonStore` to save files to 370 let mut home = dirs::home_dir().expect("no home directory found"); 371 home.push("ipfs_bot"); 372 fs::create_dir_all(&home).unwrap(); 373 374 let client_config = ClientConfig::new() 375 .store_path(&home) 376 .passphrase(password.clone()); 377 378 let homeserver_url = Url::parse(&homeserver_url).expect("Couldn't parse the homeserver URL"); 379 // create a new Client with the given homeserver url and config 380 let mut client = Client::new_with_config(homeserver_url, client_config).unwrap(); 381 382 let mut session = home.clone(); 383 session.push("session.json"); 384 if session.exists() { 385 let f = OpenOptions::new().read(true).open(&session).unwrap(); 386 let json: Session = serde_json::from_reader(f).expect("file should be proper JSON"); 387 let session = SDKSession { 388 access_token: json.access_token, 389 user_id: UserId::try_from(json.user_id).unwrap(), 390 device_id: json.device_id, 391 }; 392 client.restore_login(session).await.unwrap(); 393 } else { 394 let f = OpenOptions::new() 395 .read(true) 396 .write(true) 397 .create(true) 398 .open(&session) 399 .unwrap(); 400 401 let login_response = client 402 .login( 403 username.clone(), 404 password, 405 None, 406 Some("ipfs bot".to_string()), 407 ) 408 .await?; 409 410 let session = Session { 411 access_token: login_response.access_token, 412 user_id: login_response.user_id.to_string(), 413 device_id: login_response.device_id, 414 }; 415 416 serde_json::to_writer(&f, &session).unwrap(); 417 } 418 419 println!("logged in as {}", username); 420 421 // add our CommandBot to be notified of incoming messages, we do this after the initial 422 // sync to avoid responding to messages before the bot was running. 423 client 424 .add_event_emitter(Box::new(CommandBot::new(client.clone(), Config::load()))) 425 .await; 426 427 // since we called sync before we `sync_forever` we must pass that sync token to 428 // `sync_forever` 429 let settings = SyncSettings::default(); 430 // this keeps state from the server streaming in to CommandBot via the EventEmitter trait 431 client.sync_forever(settings, |_| async {}).await; 432 433 Ok(()) 434 } 435 436 #[tokio::main] 437 async fn main() -> Result<(), matrix_sdk::Error> { 438 let subscriber = FmtSubscriber::builder() 439 // all spans/events with a level higher than TRACE (e.g, debug, info, warn, etc.) 440 // will be written to stdout. 441 .with_max_level(Level::DEBUG) 442 // completes the builder. 443 .finish(); 444 445 tracing::subscriber::set_global_default(subscriber).expect("setting default subscriber failed"); 446 447 let (homeserver_url, username, password) = 448 match (env::args().nth(1), env::args().nth(2), env::args().nth(3)) { 449 (Some(a), Some(b), Some(c)) => (a, b, c), 450 _ => { 451 eprintln!( 452 "Usage: {} <homeserver_url> <username> <password>", 453 env::args().next().unwrap() 454 ); 455 exit(1) 456 } 457 }; 458 459 login_and_sync(homeserver_url, username, password).await?; 460 Ok(()) 461 }