Repository navigation
fix(gateway): replay the transcript spool in drop order with full fidelity on restart (salvage #122718) - #133932
Conversation
… on restart
recover_pending_to_db walks sorted(flush_dir.glob("*.json")) to replay
spool files left behind by a restart. Spool files are named
pending-<uuid4>.json by _write_payload, so that sort is effectively
random: a transcript recovered after a crash comes back scrambled.
The damage is permanent rather than cosmetic. SessionDB restores a
conversation with "ORDER BY id" — AUTOINCREMENT insertion order, never
timestamp — precisely so a non-monotonic clock cannot sort an assistant
tool_calls row after its tool response, "breaking tool-call/response
adjacency and triggering an HTTP 400 on replay". Replaying the spool in
filename order writes exactly that inversion into the row ids, so the
session errors out on the next turn.
The payloads already carry the ordering fields: spool_dropped_transcript_
message stamps every one with "ts" and a monotonic "seq". The sibling
consumer of the same spool in this file, drain_transcript_spool, already
sorts on them via sorted(entries, key=lambda e: e[:3]); only the restart
path ignored them. Order on (ts, seq, filename) so the live drain and the
cross-restart drain agree.
Ordering by the payload rather than by a new file-naming scheme also
recovers spool files that were written before this change, which are the
ones an affected user already has on disk.
Payloads that cannot be parsed keep their previous treatment: they sort
last, are handed back unparsed, and are re-read inside the loop so the
existing handler reports and preserves them unchanged.
(cherry picked from commit 0e46604)
(cherry picked from commit c043f01)
…pooled messages The cross-restart replay of a cap-dropped transcript message forwarded only four fields to append_message — session_id, role, content, timestamp — and then unlinked the spool file, making the loss permanent. Everything else on the message was discarded: tool_calls, tool_call_id, tool_name, the reasoning columns, codex items, platform_message_id, observed, and the api_content sidecar. append_message already accepts all of them, and the payload already carries them, because spool_dropped_transcript_message writes the full transcript message dict. Two of the dropped fields are load-bearing rather than decorative. Losing tool_calls/tool_call_id orphans a tool result from the call it answers, which is the adjacency SessionDB's "ORDER BY id" exists to protect. Losing api_content contradicts the requirement stated at the live writer, that the sidecar "must survive any gateway-side persistence path or the next turn's replay diverges at this row" — and this is such a path. The same spool is drained two ways: by drain_transcript_spool during live operation, which replays through SessionStore._append_transcript_message and keeps every field, and by this function after a restart. Same files, same payloads, different fidelity. Mirror _append_transcript_message field for field so the outcome no longer depends on whether the gateway happened to restart, including its role gate on the assistant-only reasoning columns and its message_id fallback for platform_message_id. Fields are whitelisted explicitly rather than splatted from the message dict: the payload is arbitrary JSON from disk, and an unexpected key would raise TypeError and abort the whole recovery pass. content is now passed through as-is instead of being coerced to "", since an assistant tool-call row legitimately has no content. (cherry picked from commit 99ef8fb) (cherry picked from commit 27ee00d)
…pend Ordering the spool files fixes the happy path only. If append_message fails partway through — the DB is still unhealthy, which is the situation that produced the spool in the first place — the loop logs the failure, keeps that file for a later retry, and then carries on and writes the messages that come after it. Those later messages land now; the failed one lands on some future start. Because SessionDB orders a conversation by AUTOINCREMENT id, the retried message then gets a HIGHER row id than the messages it originally preceded, and the inversion this pass just prevented is written to disk anyway — this time permanently, since both files are gone. drain_transcript_spool, the live drain of the same spool, already states the rule: "On the first replay failure the drain stops and remaining files are kept for the next attempt (the DB is likely still unhealthy)." Apply it here too. The block is per-session rather than global. Replay order is only defined within a session, and this function drains every session's spool in one pass, so a single unhealthy session must not strand the others' messages on disk. Non-transcript pending payloads are unaffected. (cherry picked from commit 5ea5104) (cherry picked from commit a959154)
Names the spool files so that filename order is the exact reverse of drop order, which is what uuid4 names produce on average, and asserts the replay comes back in drop order. A second case gives three payloads the same one-second ts so only the monotonic seq can separate them. Both fail on the previous sorted(glob()) implementation. A third case pins the pre-existing treatment of a corrupt payload — reported by the loop's own handler and left on disk — so the new ordering pass cannot silently swallow one. Lives in its own file rather than tests/gateway/test_shutdown_flush.py to keep the restart-recovery cases together with the spool fixtures they need. (cherry picked from commit 555d871) (cherry picked from commit c372b25)
…overy Round-trips an assistant tool-call row carrying every field the live writer persists and asserts each one reaches append_message: tool_calls, tool_call_id, tool_name, the reasoning and codex columns, platform_message_id, observed, timestamp, and the api_content sidecar. Also pins content=None passing through uncoerced, since an assistant tool-call row has no content, and the role gate that keeps the assistant-only reasoning columns off a user row. A separate case covers a failed replay: with the first of two messages for one session rejected, neither is written and both spool files stay on disk, while a third message for an unrelated session still recovers. All but the two behaviour-preservation cases fail on the previous implementation. (cherry picked from commit 869acc6) (cherry picked from commit e02d9fc)
… absent
_transcript_append_kwargs chose the spooled message's timestamp with
`message.get("timestamp") or payload.get("ts")`, carried over from the
code this PR replaces. Truthiness is the wrong test: epoch 0 is a valid
timestamp and would be silently rewritten to the payload's ts, which is
the fidelity loss this PR exists to remove.
Fall back only when the field is genuinely absent.
(cherry picked from commit 62206b0)
(cherry picked from commit a7b6828)
…uilder Review follow-ups (teknium1): - One kwargs builder, two callers. transcript_append_kwargs() in gateway/session_transcript.py builds the append_message kwargs, and both _append_transcript_message and the spool's restart replay call it, so the restart writer cannot lag the live writer again (display_kind and display_metadata were added after NousResearch#84785 was written). fallback_ts applies only to a message with no timestamp; the live writer passes none, so its rows are byte-identical. - A failed replay says what it held back: one INFO line per pass with the file count per blocked session, so a spool dir that never empties has a reason in the log. - seq is per process and ts has one-second resolution; the sort key now says so where it reads seq. - Tests fold to three invariants: drop order beats name order (seq breaks a same-second tie), a replayed row is the row the live writer writes, and a failed replay holds back that session's later messages only. The two that already passed on main (a corrupt file is kept; the payload-clock fallback) are dropped; the fallback is asserted inside the fidelity test. (cherry picked from commit e2bacb0)
_order_flush_files runs before recover_pending_to_db's per-file error handling, so float(10**400) raising OverflowError stopped every file from being recovered. Unusable ordering values (too large for a float, NaN, Infinity) now sort first like any other non-numeric value. (cherry picked from commit 50c9b97)
…iles The code-health gate (BLE001) rejects the blind except added to _order_flush_files. A file that cannot be read or parsed still sorts last with a None payload, so malformed-file isolation is unchanged.
BearHuddleston
left a comment
There was a problem hiding this comment.
Reviewed at e2154ae536. The drop-order sort and the shared row builder work as described. I spooled 12 files through the real spool_dropped_transcript_message into a real SessionDB in a temp HERMES_HOME. main restored them in filename order with tool_calls, tool_call_id and reasoning NULL. This head restores drop order and keeps those fields. test_shutdown_flush_recovery.py fails 3 of 4 on the base and passes 4/4 here. The six spool/flush test files pass (28 tests) with scripts/run_tests.sh.
One gap in the hold-back; see the inline comment.
teknium1
left a comment
There was a problem hiding this comment.
PR Review — #133932 (head e2154ae5365)
Verdict: correct, minimal, ready for the maintainer's merge call. This is #122718 (@jonpol01's salvage of @briandevans's #84785) cherry-picked onto current main with authorship kept, plus one BLE001 follow-up. Every ask from the earlier reviews on #122718 (shared row builder, held-back log line, folded tests, seq-is-per-process comment, _sort_number overflow/NaN guard) is in this head; the only new delta is the narrowed except (OSError, ValueError, RecursionError) in _order_flush_files, which keeps the same "unparseable sorts last, caller re-reads and reports" behaviour.
Premise re-confirmed on origin/main (0dbaf33f67a): gateway/shutdown_flush.py::recover_pending_to_db still walks sorted(glob("*.json")) over pending-<uuid4>.json names and _recover_one_payload still appends only role/content or ""/timestamp, while drain_transcript_spool sorts on (ts, seq, name) and the live writer persists the full row. The two consumers of one spool disagree; this PR makes them share transcript_append_kwargs.
Live repro (real SessionDB, throwaway HERMES_HOME, three cap-drop files named to sort opposite to drop order, user → assistant(tool_calls, reasoning) → tool):
| rows by AUTOINCREMENT id | tool_calls / tool_call_id / reasoning | |
|---|---|---|
| main | tool, assistant, user |
all NULL, assistant content '' |
| this head | user, assistant, tool |
intact; assistant content NULL as the live writer writes it |
Tests: tests/gateway/test_shutdown_flush_recovery.py — 3 of 4 red with main's shutdown_flush.py + session_transcript.py swapped in (order, fidelity, hold-back), 4/4 green on the head. Merged tree (merge-tree of head onto 0dbaf33f67a, clean): 7 spool/flush/transcript files, all green (the 3 reds in test_shutdown_inflight_transcript_flush.py / test_load_transcript_db_only.py in a .worktrees/ checkout are the known manifest.json home-I/O guard trip, identical on pristine main, green with a fake HOME). python -m scripts.code_health --base origin/main on the merged tree: 0 blocking. ruff clean, git diff --check clean.
Non-blocking observations
- ◇ The hold-back applies to cap-drop payloads (keyed by
session_id) only. A failedappend_messagefor a shutdown pending-text payload still lets that routing key's later files land first. Different payload kind and pre-existing; worth a follow-up, not this PR. - ◇
blocked_sessions[sid]counts the failed file plus every later one, so "Held back 2" for one failure + one later file. Reads fine, just noting the arithmetic matches the test's expectation on purpose. - ◇ Overlap unchanged from #122718's thread: #117310 edits the same
_recover_one_payloadbranch (owning-store routing); whichever lands second rebases. #132680 (@Sahilvishnaliya, Oct 4) re-implements the same three fixes; this PR credits the earlier work, so #132680 closes as redundant when this lands, together with #122718 and #84785.
Priority: agree with the body's P1 reading (DB outage + restart to trigger; permanent damage when it does), which matches the maintainer's earlier call on #123658. Merge decision stays with @teknium1.
…ext live write recover_pending_to_db held back a session's later spool files after a failed replay, but never told the running gateway. The live writer drains the spool only for sessions in SessionStore._spooled_drop_sessions, which only this process's own spooling filled, so the session's next live row landed ahead of the held-back files until the next restart. recover_pending_to_db now adds the sessions it held back to a caller-supplied set, and _recover_pending_flushes passes the store's own drain set (exposed as spooled_drop_sessions()), so the next live write replays them first. Found in review by @BearHuddleston.
|
Re-verified at Live, real
|
keeltrace
left a comment
There was a problem hiding this comment.
AI-assisted follow-up review on current head 36f27f287f0.
One still-novel sibling-path ordering gap remains in the ordinary shutdown FIFO path. This is separate from BearHuddleston's held-back/live-write finding (fixed in this head) and from the failed ordinary-payload append case Teknium already noted.
flush_pending_to_file() writes the adapter's pending-slot head with ts but no seq. Shutdown then writes that same session's FIFO tail with flush_overflow_to_file(), whose first item gets seq=0. When both flush calls happen within the same int(time.time()) second — the normal case — _order_flush_files() maps both keys to the same (ts, 0.0) because _sort_number(None) == _sort_number(0) == 0.0. The tie is then broken by pending-<uuid4>.json, so the first overflow item can replay before the slot head.
That violates the contract already documented in flush_overflow_to_file: the adapter slot is the queue head and queued_events is its tail. Since recovered rows are ordered by AUTOINCREMENT id, a successful restart replay can permanently persist B, A instead of A, B even with no DB failure at all.
A focused regression would freeze time.time(), flush one pending-slot event A plus overflow B/C for the same session, make the filenames sort adversarially, then assert recovery is A/B/C. The production fix needs a shared per-session ordering domain across the slot and overflow — e.g. give the slot head a sequence that sorts before overflow seq 0, rather than letting both normalize to zero.
I searched the current #133932 discussion plus the related #100972 / #105270 work; I did not find this head-vs-first-tail tie raised elsewhere. Given Teknium's scope call on the adjacent ordinary-payload failure case, I'm leaving this as a COMMENT rather than blocking #123658's cap-drop fix.
keeltrace
left a comment
There was a problem hiding this comment.
AI-assisted follow-up review on current head 36f27f287f0.
One additional reliability issue from the new restart sorter: _order_flush_files() now materializes the entire recovery spool in memory before replay begins.
Before this change, recover_pending_to_db() kept a sorted list of paths and parsed/processed one JSON payload at a time. _order_flush_files() instead json.loads() every *.json file and retains each full payload in entries until all files have been read and sorted.
That changes peak recovery memory from roughly O(largest payload + path list) to O(total bytes in the spool).
The transcript spool itself does not appear to have a file-count or byte-retention bound. _MAX_PENDING_PER_SESSION = 200 only bounds the in-memory retry queue; once it is exceeded, each additional oldest message is written to pending_messages/. During a prolonged DB outage — exactly the failure mode this restart recovery handles — the directory can therefore accumulate an arbitrarily large set of full transcript messages, including potentially large tool/result payloads.
That means a large but otherwise valid backlog can make restart recovery exhaust memory or spend substantial memory materializing the complete backlog before it has recovered even one row. The old implementation did not have that aggregate-memory behavior.
I would keep only the ordering metadata plus path in the first pass, sort those small records, and re-read/parse each payload when it is actually replayed. That costs a second read but restores bounded payload memory. A regression can generate many large spool payloads and verify the ordering pass does not retain their full decoded objects simultaneously.
I searched the current #133932 discussion plus #122718 / #84785 and adjacent spool-recovery work and did not find this memory-shape issue raised previously.
This is independent of the FIFO head/tail ordering comment already left on this PR.
Retain only sort keys and paths; re-read payloads one at a time for replay. Regression with 6 MiB of spool data fails on the previous head at 8.1 MiB peak allocations and passes with the metadata-only index. Review finding by @keeltrace.
Treat an absent sequence as -1 during restart sorting, before the overflow tail starting at zero. This also repairs ordering for existing spool files, without requiring a new writer format. Frozen-clock adversarial-filename regression reproduces B,A,C before the fix and A,B,C afterwards. Review finding by @keeltrace.
dokterdok
left a comment
There was a problem hiding this comment.
Reviewed at 36f27f287f0. The drop-order sort, the shared row builder and the held-back hand-off work as described: test_shutdown_flush_recovery.py and test_shutdown_flush.py pass, 16/16. I found one more ordering gap. It is separate from @BearHuddleston's held-back case and from @keeltrace's same-second tie.
P1: boot recovery runs after the startup turns, so a resumed session's new rows land ahead of its spooled rows
The spool replay only starts once runner.start() has returned (run.py L5964-L5977). By then, start() has already:
- run the boot auto-resume turns (
run_startup.pyL1528-L1529); - drained the inbound queue it held during startup (L262).
Those turns write through append_to_transcript. The live writer drains the spool only for sessions in _spooled_drop_sessions (session_transcript.py L374-L378), and that set is empty in a fresh process. So the turn's rows are written first, and the previous process's spooled rows are replayed after them, with higher row ids.
The sessions most likely to have spooled rows are the very ones auto-resumed at boot: their turn was running during the outage, and their backlog was spooled at shutdown (#114266).
Repro on this head, with the same harness as test_held_back_files_drain_before_that_sessions_next_live_write (the store's real transcript path, then _recover_pending_flushes):
- Spool
old0(ts 100, seq 0) andold1(ts 100, seq 1) forsess-1, as the previous process would. - Call
store.append_to_transcript("sess-1", {"role": "user", "content": "resumed-turn"}), as the resumed turn does insidestart(). - Run
_recover_pending_flushes(runner), asmain()does afterstart()returns.
The result is ['resumed-turn', 'old0', 'old1'], permanently. Registering sess-1 in spooled_drop_sessions() before step 2 gives ['old0', 'old1', 'resumed-turn'].
This is not a regression, since main has the same startup order. But it is the same damage this PR fixes, and the PR's ordering guarantee only holds if recovery runs before a session's first live write.
Either of these closes it with a small change:
- Run
_recover_pending_flushesinsidestart()just before_schedule_resume_pending_sessions(). Inbound is still queued at that point and no turn has run, so the replay always comes first. - Or, at that same point, register every session that has cap-drop files on disk in the store's drain set (per profile home when multiplexed). The first live write for such a session then drains its files in order, and the later boot pass finds nothing left for it.
The regression test is the repro above: a live write that happens before _recover_pending_flushes must still end up as old0, old1, <live>.
Found by reading only, not reproduced: nothing serializes the boot pass with the live drain, because recover_pending_to_db doesn't take the store's transcript drain lock. If this process also spools for a session while starting (the database is still down at boot), the boot pass and that session's live drain can read the same file and both append it. Running recovery before the resume turns, as in the first option, removes this window too.
ehz0ah
left a comment
There was a problem hiding this comment.
Request changes conclusion on head 36f27f287f0.
This PR adds durable recovery for transcript rows that could not reach the database. gateway/shutdown_flush.py replays disk spool files in a stable order and keeps the full row shape. gateway/run.py records held-back sessions in the live store. gateway/session_transcript.py then drains that backlog before the next live write. The new tests cover the normal restart and successful handoff paths.
The main path works, but a second drain failure still lets a newer live row reach the database first. I reproduced this on the reviewed head and on its clean composition with current main. The attempts were old0, old0, then new-live, and the old spool remained. The five submitted recovery tests pass. All required checks are green.
The inline comment describes the missing failure path. Please add that case before merge.
…s fail to drain The live writer replayed a session's spool before each write but ignored a failed replay and wrote the new row anyway, so a second failure (boot hold-back, then the pre-write drain) gave the live row a lower row id than the older spooled rows. _drain_spooled_drops now reports files left on disk, and the write takes the existing retry path instead, leaving the row queued in memory behind them. Review finding by @ehz0ah.
…ueued inbound Boot recovery ran in main() only after runner.start() returned, by which point start() had already run the auto-resume turns and drained the inbound queued during startup. Those turns write live transcript rows, and a fresh store does not know which sessions have spooled rows, so the previous run's spool landed after them. _start_finish_wiring now runs _recover_pending_flushes just before _schedule_resume_pending_sessions(), and failures log at WARNING instead of being swallowed silently. Review finding by @dokterdok.
… a live write Gating the live write on a pending spool raised a generic RuntimeError, hiding the replay's real failure from the FTS-rebuild, replaced-database and compression handlers, so a repairable FTS fault became a permanent stall. _drain_spooled_drops now returns the replay's own exception, the caller raises it, and after a successful FTS rebuild the loop drains the older spool again before writing the live row. Found in independent review.
unsupportedpastels
left a comment
There was a problem hiding this comment.
Follow-up pushed: 36f27f287f0..376d2e134eb.
@keeltrace, both findings are fixed:
- Memory (d4d5a8e):
_order_flush_fileskeeps only the sort key and path, andrecover_pending_to_dbre-reads each payload when it replays it.test_recovery_payload_memory_does_not_scale_with_backlogspools 6 MiB and asserts peak traced allocations stay under 4 MiB. Before the fix the peak was 8.1 MiB. - Slot head vs. first overflow item (4d7d3c0): a payload with no
seqnow sorts as -1, ahead of overflow seq 0. The change is on the reading side, so files already on disk also recover in order.test_same_second_shutdown_slot_precedes_overflowfreezes the clock and uses adversarial filenames:B, A, Cbefore the fix,A, B, Cafter.
@dokterdok, fixed in 7a54b7f with your first option. _start_finish_wiring now runs _recover_pending_flushes just before _schedule_resume_pending_sessions(), so the replay comes before the resume turns and the startup inbound drain. I removed the call after runner.start() in main() (gateway/run.py is 7 lines shorter). A recovery failure now logs at WARNING instead of being swallowed. test_boot_recovery_runs_before_resume_turns_and_queued_inbound pins the order. It's a call-order test with stubs, and I didn't separately reproduce the double replay you found by reading. Neither a resume turn nor a queued inbound turn can now write before the boot pass has finished. If startup aborts before that point, the files stay on disk for the next start.
@ehz0ah's two-failure case is answered in the inline thread (02938cb, 376d2e1).
@teknium1, the hold-back for ordinary shutdown pending-text payloads is still a follow-up, as you scoped it.
Verification: each new test fails without its fix. The full tests/gateway + tests/hermes_state run at 7a54b7f gave 10,797 passed, 0 failed, 76 skipped. After 376d2e1, the 73 related files gave 978 passed. Ruff, code health (0 blocking) and git diff --check are clean.
kshitijk4poor
left a comment
There was a problem hiding this comment.
Reviewed at 376d2e134e (re-run after the 5 commits that landed during review).
The ordering and field-fidelity fix is correct. Real SessionDB, real spool_dropped_transcript_message, temp HERMES_HOME, 12 spooled messages including tool calls:
| order | tool_calls / tool_call_id / reasoning | |
|---|---|---|
origin/main |
shuffled | all NULL |
| this head | drop order | intact |
The new commits fix four things I had queued against 36f27f287f0, and each held up when I re-ran it:
- Memory: 8 agent-history snapshots, 43.6 MB on disk. The
recover_pending_to_dbpeak went from 86.5 MB at36f27f2to 25.7 MB here, the same asmain. - Slot-head/overflow tie:
seqnow defaults to-1. - Second-drain-failure inversion: the live row stays queued.
- Boot order: recovery now runs before resume turns and queued inbound.
9 spool/flush/transcript/session test files, 109 passed. The one red CI job, test_pty_keepalive_ws.py::test_legacy_channel_connect_pops_marker_on_disconnect, is in the PTY dashboard code and unrelated to this PR.
Blocking: one malformed spooled row now stalls every write for its session, permanently
transcript_append_kwargs now passes tool_call_id, tool_name and the reasoning fields to sqlite as they are; non-strings get through _scrub_surrogates unchanged. A spooled message with a dict or list in one of those fields therefore raises Error binding parameter N: type 'dict' is not supported. main replayed only role/content/timestamp, so the same file recovered (lossily) and the session carried on.
On this head, the boot pass holds the session back. Then, because of 02938cb9d5, every live write for that session waits behind the poisoned file. Probe: real SessionDB plus a real SessionStore, one spooled {"role": "tool", "tool_call_id": {"bad": 1}}, then 300 live append_to_transcript calls:
| DB rows for the session | spool files | |
|---|---|---|
main (recovery only) |
poisoned row recovered, 0 left | 0 |
| this head | 0 | 301 |
Every new message goes to disk. The escalation path logs one ERROR per write ("session is stalled and needs operator attention"), and a restart can't clear it because the boot pass hits the same file first. With the new "wait while the spool fails" rule, a deterministic rejection now costs the whole session. Before, it lost one row's fields.
Fix: tell deterministic rejections apart from transient ones. Coerce non-scalar fields to str/None in transcript_append_kwargs so these rows can't fail to bind. For any other row-specific rejection (for example sqlite3.InterfaceError/IntegrityError, as opposed to OperationalError: database is locked), quarantine that file (rename it to *.bad and log once) instead of holding the session back. Please add a test that pins "a poisoned spool row does not stop later live rows from landing".
Should fix in the same pass
- Two spool orderers that disagree.
drain_transcript_spoolstill sorts on rawpayload.get("ts", 0), payload.get("seq", 0)(shutdown_flush.py:166). A stringtsraisesTypeErrorthere;NaNscrambles the sort. Those are the cases_sort_numberhandles for recovery, and now also theseqdefault, which is-1in recovery but0in the live drain. The_order_flush_filesdocstring says the two mirror each other. One shared_spool_sort_key(payload, name)used by both makes that true. Since the live drain now gates every live write, aTypeErrorthere matters more than it did. SessionStore.spooled_drop_sessions()hands out the internal mutable set, andrecover_pending_to_dbmutates it from another module without the store's lock. A narrowmark_spooled_drop_sessions(ids)on the store, withrecover_pending_to_dbreturning the held-back ids instead of filling a set passed in, keeps the store owning its state.run.py::_recover_pending_flushesstill has the 119-character tuple assignment that keeps the line count flat under theFILE_LINEScap. Its only caller is nowrun_startup.py, so it can move intogateway/shutdown_flush.py(orrun_startup.py) and actually bringrun.pydown.- The late
transcript_append_kwargsimport sits inside thetrythat marks the session blocked, so anImportErrorwould be counted as a failed DB replay. Hoist it above thetry. blocked_sessionsisOptional, but its only caller always passes it. Make it required and drop the twois not Noneguards.
Nits
- Make
"shutdown-with-unpersisted-agent-history"a constant next toTRANSCRIPT_CAP_DROP_REASON; it is a literal twice. - The "utf-8-sig: our BOM-tolerant read fix" comment and the "Recovery used to walk…" history paragraph in the
_order_flush_filesdocstring are commit-message material. - Three tests now build
SessionStorewithobject.__new__plus 5 hand-set private attributes. One shared fixture in the test file would stop them breaking on unrelated attribute additions. test_a_replayed_row_is_the_row_the_live_writer_writeswould be stronger if it compared recovery's kwargs to what_append_transcript_messageactually sends, rather than re-listing the fields.
Checked and fine
- Multiplexed profiles: live profile turns run inside
_async_profile_runtime_scope, andAsyncSessionStoreusesasyncio.to_thread, which copies the context. So the live drain resolves the samepending_messages/dir the per-profile boot pass read. content=Noneinstead of""for cap-drop rows: the column is nullable, the FTS trigger handles NULL, and the live writer already writesNonefor tool-call rows.- The FTS-repair
continuein_append_to_transcript_serializedre-drains before the live row.test_repairable_spool_replay_failure_still_reaches_fts_rebuildcovers it.
transcript_append_kwargs passed tool_call_id, tool_name, role, reasoning and platform_message_id to sqlite as they were, so a spooled dict or list in one of them failed to bind on every replay. With live writes now waiting behind a failing spool, that one row stopped the whole session persisting, across restarts too. Those fields are now made bindable, and a row the database still rejects on its own values is moved to *.json.bad and logged once instead of holding the session back, both at boot and in the live drain. Also from the same review: the live drain and the restart pass share one sort key (_spool_sort_key), so a string or NaN ts no longer raises in the live drain; boot recovery returns the sessions it held back and the store marks them under its lock (mark_spooled_drop_sessions) instead of handing out its set; _recover_pending_flushes moves out of gateway/run.py as shutdown_flush.recover_gateway_pending; the transcript_append_kwargs import sits outside the try; blocked_sessions is required; the agent-history reason is a constant; the recovery tests share one store fixture and compare replayed rows with the live writer's. Review findings by @kshitijk4poor.
is_row_rejection treated every sqlite3.InterfaceError as a bad row, but SessionDB retries InterfaceError('no more rows available') as WAL contention and re-raises it once its patience runs out, so a valid row could be moved out of replay for good. Among sqlite errors only a 'binding parameter' message now counts; OverflowError does too. _bindable_text also passed ints outside sqlite's 64-bit range, which raise OverflowError even for a TEXT column and stalled the session the same way a dict did; they are now stored as text.
unsupportedpastels
left a comment
There was a problem hiding this comment.
Thanks @kshitijk4poor, the poisoned-row probe was spot on. Everything you raised is addressed in ab8d36830c3 and 575b4dd70d2.
Blocking: one malformed spooled row stalled its session
I reproduced your probe first on 376d2e134eb: after a restart and 5 live writes, the session had 0 rows and 6 spool files. Fixed in two layers:
- Bindable fields.
transcript_append_kwargsnow passesrole,tool_name,tool_call_id,reasoning,reasoning_contentandplatform_message_idthrough_bindable_text. A dict or list becomes its JSON text, and so does an int outside sqlite's signed 64-bit range (2**100raisesOverflowErroreven into a TEXT column). - Quarantine for deterministic rejections.
is_row_rejection()covers sqliteInterfaceError/ProgrammingErrorwhose message is about binding a parameter, plusTypeError/ValueError/AttributeError/OverflowErrorraised on the row's own fields. That file is renamed to*.json.bad, logged once, and the session's later rows carry on. This happens both in the boot pass and in the live drain.
Two places where I went narrower than your suggestion, deliberately:
- Not every
InterfaceError.SessionDB._execute_writeretriesInterfaceError("no more rows available")as WAL contention and re-raises it unchanged when its patience runs out (hermes_state.py:1064). Treating it as a rejection would quarantine a valid row for good, so only binding messages count. - Not
IntegrityError. Constraint failures here can be session-wide (for example a missing session row), and quarantining them would pull every good row out of automatic replay. They stay retryable, as before.
Tests: test_a_poisoned_spool_row_does_not_stop_later_live_rows (real SessionDB and SessionStore; dict and 2**100 tool_call_id rows, then a live write: result, older, big id, live, nothing left on disk or queued), test_a_row_the_database_rejects_is_quarantined_not_retried (boot and live drain), test_transient_no_more_rows_error_keeps_the_row_for_retry, and test_only_errors_about_the_rows_own_values_are_rejections.
Should-fix items
- One orderer.
_spool_sort_key(payload, name)is shared bydrain_transcript_spooland_order_flush_files. A string or NaNtsno longer raises in the live drain, andseqdefaults to-1in both.test_live_drain_and_restart_pass_order_malformed_files_alikeconfirms that both passes produce the same order. - Store owns its set.
recover_pending_spoolreturns(recovered, held_back), andSessionStore.mark_spooled_drop_sessions(ids)records them under_transcript_retry_lock. The other two touches of that set now take the lock too.spooled_drop_sessions()is gone. - Out of
run.py._recover_pending_flushesis nowshutdown_flush.recover_gateway_pending, andrun.pyis 29 lines shorter. - Import hoisted above the
trythat marks the session blocked. blocked_sessionsis required, and theis not Noneguards are removed.
Nits
AGENT_HISTORY_REASONis a constant next toTRANSCRIPT_CAP_DROP_REASON.- The history paragraph and the BOM comment are out of the docstrings.
- One
make_storefixture builds a realSessionStorethrough its constructor. test_a_replayed_row_is_the_row_the_live_writer_writesnow compares recovery'sappend_messagekwargs with what_append_transcript_messageactually sends. Only the timestamp differs: a message with none takes the payload clock.
One side effect: moving the per-file catch-all into the new recover_pending_spool made code health flag it as new BLE001. It carries a # health: allow waiver because the behaviour is unchanged. That handler is what lets one bad file skip only itself.
Verification
tests/gateway+tests/hermes_state: 10,803 passed, 0 failed.- Ruff, code health (0 blocking) and
git diff --checkare clean. - The red
test_pty_keepalive_ws.pyjob is the one you already identified as unrelated. I haven't re-checked it on this head.
|
| except Exception as exc: | ||
| logger.warning("Failed to drain transcript spool for %s: %s", session_id, exc) | ||
| return None |
There was a problem hiding this comment.
If a spool row is saved but its file cannot be deleted, _drain_spooled_drops returns None after logging the error. The live writer then saves a newer row while the old file remains. On the next restart, the old row can be saved again, leaving a duplicate out of order. Keep the live row queued when file cleanup fails.
Prompt To Fix With AI
This is a comment left during a code review.
Path: gateway/session_transcript.py
Line: 431-433
Comment:
**Old spool row repeats**
If a spool row is saved but its file cannot be deleted, `_drain_spooled_drops` returns `None` after logging the error. The live writer then saves a newer row while the old file remains. On the next restart, the old row can be saved again, leaving a duplicate out of order. Keep the live row queued when file cleanup fails.
---
For each issue above, determine whether it is valid and should be fixed. If so, fix it directly.There was a problem hiding this comment.
Partly fixed in 0b9659c7887. Within a running process, a spool file whose row was written but could not be deleted is no longer replayed again, and replay carries on with the files after it, so the live row no longer overtakes them. Covered by test_a_replayed_file_that_cannot_be_deleted_is_not_replayed_again and test_boot_recovery_counts_and_never_repeats_a_file_it_could_not_delete.
The restart duplicate you describe remains: nothing on disk records that the row was written, so the next start writes it once more. We are keeping fail-open behaviour rather than holding newer messages back, and the error log names the file to delete by hand. Keeping the live row queued would not prevent the restart duplicate either.
…eleted When a spooled row was written but its file could not be deleted, the OSError aborted the drain: the files after it were left for later, so the live row landed ahead of them, and because the file stayed, the next live write replayed and wrote the same row again, once per write. The boot pass had the same repeat. _remove_replayed now logs the failure, remembers the path so neither replay pass in this process touches it again, and lets replay go on with the later files. A restart still writes that row once more, since nothing on disk marks it; the log says to delete the file. Live rows are not held back for this: an unwritable spool directory stays fail-open. Review finding by Greptile.
unsupportedpastels
left a comment
There was a problem hiding this comment.
Merged. Thanks to everyone who reviewed this:
- @kshitijk4poor for the poisoned-row probe that found the session stall.
- @ehz0ah for the second-drain-failure inversion and the hidden replay error.
- @BearHuddleston for spotting that live writes could overtake boot held-back files.
- @keeltrace for the memory and shutdown-order findings.
- @dokterdok for the boot-order race.
- @teknium1 for the priority call and the shutdown-message follow-up.
- Greptile for the spool cleanup edge cases.
The commits from @briandevans (#84785) and @jonpol01 (#122718) keep their original authorship on main.
Proposed priority: P1
The issue is labeled P0. teknium assessed it as P1 on #123658: it needs a DB outage plus a gateway restart before the live drain recovers the spool. When it does trigger, the damage is permanent. The recovered transcript is reordered, tool calls can be separated from their results, and tool-call/reasoning fields are lost.
What does this PR do?
After a restart, recovered transcript spool files are now replayed in drop order and with full fidelity, the same way the live
drain_transcript_spoolpath replays them. A failed append holds back that session's later files instead of letting them land ahead of it.Root cause
recover_pending_to_db/_recover_one_payloadingateway/shutdown_flush.pywere written separately from the live drain and never picked up its behaviour:sorted(glob("*.json")), but spool files are namedpending-{uuid4}.json, so that order is random.SessionDBrestores by AUTOINCREMENT id, so the shuffle is permanent. The payloads already carryts/seq, which the live drain sorts on.role,content or ""andtimestampwere appended before the file was deleted.tool_calls,tool_call_id,tool_name, reasoning, platform ids, theapi_contentsidecar and the display fields were lost.Related Issue
Fixes #123658
Salvages #122718, #84785
Type of Change
Changes Made
The 8 commits from #122718 are cherry-picked onto current
mainwithout changes, keeping their authorship. One follow-up commit is added.gateway/shutdown_flush.py:_order_flush_filesreplays in(ts, seq, filename)order, and unparseable files sort last._sort_numbermakes a corrupt, overflowing or NaNts/seqsort first instead of aborting the pass. A failed append holds back that session's later files, and other sessions still recover. Held-back files are reported in one log line. A comment notes thatseqis per process.gateway/session_transcript.py: one sharedtranscript_append_kwargsrow builder is used by both the live writer and the restart writer, so their field lists can't drift apart.tests/gateway/test_shutdown_flush_recovery.py: 4 tests covering the invariants teknium asked for on fix(gateway): replay the transcript spool in drop order with full fidelity on restart (salvage #84785) #122718: drop order, row parity with the live writer, and holding back only the failed session's files._order_flush_filescaught a blindexcept Exception, which the code-health gate blocks (BLE001). It now catches(OSError, ValueError, RecursionError). Unreadable and malformed files are still sorted last with no payload, as before.Relationship to existing work
This is #122718's fix unchanged, rebased onto current
main, plus the one code-health follow-up it needed to pass CI. It is offered as a vehicle for that work, not a replacement for the authors' design. #122718 itself salvaged @briandevans's #84785, which built on the diagnosis in @686f6c61's #78323.Credit to @briandevans (original implementation, #84785) and @jonpol01 (salvage, shared row builder, overflow fix, #122718).
Related but separate, and not addressed here: #129835 (versioned transcript codec), #117310 (routing recovery to the owning store), #133366 (sessions that no longer exist).
How to Test
main, 3 of the 4 tests intests/gateway/test_shutdown_flush_recovery.pyfail: drop order, row parity, and holding back failed sessions.scripts/run_tests.sh tests/gateway/test_shutdown_flush_recovery.py tests/gateway/test_shutdown_flush.pygives 15 passed.Checklist
Code
scripts/run_tests.sh(full suite left to CI)python scripts/check --only healthpasses (0 blocking), as doruff checkandgit diff --checkDocumentation & Housekeeping
Co-authored-by: briandevans 252620095+briandevans@users.noreply.github.com
Co-authored-by: John Paul Soliva soliva.johnpaul@icloud.com
Infographic