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}