use tracing::{debug, error, info, warn};
use matrix_sdk::{
config::SyncSettings,
deserialized_responses::TimelineEventKind,
event_handler::Ctx,
room::MessagesOptions,
room::Room,
ruma::{
api::client::{
filter::{FilterDefinition, RoomEventFilter, RoomFilter},
sync::sync_events::v3::Filter,
},
events::room::encrypted::{
OriginalSyncRoomEncryptedEvent,
SyncRoomEncryptedEvent,
},
events::room::message::{
AudioMessageEventContent,
FileMessageEventContent,
ImageMessageEventContent,
MessageType,
NoticeMessageEventContent,
RedactedSyncRoomMessageEvent,
RoomMessageEventContent,
SyncRoomMessageEvent,
TextMessageEventContent,
VideoMessageEventContent,
},
events::room::redaction::{
OriginalSyncRoomRedactionEvent, RedactedSyncRoomRedactionEvent, SyncRoomRedactionEvent,
},
events::{
AnyMessageLikeEvent,
AnyTimelineEvent,
MessageLikeEvent,
OriginalSyncMessageLikeEvent,
SyncMessageLikeEvent,
},
OwnedRoomId,
OwnedUserId,
RoomId,
UInt,
},
Client,
};
use crate::{Error, Output};
fn handle_originalsyncmessagelikeevent(
ev: &OriginalSyncMessageLikeEvent<RoomMessageEventContent>,
room_id: &OwnedRoomId,
context: &Ctx<EvHandlerContext>,
) {
debug!(
"New message: {:?} from sender {:?}, room {:?}, event_id {:?}",
ev.content,
ev.sender,
room_id, ev.event_id,
);
if context.whoami != ev.sender || context.listen_self {
match ev.content.msgtype.to_owned() {
MessageType::Text(textmessageeventcontent) => {
let TextMessageEventContent {
body,
formatted,
..
} = textmessageeventcontent;
println!(
"Message: type Text: body {:?}, room {:?}, sender {:?}, event id {:?}, formatted {:?}, ",
body, room_id, ev.sender, ev.event_id, formatted,
);
}
MessageType::File(filemessageeventcontent) => {
let FileMessageEventContent {
body,
filename,
source,
info,
..
} = filemessageeventcontent;
println!(
"Message: type File: body {:?}, room {:?}, sender {:?}, event id {:?}, filename {:?}, source {:?}, info {:?}",
body, room_id, ev.sender, ev.event_id, filename, source, info,
);
}
MessageType::Image(imagemessageeventcontent) => {
let ImageMessageEventContent {
body, source, info, ..
} = imagemessageeventcontent;
println!(
"Message: type Image: body {:?}, room {:?}, sender {:?}, event id {:?}, source {:?}, info {:?}",
body, room_id, ev.sender, ev.event_id, source, info,
);
}
MessageType::Audio(audiomessageeventcontent) => {
let AudioMessageEventContent {
body, source, info, ..
} = audiomessageeventcontent;
println!(
"Message: type Audio: body {:?}, room {:?}, sender {:?}, event id {:?}, source {:?}, info {:?}",
body, room_id, ev.sender, ev.event_id, source, info,
);
}
MessageType::Video(videomessageeventcontent) => {
let VideoMessageEventContent {
body, source, info, ..
} = videomessageeventcontent;
println!(
"Message: type Video: body {:?}, room {:?}, sender {:?}, event id {:?}, source {:?}, info {:?}",
body, room_id, ev.sender, ev.event_id, source, info,
);
}
MessageType::Notice(noticemessageeventcontent) => {
let NoticeMessageEventContent {
body, formatted, ..
} = noticemessageeventcontent;
println!(
"Message: type Notice: body {:?}, room {:?}, sender {:?}, event id {:?}, formatted {:?}, ",
body, room_id, ev.sender, ev.event_id, formatted,
);
}
_ => {
debug!("Not handling this event: {:?}", ev);
warn!(
"Not handling this message type. Not implemented yet. {:?}",
ev
);
}
}
} else {
debug!("Skipping message from itself because --listen-self is not set.");
}
}
async fn handle_redactedsyncroommessageevent(
ev: RedactedSyncRoomMessageEvent,
room: Room,
_client: Client,
context: Ctx<EvHandlerContext>,
) {
debug!(
"Received a message for RedactedSyncRoomMessageEvent. {:?}",
ev
);
if context.whoami == ev.sender && !context.listen_self {
debug!("Skipping message from itself because --listen-self is not set.");
return;
}
if !context.output.is_text() {
let j = match serde_json::to_string(&ev.content) {
Ok(jsonstr) => {
let mut s = jsonstr;
s.insert_str(s.len() - 1, ",\"event_id\":\"\"");
s.insert_str(s.len() - 2, ev.event_id.as_str());
s.insert_str(s.len() - 1, ",\"sender\":\"\"");
s.insert_str(s.len() - 2, ev.sender.as_str());
s.insert_str(s.len() - 1, ",\"origin_server_ts\":\"\"");
s.insert_str(s.len() - 2, &ev.origin_server_ts.0.to_string());
s.insert_str(s.len() - 1, ",\"room_id\":\"\"");
s.insert_str(s.len() - 2, room.room_id().as_str());
s
}
Err(e) => e.to_string(),
};
println!("{}", j);
return;
}
debug!(
"Received a message for RedactedSyncRoomMessageEvent. Not implemented yet for text format, try --output json. {:?}",
ev
);
}
fn handle_originalsyncroomredactionevent(ev: OriginalSyncRoomRedactionEvent, room: Room) {
debug!(
"Received a message for OriginalSyncRoomRedactionEvent. {:?}",
ev
);
let j = match serde_json::to_string(&ev.content) {
Ok(jsonstr) => {
let mut s = jsonstr;
s.insert_str(s.len() - 1, ",\"event_id\":\"\"");
s.insert_str(s.len() - 2, ev.event_id.as_str());
s.insert_str(s.len() - 1, ",\"sender\":\"\"");
s.insert_str(s.len() - 2, ev.sender.as_str());
s.insert_str(s.len() - 1, ",\"origin_server_ts\":\"\"");
s.insert_str(s.len() - 2, &ev.origin_server_ts.0.to_string());
s.insert_str(s.len() - 1, ",\"room_id\":\"\"");
s.insert_str(s.len() - 2, room.room_id().as_str());
s
}
Err(e) => e.to_string(),
};
println!("{}", j);
}
fn handle_redactedsyncroomredactionevent(ev: RedactedSyncRoomRedactionEvent, room: Room) {
debug!(
"Received a message for RedactedSyncRoomRedactionEvent. {:?}",
ev
);
let j = match serde_json::to_string(&ev.content) {
Ok(jsonstr) => {
let mut s = jsonstr;
s.insert_str(s.len() - 1, ",\"event_id\":\"\"");
s.insert_str(s.len() - 2, ev.event_id.as_str());
s.insert_str(s.len() - 1, ",\"sender\":\"\"");
s.insert_str(s.len() - 2, ev.sender.as_str());
s.insert_str(s.len() - 1, ",\"origin_server_ts\":\"\"");
s.insert_str(s.len() - 2, &ev.origin_server_ts.0.to_string());
s.insert_str(s.len() - 1, ",\"room_id\":\"\"");
s.insert_str(s.len() - 2, room.room_id().as_str());
s
}
Err(e) => e.to_string(),
};
println!("{}", j);
}
async fn handle_syncroomredactedevent(
ev: SyncRoomRedactionEvent,
room: Room,
_client: Client,
context: Ctx<EvHandlerContext>,
) {
debug!("Received a message for SyncRoomRedactionEvent. {:?}", ev);
if context.whoami == ev.sender() && !context.listen_self {
debug!("Skipping message from itself because --listen-self is not set.");
return;
}
if !context.output.is_text() {
match ev {
SyncRoomRedactionEvent::Original(evi) => {
handle_originalsyncroomredactionevent(evi, room)
}
SyncRoomRedactionEvent::Redacted(evi) => {
handle_redactedsyncroomredactionevent(evi, room)
}
}
return;
}
debug!(
"Received a message for SyncRoomRedactionEvent. Not implemented yet for text format, try --output json. {:?}",
ev
);
}
async fn handle_syncroomencryptedevent(
ev: SyncRoomEncryptedEvent,
room: Room,
_client: Client,
context: Ctx<EvHandlerContext>,
) {
debug!("Received a SyncRoomEncryptedEvent message {:?}", ev);
if context.whoami == ev.sender() && !context.listen_self {
debug!("Skipping message from itself because --listen-self is not set.");
return;
}
match ev {
SyncMessageLikeEvent::Original(originalmessagelikeevent) => {
debug!(
"New message: {:?} from sender {:?}, room {:?}, event_id {:?}",
originalmessagelikeevent.content,
originalmessagelikeevent.sender,
room.room_id(), originalmessagelikeevent.event_id,
);
}
_ => {
debug!(
"Received a message for RedactedSyncMessageLikeEvent. Not implemented yet. {:?}",
ev
);
}
}
warn!("Decryption attempt not implemented yet.");
}
async fn handle_originalsyncroomencryptedevent(
ev: OriginalSyncRoomEncryptedEvent,
room: Room,
_client: Client,
context: Ctx<EvHandlerContext>,
) {
debug!("Received a OriginalSyncRoomEncryptedEvent message {:?}", ev);
if context.whoami == ev.sender && !context.listen_self {
debug!("Skipping message from itself because --listen-self is not set.");
return;
}
debug!(
"New message: {:?} from sender {:?}, room {:?}, event_id {:?}",
ev.content,
ev.sender,
room.room_id(), ev.event_id,
);
warn!("Decryption attempt not implemented yet.");
}
async fn handle_syncroommessageevent(
ev: SyncRoomMessageEvent,
room: Room,
_client: Client,
context: Ctx<EvHandlerContext>,
) {
debug!("Received a message for event SyncRoomMessageEvent {:?}", ev);
if context.whoami == ev.sender() && !context.listen_self {
debug!("Skipping message from itself because --listen-self is not set.");
return;
}
match ev {
SyncMessageLikeEvent::Original(originalmessagelikeevent) => {
if !context.output.is_text() {
let j = match serde_json::to_string(&originalmessagelikeevent.content) {
Ok(jsonstr) => {
let mut s = jsonstr;
s.insert_str(s.len() - 1, ",\"event_id\":\"\"");
s.insert_str(s.len() - 2, originalmessagelikeevent.event_id.as_str());
s.insert_str(s.len() - 1, ",\"sender\":\"\"");
s.insert_str(s.len() - 2, originalmessagelikeevent.sender.as_str());
s.insert_str(s.len() - 1, ",\"origin_server_ts\":\"\"");
s.insert_str(
s.len() - 2,
&originalmessagelikeevent.origin_server_ts.0.to_string(),
);
s.insert_str(s.len() - 1, ",\"room_id\":\"\"");
s.insert_str(s.len() - 2, room.room_id().as_str());
s
}
Err(e) => e.to_string(),
};
println!("{}", j);
return;
}
handle_originalsyncmessagelikeevent(
&originalmessagelikeevent,
&RoomId::parse(room.room_id()).unwrap(),
&context,
);
}
_ => {
debug!(
"Received a message for RedactedSyncMessageLikeEvent. Not implemented yet. {:?}",
ev
);
}
};
}
#[derive(Clone, Debug)]
struct EvHandlerContext {
whoami: OwnedUserId,
listen_self: bool,
output: Output,
}
pub(crate) async fn listen_once(
client: &Client,
listen_self: bool, whoami: OwnedUserId,
output: Output,
) -> Result<(), Error> {
info!(
"mclient::listen_once(): listen_self {}, room {}",
listen_self, "all"
);
let context = EvHandlerContext {
whoami,
listen_self,
output,
};
client.add_event_handler_context(context.clone());
client.add_event_handler(|ev: SyncRoomMessageEvent, room: Room,
client: Client, context: Ctx<EvHandlerContext>| async move {
tokio::spawn(handle_syncroommessageevent(ev, room, client, context));
});
client.add_event_handler(
|ev: RedactedSyncRoomMessageEvent,
room: Room,
client: Client,
context: Ctx<EvHandlerContext>| async move {
tokio::spawn(handle_redactedsyncroommessageevent(
ev, room, client, context,
));
},
);
client.add_event_handler(
|ev: SyncRoomRedactionEvent,
room: Room,
client: Client,
context: Ctx<EvHandlerContext>| async move {
tokio::spawn(handle_syncroomredactedevent(ev, room, client, context));
},
);
client.add_event_handler(
|ev: OriginalSyncRoomEncryptedEvent,
room: Room,
client: Client,
context: Ctx<EvHandlerContext>| async move {
tokio::spawn(handle_originalsyncroomencryptedevent(
ev, room, client, context,
));
},
);
client.add_event_handler(
|ev: SyncRoomEncryptedEvent,
room: Room,
client: Client,
context: Ctx<EvHandlerContext>| async move {
tokio::spawn(handle_syncroomencryptedevent(ev, room, client, context));
},
);
info!("Ready and getting messages from server...");
let settings = SyncSettings::default();
client.sync_once(settings).await?;
Ok(())
}
pub(crate) async fn listen_forever(
client: &Client,
listen_self: bool, whoami: OwnedUserId,
output: Output,
) -> Result<(), Error> {
info!(
"mclient::listen_forever(): listen_self {}, room {}",
listen_self, "all"
);
let context = EvHandlerContext {
whoami,
listen_self,
output,
};
client.add_event_handler_context(context.clone());
client.add_event_handler(
|ev: SyncRoomMessageEvent, room: Room, client: Client, context: Ctx<EvHandlerContext>| async move {
tokio::spawn(handle_syncroommessageevent(ev, room, client, context));
});
client.add_event_handler(
|ev: SyncRoomEncryptedEvent, room: Room, client: Client, context: Ctx<EvHandlerContext>| async move {
tokio::spawn(handle_syncroomencryptedevent(ev, room, client, context));
});
client.add_event_handler(
|ev: OriginalSyncRoomEncryptedEvent,
room: Room,
client: Client,
context: Ctx<EvHandlerContext>| async move {
tokio::spawn(handle_originalsyncroomencryptedevent(
ev, room, client, context,
));
},
);
client.add_event_handler(
|ev: RedactedSyncRoomMessageEvent,
room: Room,
client: Client,
context: Ctx<EvHandlerContext>| async move {
tokio::spawn(handle_redactedsyncroommessageevent(
ev, room, client, context,
));
},
);
client.add_event_handler(|ev: SyncRoomRedactionEvent,
room: Room, client: Client, context: Ctx<EvHandlerContext>| async move {
tokio::spawn(handle_syncroomredactedevent(ev, room, client, context));
});
info!("Ready and waiting for messages ...");
info!("Once done listening, kill the process manually with Control-C.");
let settings = SyncSettings::default();
match client.sync(settings).await {
Ok(()) => Ok(()),
Err(e) => {
error!("Event loop reported: {:?}", e);
Ok(())
}
}
}
#[allow(dead_code)]
fn print_type_of<T>(_: &T) {
println!("{}", std::any::type_name::<T>())
}
pub(crate) async fn listen_tail(
client: &Client,
roomnames: &Vec<String>, number: u64, listen_self: bool, whoami: OwnedUserId,
output: Output,
) -> Result<(), Error> {
info!(
"mclient::listen_tail(): listen_self {}, roomnames {:?}",
listen_self, roomnames
);
if roomnames.is_empty() {
return Err(Error::MissingRoom);
}
info!("Ready and getting messages from server ...");
let mut roomids: Vec<OwnedRoomId> = Vec::new();
for roomname in roomnames {
roomids.push(match RoomId::parse(roomname.clone()) {
Ok(id) => id,
Err(ref e) => {
error!(
"Error: invalid room id {:?}. Error reported is {:?}.",
roomname, e
);
continue;
}
});
}
let ownedroomidvecoption: Option<Vec<OwnedRoomId>> = Some(roomids.clone());
let mut filter = FilterDefinition::default();
let mut roomfilter = RoomFilter::empty();
roomfilter.rooms = ownedroomidvecoption;
filter.room = roomfilter;
let mut roomtimeline = RoomEventFilter::empty();
roomtimeline.limit = UInt::new(number);
filter.room.timeline = roomtimeline;
let context = EvHandlerContext {
whoami: whoami.clone(),
listen_self,
output,
};
let ctx = Ctx(context);
let mut err_count = 0u32;
for roomid in roomids.iter() {
let mut options = MessagesOptions::backward(); options.limit = UInt::new(number).unwrap();
let jroom = client.get_room(roomid.clone().as_ref()).unwrap();
let msgs = jroom.messages(options).await;
let chunk = msgs.unwrap().chunk;
for index in 0..chunk.len() {
debug!(
"processing message {:?} out of {:?}",
index + 1,
chunk.len()
);
let anytimelineevent = &chunk[chunk.len() - 1 - index];
let rawevent = match &anytimelineevent.kind {
TimelineEventKind::Decrypted(decrypted) => &decrypted.event,
TimelineEventKind::PlainText { event } => {
event.cast_ref_unchecked::<AnyTimelineEvent>()
}
_ => {
let utdraw = anytimelineevent.raw();
let event_id = anytimelineevent
.event_id()
.map(|e| e.to_string())
.unwrap_or_default();
let sender = utdraw
.deserialize()
.ok()
.map(|e| e.sender().to_string())
.unwrap_or_default();
if !output.is_text() {
println!("{}", utdraw.json());
} else {
println!(
"Message: type Encrypted: room {:?}, sender {:?}, event_id {:?}, message could not be decrypted",
roomid, sender, event_id,
);
}
err_count += 1;
continue;
}
};
match rawevent.deserialize().unwrap() {
AnyTimelineEvent::MessageLike(anymessagelikeevent) => {
debug!("value: {:?}", anymessagelikeevent);
match anymessagelikeevent {
AnyMessageLikeEvent::RoomMessage(messagelikeevent) => {
debug!("value: {:?}", messagelikeevent);
match messagelikeevent {
MessageLikeEvent::Original(originalmessagelikeevent) => {
let room_id = originalmessagelikeevent.room_id.clone();
let originalsyncmessagelikeevent =
OriginalSyncMessageLikeEvent::from(
originalmessagelikeevent,
);
handle_originalsyncmessagelikeevent(
&originalsyncmessagelikeevent,
&room_id,
&ctx,
);
}
_ => {
warn!("RoomMessage type is not handled. Not implemented yet.");
err_count += 1;
}
}
}
AnyMessageLikeEvent::RoomEncrypted(messagelikeevent) => {
warn!(
"Event of type RoomEncrypted received: {:?}",
messagelikeevent
);
match messagelikeevent {
MessageLikeEvent::Original(originalmessagelikeevent) => {
debug!(
"New message: {:?} from sender {:?}, room {:?}, event_id {:?}",
originalmessagelikeevent.content,
originalmessagelikeevent.sender,
originalmessagelikeevent.room_id,
originalmessagelikeevent.event_id,
);
if whoami != originalmessagelikeevent.sender || listen_self {
println!(
"Message: type Encrypted: body {:?}, room {:?}, sender {:?}, event_id {:?}, message could not be decrypted",
originalmessagelikeevent.content, originalmessagelikeevent.room_id, originalmessagelikeevent.sender, originalmessagelikeevent.event_id,
);
} else {
debug!("Skipping message from itself because --listen-self is not set.");
}
}
_ => {
warn!("RoomMessage type is not handled. Not implemented yet.");
err_count += 1;
}
}
}
AnyMessageLikeEvent::RoomRedaction(messagelikeevent) => {
warn!("Event of type RoomRedaction received. Not implemented yet. value: {:?}", messagelikeevent);
err_count += 1;
}
_ => {
warn!("MessageLike type is not handle. Not implemented yet.");
err_count += 1;
}
}
}
_ => debug!("State event, not interested in that."),
}
}
}
if err_count != 0 {
Err(Error::NotImplementedYet)
} else {
Ok(())
}
}
pub(crate) async fn listen_all(
client: &Client,
roomnames: &Vec<String>, listen_self: bool, whoami: OwnedUserId,
output: Output,
) -> Result<(), Error> {
if roomnames.is_empty() {
return Err(Error::MissingRoom);
}
info!(
"mclient::listen_all(): listen_self {}, roomnames {:?}",
listen_self, roomnames
);
let context = EvHandlerContext {
whoami,
listen_self,
output,
};
client.add_event_handler_context(context.clone());
client.add_event_handler(|ev: SyncRoomMessageEvent, room: Room, client: Client, context: Ctx<EvHandlerContext>| async move {
tokio::spawn(handle_syncroommessageevent(ev, room, client, context));
});
client.add_event_handler(
|ev: RedactedSyncRoomMessageEvent,
room: Room,
client: Client,
context: Ctx<EvHandlerContext>| async move {
tokio::spawn(handle_redactedsyncroommessageevent(
ev, room, client, context,
));
},
);
client.add_event_handler(|ev: SyncRoomRedactionEvent,
room: Room, client: Client, context: Ctx<EvHandlerContext>| async move {
tokio::spawn(handle_syncroomredactedevent(ev, room, client, context));
});
info!("Ready and waiting for messages ...");
let mut roomids: Vec<OwnedRoomId> = Vec::new();
for roomname in roomnames {
roomids.push(RoomId::parse(roomname.clone()).unwrap());
}
let ownedroomidvecoption: Option<Vec<OwnedRoomId>> = Some(roomids);
let mut filter = FilterDefinition::default();
let mut roomfilter = RoomFilter::empty();
roomfilter.rooms = ownedroomidvecoption;
filter.room = roomfilter;
let mut err_count = 0u32;
let filterclone = filter.clone();
let sync_settings = SyncSettings::default().filter(Filter::FilterDefinition(filterclone));
match client.sync_once(sync_settings).await {
Ok(response) => debug!("listen_all successful {:?}", response),
Err(ref e) => {
err_count += 1;
error!("listen_all returned error {:?}", e);
}
}
if err_count != 0 {
Err(Error::ListenFailed)
} else {
Ok(())
}
}