Skip to content

feat(pull): route interactive producers onto the durable queue on pull pilots (#3114) - #3124

Merged
vybe merged 9 commits into
devfrom
feature/3114-pull-route-interactive
Oct 2, 2026
Merged

vybe merged 9 commits into
devfrom
feature/3114-pull-route-interactive

Conversation

@obasilakis

Copy link
Copy Markdown
Contributor

Summary

On a pull-pilot agent every interactive execute_task caller 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 /chat path is the one remaining push path; it is the next PR. Non-pilot agents are unchanged.

  • Predicate: pull_owns_dispatch covers PULL_REACHABLE_NON_AUTONOMOUS (interactive triggers except chat, plus validation). _AUTONOMOUS_TRIGGERS (alert semantics) is unchanged.
  • Wait: dispatch_and_await_terminal waits 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.
  • Carried through the queue: server-built conversation_key (session:, public:, channel:, room:, paid:), persist_session, images, schedule context, attempt. The claim-time prompt carries the same provenance as push.
  • Pull sink + result: compact events stored, sync waiter woken, push-shaped result rebuilt from the row (error code parsed from the [code] prefix and removed from the text, cancelled turns keep their partial reply, dispatch sentinels never pass as a session id).
  • Streams: chat, public and portal stream proxies hold the SSE connection while the turn is queued.
  • Session lock: ResumeLock and 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.
  • Review note from feat(pull): interactive turns first, one turn per conversation (#2842, #2843) #3110: claim_next_queued logs 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_worker reads persist_session and 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)

  • Routing: Session tab, public link, /task manual and /task MCP on the pilot each produced a row claimed by a worker (lease + claim token set), status success, and the HTTP caller got the reply.
  • Parity: 9 webhook turns (~45s) queued, all 3 workers busy, then one interactive turn. Pilot wait to claim 39.3s, total 42.6s; claimed 54 ms after the first worker freed, ahead of 6 queued webhook rows. Same agent shape on push: 42.0s, total 46.1s. Both waits equal the remainder of the running batch turn.
  • Transcript burst: 4 concurrent turns of one Session conversation ran one at a time (no overlap), each reply recalled every earlier number, and the JSONL holds 5 user + 5 assistant entries strictly alternating.
  • Non-pilot: a Session turn was pushed (no lease, conversation_key NULL).

Known gaps (follow-ups)

  • On a pilot the portal does not show an in-turn subscription switch (the pull sink switches asynchronously).
  • Anonymous public-link turns on a pilot count toward the agent's max_backlog_depth, which scheduled work also uses (see the /cso report appendix).
  • Failed rows keep base64 images in backlog_metadata until 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)
  • Flipped pins: test_1766, test_2048, test_2391; pull, backlog, canary and test_679_callers suites: 514 passed
  • Full unit suite: 19,556 passed; the 26 failures/errors reproduce on clean dev (IPv6/SSRF tests under local Python 3.11, and a missing local opentelemetry-instrumentation-httpx package)
  • /cso --diff: 0 findings

Fixes #3114

🤖 Generated with Claude Code

…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>
@obasilakis

Copy link
Copy Markdown
Contributor Author

Pushed b043084: Workspace paths on pull pilots, from an audit of every Workspace action that runs or controls an agent turn.

  • Delegated children report back: the pull sink now sends the completion report the push terminals send, so a pulled child's result reaches the portal or room thread that delegated it. That path was already broken for the pulled agent trigger before this PR.
  • Rooms: the claim wait is capped at 300s and the working marker covers it. Wakes of one agent in one room are serialized with a per-(room, agent) lock (re-entrant within a mention chain), so a second wake reads the state the first left.
  • Unclaimed turns: a turn nobody claims in time is stored failed with the capacity error. A caller that went away still leaves it cancelled.
  • Stop right after a claim: the worker registers the claimed turn as pending, skips it if it was terminated first, and reports it cancelled. This is part of the base-image rebuild already required.

Tests: tests/unit/test_3114_pull_route_workspace.py (10). The related suites (62 files) pass with 1,853 tests, and all 12 mutations across the five items go red.

Separate private PR: abilityai/trinity-enterprise#740. The skill runner reads the reply of a pulled agent turn straight away, so every run came back empty.

obasilakis and others added 2 commits October 1, 2026 11:48
… 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
@vybe

vybe commented Oct 1, 2026

Copy link
Copy Markdown
Contributor

merge-train: pushed a merge of dev into this branch (52dba43) to clear the conflict in task_execution_service.py. One hunk: execute_task's signature where dev added chain_depth (#2973) and this PR adds conversation_key — kept both. No other changes.

@vybe

vybe commented Oct 1, 2026

Copy link
Copy Markdown
Contributor

merge-train: not on today's train — validated at lane C (/validate-pr + /review + /cso --diff): Tier 1 green, no CRITICAL, CSO clean, and the claim-side serialisation is executed by real DB tests. Four findings, the first mechanical, the rest yours:

  1. Unexecuted link, execute_task → queue payload. Deleting conversation_key=conversation_key from the build_pull_queue_payload call (task_execution_service.py ~1547) leaves all 139 tests in the PR's suites green — every producer's key would be silently dropped and the claim guard would stop serialising rooms, public links and channels. One test that calls execute_task on a pilot and asserts the row's conversation_key closes it.
  2. Room wake lock bypass after a disconnect (shared_sessions/service.py, the _held_wake_locks contextvar). Nested wakes run inside asyncio.shield; when the poster disconnects the outer wake releases the lock but the orphaned nested task still carries the "held" marker and skips it, so a fresh wake of the same agent runs beside it. On a pilot the room: claim guard still serialises; on push agents nothing does — and the lock changes room behaviour for non-pilot agents, which the PR body says are unchanged.
  3. Failed /task turns on a pilot lose their error code. manual/mcp turns now go through _dispatch_sync_backlog, which rebuilds the result from the row without parsing the [code] prefix — failures come back 503 instead of the producer's status. The result_from_execution_row parsing fix does not cover this path.
  4. Low: a paid caller that disconnects after the claim leaves the turn running with no settlement; only the pre-claim disconnect cancels.

Labelled status-needs-fix; your next push clears it. Rides the next train once fixed. (The merge of dev I pushed earlier, 52dba43, stays — it only resolved the chain_depth/conversation_key signature conflict.)

@vybe vybe added the status-needs-fix PR has an unaddressed review/validation finding; cleared by the author's next push (#2815) label Oct 1, 2026
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
@github-actions github-actions Bot removed the status-needs-fix PR has an unaddressed review/validation finding; cleared by the author's next push (#2815) label Oct 1, 2026
obasilakis and others added 2 commits October 1, 2026 19:20
…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>
@obasilakis

Copy link
Copy Markdown
Contributor Author

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 queued (client_portal/work/service.py::awaits_claim, lease_expires_at added to get_fleet_executions), the card's pre-feed placeholder reads queued until the agent streams, and the clock restarts on queued → running. Tests: tests/unit/test_ent525_portal_work.py::test_unclaimed_pilot_row_reads_queued, src/frontend/tests/unit/portalWork.spec.js (pending turn on a pull agent). #3145 rebased onto it.

@vybe

vybe commented Oct 1, 2026

Copy link
Copy Markdown
Contributor

merge-train: not on this train. It rides the next one once fixed. I've set status-needs-fix, and your next push clears it.

Needs your call: shared_sessions/service.py _wake_agent (~L1124) takes room_wake_lock:{room}:{agent} for every agent, pilot or not. That contradicts the body's "non-pilot agents are unchanged":

  • With Redis available, concurrent wakes of one agent in one room are serialised for up to 360 s, after which the second is dropped with "still busy".
  • Two crossing chains (one holds X and wakes Y, the other holds Y and wakes X) stall for 360 s and then drop a wake. The contextvar only covers a single chain.

Either gate the lock on pull pilots, or keep it and say so in the body.

Mechanical, and also red on Tier 1: regression diff reports 2 new failures, test_1474_read_boundary_z::test_fleet_executions_normalizes_naive and test_2467_turn_integrity::...test_fleet_list_returns_turn_integrity. db/schedules/stats.py:416 now selects lease_expires_at, and those tests' hand-built tables lack that column. Add it to both test schemas.

Smaller:

  • portalWork.js pendingTurnOutcome shows "Waiting for a slot" for non-pilots too.
  • A claim-timeout FAILED row and a caller-away CANCELLED row emit no terminal event.
  • Every end-to-end test stubs execute_task to return QUEUED. One test that drives the real execute_task → capacity.acquire → enqueue path for an interactive trigger would close that gap.

#3145 is stacked on this PR and waits with it.

@vybe vybe added the status-needs-fix PR has an unaddressed review/validation finding; cleared by the author's next push (#2815) label Oct 1, 2026
…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>
@github-actions github-actions Bot removed the status-needs-fix PR has an unaddressed review/validation finding; cleared by the author's next push (#2815) label Oct 2, 2026
@obasilakis

Copy link
Copy Markdown
Contributor Author

Pushed 3ba55d2 for the merge-train review:

  • Room wake lock: taken only when a pull pilot owns the room trigger. Push agents wake as before, so "non-pilot agents are unchanged" holds. Test: a push agent's wake takes no lock.
  • regression diff: lease_expires_at added to the hand-built schemas in test_1474 and test_2467. The check passes.
  • "Waiting for a slot" on non-pilots: the portal roster carries pulls_turns, and the pending card reads queued only when it is true.
  • Terminal event on claim-wait rows: the unclaimed FAILED and caller-away CANCELLED winners close their dispatch activity and emit the terminal event.
  • Real-path test: real execute_task → CapacityManager.acquire → BacklogService.enqueue for room/public/slack, asserting the row's conversation_key. This also covers item 1 from the earlier review.
  • Earlier item 3: the /task backlog path parses the [code] prefix, so a pilot's failed turn answers the producer's status.

Earlier item 2 now affects pilots only, where the room: claim guard serialises. Item 4 (paid disconnect after claim) is unchanged.

Mutating each fix turns its test red.

@vybe

vybe commented Oct 2, 2026

Copy link
Copy Markdown
Contributor

merge-train: validated READY on 3ba55d2ef and merging today. /validate-pr, /review and /cso --diff all clear: every point in your 10-02 reply is true in the code, and each is pinned by a test that goes red under mutation (the pilot gate on the room wake lock, conversation_key through the real execute_task → CapacityManager.acquire → BacklogService.enqueue path, _split_error_code on the backlog path, the unclaimed-turn CAS winner closing its activity and emitting the terminal event).

Follow-ups worth an issue, none blocking — all are pre-existing classes that this PR widens rather than regressions:

  1. HTTP callers can't actually cancel a queued turn once the client drops (task_execution_service.py:584-592). The CancelledError branch is the only thing that cancels a still-queued row, and Starlette's request_response carries no disconnect watcher, so a plain endpoint coroutine runs to completion after the client is gone. On a pilot at capacity, an x402 client that times out at 60 s leaves the handler waiting up to claim_budget, a worker claims, the turn succeeds, and paid_chat settles for a reply nobody received. Same shape holds a Session ResumeLock / portal in-flight marker for the whole queued wait.
  2. Default claim_budget is one full agent timeout (:576-577) for every caller but rooms, which cap at 300 s. Upstream proxies cut long before that, feeding (1). A per-caller cap or an interactive default well below the agent timeout would help.
  3. LOW: stream_execution_log (routers/chat.py:991-996, 1023-1036) authorises on name but never checks execution.agent_name == name. Before this PR the agent-side 404 made that harmless; now wait_while_queued / execution_is_running read the row with no agent predicate, so a reader of pilot A holding B's execution id gets : queued ticks while B's row is queued. Prerequisite is hard (ids are token_urlsafe(16)), and every sibling handler binds — one uniform-404 check before proxy_stream() plus a second-agent test closes it.
  4. images stay in backlog_metadata on FAILED rows until the 90-day prune (acknowledged in the body); NULLing images on any terminal would close it.
  5. The two 1 s pollers (_wait_until_claimed, wait_while_queued) do one sync DB read per waiting caller per second on the event loop.

Ops note from the body stands: pilot agents need the rebuilt base image before this reaches them.

@vybe vybe left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

merge-train: validated READY on 3ba55d2 — /validate-pr, /review and /cso --diff clear; every point in the 10-02 reply verified in code and pinned by a mutation-red test. Follow-ups noted in the comment above.

@vybe
vybe merged commit 170d870 into dev Oct 2, 2026
27 checks passed
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

3 participants