bamboo_engine/runtime/execution/
session_events.rs1use 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
17pub const SESSION_EVENT_CHANNEL_CAPACITY: usize = 1000;
19
20pub 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 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#[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 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}