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