diff --git a/event_bus/src/bus.rs b/event_bus/src/bus.rs index 64a788b..883ccc1 100644 --- a/event_bus/src/bus.rs +++ b/event_bus/src/bus.rs @@ -4,9 +4,10 @@ use std::sync::Arc; use parking_lot::RwLock; use std::collections::HashMap; +use std::iter; use tokio::sync::broadcast; use tokio::task::JoinHandle; -use tracing::log::kv::{Key, Value}; +// use tracing::log::kv::{Key, Value}; use tracing::{debug, trace, warn}; use uuid::Uuid; @@ -17,7 +18,7 @@ pub type AnyEvent = Arc; const DEFAULT_CAPACITY: usize = 64; #[derive(Debug, Clone, PartialEq, Eq)] -enum ScopeValue { +pub enum ScopeValue { String(String), Uuid(Uuid), } @@ -29,10 +30,36 @@ impl ScopeValue { } } } + #[derive(Debug, Clone, PartialEq, Eq)] -struct Scope { - key: String, - value: ScopeValue, +pub struct Scope { + pub key: String, + pub value: ScopeValue, +} + +impl Scope { + pub fn new(key: impl Into, value: ScopeValue) -> Self { + Self { + key: key.into(), + value, + } + } + + pub fn uuid(key: impl Into, value: Uuid) -> Self { + Self::new(key, ScopeValue::Uuid(value)) + } + + pub fn string(key: impl Into, value: impl Into) -> Self { + Self::new(key, ScopeValue::String(value.into())) + } +} +impl IntoIterator for Scope { + type Item = Scope; + type IntoIter = iter::Once; + + fn into_iter(self) -> Self::IntoIter { + iter::once(self) + } } /// The central event bus. @@ -157,14 +184,7 @@ impl EventBus { trace!(topic, "Emitting event"); let event: AnyEvent = Arc::new(event); - if let Some(tx) = self.channels.read().get(topic) { - let receiver_count = tx.receiver_count(); - let _ = tx.send(Arc::clone(&event)); - trace!( - topic, - receiver_count, "Event delivered to exact-topic channel" - ); - } + self.emit_arc(topic, event); } // todo : undocumented... diff --git a/event_bus/src/lib.rs b/event_bus/src/lib.rs index 8397568..e17a115 100644 --- a/event_bus/src/lib.rs +++ b/event_bus/src/lib.rs @@ -43,7 +43,7 @@ macro_rules! match_event { } mod bus; -pub use bus::{AnyEvent, EventBus}; +pub use bus::{AnyEvent, EventBus, Scope, ScopeValue}; #[cfg(test)] mod tests; diff --git a/src/core/permission_sync.rs b/src/core/permission_sync.rs index 401fd0c..641575a 100644 --- a/src/core/permission_sync.rs +++ b/src/core/permission_sync.rs @@ -41,4 +41,6 @@ impl PermissionSyncService { event_bus, } } + + pub fn listen(&self) {} } diff --git a/src/domain/events/channel.rs b/src/domain/events/channel.rs new file mode 100644 index 0000000..00522a0 --- /dev/null +++ b/src/domain/events/channel.rs @@ -0,0 +1,17 @@ +use crate::models::prelude::Channel; +use uuid::Uuid; + +pub struct ChannelCreated { + server_id: Uuid, + channel: Channel, +} + +pub struct ChannelUpdated { + server_id: Uuid, + channel: Channel, +} + +pub struct ChannelDeleted { + server_id: Uuid, + channel: Channel, +} diff --git a/src/domain/events/message.rs b/src/domain/events/message.rs new file mode 100644 index 0000000..19829b9 --- /dev/null +++ b/src/domain/events/message.rs @@ -0,0 +1,20 @@ +use crate::models::prelude::Message; +use uuid::Uuid; + +pub struct MessageCreated { + server_id: Option, + channel_id: Uuid, + message: Message, +} + +pub struct MessageUpdated { + server_id: Option, + channel_id: Uuid, + message: Message, +} + +pub struct MessageDeleted { + server_id: Option, + channel_id: Uuid, + message: Message, +} diff --git a/src/domain/events/mod.rs b/src/domain/events/mod.rs new file mode 100644 index 0000000..0e20bb8 --- /dev/null +++ b/src/domain/events/mod.rs @@ -0,0 +1,3 @@ +pub mod channel; +pub mod message; +pub mod server; diff --git a/src/domain/events/server.rs b/src/domain/events/server.rs new file mode 100644 index 0000000..8325137 --- /dev/null +++ b/src/domain/events/server.rs @@ -0,0 +1,13 @@ +use crate::models::prelude::Server; + +pub struct ServerCreated { + server: Server, +} + +pub struct ServerUpdated { + server: Server, +} + +pub struct ServerDeleted { + server: Server, +} diff --git a/src/domain/mod.rs b/src/domain/mod.rs new file mode 100644 index 0000000..a9970c2 --- /dev/null +++ b/src/domain/mod.rs @@ -0,0 +1 @@ +pub mod events; diff --git a/src/lib.rs b/src/lib.rs index d9f0858..264b1ca 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -11,3 +11,5 @@ pub mod udp; pub mod auth; pub mod metrics; + +pub mod domain; diff --git a/src/repositories/message.rs b/src/repositories/message.rs index 5c1601c..aacf878 100644 --- a/src/repositories/message.rs +++ b/src/repositories/message.rs @@ -1,6 +1,7 @@ use super::types::MessageFilter; -use crate::models::message; +use crate::models::{channel, message}; use crate::repositories::{AnyResult, RepositoryContext}; +use event_bus::Scope; use sea_orm::{ActiveModelTrait, ColumnTrait, EntityTrait, QueryFilter, QueryOrder, QuerySelect}; use std::sync::Arc; @@ -53,7 +54,32 @@ impl MessageRepository { pub async fn create(&self, active: message::ActiveModel) -> AnyResult { let message = active.insert(&self.context.db).await?; - self.context.events.emit("message_created", message.clone()); + + // self.context.events.emit("message_created", message.clone()); + + // todo : test + // Ici l'évènement est déclencher sur les topic suivant : + // message_created + // channel:_channel_uuid_:message_created + // si server : server:_server_uuid_:message_created + // scoped event + let mut scopes: Vec = Vec::new(); + scopes.push(Scope::uuid("channel", message.channel_id)); + // retrieve related channel and server + let server_id: Option = channel::Entity::find_by_id(message.channel_id) + .select_only() + .column(channel::Column::ServerId) + .into_tuple::>() + .one(&self.context.db) + .await? + .flatten(); + + if let Some(server_id) = server_id { + scopes.push(Scope::uuid("server", server_id)); + } + self.context + .events + .emit_scoped("message_created", scopes, message.clone()); Ok(message) }