matrix.rs (12954B)
1 use std::path::Path; 2 3 use crate::indradb_utils::{BulkInserter, MessagesMap, UUIDEventMapType, UUIDRoomMapType}; 4 use color_eyre::Result; 5 use futures::StreamExt; 6 use matrix_sdk::{ 7 config::SyncSettings, 8 ruma::{ 9 events::{ 10 room::message::MessageType, AnySyncMessageLikeEvent, AnySyncTimelineEvent, 11 SyncMessageLikeEvent, 12 }, 13 OwnedUserId, 14 }, 15 Client, Session, 16 }; 17 use tracing::{error, info}; 18 19 struct Identifiers { 20 room_type: utils::indradb::Identifier, 21 room_id_type: utils::indradb::Identifier, 22 room_name_type: utils::indradb::Identifier, 23 room_topic_type: utils::indradb::Identifier, 24 text_message_event_type: utils::indradb::Identifier, 25 notice_message_event_type: utils::indradb::Identifier, 26 event_id_type: utils::indradb::Identifier, 27 event_in_room_type: utils::indradb::Identifier, 28 } 29 30 pub struct IndexerBot { 31 client: Client, 32 indexer_client: utils::indradb_proto::Client, 33 message_map: MessagesMap, 34 identifiers: Identifiers, 35 } 36 37 impl IndexerBot { 38 pub fn access_token(&self) -> Option<String> { 39 self.client.access_token() 40 } 41 42 pub fn device_id(&self) -> Option<String> { 43 self.client.device_id().map(ToString::to_string) 44 } 45 async fn get_client(homeserver_url: String) -> Result<Client> { 46 let mut client_builder = Client::builder().homeserver_url(homeserver_url); 47 client_builder = client_builder.sled_store(Path::new("./matrix_data"), None)?; 48 49 Ok(client_builder.build().await?) 50 } 51 52 async fn get_indexer_client( 53 endpoint: String, 54 ) -> Result<(utils::indradb_proto::Client, Identifiers)> { 55 info!("Trying to connect to indradb"); 56 let mut indexer_client = utils::get_client_retrying(endpoint).await?; 57 indexer_client.ping().await?; 58 let room_type = utils::indradb::Identifier::new("matrix_room")?; 59 let room_id_type = utils::indradb::Identifier::new("room_id")?; 60 let room_name_type = utils::indradb::Identifier::new("room_name")?; 61 let room_topic_type = utils::indradb::Identifier::new("room_topic")?; 62 let text_message_event_type = utils::indradb::Identifier::new("text_message_event")?; 63 let notice_message_event_type = utils::indradb::Identifier::new("notice_message_event")?; 64 let event_id_type = utils::indradb::Identifier::new("event_id")?; 65 let event_in_room_type = utils::indradb::Identifier::new("event_in_room")?; 66 indexer_client.index_property(room_id_type).await?; 67 indexer_client.index_property(room_name_type).await?; 68 indexer_client.index_property(room_topic_type).await?; 69 indexer_client.index_property(event_id_type).await?; 70 indexer_client 71 .index_property(utils::indradb::Identifier::new("text_message_body")?) 72 .await?; 73 indexer_client 74 .index_property(utils::indradb::Identifier::new( 75 "text_message_formatted_body", 76 )?) 77 .await?; 78 info!("Connected to indradb"); 79 80 Ok(( 81 indexer_client, 82 Identifiers { 83 room_type, 84 room_id_type, 85 room_name_type, 86 room_topic_type, 87 text_message_event_type, 88 notice_message_event_type, 89 event_id_type, 90 event_in_room_type, 91 }, 92 )) 93 } 94 95 pub async fn new( 96 homeserver_url: String, 97 user_id: String, 98 password: String, 99 indra_endpoint: String, 100 ) -> Result<Self> { 101 let client = IndexerBot::get_client(homeserver_url).await?; 102 client 103 .login_username(&user_id, &password) 104 .initial_device_display_name("Knowledge Indexer bot") 105 .send() 106 .await?; 107 108 let (indexer_client, identifiers) = IndexerBot::get_indexer_client(indra_endpoint).await?; 109 110 let client_clone = client.clone(); 111 tokio::spawn(async move { 112 let settings = SyncSettings::default(); 113 client_clone 114 .sync(settings) 115 .await 116 .expect("Failed to start matrix sync"); 117 }); 118 119 Ok(IndexerBot { 120 client, 121 indexer_client, 122 message_map: MessagesMap::default(), 123 identifiers, 124 }) 125 } 126 127 pub async fn relogin( 128 homeserver_url: String, 129 user_id: String, 130 access_token: String, 131 device_id: String, 132 indra_endpoint: String, 133 ) -> Result<Self> { 134 let client = IndexerBot::get_client(homeserver_url).await?; 135 client 136 .restore_login(Session { 137 access_token, 138 device_id: device_id.into(), 139 refresh_token: None, 140 user_id: OwnedUserId::try_from(user_id)?, 141 }) 142 .await?; 143 144 let (indexer_client, identifiers) = IndexerBot::get_indexer_client(indra_endpoint).await?; 145 146 let client_clone = client.clone(); 147 tokio::spawn(async move { 148 let settings = SyncSettings::default(); 149 client_clone 150 .sync(settings) 151 .await 152 .expect("Failed to start matrix sync"); 153 }); 154 155 Ok(IndexerBot { 156 client, 157 indexer_client, 158 message_map: MessagesMap::default(), 159 identifiers, 160 }) 161 } 162 163 // FIXME:_split into multiple functions 164 #[allow(clippy::too_many_lines)] 165 pub async fn start_processing(&mut self) -> Result<()> { 166 let mut inserter = BulkInserter::new(self.indexer_client.clone()); 167 168 info!("Got bulk inserter. Starting sync"); 169 170 let mut sync_stream = Box::pin(self.client.sync_stream(SyncSettings::default()).await); 171 172 info!("Sync obtained. Starting to process sync stream"); 173 while let Some(Ok(response)) = sync_stream.next().await { 174 for (ref room_id, room) in response.rooms.join { 175 for e in &room.timeline.events { 176 let room_uuid = if let Some(room) = self.client.get_joined_room(room_id) { 177 self.message_map.insert_room( 178 room_id.clone(), 179 crate::indradb_utils::RoomProperties { 180 name: room.name(), 181 topic: room.topic(), 182 }, 183 ) 184 } else { 185 self.message_map.insert_room( 186 room_id.clone(), 187 crate::indradb_utils::RoomProperties { 188 name: None, 189 topic: None, 190 }, 191 ) 192 }; 193 194 match e.event.deserialize() { 195 Ok(AnySyncTimelineEvent::MessageLike( 196 AnySyncMessageLikeEvent::RoomMessage(event), 197 )) => { 198 if let SyncMessageLikeEvent::Original(message) = event { 199 match message.content.msgtype { 200 MessageType::Text(message_content) => { 201 self.message_map.insert_event( 202 message.event_id, 203 room_uuid, 204 self.identifiers.text_message_event_type, 205 crate::indradb_utils::EventProperties::TextMessage( 206 message_content.body, 207 message_content 208 .formatted 209 .clone() 210 .map(|x| x.format.to_string()), 211 message_content.formatted.map(|x| x.body), 212 ), 213 ); 214 } 215 MessageType::Notice(message_content) => { 216 self.message_map.insert_event( 217 message.event_id, 218 room_uuid, 219 self.identifiers.notice_message_event_type, 220 crate::indradb_utils::EventProperties::TextMessage( 221 message_content.body, 222 message_content 223 .formatted 224 .clone() 225 .map(|x| x.format.to_string()), 226 message_content.formatted.map(|x| x.body), 227 ), 228 ); 229 } 230 _ => {} 231 } 232 } 233 } 234 // TODO: index space hierachy 235 Ok( 236 AnySyncTimelineEvent::MessageLike(_) | AnySyncTimelineEvent::State(_), 237 ) => {} 238 Err(e) => { 239 error!("Error deserializing event: {}", e); 240 } 241 } 242 } 243 } 244 245 // Push to indexer after we preprocessed it 246 for UUIDRoomMapType { 247 room_id, 248 uuid, 249 room_properties, 250 } in &self.message_map.room_list 251 { 252 inserter 253 .push(utils::indradb::BulkInsertItem::Vertex( 254 utils::indradb::Vertex::with_id(*uuid, self.identifiers.room_type), 255 )) 256 .await?; 257 inserter 258 .push(utils::indradb::BulkInsertItem::VertexProperty( 259 *uuid, 260 self.identifiers.room_id_type, 261 serde_json::Value::String(room_id.to_string()).into(), 262 )) 263 .await?; 264 if let Some(room_name) = &room_properties.name { 265 inserter 266 .push(utils::indradb::BulkInsertItem::VertexProperty( 267 *uuid, 268 self.identifiers.room_name_type, 269 serde_json::Value::String(room_name.to_string()).into(), 270 )) 271 .await?; 272 } 273 if let Some(room_topic) = &room_properties.topic { 274 inserter 275 .push(utils::indradb::BulkInsertItem::VertexProperty( 276 *uuid, 277 self.identifiers.room_topic_type, 278 serde_json::Value::String(room_topic.to_string()).into(), 279 )) 280 .await?; 281 } 282 } 283 for UUIDEventMapType { 284 event_id, 285 uuid, 286 event_type, 287 event_properties, 288 } in &self.message_map.message_list 289 { 290 inserter 291 .push(utils::indradb::BulkInsertItem::Vertex( 292 utils::indradb::Vertex::with_id(*uuid, *event_type), 293 )) 294 .await?; 295 inserter 296 .push(utils::indradb::BulkInsertItem::VertexProperty( 297 *uuid, 298 self.identifiers.event_id_type, 299 serde_json::Value::String(event_id.to_string()).into(), 300 )) 301 .await?; 302 for event_property in event_properties 303 .as_vec(*uuid) 304 .expect("Unable to convert to indradb properties") 305 { 306 inserter.push(event_property).await?; 307 } 308 } 309 310 for (event_uuid, room_uuid) in &self.message_map.room_event_links { 311 inserter 312 .push(utils::indradb::BulkInsertItem::Edge( 313 utils::indradb::Edge::new( 314 *event_uuid, 315 self.identifiers.event_in_room_type, 316 *room_uuid, 317 ), 318 )) 319 .await?; 320 } 321 inserter.flush().await?; 322 } 323 Ok(()) 324 } 325 }