matrix-ipfs-bot

git clone git://archive.git.mtrnord.blog/MTRNord/matrix-ipfs-bot.git
Log | Files | Refs | LICENSE

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 }