Skip to main content

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}