feat: add main features
Agents for developpement added + frontend add + backend added. Git viewer created + agent and template creator + layout and project creator
This commit is contained in:
86
crates/infrastructure/src/eventbus/mod.rs
Normal file
86
crates/infrastructure/src/eventbus/mod.rs
Normal file
@ -0,0 +1,86 @@
|
||||
//! [`TokioBroadcastEventBus`] — in-process [`EventBus`] backed by
|
||||
//! [`tokio::sync::broadcast`].
|
||||
//!
|
||||
//! `publish` fans out to all current subscribers. `subscribe` returns the
|
||||
//! domain's [`EventStream`] (a `Box<dyn Iterator>`): a **blocking** iterator
|
||||
//! that pulls events off the broadcast receiver. The Tauri event relay drives
|
||||
//! it from a dedicated thread / blocking task, so blocking is acceptable and
|
||||
//! keeps the domain port signature (`-> EventStream`) intact without forcing it
|
||||
//! to become async.
|
||||
|
||||
use domain::events::DomainEvent;
|
||||
use domain::ports::{EventBus, EventStream};
|
||||
use tokio::sync::broadcast;
|
||||
|
||||
/// Default capacity of the broadcast ring buffer.
|
||||
const DEFAULT_CAPACITY: usize = 1024;
|
||||
|
||||
/// An in-process event bus relaying [`DomainEvent`]s to all subscribers via a
|
||||
/// Tokio broadcast channel.
|
||||
#[derive(Clone)]
|
||||
pub struct TokioBroadcastEventBus {
|
||||
sender: broadcast::Sender<DomainEvent>,
|
||||
}
|
||||
|
||||
impl TokioBroadcastEventBus {
|
||||
/// Creates a bus with the default buffer capacity.
|
||||
#[must_use]
|
||||
pub fn new() -> Self {
|
||||
Self::with_capacity(DEFAULT_CAPACITY)
|
||||
}
|
||||
|
||||
/// Creates a bus with an explicit buffer capacity.
|
||||
#[must_use]
|
||||
pub fn with_capacity(capacity: usize) -> Self {
|
||||
let (sender, _rx) = broadcast::channel(capacity);
|
||||
Self { sender }
|
||||
}
|
||||
|
||||
/// Returns a raw async receiver, useful for relays that prefer to consume
|
||||
/// events on the Tokio runtime rather than via the blocking [`EventStream`].
|
||||
#[must_use]
|
||||
pub fn raw_receiver(&self) -> broadcast::Receiver<DomainEvent> {
|
||||
self.sender.subscribe()
|
||||
}
|
||||
}
|
||||
|
||||
impl Default for TokioBroadcastEventBus {
|
||||
fn default() -> Self {
|
||||
Self::new()
|
||||
}
|
||||
}
|
||||
|
||||
impl EventBus for TokioBroadcastEventBus {
|
||||
fn publish(&self, event: DomainEvent) {
|
||||
// A send error only means there are currently no subscribers; that is
|
||||
// not an error condition for a fire-and-forget bus.
|
||||
let _ = self.sender.send(event);
|
||||
}
|
||||
|
||||
fn subscribe(&self) -> EventStream {
|
||||
Box::new(BroadcastIter {
|
||||
rx: self.sender.subscribe(),
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
/// Blocking iterator adapter over a broadcast receiver. `next()` blocks until an
|
||||
/// event arrives; it ends when the channel is closed, and skips past `Lagged`
|
||||
/// notices (dropping the lagged count and continuing).
|
||||
struct BroadcastIter {
|
||||
rx: broadcast::Receiver<DomainEvent>,
|
||||
}
|
||||
|
||||
impl Iterator for BroadcastIter {
|
||||
type Item = DomainEvent;
|
||||
|
||||
fn next(&mut self) -> Option<Self::Item> {
|
||||
loop {
|
||||
match self.rx.blocking_recv() {
|
||||
Ok(event) => return Some(event),
|
||||
Err(broadcast::error::RecvError::Lagged(_)) => continue,
|
||||
Err(broadcast::error::RecvError::Closed) => return None,
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user