Repository navigation
feat(pull): /chat on the durable queue for pull pilots (#3127) - #3145
Conversation
4bf1eab to
25d2096
Compare
7b3238a to
fb70e5d
Compare
|
✅ Alembic head check clear — merging this PR into Previously flagged; resolved. Advisory — this check does not block merge. · head_sha: |
|
✅ Nightly unit-suite clean when this PR is merged into |
|
Resolve by merging |
|
merge-train (2026-10-04): not on this train. It rides the next one once fixed. This PR was not validated today because the diff will change with the rework below. The branch is 45 commits behind
I've set |
On a pull-pilot agent, POST /api/agents/{name}/chat (MCP chat_with_agent,
the connector tool and the trinity CLI) is admitted onto the durable
queue and claimed by a worker. Every interactive trigger is now
pull-owned on a pilot.
Memory on a pilot is one Claude conversation per chat_sessions row,
i.e. per (agent, user), resumed by id through run_resumable_turn
(ResumeLock, persist_session, cold retry on resume-not-found). The id is
cached on chat_sessions.cached_claude_session_id (SQLite migration +
Alembic 0084) and joins the session reaper's keep set.
The post-turn work /chat did inline runs in the caller after the
terminal: assistant message, collaboration close, response shape,
idempotency complete, timeout receipt. GET /chat/history serves the
caller's session from the database on a pilot; DELETE /chat/history
also clears the cached ids. execute_task carries chain_depth so a
cold-retry row keeps its depth.
Non-pilot agents are unchanged.
Fixes #3127
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
DELETE /chat/history forgot the cached Claude ids but left the sessions active, so GET /chat/history still showed the old conversation. The sessions that carry a cached id are now closed with it; sessions /chat never used stay open. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
… 0085 (#3127) Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
… 0088 (#3127) Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
fb70e5d to
f6c0924
Compare
|
merge-train: not on today's train. It rides the next one once fixed. Merged with current The new
Adding it to Reproduce: merge |
- capacity_manager: take dev's earlier pull_exclusive placement (#2514); keep the #3127 comment that /chat is pull-owned on pilots. - Rechain the chat_sessions Alembic revision as 0090_chat_session_claude_id off dev's head 0089_supersede_queue_flood_backlog; SQLite entry ordered after supersede_queue_flood_backlog. - run_pulled_chat_turn passes request_text=request.message to the resumable turn. /chat runs no admission-seam skill gate, so the executor backstop must gate on the caller's own words (ent#751). Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
…3247) CP1 of the replace-a-pending-ask slice. `operator_queue` gains two nullable TEXT link columns: `replaces` on the successor (the predecessor row's uuid) and `replaced_by` on the predecessor (the successor row's uuid). Both are stamped in the one per-agent locked transaction the next checkpoint adds, so each row is self-describing on every surface that holds only one of the pair. Both schema tracks (Invariant #9): SQLite entry `operator_queue_replace` and Alembic `0090_operator_queue_replace` chained after this branch's head `0089_supersede_queue_flood_backlog` — `0090` is also taken by the independent #3246 branch and open PR #3145, so expect a renumber at merge; `check_alembic_heads.py` names the fork. `schema.py` and `tables.py` DDL updated (the `disposed_by` comment gains `agent`); `_row_to_item` / `_SELECT_COLS` carry the columns; `_MACHINE_ROW_FIELDS` and the ent#715 `MACHINE_ROW_KEYS` pin are extended in the same commit. No index. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
AndriiPasternak31
left a comment
There was a problem hiding this comment.
Thanks for this. The pull routing is clean: I couldn't find a double-execution path (run_chat_turn branches only on capacity is None, and execute_task re-decides with the same env-only predicate). The only caller-side terminal write is on a never-enqueued RUNNING row, so it can't race the token-gated sink, and the slot and idempotency accounting balance. Tier 1 is green and check_alembic_heads passes. I ran test_3127 + the flipped pins (164 passed) and the 1804/1578/2806/2842/2889/ent751/3114 guards (297 passed).
Two things I'd like fixed before merge. Both are invisible to the suite because dispatch_and_await_terminal is mocked in every turn test:
1. Skill gate runs twice on a pilot /chat (ent#751). The comment at chat_execution_service.py:1011 says "/chat runs no admission-seam gate", but admit_chat_request calls skill_gate_service.enforce at dispatch_admission_service.py:326, before the pilot branch at :387. Without gate_checked=True, execute_task's backstop gates again with a requester rebuilt from the row (requester_for_dispatch), and that requester never has is_person=True. I checked this with the real enforce: the admission requester self-approves the owner, and the backstop requester does not. So on a pilot, an approver's own gated request turns into a pending approval: the row is SKIPPED, a second ask is raised, the caller gets a 202, and the user message has no reply. The pilot branch also returns before audit_self_approved (:431). Fix: pass gate_checked=True (keep request_text), call audit_self_approved in the pilot branch, and assert gate_checked in test_pilot_turn_dispatches_through_the_queue. This is the first branch of vybe's 10-05 note.
2. Agent-to-agent /chat (trigger agent) gets no claim priority and no claim budget. "agent" is not in INTERACTIVE_TRIGGERS, so _CLAIM_WAITING_TRIGGERS (task_execution_service.py:720/811) skips the phase-1 wait. The turn queues behind batch work, and the only bound is timeout + 120s from enqueue. With a deep queue that ends in the 504 queued_timeout receipt while the row is still queued. It then runs for nobody: no assistant chat_messages row and no cached Claude id. The collaboration activity also stays started until the 120-min backstop, because build_pull_queue_payload sends collaboration_activity_id=None. That's the #1804 symptom. The a2a note in pull_pilot.py states the rule: a blocked caller earns priority and the budget. Could the pulled /chat opt into claim semantics explicitly (rather than by trigger), and could the collaboration activity id ride the payload so the sink closes it? A test with the real adapter for trigger agent would pin it.
Smaller items, non-blocking:
- Post-turn work runs in the caller, so it's lost on a backend restart or a 504. Issue scope item 3 said "moves to the terminal", and #3227 puts post-turn delivery in the sink. No textual conflict between the two PRs, but worth converging once #3227 lands (carry
chat_session_idandcollaboration_activity_idin the metadata). - Same-user concurrency: the ResumeLock waits 30s and then the call gets a 429 and the row is FAILED. Agent keys resolve to the owner, so every agent of one owner shares one session with the target, and parallel
chat_with_agentcalls fail where push serialized them. Consider a pilot lock wait of at least one turn, as rooms do. cached_uuidis read before the lock (:1005), so two concurrent first turns both run cold and the second overwrites the first's id.- Lock-busy path: the
update_execution_statusbool is ignored, noagent.task.failedis emitted (MCP's receipt goes out at 25s and this FAILED lands at 30s), and the 429 has noX-Trinity-Error-Code: capacity. - Backlog full on a pilot
/chatmaps to 503 with no code, because "Agent backlog full…" doesn't match "at capacity". Push answered 429/capacity. - The Chat tab writes into the same
chat_sessionsrow, so "New chat" in the UI resets MCP/CLI memory, and the owner reset closes Chat-tab sessions. - Docs: the
chat_with_agentdescription inmcp-server/src/tools/chat.ts:395-399still says the session is shared by every caller and restarts on a model change.persistent-chat-tracking.md,mcp-orchestration.mdandagent-to-agent-collaboration.mdaren't updated. The PR body and the CSO report still say Alembic0086and "stacked on #3124". - Alembic: #3255 and #3256 also add
0090_*off0089_supersede_queue_flood_backlog. The tables are disjoint, so this is a mechanical re-parent for whoever lands second.
- Skill gate runs once: the pulled turn passes gate_checked=True, and the pilot admission branch audits a self-approved gate like the push branch. - Agent-to-agent /chat (trigger agent) on a pilot: run_resumable_turn opts into the claim phase (caller_waiting=True), the claim orders session:-keyed rows with interactive turns, and the collaboration activity rides the queue payload so the pull sink closes it. - Resume lock waits one turn's lock TTL; lock-busy closes the activity and emits the terminal event only on a CAS win, and answers 429 capacity. - Full backlog returns FAILED/CAPACITY and maps to 429 capacity. - Docs: chat_with_agent description, execution.md, three feature flows. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
- Pull sink: take dev's post-turn delivery (#2329), which closes the collaboration activity from the queued row; drop the separate close. - Pilot /chat admission hands the gate decision on (ent#752), so the row's setup records the self-approval; the admission-side audit is dropped. - A self-approved pulled turn (isolated_session) starts cold and its Claude id is not cached. - chat_with_agent description: dev's shortened text (2,048-char cap). Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
…3247) CP1 of the replace-a-pending-ask slice. `operator_queue` gains two nullable TEXT link columns: `replaces` on the successor (the predecessor row's uuid) and `replaced_by` on the predecessor (the successor row's uuid). Both are stamped in the one per-agent locked transaction the next checkpoint adds, so each row is self-describing on every surface that holds only one of the pair. Both schema tracks (Invariant #9): SQLite entry `operator_queue_replace` and Alembic `0090_operator_queue_replace` chained after this branch's head `0089_supersede_queue_flood_backlog` — `0090` is also taken by the independent #3246 branch and open PR #3145, so expect a renumber at merge; `check_alembic_heads.py` names the fork. `schema.py` and `tables.py` DDL updated (the `disposed_by` comment gains `agent`); `_row_to_item` / `_SELECT_COLS` carry the columns; `_MACHINE_ROW_FIELDS` and the ent#715 `MACHINE_ROW_KEYS` pin are extended in the same commit. No index. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
…3247) CP1 of the replace-a-pending-ask slice. `operator_queue` gains two nullable TEXT link columns: `replaces` on the successor (the predecessor row's uuid) and `replaced_by` on the predecessor (the successor row's uuid). Both are stamped in the one per-agent locked transaction the next checkpoint adds, so each row is self-describing on every surface that holds only one of the pair. Both schema tracks (Invariant #9): SQLite entry `operator_queue_replace` and Alembic `0090_operator_queue_replace` chained after this branch's head `0089_supersede_queue_flood_backlog` — `0090` is also taken by the independent #3246 branch and open PR #3145, so expect a renumber at merge; `check_alembic_heads.py` names the fork. `schema.py` and `tables.py` DDL updated (the `disposed_by` comment gains `agent`); `_row_to_item` / `_SELECT_COLS` carry the columns; `_MACHINE_ROW_FIELDS` and the ent#715 `MACHINE_ROW_KEYS` pin are extended in the same commit. No index. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
…3127) — mechanical, per the merge-train note on the PR #3255, #3145 and #3269 each added an 0090 revision off 0089_supersede_queue_flood_backlog, which is the #2068 two-heads fork once two of them land. The tables are disjoint (operator_queue vs chat_sessions), so this is a re-parent: 0090_chat_session_claude_id becomes 0091_chat_session_claude_id with down_revision 0090_platform_alert_subjects. The SQLite entry is keyed by name and needs nothing; only the two doc references to the id change. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
…#3265) — mechanical, per the merge-train note on the PR #3255, #3145 and this PR each added an 0090 revision off 0089_supersede_queue_flood_backlog (the #2068 two-heads fork). The tables are disjoint (operator_queue, chat_sessions, enterprise_portal_messages), so this is a re-parent: 0090_portal_messages_attachments becomes 0092_portal_messages_attachments with down_revision 0091_chat_session_claude_id. The SQLite entry is keyed by name and needs nothing. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
|
merge-train (2026-10-06): riding this train. Mechanical changes pushed to your branch:
Expected until #3255 lands: The validation found both blocking review items resolved by |
Merge train 2026-10-06: both blocking items verified resolved by 7902fe8 (mutation-checked: gate_checked, caller_waiting, collaboration_activity_id each red when removed). Dismissed per the operator's Gate 1 call; open items are posted on the PR.
AndriiPasternak31
left a comment
There was a problem hiding this comment.
Thanks for the quick turnaround. I re-checked against b8e24ee. The current head 2082b2f only renumbers the migration (0090 → 0091), so all of this still applies. @vybe, this is on today's train. Could you hold it until the item below is fixed?
Both blockers are fixed, and the tests have teeth:
- Skill gate runs once: gate_checked=True on the pulled turn, and the self-approval now rides admission.gate into the row's setup (ent#752). Reverting either one fails a test.
- Agent-to-agent /chat: caller_waiting=True gives it the claim phase, the session: prefix gives it interactive claim order, and the collaboration activity id reaches the queued row, where #3227's sink delivery closes it. Only the CAS winner closes it. The caller's later close is a no-op under the #1804 lattice. Session tab and Workspace already wait for a claim (
session/publictriggers), so caller_waiting doesn't change them. - Lock-busy CAS/event/429-capacity and backlog-full → 429 are fixed too. Alembic is one head on the merged tree. The merge resolution with #3227 looks right: no double chat-message save.
One thing I'd like changed before merge:
- The longer resume-lock wait runs into the #106 no-session sweep. While a second turn of the same session waits, its admission row is
runningwith no lease and no Claude session, so the 5-minute cleanup fails it after 60s ("Silent launch failure"). When the lock frees, the enqueue CAS refuses the row and the caller gets 429 "Agent backlog full" for a turn that never ran. I reproduced the sweep half against the real db layer. Before this round the same case got an honest 429 after 30s. Either keep the default 30s wait for /chat, or keep the waiting row out of the sweeps (create or enqueue it after the lock). The queue's per-conversation claim guard already serializes execution.
Smaller, non-blocking:
- The merge took dev's shorter chat_with_agent text and dropped your pilot note. chat.ts still says the session is shared by every caller and restarts on a model change. The
parallelparameter description isn't under the 2,048-char cap, so the note can go there. - execution.md still says admission audits the self-approval on the pilot branch. Since the merge, the row's setup does that.
- The CAPACITY code on CapacityFull now applies to every execute_task caller (e.g. the Session tab at capacity on push becomes 429 instead of 502, and fan-out reports
capacity). Probably right, but worth a line in the PR body. - Three links have no test that fails when they're reverted: execute_task forwarding collaboration_activity_id into the payload, caller_waiting on the cold retry, and the backlog-full producer setting the code. A test that pushes an agent /chat row through apply_task_result and checks the activity ends completed would cover the first.
- The assistant message and the Claude id are still lost on a 504 (documented). Fine as a follow-up.
A longer wait outlived the #106 no-session sweep: the waiting turn's admission row is running with no Claude session, so the sweep failed it at 60s and the enqueue CAS then refused it, answering 429 for a turn that never ran. Drop the lock_wait plumbing; a test pins the default wait under NO_SESSION_TIMEOUT_SECONDS. chat_with_agent's parallel description notes per-user sessions on a pilot; execution.md says the self-approval audit rides the row's setup. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
|
@AndriiPasternak31 addressed in b829ce1:
Merged Not done here: tests for the three unpinned links (payload |
…nical #3255 landed with 0090_platform_alert_subjects, this branch's 0091 parent, so the Alembic line is one head again. The SQLite MIGRATIONS tail kept both entries. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
AndriiPasternak31
left a comment
There was a problem hiding this comment.
Re-checked at b829ce1, plus the merge-train merge 3783188. F1 is fixed. Approving.
F1 (lock wait vs the #106 sweep)
- /chat on a pilot is back to the default 30s resume-lock wait. The no-session sweep and the watchdog both leave rows younger than 60s alone, so a second turn on the same session now gets the honest "another turn is in progress" 429 at ~30s, before either of them can touch its row.
- I ran this against the real db layer: real ResumeLock (fakeredis), real mark_no_session_executions_failed every 2s, real enqueue CAS.
- b829ce1: 32s, 0 rows swept, lock-busy 429.
- With the fix reverted: 72s, the row is reaped as "Silent launch failure", then a capacity 429.
- Reverting the fix turns test_lock_wait_stays_under_the_no_session_sweep red. It checks the wiring plus a constant (30 < 60), which is enough for this fix.
- C1/C2 and the lock-busy CAS/event/429 path are unchanged. All 16 of my round-2 mutations give the same results as before.
Migrations
- #3255 is in, and 3783188 resolves the SQLite MIGRATIONS tail keeping both entries. On the Alembic dir at 3783188, check_alembic_heads gives 92 revisions and 1 head (0091_chat_session_claude_id). schema-parity, pg-migrations and regression diff were red only because 0090 was missing, so they should clear on this run. CI was still running when I checked.
Non-blocking, fine as a follow-up (no need to hold the train for them):
- Three links still have no test that fails when they're reverted: execute_task forwarding collaboration_activity_id into the payload, caller_waiting on the cold retry, and the backlog-full producer setting error_code=CAPACITY. There is also no test that runs a queued agent /chat row through apply_task_result and checks that the collaboration activity ends completed.
- chat.ts: the pilot note comes right after "The session restarts after a model change". On a pilot the cached Claude id survives a model change, so that sentence doesn't apply there. Something like "On a pull-pilot agent each calling user has their own session, which a model change does not restart."
- PR body: "Converge with #3227 … once it lands" is stale (#3227 is in).
Summary
The last push path on a pull-pilot agent.
POST /api/agents/{name}/chatis called by MCPchat_with_agent, the agent-to-agent connector tool, and thetrinity chatCLI. On a pilot it is now admitted onto the durable queue and claimed by a worker, so every interactive trigger is pull-owned on a pilot. The web UI chat panel already went throughPOST /task, which #3124 routes.session_turn_service.run_resumable_turnwith keysession:chat:<chat_sessions.id>, resumingchat_sessions.cached_claude_session_id. That gives the ResumeLock,persist_session, and a cold retry on resume-not-found. Off a pilot/chatkeeps the agent container's single shared session.gate_checked=True, so the executor backstop does not gate a second time./chat(triggeragent):run_resumable_turnopts into the claim phase (caller_waiting=True), the pull claim orderssession:-keyed rows with interactive turns, and the collaboration activity rides the queue payload so the pull sink closes it. A full backlog answers 429capacity./chatdid inline runs in the caller after the terminal: assistantchat_messagesrow (secrets scrubbed), activity closes, the same response shape, idempotency complete, a timeout receipt + 504, andResumeLockBusy(after the default 30s lock wait, which stays under the 60s no-session sweep) → FAILED + 429capacity.CapacityFullfromexecute_tasknow carries error codecapacityfor every caller: a Session tab turn at capacity on push answers 429 (was 502), and fan-out reportscapacity.GET /chat/historyon a pilot serves the caller's session from the database.DELETE /chat/historycloses the sessions/chatused and forgets their ids.chat_sessions.cached_claude_session_id(SQLite migration + Alembic0091, re-parented by the 2026-10-06 merge train onto fix(operator-queue): platform alerts are conditions — one pending row per subject (#3246, part 1) #3255's0090_platform_alert_subjects), which joins the session reaper's keep set.execute_task(chain_depth=)keeps the depth on a cold-retry row./chathandler is not cancelled when the client goes away. The MCP 25s abort therefore leaves the queued turn running and its receipt valid.Live verification (isolated sibling stack, real Claude turns, 3-worker pilot + non-pilot)
/chaton the pilot: 3 turns from one user were pulled on one Claude session and recalled earlier numbers ("17, 29"). A second user got "NONE" on their own key and session. History returned each caller's own messages. A reset made the next turn start cold, and history is empty after it./chatwas pushed with memory kept, as before.a04d6e8ce) and here (883656752).Known gaps
cached_uuidis read before the resume lock, so two concurrent first turns both run cold and the second id wins.chat_sessionsrow: its "New chat" resets MCP/CLI memory, and the owner reset closes Chat-tab sessions./chatturn gets the task-mode platform prompt, where push gives it the chat mode.output_tokensis not recorded on pulled/chatmessages.GET /chat/sessionstill reads the agent's shared session.Test Plan
cd tests && pytest unit/test_3127_pull_chat.py -v(21 tests; 11 key-behaviour mutations + the reset mutation go red)test_1766,test_2391,test_3114_pull_route_interactive,test_2610check_alembic_heads: one head (0090_chat_session_claude_id; now0091_chat_session_claude_idafter the merge-train re-parent)/cso --diff: 0 findingstest_3127_pull_chat.py,test_3114_pull_route_interactive.py,test_2842_2843_pull_claim_order.py(72 passed); related guard suites show no new failuresFixes #3127
🤖 Generated with Claude Code