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::{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}