trusty_memory/http_server.rs
1//! HTTP/SSE daemon surface: port binding, address discovery, and serving.
2//!
3//! Why: the axum HTTP server, its dynamic-port binding, the `http_addr`
4//! discovery-file plumbing, and the SSE stream are a cohesive "how the daemon
5//! is reachable" concern that is orthogonal to the `AppState` model in
6//! `lib.rs`. Splitting it here keeps `lib.rs` under the SLOC cap and lets the
7//! whole HTTP surface (all `axum-server`-gated) sit behind one module boundary.
8//! What: exports `DEFAULT_HTTP_PORT`, `http_addr_path`, `bind_dynamic_port`,
9//! `is_data_dir_override_active`, and the `run_http*` serving entry points
10//! (re-exported at the crate root so existing `trusty_memory::run_http_on`
11//! paths are unchanged).
12//! Test: `lib_tests` covers the address-file helpers and port binding; the
13//! serving entry points are covered by `web::tests` + manual `curl`.
14
15use crate::AppState;
16use anyhow::Result;
17use std::net::SocketAddr;
18use std::path::{Path, PathBuf};
19
20#[cfg(feature = "axum-server")]
21use tracing::info;
22
23/// Preferred starting port for the trusty-memory HTTP daemon.
24///
25/// Why: keeps the well-known default stable for clients that have hard-coded
26/// `127.0.0.1:7070` in their configuration, while still allowing dynamic
27/// walking when the port is in use (`DYNAMIC_PORT_RANGE` ports starting here).
28/// What: `7070` — historic default, matches the launchd plist's prior value.
29/// Test: covered indirectly by `bind_dynamic_port_returns_listener`.
30pub const DEFAULT_HTTP_PORT: u16 = 7070;
31
32/// Number of consecutive ports `bind_dynamic_port` walks before falling back
33/// to the OS-assigned port. Matches the trusty-search convention.
34const DYNAMIC_PORT_RANGE: u16 = 10;
35
36/// Path to the canonical address-discovery file for the trusty-memory daemon.
37///
38/// Why: clients (CLI, MCP tools, dashboards) need to find the running daemon
39/// without configuration when the port was selected dynamically. Using
40/// `trusty_common::resolve_data_dir` aligns this path with the location
41/// that `trusty_common::read_daemon_addr("trusty-memory")` reads from, so
42/// `prompt-context`, `doctor`, and `start`'s probe all find the running daemon.
43/// The old `~/.trusty-memory/http_addr` path and the new
44/// `~/Library/Application Support/trusty-memory/http_addr` (macOS) path were
45/// divergent — the daemon wrote one; readers expected the other.
46/// What: returns `{resolve_data_dir("trusty-memory")}/http_addr`, or `None` if
47/// the data dir cannot be resolved (locked-down container, no passwd entry).
48/// Test: `http_addr_path_uses_resolve_data_dir`.
49pub fn http_addr_path() -> Option<PathBuf> {
50 trusty_common::resolve_data_dir("trusty-memory")
51 .ok()
52 .map(|d| d.join("http_addr"))
53}
54
55/// Bind a `TcpListener` to `127.0.0.1`, dynamically selecting a port.
56///
57/// Why: the historic default `7070` is convenient for clients but a stale
58/// process or a second daemon must not produce a noisy failure. Walking
59/// `DEFAULT_HTTP_PORT..DEFAULT_HTTP_PORT+DYNAMIC_PORT_RANGE` first preserves
60/// backwards compatibility for the common case; OS-assigned fallback (`:0`)
61/// guarantees the daemon always comes up even when every preferred port is
62/// busy.
63/// What: returns the first successful `TcpListener` (7070..=7079, then
64/// OS-assigned); caller inspects `local_addr()` to learn the chosen port.
65/// Test: `bind_dynamic_port_returns_listener` confirms it always binds *some*
66/// port even after another listener occupies the preferred one.
67pub async fn bind_dynamic_port() -> Result<tokio::net::TcpListener> {
68 let preferred: SocketAddr = SocketAddr::from(([127, 0, 0, 1], DEFAULT_HTTP_PORT));
69 // First: walk the preferred range (7070..=7079).
70 if let Ok(listener) =
71 trusty_common::bind_with_auto_port(preferred, DYNAMIC_PORT_RANGE - 1).await
72 {
73 return Ok(listener);
74 }
75 // Last resort: ask the kernel for any free port. `bind_with_auto_port`
76 // with `:0` resolves immediately to the OS-assigned port.
77 tracing::warn!(
78 "all ports {DEFAULT_HTTP_PORT}..{} in use; requesting OS-assigned port",
79 DEFAULT_HTTP_PORT + DYNAMIC_PORT_RANGE - 1
80 );
81 let any: SocketAddr = SocketAddr::from(([127, 0, 0, 1], 0));
82 trusty_common::bind_with_auto_port(any, 0).await
83}
84
85/// Write the bound `host:port` to `~/.trusty-memory/http_addr` atomically.
86///
87/// Why: clients must read the file mid-write without observing a partial
88/// value. Writing to a `.tmp` sibling and renaming over the target gives
89/// POSIX atomicity, matching the trusty-search implementation.
90/// What: creates the parent directory if missing; writes `addr` followed by a
91/// trailing newline (avoids the "no newline at end of file" warnings from
92/// `cat`); renames `.tmp` → `http_addr`. Best-effort: I/O errors are
93/// returned to the caller so `run_http_on` can log without panicking.
94/// Test: `http_addr_file_round_trip_via_helpers`.
95#[cfg(feature = "axum-server")]
96pub(crate) fn write_http_addr_file(path: &Path, addr: &SocketAddr) -> std::io::Result<()> {
97 use std::io::Write;
98 if let Some(parent) = path.parent() {
99 std::fs::create_dir_all(parent)?;
100 }
101 let tmp = path.with_extension("addr.tmp");
102 {
103 let mut f = std::fs::File::create(&tmp)?;
104 writeln!(f, "{addr}")?;
105 f.sync_all()?;
106 }
107 std::fs::rename(&tmp, path)?;
108 Ok(())
109}
110
111/// Return `true` when a non-default data directory is in effect.
112///
113/// Why (issue #880): two startup side-effects must be suppressed when the
114/// daemon runs with an isolated/overridden data root:
115/// 1. The legacy `~/.trusty-memory/http_addr` dotfile write — it would
116/// overwrite the real production daemon's discovery file with the isolated
117/// instance's throwaway address.
118/// 2. The startup pin-scan — it reads project pin files from the **real**
119/// user environment (~/Projects, ~/Developer, …) and imports palaces from
120/// the real environment into the isolated data root, defeating isolation.
121///
122/// A "non-default data dir" means `TRUSTY_DATA_DIR_OVERRIDE` is set to a
123/// non-empty, non-whitespace value. Empty or whitespace-only values are
124/// treated as unset (same rule as `resolve_data_dir`), so an accidental blank
125/// env var does not suppress the dotfile write on real production instances.
126/// What: reads `TRUSTY_DATA_DIR_OVERRIDE`; returns `true` when it contains a
127/// non-empty, non-whitespace string. Returns `false` otherwise.
128/// Test: `is_data_dir_override_active_when_set`,
129/// `is_data_dir_override_inactive_when_unset`,
130/// `is_data_dir_override_inactive_when_blank`.
131#[inline]
132pub fn is_data_dir_override_active() -> bool {
133 matches!(
134 std::env::var(trusty_common::DATA_DIR_OVERRIDE_ENV),
135 Ok(v) if !v.trim().is_empty()
136 )
137}
138
139/// Resolve the dotfile discovery path `~/.trusty-memory/http_addr`.
140///
141/// Why (issue #498): external tooling such as claude-mpm's `migrate_trusty_autodetect`
142/// reads `~/.trusty-memory/http_addr` to find the running daemon's port. On
143/// macOS, `resolve_data_dir("trusty-memory")` returns
144/// `~/Library/Application Support/trusty-memory/`, not `~/.trusty-memory/`,
145/// so the daemon was writing to the OS-standard location while readers expected
146/// the dotfile location. Writing to both locations keeps every reader happy
147/// regardless of which convention they follow.
148///
149/// Fix #880: returns `None` when `TRUSTY_DATA_DIR_OVERRIDE` is active so an
150/// isolated instance (test rig, CI, parallel run) never overwrites the real
151/// production daemon's discovery dotfile.
152///
153/// What: returns `$HOME/.trusty-memory/http_addr` in the default (production)
154/// case, or `None` when `dirs::home_dir()` is unavailable OR when a data-dir
155/// override is active (see `is_data_dir_override_active`).
156/// Test: `dotfile_http_addr_path_uses_home_dir`,
157/// `dotfile_suppressed_when_override_active`.
158#[cfg(feature = "axum-server")]
159pub(crate) fn dotfile_http_addr_path() -> Option<PathBuf> {
160 // Fix #880: never write to the shared dotfile when an override is active.
161 if is_data_dir_override_active() {
162 return None;
163 }
164 dirs::home_dir().map(|h| h.join(".trusty-memory").join("http_addr"))
165}
166
167/// Run the optional HTTP/SSE + web admin server.
168///
169/// Why: A long-running daemon mode lets non-stdio clients (browsers, curl,
170/// future remote agents) hit `/health`, the `/api/v1/*` REST surface, and the
171/// embedded admin SPA. The Unix-domain-socket transport and the
172/// `trusty-memory-mcp-bridge` binary were removed in PR3 of the #914
173/// stdio-cutover epic; the canonical MCP integration is now
174/// `trusty-memory serve --stdio` (PR1 #919).
175/// What: axum router built from `web::router()` plus a `/sse` stub for the
176/// existing MCP-over-SSE clients. Caller provides a pre-bound listener so
177/// port auto-detection lives at the call site. Before accepting connections
178/// the daemon stamps the bound `host:port` onto `AppState.bound_addr` and
179/// writes `~/.trusty-memory/http_addr` so clients can discover the live port.
180/// On shutdown the file is removed best-effort (a stale file with the wrong
181/// port is worse than a missing one).
182/// Test: `cargo test -p trusty-memory web::tests` exercises the router shape;
183/// manual: `curl http://127.0.0.1:<port>/health` returns `ok` with `addr`.
184#[cfg(feature = "axum-server")]
185pub async fn run_http_on(state: AppState, listener: tokio::net::TcpListener) -> Result<()> {
186 use axum::routing::get;
187
188 // Issue #35: recompute the `data_root` disk footprint every 10 s on a
189 // background task so `GET /health` reports `disk_bytes` without doing a
190 // recursive directory walk on the request path.
191 spawn_disk_size_ticker(state.clone());
192
193 // Issue #228: emit aggregate `StatusChanged` on a fixed cadence rather
194 // than on every drawer write. The previous design called
195 // `aggregate_status_event` from every `memory_remember` / `memory_note`
196 // / `memory_forget` (and the matching HTTP handlers), each of which
197 // walked the data root + opened every palace handle. Coalescing the
198 // emit to a 30 s ticker keeps dashboards live without dragging an
199 // O(N palaces) recompute onto the write hot path.
200 spawn_status_event_ticker(state.clone());
201
202 // Capture and advertise the bound address BEFORE serving so the first
203 // request handler — and the http_addr discovery file — see the real port
204 // even if `local_addr()` would otherwise be racy.
205 let local = listener.local_addr().ok();
206 let (written_path, written_dotfile_path) = if let Some(a) = local {
207 // Stash on state for handlers (e.g. /health) to surface.
208 let _ = state.bound_addr.set(a);
209 info!("HTTP server listening on http://{a}");
210 eprintln!("HTTP server listening on http://{a}");
211 // Primary: write to the OS-standard data dir (`~/Library/Application
212 // Support/trusty-memory/http_addr` on macOS, `~/.local/share/…` on
213 // Linux). This is what `trusty_common::read_daemon_addr` reads.
214 // Best-effort: a missing $HOME or read-only fs is non-fatal.
215 let primary = match http_addr_path() {
216 Some(p) => match write_http_addr_file(&p, &a) {
217 Ok(()) => {
218 info!("wrote daemon address to {}", p.display());
219 Some(p)
220 }
221 Err(e) => {
222 tracing::warn!("could not write {}: {e}", p.display());
223 None
224 }
225 },
226 None => {
227 tracing::warn!("no $HOME — skipping http_addr discovery file");
228 None
229 }
230 };
231 // Issue #498: also write to `~/.trusty-memory/http_addr` so external
232 // tools (e.g. claude-mpm's `migrate_trusty_autodetect`) that read the
233 // dotfile path can discover the daemon's port. On macOS the OS-standard
234 // path differs from the dotfile path; writing both ensures consumers
235 // using either convention find the file. Best-effort: failures are
236 // logged but do not block startup.
237 let dotfile = match dotfile_http_addr_path() {
238 Some(p) => match write_http_addr_file(&p, &a) {
239 Ok(()) => {
240 info!("wrote daemon address to dotfile {}", p.display());
241 Some(p)
242 }
243 Err(e) => {
244 tracing::warn!("could not write dotfile {}: {e}", p.display());
245 None
246 }
247 },
248 None => None,
249 };
250 (primary, dotfile)
251 } else {
252 (None, None)
253 };
254
255 // Keep a handle to the BM25 supervisor (if any) so we can call
256 // `shutdown()` on the exit path. Cloning here is cheap (`Arc`) and
257 // detaches the lifetime of the supervisor from the `state` move into
258 // the router below.
259 let bm25_supervisor = state.bm25_supervisor.clone();
260
261 let app = crate::web::router()
262 .route("/sse", get(sse_handler))
263 .with_state(state);
264
265 // Why (issue #534): bare axum::serve exits only on an internal error; SIGTERM
266 // (launchctl bootout) would kill the process before the cleanup below had a
267 // chance to run, leaving stale addr/socket files behind and dropping any
268 // in-flight request without draining. `with_graceful_shutdown` installs a
269 // SIGTERM + SIGINT watcher; when either fires axum stops accepting new
270 // connections, drains active requests, then returns here so cleanup runs.
271 let serve_result = axum::serve(listener, app)
272 .with_graceful_shutdown(trusty_common::shutdown_signal())
273 .await;
274
275 // Best-effort cleanup: remove `http_addr` files so stale clients fail fast
276 // instead of timing out against a dead port. Remove both the OS-standard
277 // path and the dotfile path (#498).
278 if let Some(p) = written_path.as_ref() {
279 let _ = std::fs::remove_file(p);
280 }
281 if let Some(p) = written_dotfile_path.as_ref() {
282 let _ = std::fs::remove_file(p);
283 }
284
285 // Issue #193: gracefully reap every spawned BM25 daemon before the
286 // process exits so each one gets a chance to flush its snapshot and
287 // unlink its socket. `kill_on_drop=true` on the children would
288 // SIGKILL them on Drop anyway, but that skips the daemon's own
289 // shutdown sequence and leaves stale sockets behind.
290 if let Some(supervisor) = bm25_supervisor {
291 supervisor.shutdown().await;
292 }
293
294 serve_result?;
295 Ok(())
296}
297
298/// Convenience: bind `addr` and serve via [`run_http_on`].
299#[cfg(feature = "axum-server")]
300pub async fn run_http(state: AppState, addr: std::net::SocketAddr) -> Result<()> {
301 let listener = tokio::net::TcpListener::bind(addr).await?;
302 run_http_on(state, listener).await
303}
304
305/// Convenience: bind dynamically (7070..=7079, OS fallback) and serve.
306///
307/// Why: `trusty-memory serve` with no `--http` flag is the canonical
308/// launchd-managed daemon entry point. Dynamic binding lets a stale daemon
309/// or a hand-spawned `serve --http 127.0.0.1:7070` coexist without breaking
310/// the launchd-managed instance.
311/// What: calls [`bind_dynamic_port`] then [`run_http_on`].
312/// Test: integration via `trusty-memory serve` + `cat ~/.trusty-memory/http_addr`.
313#[cfg(feature = "axum-server")]
314pub async fn run_http_dynamic(state: AppState) -> Result<()> {
315 let listener = bind_dynamic_port().await?;
316 run_http_on(state, listener).await
317}
318
319/// Spawn a background ticker that recomputes the `data_root` disk footprint
320/// every 10 seconds and stores it in `state.disk_bytes` (issue #35).
321///
322/// Why: `GET /health` reports `disk_bytes`. Walking the data directory on
323/// every health request would turn a frequent health poll into unbounded
324/// recursive I/O. Computing it off the request path on a fixed cadence keeps
325/// `/health` cheap and bounds the staleness to ~10 s — fine for an
326/// at-a-glance footprint figure.
327/// What: spawns a detached tokio task. `AppState` is cheap to `Clone` (all
328/// `Arc` fields), so the task holds a full clone; the daemon process lives
329/// for the lifetime of the server anyway, so no `Weak` downgrade is needed.
330/// Each tick runs the blocking directory walk on `spawn_blocking` so it never
331/// stalls the async runtime, then stores the byte total atomically.
332/// Test: `health_endpoint_includes_resource_fields` asserts the field shape;
333/// the ticker cadence is not unit-tested (timing-dependent).
334#[cfg(feature = "axum-server")]
335fn spawn_disk_size_ticker(state: AppState) {
336 tokio::spawn(async move {
337 let mut interval = tokio::time::interval(std::time::Duration::from_secs(10));
338 loop {
339 interval.tick().await;
340 let dir = state.data_root.clone();
341 // The directory walk is blocking filesystem I/O — run it on the
342 // blocking pool so it never parks an async worker thread.
343 let bytes = tokio::task::spawn_blocking(move || {
344 trusty_common::sys_metrics::dir_size_bytes(&dir)
345 })
346 .await
347 .unwrap_or(0);
348 state
349 .disk_bytes
350 .store(bytes, std::sync::atomic::Ordering::Relaxed);
351 }
352 });
353}
354
355/// Interval between aggregate-status snapshot emits on the SSE bus.
356///
357/// Why (issue #228): mutations used to fire `StatusChanged` synchronously on
358/// the write path, which forced an O(N palaces) sum of drawer / vector / KG
359/// counts on every `memory_remember`. Coalescing into a fixed-cadence ticker
360/// lets dashboards stay current (a 30 s lag is invisible at human scale)
361/// while keeping the write path free of aggregate work.
362/// What: 30 seconds — short enough that the operator UI doesn't feel stale
363/// between manual writes, long enough that the recompute cost (in-memory
364/// registry walk plus the redb `count_active_triples` per palace) is a
365/// rounding error on the daemon's CPU budget.
366/// Test: covered indirectly — the math has not changed, only the cadence.
367#[allow(dead_code)]
368const STATUS_EVENT_TICK_SECS: u64 = 30;
369
370/// Spawn a background ticker that emits `DaemonEvent::StatusChanged` every
371/// [`STATUS_EVENT_TICK_SECS`] seconds (issue #228).
372///
373/// Why: replaces the per-write `state.emit(self.aggregate_status_event())`
374/// call sites that used to recompute the aggregate every time a drawer was
375/// created or deleted. Walking N palaces on every write blocks the async
376/// runtime; coalescing the emit onto a ticker keeps dashboards up-to-date
377/// without that cost.
378/// What: spawns a detached tokio task that holds a full `AppState` clone
379/// (cheap — every field is `Arc`-backed) and ticks every
380/// [`STATUS_EVENT_TICK_SECS`] seconds. Each tick computes
381/// `MemoryService::aggregate_status_event` (which now iterates the
382/// in-memory registry, not disk) and broadcasts it via `state.emit`. If
383/// no SSE subscribers are connected the broadcast `send` is a cheap no-op,
384/// so the ticker imposes no cost when nobody is listening.
385/// Test: not unit-tested (timing-dependent fire-and-forget); the underlying
386/// `aggregate_status_event` math is exercised by the existing
387/// `status_endpoint_returns_payload` path.
388#[allow(dead_code)]
389fn spawn_status_event_ticker(state: AppState) {
390 tokio::spawn(async move {
391 let mut interval =
392 tokio::time::interval(std::time::Duration::from_secs(STATUS_EVENT_TICK_SECS));
393 // The first tick fires immediately, which is fine: it gives SSE
394 // subscribers a baseline `StatusChanged` shortly after they connect.
395 loop {
396 interval.tick().await;
397 let event = crate::service::MemoryService::new(state.clone()).aggregate_status_event();
398 state.emit(event);
399 }
400 });
401}
402
403/// Live SSE event stream — pushes `DaemonEvent` frames to dashboard clients.
404///
405/// Why: The dashboard subscribes once and reacts to live pushes (palace
406/// created, drawer added/deleted, dream completed, status changed) instead of
407/// polling `/api/v1/*` endpoints.
408/// What: Subscribes to `state.events`, emits an initial `connected` frame,
409/// then forwards every `DaemonEvent` as `data: <json>\n\n`. Lagged
410/// subscribers receive a `lag` frame indicating skipped events; channel
411/// closure ends the stream.
412/// Test: `web::tests::sse_stream_emits_palace_created` (covers subscribe +
413/// emit + receive); manual: `curl -N http://.../sse`.
414#[cfg(feature = "axum-server")]
415pub(crate) async fn sse_handler(
416 axum::extract::State(state): axum::extract::State<AppState>,
417) -> impl axum::response::IntoResponse {
418 use futures::StreamExt;
419 use tokio_stream::wrappers::BroadcastStream;
420
421 let rx = state.events.subscribe();
422 let initial = futures::stream::once(async {
423 Ok::<axum::body::Bytes, std::io::Error>(axum::body::Bytes::from(
424 "data: {\"type\":\"connected\"}\n\n",
425 ))
426 });
427 let events = BroadcastStream::new(rx).map(|res| {
428 let frame = match res {
429 Ok(event) => match serde_json::to_string(&event) {
430 Ok(json) => format!("data: {json}\n\n"),
431 Err(e) => format!("data: {{\"type\":\"error\",\"message\":\"{e}\"}}\n\n"),
432 },
433 Err(tokio_stream::wrappers::errors::BroadcastStreamRecvError::Lagged(n)) => {
434 format!("data: {{\"type\":\"lag\",\"skipped\":{n}}}\n\n")
435 }
436 };
437 Ok::<axum::body::Bytes, std::io::Error>(axum::body::Bytes::from(frame))
438 });
439 let stream = initial.chain(events);
440
441 axum::response::Response::builder()
442 .header("Content-Type", "text/event-stream")
443 .header("Cache-Control", "no-cache")
444 .header("X-Accel-Buffering", "no")
445 .body(axum::body::Body::from_stream(stream))
446 .expect("valid SSE response") // Why: invariant — SSE headers are compile-time constants; builder cannot fail
447}