From ccb48156d02c4dd8b405792ee1708e5f17e9840f Mon Sep 17 00:00:00 2001 From: John Ky Date: Thu, 1 Oct 2026 20:19:48 +1000 Subject: [PATCH 1/5] feat(sessions,daemon,vscode,cli,lib,docs,changelog): add diagnostics Make session ingestion, registry decisions, subscription delivery, and companion rejection paths observable without changing session behavior. Add daemon tracing settings and a bounded opt-in wrapper metadata sink. Keep conversation content out of diagnostics and byte forwarding free of file I/O. Document troubleshooting and add regression coverage. Closes #1447 --- CHANGELOG.md | 1 + docs/adrs/adr-0057.md | 13 + docs/configuration.md | 29 ++ docs/sessions-service.md | 137 ++++++++++ editors/vscode/src/extension.ts | 15 +- editors/vscode/src/replyDiagnostics.test.ts | 25 ++ editors/vscode/src/replyDiagnostics.ts | 21 ++ editors/vscode/src/subscription.test.ts | 31 +++ editors/vscode/src/subscription.ts | 12 + src/cli/claude_wrap.rs | 286 +++++++++++++++++++- src/cli/claude_wrap/diagnostics.rs | 182 +++++++++++++ src/cli/sessions.rs | 88 +++++- src/daemon/server.rs | 228 +++++++++++++--- src/daemon/services/sessions.rs | 10 +- src/main.rs | 68 ++++- src/sessions.rs | 139 +++++++++- src/sessions/pid_watcher.rs | 1 + src/sessions/stream.rs | 72 ++++- src/sessions/watcher.rs | 70 ++++- src/test_support.rs | 19 ++ src/utils/settings.rs | 44 +++ 21 files changed, 1391 insertions(+), 100 deletions(-) create mode 100644 editors/vscode/src/replyDiagnostics.test.ts create mode 100644 editors/vscode/src/replyDiagnostics.ts create mode 100644 src/cli/claude_wrap/diagnostics.rs diff --git a/CHANGELOG.md b/CHANGELOG.md index 6cdb4c0ee..157b32b37 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -13,6 +13,7 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 - **Terminal in Editor profile** ([#1683](https://github.com/rust-works/omni-dev/issues/1683)): the VS Code terminal **+ ▾** menu gains an entry that creates a fresh, focused shell terminal as an editor tab using the normal shell and working-directory configuration. - **Jev route top-two margin close calls** ([#2050](https://github.com/rust-works/omni-dev/issues/2050)): `ai jev route` flags model-class stages when confidence is below `--close-call` or their top-two usable tier probabilities differ by less than `--close-call-margin` (default `0.2`). JSON, YAML, and text share the same flags. Set the margin to `0` to retain confidence-only behavior; effort advice keeps its existing confidence rule. Partial maps without two usable offered probabilities fall back to confidence. +- **Session diagnostics** ([#1447](https://github.com/rust-works/omni-dev/issues/1447)): configurable `daemon.log_level`, registry/feed and subscription diagnostics, VS Code rejection/frame reporting, and opt-in metadata-only `OMNI_DEV_CLAUDE_WRAP_LOG` files make stale or missing session cues diagnosable without persisting conversation content. See [session troubleshooting](docs/sessions-service.md#troubleshooting). - **`ai jev route` reports which stage supplies the class** ([#2051](https://github.com/rust-works/omni-dev/issues/2051), #1823 step 3): each provider's result gains `class_from` (`design` or `implement`), and text output annotates the class, as `opus (from design)` or `class: (from implementation)`, for every ladder. `class` is unchanged. A tie goes to `implement`, so `design` means design alone raised the class above what implementation chose. The source is derived from the existing stage answers, so there is no new Jev question or call. See [docs/jev.md](docs/jev.md#route). - **`ai jev` rejects oversized `choice` and `score` questions locally** ([#1774](https://github.com/rust-works/omni-dev/issues/1774)): `ai jev choice`, `score` and every spec in an `ask` questions file now fail before any request when a `choice` has more than 255 options or a `score` more than 10 levels, with an error naming the question, instead of waiting for an API rejection whose body does not say which question was too big. The caps belong to a model version, so they apply only to `jev-latest` and `jev-1.13`/`jev-1.13.*`; any other `--jev-model` skips them and the minimums still apply to every model. See [docs/jev.md](docs/jev.md#keep-the-state-small-and-relevant). - **`ai claude skills {status,sync,clean}` accept `-C/--repo`** ([#1781](https://github.com/rust-works/omni-dev/issues/1781), follow-up to [#1778](https://github.com/rust-works/omni-dev/issues/1778)): `-C/--repo ` now replaces the current directory as the default `sync` source and `status`/`clean` target, and a relative `--source`/`--target` resolves against it, as with `git -C`. Before #1778 the root-global flag parsed but was silently ignored here; since #1778 these commands rejected it. See [docs/user-guide.md](docs/user-guide.md#ai-claude-skills--distribute-skills-across-repositories). diff --git a/docs/adrs/adr-0057.md b/docs/adrs/adr-0057.md index 8eae53709..6c823df59 100644 --- a/docs/adrs/adr-0057.md +++ b/docs/adrs/adr-0057.md @@ -250,3 +250,16 @@ coverage for anything this wrapper does not launch. See [ADR-0039](adr-0039.md) for the daemon framework, and [docs/sessions-service.md](../sessions-service.md) for the operator-facing guide and the wrapper's setup. + + +## Amendment: opt-in diagnostic metadata (#1447) + +Feed 4 keeps conversation content unpersisted. `OMNI_DEV_CLAUDE_WRAP_LOG` +explicitly opts into a user-controlled metadata-only diagnostic file containing +process/session identity, state/report outcomes, and aggregate parser/tee drop +counts. Raw stream lines, messages, tool inputs/results, and daemon error text +are excluded. File I/O runs on an independent bounded writer queue, never on +the byte-forwarding path; failures remain fail-open and shutdown flushing is +bounded. The default remains no diagnostic persistence. See the +[troubleshooting guide](../sessions-service.md#troubleshooting) for enabling, +sharing, and deleting these logs. diff --git a/docs/configuration.md b/docs/configuration.md index 889ba64bc..ed2939263 100644 --- a/docs/configuration.md +++ b/docs/configuration.md @@ -1187,3 +1187,32 @@ Configure one token mode only. Both token keys accept `_FILE` and `_COMMAND`. Process sources override the selected settings profile; a selected profile does not inherit the base `env`. See [Atlassian PAT setup and compatibility](user-guide.md#serverdata-center-personal-access-tokens) for login, verification, mode switching, context-path examples and limitations. + +## Daemon and MCP tracing settings + +The per-user `$HOME/.omni-dev/settings.json` file accepts tracing defaults +independently of project configuration: + +```json +{ + "daemon": { "log_level": "info,omni_dev::sessions=debug,omni_dev::daemon=debug" }, + "mcp": { "log_level": "warn" } +} +``` + +For `omni-dev daemon run`, the first valid filter wins: process `RUST_LOG`, +`daemon.log_level`, then `info`. Directives can be plain levels (`debug`) or +comma-separated module filters. Invalid directives fall through to the next +source with a startup warning; malformed/unreadable settings warn and use +built-in defaults. A missing file or missing `daemon` block is normal. +An empty directive is a valid filter that enables no events. + +The daemon reads this file at startup, including launchd/socket-activated and +systemd launches that do not inherit a shell's environment. Restart the daemon +after changing it (`omni-dev daemon restart`). Other CLI commands retain +`RUST_LOG` → `warn` and do not use `daemon.log_level`. + +`mcp.log_level` configures the separate MCP server (valid `RUST_LOG` takes +precedence there too); it does not control the daemon or hook subprocesses. +See [session troubleshooting](sessions-service.md#troubleshooting) for targeted +filters, the wrapper's opt-in metadata file, and where to read the logs. diff --git a/docs/sessions-service.md b/docs/sessions-service.md index c1448d045..c761c7187 100644 --- a/docs/sessions-service.md +++ b/docs/sessions-service.md @@ -912,3 +912,140 @@ always includes it. - **Windows** support waits on the broader daemon Windows work (#1363); the hook sink and transcript scheme are already portable, only the socket transport is Unix-only. + + +## Troubleshooting + +Session diagnostics let you follow a sighting through ingestion, the registry, +the daemon subscription, and the VS Code companion. They preserve existing +state and attribution rules; an attribution mismatch remains a separate bug +to investigate using the evidence. + +### Enable and read daemon diagnostics + +Set `daemon.log_level` in `$HOME/.omni-dev/settings.json`: + +```json +{ + "daemon": { + "log_level": "info,omni_dev::sessions=debug,omni_dev::daemon=debug" + } +} +``` + +Restart the daemon, then follow its log: + +```bash +omni-dev daemon restart +omni-dev daemon logs --follow +``` + +The log reader supports foreground/background file-log launches and launchd; +a foreground `daemon run` emits tracing to stderr. For a systemd user unit, +read the journal with `journalctl --user -u omni-dev.service -f` instead. +See [configuration](configuration.md#daemon-and-mcp-tracing-settings) for +precedence and fallback behavior. `RUST_LOG` in the daemon process overrides +the settings filter; shell exports are not inherited by socket activation. + +Use `omni_dev::daemon::server=trace` to see each subscription sample's +`change_notification` or `periodic_tick` trigger and `pushed` versus +`suppressed_identical` result. Initial snapshots and terminal counts are +accounted for too. A subscription ends with a distinct `ClientCancel`, +`ClientEof`, `ReadDecodeError`, `DaemonShutdown`, or `WriteFailure` reason. +Cancellation and EOF are normal debug events; failed writes and service +rejections are warnings visible at the default `info` level. + +Registry `session_observed` records include identity, agent/PID/event, +old/new states, an outcome, reap count, and the actual `bumped` decision. +`heartbeat_only` refreshes liveness without triggering a change notification; +`replaced_pid_ignored` explains a late event from the process an in-place resume +replaced. `ended_passive_ignored` prevents passive sightings reviving ended +sessions. `session_end` distinguishes unknown, already ended, ignored, and +newly ended sessions. Window reports and unregisters expose their bump decision. + +Positive TTL reap counts and capacity evictions explain disappearing entries. +At trace level, individual reaps identify `session_ttl`, `ended_ttl`, or +`window_ttl`; `session_process_ended` identifies PID-based cleanup. Busy +stream-wrapper sessions keep alive, while idle ones can age out. A +`session_attribution_miss` reports cwd and the live-window count; trace-level +candidate folders help reveal lexical mismatches such as `/tmp` versus +`/private/tmp`. Matching still uses the existing folder specificity, +registration time, and key ordering, without canonicalization. + +Feed 2 emits one scan summary containing candidate/sighting and skip/error +counts. A blocking scan failure warns that tracking state was reset. Repository +enrichment task failures warn separately from a normal nonrepository cwd. + +### Hook diagnostics + +Hook sinks are short-lived CLI processes; daemon settings do not set their +filter. To investigate hooks, run the configured hook command with +`RUST_LOG=omni_dev::cli::sessions=debug` in its own environment. Diagnostics go +to stderr, never stdout, and do not change exit 0 or the report timeout. +`session_hook_skipped` identifies read/JSON/identity/event/socket gates; +`session_hook_report` distinguishes delivery, timeout, transport failure, and +daemon rejection. Notification diagnostics contain the classification and +presence of type/message fields, never the raw notification text. Claude and +Codex sightings are tagged separately by the shared hook sink. + +### Opt-in wrapper metadata file + +For the Claude stream wrapper, set `OMNI_DEV_CLAUDE_WRAP_LOG` to an absolute +file path in the environment of the process launching the wrapper. For example, +when launching a fresh VS Code process from a shell: + +```bash +OMNI_DEV_CLAUDE_WRAP_LOG="$HOME/claude-wrap-diagnostics.jsonl" code . +``` + +An already running VS Code process may retain its earlier environment; fully +quit/relaunch it and start a new wrapped Claude process to apply the change. +The variable does not enable or install the wrapper by itself: use the Feed 4 +installation instructions above first. + +The file contains newline-delimited metadata: process start/exit, learned +session identity/cwd/model, state reports, report outcome codes, and periodic +plus final cumulative diagnostic summaries. `tee_full` identifies dropped +observer lines when a slow daemon backs up the bounded tee; `tee_closed`, +`tee_oversize`, and `tee_non_utf8` distinguish other drops. Tracker counters +identify parse failure, unknown control shapes, missing permission IDs, and +permission-cap drops. Known unrelated control requests are ignored normally. +The summary also counts drops from the independent diagnostic queue. + +When unset or empty, no file, writer, or diagnostic counters are created. +When enabled, a separate thread appends records through a bounded nonblocking +queue; byte pumps only update counters. New files use `0600`; symlinks and +nonregular targets are refused, and existing file permissions are left intact. +Open/write failures silently disable writing without preventing Claude from +launching. Shutdown waits at most 200 ms for diagnostics; a killed process, +stuck disk, or full queue can lose records. Files append across processes and +include session/PID metadata where known; there is no automatic rotation. + +Conversation messages, tool inputs/results, raw stdio/hook payloads, and daemon +error text are never written to this wrapper file. Paths and identifiers are +still personal metadata: inspect logs before sharing them. Remove the variable +and launch new wrapper processes to disable logging; delete the file when the +investigation is finished. This is an explicit metadata-only persistence +exception to Feed 4's normal no-persistence rule. + +### VS Code output and tracing a stale cue + +Open **View → Output → omni-dev**. One-shot requests now report daemon +rejections as well as transport failures. A rejected/unreachable session-window +report includes the window key; invalid subscription frames report a reason +without dumping the frame, preserving the last valid snapshot. A rejected +subscription continues to use the existing unsupported/fallback behavior. + +Compare evidence in this order: + +1. Did the hook/watcher/wrapper deliver a sighting, or record a gate/drop? +2. Did `session_observed` accept it, change state/metadata, and bump? +3. Did the subscription wake and push, or suppress an identical snapshot? +4. Was the window report accepted and cwd matched to a live window? +5. Did the companion reject the frame or report a connection problem? +6. Did TTL/PID cleanup or capacity eviction remove the session afterward? + +Sessions outside tracked worktrees legitimately contribute no row count. +The pure tally function remains unchanged and does not log every unmatched +session. Gather logs for the narrow reproduction window; use trace only when +debug cannot identify the failing hop. diff --git a/editors/vscode/src/extension.ts b/editors/vscode/src/extension.ts index 6f701224d..8c2a19743 100644 --- a/editors/vscode/src/extension.ts +++ b/editors/vscode/src/extension.ts @@ -38,6 +38,7 @@ import { treeEnvelope, unregisterEnvelope, } from "./socket"; +import { sendWithDiagnostics } from "./replyDiagnostics"; import { runGh } from "./gh"; import { PullRequest, parsePrList, prFallbackBadge, prListArgsForRepo } from "./github"; import { countClaudeTabs, countClaudeTerminals } from "./claudeEmbeddings"; @@ -257,15 +258,9 @@ function registerPayload(): RegisterPayload { * when the daemon was unreachable. `timeoutMs` overrides the default for a * long-running op (the `close` execute waits on windows closing). */ -async function send(envelope: Envelope, timeoutMs?: number): Promise { - try { - return await sendEnvelope(socketPath(), envelope, timeoutMs); - } catch (err) { - output?.appendLine( - `${envelope.op} skipped: ${err instanceof Error ? err.message : String(err)}`, - ); - return undefined; - } +async function send(envelope: Envelope, timeoutMs?: number, context?: string): Promise { + return sendWithDiagnostics(envelope, () => sendEnvelope(socketPath(), envelope, timeoutMs), + (message) => output?.appendLine(message), context); } /** @@ -486,7 +481,7 @@ function claudeEmbeddings(): { tabs: number; terminals: number } { async function reportSessionWindow(): Promise { const folders = (vscode.workspace.workspaceFolders ?? []).map((f) => f.uri.fsPath); const { tabs, terminals } = claudeEmbeddings(); - await send(sessionWindowEnvelope({ key: windowKey, folders, tabs, terminals })); + await send(sessionWindowEnvelope({ key: windowKey, folders, tabs, terminals }), undefined, `window=${windowKey}`); } async function heartbeat(): Promise { diff --git a/editors/vscode/src/replyDiagnostics.test.ts b/editors/vscode/src/replyDiagnostics.test.ts new file mode 100644 index 000000000..cda874feb --- /dev/null +++ b/editors/vscode/src/replyDiagnostics.test.ts @@ -0,0 +1,25 @@ +import assert from "node:assert/strict"; +import { test } from "node:test"; +import { sendWithDiagnostics } from "./replyDiagnostics"; + +test("request diagnostics preserve rejection and include window context once", async () => { + const lines: string[] = []; + const reply = { ok: false, error: "invalid window" }; + const received = await sendWithDiagnostics( + { service: "sessions", op: "window", payload: {} }, + async () => reply, (line) => lines.push(line), "window=abc", + ); + assert.equal(received, reply); + assert.deepEqual(lines, ["sessions/window (window=abc) rejected: invalid window"]); +}); + +test("request diagnostics swallow transport throws and leave success silent", async () => { + const lines: string[] = []; + const env = { service: "sessions", op: "list", payload: {} }; + assert.equal(await sendWithDiagnostics(env, async () => { throw new Error("offline"); }, + (line) => lines.push(line)), undefined); + assert.equal(lines.length, 1); + const reply = { ok: true, payload: {} }; + assert.equal(await sendWithDiagnostics(env, async () => reply, (line) => lines.push(line)), reply); + assert.equal(lines.length, 1); +}); diff --git a/editors/vscode/src/replyDiagnostics.ts b/editors/vscode/src/replyDiagnostics.ts new file mode 100644 index 000000000..e9867383b --- /dev/null +++ b/editors/vscode/src/replyDiagnostics.ts @@ -0,0 +1,21 @@ +import { Envelope, Reply } from "./socket"; + +/** Preserve fail-open replies while making every rejection visible in output. */ +export async function sendWithDiagnostics( + envelope: Envelope, + request: () => Promise, + log: (message: string) => void, + context?: string, +): Promise { + const label = `${envelope.service ?? "daemon"}/${envelope.op}${context ? ` (${context})` : ""}`; + try { + const reply = await request(); + if (!reply.ok) { + log(`${label} rejected: ${reply.error ?? "daemon refused the request"}`); + } + return reply; + } catch (err) { + log(`${label} skipped: ${err instanceof Error ? err.message : String(err)}`); + return undefined; + } +} diff --git a/editors/vscode/src/subscription.test.ts b/editors/vscode/src/subscription.test.ts index 7b68f938e..fac2eb8bb 100644 --- a/editors/vscode/src/subscription.test.ts +++ b/editors/vscode/src/subscription.test.ts @@ -330,3 +330,34 @@ test("sessions subscribe: without onUnsupported an error reply is ignored, as be // behaviour `TreeSubscription` still relies on. assert.equal(received[0].sessions[0].session_id, "s1"); }); + + +test("sessions subscribe: invalid frames report errors and retain the last valid snapshot", async (t: TestContext) => { + const socketPath = tempSocketPath(); + const srv = trackingServer((conn) => { + conn.on("data", () => { + for (const frame of ["null", "{}", '{"ok":true,"payload":{}}', + sessionsLine([]), '{"ok":true,"payload":{"wrong":[]}}', + '{"ok":false,"error":"refused"}', "not json"]) { + conn.write(frame.endsWith("\n") ? frame : frame + "\n"); + } + }); + }); + await srv.listen(socketPath); + const received: SessionsSnapshot[] = []; + const errors: string[] = []; + const statuses: boolean[] = []; + const sub = new SessionsSubscription(socketPath, { + onSnapshot: (snapshot) => received.push(snapshot), + onError: (message) => errors.push(message), + onStatus: (status) => statuses.push(status), + }); + t.after(() => { sub.close(); srv.close(); }); + sub.start(); + await waitFor(() => errors.length === 6); + assert.equal(received.length, 1); + assert.deepEqual(statuses, [true]); + assert.match(errors[0], /invalid snapshot envelope/); + assert.match(errors[2], /invalid snapshot payload/); + assert.equal(errors[4], "refused"); +}); diff --git a/editors/vscode/src/subscription.ts b/editors/vscode/src/subscription.ts index 66f111540..ba7eb6bda 100644 --- a/editors/vscode/src/subscription.ts +++ b/editors/vscode/src/subscription.ts @@ -191,6 +191,10 @@ export class DaemonSubscription { this.onError?.(`malformed snapshot: ${err instanceof Error ? err.message : String(err)}`); return; } + if (!reply || typeof reply !== "object" || typeof reply.ok !== "boolean") { + this.onError?.("invalid snapshot envelope"); + return; + } // An explicit error reply means the daemon is up but will not serve this op // — almost always a daemon too old to know it. Hand that to the caller so it // can degrade, and stop: reconnecting would only re-earn the same refusal. @@ -200,6 +204,14 @@ export class DaemonSubscription { this.onUnsupported(message); return; } + if (!reply.ok) { + this.onError?.(reply.error ?? "daemon refused the subscription"); + return; + } + if (!reply.payload || typeof reply.payload !== "object" || !this.isSnapshot(reply.payload)) { + this.onError?.("invalid snapshot payload for this subscription"); + return; + } // Ignore anything that is not a well-formed snapshot for this stream; a // fresh snapshot is the only frame the stream should carry. if (reply.ok && reply.payload && this.isSnapshot(reply.payload)) { diff --git a/src/cli/claude_wrap.rs b/src/cli/claude_wrap.rs index 2397a8864..c9434a62e 100644 --- a/src/cli/claude_wrap.rs +++ b/src/cli/claude_wrap.rs @@ -17,7 +17,8 @@ //! full, and every reporting error is swallowed exactly as the `sessions hook` //! sink swallows its own. //! -//! Nothing is logged or persisted. The observer reads only the state, the +//! Conversation content is never logged or persisted. An optional metadata-only +//! file sink is enabled by `OMNI_DEV_CLAUDE_WRAP_LOG`; it never runs on the byte pump. The observer reads only the state, the //! `session_id`, the `cwd` and the model out of the stream, and reports them to //! the daemon's existing `0600` Unix socket. @@ -42,6 +43,9 @@ use crate::sessions::stream::{Direction, StreamTracker}; /// The `sessions` service routing key on the daemon control socket. const SERVICE: &str = "sessions"; +mod diagnostics; +use diagnostics::{increment, Diagnostics}; + /// How long a fire-and-forget report waits for the daemon before giving up. /// Short, and on a task that no byte ever waits behind. const REPORT_TIMEOUT: Duration = Duration::from_secs(2); @@ -147,12 +151,13 @@ fn exec_replace(program: &str, args: &[String]) -> anyhow::Error { /// Wraps `program`, joining it to this process's own stdin and stdout. async fn wrap(program: &str, args: &[String], socket: Option) -> Result { - wrap_io( + wrap_io_diagnostics( program, args, socket, tokio::io::stdin(), tokio::io::stdout(), + std::env::var_os("OMNI_DEV_CLAUDE_WRAP_LOG").map(PathBuf::from), ) .await } @@ -164,6 +169,7 @@ async fn wrap(program: &str, args: &[String], socket: Option) -> Result /// Generic over the two endpoints so the wrapping can be tested against a real /// child process without touching the test runner's own stdio — which is also /// the only way to assert that the forwarding is byte-for-byte lossless. +#[cfg(test)] async fn wrap_io( program: &str, args: &[String], @@ -175,6 +181,23 @@ where R: AsyncRead + Unpin + Send + 'static, W: AsyncWrite + Unpin + Send + 'static, { + wrap_io_diagnostics(program, args, socket, input, output, None).await +} + +async fn wrap_io_diagnostics( + program: &str, + args: &[String], + socket: Option, + input: R, + output: W, + log_path: Option, +) -> Result +where + R: AsyncRead + Unpin + Send + 'static, + W: AsyncWrite + Unpin + Send + 'static, +{ + let (diagnostics, log_done) = Diagnostics::open(log_path.as_deref()); + diagnostics.record(|| json!({"event":"process_start"})); // stderr is inherited and env is left untouched, so the child sees exactly // the environment it would have without the wrapper. The child is // deliberately *not* put in its own process group (unlike the managed @@ -200,31 +223,34 @@ where let (tee, lines) = mpsc::channel::<(Direction, String)>(TEE_CAPACITY); let (title_tx, title_rx) = watch::channel::>(None); - let observer = tokio::spawn(observe( + let observer = tokio::spawn(observe_diagnostics( lines, socket, KEEPALIVE_INTERVAL, title_tx, child_pid, + diagnostics.clone(), )); let signals = tokio::spawn(forward_signals(child_pid)); // The title rewrite only ever applies to the FromClaude direction — it // rewrites what Claude asserts about itself, never what we send it. let from_child_title_rx = title_rewrite_enabled().then_some(title_rx); - let from_child = tokio::spawn(pump( + let from_child = tokio::spawn(pump_diagnostics( child_stdout, output, Direction::FromClaude, tee.clone(), from_child_title_rx, + diagnostics.clone(), )); - let to_child = tokio::spawn(pump( + let to_child = tokio::spawn(pump_diagnostics( input, child_stdin, Direction::ToClaude, tee.clone(), None, + diagnostics.clone(), )); drop(tee); @@ -244,7 +270,13 @@ where .await .context("failed to wait for the wrapped process")?; let _ = observer.await; - Ok(exit_code(status)) + let code = exit_code(status); + diagnostics.record(|| json!({"event":"process_exit", "pid":child_pid, "code":code})); + drop(diagnostics); + if let Some(done) = log_done { + let _ = tokio::time::timeout(Duration::from_millis(200), done).await; + } + Ok(code) } /// Copies `reader` to `writer` byte-for-byte, teeing complete lines to the @@ -275,12 +307,35 @@ where /// updates it, this branch simply never fires — it cannot hang the pump. /// /// [`try_send`]: tokio::sync::mpsc::Sender::try_send +#[cfg(test)] async fn pump( + reader: R, + writer: W, + direction: Direction, + tee: mpsc::Sender<(Direction, String)>, + title_rx: Option>>, +) where + R: AsyncRead + Unpin, + W: AsyncWrite + Unpin, +{ + pump_diagnostics( + reader, + writer, + direction, + tee, + title_rx, + Diagnostics::default(), + ) + .await; +} + +async fn pump_diagnostics( mut reader: R, mut writer: W, direction: Direction, tee: mpsc::Sender<(Direction, String)>, mut title_rx: Option>>, + diagnostics: Diagnostics, ) where R: AsyncRead + Unpin, W: AsyncWrite + Unpin, @@ -316,7 +371,7 @@ async fn pump( if write_result.is_err() || writer.flush().await.is_err() { break; } - tee_chunk(chunk, direction, &tee, &mut line); + tee_chunk_diagnostics(chunk, direction, &tee, &mut line, &diagnostics); } changed = title_changed => { match changed { @@ -570,18 +625,40 @@ fn title_prefix_for_model(model_id: &str) -> String { /// Reassembles newline-delimited lines out of a forwarded chunk and offers each /// to the observer, dropping any that is over-long, not UTF-8, or arrives while /// the channel is full. +#[cfg(test)] fn tee_chunk( chunk: &[u8], direction: Direction, tee: &mpsc::Sender<(Direction, String)>, line: &mut Vec, +) { + tee_chunk_diagnostics(chunk, direction, tee, line, &Diagnostics::default()); +} + +fn tee_chunk_diagnostics( + chunk: &[u8], + direction: Direction, + tee: &mpsc::Sender<(Direction, String)>, + line: &mut Vec, + diagnostics: &Diagnostics, ) { for &byte in chunk { if byte == b'\n' { if line.len() < MAX_LINE_BYTES { if let Ok(text) = std::str::from_utf8(line) { - let _ = tee.try_send((direction, text.to_string())); + if let Err(error) = tee.try_send((direction, text.to_string())) { + if let Some(c) = diagnostics.counts() { + match error { + mpsc::error::TrySendError::Full(_) => increment(&c.full), + mpsc::error::TrySendError::Closed(_) => increment(&c.closed), + } + } + } + } else if let Some(c) = diagnostics.counts() { + increment(&c.non_utf8); } + } else if let Some(c) = diagnostics.counts() { + increment(&c.oversize); } line.clear(); } else if line.len() < MAX_LINE_BYTES { @@ -600,18 +677,41 @@ fn tee_chunk( /// cannot end the resumed session (#1948). It is also the seed for pid-based /// liveness (#1916), whose identity-token reading is entirely the daemon's own /// [`pid_watcher`](crate::sessions::pid_watcher) work, never this wrapper's. +#[cfg(test)] async fn observe( + lines: mpsc::Receiver<(Direction, String)>, + socket: Option, + every: Duration, + title_tx: watch::Sender>, + child_pid: Option, +) { + observe_diagnostics( + lines, + socket, + every, + title_tx, + child_pid, + Diagnostics::default(), + ) + .await; +} + +async fn observe_diagnostics( mut lines: mpsc::Receiver<(Direction, String)>, socket: Option, every: Duration, title_tx: watch::Sender>, child_pid: Option, + diagnostics: Diagnostics, ) { let Ok(socket) = server::resolve_socket(socket) else { + diagnostics + .record(|| json!({"event":"observer_stopped", "outcome":"socket_resolution_failed"})); return; }; let mut tracker = StreamTracker::new(); let mut last_model: Option = None; + let mut last_identity = None; let mut keepalive = tokio::time::interval(every); // The first tick of a tokio interval completes immediately; consume it so // the keep-alive does not fire before anything has been observed. @@ -623,7 +723,10 @@ async fn observe( Some((direction, text)) => tracker.observe_line(direction, &text), None => break, }, - _ = keepalive.tick() => tracker.keepalive(), + _ = keepalive.tick() => { + diagnostics.summary(tracker.diagnostics()); + tracker.keepalive() + }, }; // Publish the classified title prefix whenever the model changes — // independent of `observed`, since a mid-session model switch does @@ -636,8 +739,30 @@ async fn observe( } if let Some(mut request) = observed { request.pid = child_pid; + if diagnostics.enabled() { + let identity = ( + request.session_id.clone(), + request.cwd.clone(), + request.model.clone(), + ); + if last_identity.as_ref() != Some(&identity) { + diagnostics.record(|| { + json!({"event":"session_identity", "session_id":identity.0, + "cwd":identity.1, "model":identity.2, "pid":child_pid}) + }); + last_identity = Some(identity); + } + } + diagnostics.record(|| { + json!({"event":"state_report", "session_id":request.session_id, + "state":request.event, "pid":child_pid}) + }); if let Ok(payload) = serde_json::to_value(request) { - report(&socket, "observe", payload).await; + report_diagnostics(&socket, "observe", payload, &diagnostics).await; + } else { + diagnostics.record( + || json!({"event":"report", "op":"observe", "outcome":"serialization_failed"}), + ); } } } @@ -647,15 +772,25 @@ async fn observe( if let Some(pid) = child_pid { payload["pid"] = Value::from(pid); } - report(&socket, "end", payload).await; + report_diagnostics(&socket, "end", payload, &diagnostics).await; } + diagnostics.summary(tracker.diagnostics()); } /// Sends one bounded, fire-and-forget op to the daemon's `sessions` service, /// swallowing every failure — a missing or wedged daemon must be a silent no-op. -async fn report(socket: &Path, op: &str, payload: Value) { +async fn report_diagnostics(socket: &Path, op: &str, payload: Value, diagnostics: &Diagnostics) { let envelope = DaemonEnvelope::service(SERVICE, op, payload); - let _ = tokio::time::timeout(REPORT_TIMEOUT, DaemonClient::new(socket).request(envelope)).await; + let outcome = + match tokio::time::timeout(REPORT_TIMEOUT, DaemonClient::new(socket).request(envelope)) + .await + { + Err(_) => "timeout", + Ok(Err(_)) => "transport_failed", + Ok(Ok(reply)) if !reply.ok => "daemon_rejected", + Ok(Ok(_)) => "delivered", + }; + diagnostics.record(|| json!({"event":"report", "op":op, "outcome":outcome})); } /// Relays `SIGINT`/`SIGTERM` to the child, so a caller that signals the wrapper @@ -701,6 +836,131 @@ mod tests { use super::*; + #[tokio::test] + async fn diagnostic_tee_counts_each_drop_reason() { + use std::sync::atomic::Ordering; + let dir = tempfile::tempdir().unwrap(); + let (diagnostics, done) = Diagnostics::open(Some(&dir.path().join("log"))); + let (tee, lines) = mpsc::channel(1); + let mut buffer = Vec::new(); + tee_chunk_diagnostics( + b"a\nb\n", + Direction::FromClaude, + &tee, + &mut buffer, + &diagnostics, + ); + drop(lines); + tee_chunk_diagnostics( + b"c\n", + Direction::FromClaude, + &tee, + &mut buffer, + &diagnostics, + ); + tee_chunk_diagnostics( + b"\xff\n", + Direction::FromClaude, + &tee, + &mut buffer, + &diagnostics, + ); + tee_chunk_diagnostics( + &vec![b'x'; MAX_LINE_BYTES], + Direction::FromClaude, + &tee, + &mut buffer, + &diagnostics, + ); + tee_chunk_diagnostics( + b"\n", + Direction::FromClaude, + &tee, + &mut buffer, + &diagnostics, + ); + let counts = diagnostics.counts().unwrap(); + assert_eq!(counts.full.load(Ordering::Relaxed), 1); + assert_eq!(counts.closed.load(Ordering::Relaxed), 1); + assert_eq!(counts.oversize.load(Ordering::Relaxed), 1); + assert_eq!(counts.non_utf8.load(Ordering::Relaxed), 1); + drop(diagnostics); + done.unwrap().await.unwrap(); + } + + #[tokio::test] + async fn enabled_diagnostics_preserve_bytes_and_exclude_conversation_content() { + let (dir, socket, _seen) = fake_daemon(); + let log = dir.path().join("wrapper.jsonl"); + let output = Sink::default(); + let stream = concat!( + r#"{"type":"system","subtype":"init","session_id":"private-test","cwd":"/project","model":"test"}"#, + "\n", + r#"{"type":"assistant","message":{"content":"CONVERSATION_SECRET"}}"#, + "\n", + r#"{"type":"control_request","request_id":"perm","request":{"subtype":"can_use_tool","input":"TOOL_SECRET"}}"#, + "\n", + r#"{"type":"result"}"#, + "\n" + ); + let code = wrap_io_diagnostics( + "/bin/cat", + &[], + Some(socket), + stream.as_bytes(), + output.clone(), + Some(log.clone()), + ) + .await + .unwrap(); + assert_eq!(code, 0); + assert_eq!(output.contents(), stream); + let text = std::fs::read_to_string(log).unwrap(); + assert!(text.contains("session_identity")); + assert!(text.contains("process_exit")); + assert!(text.contains("state_report")); + assert!(!text.contains("CONVERSATION_SECRET")); + assert!(!text.contains("TOOL_SECRET")); + for line in text.lines() { + serde_json::from_str::(line).unwrap(); + } + } + + #[tokio::test] + async fn diagnostic_report_rejection_does_not_persist_daemon_error_content() { + let dir = tempfile::tempdir().unwrap(); + let socket = dir.path().join("d.sock"); + let listener = UnixListener::bind(&socket).unwrap(); + let server = tokio::spawn(async move { + let (stream, _) = listener.accept().await.unwrap(); + let mut framed = + tokio_util::codec::Framed::new(stream, tokio_util::codec::LinesCodec::new()); + use futures::{SinkExt, StreamExt}; + framed.next().await.unwrap().unwrap(); + framed + .send(r#"{"ok":false,"error":"DAEMON_CONTENT_SECRET"}"#.to_string()) + .await + .unwrap(); + }); + let path = dir.path().join("log"); + let (diagnostics, done) = Diagnostics::open(Some(&path)); + report_diagnostics(&socket, "observe", json!({}), &diagnostics).await; + server.await.unwrap(); + report_diagnostics( + &dir.path().join("missing.sock"), + "observe", + json!({}), + &diagnostics, + ) + .await; + drop(diagnostics); + done.unwrap().await.unwrap(); + let text = std::fs::read_to_string(path).unwrap(); + assert!(text.contains("daemon_rejected")); + assert!(text.contains("transport_failed")); + assert!(!text.contains("DAEMON_CONTENT_SECRET")); + } + /// An [`AsyncWrite`] that appends into a shared buffer, so a test can assert /// on exactly the bytes the wrapper forwarded. #[derive(Clone, Default)] diff --git a/src/cli/claude_wrap/diagnostics.rs b/src/cli/claude_wrap/diagnostics.rs new file mode 100644 index 000000000..1bf71ddc5 --- /dev/null +++ b/src/cli/claude_wrap/diagnostics.rs @@ -0,0 +1,182 @@ +//! Opt-in metadata logging, isolated from forwarding and global tracing. +use std::io::Write; +use std::os::unix::fs::OpenOptionsExt; +use std::path::Path; +use std::sync::atomic::{AtomicU64, Ordering}; +use std::sync::{mpsc, Arc}; + +use serde_json::{json, Value}; + +const CAPACITY: usize = 128; + +#[derive(Default)] +pub(super) struct Counts { + pub full: AtomicU64, + pub closed: AtomicU64, + pub oversize: AtomicU64, + pub non_utf8: AtomicU64, + pub diagnostic_drops: AtomicU64, +} + +/// Cloning is cheap; a disabled sink has no channel, worker, or counters. +#[derive(Clone, Default)] +pub(super) struct Diagnostics { + sender: Option>, + counts: Option>, +} + +impl Diagnostics { + pub fn open(path: Option<&Path>) -> (Self, Option>) { + let Some(path) = path.filter(|p| !p.as_os_str().is_empty()) else { + return (Self::default(), None); + }; + let path = path.to_path_buf(); + let (sender, receiver) = mpsc::sync_channel::(CAPACITY); + let (done_tx, done_rx) = tokio::sync::oneshot::channel(); + // All disk work lives on this thread, including open. Even a stuck + // filesystem cannot hold up byte forwarding or process launch. + let worker = std::thread::Builder::new() + .name("claude-wrap-log".into()) + .spawn(move || { + let file = std::fs::OpenOptions::new() + .create(true) + .append(true) + .mode(0o600) + .custom_flags(nix::libc::O_NOFOLLOW | nix::libc::O_NONBLOCK) + .open(path); + if let Ok(mut file) = file { + // Refuse devices/FIFOs, and never follow a symlink log target. + if file.metadata().is_ok_and(|m| m.is_file()) { + for record in receiver { + if serde_json::to_writer(&mut file, &record).is_err() + || file.write_all(b"\n").is_err() + { + break; + } + } + let _ = file.flush(); + } + } + let _ = done_tx.send(()); + }); + if worker.is_err() { + return (Self::default(), None); + } + ( + Self { + sender: Some(sender), + counts: Some(Arc::new(Counts::default())), + }, + Some(done_rx), + ) + } + + pub fn enabled(&self) -> bool { + self.sender.as_ref().is_some_and(|_| self.counts.is_some()) + } + + pub fn counts(&self) -> Option<&Counts> { + self.counts.as_deref() + } + + /// Only call with allowlisted metadata; never stream lines or daemon errors. + pub fn record(&self, record: impl FnOnce() -> Value) { + if let Some(sender) = &self.sender { + if sender.try_send(record()).is_err() { + if let Some(counts) = &self.counts { + increment(&counts.diagnostic_drops); + } + } + } + } + + pub fn summary(&self, tracker: crate::sessions::stream::StreamDiagnostics) { + self.record(|| { + let load = |counter: &AtomicU64| counter.load(Ordering::Relaxed); + let c = self.counts.as_deref(); + json!({"event":"diagnostic_summary", "tracker":tracker, + "tee_full":c.map(|c| load(&c.full)), "tee_closed":c.map(|c| load(&c.closed)), + "tee_oversize":c.map(|c| load(&c.oversize)), "tee_non_utf8":c.map(|c| load(&c.non_utf8)), + "diagnostic_drops":c.map(|c| load(&c.diagnostic_drops))}) + }); + } +} + +pub(super) fn increment(counter: &AtomicU64) { + let _ = counter.fetch_update(Ordering::Relaxed, Ordering::Relaxed, |n| { + Some(n.saturating_add(1)) + }); +} + +#[cfg(test)] +#[allow(clippy::unwrap_used)] +mod tests { + use super::*; + + #[test] + fn diagnostic_queue_saturation_is_bounded_and_counted() { + let (sender, _receiver) = mpsc::sync_channel(1); + let sink = Diagnostics { + sender: Some(sender), + counts: Some(Arc::new(Counts::default())), + }; + for _ in 0..100 { + sink.record(|| json!({"event":"test"})); + } + assert_eq!( + sink.counts() + .unwrap() + .diagnostic_drops + .load(Ordering::Relaxed), + 99 + ); + } + + #[tokio::test] + async fn disabled_sink_does_not_evaluate_records() { + let (sink, done) = Diagnostics::open(None); + assert!(!sink.enabled()); + assert!(done.is_none()); + sink.record(|| panic!("disabled sink formatted a record")); + } + + #[tokio::test] + async fn sink_writes_metadata_with_private_permissions() { + use std::os::unix::fs::PermissionsExt; + let dir = tempfile::tempdir().unwrap(); + let path = dir.path().join("metadata.jsonl"); + let (sink, done) = Diagnostics::open(Some(&path)); + sink.record(|| json!({"event":"process_exit", "code":0})); + drop(sink); + tokio::time::timeout(std::time::Duration::from_secs(2), done.unwrap()) + .await + .unwrap() + .unwrap(); + let text = std::fs::read_to_string(&path).unwrap(); + assert!(text.contains("process_exit")); + assert_eq!( + std::fs::metadata(&path).unwrap().permissions().mode() & 0o777, + 0o600 + ); + } + + #[tokio::test] + async fn failed_sink_is_nonblocking() { + let dir = tempfile::tempdir().unwrap(); + let (sink, done) = Diagnostics::open(Some(dir.path())); + tokio::time::timeout(std::time::Duration::from_secs(2), done.unwrap()) + .await + .unwrap() + .unwrap(); + for _ in 0..1000 { + sink.record(|| json!({"event":"test"})); + } + assert_eq!( + sink.counts() + .unwrap() + .diagnostic_drops + .load(Ordering::Relaxed), + 1000 + ); + } +} diff --git a/src/cli/sessions.rs b/src/cli/sessions.rs index be30842fd..9f44dcdf8 100644 --- a/src/cli/sessions.rs +++ b/src/cli/sessions.rs @@ -242,6 +242,7 @@ impl HookCommand { let pid = agent_pid(std::os::unix::process::parent_id()); let mut input = String::new(); if std::io::stdin().read_to_string(&mut input).is_err() { + tracing::debug!(agent = ?self.agent, outcome = "stdin_read_failed", "session_hook_skipped"); return Ok(()); } self.report(&input, pid).await; @@ -252,18 +253,26 @@ impl HookCommand { /// out so tests can exercise the send path against a fake socket. async fn report(&self, input: &str, pid: Option) { let Ok(hook) = serde_json::from_str::(input) else { + tracing::debug!(agent = ?self.agent, outcome = "invalid_json", "session_hook_skipped"); return; }; let Some((op, payload)) = hook.to_op(self.agent, pid) else { return; }; let Ok(socket) = server::resolve_socket(self.socket.clone()) else { + tracing::debug!(agent = ?self.agent, outcome = "socket_resolution_failed", "session_hook_skipped"); return; }; - // Bounded, and every failure ignored: the daemon may be down, and that - // must be a silent no-op. let env = DaemonEnvelope::service(SERVICE, op, payload); - let _ = tokio::time::timeout(HOOK_TIMEOUT, DaemonClient::new(&socket).request(env)).await; + let outcome = + match tokio::time::timeout(HOOK_TIMEOUT, DaemonClient::new(&socket).request(env)).await + { + Err(_) => "timeout", + Ok(Err(_)) => "transport_failed", + Ok(Ok(reply)) if !reply.ok => "daemon_rejected", + Ok(Ok(_)) => "delivered", + }; + tracing::debug!(agent = ?self.agent, session_id = ?hook.session_id, ?pid, op, outcome, "session_hook_report"); } } @@ -326,8 +335,18 @@ impl HookPayload { /// entirely the daemon's own [`pid_watcher`](crate::sessions::pid_watcher) /// work, never this sink's. fn to_op(&self, agent: HookAgent, pid: Option) -> Option<(&'static str, Value)> { - let session_id = self.session_id.clone().filter(|s| !s.trim().is_empty())?; - let event_name = self.hook_event_name.as_deref()?; + let Some(session_id) = self.session_id.clone().filter(|s| !s.trim().is_empty()) else { + tracing::debug!( + ?agent, + outcome = "missing_session_id", + "session_hook_skipped" + ); + return None; + }; + let Some(event_name) = self.hook_event_name.as_deref() else { + tracing::debug!(?agent, outcome = "missing_event", "session_hook_skipped"); + return None; + }; if event_name == "SessionEnd" { let mut payload = json!({ "session_id": session_id }); if let Some(reason) = self.reason.as_ref().or(self.message.as_ref()) { @@ -338,18 +357,22 @@ impl HookPayload { } return Some(("end", payload)); } - let mut event = match agent { + let mapped = match agent { HookAgent::Claude => session_event_for( event_name, self.source.as_deref(), self.notification_type.as_deref(), self.message.as_deref(), - )?, + ), HookAgent::Codex => codex_session_event_for( event_name, self.source.as_deref(), self.tool_name.as_deref(), - )?, + ), + }; + let Some(mut event) = mapped else { + tracing::debug!(?agent, outcome = "unmapped_event", "session_hook_skipped"); + return None; }; let agent_id = (agent == HookAgent::Claude) .then(|| self.agent_id.clone()) @@ -369,7 +392,17 @@ impl HookPayload { model: self.model.clone(), pid, }; - Some(("observe", serde_json::to_value(request).ok()?)) + match serde_json::to_value(request) { + Ok(payload) => Some(("observe", payload)), + Err(_) => { + tracing::debug!( + ?agent, + outcome = "serialization_failed", + "session_hook_skipped" + ); + None + } + } } } @@ -430,7 +463,14 @@ fn session_event_for( "PermissionRequest" => SessionEvent::Notification(NotificationKind::PermissionPrompt), "Elicitation" => SessionEvent::Notification(NotificationKind::AgentNeedsInput), "Notification" => { - SessionEvent::Notification(classify_notification(notification_type, message)) + let kind = classify_notification(notification_type, message); + tracing::debug!( + ?kind, + has_notification_type = notification_type.is_some(), + has_message = message.is_some(), + "session_notification_classified" + ); + SessionEvent::Notification(kind) } _ => return None, }) @@ -1739,6 +1779,34 @@ fn age_secs(ts: Option<&str>) -> i64 { #[cfg(test)] #[allow(clippy::unwrap_used, clippy::expect_used)] mod tests { + #[tokio::test] + async fn hook_diagnostics_identify_gates_without_logging_payload_content() { + let dir = tempfile::tempdir().unwrap(); + let command = HookCommand { + socket: Some(dir.path().join("missing.sock")), + agent: HookAgent::Claude, + }; + let (_, logs) = crate::test_support::capture_future_at(tracing::Level::DEBUG, async { + for input in ["not json HOOK_CONTENT_SECRET", "{}", + r#"{"session_id":"test"}"#, + r#"{"session_id":"test","hook_event_name":"future-event"}"#, + r#"{"session_id":"test","hook_event_name":"Notification","message":"HOOK_CONTENT_SECRET"}"#] { + command.report(input, None).await; + } + }).await; + for reason in [ + "invalid_json", + "missing_session_id", + "missing_event", + "unmapped_event", + "transport_failed", + ] { + assert!(logs.contains(reason), "missing {reason}: {logs}"); + } + assert!(logs.contains("session_notification_classified")); + assert!(!logs.contains("HOOK_CONTENT_SECRET")); + } + use super::*; use crate::sessions::SessionState; diff --git a/src/daemon/server.rs b/src/daemon/server.rs index 30c766a36..14c4c2d35 100644 --- a/src/daemon/server.rs +++ b/src/daemon/server.rs @@ -332,7 +332,7 @@ async fn handle_connection( if let Some(name) = envelope.service.as_deref() { if name != DAEMON_SERVICE { if let Some(stream) = registry.subscribe(name, &envelope.op, &envelope.payload) { - run_stream(&mut framed, stream, &shutdown).await; + run_stream(&mut framed, stream, &shutdown, name, &envelope.op).await; return; } } @@ -355,46 +355,91 @@ async fn handle_connection( /// The subscription owns the connection for its lifetime: any further inbound /// line is treated as an explicit cancel and ends the stream, matching the /// one-op-per-connection the companion uses (a dedicated subscribe socket). +/// Reasons a subscription ends; no incoming line content is retained. +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +enum StreamEnd { + ClientCancel, + ClientEof, + ReadDecodeError, + DaemonShutdown, + WriteFailure, +} + +fn inbound_end(line: Option>) -> StreamEnd { + match line { + Some(Ok(_)) => StreamEnd::ClientCancel, + None => StreamEnd::ClientEof, + Some(Err(_)) => StreamEnd::ReadDecodeError, + } +} + async fn run_stream( framed: &mut Framed, mut stream: Box, shutdown: &CancellationToken, -) { - // Initial snapshot up front. The stream's change source was captured when it - // was built (before this snapshot), so the loop below only pushes deltas — - // and any change racing this initial sample is caught by the first wakeup. + service: &str, + op: &str, +) -> StreamEnd { + tracing::debug!(service, op, "subscription_started"); + let mut pushed = 0_u64; + let mut suppressed = 0_u64; let mut last = stream.snapshot().await; - if !send_reply(framed, DaemonReply::ok(last.clone())).await { - return; - } - - // `interval` fires immediately on the first `tick()`; consume that so the - // periodic re-sample starts one full interval out. - let mut tick = tokio::time::interval(stream_tick()); - tick.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip); - tick.tick().await; - - loop { - tokio::select! { - () = stream.changed() => {} - _ = tick.tick() => {} - // Reading `framed` serves double duty and every outcome ends the - // stream: an inbound line is an explicit cancel, `None` is the client - // hanging up, and an `Err` is a read/decode error. `Framed`'s decode - // buffer lives in the codec, not this future, so cancelling this arm - // mid-poll loses no buffered bytes. - _ = framed.next() => break, - () = shutdown.cancelled() => break, - } - // Any wakeup means "maybe changed": re-sample and push only a real delta. - let snap = stream.snapshot().await; - if snap != last { - if !send_reply(framed, DaemonReply::ok(snap.clone())).await { - break; + let end = if !send_reply(framed, DaemonReply::ok(last.clone())).await { + StreamEnd::WriteFailure + } else { + pushed += 1; + tracing::trace!( + service, + op, + trigger = "initial", + outcome = "pushed", + "subscription_sample" + ); + let mut tick = tokio::time::interval(stream_tick()); + tick.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip); + tick.tick().await; + loop { + let trigger = tokio::select! { + () = stream.changed() => "change_notification", + _ = tick.tick() => "periodic_tick", + line = framed.next() => { + // A line is cancel, EOF is hangup, and a codec error is a + // broken read. Never log the line (it may contain content). + if let Some(Err(error)) = &line { + tracing::debug!(service, op, %error, "subscription_read_failed"); + } + break inbound_end(line); + }, + () = shutdown.cancelled() => break StreamEnd::DaemonShutdown, + }; + let snap = stream.snapshot().await; + if snap != last { + if !send_reply(framed, DaemonReply::ok(snap.clone())).await { + break StreamEnd::WriteFailure; + } + pushed = pushed.saturating_add(1); + last = snap; + tracing::trace!( + service, + op, + trigger, + outcome = "pushed", + "subscription_sample" + ); + } else { + suppressed = suppressed.saturating_add(1); + tracing::trace!( + service, + op, + trigger, + outcome = "suppressed_identical", + "subscription_sample" + ); } - last = snap; } - } + }; + tracing::debug!(service, op, outcome = ?end, pushed, suppressed, "subscription_ended"); + end } /// Encodes and writes one reply line. Returns `false` when the connection @@ -408,7 +453,7 @@ async fn send_reply(framed: &mut Framed, reply: DaemonRe } }; if let Err(e) = framed.send(encoded).await { - tracing::debug!("daemon client write failed: {e}"); + tracing::warn!(error = %e, "daemon_client_write_failed"); return false; } true @@ -423,7 +468,13 @@ async fn dispatch_envelope( shutdown: &CancellationToken, ) -> DaemonReply { match envelope.service.as_deref() { - None | Some(DAEMON_SERVICE) => handle_builtin(&envelope.op, registry, shutdown).await, + None | Some(DAEMON_SERVICE) => { + let reply = handle_builtin(&envelope.op, registry, shutdown).await; + if !reply.ok { + tracing::warn!(service = DAEMON_SERVICE, op = %envelope.op, error = ?reply.error, "daemon_dispatch_failed"); + } + reply + } Some(name) => { // Correlate any HTTP the service issues to the originating client's // invocation, when it threaded its id across the socket (#1198). @@ -439,7 +490,10 @@ async fn dispatch_envelope( // query failed: snowflake server error (000630): …") so the // client can see the underlying cause, not just the top-level // wrapper. - Err(e) => DaemonReply::err(format!("{e:#}")), + Err(e) => { + tracing::warn!(service = name, op = %envelope.op, error = %format!("{e:#}"), "daemon_dispatch_failed"); + DaemonReply::err(format!("{e:#}")) + } } } } @@ -502,6 +556,86 @@ pub fn resolve_socket(socket: Option) -> Result { #[cfg(test)] #[allow(clippy::unwrap_used, clippy::expect_used)] mod tests { + + #[tokio::test] + async fn subscriptions_report_all_terminal_causes() { + use tokio::io::AsyncWriteExt; + for expected in [ + StreamEnd::ClientCancel, + StreamEnd::ClientEof, + StreamEnd::ReadDecodeError, + StreamEnd::DaemonShutdown, + StreamEnd::WriteFailure, + ] { + let (mut client, server) = UnixStream::pair().unwrap(); + let (_tx, rx) = watch::channel(0_u64); + let fake = FakeStream { + rx, + snap: Arc::new(StdMutex::new(json!({"n":0}))), + }; + let token = CancellationToken::new(); + let shutdown = token.clone(); + let task = tokio::spawn(async move { + let mut framed = Framed::new(server, LinesCodec::new_with_max_length(64)); + crate::test_support::capture_future_at( + tracing::Level::TRACE, + run_stream( + &mut framed, + Box::new(fake), + &shutdown, + "sessions", + "subscribe", + ), + ) + .await + }); + if expected == StreamEnd::WriteFailure { + drop(client); + } else { + let mut reader = BufReader::new(&mut client); + let _ = read_reply(&mut reader).await; + drop(reader); + match expected { + StreamEnd::ClientCancel => { + client.write_all(b"cancel\n").await.unwrap(); + } + StreamEnd::ClientEof => { + client.shutdown().await.unwrap(); + } + StreamEnd::ReadDecodeError => { + client.write_all(&[b'x'; 65]).await.unwrap(); + } + StreamEnd::DaemonShutdown => token.cancel(), + StreamEnd::WriteFailure => unreachable!(), + } + } + let (cause, logs) = tokio::time::timeout(Duration::from_secs(2), task) + .await + .unwrap() + .unwrap(); + assert_eq!(cause, expected); + assert!(logs.contains("subscription_started")); + assert!(logs.contains("subscription_ended")); + assert!(logs.contains(&format!("outcome={expected:?}"))); + if expected != StreamEnd::WriteFailure { + assert!(logs.contains("pushed=1")); + } + } + } + + #[test] + fn inbound_subscription_end_reasons_are_distinct() { + assert_eq!( + super::inbound_end(Some(Ok("cancel".into()))), + super::StreamEnd::ClientCancel + ); + assert_eq!(super::inbound_end(None), super::StreamEnd::ClientEof); + assert_eq!( + super::inbound_end(Some(Err(LinesCodecError::MaxLineLengthExceeded))), + super::StreamEnd::ReadDecodeError + ); + } + use super::*; #[test] @@ -650,7 +784,14 @@ mod tests { let server_task = tokio::spawn(async move { let mut framed = Framed::new(server, LinesCodec::new_with_max_length(MAX_LINE_BYTES)); - run_stream(&mut framed, Box::new(fake), &server_shutdown).await; + run_stream( + &mut framed, + Box::new(fake), + &server_shutdown, + "fake", + "subscribe", + ) + .await; }); let mut reader = BufReader::new(client); @@ -690,7 +831,14 @@ mod tests { let server_task = tokio::spawn(async move { let mut framed = Framed::new(server, LinesCodec::new_with_max_length(MAX_LINE_BYTES)); - run_stream(&mut framed, Box::new(fake), &server_shutdown).await; + run_stream( + &mut framed, + Box::new(fake), + &server_shutdown, + "fake", + "subscribe", + ) + .await; }); let mut reader = BufReader::new(&mut client); @@ -814,7 +962,7 @@ mod tests { let mut framed = Framed::new(server, LinesCodec::new_with_max_length(MAX_LINE_BYTES)); tokio::time::timeout( Duration::from_secs(2), - run_stream(&mut framed, Box::new(fake), &shutdown), + run_stream(&mut framed, Box::new(fake), &shutdown, "fake", "subscribe"), ) .await .expect("run_stream should return promptly when the initial send fails"); diff --git a/src/daemon/services/sessions.rs b/src/daemon/services/sessions.rs index 20d4d9705..c2eac9a65 100644 --- a/src/daemon/services/sessions.rs +++ b/src/daemon/services/sessions.rs @@ -134,9 +134,15 @@ impl DaemonService for SessionsService { // a caller already supplied `repo`. if req.repo.is_none() { if let Some(cwd) = req.cwd.clone() { - req.repo = tokio::task::spawn_blocking(move || repo_name_for(&cwd)) + req.repo = match tokio::task::spawn_blocking(move || repo_name_for(&cwd)) .await - .unwrap_or_default(); + { + Ok(repo) => repo, + Err(error) => { + tracing::warn!(session_id = %req.session_id, %error, "sessions_repo_enrichment_failed"); + None + } + }; } } self.registry.observe(req); diff --git a/src/main.rs b/src/main.rs index 2f6f74858..c1e7ce20c 100644 --- a/src/main.rs +++ b/src/main.rs @@ -124,13 +124,56 @@ fn default_filter(daemon_run: bool) -> &'static str { /// daemon/debug logs off stdout. The default level when `RUST_LOG` is unset is /// [`default_filter`]; `RUST_LOG` still overrides it. fn init_tracing(daemon_run: bool) { + // Loading through the normal warning loader here would lose its warning: + // there is no subscriber yet. Emit once after installation instead. + let (daemon_filter, settings_error) = if daemon_run { + match omni_dev::utils::settings::Settings::load() { + Ok(settings) => (settings.daemon.log_level, None), + Err(error) => (None, Some(format!("{error:#}"))), + } + } else { + (None, None) + }; + let env_filter = std::env::var("RUST_LOG").ok(); + let (filter, rejected) = + resolve_filter(daemon_run, env_filter.as_deref(), daemon_filter.as_deref()); tracing_subscriber::fmt() .with_writer(std::io::stderr) - .with_env_filter( - tracing_subscriber::EnvFilter::try_from_default_env() - .unwrap_or_else(|_| tracing_subscriber::EnvFilter::new(default_filter(daemon_run))), - ) + .with_env_filter(filter) .init(); + for source in rejected { + tracing::warn!( + source, + "invalid tracing directive; using next filter source" + ); + } + if let Some(error) = settings_error { + omni_dev::utils::settings::Settings::warn_bootstrap_failure(&error); + } +} + +/// Resolves injected directives without a global subscriber or environment mutation. +fn resolve_filter( + daemon_run: bool, + environment: Option<&str>, + daemon: Option<&str>, +) -> (tracing_subscriber::EnvFilter, Vec<&'static str>) { + let mut rejected = Vec::new(); + for (source, directive) in [ + ("RUST_LOG", environment), + ("daemon.log_level", daemon.filter(|_| daemon_run)), + ] { + if let Some(directive) = directive { + match tracing_subscriber::EnvFilter::try_new(directive) { + Ok(filter) => return (filter, rejected), + Err(_) => rejected.push(source), + } + } + } + ( + tracing_subscriber::EnvFilter::new(default_filter(daemon_run)), + rejected, + ) } /// Prints an error and its source chain to stderr, then exits non-zero. @@ -171,6 +214,23 @@ mod tests { assert!(!is_daemon_run(&[])); } + #[test] + fn tracing_filter_precedence_and_fallbacks() { + for (daemon_run, env, daemon, expected, rejected) in [ + (true, Some("trace"), Some("debug"), "trace", 0), + (true, None, Some("debug"), "debug", 0), + (true, None, None, "info", 0), + (false, None, Some("debug"), "warn", 0), + (true, Some("["), Some("debug"), "debug", 1), + (true, None, Some("["), "info", 1), + (false, Some("["), Some("debug"), "warn", 1), + ] { + let (filter, errors) = resolve_filter(daemon_run, env, daemon); + assert_eq!(filter.to_string(), expected); + assert_eq!(errors.len(), rejected); + } + } + #[test] fn default_filter_is_info_only_for_daemon_run() { assert_eq!(default_filter(true), "info"); diff --git a/src/sessions.rs b/src/sessions.rs index 9b0003956..87f4e3557 100644 --- a/src/sessions.rs +++ b/src/sessions.rs @@ -695,10 +695,16 @@ impl SessionsRegistry { .agent_id .as_deref() .filter(|id| req.agent == Agent::Claude && !id.trim().is_empty()); + let session_id = req.session_id.clone(); + let event = req.event; + let agent = req.agent; + let pid = req.pid; let now = Utc::now(); - let changed = { + let (changed, old_state, new_state, outcome, reaped) = { let mut sessions = self.lock_sessions(); let reaped = reap_sessions(&mut sessions, self.session_ttl, self.ended_ttl, now); + let old_state = sessions.get(&session_id).map(|entry| entry.state); + let mut outcome = "created"; let mutated = match sessions.get_mut(&req.session_id) { // A passive re-sighting (the Codex rollout watcher's heartbeat) // must not refresh an ended session, or it would outlive its @@ -708,6 +714,7 @@ impl SessionsRegistry { && (req.event == SessionEvent::TranscriptDiscovered || agent_id.is_some()) => { + outcome = "ended_passive_ignored"; false } // A straggling sighting from a process this session's resume @@ -720,6 +727,7 @@ impl SessionsRegistry { .pid .is_some_and(|pid| entry.replaced_pids.contains(&pid)) => { + outcome = "replaced_pid_ignored"; false } Some(entry) => { @@ -737,7 +745,16 @@ impl SessionsRegistry { let filled_transcript = fill(&mut entry.transcript_path, req.transcript_path); let filled_repo = fill(&mut entry.repo, req.repo); let filled_model = fill(&mut entry.model, req.model); - state_changed || filled_cwd || filled_transcript || filled_repo || filled_model + let metadata_changed = + filled_cwd || filled_transcript || filled_repo || filled_model; + outcome = if state_changed { + "state_changed" + } else if metadata_changed { + "metadata_enriched" + } else { + "heartbeat_only" + }; + state_changed || metadata_changed } None => { if sessions.len() >= MAX_SESSIONS { @@ -776,11 +793,19 @@ impl SessionsRegistry { true } }; - mutated || reaped > 0 + ( + mutated || reaped > 0, + old_state, + sessions.get(&session_id).map(|entry| entry.state), + outcome, + reaped, + ) }; if changed { self.bump(); } + tracing::debug!(%session_id, ?agent, ?pid, ?agent_id, ?event, ?old_state, ?new_state, outcome, + reaped, bumped = changed, "session_observed"); } /// Marks a session ended (`SessionEnd`), so `list` shows it as `ended` for a @@ -798,19 +823,26 @@ impl SessionsRegistry { /// longer owns. pub fn end(&self, session_id: &str, _reason: Option<&str>, pid: Option) -> bool { let now = Utc::now(); - let (known, reaped) = { + let (known, reaped, old_state, outcome) = { let mut sessions = self.lock_sessions(); let reaped = reap_sessions(&mut sessions, self.session_ttl, self.ended_ttl, now); + let old_state = sessions.get(session_id).map(|entry| entry.state); + let mut outcome = "unknown"; let known = match sessions.get_mut(session_id) { // Already ended (a hook and a watcher can both end it): leave the // linger window alone rather than restart it. - Some(entry) if entry.state == SessionState::Ended => (true, false), + Some(entry) if entry.state == SessionState::Ended => { + outcome = "already_ended"; + (true, false) + } // The replaced process's late `SessionEnd`: the resumed session // lives on under its new process. Some(entry) if pid.is_some_and(|pid| entry.replaced_pids.contains(&pid)) => { + outcome = "replaced_pid_ignored"; (true, false) } Some(entry) => { + outcome = "ended"; entry.state = SessionState::Ended; entry.last_event = SessionEvent::Stop; entry.last_seen = now; @@ -818,7 +850,7 @@ impl SessionsRegistry { } None => (false, false), }; - (known, reaped) + (known, reaped, old_state, outcome) }; let (known, flipped) = known; // A known session flipped to `ended`; otherwise only this call's inline @@ -826,6 +858,15 @@ impl SessionsRegistry { if flipped || reaped > 0 { self.bump(); } + tracing::debug!( + session_id, + ?pid, + ?old_state, + outcome, + reaped, + bumped = flipped || reaped > 0, + "session_end" + ); known } @@ -884,10 +925,17 @@ impl SessionsRegistry { /// window sends, which would otherwise put a permanent push floor under the /// daemon proportional to the window count. pub fn report_window(&self, report: WindowReport) { + let window_key = report.key.clone(); + let folder_count = report.folders.len(); let now = Utc::now(); - let changed = { + let (changed, outcome, reaped) = { let mut windows = self.lock_windows(); let reaped = reap_windows(&mut windows, self.window_ttl, now); + let outcome = if windows.contains_key(&report.key) { + "refresh" + } else { + "registered" + }; let (mutated, registered_at) = if let Some(previous) = windows.get(&report.key) { ( previous.report.folders != report.folders @@ -908,11 +956,20 @@ impl SessionsRegistry { registered_at, }, ); - mutated || reaped > 0 + ( + mutated || reaped > 0, + if outcome == "refresh" && mutated { + "embedding_changed" + } else { + outcome + }, + reaped, + ) }; if changed { self.bump(); } + tracing::debug!(%window_key, folder_count, outcome, reaped, bumped = changed, "session_window_reported"); } /// Drops a companion window-embedding report (the window closed). Returns @@ -925,6 +982,12 @@ impl SessionsRegistry { if removed { self.bump(); } + tracing::debug!( + window_key = key, + removed, + bumped = removed, + "session_window_unregistered" + ); removed } @@ -1057,13 +1120,20 @@ fn fill(slot: &mut Option, incoming: Option) -> bool { /// [`Source::Terminal`]. fn resolve_source(cwd: Option<&Path>, windows: &[WindowEntry]) -> Source { let Some(cwd) = cwd else { + tracing::trace!(outcome = "missing_cwd", "session_attribution_miss"); return Source::Terminal; }; match pick_window(cwd, windows) { Some(window) => Source::VsCode { window_key: window.key.clone(), }, - None => Source::Terminal, + None => { + tracing::debug!(cwd = %cwd.display(), window_count = windows.len(), "session_attribution_miss"); + for window in windows { + tracing::trace!(window_key = %window.report.key, folders = ?window.report.folders, "session_attribution_candidate"); + } + Source::Terminal + } } } @@ -1142,9 +1212,17 @@ fn reap_sessions( } else { session_max }; - (now - e.last_seen).num_seconds() <= max_age + let keep = (now - e.last_seen).num_seconds() <= max_age; + if !keep { + tracing::trace!(session_id = %e.session_id, reason = if e.state == SessionState::Ended { "ended_ttl" } else { "session_ttl" }, "session_reaped"); + } + keep }); - before - sessions.len() + let count = before - sessions.len(); + if count > 0 { + tracing::debug!(count, "sessions_reaped"); + } + count } /// Removes window-embedding reports last refreshed longer than `ttl` ago. @@ -1155,8 +1233,18 @@ fn reap_windows( ) -> usize { let max_age = ttl.as_secs() as i64; let before = windows.len(); - windows.retain(|_, e| (now - e.last_seen).num_seconds() <= max_age); - before - windows.len() + windows.retain(|key, e| { + let keep = (now - e.last_seen).num_seconds() <= max_age; + if !keep { + tracing::trace!(window_key = %key, reason = "window_ttl", "session_window_reaped"); + } + keep + }); + let count = before - windows.len(); + if count > 0 { + tracing::debug!(count, "session_windows_reaped"); + } + count } /// Removes the session with the oldest `last_seen` (ties broken by lowest @@ -1173,6 +1261,7 @@ fn evict_oldest_session(sessions: &mut HashMap) { .map(|e| e.session_id.clone()); if let Some(key) = oldest { sessions.remove(&key); + tracing::debug!(session_id = %key, reason = "capacity", "session_evicted"); } } @@ -1185,6 +1274,7 @@ fn evict_oldest_window(windows: &mut HashMap) { .map(|(k, _)| k.clone()); if let Some(key) = oldest { windows.remove(&key); + tracing::debug!(window_key = %key, reason = "capacity", "session_window_evicted"); } } @@ -1193,6 +1283,29 @@ fn evict_oldest_window(windows: &mut HashMap) { mod tests { use super::*; + #[test] + fn registry_diagnostics_distinguish_transition_heartbeat_and_unknown_end() { + let registry = SessionsRegistry::new(); + let logs = crate::test_support::capture_at(tracing::Level::DEBUG, || { + registry.observe(observe_request( + "diagnostics", + SessionEvent::UserPromptSubmit, + Some("/project"), + )); + registry.observe(observe_request( + "diagnostics", + SessionEvent::UserPromptSubmit, + Some("/project"), + )); + registry.end("missing", None, None); + }); + assert!(logs.contains("state_changed") || logs.contains("created")); + assert!(logs.contains("heartbeat_only")); + assert!(logs.contains("unknown")); + assert!(logs.contains("bumped=true")); + assert!(logs.contains("bumped=false")); + } + fn observe_request(session_id: &str, event: SessionEvent, cwd: Option<&str>) -> ObserveRequest { ObserveRequest { agent_id: None, diff --git a/src/sessions/pid_watcher.rs b/src/sessions/pid_watcher.rs index d72a73d8d..9c21c1879 100644 --- a/src/sessions/pid_watcher.rs +++ b/src/sessions/pid_watcher.rs @@ -113,6 +113,7 @@ impl Action { fn apply(self, registry: &SessionsRegistry, now: DateTime) { match self { Self::End { session_id } => { + tracing::debug!(%session_id, reason = "pid_liveness", "session_process_ended"); registry.end(&session_id, Some("pid liveness watcher"), None); } Self::Confirm { diff --git a/src/sessions/stream.rs b/src/sessions/stream.rs index ba6336a98..ee9ebb819 100644 --- a/src/sessions/stream.rs +++ b/src/sessions/stream.rs @@ -114,6 +114,19 @@ struct ControlBody { model: Option, } +/// Content-free cumulative protocol diagnostics for one wrapped process. +#[derive(Debug, Default, Clone, Copy, PartialEq, Eq, serde::Serialize)] +pub struct StreamDiagnostics { + /// Nonempty lines that could not be parsed after terminal-prefix recovery. + pub parse_failures: u64, + /// Control requests whose subtype/body is missing or unknown. + pub unknown_controls: u64, + /// Permission prompts without a usable correlation id. + pub missing_permission_ids: u64, + /// Permission prompts dropped at the bounded pending-request ceiling. + pub permission_cap_drops: u64, +} + /// The authoritative session-state machine over one wrapped `claude` process. /// /// Feed it every line of both stdio directions; it returns an [`ObserveRequest`] @@ -142,6 +155,7 @@ pub struct StreamTracker { /// independently: an ordinary turn changes state without the model, while /// a mid-turn `set_model` changes the model without the state. reported: Option<(SessionState, Option)>, + diagnostics: StreamDiagnostics, } impl StreamTracker { @@ -160,9 +174,16 @@ impl StreamTracker { base: SessionState::Idle, pending: HashSet::new(), reported: None, + diagnostics: StreamDiagnostics::default(), } } + /// Returns cumulative metadata-only protocol diagnostic counts. + #[must_use] + pub fn diagnostics(&self) -> StreamDiagnostics { + self.diagnostics + } + /// The session id, once a line has carried one. #[must_use] pub fn session_id(&self) -> Option<&str> { @@ -201,7 +222,13 @@ impl StreamTracker { // line here is a JSON object, so skipping to the first `{` recovers // it without needing to understand escape-sequence structure at all. let json = line.find('{').map_or(line, |start| &line[start..]); - let parsed: StreamLine = serde_json::from_str(json).ok()?; + let parsed: StreamLine = match serde_json::from_str(json) { + Ok(parsed) => parsed, + Err(_) => { + self.diagnostics.parse_failures = self.diagnostics.parse_failures.saturating_add(1); + return None; + } + }; self.absorb_identity(direction, &parsed); self.apply(direction, &parsed); self.emit_if_changed() @@ -322,14 +349,25 @@ impl StreamTracker { /// state signal and is ignored. fn open_permission(&mut self, parsed: &StreamLine) { let body = parsed.request.as_ref(); - if body.and_then(|b| b.subtype.as_deref()) != Some("can_use_tool") { - return; + match body.and_then(|b| b.subtype.as_deref()) { + Some("can_use_tool") => {} + Some("initialize" | "hook_callback" | "mcp_message" | "set_model") => return, + _ => { + self.diagnostics.unknown_controls = + self.diagnostics.unknown_controls.saturating_add(1); + return; + } } let Some(id) = correlation_id(parsed, body) else { + self.diagnostics.missing_permission_ids = + self.diagnostics.missing_permission_ids.saturating_add(1); return; }; if self.pending.len() < MAX_PENDING_PERMISSIONS { self.pending.insert(id); + } else { + self.diagnostics.permission_cap_drops = + self.diagnostics.permission_cap_drops.saturating_add(1); } } @@ -401,6 +439,34 @@ fn correlation_id(parsed: &StreamLine, body: Option<&ControlBody>) -> Option bool { /// unit-tested against a temp directory. fn scan(root: &Path, state: &mut ScanState, now: SystemTime) -> Vec { let mut sightings = Vec::new(); + let mut scanned = 0_u64; + let mut io_errors = 0_u64; + let mut invalid_names = 0_u64; + let mut inactive = 0_u64; + let mut unchanged = 0_u64; let Ok(project_dirs) = std::fs::read_dir(root) else { + tracing::debug!(outcome = "root_unreadable", "session_transcript_scan"); return sightings; }; - for project in project_dirs.flatten() { + for project in project_dirs { + let project = match project { + Ok(project) => project, + Err(_) => { + io_errors += 1; + continue; + } + }; let Ok(files) = std::fs::read_dir(project.path()) else { + io_errors += 1; continue; }; - for file in files.flatten() { + for file in files { + let file = match file { + Ok(file) => file, + Err(_) => { + io_errors += 1; + continue; + } + }; let path = file.path(); if path.extension().and_then(|e| e.to_str()) != Some("jsonl") { continue; } + scanned += 1; let Some(session_id) = path .file_stem() .and_then(|s| s.to_str()) .filter(|s| !s.is_empty()) .map(str::to_string) else { + invalid_names += 1; continue; }; let Ok(meta) = file.metadata() else { + io_errors += 1; continue; }; let size = meta.len(); - let recent = meta.modified().is_ok_and(|m| is_recent(m, now)); + let recent = match meta.modified() { + Ok(modified) => is_recent(modified, now), + Err(_) => { + io_errors += 1; + false + } + }; let previous = state.insert(path.clone(), size); if !recent { + inactive += 1; // Record the size (for a future growth comparison) but do not // announce an inactive session. continue; @@ -159,7 +190,10 @@ fn scan(root: &Path, state: &mut ScanState, now: SystemTime) -> Vec { let event = match previous { None => SessionEvent::TranscriptDiscovered, Some(prev) if size > prev => SessionEvent::TranscriptGrew, - Some(_) => continue, + Some(_) => { + unchanged += 1; + continue; + } }; sightings.push(Sighting { session_id, @@ -168,6 +202,15 @@ fn scan(root: &Path, state: &mut ScanState, now: SystemTime) -> Vec { }); } } + tracing::debug!( + scanned, + sightings = sightings.len(), + io_errors, + invalid_names, + inactive, + unchanged, + "session_transcript_scan" + ); sightings } @@ -197,7 +240,10 @@ pub fn spawn(registry: Arc, token: CancellationToken) -> JoinH (owned_state, sightings) }) .await - .unwrap_or_else(|_| (ScanState::new(), Vec::new())); + .unwrap_or_else(|error| { + tracing::warn!(%error, outcome = "state_reset", "session_transcript_scan_failed"); + (ScanState::new(), Vec::new()) + }); state = returned_state; for sighting in sightings { registry.observe(sighting.into_observe()); @@ -216,6 +262,20 @@ mod tests { use super::*; use std::io::Write; + #[test] + fn scan_diagnostics_summarize_candidates_and_unchanged_files() { + let dir = tempfile::tempdir().unwrap(); + write_transcript(dir.path(), "project", "session", b"hello"); + let mut state = ScanState::new(); + let logs = crate::test_support::capture_at(tracing::Level::DEBUG, || { + assert_eq!(scan(dir.path(), &mut state, SystemTime::now()).len(), 1); + assert!(scan(dir.path(), &mut state, SystemTime::now()).is_empty()); + }); + assert_eq!(logs.matches("session_transcript_scan").count(), 2); + assert!(logs.contains("scanned=1")); + assert!(logs.contains("unchanged=1")); + } + /// Creates `root//.jsonl` with `contents`, returning its path. fn write_transcript(root: &Path, project: &str, session: &str, contents: &[u8]) -> PathBuf { let dir = root.join(project); diff --git a/src/test_support.rs b/src/test_support.rs index 3ce3b2388..d81c01852 100644 --- a/src/test_support.rs +++ b/src/test_support.rs @@ -150,6 +150,25 @@ pub(crate) fn capture_at(level: tracing::Level, f: impl FnOnce()) -> String { logs } +/// Captures events from an async future on every poll, even across worker threads. +/// Spawned child tasks still need their own subscriber; this does not install a +/// global subscriber or hold a thread-local guard across an await. +pub(crate) async fn capture_future_at( + level: tracing::Level, + future: F, +) -> (F::Output, String) { + use tracing::instrument::WithSubscriber; + let writer = CaptureWriter::default(); + let subscriber = tracing_subscriber::fmt() + .with_max_level(level) + .with_ansi(false) + .with_writer(writer.clone()) + .finish(); + let result = future.with_subscriber(subscriber).await; + let logs = String::from_utf8_lossy(&writer.0.lock().unwrap()).into_owned(); + (result, logs) +} + pub(crate) mod failing_io { //! Writer fixture that always returns `ErrorKind::Other` from //! `write` and `flush`. Used to drive `?`-propagation Err branches diff --git a/src/utils/settings.rs b/src/utils/settings.rs index a606f21a1..cfc115606 100644 --- a/src/utils/settings.rs +++ b/src/utils/settings.rs @@ -372,6 +372,15 @@ pub struct LeaseSettings { pub allow_headless: bool, } +/// Defaults for the long-lived daemon, independent of the launcher's environment. +#[derive(Debug, Default, Deserialize)] +pub struct DaemonSettings { + /// Tracing directive (e.g. `"info"` or `"omni_dev::sessions=debug"`). + /// A valid `RUST_LOG` overrides this; the built-in fallback is `"info"`. + #[serde(default)] + pub log_level: Option, +} + /// Settings loaded from $HOME/.omni-dev/settings.json. #[derive(Debug, Default, Deserialize)] pub struct Settings { @@ -390,6 +399,10 @@ pub struct Settings { #[serde(default)] pub mcp: McpSettings, + /// Daemon tracing defaults; an absent section preserves the built-in filter. + #[serde(default)] + pub daemon: DaemonSettings, + /// Named Gmail accounts (issue #1500); an absent block yields /// [`GmailSettings::default`], which is an empty account map. #[serde(default)] @@ -633,6 +646,22 @@ impl Settings { Self::load_or_warn_default().mcp } + /// Loads daemon defaults with the shared warn-and-default contract. + pub fn load_daemon() -> DaemonSettings { + Self::load_or_warn_default().daemon + } + + /// Records a bootstrap settings failure after tracing has been installed. + /// Shares deduplication with later settings reads in the same process. + pub fn warn_bootstrap_failure(message: &str) { + if LOAD_WARN_DEDUP.observe(Some(message)) { + tracing::warn!( + "{message}; falling back to default settings for this invocation — \ + any settings.json configuration is being ignored" + ); + } + } + /// Loads settings from a specific path. pub fn load_from_path>(path: P) -> Result { let path = path.as_ref(); @@ -1404,6 +1433,21 @@ pub fn get_env_vars(keys: &[&str]) -> Result { #[cfg(test)] #[allow(clippy::unwrap_used, clippy::expect_used)] mod tests { + #[test] + fn daemon_settings_are_optional_and_accept_directives() { + for json in ["{}", r#"{"daemon":{}}"#] { + let settings: super::Settings = serde_json::from_str(json).unwrap(); + assert!(settings.daemon.log_level.is_none()); + } + let settings: super::Settings = + serde_json::from_str(r#"{"daemon":{"log_level":"info,omni_dev::sessions=debug"}}"#) + .unwrap(); + assert_eq!( + settings.daemon.log_level.as_deref(), + Some("info,omni_dev::sessions=debug") + ); + } + use super::*; use crate::test_support::env::MapEnv; use std::env; From 1b4415d2f283528e84618b882faacf2dada11f13 Mon Sep 17 00:00:00 2001 From: John Ky Date: Thu, 1 Oct 2026 20:26:56 +1000 Subject: [PATCH 2/5] fix(cli,daemon,docs): harden wrapper diagnostic appends Serialize complete records before appending so simultaneous wrappers cannot interleave JSON tokens. Include timestamps and wrapper PID for correlation, and stop formatting records after the writer has failed. Add a concurrent-append regression and repair stream lifecycle docs. --- docs/sessions-service.md | 6 ++- src/cli/claude_wrap/diagnostics.rs | 71 +++++++++++++++++++++++++----- src/daemon/server.rs | 20 ++++----- 3 files changed, 74 insertions(+), 23 deletions(-) diff --git a/docs/sessions-service.md b/docs/sessions-service.md index c761c7187..bb0778498 100644 --- a/docs/sessions-service.md +++ b/docs/sessions-service.md @@ -1012,14 +1012,16 @@ identify parse failure, unknown control shapes, missing permission IDs, and permission-cap drops. Known unrelated control requests are ignored normally. The summary also counts drops from the independent diagnostic queue. -When unset or empty, no file, writer, or diagnostic counters are created. +When unset or empty, no file, writer thread, tee-counter allocation, or +diagnostic formatting is created. The pure tracker still counts protocol drift. When enabled, a separate thread appends records through a bounded nonblocking queue; byte pumps only update counters. New files use `0600`; symlinks and nonregular targets are refused, and existing file permissions are left intact. Open/write failures silently disable writing without preventing Claude from launching. Shutdown waits at most 200 ms for diagnostics; a killed process, stuck disk, or full queue can lose records. Files append across processes and -include session/PID metadata where known; there is no automatic rotation. +include timestamps, wrapper PID, and session/child PID metadata where known; +there is no automatic rotation. Conversation messages, tool inputs/results, raw stdio/hook payloads, and daemon error text are never written to this wrapper file. Paths and identifiers are diff --git a/src/cli/claude_wrap/diagnostics.rs b/src/cli/claude_wrap/diagnostics.rs index 1bf71ddc5..120497c22 100644 --- a/src/cli/claude_wrap/diagnostics.rs +++ b/src/cli/claude_wrap/diagnostics.rs @@ -2,7 +2,7 @@ use std::io::Write; use std::os::unix::fs::OpenOptionsExt; use std::path::Path; -use std::sync::atomic::{AtomicU64, Ordering}; +use std::sync::atomic::{AtomicBool, AtomicU64, Ordering}; use std::sync::{mpsc, Arc}; use serde_json::{json, Value}; @@ -16,6 +16,7 @@ pub(super) struct Counts { pub oversize: AtomicU64, pub non_utf8: AtomicU64, pub diagnostic_drops: AtomicU64, + stopped: AtomicBool, } /// Cloning is cheap; a disabled sink has no channel, worker, or counters. @@ -33,6 +34,8 @@ impl Diagnostics { let path = path.to_path_buf(); let (sender, receiver) = mpsc::sync_channel::(CAPACITY); let (done_tx, done_rx) = tokio::sync::oneshot::channel(); + let counts = Arc::new(Counts::default()); + let worker_counts = counts.clone(); // All disk work lives on this thread, including open. Even a stuck // filesystem cannot hold up byte forwarding or process launch. let worker = std::thread::Builder::new() @@ -48,15 +51,21 @@ impl Diagnostics { // Refuse devices/FIFOs, and never follow a symlink log target. if file.metadata().is_ok_and(|m| m.is_file()) { for record in receiver { - if serde_json::to_writer(&mut file, &record).is_err() - || file.write_all(b"\n").is_err() - { + // Serialize the entire record before the append: + // multiple wrappers may share this file. Token-sized + // writes from to_writer could interleave their JSON. + let Ok(mut line) = serde_json::to_vec(&record) else { + break; + }; + line.push(b'\n'); + if file.write_all(&line).is_err() { break; } } let _ = file.flush(); } } + worker_counts.stopped.store(true, Ordering::Relaxed); let _ = done_tx.send(()); }); if worker.is_err() { @@ -65,14 +74,18 @@ impl Diagnostics { ( Self { sender: Some(sender), - counts: Some(Arc::new(Counts::default())), + counts: Some(counts), }, Some(done_rx), ) } pub fn enabled(&self) -> bool { - self.sender.as_ref().is_some_and(|_| self.counts.is_some()) + self.sender.is_some() + && self + .counts + .as_ref() + .is_some_and(|c| !c.stopped.load(Ordering::Relaxed)) } pub fn counts(&self) -> Option<&Counts> { @@ -81,8 +94,22 @@ impl Diagnostics { /// Only call with allowlisted metadata; never stream lines or daemon errors. pub fn record(&self, record: impl FnOnce() -> Value) { + if !self.enabled() { + return; + } if let Some(sender) = &self.sender { - if sender.try_send(record()).is_err() { + let mut record = record(); + if let Some(fields) = record.as_object_mut() { + fields.insert("wrapper_pid".into(), json!(std::process::id())); + fields.insert( + "timestamp_ms".into(), + json!(std::time::SystemTime::now() + .duration_since(std::time::UNIX_EPOCH) + .unwrap_or_default() + .as_millis()), + ); + } + if sender.try_send(record).is_err() { if let Some(counts) = &self.counts { increment(&counts.diagnostic_drops); } @@ -160,6 +187,29 @@ mod tests { ); } + #[tokio::test] + async fn concurrent_sinks_append_complete_records() { + let dir = tempfile::tempdir().unwrap(); + let path = dir.path().join("shared.jsonl"); + let (a, a_done) = Diagnostics::open(Some(&path)); + let (b, b_done) = Diagnostics::open(Some(&path)); + for n in 0..50 { + a.record(|| json!({"event":"a", "n":n})); + b.record(|| json!({"event":"b", "n":n})); + } + drop(a); + drop(b); + a_done.unwrap().await.unwrap(); + b_done.unwrap().await.unwrap(); + let text = std::fs::read_to_string(path).unwrap(); + assert_eq!(text.lines().count(), 100); + for line in text.lines() { + let value: Value = serde_json::from_str(line).unwrap(); + assert!(value["wrapper_pid"].is_number()); + assert!(value["timestamp_ms"].is_number()); + } + } + #[tokio::test] async fn failed_sink_is_nonblocking() { let dir = tempfile::tempdir().unwrap(); @@ -168,15 +218,14 @@ mod tests { .await .unwrap() .unwrap(); - for _ in 0..1000 { - sink.record(|| json!({"event":"test"})); - } + assert!(!sink.enabled()); + sink.record(|| panic!("failed sink formatted a record")); assert_eq!( sink.counts() .unwrap() .diagnostic_drops .load(Ordering::Relaxed), - 1000 + 0 ); } } diff --git a/src/daemon/server.rs b/src/daemon/server.rs index 14c4c2d35..0a4c1a84f 100644 --- a/src/daemon/server.rs +++ b/src/daemon/server.rs @@ -345,16 +345,6 @@ async fn handle_connection( } } -/// Drives a push subscription over `framed` until the client goes away or the -/// daemon shuts down. Sends an initial snapshot, then re-samples the stream on -/// each change notification and on a periodic [`stream_tick`], pushing **only** -/// snapshots that differ from the last one sent — so identical frames are never -/// duplicated (the acceptance criterion). Mirrors the browser bridge's -/// `start_stream` coalescing shape, but on the control socket. -/// -/// The subscription owns the connection for its lifetime: any further inbound -/// line is treated as an explicit cancel and ends the stream, matching the -/// one-op-per-connection the companion uses (a dedicated subscribe socket). /// Reasons a subscription ends; no incoming line content is retained. #[derive(Debug, Clone, Copy, PartialEq, Eq)] enum StreamEnd { @@ -373,6 +363,16 @@ fn inbound_end(line: Option>) -> StreamEnd { } } +/// Drives a push subscription over `framed` until the client goes away or the +/// daemon shuts down. Sends an initial snapshot, then re-samples the stream on +/// each change notification and on a periodic [`stream_tick`], pushing **only** +/// snapshots that differ from the last one sent — so identical frames are never +/// duplicated (the acceptance criterion). Mirrors the browser bridge's +/// `start_stream` coalescing shape, but on the control socket. +/// +/// The subscription owns the connection for its lifetime: any further inbound +/// line is treated as an explicit cancel and ends the stream, matching the +/// one-op-per-connection the companion uses (a dedicated subscribe socket). async fn run_stream( framed: &mut Framed, mut stream: Box, From d3390ff17c249782306f95fe10c14f124464f1d8 Mon Sep 17 00:00:00 2001 From: John Ky Date: Thu, 1 Oct 2026 20:57:49 +1000 Subject: [PATCH 3/5] style(sessions,cli): satisfy diagnostic branch lints Use let-else and if-let in diagnostic error branches, following the repository pedantic Clippy rules. Preserve outcomes and counters. --- src/cli/sessions.rs | 21 ++++++++++----------- src/sessions.rs | 17 ++++++++--------- src/sessions/stream.rs | 9 +++------ src/sessions/watcher.rs | 29 +++++++++++------------------ 4 files changed, 32 insertions(+), 44 deletions(-) diff --git a/src/cli/sessions.rs b/src/cli/sessions.rs index 9f44dcdf8..eeb50cf12 100644 --- a/src/cli/sessions.rs +++ b/src/cli/sessions.rs @@ -392,16 +392,15 @@ impl HookPayload { model: self.model.clone(), pid, }; - match serde_json::to_value(request) { - Ok(payload) => Some(("observe", payload)), - Err(_) => { - tracing::debug!( - ?agent, - outcome = "serialization_failed", - "session_hook_skipped" - ); - None - } + if let Ok(payload) = serde_json::to_value(request) { + Some(("observe", payload)) + } else { + tracing::debug!( + ?agent, + outcome = "serialization_failed", + "session_hook_skipped" + ); + None } } } @@ -1786,7 +1785,7 @@ mod tests { socket: Some(dir.path().join("missing.sock")), agent: HookAgent::Claude, }; - let (_, logs) = crate::test_support::capture_future_at(tracing::Level::DEBUG, async { + let ((), logs) = crate::test_support::capture_future_at(tracing::Level::DEBUG, async { for input in ["not json HOOK_CONTENT_SECRET", "{}", r#"{"session_id":"test"}"#, r#"{"session_id":"test","hook_event_name":"future-event"}"#, diff --git a/src/sessions.rs b/src/sessions.rs index 87f4e3557..fcf91b270 100644 --- a/src/sessions.rs +++ b/src/sessions.rs @@ -1123,17 +1123,16 @@ fn resolve_source(cwd: Option<&Path>, windows: &[WindowEntry]) -> Source { tracing::trace!(outcome = "missing_cwd", "session_attribution_miss"); return Source::Terminal; }; - match pick_window(cwd, windows) { - Some(window) => Source::VsCode { + if let Some(window) = pick_window(cwd, windows) { + Source::VsCode { window_key: window.key.clone(), - }, - None => { - tracing::debug!(cwd = %cwd.display(), window_count = windows.len(), "session_attribution_miss"); - for window in windows { - tracing::trace!(window_key = %window.report.key, folders = ?window.report.folders, "session_attribution_candidate"); - } - Source::Terminal } + } else { + tracing::debug!(cwd = %cwd.display(), window_count = windows.len(), "session_attribution_miss"); + for window in windows { + tracing::trace!(window_key = %window.report.key, folders = ?window.report.folders, "session_attribution_candidate"); + } + Source::Terminal } } diff --git a/src/sessions/stream.rs b/src/sessions/stream.rs index ee9ebb819..d3a6745be 100644 --- a/src/sessions/stream.rs +++ b/src/sessions/stream.rs @@ -222,12 +222,9 @@ impl StreamTracker { // line here is a JSON object, so skipping to the first `{` recovers // it without needing to understand escape-sequence structure at all. let json = line.find('{').map_or(line, |start| &line[start..]); - let parsed: StreamLine = match serde_json::from_str(json) { - Ok(parsed) => parsed, - Err(_) => { - self.diagnostics.parse_failures = self.diagnostics.parse_failures.saturating_add(1); - return None; - } + let Ok(parsed) = serde_json::from_str::(json) else { + self.diagnostics.parse_failures = self.diagnostics.parse_failures.saturating_add(1); + return None; }; self.absorb_identity(direction, &parsed); self.apply(direction, &parsed); diff --git a/src/sessions/watcher.rs b/src/sessions/watcher.rs index 175d477c1..5ddc0ff98 100644 --- a/src/sessions/watcher.rs +++ b/src/sessions/watcher.rs @@ -135,24 +135,18 @@ fn scan(root: &Path, state: &mut ScanState, now: SystemTime) -> Vec { return sightings; }; for project in project_dirs { - let project = match project { - Ok(project) => project, - Err(_) => { - io_errors += 1; - continue; - } + let Ok(project) = project else { + io_errors += 1; + continue; }; let Ok(files) = std::fs::read_dir(project.path()) else { io_errors += 1; continue; }; for file in files { - let file = match file { - Ok(file) => file, - Err(_) => { - io_errors += 1; - continue; - } + let Ok(file) = file else { + io_errors += 1; + continue; }; let path = file.path(); if path.extension().and_then(|e| e.to_str()) != Some("jsonl") { @@ -173,12 +167,11 @@ fn scan(root: &Path, state: &mut ScanState, now: SystemTime) -> Vec { continue; }; let size = meta.len(); - let recent = match meta.modified() { - Ok(modified) => is_recent(modified, now), - Err(_) => { - io_errors += 1; - false - } + let recent = if let Ok(modified) = meta.modified() { + is_recent(modified, now) + } else { + io_errors += 1; + false }; let previous = state.insert(path.clone(), size); if !recent { From 1e73eb2d8221ffcad1392d53541fa68ab5c620fb Mon Sep 17 00:00:00 2001 From: John Ky Date: Thu, 1 Oct 2026 22:25:19 +1000 Subject: [PATCH 4/5] test(sessions,cli,daemon,lib): cover the diagnostic logging branches Close the remaining coverage gaps in the diagnostics added for #1447, mostly by testing the new log lines and failure arms directly rather than excluding them. Pull a few small seams out of the bodies that held the untestable branches, without changing behavior: - Diagnostics::drain takes the writer and receiver, so a failing writer can drive the abandon-the-log path. - HookCommand::read_input takes a reader, so a broken or non-UTF-8 stdin can be fed in. - watcher::session_id_of isolates the file-stem rule behind scan. - main::daemon_filter_from and emit_bootstrap_warnings split the bootstrap settings load from the warnings it owes once tracing exists. - Settings::warn_settings_fallback takes its dedup, so the once-only warning is asserted on a private one. Mark the arms that only an OS fault can reach (a thread that will not spawn, a read_dir error entry, a missing data directory) with coverage ignores that state why a test cannot reproduce them. Make capture_at and capture_future_at deterministic. tracing-core recomputes a callsite's cached interest against only the calling thread's default while a single dispatcher is registered, so a thread with no subscriber could cache "never" and silently drop the capturing test's event. Keeping one extra dispatcher registered for the life of the process keeps it on the path that consults every live dispatcher. --- src/cli/claude_wrap.rs | 22 ++++++ src/cli/claude_wrap/diagnostics.rs | 93 ++++++++++++++++--------- src/cli/sessions.rs | 105 +++++++++++++++++++++++++++-- src/daemon/server.rs | 35 ++++++++++ src/main.rs | 73 ++++++++++++++++++-- src/sessions.rs | 34 +++++++++- src/sessions/pid_watcher.rs | 44 ++++++++++++ src/sessions/watcher.rs | 66 ++++++++++++++++-- src/test_support.rs | 20 ++++++ src/utils/settings.rs | 85 +++++++++++++++++++---- 10 files changed, 513 insertions(+), 64 deletions(-) diff --git a/src/cli/claude_wrap.rs b/src/cli/claude_wrap.rs index c9434a62e..56e669359 100644 --- a/src/cli/claude_wrap.rs +++ b/src/cli/claude_wrap.rs @@ -151,6 +151,7 @@ fn exec_replace(program: &str, args: &[String]) -> anyhow::Error { /// Wraps `program`, joining it to this process's own stdin and stdout. async fn wrap(program: &str, args: &[String], socket: Option) -> Result { + // omni-dev: coverage ignore reason="process-bound wiring shell: it joins this process's own stdin/stdout and reads the real OMNI_DEV_CLAUDE_WRAP_LOG, which a test must not take over (see the note on run/wrap in the tests module); wrap_io_diagnostics beneath it is covered directly" wrap_io_diagnostics( program, args, @@ -160,6 +161,7 @@ async fn wrap(program: &str, args: &[String], socket: Option) -> Result std::env::var_os("OMNI_DEV_CLAUDE_WRAP_LOG").map(PathBuf::from), ) .await + // omni-dev: coverage end } /// Spawns the child with piped stdio, pumps `input` and `output` through it @@ -705,9 +707,11 @@ async fn observe_diagnostics( diagnostics: Diagnostics, ) { let Ok(socket) = server::resolve_socket(socket) else { + // omni-dev: coverage ignore reason="resolve_socket fails only when the platform has no data directory to put the default socket in (no resolvable home), which a test cannot reproduce on macOS or Linux; the fail-open return is what keeps the wrapper forwarding" diagnostics .record(|| json!({"event":"observer_stopped", "outcome":"socket_resolution_failed"})); return; + // omni-dev: coverage end }; let mut tracker = StreamTracker::new(); let mut last_model: Option = None; @@ -760,9 +764,11 @@ async fn observe_diagnostics( if let Ok(payload) = serde_json::to_value(request) { report_diagnostics(&socket, "observe", payload, &diagnostics).await; } else { + // omni-dev: coverage ignore reason="to_value on an ObserveRequest fails only for a non-UTF-8 cwd, and the tracker takes cwd from a JSON string, so it is always UTF-8; the arm exists so a future non-string field cannot silently drop the report" diagnostics.record( || json!({"event":"report", "op":"observe", "outcome":"serialization_failed"}), ); + // omni-dev: coverage end } } } @@ -961,6 +967,22 @@ mod tests { assert!(!text.contains("DAEMON_CONTENT_SECRET")); } + #[tokio::test(start_paused = true)] + async fn diagnostic_report_gives_up_on_a_wedged_daemon() { + let dir = tempfile::tempdir_in("/tmp").unwrap(); + let socket = dir.path().join("d.sock"); + // Bound but never accepted: the connect lands in the backlog and the + // request in the socket buffer, so only the report timeout can end the + // exchange. With the clock paused, that two-second wait costs no wall time. + let _wedged = UnixListener::bind(&socket).unwrap(); + let path = dir.path().join("log"); + let (diagnostics, done) = Diagnostics::open(Some(&path)); + report_diagnostics(&socket, "observe", json!({}), &diagnostics).await; + drop(diagnostics); + done.unwrap().await.unwrap(); + assert!(std::fs::read_to_string(path).unwrap().contains("timeout")); + } + /// An [`AsyncWrite`] that appends into a shared buffer, so a test can assert /// on exactly the bytes the wrapper forwarded. #[derive(Clone, Default)] diff --git a/src/cli/claude_wrap/diagnostics.rs b/src/cli/claude_wrap/diagnostics.rs index 120497c22..4e11f77d9 100644 --- a/src/cli/claude_wrap/diagnostics.rs +++ b/src/cli/claude_wrap/diagnostics.rs @@ -50,26 +50,14 @@ impl Diagnostics { if let Ok(mut file) = file { // Refuse devices/FIFOs, and never follow a symlink log target. if file.metadata().is_ok_and(|m| m.is_file()) { - for record in receiver { - // Serialize the entire record before the append: - // multiple wrappers may share this file. Token-sized - // writes from to_writer could interleave their JSON. - let Ok(mut line) = serde_json::to_vec(&record) else { - break; - }; - line.push(b'\n'); - if file.write_all(&line).is_err() { - break; - } - } - let _ = file.flush(); + drain(&receiver, &mut file); } } worker_counts.stopped.store(true, Ordering::Relaxed); let _ = done_tx.send(()); }); if worker.is_err() { - return (Self::default(), None); + return (Self::default(), None); // omni-dev: coverage ignore-line reason="Builder::spawn fails only when the OS refuses a new thread (resource exhaustion), which a test cannot provoke; the fail-open return is what keeps the wrapper forwarding" } ( Self { @@ -94,25 +82,23 @@ impl Diagnostics { /// Only call with allowlisted metadata; never stream lines or daemon errors. pub fn record(&self, record: impl FnOnce() -> Value) { - if !self.enabled() { + let Some(sender) = self.sender.as_ref().filter(|_| self.enabled()) else { return; + }; + let mut record = record(); + if let Some(fields) = record.as_object_mut() { + fields.insert("wrapper_pid".into(), json!(std::process::id())); + fields.insert( + "timestamp_ms".into(), + json!(std::time::SystemTime::now() + .duration_since(std::time::UNIX_EPOCH) + .unwrap_or_default() + .as_millis()), + ); } - if let Some(sender) = &self.sender { - let mut record = record(); - if let Some(fields) = record.as_object_mut() { - fields.insert("wrapper_pid".into(), json!(std::process::id())); - fields.insert( - "timestamp_ms".into(), - json!(std::time::SystemTime::now() - .duration_since(std::time::UNIX_EPOCH) - .unwrap_or_default() - .as_millis()), - ); - } - if sender.try_send(record).is_err() { - if let Some(counts) = &self.counts { - increment(&counts.diagnostic_drops); - } + if sender.try_send(record).is_err() { + if let Some(counts) = &self.counts { + increment(&counts.diagnostic_drops); } } } @@ -129,6 +115,22 @@ impl Diagnostics { } } +/// Appends each record to `out` as one line until the channel closes or a write +/// fails; a failed write abandons the log rather than the wrapper's forwarding. +fn drain(receiver: &mpsc::Receiver, out: &mut impl Write) { + for record in receiver { + // Serialize the entire record before the append: multiple wrappers may + // share this file, and token-sized writes from `to_writer` could + // interleave their JSON. `Value`'s `Display` cannot fail. + let mut line = record.to_string().into_bytes(); + line.push(b'\n'); + if out.write_all(&line).is_err() { + break; + } + } + let _ = out.flush(); +} + pub(super) fn increment(counter: &AtomicU64) { let _ = counter.fetch_update(Ordering::Relaxed, Ordering::Relaxed, |n| { Some(n.saturating_add(1)) @@ -210,6 +212,35 @@ mod tests { } } + #[test] + fn drain_abandons_the_log_when_a_write_fails() { + let (sender, receiver) = mpsc::sync_channel(4); + sender.send(json!({"event":"first"})).unwrap(); + sender.send(json!({"event":"second"})).unwrap(); + drop(sender); + // The first failed write ends the drain; it must neither panic nor + // keep retrying the remaining records. + drain( + &receiver, + &mut crate::test_support::failing_io::FailingWriter, + ); + assert_eq!(receiver.try_iter().count(), 1); + } + + #[test] + fn drain_writes_one_complete_line_per_record() { + let (sender, receiver) = mpsc::sync_channel(4); + sender.send(json!({"event":"first"})).unwrap(); + sender.send(json!({"event":"second","n":2})).unwrap(); + drop(sender); + let mut out = Vec::new(); + drain(&receiver, &mut out); + assert_eq!( + String::from_utf8(out).unwrap(), + "{\"event\":\"first\"}\n{\"event\":\"second\",\"n\":2}\n" + ); + } + #[tokio::test] async fn failed_sink_is_nonblocking() { let dir = tempfile::tempdir().unwrap(); diff --git a/src/cli/sessions.rs b/src/cli/sessions.rs index eeb50cf12..388d9acb5 100644 --- a/src/cli/sessions.rs +++ b/src/cli/sessions.rs @@ -240,13 +240,22 @@ impl HookCommand { // Read before stdin, while the agent that spawned this hook is surely // still alive to be our parent. let pid = agent_pid(std::os::unix::process::parent_id()); + if let Some(input) = self.read_input(std::io::stdin()) { + self.report(&input, pid).await; + } + Ok(()) + } + + /// Reads the whole hook payload from `reader`, or `None` — after a debug + /// line — when it cannot be read (an I/O error, or bytes that are not + /// UTF-8). Split out so tests can feed it a reader that fails. + fn read_input(&self, mut reader: impl Read) -> Option { let mut input = String::new(); - if std::io::stdin().read_to_string(&mut input).is_err() { + if reader.read_to_string(&mut input).is_err() { tracing::debug!(agent = ?self.agent, outcome = "stdin_read_failed", "session_hook_skipped"); - return Ok(()); + return None; } - self.report(&input, pid).await; - Ok(()) + Some(input) } /// Parses the hook JSON, maps it to an op, and best-effort sends it. Split @@ -260,8 +269,10 @@ impl HookCommand { return; }; let Ok(socket) = server::resolve_socket(self.socket.clone()) else { + // omni-dev: coverage ignore reason="resolve_socket fails only when the platform has no data directory to put the default socket in (no resolvable home), which a test cannot reproduce on macOS or Linux; the sink is fail-open by design, so this is the same silent return as every other skipped hook" tracing::debug!(agent = ?self.agent, outcome = "socket_resolution_failed", "session_hook_skipped"); return; + // omni-dev: coverage end }; let env = DaemonEnvelope::service(SERVICE, op, payload); let outcome = @@ -395,12 +406,14 @@ impl HookPayload { if let Ok(payload) = serde_json::to_value(request) { Some(("observe", payload)) } else { + // omni-dev: coverage ignore reason="to_value on an ObserveRequest fails only for a non-UTF-8 cwd, and the hook payload's cwd is deserialized from a JSON string, so it is always UTF-8; the arm exists so a future non-string field cannot silently drop the report" tracing::debug!( ?agent, outcome = "serialization_failed", "session_hook_skipped" ); None + // omni-dev: coverage end } } } @@ -463,10 +476,12 @@ fn session_event_for( "Elicitation" => SessionEvent::Notification(NotificationKind::AgentNeedsInput), "Notification" => { let kind = classify_notification(notification_type, message); + let has_notification_type = notification_type.is_some(); + let has_message = message.is_some(); tracing::debug!( ?kind, - has_notification_type = notification_type.is_some(), - has_message = message.is_some(), + has_notification_type, + has_message, "session_notification_classified" ); SessionEvent::Notification(kind) @@ -3667,6 +3682,84 @@ mod tests { cmd.report(r#"{"hook_event_name":"Stop"}"#, None).await; // no session_id → no op } + #[test] + fn hook_input_that_cannot_be_read_is_skipped_with_a_debug_line() { + /// A reader whose every read fails, like a stdin closed under the sink. + struct Broken; + impl Read for Broken { + fn read(&mut self, _: &mut [u8]) -> std::io::Result { + Err(std::io::Error::other("stdin went away")) + } + } + let cmd = HookCommand { + socket: None, + agent: HookAgent::Claude, + }; + let logs = crate::test_support::capture_at(tracing::Level::DEBUG, || { + // Bytes that are not UTF-8 fail `read_to_string`, as does an I/O error. + assert_eq!(cmd.read_input(&b"\xff\xfe"[..]), None); + assert_eq!(cmd.read_input(Broken), None); + }); + assert_eq!(logs.matches("stdin_read_failed").count(), 2, "{logs}"); + assert_eq!( + cmd.read_input(&b"{\"session_id\":\"s1\"}"[..]).as_deref(), + Some("{\"session_id\":\"s1\"}") + ); + } + + #[tokio::test] + async fn hook_report_logs_a_daemon_rejection_without_its_content() { + let (_dir, sock, server) = fake_daemon(json!({"ok": false, "error": "HOOK_REPLY_SECRET"})); + let cmd = HookCommand { + socket: Some(sock), + agent: HookAgent::Claude, + }; + let ((), logs) = crate::test_support::capture_future_at( + tracing::Level::DEBUG, + cmd.report(r#"{"session_id":"s1","hook_event_name":"Stop"}"#, None), + ) + .await; + server.await.unwrap(); + assert!(logs.contains("daemon_rejected"), "{logs}"); + assert!(!logs.contains("HOOK_REPLY_SECRET"), "{logs}"); + } + + #[tokio::test(start_paused = true)] + async fn hook_report_gives_up_on_a_wedged_daemon() { + let dir = tempfile::tempdir_in("/tmp").unwrap(); + let sock = dir.path().join("d.sock"); + // Bound but never accepted: the connect lands in the backlog and the + // request in the socket buffer, so only the hook timeout can end the + // exchange. With the clock paused, that two-second wait costs no wall time. + let _wedged = tokio::net::UnixListener::bind(&sock).unwrap(); + let cmd = HookCommand { + socket: Some(sock), + agent: HookAgent::Claude, + }; + let ((), logs) = crate::test_support::capture_future_at( + tracing::Level::DEBUG, + cmd.report(r#"{"session_id":"s1","hook_event_name":"Stop"}"#, None), + ) + .await; + assert!(logs.contains("timeout"), "{logs}"); + } + + #[tokio::test] + async fn hook_report_logs_a_delivered_event() { + let (_dir, sock, server) = fake_daemon(json!({"ok": true, "payload": {}})); + let cmd = HookCommand { + socket: Some(sock), + agent: HookAgent::Claude, + }; + let ((), logs) = crate::test_support::capture_future_at( + tracing::Level::DEBUG, + cmd.report(r#"{"session_id":"s1","hook_event_name":"Stop"}"#, None), + ) + .await; + server.await.unwrap(); + assert!(logs.contains("delivered"), "{logs}"); + } + /// Spawns a minimal fake daemon on a short-path Unix socket that answers one /// request with `reply`. Returns the temp dir (kept alive), the socket path, /// and the server task. diff --git a/src/daemon/server.rs b/src/daemon/server.rs index 0a4c1a84f..d0b7968d5 100644 --- a/src/daemon/server.rs +++ b/src/daemon/server.rs @@ -946,6 +946,41 @@ mod tests { .unwrap(); } + /// A pushed delta that cannot be written ends the stream as a write failure. + /// The server end shuts down its own write half after the initial snapshot, + /// so the next push fails while the read side is still open — dropping the + /// client instead would race the failing write against the EOF it also causes. + #[tokio::test] + async fn run_stream_ends_when_a_pushed_delta_cannot_be_written() { + let (mut client, server) = UnixStream::pair().unwrap(); + let server = server.into_std().unwrap(); + let hangup = server.try_clone().unwrap(); + let server = UnixStream::from_std(server).unwrap(); + let (tx, rx) = watch::channel(0u64); + let snap = Arc::new(StdMutex::new(json!({ "n": 0 }))); + let fake = FakeStream { + rx, + snap: Arc::clone(&snap), + }; + let shutdown = CancellationToken::new(); + let task = tokio::spawn(async move { + let mut framed = Framed::new(server, LinesCodec::new_with_max_length(MAX_LINE_BYTES)); + run_stream(&mut framed, Box::new(fake), &shutdown, "fake", "subscribe").await + }); + + let mut reader = BufReader::new(&mut client); + read_reply(&mut reader).await; + hangup.shutdown(std::net::Shutdown::Write).unwrap(); + *snap.lock().unwrap() = json!({ "n": 1 }); + tx.send(1).unwrap(); + + let end = tokio::time::timeout(Duration::from_secs(2), task) + .await + .expect("run_stream should end when a push cannot be written") + .unwrap(); + assert_eq!(end, StreamEnd::WriteFailure); + } + /// `run_stream` returns immediately when even the initial snapshot cannot be /// sent (the client is already gone) rather than entering the select loop. #[tokio::test] diff --git a/src/main.rs b/src/main.rs index c1e7ce20c..62d687eca 100644 --- a/src/main.rs +++ b/src/main.rs @@ -127,10 +127,7 @@ fn init_tracing(daemon_run: bool) { // Loading through the normal warning loader here would lose its warning: // there is no subscriber yet. Emit once after installation instead. let (daemon_filter, settings_error) = if daemon_run { - match omni_dev::utils::settings::Settings::load() { - Ok(settings) => (settings.daemon.log_level, None), - Err(error) => (None, Some(format!("{error:#}"))), - } + daemon_filter_from(omni_dev::utils::settings::Settings::load()) } else { (None, None) }; @@ -141,6 +138,23 @@ fn init_tracing(daemon_run: bool) { .with_writer(std::io::stderr) .with_env_filter(filter) .init(); + emit_bootstrap_warnings(&rejected, settings_error.as_deref()); +} + +/// Splits a bootstrap settings load into the daemon's configured log filter and, +/// when the load failed, the error to report once tracing is installed. +fn daemon_filter_from( + loaded: anyhow::Result, +) -> (Option, Option) { + match loaded { + Ok(settings) => (settings.daemon.log_level, None), + Err(error) => (None, Some(format!("{error:#}"))), + } +} + +/// Reports what could not be said before the subscriber existed: each rejected +/// filter directive, and the settings failure the bootstrap load hit. +fn emit_bootstrap_warnings(rejected: &[&str], settings_error: Option<&str>) { for source in rejected { tracing::warn!( source, @@ -148,7 +162,7 @@ fn init_tracing(daemon_run: bool) { ); } if let Some(error) = settings_error { - omni_dev::utils::settings::Settings::warn_bootstrap_failure(&error); + omni_dev::utils::settings::Settings::warn_bootstrap_failure(error); } } @@ -231,6 +245,55 @@ mod tests { } } + #[test] + fn a_failed_bootstrap_settings_load_is_carried_for_later_reporting() { + let (filter, error) = daemon_filter_from(Err(anyhow::anyhow!("settings.json is broken"))); + assert_eq!(filter, None); + assert_eq!(error.as_deref(), Some("settings.json is broken")); + + let settings: omni_dev::utils::settings::Settings = + serde_json::from_str(r#"{"daemon":{"log_level":"debug"}}"#).unwrap(); + assert_eq!( + daemon_filter_from(Ok(settings)), + (Some("debug".to_string()), None) + ); + } + + #[test] + fn bootstrap_warnings_name_each_rejected_source_and_the_settings_failure() { + use std::io::{Read, Seek}; + use std::sync::Arc; + + // A temp file is the capture buffer: `&File` is already `Write`, so no + // custom writer is needed. + let file = Arc::new(tempfile::tempfile().unwrap()); + let subscriber = tracing_subscriber::fmt() + .with_max_level(tracing::Level::WARN) + .with_ansi(false) + .with_writer(Arc::clone(&file)) + .finish(); + tracing::subscriber::with_default(subscriber, || { + // Nothing failed, so nothing is said. + emit_bootstrap_warnings(&[], None); + assert_eq!(file.metadata().unwrap().len(), 0); + + emit_bootstrap_warnings( + &["RUST_LOG", "daemon.log_level"], + Some("main-test bootstrap settings failure"), + ); + }); + let mut logs = String::new(); + (&*file).rewind().unwrap(); + (&*file).read_to_string(&mut logs).unwrap(); + assert!(logs.contains("invalid tracing directive"), "{logs}"); + assert!(logs.contains("RUST_LOG"), "{logs}"); + assert!(logs.contains("daemon.log_level"), "{logs}"); + assert!( + logs.contains("main-test bootstrap settings failure"), + "{logs}" + ); + } + #[test] fn default_filter_is_info_only_for_daemon_run() { assert_eq!(default_filter(true), "info"); diff --git a/src/sessions.rs b/src/sessions.rs index fcf91b270..c16f3d2c6 100644 --- a/src/sessions.rs +++ b/src/sessions.rs @@ -855,7 +855,8 @@ impl SessionsRegistry { let (known, flipped) = known; // A known session flipped to `ended`; otherwise only this call's inline // reap could have changed anything. - if flipped || reaped > 0 { + let bumped = flipped || reaped > 0; + if bumped { self.bump(); } tracing::debug!( @@ -864,7 +865,7 @@ impl SessionsRegistry { ?old_state, outcome, reaped, - bumped = flipped || reaped > 0, + bumped, "session_end" ); known @@ -1297,14 +1298,39 @@ mod tests { Some("/project"), )); registry.end("missing", None, None); + registry.end("diagnostics", None, None); }); assert!(logs.contains("state_changed") || logs.contains("created")); assert!(logs.contains("heartbeat_only")); assert!(logs.contains("unknown")); + assert!(logs.contains("outcome=\"ended\""), "{logs}"); assert!(logs.contains("bumped=true")); assert!(logs.contains("bumped=false")); } + #[test] + fn attribution_misses_are_logged_with_each_candidate_window() { + let windows = vec![window_entry("w1", 20), window_entry("w2", 5)]; + let logs = crate::test_support::capture_at(tracing::Level::TRACE, || { + assert_eq!( + resolve_source(Some(Path::new("/elsewhere/x")), &windows), + Source::Terminal + ); + assert_eq!(resolve_source(None, &windows), Source::Terminal); + }); + assert!(logs.contains("session_attribution_miss"), "{logs}"); + assert!(logs.contains("missing_cwd"), "{logs}"); + assert_eq!( + logs.matches("session_attribution_candidate").count(), + 2, + "{logs}" + ); + assert!( + logs.contains("window_key=w1") && logs.contains("window_key=w2"), + "{logs}" + ); + } + fn observe_request(session_id: &str, event: SessionEvent, cwd: Option<&str>) -> ObserveRequest { ObserveRequest { agent_id: None, @@ -1868,6 +1894,10 @@ mod tests { assert!(!sessions.contains_key("older")); assert!(sessions.contains_key("young")); assert!(sessions.contains_key("old")); + // An empty map is a no-op, not a panic. + let mut empty: HashMap = HashMap::new(); + evict_oldest_session(&mut empty); + assert!(empty.is_empty()); } #[test] diff --git a/src/sessions/pid_watcher.rs b/src/sessions/pid_watcher.rs index 9c21c1879..1fed626fa 100644 --- a/src/sessions/pid_watcher.rs +++ b/src/sessions/pid_watcher.rs @@ -245,6 +245,7 @@ pub(crate) fn spawn(registry: Arc, token: CancellationToken) - #[allow(clippy::unwrap_used, clippy::expect_used)] mod tests { use super::*; + use crate::sessions::{Agent, ObserveRequest, SessionEvent, SessionState}; fn candidate( session_id: &str, @@ -487,6 +488,49 @@ mod tests { assert_eq!(*calls.lock().unwrap(), 1); } + /// Registers `session_id` as a prompted session owned by `pid`. + fn observe_owned_by(registry: &SessionsRegistry, session_id: &str, pid: u32) { + registry.observe(ObserveRequest { + agent_id: None, + session_id: session_id.to_string(), + cwd: None, + transcript_path: None, + event: SessionEvent::UserPromptSubmit, + repo: None, + model: None, + agent: Agent::Claude, + pid: Some(pid), + }); + } + + #[test] + fn an_end_action_ends_the_session_and_logs_why() { + let registry = SessionsRegistry::new(); + observe_owned_by(®istry, "s1", 100); + let logs = crate::test_support::capture_at(tracing::Level::DEBUG, || { + Action::End { + session_id: "s1".to_string(), + } + .apply(®istry, Utc::now()); + }); + assert!(logs.contains("session_process_ended"), "{logs}"); + assert!(logs.contains("pid_liveness"), "{logs}"); + assert_eq!(registry.list()[0].state, SessionState::Ended); + } + + #[test] + fn a_confirm_action_captures_the_pid_identity_token() { + let registry = SessionsRegistry::new(); + observe_owned_by(®istry, "s1", 100); + Action::Confirm { + session_id: "s1".to_string(), + pid_start: "tok".to_string(), + } + .apply(®istry, Utc::now()); + let candidates = registry.pid_liveness_candidates(); + assert_eq!(candidates[0].pid_start.as_deref(), Some("tok")); + } + #[test] fn a_pid_no_longer_among_candidates_is_forgotten() { let mut confirmed = HashSet::from([100, 200]); diff --git a/src/sessions/watcher.rs b/src/sessions/watcher.rs index 5ddc0ff98..67654c13a 100644 --- a/src/sessions/watcher.rs +++ b/src/sessions/watcher.rs @@ -108,6 +108,15 @@ fn is_recent(modified: SystemTime, now: SystemTime) -> bool { } } +/// The session id a transcript path names: its file stem, when that is +/// non-empty UTF-8 (a session id is a UUID, so anything else is not one). +fn session_id_of(path: &Path) -> Option { + path.file_stem() + .and_then(|s| s.to_str()) + .filter(|s| !s.is_empty()) + .map(str::to_string) +} + /// Scans `root` for `*.jsonl` transcripts and returns the sightings since the /// previous scan, updating `state` (path → last size) in place. /// @@ -136,8 +145,10 @@ fn scan(root: &Path, state: &mut ScanState, now: SystemTime) -> Vec { }; for project in project_dirs { let Ok(project) = project else { + // omni-dev: coverage ignore reason="read_dir yields an Err entry only on an I/O fault (EIO, a vanished directory mid-iteration), which a test cannot provoke; the scan counts it and moves on, the same handling as the unreadable project directory below" io_errors += 1; continue; + // omni-dev: coverage end }; let Ok(files) = std::fs::read_dir(project.path()) else { io_errors += 1; @@ -145,33 +156,34 @@ fn scan(root: &Path, state: &mut ScanState, now: SystemTime) -> Vec { }; for file in files { let Ok(file) = file else { + // omni-dev: coverage ignore reason="read_dir yields an Err entry only on an I/O fault (EIO, a vanished directory mid-iteration), which a test cannot provoke; the scan counts it and moves on, the same handling as the unreadable project directory above" io_errors += 1; continue; + // omni-dev: coverage end }; let path = file.path(); if path.extension().and_then(|e| e.to_str()) != Some("jsonl") { continue; } scanned += 1; - let Some(session_id) = path - .file_stem() - .and_then(|s| s.to_str()) - .filter(|s| !s.is_empty()) - .map(str::to_string) - else { + let Some(session_id) = session_id_of(&path) else { invalid_names += 1; continue; }; let Ok(meta) = file.metadata() else { + // omni-dev: coverage ignore reason="DirEntry::metadata does not follow symlinks, so it fails only when the entry vanished between read_dir and the stat, a race a test cannot provoke; the scan counts it and moves on" io_errors += 1; continue; + // omni-dev: coverage end }; let size = meta.len(); let recent = if let Ok(modified) = meta.modified() { is_recent(modified, now) } else { + // omni-dev: coverage ignore reason="Metadata::modified fails only on a platform with no mtime, and every supported one (Linux, macOS) has it" io_errors += 1; false + // omni-dev: coverage end }; let previous = state.insert(path.clone(), size); if !recent { @@ -195,9 +207,10 @@ fn scan(root: &Path, state: &mut ScanState, now: SystemTime) -> Vec { }); } } + let found = sightings.len(); tracing::debug!( scanned, - sightings = sightings.len(), + sightings = found, io_errors, invalid_names, inactive, @@ -234,9 +247,11 @@ pub fn spawn(registry: Arc, token: CancellationToken) -> JoinH }) .await .unwrap_or_else(|error| { + // omni-dev: coverage ignore reason="spawn_blocking's JoinError needs the scan closure to panic or the runtime to shut down mid-scan; scan has no panicking path, and the runtime outlives the watcher, whose token is cancelled first" tracing::warn!(%error, outcome = "state_reset", "session_transcript_scan_failed"); (ScanState::new(), Vec::new()) }); + // omni-dev: coverage end state = returned_state; for sighting in sightings { registry.observe(sighting.into_observe()); @@ -269,6 +284,43 @@ mod tests { assert!(logs.contains("unchanged=1")); } + #[test] + fn a_session_id_is_the_non_empty_utf8_file_stem() { + assert_eq!( + session_id_of(Path::new("/p/abc-123.jsonl")).as_deref(), + Some("abc-123") + ); + // No file name at all, so no stem. + assert_eq!(session_id_of(Path::new("/")), None); + } + + #[cfg(unix)] + #[test] + fn scan_counts_a_non_utf8_transcript_name_and_skips_it() { + use std::os::unix::ffi::OsStrExt; + let dir = tempfile::tempdir().unwrap(); + let project = dir.path().join("project"); + std::fs::create_dir_all(&project).unwrap(); + let name = std::ffi::OsStr::from_bytes(b"\xff\xfe.jsonl"); + // macOS filesystems refuse names that are not valid UTF-8; Linux accepts them. + if std::fs::write(project.join(name), b"x").is_err() { + return; + } + let mut state = ScanState::new(); + let logs = crate::test_support::capture_at(tracing::Level::DEBUG, || { + assert!(scan(dir.path(), &mut state, SystemTime::now()).is_empty()); + }); + assert!(logs.contains("invalid_names=1"), "{logs}"); + } + + #[cfg(unix)] + #[test] + fn a_non_utf8_file_name_is_not_a_session_id() { + use std::os::unix::ffi::OsStrExt; + let path = Path::new(std::ffi::OsStr::from_bytes(b"/p/\xff\xfe.jsonl")); + assert_eq!(session_id_of(path), None); + } + /// Creates `root//.jsonl` with `contents`, returning its path. fn write_transcript(root: &Path, project: &str, session: &str, contents: &[u8]) -> PathBuf { let dir = root.join(project); diff --git a/src/test_support.rs b/src/test_support.rs index d81c01852..f490ac9cb 100644 --- a/src/test_support.rs +++ b/src/test_support.rs @@ -133,12 +133,31 @@ impl<'a> tracing_subscriber::fmt::MakeWriter<'a> for CaptureWriter { } } +/// Keeps one extra dispatcher registered for the life of the process, so that +/// [`capture_at`] and [`capture_future_at`] cannot lose events to a race. +/// +/// `tracing-core` caches each callsite's interest. While only *one* dispatcher is +/// registered it recomputes that interest against the *calling thread's* default +/// instead of every registered dispatcher. In the parallel test binary that +/// lets a thread with no subscriber installed — the first to reach a callsite +/// while another thread's capture subscriber is the only one live — cache +/// "never" for it, and the capturing thread's event is then silently dropped: +/// a flake in the capturing test, and a macro line that reports uncovered. A +/// second, never-dropped registration keeps `tracing-core` on the path that +/// consults every live dispatcher. It is never installed as a default, so it +/// changes what no thread logs. +fn keep_every_dispatcher_consulted() { + static KEEP: std::sync::OnceLock = std::sync::OnceLock::new(); + KEEP.get_or_init(|| tracing::Dispatch::new(tracing::subscriber::NoSubscriber::default())); +} + /// Runs `f` under a thread-local subscriber that captures every event at /// `level` or above, and returns everything it logged. `f` must be fully /// synchronous on this thread. The one shared home for this capture /// pattern (issue #1744); per-module `capture_info`/`capture_warnings` /// helpers are thin aliases over it. pub(crate) fn capture_at(level: tracing::Level, f: impl FnOnce()) -> String { + keep_every_dispatcher_consulted(); let writer = CaptureWriter::default(); let subscriber = tracing_subscriber::fmt() .with_max_level(level) @@ -158,6 +177,7 @@ pub(crate) async fn capture_future_at( future: F, ) -> (F::Output, String) { use tracing::instrument::WithSubscriber; + keep_every_dispatcher_consulted(); let writer = CaptureWriter::default(); let subscriber = tracing_subscriber::fmt() .with_max_level(level) diff --git a/src/utils/settings.rs b/src/utils/settings.rs index cfc115606..cb1573ee2 100644 --- a/src/utils/settings.rs +++ b/src/utils/settings.rs @@ -563,6 +563,18 @@ impl LoadWarnDedup { static LOAD_WARN_DEDUP: LoadWarnDedup = LoadWarnDedup::new(); +/// Warns that settings failed to load and defaults are being used, unless +/// `dedup` has already reported this exact failure. Takes the dedup as a +/// parameter so a test can assert the once-only behaviour on a private one. +fn warn_settings_fallback(dedup: &LoadWarnDedup, message: &str) { + if dedup.observe(Some(message)) { + tracing::warn!( + "{message}; falling back to default settings for this invocation — \ + any settings.json configuration is being ignored" + ); + } +} + /// Whether `key` is a registered secret, so it has a `_FILE` companion. fn is_secret_env_var(key: &str) -> bool { secret_env::SECRET_ENV_VARS.contains(&key) @@ -625,13 +637,7 @@ impl Settings { settings } Err(e) => { - let message = format!("{e:#}"); - if LOAD_WARN_DEDUP.observe(Some(&message)) { - tracing::warn!( - "{message}; falling back to default settings for this invocation — \ - any settings.json configuration is being ignored" - ); - } + warn_settings_fallback(&LOAD_WARN_DEDUP, &format!("{e:#}")); Self::default() } } @@ -654,12 +660,7 @@ impl Settings { /// Records a bootstrap settings failure after tracing has been installed. /// Shares deduplication with later settings reads in the same process. pub fn warn_bootstrap_failure(message: &str) { - if LOAD_WARN_DEDUP.observe(Some(message)) { - tracing::warn!( - "{message}; falling back to default settings for this invocation — \ - any settings.json configuration is being ignored" - ); - } + warn_settings_fallback(&LOAD_WARN_DEDUP, message); } /// Loads settings from a specific path. @@ -1540,6 +1541,64 @@ mod tests { assert!(logs.contains("settings.json"), "{logs}"); } + #[test] + fn load_daemon_reads_the_daemon_section_of_settings_json() { + let guard = crate::drive::test_support::EnvGuard::take(); + let dir = guard.clear_credentials(); + let settings_dir = dir.path().join(".omni-dev"); + fs::create_dir_all(&settings_dir).unwrap(); + fs::write( + settings_dir.join("settings.json"), + r#"{"daemon":{"log_level":"omni_dev::sessions=debug"}}"#, + ) + .unwrap(); + + assert_eq!( + Settings::load_daemon().log_level.as_deref(), + Some("omni_dev::sessions=debug") + ); + } + + #[test] + fn load_daemon_warns_and_falls_back_when_settings_json_fails_to_parse() { + let guard = crate::drive::test_support::EnvGuard::take(); + let dir = guard.clear_credentials(); + let settings_dir = dir.path().join(".omni-dev"); + fs::create_dir_all(&settings_dir).unwrap(); + fs::write(settings_dir.join("settings.json"), "{not valid json").unwrap(); + + let logs = crate::test_support::capture_at(tracing::Level::WARN, || { + assert!(Settings::load_daemon().log_level.is_none()); + }); + assert!(logs.contains("settings.json"), "{logs}"); + } + + #[test] + fn a_settings_fallback_is_warned_about_once_per_distinct_failure() { + let dedup = LoadWarnDedup::new(); + let logs = crate::test_support::capture_at(tracing::Level::WARN, || { + warn_settings_fallback(&dedup, "broken A"); + warn_settings_fallback(&dedup, "broken A"); + warn_settings_fallback(&dedup, "broken B"); + }); + assert_eq!(logs.matches("broken A").count(), 1, "{logs}"); + assert_eq!(logs.matches("broken B").count(), 1, "{logs}"); + assert!(logs.contains("falling back to default settings"), "{logs}"); + } + + #[test] + fn warn_bootstrap_failure_reports_the_failure_it_is_handed() { + // The bootstrap load in `main` fails before tracing exists, so it hands + // its error over afterwards. Repeat suppression is asserted above on a + // private dedup: doing it here would race every parallel test whose + // successful load resets the process-wide one. + let message = "Failed to parse settings file: /bootstrap-test/settings.json"; + let logs = crate::test_support::capture_at(tracing::Level::WARN, || { + Settings::warn_bootstrap_failure(message); + }); + assert!(logs.contains("/bootstrap-test/settings.json"), "{logs}"); + } + #[test] fn settings_get_env_var() { // Create a temporary directory (use current dir to avoid TMPDIR issues in tarpaulin) From 3978c87d347a6bd49c4e1728e5220f57fe584b33 Mon Sep 17 00:00:00 2001 From: John Ky Date: Fri, 2 Oct 2026 00:18:36 +1000 Subject: [PATCH 5/5] fix(cli): avoid deprecated fetch_update in wrapper counters Clippy 1.99 deprecates `AtomicU64::fetch_update` in favour of `try_update`, which failed the Clippy CI job under `-D warnings`. `try_update` is newer than the 1.88 MSRV, so use a `compare_exchange_weak` loop instead: same saturating semantics, no deprecation on any toolchain. Add tests for the saturation bound and for lost updates under contention (which also exercises the retry path). --- src/cli/claude_wrap/diagnostics.rs | 42 +++++++++++++++++++++++++++--- 1 file changed, 39 insertions(+), 3 deletions(-) diff --git a/src/cli/claude_wrap/diagnostics.rs b/src/cli/claude_wrap/diagnostics.rs index 4e11f77d9..9fc9ed580 100644 --- a/src/cli/claude_wrap/diagnostics.rs +++ b/src/cli/claude_wrap/diagnostics.rs @@ -131,10 +131,19 @@ fn drain(receiver: &mpsc::Receiver, out: &mut impl Write) { let _ = out.flush(); } +/// Saturating increment. A `compare_exchange` loop rather than `fetch_update` +/// (deprecated on Rust 1.99+) or its `try_update` replacement (newer than the +/// MSRV). pub(super) fn increment(counter: &AtomicU64) { - let _ = counter.fetch_update(Ordering::Relaxed, Ordering::Relaxed, |n| { - Some(n.saturating_add(1)) - }); + let mut current = counter.load(Ordering::Relaxed); + while let Err(actual) = counter.compare_exchange_weak( + current, + current.saturating_add(1), + Ordering::Relaxed, + Ordering::Relaxed, + ) { + current = actual; + } } #[cfg(test)] @@ -161,6 +170,33 @@ mod tests { ); } + #[test] + fn increment_saturates_instead_of_wrapping() { + let counter = AtomicU64::new(u64::MAX - 1); + increment(&counter); + increment(&counter); + assert_eq!(counter.load(Ordering::Relaxed), u64::MAX); + } + + #[test] + fn increment_loses_no_updates_under_contention() { + let counter = Arc::new(AtomicU64::new(0)); + let threads: Vec<_> = (0..8) + .map(|_| { + let counter = Arc::clone(&counter); + std::thread::spawn(move || { + for _ in 0..1_000 { + increment(&counter); + } + }) + }) + .collect(); + for thread in threads { + thread.join().unwrap(); + } + assert_eq!(counter.load(Ordering::Relaxed), 8_000); + } + #[tokio::test] async fn disabled_sink_does_not_evaluate_records() { let (sink, done) = Diagnostics::open(None);