Skip to main content

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::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 supervisor (if any) so we can call
263    // `shutdown()` on the exit path. Cloning here is cheap (`Arc`) and
264    // detaches the lifetime of the supervisor from the `state` move into
265    // the router below.
266    let bm25_supervisor = state.bm25_supervisor.clone();
267
268    // #3304: trust the resolved bind address as a self-origin (non-loopback
269    // binds only; `from_bind_addrs` drops loopback) for the router-wide write
270    // guard, so a non-loopback bind still serves its own write UI. The `/sse`
271    // route is now registered INSIDE `router_with_self_origins` (before the
272    // guard/middleware layer) rather than chained on here — see #3304 review:
273    // routes added after `.layer()` get no middleware per axum's contract.
274    let self_origins = match local {
275        Some(a) => trusty_common::server::SelfOrigins::from_bind_addrs(&[a]),
276        None => trusty_common::server::SelfOrigins::default(),
277    };
278    let app = crate::web::router_with_self_origins(self_origins).with_state(state);
279
280    // Why (issue #534): bare axum::serve exits only on an internal error; SIGTERM
281    // (launchctl bootout) would kill the process before the cleanup below had a
282    // chance to run, leaving stale addr/socket files behind and dropping any
283    // in-flight request without draining. `with_graceful_shutdown` installs a
284    // SIGTERM + SIGINT watcher; when either fires axum stops accepting new
285    // connections, drains active requests, then returns here so cleanup runs.
286    let serve_result = axum::serve(listener, app)
287        .with_graceful_shutdown(trusty_common::shutdown_signal())
288        .await;
289
290    // Best-effort cleanup: remove `http_addr` files so stale clients fail fast
291    // instead of timing out against a dead port. Remove both the OS-standard
292    // path and the dotfile path (#498).
293    if let Some(p) = written_path.as_ref() {
294        let _ = std::fs::remove_file(p);
295    }
296    if let Some(p) = written_dotfile_path.as_ref() {
297        let _ = std::fs::remove_file(p);
298    }
299
300    // Issue #193: gracefully reap every spawned BM25 daemon before the
301    // process exits so each one gets a chance to flush its snapshot and
302    // unlink its socket. `kill_on_drop=true` on the children would
303    // SIGKILL them on Drop anyway, but that skips the daemon's own
304    // shutdown sequence and leaves stale sockets behind.
305    if let Some(supervisor) = bm25_supervisor {
306        supervisor.shutdown().await;
307    }
308
309    serve_result?;
310    Ok(())
311}
312
313/// Convenience: bind `addr` and serve via [`run_http_on`].
314#[cfg(feature = "axum-server")]
315pub async fn run_http(state: AppState, addr: std::net::SocketAddr) -> Result<()> {
316    let listener = tokio::net::TcpListener::bind(addr).await?;
317    run_http_on(state, listener).await
318}
319
320/// Convenience: bind dynamically (7070..=7079, OS fallback) and serve.
321///
322/// Why: `trusty-memory serve` with no `--http` flag is the canonical
323/// launchd-managed daemon entry point. Dynamic binding lets a stale daemon
324/// or a hand-spawned `serve --http 127.0.0.1:7070` coexist without breaking
325/// the launchd-managed instance.
326/// What: calls [`bind_dynamic_port`] then [`run_http_on`].
327/// Test: integration via `trusty-memory serve` + `cat ~/.trusty-memory/http_addr`.
328#[cfg(feature = "axum-server")]
329pub async fn run_http_dynamic(state: AppState) -> Result<()> {
330    let listener = bind_dynamic_port().await?;
331    run_http_on(state, listener).await
332}
333
334/// Interval between recomputations of the data-root disk footprint.
335///
336/// Why (#4764): see the cadence note inside [`spawn_disk_size_ticker`].
337#[cfg(feature = "axum-server")]
338const DISK_SIZE_INTERVAL: std::time::Duration = std::time::Duration::from_secs(60);
339
340/// Spawn a background ticker that recomputes the `data_root` disk footprint
341/// on a fixed cadence and stores it in `state.disk_bytes` (issue #35).
342///
343/// Why: `GET /health` reports `disk_bytes`. Walking the data directory on
344/// every health request would turn a frequent health poll into unbounded
345/// recursive I/O. Computing it off the request path on a fixed cadence keeps
346/// `/health` cheap and bounds the staleness to [`DISK_SIZE_INTERVAL`] — fine
347/// for an at-a-glance footprint figure.
348/// What: spawns a detached tokio task. `AppState` is cheap to `Clone` (all
349/// `Arc` fields), so the task holds a full clone; the daemon process lives
350/// for the lifetime of the server anyway, so no `Weak` downgrade is needed.
351/// Each tick runs the blocking directory walk on `spawn_blocking` so it never
352/// stalls the async runtime, then stores the byte total atomically.
353/// Test: `health_endpoint_includes_resource_fields` asserts the field shape;
354/// the ticker cadence is not unit-tested (timing-dependent).
355#[cfg(feature = "axum-server")]
356fn spawn_disk_size_ticker(state: AppState) {
357    tokio::spawn(async move {
358        // #4764: 10 s → 60 s, matching trusty-search. trusty-search aborted 40
359        // times on this shared walk; trusty-memory runs the identical ticker
360        // against the identical function and shares the exposure, so it gets
361        // the identical cadence reduction.
362        let mut interval = tokio::time::interval(DISK_SIZE_INTERVAL);
363        loop {
364            interval.tick().await;
365            let dir = state.data_root.clone();
366            // The directory walk is blocking filesystem I/O — run it on the
367            // blocking pool so it never parks an async worker thread.
368            let bytes = tokio::task::spawn_blocking(move || {
369                trusty_common::sys_metrics::dir_size_bytes(&dir)
370            })
371            .await
372            .unwrap_or(0);
373            state
374                .disk_bytes
375                .store(bytes, std::sync::atomic::Ordering::Relaxed);
376        }
377    });
378}
379
380/// Interval between aggregate-status snapshot emits on the SSE bus.
381///
382/// Why (issue #228): mutations used to fire `StatusChanged` synchronously on
383/// the write path, which forced an O(N palaces) sum of drawer / vector / KG
384/// counts on every `memory_remember`. Coalescing into a fixed-cadence ticker
385/// lets dashboards stay current (a 30 s lag is invisible at human scale)
386/// while keeping the write path free of aggregate work.
387/// What: 30 seconds — short enough that the operator UI doesn't feel stale
388/// between manual writes, long enough that the recompute cost (in-memory
389/// registry walk plus the redb `count_active_triples` per palace) is a
390/// rounding error on the daemon's CPU budget.
391/// Test: covered indirectly — the math has not changed, only the cadence.
392#[allow(dead_code)]
393const STATUS_EVENT_TICK_SECS: u64 = 30;
394
395/// Spawn a background ticker that emits `DaemonEvent::StatusChanged` every
396/// [`STATUS_EVENT_TICK_SECS`] seconds (issue #228).
397///
398/// Why: replaces the per-write `state.emit(self.aggregate_status_event())`
399/// call sites that used to recompute the aggregate every time a drawer was
400/// created or deleted. Walking N palaces on every write blocks the async
401/// runtime; coalescing the emit onto a ticker keeps dashboards up-to-date
402/// without that cost.
403/// What: spawns a detached tokio task that holds a full `AppState` clone
404/// (cheap — every field is `Arc`-backed) and ticks every
405/// [`STATUS_EVENT_TICK_SECS`] seconds. Each tick computes
406/// `MemoryService::aggregate_status_event` (which now iterates the
407/// in-memory registry, not disk) and broadcasts it via `state.emit`. If
408/// no SSE subscribers are connected the broadcast `send` is a cheap no-op,
409/// so the ticker imposes no cost when nobody is listening.
410/// Test: not unit-tested (timing-dependent fire-and-forget); the underlying
411/// `aggregate_status_event` math is exercised by the existing
412/// `status_endpoint_returns_payload` path.
413#[allow(dead_code)]
414fn spawn_status_event_ticker(state: AppState) {
415    tokio::spawn(async move {
416        let mut interval =
417            tokio::time::interval(std::time::Duration::from_secs(STATUS_EVENT_TICK_SECS));
418        // The first tick fires immediately, which is fine: it gives SSE
419        // subscribers a baseline `StatusChanged` shortly after they connect.
420        loop {
421            interval.tick().await;
422            let event = crate::service::MemoryService::new(state.clone()).aggregate_status_event();
423            state.emit(event);
424        }
425    });
426}
427
428/// Live SSE event stream — pushes `DaemonEvent` frames to dashboard clients.
429///
430/// Why: The dashboard subscribes once and reacts to live pushes (palace
431/// created, drawer added/deleted, dream completed, status changed) instead of
432/// polling `/api/v1/*` endpoints.
433/// What: Subscribes to `state.events`, emits an initial `connected` frame,
434/// then forwards every `DaemonEvent` as `data: <json>\n\n`. Lagged
435/// subscribers receive a `lag` frame indicating skipped events; channel
436/// closure ends the stream.
437/// Test: `web::tests::sse_stream_emits_palace_created` (covers subscribe +
438/// emit + receive); manual: `curl -N http://.../sse`.
439#[cfg(feature = "axum-server")]
440pub(crate) async fn sse_handler(
441    axum::extract::State(state): axum::extract::State<AppState>,
442) -> impl axum::response::IntoResponse {
443    use futures::StreamExt;
444    use tokio_stream::wrappers::BroadcastStream;
445
446    let rx = state.events.subscribe();
447    let initial = futures::stream::once(async {
448        Ok::<axum::body::Bytes, std::io::Error>(axum::body::Bytes::from(
449            "data: {\"type\":\"connected\"}\n\n",
450        ))
451    });
452    let events = BroadcastStream::new(rx).map(|res| {
453        let frame = match res {
454            Ok(event) => match serde_json::to_string(&event) {
455                Ok(json) => format!("data: {json}\n\n"),
456                Err(e) => format!("data: {{\"type\":\"error\",\"message\":\"{e}\"}}\n\n"),
457            },
458            Err(tokio_stream::wrappers::errors::BroadcastStreamRecvError::Lagged(n)) => {
459                format!("data: {{\"type\":\"lag\",\"skipped\":{n}}}\n\n")
460            }
461        };
462        Ok::<axum::body::Bytes, std::io::Error>(axum::body::Bytes::from(frame))
463    });
464    let stream = initial.chain(events);
465
466    axum::response::Response::builder()
467        .header("Content-Type", "text/event-stream")
468        .header("Cache-Control", "no-cache")
469        .header("X-Accel-Buffering", "no")
470        .body(axum::body::Body::from_stream(stream))
471        .expect("valid SSE response") // Why: invariant — SSE headers are compile-time constants; builder cannot fail
472}