Skip to main content

bamboo_engine/runtime/execution/
session_events.rs

1//! Session-scoped event sender management.
2//!
3//! Provides long-lived broadcast senders for session event streams.
4//! Unlike runner-scoped senders (which exist only during agent execution),
5//! these persist for the lifetime of the session.
6
7use std::collections::HashMap;
8use std::sync::Arc;
9
10use tokio::sync::{broadcast, RwLock};
11
12use bamboo_agent_core::AgentEvent;
13
14use super::event_forwarder::AccountFeedInbox;
15use super::runner_state::AgentRunner;
16
17/// Default broadcast channel capacity for session event senders.
18pub const SESSION_EVENT_CHANNEL_CAPACITY: usize = 1000;
19
20/// Get or create a broadcast sender for the given session.
21///
22/// If a sender already exists in the map, returns a clone of it.
23/// Otherwise creates a new one with [`SESSION_EVENT_CHANNEL_CAPACITY`]
24/// capacity and inserts it.
25pub async fn get_or_create_event_sender(
26    senders: &Arc<RwLock<HashMap<String, broadcast::Sender<AgentEvent>>>>,
27    session_id: &str,
28) -> broadcast::Sender<AgentEvent> {
29    {
30        let read = senders.read().await;
31        if let Some(sender) = read.get(session_id) {
32            return sender.clone();
33        }
34    }
35
36    let mut write = senders.write().await;
37    // Double-check after acquiring write lock.
38    if let Some(sender) = write.get(session_id) {
39        return sender.clone();
40    }
41
42    let (sender, _) = broadcast::channel(SESSION_EVENT_CHANNEL_CAPACITY);
43    write.insert(session_id.to_string(), sender.clone());
44    sender
45}
46
47/// Atomic cache + account-feed + broadcast publisher for replayable session
48/// events produced outside an agent event forwarder.
49///
50/// Subscribers establish their receiver and clone `last_critical_events` while
51/// holding the matching runner read lock. Keeping cache mutation and broadcast
52/// inside the write lock means each publication appears either in that replay
53/// snapshot or in the live receiver, never both and never neither.
54#[derive(Clone)]
55pub(crate) struct ReplayableSessionEventPublisher {
56    runners: Arc<RwLock<HashMap<String, AgentRunner>>>,
57    senders: Arc<RwLock<HashMap<String, broadcast::Sender<AgentEvent>>>>,
58    account_feed_inbox: Option<AccountFeedInbox>,
59}
60
61impl ReplayableSessionEventPublisher {
62    pub(crate) fn new(
63        runners: Arc<RwLock<HashMap<String, AgentRunner>>>,
64        senders: Arc<RwLock<HashMap<String, broadcast::Sender<AgentEvent>>>>,
65        account_feed_inbox: Option<AccountFeedInbox>,
66    ) -> Self {
67        Self {
68            runners,
69            senders,
70            account_feed_inbox,
71        }
72    }
73
74    pub(crate) async fn publish(&self, session_id: &str, event: AgentEvent) {
75        debug_assert!(event.is_replayable_session_state());
76        // Resolve the sender before the runner lock. The SSE/WS subscription
77        // boundary uses the same sender-then-runner lock order.
78        let sender = get_or_create_event_sender(&self.senders, session_id).await;
79        let mut runners = self.runners.write().await;
80        if let Some(runner) = runners.get_mut(session_id) {
81            runner.push_critical_event(event.clone());
82        }
83        if let Some(inbox) = &self.account_feed_inbox {
84            if event.is_durable_change() {
85                let route_session_id = event.session_id().unwrap_or(session_id);
86                let _ = inbox.try_send((Some(route_session_id.to_string()), event.clone()));
87            }
88        }
89        let _ = sender.send(event);
90    }
91}