From afed5c01b4d3c923788545d2002d589ed8cde85c Mon Sep 17 00:00:00 2001 From: Gabriel Kaszewski Date: Sun, 12 Jul 2026 02:54:43 +0200 Subject: [PATCH] adapter-event-publisher: broadcast channel event bus --- Cargo.toml | 2 +- crates/adapters/event-publisher/Cargo.toml | 10 +++++ crates/adapters/event-publisher/src/lib.rs | 45 ++++++++++++++++++++++ 3 files changed, 56 insertions(+), 1 deletion(-) create mode 100644 crates/adapters/event-publisher/Cargo.toml create mode 100644 crates/adapters/event-publisher/src/lib.rs diff --git a/Cargo.toml b/Cargo.toml index 0c3ebb0..d2d093b 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -1,5 +1,5 @@ [workspace] -members = ["crates/domain", "crates/application", "crates/api-types", "crates/infra-wiring", "crates/adapters/adapter-common", "crates/adapters/sqlite", "crates/adapters/postgres", "crates/adapters/auth", "crates/adapters/jellyfin", "crates/adapters/local-files"] +members = ["crates/domain", "crates/application", "crates/api-types", "crates/infra-wiring", "crates/adapters/adapter-common", "crates/adapters/sqlite", "crates/adapters/postgres", "crates/adapters/auth", "crates/adapters/jellyfin", "crates/adapters/local-files", "crates/adapters/event-publisher"] exclude = ["k-tv-backend", "k-tv-frontend"] resolver = "2" diff --git a/crates/adapters/event-publisher/Cargo.toml b/crates/adapters/event-publisher/Cargo.toml new file mode 100644 index 0000000..c812897 --- /dev/null +++ b/crates/adapters/event-publisher/Cargo.toml @@ -0,0 +1,10 @@ +[package] +name = "adapter-event-publisher" +version = "0.1.0" +edition = "2024" + +[dependencies] +domain = { workspace = true } +async-trait = { workspace = true } +tokio = { workspace = true } +tracing = { workspace = true } diff --git a/crates/adapters/event-publisher/src/lib.rs b/crates/adapters/event-publisher/src/lib.rs new file mode 100644 index 0000000..aefcb43 --- /dev/null +++ b/crates/adapters/event-publisher/src/lib.rs @@ -0,0 +1,45 @@ +use async_trait::async_trait; +use domain::errors::{DomainError, DomainResult}; +use domain::events::DomainEvent; +use domain::ports::events::{EventConsumer, EventPublisher}; +use tokio::sync::broadcast; + +pub struct ChannelEventBus { + tx: broadcast::Sender, +} + +impl ChannelEventBus { + pub fn new(capacity: usize) -> Self { + let (tx, _) = broadcast::channel(capacity); + Self { tx } + } + + pub fn subscriber(&self) -> broadcast::Receiver { + self.tx.subscribe() + } + + pub fn sender(&self) -> broadcast::Sender { + self.tx.clone() + } +} + +#[async_trait] +impl EventPublisher for ChannelEventBus { + async fn publish(&self, event: DomainEvent) -> DomainResult<()> { + let _ = self.tx.send(event); // Ok to drop if no receivers + Ok(()) + } +} + +#[async_trait] +impl EventConsumer for ChannelEventBus { + async fn recv(&self) -> DomainResult { + // Note: This creates a new subscriber each call — for real use, + // the presentation layer should hold a receiver from subscriber() + // This impl exists to satisfy the port trait + let mut rx = self.tx.subscribe(); + rx.recv() + .await + .map_err(|e| DomainError::InfrastructureError(e.to_string())) + } +}