Repository navigation
Faster request path, history and macOS file watching; incremental transcripts and dry-run shares, protocol 41 - #298
Merged
Merged
Conversation
…onversations The desktop app's dry-run dialog worked out each member's share itself (weight x balance factor, skipping tripped upstreams). That is the data plane's algorithm copied into the UI, and the copies had already drifted: it ignored members whose concurrency limit is full, which the engine leaves out of the round, and it multiplied the raw factor where the engine rounds the effective weight to thousandths with a floor of one. tw-engine gains `weighted::shares`, built on the same `round` that picks the next leader: each member's effective weight over the round's total, None for members sitting the round out (paused or busy, unless all are). Smooth weighted round-robin leads each member exactly its effective weight times per cycle, so this is the long-run split the data plane produces from the current numbers. `tw_api::DryRunCandidate` gains `share: Option<f64>` (0..=1), present only for load-balance groups; the shares of one group sum to 1. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Over the size cap, GC deleted whole days oldest-first and did not stop at the newest one. A single day can exceed body_max_bytes (each request keeps up to 4 MiB of request and 4 MiB of response), and then every hourly pass wiped today's directory, including the session being viewed and the request that had just finished. Older days still go whole. Once only the newest day is left over the cap, its oldest requests are removed, all files of one request together (request id order is start order), until the total fits. Bodies are now written to a temporary file and renamed into place, with the `.len` marker written first. Readers (detail, transcript, search) do not take the recorder's lock, so they could read a body halfway through being written and take it for an unreadable or complete one. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Measured on a synthetic 270k-row, 222 MB request database (90 days, release build; `history_cost` reproduces it): - sessions(None, 200): 160-200 ms -> 12 ms. It grouped every row of the retention period and sorted the groups. It now walks a partial index (at_ms, session) newest-first until it has `limit` sessions, checks a session's first turn only when its last one is after the window, and aggregates just those sessions through the session index. A test compares it with the old grouping query across windows and limits. - Db::session(id): direct lookup for the session detail (was: take the newest 500 sessions and search them). - last_seen_by_client: 78 ms -> 0.02 ms, a loose index scan over a partial (client, at_ms) index. The index is partial (`client > ''`) so the windowed GROUP BY client queries keep walking the time index instead of scanning it in key order. - route_hits: 22 ms -> 6 ms, and it no longer reads 7.6 MB of routing JSON: route, rule, rewrites and denial live in an expression index (json_extract, written at insert), so the query reads only the index. Rows whose routing does not parse stay out of it and are still written. The indexes are created with CREATE INDEX IF NOT EXISTS on every open, not through SCHEMA: they do not change what a row looks like, other versions read and write the database as before, and a SCHEMA bump would wipe everyone's request history. First open of an existing 270k-row database builds them in about a second; the file grows ~15%. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
… the async workers
GET /sessions/{id} took the newest 500 sessions and searched them, so
an older session listed further down answered 404. It now asks for that
session directly (Db::session).
The session list, session detail, route stats, storage status and the
overview's body total now run on the blocking pool (`on_store`): the
queries and the walk over every stored body file used to run on a
runtime worker, and the body walk did so while holding the recorder's
lock. The body total no longer takes the lock at all; the body
directory is only a path.
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
…es without the lock The recorder ran as tokio tasks doing synchronous work on runtime workers under one tokio mutex: SQLite inserts per event, body writes of up to 4 MiB plus chmods, plugin-run inserts. The hourly GC held that same lock while deleting day directories, which can take seconds. - Events, bodies and plugin runs are each handled on a named OS thread (tw-recorder, tw-bodies, tw-plugin-runs). The event thread takes the lock once per batch of up to 64 events. Body writes no longer take the lock at all; the body directory is only a path. - The recorder no longer subscribes to the UI's 1024-slot broadcast, where lagging meant those requests were missing from history for good. EventBus::record_feed gives it its own bounded channel (16384 slots, room for ~2000 requests of backlog). emit() try_sends into it, so forwarding still never waits; when it is full the event is dropped and counted, and the recorder logs how many. - task::gc deletes expired and over-cap bodies without the lock, then takes it only to prune rows; twcore runs it on the blocking pool. Recorder::gc became Recorder::prune (rows only). - key_limits::follow keeps the runtime handle it was attached on: settling now happens on the recorder thread, where Handle::try_current() fails, so a period change found while settling would never have been read back from the records. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
notify's default macOS backend is FSEvents, which always watches the whole tree below a path. `RecursiveMode::NonRecursive` only drops the deeper events inside our process, after the kernel has already woken us, with zero latency and one event per file. The desktop app watches `~` itself (that is where `~/.claude.json` lives): an audit counted 7-14 raw events a second under `~` while idle, none of them relevant, and writes under `~/Library/Caches` cost the app measurable CPU. tw-watch also backs core's config and plugin hot-reload. On macOS tw-watch now does it with kqueue: - each watched directory gets an O_EVTONLY vnode watch, so only its own entries being added, removed or renamed wake us; on a wake it lists the directory and fires when a relevant entry appeared, disappeared or was replaced (inode changed); - each relevant regular file (or symlink to one) also gets a watch, for in-place writes, appends and truncation, which do not touch the directory; after an atomic replace the directory listing re-arms it on the new inode; - attribute-only changes count only when size or mtime moved (truncate reports nothing else; chmod and reads do not count); - a watched directory that is removed or renamed away counts, and its parent is watched until the name reappears, so a replaced or recreated directory keeps being watched, as with FSEvents. notify's own kqueue backend (`macos_kqueue`) is not usable here: on the first change in a non-recursively watched directory it treats the first entry it has not seen as new and watches it recursively, which for `~` means opening all of `~/Library`. Each watch costs a file descriptor and apps started from Finder get a soft limit of 256, so all watches in a process together hold at most half the soft limit (capped at 1024). A watch that would exceed it falls back to FSEvents as a whole, at start or later, rather than silently missing in-place writes. Measured with a probe binary (release build, 200 ms debounce): - 1,500 writes 3 levels below a watched directory over ~9 s: before 51-54 ms CPU and ~890 context switches, after 0.4-0.5 ms and 9-13; both reported the one relevant change afterwards. - watching `~` for 30 s idle, both probes side by side: before 8.8-11.1 ms CPU and ~233 context switches, after 0.1-0.2 ms and 3. Linux (inotify) and Windows (ReadDirectoryChangesW) keep using notify unchanged. The public API is unchanged. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Every GET /sessions/{id}/transcript re-read and re-parsed every stored
body of the session, and the app asks again on every new turn: for a
200-turn session of ~1.5 MB requests that is ~300 MB per turn. Measured
with `transcript_cost` (300 turns, 367 MB of bodies, release build):
the whole transcript took ~830 ms; with one new turn it now takes
6.4 ms, and asking from `settled_turns` takes 0.07 ms and sends 27 KB
instead of 1.5 MB.
tw-store: transcript::Cache (one per Recorder, Recorder::transcripts)
keeps, for the last four sessions read, the request ids read, the chain
state after them and the unmasked turns, and continues from there when
those ids are still the session's first rows. A middle insertion or a
pruned head fails that check and reads from the start. Only the settled
prefix is kept: a turn that ended under two minutes ago with a gap may
still get its body, so it and everything after it are read again. GC
invalidates the cache whenever it deleted a body (a generation counter
keeps a read that raced the GC from being stored).
tw-api, protocol addition (no version bump here):
- SessionTranscript takes TranscriptQuery { from_turn?: u32 }. Present:
`turns` holds only the turns with index >= from_turn (empty past the
end); absent or 0: the whole transcript, as before.
- Transcript gains `total_turns` (number of turns in the session) and
`settled_turns` (leading turns that will not change any more: bodies
final, not rewritten by a later turn's tool-call ids, and no running
request of the session will land before them). The app keeps its
first `settled_turns` turns and asks again with that as from_turn.
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
The local-answer check looks for five markers in every request body before deciding whether to parse it, and it did so with `windows().any()`: about 16 ms of synchronous CPU for a 2 MiB Claude Code request, on the async worker, before anything else happened. `strip_carried` did the same on every same-format hop, and the usage sniffer on every streamed chunk. memchr is already in the tree through serde_json and regex, so adding it as a workspace dependency compiles nothing new. memmem gives the same first match: 0.26 ms for the five markers over 2 MiB, 0.04 ms for `strip_carried`'s. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
The redaction scan runs over every request body even in observe mode (22 ms per 2 MiB, on the request path and again before storage). About a quarter of it was SipHash: three rules were looked up in the enabled set for every token. Another sixth was memcmp: every token was compared with twenty-odd key prefixes. The rest was the tokenizer decoding the body character by character. - Ask whether the per-token rules are on once per scan, as `want_jwt` already did. - Reject a token before the prefix loop when it is shorter than the shortest key any rule can match or starts with a byte no prefix starts with. Both bounds are computed from BUILTINS at compile time, and a test checks the shortcut against the full comparison for every rule set. - Walk the body by byte. Tokens, backslashes and quotes are all ASCII and every byte of a multi-byte character is >= 0x80, so cutting by byte gives the same tokens; a test compares it with the old char-based cutter on random text with escapes, CJK, emoji and combining marks. - Find PEM headers and `://` with memmem. - `look` in enforce mode built a fully replaced copy of the body only to keep its ledger. Number the values without writing the copy. `look_hits` also hands back the hits it found, so a caller that scans the same bytes again (the desktop gateway, before storing the body) can reuse them. Results are unchanged: 2 MiB Claude Code-shaped body, 29-54 ms -> 3.7 ms; 8 MiB, 71-92 ms -> 15 ms. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
With the default rules (two code-point rules, three phrases, observe mode) the content filter cost 13-16 ms of a 2 MiB Claude Code request besides parsing it. Claude Code's file reads put `→` after every line number, so nearly every tool result is non-ASCII and missed the existing all-ASCII shortcut: - A code-point rule whose code points are all above U+007F cannot match an ASCII byte. Skip ASCII eight bytes at a time and decode only the other characters. - The lowercase copy for phrase rules came from `str::to_lowercase`, which goes character by character after the first non-ASCII one. Lowercase ASCII runs in bulk and the other characters one by one; `Σ`, the only context-dependent mapping, hands the whole text back to the standard library. - Find phrases with memmem, non-overlapping like `match_indices`. Tests compare each against the old way on random text. Content filter on a 2 MiB body: 15-18 ms -> 4.8 ms including the parse (1.9 ms). Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
A Claude Code request body was parsed three times before its first hop: for the routing facts, for the content filter, and again on every hop to drop `metadata.user_id` -- which Claude Code always sends and `forward_client_identity` is off for by default. The local-answer check parsed it a fourth time whenever the body mentions `haiku`, which Claude Code's Task tool schema lists among its models. The pipeline now parses the body once and hands the value to the local-answer check, the routing facts and the content filter (`tw_guard::content::screen_value`, which only copies the value when it has to strip something). The identity is cut out of the body bytes instead of parsing and writing it back. Besides the cost (4-6.6 ms and a full copy per hop at 2 MiB), writing it back was a fidelity bug: serde_json without `preserve_order` sorts every object's keys, so every Claude Code request reached the upstream with its tool schemas, tool inputs and top level reordered, which also changes the bytes an upstream's prompt cache keys on. Now the forwarded body is the client's body less that one member (and its comma, or the whole `metadata` once it is empty). Keys are compared decoded, as serde_json does. A key written twice falls back to the old parse and write, since serde_json keeps the last one. The cut relies on the body being valid JSON, which the pipeline already knows from its parse (`Reading::json`); a body that does not parse is still forwarded untouched. A random test checks that the cut means what the old rewrite meant, and a gateway test that a 1.3 MB Claude Code body reaches the upstream byte for byte less `metadata`. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
The request body was scanned for secrets on the request path and then, for storage, scanned again in full on the blocking pool (the scan plus `mask_body` was 29 ms per 2 MiB). In enforce mode each hop scanned again too, even when the hop sends exactly the client's bytes. - The hits from the request-path look travel with the stored request body (`BodyRecord::found`). They are used only when the record holds the whole body and it is UTF-8, so the text masked before storage is the text that was scanned; a body cut to the storage window is scanned again as before. - In enforce mode a hop whose body is the very buffer the client sent (same-format passthrough, nothing rewritten) is replaced from those hits (`guard::replace_found`). Any other hop body -- translated, rewritten by a rule or plugin, identity cut -- is scanned as before. Same rules, same bytes, same hits: tests check both paths against a fresh scan. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Parsing, the local-answer check, the routing facts, the content filter, the secret scan and each hop's rewrite all ran inline on the tokio worker that owns the connection. Every other connection on that worker -- a streaming answer being relayed -- stopped for as long: with one worker, SSE gaps of 470-610 ms while 8 MiB requests came in. For bodies of 1 MiB or more these steps now run in `block_in_place`, which hands the worker's other tasks to another thread for the duration. Smaller bodies stay inline, where the hand-off would cost more than the work. On a current-thread runtime (the tests) `block_in_place` panics, so it is used only on a multi-thread one. With one gateway worker and 8 MiB requests arriving alongside an SSE relay, the longest gap between relayed chunks fell from 470-610 ms to 17-39 ms, and the gateway took 45-55 such requests in the time it used to take 8-12. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
When a routing rule (or an alias rename) changes the model, the output limit or turns thinking off, `forward::apply_set` parsed the body and wrote it back whole. serde_json without `preserve_order` sorts every object's keys, so the upstream got the client's tool schemas, tool inputs and top level reordered -- a different byte string for its prompt cache on every such request, the very cache loss the module warns about. The rewrite now diffs the parsed body before and after the change and edits the original text: a changed value is replaced in place (objects on both sides are compared one level further in, so Gemini's `generationConfig` keeps its other keys and their order), a removed member is cut with its comma, an added one is appended to its object. Everything else -- order, whitespace, escapes -- is the client's. A rewrite that changes nothing sends the original buffer. Only when a key is written twice in an object the change reaches does it fall back to writing the body back, since serde_json keeps the last one. The member scanner moves from `egress` to a shared `splice` module. Why not turn on `preserve_order`: it is a feature of the shared serde_json, so it would also reorder every JSON the gateway, the control-plane API and ThinkWatch Enterprise produce (tw-dialect and tw-guard inherit it), and `Map::remove` would become `swap_remove`, moving the last key into the removed one's place -- removing `thinking` would still reorder the body. The splice keeps the change on the one path that rewrites, which is off by default. Tests: the exact bytes for an Anthropic body with a nested tool schema, an appended limit and Gemini's nested config; a random test over five formats and random rewrites checks that the result means what writing it back meant and that untouched members keep their bytes and order. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
…est skips bodies still being written fetch_update is deprecated on the current stable toolchain. Bodies are written to a .tmp file and renamed into place, so a directory listing can see a file that is gone by the time it is read (seen on Windows). Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Merged
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
A performance and robustness round across the request path, request history, sessions and file watching, plus two bugs found on the way. Control-plane protocol 40 → 41.
Request path (default observe mode, Claude Code headers, client send → upstream handler, release build)
strip_carriedand the usage sniffer search withmemchr::memmeminstead ofwindows().any()(16 ms → 0.26 ms per 2 MiB).lookno longer builds a replaced copy only to get the ledger (flow::look_hits), and the storage copy reuses the request-path hits. Equivalence tests compare old and new tokenization on random text with escapes, CJK, emoji and combining marks.content::screen_value). The content filter skips ASCII when matching code points and finds phrases with memmem (15–18 ms → 4.8 ms per 2 MiB).block_in_place, so streams relayed on the same worker no longer stall: the longest gap between relayed chunks during 8 MiB requests fell from 470–610 ms to 17–39 ms.metadata.user_id(every Claude Code request) and a routing rule or alias changing model, max_tokens or thinking re-serialized the body with every object's keys sorted, tool schemas and tool inputs included. Both now splice the bytes (splice.rs); only a duplicate key in a touched object falls back to the full rewrite. A random test over all five formats checks the splice means what the old rewrite meant.Request history and sessions
body_max_bytes, the hourly pass deleted the whole current day, including the session being viewed. Older days are still deleted whole; the newest day is trimmed from its oldest requests. Bodies are written to a temp file and renamed into place.last_seen_by_clientuses an index (78 ms → 0.02 ms); route stats read an expression index instead of parsing routing JSON (22 → 6 ms). Indexes are created withCREATE INDEX IF NOT EXISTSon open: no SCHEMA bump, history is kept; the first open of a large database spends about a second building them.File watching on macOS
~woke the process for every file change in the home folder. Non-recursive watches now use kqueue (kqueue.rs): a directory wakes only when its own entries change, relevant files get their own watch for in-place writes, and atomic replace, create-later and delete-recreate are followed. If a watch would exceed half the soft fd limit (capped at 1024) it falls back to FSEvents. Lite's real watch plan idle for 30 s: 105 wakeups → 4. Linux and Windows are unchanged.Protocol 41
GET /sessions/{id}/transcripttakesTranscriptQuery { from_turn };Transcriptgainstotal_turnsandsettled_turns(leading turns that will not change). A UI written for 40 that sendsnullis refused.DryRunCandidate.share: the candidate's share of a load-balance group as the engine distributes it (tw_engine::weighted::shares, on the sameround()the router uses);Noneoutside load-balance groups and for members sitting out.Enterprise: tw-dialect and tw-guard gain a direct
memchrdependency (already in their tree through serde_json). Public API changes are additions only (flow::look_hits,Look,content::screen_value);scan,scan_text,redact_text, the content filter andstrip_carriedgive the same results, covered by equivalence tests.Checks:
cargo fmt --check,cargo clippy --workspace --all-targets -D warnings,cargo test --workspace(3068 passed),cargo test -p tw-api --features tswithtsc --stricton the export,scripts/smoke.sh75/75, tw-watch clippy for Windows and Linux targets.No version bump or tag; release is held.
🤖 Generated with Claude Code