Skip to content

feat(pull): interactive turns first, one turn per conversation (#2842, #2843) - #3110

Merged
obasilakis merged 1 commit into
devfrom
feature/2842-interactive-priority
Sep 30, 2026
Merged

obasilakis merged 1 commit into
devfrom
feature/2842-interactive-priority

Conversation

@obasilakis

Copy link
Copy Markdown
Contributor

Summary

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, IntegrityError retry; conversation_key on the enqueue transition
  • db/schema.py, db/tables.py, db/migrations.py, Alembic 0083_execution_conversation_key — column + partial unique index, both tracks
  • services/pull_pilot.py (INTERACTIVE_TRIGGERS), services/pull_coordination_service.py, services/backlog_service.py
  • canary/snapshot.py, canary/invariants/b02_no_queued_without_slots_full.py
  • Docs: PULL_MIGRATION_STATUS.md, architecture/execution.md, feature-flows/persistent-task-backlog.md, learning fragment, /cso --diff report (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 PG
  • Mutations: unique index dropped → concurrent PG test red; ordering disabled → red; conversation skip disabled → red
  • Pull, backlog, canary, schema-parity suites pass on both backends
  • alembic upgrade head on a fresh PostgreSQL reaches 0083; check_alembic_heads + check_alembic_parity pass
  • Full unit suite: 24 failures, identical on clean dev (IPv6/SSRF address tests under local Python 3.11)

Fixes #2842
Fixes #2843

🤖 Generated with Claude Code

…#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 dolho 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.

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, the NOT EXISTS correlation evaluates to NULL, so ungrouped rows are never skipped. On PG, SKIP LOCKED lets two claimers lock different head rows of the same conversation. The partial unique index then makes the second commit raise IntegrityError, and the retry's NOT EXISTS sees the winner. Each attempt opens its own begin(), 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, and release_claim_to_queued all free it with no extra bookkeeping.
  • Dual track. The SQLite migration, the Alembic 0083 revision (chained off 0082_agent_sync_state_divergence), the schema.py DDL, and the tables.py Index (declared so autogenerate won't propose dropping it) all agree.
  • Scope. Keys are only stamped for pull pilots, and drain_next already 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

  1. De-piloting with keyed rows still queued. If an agent leaves the pilot while keyed rows are queued, drain_next claims with no NOT EXISTS. A queued turn whose conversation is still running then hits the index 3 times, gets None, and the head row blocks drain until that turn finishes. It clears itself, but the IntegrityError branch returns silently. A logger.info there would make that stall visible.
  2. 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.
  3. Ordering and the partial index. The CASE sort key means idx_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.
  4. Merge order. Several open PRs also add an 0083_* revision off the same 0082 (#3085, #3021/#3022, #2984). Whichever merges second needs its down_revision repointed, and alembic-head-watch will flag it on the next push to dev.
@obasilakis
obasilakis merged commit 863240f into dev Sep 30, 2026
25 checks passed
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>
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

2 participants