trusty_memory/transport/methods/activity.rs
1//! The activity history, and the live feed that replaces `/sse` (#6286).
2//!
3//! Why: `GET /api/v1/activity` seeded the console's activity feed on mount;
4//! without it the pane rendered empty until the next live event. The hook
5//! ingestion route it shared a file with is NOT folded — `hook_fired` was
6//! already a dispatcher method, and the route was its duplicate.
7//!
8//! [`activity_stream`] is what `/sse` was. The first pass at this migration
9//! retired that listener with nothing in its place, so the monitor TUI polled
10//! `memory.activity` on a 2-second tick — an event appeared up to two seconds
11//! late, and one evicted from the log between two ticks was never seen at all.
12//! This streams from the SAME `broadcast::Sender<DaemonEvent>` the SSE handler
13//! subscribed to, so an event reaches a reader as it is emitted.
14//!
15//! What: `memory.activity` (a page, with the same filters and the same clamp)
16//! and `memory.activity_stream` (every event from now on, as it happens).
17//! Test: `super::super::uds::tests` — `rpc_activity_*`.
18
19use serde::{Deserialize, Serialize};
20use serde_json::Value;
21use tokio::sync::broadcast::error::RecvError;
22use tokio::sync::mpsc;
23use trusty_common::uds::server::{RpcError, RpcStreamItems};
24
25use crate::transport::api_error::ApiError;
26use crate::{ActivityFilter, ActivitySource, AppState};
27
28use super::parse_iso_or_bad_request;
29
30/// Default page size — the console's 50-row window.
31const ACTIVITY_DEFAULT_LIMIT: usize = 50;
32
33/// Ceiling on one page.
34///
35/// Bounds both the per-request work and the frame size. 500 is large enough for
36/// ad-hoc inspection without becoming a lever.
37const ACTIVITY_MAX_LIMIT: usize = 500;
38
39/// Params for `memory.activity`. Every filter is optional and they combine
40/// with AND.
41#[derive(Debug, Default, Deserialize)]
42pub struct ActivityParams {
43 /// Page size, clamped to `[1, ACTIVITY_MAX_LIMIT]`.
44 #[serde(default)]
45 pub limit: Option<usize>,
46 /// Rows to skip.
47 #[serde(default)]
48 pub offset: Option<usize>,
49 /// Restrict to one palace.
50 #[serde(default)]
51 pub palace: Option<String>,
52 /// `http` | `mcp` | `hook`.
53 #[serde(default)]
54 pub source: Option<String>,
55 /// RFC 3339 lower bound.
56 #[serde(default)]
57 pub since: Option<String>,
58 /// RFC 3339 upper bound.
59 #[serde(default)]
60 pub until: Option<String>,
61}
62
63/// One row of the activity response.
64///
65/// The persisted entry carries `payload` as a JSON-encoded STRING so the stored
66/// schema is decoupled from `DaemonEvent`'s evolution; it is re-decoded here so
67/// the caller receives an object rather than an escaped string.
68#[derive(Debug, Serialize)]
69pub struct ActivityRow {
70 /// Monotonic row id.
71 pub id: u64,
72 /// When the event was emitted.
73 pub timestamp: chrono::DateTime<chrono::Utc>,
74 /// Which transport produced it.
75 pub source: &'static str,
76 /// The palace it concerned, when it concerned one.
77 #[serde(skip_serializing_if = "Option::is_none")]
78 pub palace_id: Option<String>,
79 /// The `DaemonEvent` variant name.
80 pub event_type: String,
81 /// The event's own body.
82 pub payload: Value,
83}
84
85/// `memory.activity` — a page of activity history (#96).
86///
87/// Answers `{entries, total, limit, offset}` so the caller can tell whether
88/// more rows exist without a second call.
89pub async fn activity(state: &AppState, params: ActivityParams) -> Result<Value, ApiError> {
90 let limit = params
91 .limit
92 .unwrap_or(ACTIVITY_DEFAULT_LIMIT)
93 .clamp(1, ACTIVITY_MAX_LIMIT);
94 let offset = params.offset.unwrap_or(0);
95
96 let source = match params.source.as_deref() {
97 Some(s) => match ActivitySource::parse(s) {
98 Some(parsed) => Some(parsed),
99 None => {
100 return Err(ApiError::bad_request(format!(
101 "unknown source '{s}'; expected one of http, mcp, hook"
102 )));
103 }
104 },
105 None => None,
106 };
107
108 let filter = ActivityFilter {
109 palace_id: params.palace.filter(|s| !s.is_empty()),
110 source,
111 since: parse_iso_or_bad_request(params.since.as_deref(), "since")?,
112 until: parse_iso_or_bad_request(params.until.as_deref(), "until")?,
113 };
114
115 let entries = state
116 .activity_log
117 .list(&filter, limit, offset)
118 .map_err(|e| ApiError::internal(format!("activity list: {e:#}")))?;
119 let total = state
120 .activity_log
121 .count()
122 .map_err(|e| ApiError::internal(format!("activity count: {e:#}")))?;
123
124 let rows: Vec<ActivityRow> = entries
125 .into_iter()
126 .map(|e| {
127 let payload = serde_json::from_str::<Value>(&e.payload)
128 .unwrap_or_else(|_| Value::String(e.payload.clone()));
129 ActivityRow {
130 id: e.id,
131 timestamp: e.timestamp,
132 source: e.source.as_str(),
133 palace_id: e.palace_id,
134 event_type: e.event_type,
135 payload,
136 }
137 })
138 .collect();
139
140 super::to_value(serde_json::json!({
141 "entries": rows,
142 "total": total,
143 "limit": limit,
144 "offset": offset,
145 }))
146}
147
148/// How many events the daemon buffers for one slow reader.
149///
150/// Why: the producer is a broadcast subscriber and the consumer is a socket
151/// write, so a reader that stalls has to be bounded somewhere. 256 is four
152/// times the broadcast channel's own capacity, which means the broadcast
153/// channel's lag guard fires first for a reader that stops draining — and lag
154/// is a frame this method can report, where a full mpsc would only block.
155const STREAM_BUFFER: usize = 256;
156
157/// `memory.activity_stream` — every event from now on, as it happens (#6286).
158///
159/// Why: this is what `/sse` was. Polling `memory.activity` on a tick — the
160/// stopgap the first migration pass left the TUI with — shows an event up to
161/// one tick late and never shows one evicted from the log between two ticks.
162///
163/// What: subscribes to the same `broadcast::Sender<DaemonEvent>` the SSE
164/// handler used and forwards each event as one `"stream":"item"` frame,
165/// carrying the same `type`-tagged body the SSE `data:` lines carried.
166///
167/// **History is NOT replayed.** The stream begins at the subscription, exactly
168/// as `/sse` did; a reader that wants what came before asks `memory.activity`.
169/// Saying so is the whole contract: a caller that assumed a replay would show
170/// an empty pane and conclude the daemon is idle.
171///
172/// **A lagged reader gets a frame, not a silent gap.** A `broadcast` receiver
173/// that falls behind drops the events it missed. Reporting that as
174/// `{"lagged": N}` is what lets a reader say so rather than render a
175/// continuous-looking feed with a hole in it.
176///
177/// The stream ends when the daemon shuts down (the broadcast sender is
178/// dropped) or when the reader disconnects, which closes the mpsc and stops
179/// the task.
180///
181/// # Errors
182///
183/// Never before the first item: the subscription cannot fail. A failure to
184/// serialise one event becomes that event's terminal error frame.
185///
186/// Test: `rpc_activity_stream_delivers_an_event_without_polling`,
187/// `rpc_activity_stream_does_not_replay_history`.
188pub async fn activity_stream(state: &AppState) -> Result<RpcStreamItems, RpcError> {
189 let mut events = state.events.subscribe();
190 let (tx, items) = mpsc::channel(STREAM_BUFFER);
191
192 tokio::spawn(async move {
193 loop {
194 let frame = match events.recv().await {
195 Ok(event) => match serde_json::to_value(&event) {
196 Ok(value) => Ok(value),
197 Err(e) => Err(RpcError::internal(format!("serialize activity event: {e}"))),
198 },
199 // The reader missed `n` events. Telling it beats a gap it
200 // cannot see.
201 Err(RecvError::Lagged(n)) => {
202 tracing::warn!(
203 target: "trusty_memory::activity_stream",
204 "an activity_stream reader lagged {n} events"
205 );
206 Ok(serde_json::json!({ "type": "lagged", "lagged": n }))
207 }
208 // The daemon is shutting down.
209 Err(RecvError::Closed) => break,
210 };
211 let terminal = frame.is_err();
212 if tx.send(frame).await.is_err() {
213 // The reader disconnected.
214 break;
215 }
216 if terminal {
217 break;
218 }
219 }
220 });
221
222 Ok(items)
223}