Skip to main content

agent_base/engine/runtime/
notify.rs

1//! Run-level notice channel: NoticeHandle + engine pump.
2//!
3//! Engine mechanisms (guard, judge, future reaper) hold a cloneable
4//! [`NoticeHandle`] and send display-only [`Notice`]s. The pump (see
5//! [`RuntimeCore::drain_notices`]) stamps each notice with the current
6//! session id and forwards it into the event stream as
7//! `RuntimeEvent::UserEvent { event: UserEvent::Notice { .. } }` — delivered
8//! to the run's `on_event` callback and dual-written to the event bus so
9//! external subscribers (bridge remote clients, persistence) see it too.
10//!
11//! Delivery is deliberately synchronous (try_recv drain at turn boundaries
12//! and before each terminal event) rather than a spawned task: the guard
13//! decides at `on_turn`, so a synchronous drain at the turn boundary
14//! guarantees the notice is delivered before the next turn emits anything —
15//! and touches none of the fan-in select loops. Terminal drains
16//! (`RunFinished`/`RunCancelled`) cover whatever the boundary drain has not
17//! delivered yet. Consequence: a notice sent mid-turn (e.g. during tool
18//! execution) surfaces at the next turn boundary, not in real time.
19//!
20//! Lifecycle: the channel is created once per engine. Any notice still
21//! pending when a new run starts is discarded (`discard_pending_notices`) so
22//! it can never be stamped with a later run's session id. Sends outside a
23//! run are therefore best-effort and may be dropped.
24
25use std::sync::Mutex;
26
27use tokio::sync::mpsc;
28
29use crate::types::{AgentResult, Notice, NoticeKind, RuntimeEvent, SessionId, UserEvent};
30
31use super::plan_runner::RuntimeCore;
32
33/// Cloneable publish handle for engine-internal notices.
34///
35/// Publishers don't know whether a session/engine exists — sending into a
36/// closed channel (run over, runtime dropped) silently discards the notice:
37/// it is a best-effort display signal and must never affect execution.
38#[derive(Clone)]
39pub struct NoticeHandle {
40    tx: mpsc::UnboundedSender<Notice>,
41}
42
43impl NoticeHandle {
44    pub(crate) fn new(tx: mpsc::UnboundedSender<Notice>) -> Self {
45        Self { tx }
46    }
47
48    /// Create a standalone handle + receiver pair.
49    ///
50    /// The runtime builds its own channel internally; this constructor exists
51    /// for hosts and tests that want to drive a guard (or any publisher)
52    /// outside a full engine build and observe the raw notices.
53    pub fn new_channel() -> (Self, mpsc::UnboundedReceiver<Notice>) {
54        let (tx, rx) = mpsc::unbounded_channel();
55        (Self { tx }, rx)
56    }
57
58    /// Send a notice. Never fails from the caller's perspective — a closed
59    /// receiver drops the notice silently.
60    ///
61    /// `source` takes `impl Into<String>` so dynamic emitters ("tool:<name>")
62    /// work alongside static ones ("guard").
63    pub fn send(&self, kind: NoticeKind, source: impl Into<String>, text: impl Into<String>) {
64        let _ = self.tx.send(Notice {
65            kind,
66            source: source.into(),
67            text: text.into(),
68        });
69    }
70}
71
72impl RuntimeCore {
73    /// Drain pending notices, stamping each with `session_id` and delivering
74    /// it to the run's event callback as a `UserEvent::Notice`.
75    ///
76    /// Called at turn boundaries (react loop, before the next turn emits) and
77    /// after the react loop completes / before every terminal event
78    /// (`RunFinished`/`RunCancelled`) — a notice is therefore always
79    /// delivered before the terminal event on every exit path, Ok or Err,
80    /// and mid-run notices surface at the next turn boundary.
81    ///
82    /// Display-only contract: a consumer callback error is logged and the
83    /// notice dropped — it never fails the run nor blocks later notices.
84    pub(crate) fn drain_notices<F>(
85        &self,
86        session_id: &SessionId,
87        on_event: &Mutex<F>,
88    ) -> AgentResult<()>
89    where
90        F: FnMut(RuntimeEvent) -> AgentResult<()>,
91    {
92        let mut rx = self.notice_rx.lock().unwrap();
93        while let Ok(notice) = rx.try_recv() {
94            let event = RuntimeEvent::UserEvent {
95                session_id: session_id.clone(),
96                event: UserEvent::Notice {
97                    kind: notice.kind,
98                    source: notice.source,
99                    text: notice.text,
100                },
101                agent_id: None,
102                trace_id: None,
103            };
104            // Dual-write: deliver to on_event for the local renderer, and emit
105            // to the event bus so external subscribers (bridge remote clients,
106            // persistence) see the notice too. The bus copy carries
107            // `agent_id: None` — the self-echo marker — so bus loopback drops
108            // it instead of re-rendering (see `is_self_echo_user_event`).
109            if let Ok(mut cb) = on_event.lock()
110                && let Err(e) = cb(event.clone())
111            {
112                tracing::warn!(error = %e, "notice consumer callback failed — notice still published to the event bus");
113            }
114            self.event_bus.emit(event);
115        }
116        Ok(())
117    }
118
119    /// Drop any notices still pending from before this run.
120    ///
121    /// Called at run entry: the notice channel is engine-level, so without
122    /// this a notice sent outside a run (or left over from an aborted one)
123    /// would surface later stamped with the wrong run's session id.
124    pub(crate) fn discard_pending_notices(&self) {
125        let mut rx = self.notice_rx.lock().unwrap();
126        while rx.try_recv().is_ok() {}
127    }
128}
129
130#[cfg(test)]
131mod tests {
132    use super::*;
133
134    // Design doc §9: receiver closed (run over) — send must not panic and
135    // must not affect the caller.
136    #[test]
137    fn send_after_receiver_closed_is_silent_no_panic() {
138        let (tx, rx) = mpsc::unbounded_channel::<Notice>();
139        let handle = NoticeHandle::new(tx);
140        drop(rx);
141        handle.send(NoticeKind::Warning, "guard", "nobody listening");
142        handle.send(NoticeKind::Info, "judge", "also dropped");
143    }
144
145    #[test]
146    fn send_stores_kind_source_text() {
147        let (tx, mut rx) = mpsc::unbounded_channel::<Notice>();
148        let handle = NoticeHandle::new(tx);
149        handle.send(
150            NoticeKind::Warning,
151            "guard",
152            "judge unparsed — treating as complete",
153        );
154        let notice = rx.try_recv().unwrap();
155        assert_eq!(notice.kind, NoticeKind::Warning);
156        assert_eq!(notice.source, "guard");
157        assert_eq!(notice.text, "judge unparsed — treating as complete");
158    }
159}