use lru::LruCache; use std::num::NonZeroUsize; use tokio::sync::mpsc::{channel, Receiver}; use twilight_model::{channel::Message, gateway::payload::incoming::MessageUpdate, id::{marker::{MessageMarker, UserMarker}, Id}, util::Timestamp}; use crate::system::types::System; type DiscordToken = String; pub struct MessageAggregator; enum IncomingMessage { Complete {message: Message, timestamp: Timestamp, seen_by: DiscordToken}, Partial {message: MessageUpdate, timestamp: Timestamp, seen_by: DiscordToken}, } pub enum MemberEvent { Message(Message, DiscordToken), GatewayConnect(DiscordToken, Id), GatewayDisconnect(DiscordToken), GatewayError(DiscordToken), } impl MessageAggregator { pub fn start(system: &System) -> Receiver { let incoming_channel = channel::(32); let outgoing_channel = channel::(32); // Start incoming task for each member for member in system.members.iter() { let followed_user = system.followed_user.clone(); let sender = incoming_channel.0.clone(); let member = member.clone(); let outgoing_channel = outgoing_channel.0.clone(); tokio::spawn(async move { loop { let next_event = {member.shard.lock().await.next_event().await}; match next_event { Err(source) => { if source.is_fatal() { let _ = outgoing_channel.send(MemberEvent::GatewayDisconnect(member.discord_token.clone())).await; return } else { let _ = outgoing_channel.send(MemberEvent::GatewayError(member.discord_token.clone())).await; } }, Ok(twilight_gateway::Event::Ready(ready)) => { let _ = outgoing_channel.send(MemberEvent::GatewayConnect(member.discord_token.clone(), ready.user.id)).await; }, Ok(twilight_gateway::Event::MessageCreate(message)) => { if message.author.id == followed_user { let _ = sender.send(IncomingMessage::Complete { message: message.0.clone(), timestamp: message.timestamp.clone(), seen_by: member.discord_token.clone() }).await; } }, Ok(twilight_gateway::Event::MessageUpdate(message)) => { if message.author.is_some() && {message.author.as_ref().unwrap().id == followed_user} { let _ = sender.send(IncomingMessage::Partial { message: (*message).clone(), timestamp: message.edited_timestamp.unwrap_or(message.timestamp.expect("No message timestamp")), seen_by: member.discord_token.clone() }).await; } }, _ => {} } } }); } // Start deduping task let mut message_cache = LruCache::, (Message, Timestamp, DiscordToken)>::new(NonZeroUsize::new(32).unwrap()); let mut receiver = incoming_channel.1; let sender = outgoing_channel.0; let system = system.clone(); tokio::spawn(async move { loop { match receiver.recv().await { None => (), Some(IncomingMessage::Complete { message, timestamp, seen_by }) => { if let None = message_cache.get(&message.id) { message_cache.put(message.id.clone(), (message.clone(), timestamp, seen_by.clone())); let _ = sender.send(MemberEvent::Message(message, seen_by)).await; } }, Some(IncomingMessage::Partial { message, timestamp, seen_by }) => { if let Some((previous_message, previous_timestamp, _)) = message_cache.get(&message.id) { if previous_timestamp.as_micros() < timestamp.as_micros() { let mut updated_message = previous_message.clone(); updated_message.content = message.content.unwrap_or(updated_message.content); message_cache.put(message.id.clone(), (updated_message.clone(), timestamp, seen_by.clone())); let _ = sender.send(MemberEvent::Message(updated_message, seen_by)).await; } } else { let client = system.members.iter().find(|m| m.discord_token == seen_by).map(|m| m.client.clone()).expect("Could not find client"); if let Ok(updated_message) = client.lock().await.message(message.channel_id, message.id).await.unwrap().model().await.map(|r|r.clone()) { message_cache.put(message.id.clone(), (updated_message.clone(), timestamp, seen_by.clone())); let _ = sender.send(MemberEvent::Message(updated_message, seen_by)).await; }; } }, }; } }); return outgoing_channel.1 } }