feat(pull): route interactive producers onto the durable queue on pull pilots (#3114) - #3124
Conversation
…l pilots (#3114) On a pull-pilot agent every interactive execute_task caller now enqueues and waits for the row's terminal: Session tab, Workspace portal, public link, Slack/Telegram/WhatsApp, rooms, paid, MCP inline auth, validation and run-now. The UI /chat path is the one remaining push path. - pull_owns_dispatch covers PULL_REACHABLE_NON_AUTONOMOUS; the alert set _AUTONOMOUS_TRIGGERS is unchanged. - dispatch_and_await_terminal waits for a claim (at most one agent timeout; unclaimed turns are cancelled and reported as capacity; a caller that goes away cancels its queued turn), then for the terminal. - Server-built conversation_key per entry point; persist_session, images, schedule context and attempt ride the queue; the claim-time prompt carries the same provenance as push. - The pull sink stores compact events and wakes the sync waiter; the row rebuild returns a push-shaped result. - Stream proxies hold while the turn is queued. - ResumeLock and portal in-flight markers are kept, with TTLs covering the claim waits. Pilot agents must run the rebuilt base image (pull_worker reads persist_session and images from the payload). Fixes #3114 Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
- Delegated children: the pull sink sends the completion report the push terminals send, so a pulled child's result reaches the portal or room thread that delegated it. - Rooms: the claim wait is capped at 300s (dispatch_and_await_terminal claim_budget) and the working marker's TTL covers it. Wakes of one agent in one room are serialized with a per-(room, agent) lock, so a second wake reads the cursor and session the first one left. - A turn nobody claims in time is stored FAILED with the capacity error; a caller that went away still leaves it CANCELLED. - Stop between claim and spawn: the worker registers the claimed turn as pending, skips it if it was terminated first, and reports it cancelled. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
|
Pushed b043084: Workspace paths on pull pilots, from an audit of every Workspace action that runs or controls an agent turn.
Tests: Separate private PR: abilityai/trinity-enterprise#740. The skill runner reads the reply of a pulled |
… pilots (#3114) From the live check of Workspace paths on a pull pilot: - The stream hold covers the moment a pre-created portal row is still running before it is enqueued (db.execution_awaits_claim), so a stream opened right after the 202 waits for the claim. - The Session tab answers 429 when the turn fails for capacity, including a pilot turn no worker claimed in time. - A turn stopped while a worker runs it reports "Execution terminated by user". Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
…e-interactive # Conflicts: # src/backend/services/task_execution_service.py
|
merge-train: not on today's train — validated at lane C (
Labelled |
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
…d; clock restarts at claim (#3114) A pull pilot's Workspace turn is pre-created `running` and enqueued a moment later, and the live card's pre-feed placeholder said `running`, so the card flashed "Working" before "Waiting for a slot". The Work projection now reports an un-leased running row on a pull-owned trigger as `queued`, and the placeholder reads `queued` until the agent streams. The card's clock restarts when the turn moves from queued to running, so "Working" counts working time. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
…ilot row test (#3114) Which triggers a pilot pulls changes across the pull-migration PRs; the pilot/push and leased/un-leased cases cover the rule. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
|
Pushed b463296 + 14faaf8 from manual testing on a preview stack: a pilot's Workspace turn showed "Working" for a moment before "Waiting for a slot", and the live card's clock kept counting from send after a worker claimed the turn. The Work projection now reports an un-leased running row on a pull-owned trigger as |
|
merge-train: not on this train. It rides the next one once fixed. I've set Needs your call:
Either gate the lock on pull pilots, or keep it and say so in the body. Mechanical, and also red on Tier 1: Smaller:
#3145 is stacked on this PR and waits with it. |
…ng (#3114) - Room wake lock is taken only when a pull pilot owns the room trigger; push agents wake as before. - Claim-wait CAS winners (unclaimed FAILED, caller-away CANCELLED) close their dispatch activity and emit the terminal event. - /task backlog path reads the pull sink's `[code]` prefix, so a pilot's failed turn answers the producer's status instead of a blanket 503. - Portal roster carries `pulls_turns`; the pending-turn card reads "Waiting for a slot" only on pull agents. - Tests: `lease_expires_at` in the hand-built schedule_executions schemas (test_1474, test_2467); real execute_task -> CapacityManager.acquire -> enqueue for room/public/slack asserting the row's conversation_key; terminal hooks on both claim-wait winners; backlog error code; push-agent wake takes no lock; push agent pending card reads "Working". Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
|
Pushed 3ba55d2 for the merge-train review:
Earlier item 2 now affects pilots only, where the Mutating each fix turns its test red. |
|
merge-train: validated READY on Follow-ups worth an issue, none blocking — all are pre-existing classes that this PR widens rather than regressions:
Ops note from the body stands: pilot agents need the rebuilt base image before this reaches them. |
Summary
On a pull-pilot agent every interactive
execute_taskcaller now puts its turn on the durable queue and waits for the row's terminal: Session tab, Workspace portal, public link, Slack/Telegram/WhatsApp, rooms, paid, MCP inline auth, validation and run-now. Fire-and-forget callers (voice/VoIP post-processing, portal capture-feedback, async run-now) enqueue and return. The UI/chatpath is the one remaining push path; it is the next PR. Non-pilot agents are unchanged.pull_owns_dispatchcoversPULL_REACHABLE_NON_AUTONOMOUS(interactive triggers exceptchat, plusvalidation)._AUTONOMOUS_TRIGGERS(alert semantics) is unchanged.dispatch_and_await_terminalwaits for a claim for at most one agent timeout. An unclaimed turn is cancelled and returned as capacity (429 at every caller), a caller that goes away cancels its queued turn, then the existing terminal wait runs from the claim. A turn that finished inside one claim poll returns at once.conversation_key(session:,public:,channel:,room:,paid:),persist_session, images, schedule context, attempt. The claim-time prompt carries the same provenance as push.[code]prefix and removed from the text, cancelled turns keep their partial reply, dispatch sentinels never pass as a session id).ResumeLockand portal in-flight markers are kept (they also keep the cached resume id fresh), with TTLs extended by two claim waits (turn + cold retry). The claim guard (feat(pull): one turn per conversation at a time — the transcript-safety half of Open Question 7 #2843) serialises rooms, public and channels.claim_next_queuedlogs when it loses the conversation race three times.Deploy note
Pilot agents must run the rebuilt base image before this is enabled on them:
pull_workerreadspersist_sessionand images from the payload. An old worker drops images and does not persist a cold Session-tab turn, so turn 2 starts cold.Live verification (isolated sibling stack, real Claude turns, 3-worker pilot vs non-pilot)
/taskmanual and/taskMCP on the pilot each produced a row claimed by a worker (lease + claim token set), status success, and the HTTP caller got the reply.conversation_keyNULL).Known gaps (follow-ups)
max_backlog_depth, which scheduled work also uses (see the/csoreport appendix).backlog_metadatauntil the 90-day prune, as they keep the message text.Test Plan
cd tests && pytest unit/test_3114_pull_route_interactive.py unit/test_3114_pull_route_callers.py -v(53 tests; key renames, cancel, lost wake-up and G-04 mutations go red)test_1766,test_2048,test_2391; pull, backlog, canary andtest_679_callerssuites: 514 passeddev(IPv6/SSRF tests under local Python 3.11, and a missing localopentelemetry-instrumentation-httpxpackage)/cso --diff: 0 findingsFixes #3114
🤖 Generated with Claude Code