feat(pull): interactive turns first, one turn per conversation (#2842, #2843) - #3110
Merged
Merged
Conversation
…#2843) The pull claim now takes interactive rows (pull_pilot.INTERACTIVE_TRIGGERS) before autonomous ones, oldest first within each group, with strict precedence and no anti-starvation rule. It also skips a row whose conversation already has a running turn; the partial unique index idx_executions_one_running_turn stops two concurrent claimers, and every status transition out of running releases it. conversation_key is stamped at enqueue on pull-pilot agents only. The backend drain still claims plain oldest-first. Canary B-02/B-08 skip a row waiting on its own conversation. Interactive producers are not routed onto the queue yet; that follows. Fixes #2842 Fixes #2843 Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
dolho
approved these changes
Sep 30, 2026
dolho
left a comment
Contributor
There was a problem hiding this comment.
Reviewed the claim path, both migration tracks, the canary changes, and the enqueue stamping. LGTM. Approving; the notes below are non-blocking.
What I checked
- Claim correctness. With
conversation_key IS NULL, theNOT EXISTScorrelation evaluates to NULL, so ungrouped rows are never skipped. On PG,SKIP LOCKEDlets two claimers lock different head rows of the same conversation. The partial unique index then makes the second commit raiseIntegrityError, and the retry'sNOT EXISTSsees the winner. Each attempt opens its ownbegin(), so nothing is retried inside an aborted transaction. On SQLite, writers are serialised, so the skip alone is enough there. - The guard can't leak. It is keyed on the row's own
status = 'running', so terminal, lease-expiry, andrelease_claim_to_queuedall free it with no extra bookkeeping. - Dual track. The SQLite migration, the Alembic
0083revision (chained off0082_agent_sync_state_divergence), theschema.pyDDL, and thetables.pyIndex(declared so autogenerate won't propose dropping it) all agree. - Scope. Keys are only stamped for pull pilots, and
drain_nextalready refuses pilots (#1766), so the backend drain never sees a keyed row. The canary B-02/B-08 skip matches the claim's skip.
Non-blocking notes
- De-piloting with keyed rows still queued. If an agent leaves the pilot while keyed rows are queued,
drain_nextclaims with noNOT EXISTS. A queued turn whose conversation is still running then hits the index 3 times, getsNone, and the head row blocks drain until that turn finishes. It clears itself, but theIntegrityErrorbranch returns silently. Alogger.infothere would make that stall visible. - Push-path running rows carry no key. Interactive turns still arrive by push, and those running rows have
conversation_key = NULL. Until the routing PR lands, the index only serialises queued-vs-queued turns of a conversation, not a pushed turn racing a queued one. The Redis session lock still covers that case today. Worth an explicit test in the routing PR. - Ordering and the partial index. The
CASEsort key meansidx_executions_queued (agent_name, queued_at)can't serve the ORDER BY, so the queue gets a sort step. That's fine at per-agent queue depths; just noting it. - Merge order. Several open PRs also add an
0083_*revision off the same0082(#3085, #3021/#3022, #2984). Whichever merges second needs itsdown_revisionrepointed, andalembic-head-watchwill flag it on the next push to dev.
6 tasks
6 tasks done
4 tasks done
vybe
pushed a commit
that referenced
this pull request
Sep 30, 2026
…il_identity (#3085) dev gained 0083_execution_conversation_key (#3110) off the same parent as this PR's 0083_ent720_email_identity — a two-head fork (#2068). The two touch disjoint tables, so the unmerged revision is re-parented: renamed to 0084_ent720_email_identity with down_revision 0083_execution_conversation_key. SQLite list keeps both entries, dev's first. Test path + down_revision pin, migrations.py docstring and the feature flow follow the rename. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
Summary
pull_pilot.INTERACTIVE_TRIGGERS) before autonomous ones, oldest first within each group. Strict precedence, no anti-starvation rule (by decision, recorded on feat(pull): interactive turns jump the queue — the starvation half of Open Question 7 #2842): steady chat that fills every worker is answered by raising the worker count. No worker is reserved (one of N idle is 33% of a 3-worker agent).conversation_keyalready has arunningrow, so a busy conversation never blocks the rest of the queue. The partial unique indexidx_executions_one_running_turnstops two concurrent claimers of one conversation; every transition out ofrunning(terminal, lease expiry, requeue) releases it, so there is no lock to leak.conversation_keyis stamped at enqueue on pull-pilot agents only. The backend drain still claims plain oldest-first.Nothing is routed onto the queue by this PR: interactive producers still push. That is the next PR, which also carries the live parity measurement and the transcript-burst check.
Changes
db/schedules/queue.py— claim ordering, conversation skip,IntegrityErrorretry;conversation_keyon the enqueue transitiondb/schema.py,db/tables.py,db/migrations.py, Alembic0083_execution_conversation_key— column + partial unique index, both tracksservices/pull_pilot.py(INTERACTIVE_TRIGGERS),services/pull_coordination_service.py,services/backlog_service.pycanary/snapshot.py,canary/invariants/b02_no_queued_without_slots_full.pyPULL_MIGRATION_STATUS.md,architecture/execution.md,feature-flows/persistent-task-backlog.md, learning fragment,/cso --diffreport (0 findings)Test Plan
cd tests && pytest unit/test_2842_2843_pull_claim_order.py -v— SQLite and PostgreSQL (TEST_POSTGRES_URL); includes an 8-thread concurrent claim on PGalembic upgrade headon a fresh PostgreSQL reaches0083;check_alembic_heads+check_alembic_paritypassdev(IPv6/SSRF address tests under local Python 3.11)Fixes #2842
Fixes #2843
🤖 Generated with Claude Code