From ac23dfa999870644c1f089b428e6f6cea9c7bf29 Mon Sep 17 00:00:00 2001 From: obasilakis Date: Thu, 3 Sep 2026 13:48:20 +0200 Subject: [PATCH] feat(pull): let scheduled work reach the durable queue (#2391) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit `task_execution_service` — the producer behind the scheduler, i.e. all cron plus webhooks, reminders, loops and fan-out — dispatched with `overflow_policy="reject"` unconditionally. `capacity_manager.acquire` only offers a row to the durable queue when its producer passed `"queue_persistent"`, so `PULL_MODE_PILOT_AGENTS` was structurally inert for the fleet's dominant traffic class no matter how it was set. #2048 made that legible and deliberately deferred fixing it; this is the deferred half. `build_pull_queue_payload` returns a `PersistentTaskPayload` when — and only when — `pull_owns_dispatch` says a pilot owns the trigger. That selects `overflow_policy="queue_persistent"`, the row is enqueued, `execute_task` returns `QUEUED` before any activity, agent call or dispatch marker, and the agent's own worker claims it. `PULL_REACHABLE_TRIGGERS` widens {agent, event} → {agent, event, schedule, webhook, reminder}. Capacity-pressure decision (AC 1): the gate is `pull_owns_dispatch` and nothing else, so scheduled dispatch is byte-for-byte unchanged for every agent outside the pilot allowlist — a fire arriving at capacity still fails fast with "Agent at capacity". An unconditional persistent queue here would have changed that fleet-wide, flag or no flag, which is the risk #2048 declined to take. Under the flag, rejection stops existing for that agent's pullable triggers (the row is never offered a slot); backpressure moves to `max_backlog_depth`, and the error names the backlog rather than parallel slots. #1083 reconciliation (AC 4): the two do not stack, and pull wins by construction rather than precedence — a pull-queued row is never dispatched, so no 202 can arrive and `async_result` is never sent. On a non-pilot agent fire-and-forget is unchanged. They already share the machinery that matters: the eid-keyed slot lease, the lease reaper, and the claim-token CAS terminal. Scheduler: `queued` is not a terminal. `_poll_execution_completion` treated anything but `running` as an outcome, so a queued row would have published `schedule_execution_completed(status=queued)`, been classified a failure and handed to `_maybe_schedule_retry` — duplicating work that was queued and about to run. Triggers left stranded: `loop`, `fan_out` and `a2a` are structural (their caller reads `result.response`; the async fan-out join is #1081 Phase 4). `operator_response` is a scope choice — it records `result.status` as the ent#329 dispatch receipt, and "queued" is not the outcome that reports. #2048's structural tripwire fired as designed and was re-derived, not weakened: it now pins both of this producer's policies AND that the wider one is reachable only through `pull_owns_dispatch` — the assertion #2048 could not make, and the one that keeps the blast radius equal to the allowlist. Docs: PULL_MIGRATION_TESTING.md §9 (M1 expectations, reach table, pilot selection — a cron-driven agent is viable now, which reverses the advice #2048 added), capacity-management.md, task-execution-service.md, architecture.md. Refs #1081, #2048, #1083. Watch #2392 before a side-effect-bearing cron pilot. Co-Authored-By: Claude Opus 5 (1M context) Claude-Session: https://claude.ai/code/session_019Tm7UEkd4G9KQZD5oeLRSa --- docs/memory/architecture.md | 11 +- .../feature-flows/capacity-management.md | 2 +- .../feature-flows/task-execution-service.md | 2 + docs/testing/PULL_MIGRATION_TESTING.md | 140 ++-- src/backend/services/pull_pilot.py | 100 ++- .../services/task_execution_service.py | 206 +++++- src/scheduler/models.py | 5 + src/scheduler/service.py | 29 +- tests/registry.json | 15 +- tests/unit/test_1766_pull_pilot_exclusive.py | 32 +- tests/unit/test_2048_pull_pilot_reach.py | 169 +++-- tests/unit/test_2391_scheduled_pull_reach.py | 614 ++++++++++++++++++ 12 files changed, 1173 insertions(+), 152 deletions(-) create mode 100644 tests/unit/test_2391_scheduled_pull_reach.py diff --git a/docs/memory/architecture.md b/docs/memory/architecture.md index faf52ce16..25fe3026f 100644 --- a/docs/memory/architecture.md +++ b/docs/memory/architecture.md @@ -1443,9 +1443,14 @@ on (ent#364 AC #5). `_VALID_TRIGGERS` (Executions filter), `_TRIGGER_BUCKETS` → its own "Operator queue" analytics bucket (unmapped triggers silently become "Other"), and `_AUTONOMOUS_TRIGGERS` (nobody reads a resume turn's reply, so an - unresolved slash command earns an alert). It is **stranded** for pull mode: - dispatched by a direct backend call, never by `POST /task`, so - `_derive_task_trigger` cannot emit it (#2048). + unresolved slash command earns an alert). It is **stranded** for pull mode. + Until #2391 that was structural — dispatched by a direct backend call, never by + `POST /task`, so `_derive_task_trigger` could not emit it (#2048). Since #2391 + gave `task_execution_service` a pilot-gated `queue_persistent` policy it is a + **choice**: the respond endpoint records `result.status` as the dispatch + receipt (the `operator_resume_dispatch` audit row above and the #525 + idempotency completion), and `queued` is not the outcome that contract + reports. It is absent from `pull_pilot.PULL_REACHABLE_TRIGGERS` deliberately. The owner-facing toggle is in `components/ReliabilityPanel.vue`. diff --git a/docs/memory/feature-flows/capacity-management.md b/docs/memory/feature-flows/capacity-management.md index a12a2e77e..a3d2ab347 100644 --- a/docs/memory/feature-flows/capacity-management.md +++ b/docs/memory/feature-flows/capacity-management.md @@ -296,7 +296,7 @@ APIs. | `/task` async | `src/backend/routers/chat.py` | `queue_persistent` | | `/task` sync long-poll | `src/backend/routers/chat.py` (waits on `sync_waiter`) | `queue_persistent` | | Terminate endpoint | `src/backend/routers/chat.py` | `force_release` | -| `TaskExecutionService` | `src/backend/services/task_execution_service.py` | `reject` (router pre-acquired) | +| `TaskExecutionService` | `src/backend/services/task_execution_service.py` | `reject` (router pre-acquired) — or `queue_persistent` when the agent is a `PULL_MODE_PILOT_AGENTS` pilot and `pull_owns_dispatch` claims the trigger (#2391; `schedule`/`webhook`/`reminder`). The pilot-gated branch is what lets scheduled work be claimed by a pull worker; every non-pilot agent is unchanged. | | Cleanup watchdog | `src/backend/services/cleanup_service.py` | `reclaim_stale` + `release_if_matches` | | Limiter refresher (parked entries, 15s tick) | `src/backend/services/agent_call_limiter.py` (`_tick_sync`) | `SlotService.renew_slot` — direct, not via the facade (#2433) | | Dispatch grant after a park ≥5s | `src/backend/services/task_execution_service.py` (`_on_dispatch_granted`) | `SlotService.renew_slot` — direct, beside `db.restamp_execution_dispatch` (#2433) | diff --git a/docs/memory/feature-flows/task-execution-service.md b/docs/memory/feature-flows/task-execution-service.md index e7a3fec15..84d69b4a1 100644 --- a/docs/memory/feature-flows/task-execution-service.md +++ b/docs/memory/feature-flows/task-execution-service.md @@ -128,6 +128,8 @@ > **Updated 2026-05-11 (#686 UC1):** Interactive `/chat` endpoint now mirrors the service's dispatched-sentinel pattern + real-UUID persistence inline in `routers/chat.py` (parallel of #279). The dispatched-sentinel mechanism is no longer exclusive to `TaskExecutionService.execute_task()`. +> **Updated 2026-09-03 (#2391):** `execute_task` is no longer a `reject`-only producer. `build_pull_queue_payload` returns a `PersistentTaskPayload` when — and only when — `pull_pilot.pull_owns_dispatch(agent, triggered_by)` is true (a `PULL_MODE_PILOT_AGENTS` agent on `schedule` / `webhook` / `reminder`), which selects `overflow_policy="queue_persistent"`; the row is enqueued, `execute_task` returns `TaskExecutionStatus.QUEUED` before any activity, agent call or dispatch marker, and the agent's own worker claims it. Everything else keeps `"reject"` byte-for-byte, so scheduled capacity semantics are unchanged for every non-pilot agent. Two hard preconditions: an existing `execution_id` (the enqueue is a CAS RUNNING→QUEUED on that row) and `slot_already_held=False` (queueing under a held slot would leak it for the lease TTL). #1083 fire-and-forget cannot stack on it — a queued row is never dispatched, so no 202 can arrive. + > **Updated 2026-04-26 (#428):** Slot acquisition/release now goes through [`CapacityManager`](capacity-management.md) (`acquire(overflow_policy="reject")` + `release()`) rather than calling `SlotService` directly. The `slot_already_held` parameter still applies — routers pre-acquire via `CapacityManager` and pass `slot_already_held=True` so the service's `finally` block remains the single release point. ## Overview diff --git a/docs/testing/PULL_MIGRATION_TESTING.md b/docs/testing/PULL_MIGRATION_TESTING.md index 1743113b1..09e707df7 100644 --- a/docs/testing/PULL_MIGRATION_TESTING.md +++ b/docs/testing/PULL_MIGRATION_TESTING.md @@ -363,11 +363,20 @@ name for `''`. ### Pre-flight — pick a viable pilot, then capture a baseline -**First, check the candidate can be piloted at all.** Only `agent` and `event` -traffic can reach the durable queue (#2048) — a cron-driven agent will read ~0 -`pulled` on M1 no matter how the flag is set. Run the pull-eligible-volume query -under M1 below and pick from its `pull_eligible` column *before* flipping -anything; a week of soak on a cron-only agent measures nothing. +**First, check the candidate can be piloted at all.** Five triggers reach the +durable queue — `agent`, `event`, `schedule`, `webhook`, `reminder` — and four +do not (`loop`, `fan_out`, `a2a`, `operator_response`). Run the +pull-eligible-volume query under M1 below and pick from its `pull_eligible` +column *before* flipping anything; a week of soak on an agent whose traffic is +all loops and fan-out measures nothing. + +> **Changed by #2391.** Before it, `task_execution_service` dispatched with +> `overflow_policy="reject"` unconditionally, so **no** cron, webhook or reminder +> row could ever be claimed — the pilot flag was inert for the fleet's dominant +> traffic class, and a cron-driven agent was explicitly not a viable pilot. It +> now dispatches with `"queue_persistent"` when `pull_owns_dispatch` says a pilot +> owns the trigger, and with `"reject"` otherwise. **A cron-driven agent is a +> viable pilot now** — see "Choosing a pilot" below. Then run M4 + M6 + M7 + M9 on the pilot for the 7 days *preceding* the flip and keep the output. Without a baseline, a 4% success-rate dip during the soak is @@ -435,8 +444,9 @@ WHERE agent_name = '' AND started_at::timestamptz > now() - interval '24 hours'; ``` -Expected after the #1766 gate: **`agent` and `event` rows ~100% `pulled`; every -other trigger 100% `pushed`.** That second half is correct behaviour, not a +Expected after the #1766 gate + #2391: **`agent`, `event`, `schedule`, `webhook` +and `reminder` rows ~100% `pulled`; `loop`, `fan_out`, `a2a` and +`operator_response` 100% `pushed`.** That second half is correct behaviour, not a fault — read the reach table below before concluding anything from it. Confirm the split per trigger with: @@ -450,65 +460,119 @@ WHERE agent_name = '' GROUP BY 1 ORDER BY 1; ``` -#### Which triggers can be pulled at all (#2048) +#### Which triggers can be pulled at all (#2048, widened by #2391) -`PULL_MODE_PILOT_AGENTS` routes **agent-to-agent work only**. `capacity_manager` -offers a row to the durable queue solely when its producer passed -`overflow_policy="queue_persistent"`, and one of the three producers does: +`capacity_manager` offers a row to the durable queue solely when its producer +passed `overflow_policy="queue_persistent"`. Two of the three producers can: | Producer | Carries | `overflow_policy` | Pullable? | |---|---|---|---| -| `task_execution_service` | scheduler — **all cron** — + fan-out | `reject` | **No** | +| `task_execution_service` | scheduler — **all cron** — webhooks, reminders, loops, fan-out, A2A, operator resumes | `queue_persistent` **when `pull_owns_dispatch` is true**, else `reject` | **Yes**, for `schedule` / `webhook` / `reminder` | | `dispatch_admission_service` | sequential `chat_with_agent`, human chat | `queue_in_memory` | **No** | -| `chat_execution_service` (`POST /task`) | parallel `chat_with_agent`, MCP/manual task | `queue_persistent` | **Yes** | +| `chat_execution_service` (`POST /task`) | parallel `chat_with_agent`, MCP/manual task | `queue_persistent` | **Yes**, for `agent` / `event` | `POST /task` can only derive `triggered_by ∈ {self_task, agent, mcp, manual, -event}`; intersect with `_AUTONOMOUS_TRIGGERS` and the pullable set is exactly -**`{agent, event}`** (`pull_pilot.PULL_REACHABLE_TRIGGERS`). `schedule`, -`webhook`, `loop`, `fan_out` and `reminder` are declared autonomous but reach -dispatch through the `reject` producer, so **they are pushed no matter what the -flag says.** +event}`, contributing `agent` + `event`. `task_execution_service` contributes the +autonomous triggers with **no synchronous result consumer** — `schedule`, +`webhook` and `reminder`, all three of which reach it from the scheduler, which +dispatches `async_mode=True` and then polls the DB for the terminal. The pullable +set is therefore **`{agent, event, schedule, webhook, reminder}`** +(`pull_pilot.PULL_REACHABLE_TRIGGERS`). + +The four that stay pushed are stranded on their **caller**, not on the +producer's policy — a queued row is claimed and run later and returns nothing for +any of them to read. Three are structural: `loop` (renders the next iteration +from `result.response`), `fan_out` (builds each `FanOutTaskResult` from the +returned result — the async join is #1081 Phase 4) and `a2a` (turns the result +into the JSON-RPC artifact it hands a remote caller). `operator_response` is out +by **choice**: it dispatches through the same producer #2391 widened, but the +respond endpoint records `result.status` as the dispatch receipt (the +`operator_resume_dispatch` audit row + the #525 idempotency completion) for a +turn that spends money on a person's answer (ent#329). Widening it is a +follow-up, not a blocker. **Reading a pushed autonomous row.** Two different situations, previously indistinguishable: -- `triggered_by` ∈ `{schedule, webhook, loop, fan_out, reminder}` → **expected.** +- `triggered_by` ∈ `{loop, fan_out, a2a, operator_response}` → **expected.** The flag is applied and correct; the dispatch topology is the limit. The backend logs this once per (agent, trigger) — grep `[#2048]` to confirm: ```bash docker logs trinity-backend 2>&1 | grep '\[#2048\]' ``` -- `triggered_by` ∈ `{agent, event}` → **a real fault.** This is the only case - where the old advice applies: check the backend actually has - `PULL_MODE_PILOT_AGENTS` in its env (and see G1 in §6 — compose must forward - it). +- `triggered_by` ∈ `{agent, event, schedule, webhook, reminder}` → **a real + fault.** This is the case where the old advice applies: check the backend + actually has `PULL_MODE_PILOT_AGENTS` in its env (and see G1 in §6 — compose + must forward it). -#### Choosing a pilot: a cron-only agent cannot be piloted +#### Choosing a pilot: cron-driven agents are viable since #2391 -Because cron cannot reach the queue, **an agent whose traffic is mostly -`schedule` will read ~0 `pulled` on M1 no matter how the flag is set, and the -soak measures nothing.** Verify pull-eligible volume *before* selecting a pilot, -not after a week of soak: +Verify pull-eligible volume *before* selecting a pilot, not after a week of soak: ```sql SELECT agent_name, COUNT(*) AS total, COUNT(*) FILTER (WHERE triggered_by = 'schedule') AS cron, - COUNT(*) FILTER (WHERE triggered_by IN ('agent','event')) AS pull_eligible + COUNT(*) FILTER (WHERE triggered_by IN + ('agent','event','schedule','webhook','reminder')) AS pull_eligible, + COUNT(*) FILTER (WHERE triggered_by IN + ('loop','fan_out','a2a','operator_response')) AS not_pullable FROM schedule_executions WHERE started_at::timestamptz > now() - interval '7 days' GROUP BY 1 ORDER BY pull_eligible DESC; ``` -Pick from the `pull_eligible` column. In practice that means a **fan-in hub fed -by parallel `chat_with_agent`** — on `eu2` exactly one agent -(`cornelius-oracle`, 329/329) had the volume, while the cron-driven oracles sat -at 2–5 and `brier-hq` at 0. - -Extending pull to cron is the issue's deferred Option 2 (give -`task_execution_service` a `queue_persistent` policy under the pilot flag); it is -**not** implemented. `tests/unit/test_2048_pull_pilot_reach.py` fails if that -producer's policy changes, so this section cannot go stale silently. +Pick from the `pull_eligible` column. **The advice this replaces was the +opposite**: before #2391 cron could not reach the queue, so a cron-driven agent +read ~0 `pulled` on M1 no matter how the flag was set, and the only viable pilot +was a **fan-in hub fed by parallel `chat_with_agent`** — on `eu2` exactly one +agent (`cornelius-oracle`, 329/329) qualified, while the cron-driven oracles sat +at 2–5 and `brier-hq` at 0. Under the same query today those cron-driven oracles +score ~128 each. Prefer one of them for the next soak arm: they exercise the +traffic class the fleet actually runs, which `cornelius-oracle` never did. + +`tests/unit/test_2048_pull_pilot_reach.py` pins both this producer's policies +*and* that the wider one is gated on `pull_owns_dispatch`, so this section +cannot go stale silently. + +#### What #2391 changed under capacity pressure — and what it did not + +The gate is `pull_owns_dispatch` and nothing else, so: + +- **Flag OFF (every agent, by default): unchanged, byte-for-byte.** The policy + stays `reject`; a scheduled fire arriving at capacity still fails fast with + `Agent at capacity (N/N parallel tasks running)` and a FAILED row. Giving this + producer an *unconditional* persistent queue would have changed that for the + whole fleet whether or not pull was enabled — the reliability-spine change + #2048 declined to bundle. +- **Flag ON: rejection stops existing for that agent's pullable triggers.** The + row is never offered a slot (`acquire`'s `pull_exclusive` branch skips the + ZADD), so "at capacity" is not a state that can reject it. Capacity becomes + physical — the agent's worker pool. Backpressure moves to + `agent_ownership.max_backlog_depth`: a full backlog still raises + `CapacityFull`, at a deeper threshold, and the row's error names the backlog + rather than parallel slots. **Watch `max_backlog_depth` on a cron-driven + pilot**: cron fires on a fixed cadence regardless of whether the agent is + keeping up, so a slow agent accumulates queue depth that a `reject` policy + used to discard. M8 (queue starvation) is the query to sample for it. +- **#1083 fire-and-forget does not stack.** `schedule` and `webhook` are eligible + for both, but a pull-queued row is never dispatched, so no 202 ACK can arrive + and `async_result` is never sent. Pull wins by construction, not precedence. + On a non-pilot agent #1083 behaves exactly as before. The two share the + machinery that matters anyway — the eid-keyed slot lease, the lease reaper, and + the `claim_token`-gated CAS terminal write. +- **The scheduler polls through `queued`.** `_poll_execution_completion` treated + anything but `running` as an outcome, so a queued row would have been published + as `schedule_execution_completed(status=queued)`, classified a failure, and + handed to `_maybe_schedule_retry` — duplicating work that was queued and about + to run. `queued` is now in `_NON_TERMINAL_POLL_STATES`. If the poll deadline + (task timeout + buffer) expires while a row is still queued, the scheduler logs + it and stops; recovery belongs to the lease reaper and the backlog + maintenance sweep, not to `cleanup_service`'s stale-`running` sweep. +- **#2392 becomes load-bearing.** This moves the fleet's dominant traffic class + onto a mechanism built on re-running the same execution, and `effect_guard` + still fails open when the agent omits the execution id. Do not run a + side-effect-bearing cron pilot before reading it. **M2 — split over time.** Confirms the flip took effect at the moment you think it did, and shows drift. diff --git a/src/backend/services/pull_pilot.py b/src/backend/services/pull_pilot.py index 57a9f4365..1041a934d 100644 --- a/src/backend/services/pull_pilot.py +++ b/src/backend/services/pull_pilot.py @@ -38,36 +38,61 @@ def is_pull_pilot_agent(agent_name: str) -> bool: return agent_name in _pilot_allowlist() -# Which autonomous triggers can STRUCTURALLY reach the durable queue (#2048). +# Which autonomous triggers can STRUCTURALLY reach the durable queue +# (#2048 named the set; #2391 widened it). # # This is dispatch topology, not policy. ``capacity_manager.acquire`` can only # hand a row to the queue when its producer passed -# ``overflow_policy="queue_persistent"``, and exactly one of the three producers -# does: +# ``overflow_policy="queue_persistent"``, and two of the three producers do: # -# task_execution_service → "reject" — scheduler/ALL cron, fan-out +# task_execution_service → "reject" | "queue_persistent" ← #2391, pilot-gated +# scheduler/ALL cron, webhooks, reminders, fan-out, +# loops, A2A, operator resumes # dispatch_admission_svc → "queue_in_memory" — sequential chat, human chat -# chat_execution_service → "queue_persistent" — POST /task ← only +# chat_execution_service → "queue_persistent" — POST /task # # ``POST /task`` derives its trigger in ``_derive_task_trigger``, which can only -# produce ``{self_task, agent, mcp, manual, event}``. Intersect that with the -# autonomous set and two survive: ``agent`` (agent-to-agent ``chat_with_agent``) -# and ``event`` (#1578's task-completion loopback, which POSTs to the same route -# under an internal-secret-gated ``X-Event-Trigger``). +# produce ``{self_task, agent, mcp, manual, event}``; from the autonomous set +# that contributes ``agent`` (agent-to-agent ``chat_with_agent``) and ``event`` +# (#1578's task-completion loopback, which POSTs to the same route under an +# internal-secret-gated ``X-Event-Trigger``). # -# ``schedule``, ``webhook``, ``loop``, ``fan_out`` and ``reminder`` are declared -# autonomous but reach dispatch through the ``"reject"`` producer, so on a -# cron-driven agent the pilot flag is INERT. That was previously invisible: -# ``pull_owns_dispatch`` is consulted behind an ``overflow_policy == -# "queue_persistent"`` short-circuit, so for a cron row it is never called at -# all, and the row takes the push path indistinguishably from the flag being -# unset. Naming the reachable subset here is what makes the gap legible; see -# ``note_unreachable_pull_trigger`` for the runtime signal at the other producer. +# ``task_execution_service`` contributes the three autonomous triggers that have +# **no synchronous result consumer** — ``schedule``, ``webhook`` and +# ``reminder``. All three reach it from the scheduler, which dispatches with +# ``async_mode=True`` and then polls ``schedule_executions`` for the terminal, so +# nobody is holding a coroutine waiting on a return value. That is exactly the +# property a claim-and-run-later queue needs, and it is the same property +# ``ASYNC_DISPATCH_ELIGIBLE_TRIGGERS`` (#1083) selects on. +# +# The four autonomous triggers that stay OUT all have a caller that reads the +# returned ``TaskExecutionResult``, so a row queued for a worker to claim later +# gives that caller nothing to read. Three are structural; one is a scope choice. +# +# loop ``loop_service`` renders the next iteration's template +# from ``result.response`` (#740). Structural. +# fan_out ``fan_out_service`` builds each ``FanOutTaskResult`` from +# the returned result. Structural — the async fan-out join +# is #1081 Phase 4, not this change. +# a2a ``routers/a2a``'s ``message/send`` consumes +# ``result.response`` to build the JSON-RPC artifact it +# hands back to a remote caller (ent#157). Structural. +# operator_response NOT structural — ``operator_resume_service`` dispatches +# through this same producer, so it *could* be queued. It is +# out by choice: the respond endpoint records +# ``result.status`` as the dispatch receipt (the +# ``operator_resume_dispatch`` audit row and the #525 +# idempotency completion), and ent#329 exists because this +# spends money on a person's answer — "queued" is not the +# outcome that contract reports. Widening it is a deliberate +# follow-up, not an oversight. # # Intersected with ``_AUTONOMOUS_TRIGGERS`` rather than replacing it, so this # stays a NARROWING of the single source of truth: dropping a trigger there # still drops it here, and widening reach is a deliberate edit to this set. -PULL_REACHABLE_TRIGGERS = frozenset({"agent", "event"}) +PULL_REACHABLE_TRIGGERS = frozenset( + {"agent", "event", "schedule", "webhook", "reminder"} +) def pull_owns_dispatch(agent_name: str, triggered_by: Optional[str]) -> bool: @@ -107,12 +132,12 @@ def pull_owns_dispatch(agent_name: str, triggered_by: Optional[str]) -> bool: # never a second copy that can drift. from services.task_execution_service import _AUTONOMOUS_TRIGGERS - # Narrowed to what dispatch can actually deliver (#2048). A no-op on - # today's runtime — this predicate is only consulted behind an - # ``overflow_policy == "queue_persistent"`` check, and that producer can - # only emit ``{agent, event}`` from the autonomous set — but it stops the - # predicate from *claiming* reach it does not have. See - # ``PULL_REACHABLE_TRIGGERS``. + # Narrowed to what dispatch can actually deliver. Since #2391 this + # predicate is also the PRODUCER-SIDE gate: ``task_execution_service`` + # asks it whether to dispatch with ``overflow_policy="queue_persistent"`` + # instead of ``"reject"``, so a False here is the exact condition under + # which scheduled capacity semantics stay byte-for-byte as they were. + # See ``PULL_REACHABLE_TRIGGERS`` for why each trigger is in or out. return triggered_by in (_AUTONOMOUS_TRIGGERS & PULL_REACHABLE_TRIGGERS) except Exception: # noqa: BLE001 — unresolvable trigger set ⇒ push, as today logger.warning( @@ -127,7 +152,7 @@ def pull_owns_dispatch(agent_name: str, triggered_by: Optional[str]) -> bool: # every cron fire, so an un-deduped warning would emit thousands of identical # lines a day and train operators to filter it out — which is how a signal # meant to be noticed becomes noise. Unbounded growth is not a concern: the key -# space is (pilot agents × 5 unreachable triggers), and only pilots ever reach it. +# space is (pilot agents × 4 unreachable triggers), and only pilots ever reach it. _UNREACHABLE_NOTED: Set[tuple] = set() @@ -141,15 +166,21 @@ def note_unreachable_pull_trigger(agent_name: str, triggered_by: Optional[str]) applied, or the trigger structurally cannot be pulled — and BOTH were silent. ``PULL_MIGRATION_TESTING.md`` §9 M1 told operators that a pushed autonomous row means "the producer gate did not engage — check the backend - actually has ``PULL_MODE_PILOT_AGENTS`` in its env", which for a ``schedule`` + actually has ``PULL_MODE_PILOT_AGENTS`` in its env", which for a stranded row sends them hunting for an env var that is present and correct. Deliberately a log line and not a metric or an operator-queue item: this is a property of the build, not an incident. It is constant for a given (agent, trigger) until the dispatch topology changes, so alerting on it - would fire forever on any cron-driven pilot. It exists so that the ONE - operator who runs M1, sees ~0 ``pulled``, and asks why gets a truthful - answer from the backend log instead of a wrong one from the runbook. + would fire forever on any pilot that runs loops or fan-out. It exists so + that the ONE operator who runs M1, sees a pushed autonomous row, and asks + why gets a truthful answer from the backend log instead of a wrong one from + the runbook. + + #2391 shrank what it fires on rather than retiring it: ``schedule``, + ``webhook`` and ``reminder`` are reachable now, so the remaining callers are + ``loop`` / ``fan_out`` / ``a2a`` / ``operator_response``, which are stranded + on a synchronous result consumer rather than on the producer's policy. Fail-safe like everything else here: a bookkeeping error must never interfere with dispatching a real execution. @@ -173,11 +204,12 @@ def note_unreachable_pull_trigger(agent_name: str, triggered_by: Optional[str]) logger.warning( "[#2048] pull pilot %r: trigger %r is autonomous but CANNOT reach the " - "durable queue — this producer dispatches with overflow_policy='reject', " - "so the row is pushed regardless of PULL_MODE_PILOT_AGENTS. The flag is " - "applied and correct; the dispatch topology is the limit. Reachable " - "triggers today: %s. A cron-only agent is not a viable soak pilot " - "(#1766) — see docs/testing/PULL_MIGRATION_TESTING.md §9.", + "durable queue — its caller consumes the execution result synchronously, " + "so this row is dispatched with overflow_policy='reject' and pushed " + "regardless of PULL_MODE_PILOT_AGENTS. The flag is applied and correct; " + "the dispatch topology is the limit. Reachable triggers today: %s. " + "Cron/webhook/reminder work IS pullable since #2391 — see " + "docs/testing/PULL_MIGRATION_TESTING.md §9.", agent_name, triggered_by, sorted(PULL_REACHABLE_TRIGGERS), ) return True diff --git a/src/backend/services/task_execution_service.py b/src/backend/services/task_execution_service.py index 7f003e57d..65cf3fa01 100644 --- a/src/backend/services/task_execution_service.py +++ b/src/backend/services/task_execution_service.py @@ -44,6 +44,7 @@ CapacityFull, CircuitOpen, EphemeralBudgetExhausted, + PersistentTaskPayload, get_capacity_manager, ) from services.dispatch_breaker import DispatchBreaker @@ -53,7 +54,7 @@ from services.platform_audit_service import AuditEventType, platform_audit_service # #2048: stdlib-only leaf by construction, so this cannot cycle back through the # capacity stack at import time (its own reference to this module is lazy). -from services.pull_pilot import note_unreachable_pull_trigger +from services.pull_pilot import note_unreachable_pull_trigger, pull_owns_dispatch from services.settings_service import settings_service from utils.credential_sanitizer import sanitize_dict, sanitize_execution_log, sanitize_response, sanitize_text from services.tool_call_summary import extract_tool_calls @@ -556,6 +557,112 @@ def dispatch_async_eligible(triggered_by: Optional[str]) -> bool: return False +def build_pull_queue_payload( + *, + agent_name: str, + triggered_by: str, + execution_id: Optional[str], + message: str, + model: Optional[str], + allowed_tools: Optional[list], + system_prompt: Optional[str], + timeout_seconds: Optional[int], + resume_session_id: Optional[str], + subscription_id: Optional[str], + source_user_id: Optional[int], + source_user_email: Optional[str], + source_agent_name: Optional[str], + slot_already_held: bool, +) -> Optional[PersistentTaskPayload]: + """The #2391 producer gate: the overflow payload that lets THIS producer put + a row on the durable queue, or ``None`` to keep today's ``"reject"`` policy. + + Until #2391 this producer dispatched with ``overflow_policy="reject"`` + unconditionally, so the scheduler's traffic — all cron, plus webhooks and + reminders — could never be claimed by a pull worker no matter how + ``PULL_MODE_PILOT_AGENTS`` was set (#2048 made that legible; it did not fix + it). Returning a payload here flips the policy to ``"queue_persistent"`` for + this one dispatch, which is what ``capacity_manager.acquire``'s pull branch + needs to hand the row to ``BacklogService`` instead of admitting it. + + **The capacity-pressure decision (#2391 AC #1).** The gate is + ``pull_owns_dispatch`` — the SAME predicate that already decides pull + ownership — and nothing else. So: + + * flag OFF for this agent (every agent, by default): ``None`` → the policy + stays ``"reject"`` and a fire arriving at capacity still fails fast with + "Agent at capacity", byte-for-byte as before. **Scheduled dispatch + semantics do not change for any agent that is not a pull pilot.** The + alternative — giving this producer an unconditional persistent queue — + would have changed the fleet's dominant traffic class under capacity + pressure whether or not pull was enabled, which is precisely the risk + #2048 declined to take. + * flag ON: the row is never offered a slot at all (``acquire``'s + ``pull_exclusive`` branch skips the ZADD), so "at capacity" stops being a + state that can reject it. Capacity becomes physical — the agent's worker + pool — which is #1081 Phase 5, pilot-scoped. Backpressure moves from + reject-at-dispatch to ``agent_ownership.max_backlog_depth``: a full + backlog still raises ``CapacityFull``, just at a deeper threshold and + with ``reason="persistent_full"``. + + **Interaction with #1083 fire-and-forget (AC #4).** Both mechanisms exist to + let a turn outlive the request, and they must not stack. Pull wins by + construction, not by precedence: a pull-queued row is never dispatched, so + there is no HTTP call for the agent to ACK with 202 and ``async_result`` is + never sent. ``dispatch_async_eligible`` is still evaluated in + ``execute_task`` but its result is unreachable once this returns a payload. + The two share the machinery that matters anyway — the eid-keyed slot lease, + the lease reaper, and the ``claim_token``-gated CAS terminal write — so the + recovery story is single, not layered. + + Two hard preconditions, both about not corrupting state we do not own: + + * ``execution_id`` must exist — ``BacklogService.enqueue`` transitions an + EXISTING row RUNNING→QUEUED under a CAS; with no row there is nothing to + transition. + * ``slot_already_held`` must be False — the caller (a drain, or the /task + router's pre-flight) owns a real slot and releases it in our ``finally``. + Queueing under a held slot would leak it for the lease TTL. + + Never raises: ``pull_owns_dispatch`` already fails safe to push, and the + request build is pure construction. The dangerous direction is queueing work + that should have been pushed. + """ + if slot_already_held or not execution_id: + return None + if not pull_owns_dispatch(agent_name, triggered_by): + return None + + # Lazy: several unit suites stub `models` with a partial module, and a + # top-level import of a name they omit would break importing this service + # entirely (the `_resolve_agent_runtime` precedent above). + from models import ParallelTaskRequest + + # `system_prompt` is deliberately the CALLER's raw override, NOT the + # composed prompt: `pull_coordination_service._compose_pull_system_prompt` + # rebuilds platform prompt + execution context around it at claim time + # (#1629), so composing here would double the platform preamble. + request = ParallelTaskRequest( + message=message, + model=model, + allowed_tools=allowed_tools, + system_prompt=system_prompt, + timeout_seconds=timeout_seconds, + async_mode=True, + resume_session_id=resume_session_id, + ) + return PersistentTaskPayload( + request=request, + effective_timeout=int(timeout_seconds or 900), + user_id=source_user_id, + user_email=source_user_email, + subscription_id=subscription_id, + x_source_agent=source_agent_name, + triggered_by=triggered_by, + collaboration_activity_id=None, + ) + + # Strong references to fire-and-forget breaker tasks. asyncio's event loop holds # only a WEAK reference to a bare ``create_task`` result, so an un-referenced task # can be garbage-collected mid-flight (the backlog drain would silently vanish). @@ -941,6 +1048,11 @@ async def execute_task( # returns 200, which falls through to the synchronous handling below. When # the agent ACKs 202 we hand the slot lease to the callback and skip the # `finally` release. Best-effort read; defaults off. + # #2391: pull and fire-and-forget do NOT stack. When the pilot flag + # owns this trigger the row is queued below and `execute_task` returns + # before any agent call, so this value is computed and never used — + # `async_result` is never sent and no 202 can arrive. Push dispatch is + # the only path on which fire-and-forget is reachable. async_dispatch = dispatch_async_eligible(triggered_by) # Set True once a 202 ACK hands the slot lease to the result callback, so # the `finally` does NOT release it (the callback/reaper owns it now). @@ -983,6 +1095,30 @@ async def execute_task( # try so no handler can NameError on a pre-dispatch exception. state = _AttemptState(start_time=datetime.utcnow()) + # ---- #2391: does this dispatch belong on the durable queue? -------- + # Evaluated here, where every field the queued row needs is in scope. + # None (the default for every non-pilot agent) keeps overflow_policy at + # "reject" and this whole path byte-for-byte as it was. Non-None means a + # pull pilot owns this trigger: the row is queued, the agent's worker + # claims it, and #1083's async ACK never comes into play because no + # dispatch happens at all. + pull_overflow_payload = build_pull_queue_payload( + agent_name=agent_name, + triggered_by=triggered_by, + execution_id=execution_id, + message=message, + model=model, + allowed_tools=allowed_tools, + system_prompt=system_prompt, + timeout_seconds=timeout_seconds, + resume_session_id=resume_session_id, + subscription_id=subscription_id, + source_user_id=source_user_id, + source_user_email=source_user_email, + source_agent_name=source_agent_name, + slot_already_held=slot_already_held, + ) + # Wrap entire execution flow to ensure execution status is updated on any failure. # This fixes issue #90 where exceptions during slot acquisition left executions # stuck in 'running' status with NULL session_id and duration_ms. @@ -997,6 +1133,7 @@ async def execute_task( breaker_enabled=breaker_enabled, slot_already_held=slot_already_held, capacity=capacity, + pull_overflow_payload=pull_overflow_payload, ) if admission_denied is not None: return admission_denied @@ -1224,14 +1361,24 @@ async def _admission_gate( breaker_enabled: bool, slot_already_held: bool, capacity, + pull_overflow_payload: Optional[PersistentTaskPayload] = None, ) -> tuple[bool, Optional[TaskExecutionResult]]: """Step 2 of execute_task: acquire the capacity slot (or refuse). - Returns ``(slot_acquired, denial)``. A non-None *denial* is the - terminal ``TaskExecutionResult`` for a CapacityFull / CircuitOpen / - EphemeralBudgetExhausted fast-fail (the FAILED row is already - written); the caller returns it verbatim. Any *other* exception - from ``capacity.acquire`` propagates, exactly as it did inline. + Returns ``(slot_acquired, denial)``. A non-None *denial* is a terminal + ``TaskExecutionResult`` the caller returns verbatim — a CapacityFull / + CircuitOpen / EphemeralBudgetExhausted fast-fail (the FAILED row is + already written), or, since #2391, a QUEUED handoff (see below). Any + *other* exception from ``capacity.acquire`` propagates, exactly as it + did inline. + + *pull_overflow_payload* (#2391) is non-None only when + ``build_pull_queue_payload`` decided this dispatch belongs to a pull + pilot's durable queue. It selects ``overflow_policy="queue_persistent"`` + instead of ``"reject"`` and is the payload ``BacklogService`` persists. + Not a "denial" in any failure sense: the row is QUEUED and a worker will + claim it, so the caller must stop — no activity, no agent call, no slot + to release. """ slot_acquired = slot_already_held # ---- 2. Acquire capacity slot ------------------------------------ @@ -1251,6 +1398,14 @@ async def _admission_gate( # from the flag being unset. Say so once per (agent, trigger). # Diagnostic only; never raises, never affects dispatch. note_unreachable_pull_trigger(agent_name, triggered_by) + # #2391: the ONLY thing that widens this producer's policy is a + # pilot-gated payload. Written as an if/else over two literals, not + # a ternary, so `test_2048_pull_pilot_reach.py`'s structural scan + # can still read both policies out of this module's source. + if pull_overflow_payload is not None: + overflow_policy = "queue_persistent" + else: + overflow_policy = "reject" try: cap_result = await capacity.acquire( agent_name=agent_name, @@ -1258,15 +1413,42 @@ async def _admission_gate( max_concurrent=max_parallel_tasks, message_preview=message[:100] if message else "", timeout_seconds=timeout_seconds, - overflow_policy="reject", + overflow_policy=overflow_policy, + overflow_payload=pull_overflow_payload, breaker_enabled=breaker_enabled, ) + if cap_result.state == "queued_persistent": + # #2391: handed to the durable queue; the agent's own worker + # claims it via GET /api/internal/next-task and reports the + # terminal through the claim-token CAS. Nothing further to do + # here — and critically no slot was ZADDed, so the `finally` + # in execute_task must not release one. + logger.info( + f"[TaskExecService] Pull pilot {agent_name}: queued " + f"execution {execution_id} (trigger={triggered_by}) for " + f"worker claim instead of pushing (#2391)" + ) + return False, TaskExecutionResult( + execution_id=execution_id or "", + status=TaskExecutionStatus.QUEUED, + response="", + ) slot_acquired = cap_result.state == "admitted" - except CapacityFull: - error_msg = ( - f"Agent at capacity ({max_parallel_tasks}/{max_parallel_tasks} " - f"parallel tasks running)" - ) + except CapacityFull as e: + # #2391: on the pull path "at capacity" means the DURABLE BACKLOG + # is full (max_backlog_depth), not that parallel slots are busy — + # the pull branch never asks for a slot. Name the real limit or + # the operator debugs the wrong number. + if getattr(e, "reason", None) == "persistent_full": + error_msg = ( + "Agent backlog full (max_backlog_depth reached); " + "queued task rejected" + ) + else: + error_msg = ( + f"Agent at capacity ({max_parallel_tasks}/{max_parallel_tasks} " + f"parallel tasks running)" + ) if execution_id: db.update_execution_status( execution_id=execution_id, diff --git a/src/scheduler/models.py b/src/scheduler/models.py index 4c97294f4..e3c70905e 100644 --- a/src/scheduler/models.py +++ b/src/scheduler/models.py @@ -13,6 +13,11 @@ class ExecutionStatus(str, Enum): """Status of a schedule execution.""" + # #2391: a dispatch the backend handed to the durable pull queue instead of + # pushing (pull-pilot agents only). NON-terminal — the agent's worker claims + # the row back to `running` — so `_poll_execution_completion` must keep + # polling through it, exactly as it does through `running`. + QUEUED = "queued" RUNNING = "running" SUCCESS = "success" FAILED = "failed" diff --git a/src/scheduler/service.py b/src/scheduler/service.py index 59e19c762..1e00e1170 100644 --- a/src/scheduler/service.py +++ b/src/scheduler/service.py @@ -76,6 +76,12 @@ def _reminder_outcome_unknown(exc: BaseException) -> bool: _POLL_DEADLINE_WHEN_NULL = 7200 +# #2391: statuses `_poll_execution_completion` must poll THROUGH rather than +# treat as an outcome. `queued` joined `running` when the scheduler's dispatch +# became able to land on the durable pull queue. +_NON_TERMINAL_POLL_STATES = (ExecutionStatus.RUNNING, ExecutionStatus.QUEUED) + + class SchedulerService: """ Manages scheduled task execution for agents. @@ -1461,7 +1467,14 @@ async def _poll_execution_completion( logger.warning(f"Execution {execution_id} not found in DB during polling (poll #{poll_count})") continue - if execution.status != ExecutionStatus.RUNNING: + # #2391: `queued` is NOT a terminal. A pull-pilot agent's scheduled + # row is handed to the durable queue and sits there until a worker + # claims it back to `running`; treating that as "completed" would + # publish a bogus schedule_execution_completed(status=queued), + # classify it as a failure (anything != success is), and hand it to + # `_maybe_schedule_retry` — a duplicate run of work that is queued + # and about to run. Poll through it like `running`. + if execution.status not in _NON_TERMINAL_POLL_STATES: logger.info( f"Execution {execution_id} completed: status={execution.status} " f"(polled {poll_count} times)" @@ -1478,7 +1491,10 @@ async def _poll_execution_completion( if poll_count % 6 == 0: # Log every ~60s at default 10s interval elapsed = int(time.monotonic() - (deadline - effective_timeout - 60)) - logger.info(f"Execution {execution_id} still running ({elapsed}s elapsed, poll #{poll_count})") + logger.info( + f"Execution {execution_id} still {execution.status} " + f"({elapsed}s elapsed, poll #{poll_count})" + ) raise Exception( f"Polling deadline exceeded for execution {execution_id} " @@ -1568,15 +1584,18 @@ async def _poll_and_finalize( logger.error(f"Background poll for execution {execution_id} failed: {e}") # Check if backend already finalized the execution current = self.db.get_execution(execution_id) - if current and current.status != ExecutionStatus.RUNNING: + if current and current.status not in _NON_TERMINAL_POLL_STATES: logger.info( f"Execution {execution_id} already finalized as '{current.status}' " "— background poll error is benign" ) else: - # Execution stuck in running state - cleanup service will recover + # #2391: name the state we actually saw. A row still `queued` is + # waiting for a pull worker, and its owner is the lease reaper / + # backlog maintenance, not `cleanup_service`'s stale-running sweep. logger.warning( - f"Execution {execution_id} may be stuck in 'running' state — " + f"Execution {execution_id} may be stuck in " + f"'{current.status if current else 'unknown'}' state — " "cleanup service will recover within 5 minutes" ) diff --git a/tests/registry.json b/tests/registry.json index 4110567c1..1321dc89a 100644 --- a/tests/registry.json +++ b/tests/registry.json @@ -1789,7 +1789,7 @@ "reliability", "pull-migration" ], - "description": "PULL_MODE_PILOT_AGENTS can only route agent-to-agent work (#2048). _AUTONOMOUS_TRIGGERS declares seven triggers autonomous and pull_owns_dispatch consulted all seven, but capacity_manager only offers a row to the durable queue when its producer passed overflow_policy='queue_persistent', and exactly one of the three producers does: task_execution_service uses 'reject' (scheduler/ALL cron + fan-out), dispatch_admission_service uses 'queue_in_memory' (sequential chat, human chat), only chat_execution_service (POST /task) uses 'queue_persistent'. So on a cron-driven agent the pilot flag is INERT and #1766's soak measures nothing - silently, because the gate sits behind an overflow_policy short-circuit, so for a schedule row pull_owns_dispatch is never called at all and the row is pushed indistinguishably from the flag being unset, while PULL_MIGRATION_TESTING.md section 9 M1 told the operator to go check an env var that is present and correct. Issue decision was Option 1 (make the reach honest, do not extend it). Corrects the issue's own table: the reachable set is {agent, event}, not {agent} - #1578's task-completion loopback POSTs to the same /task route and is tagged 'event' via an internal-secret-gated X-Event-Trigger, with no capacity bypass. Tests pin: the reach set derived from what _derive_task_trigger can emit rather than hardcoded; every autonomous trigger classified as reachable or stranded so a new one cannot join the stranded set unreviewed; the narrowing is a runtime NO-OP over the only producer that consults the predicate (it stops a false claim, it does not change a verdict); and note_unreachable_pull_trigger reports a stranded row once per (agent, trigger), never for a non-pilot / reachable / interactive trigger, and never raises. A structural guard asserts the three producers' overflow_policy literals, so implementing the deferred Option 2 (wiring cron to queue_persistent) fails this test and forces PULL_REACHABLE_TRIGGERS and section 9's prose to be revisited in the same change - a silently-stale reach set being exactly the bug. Also updates test_1766: its 'pilot owns EVERY autonomous trigger' case asserted reach the system never had (it passed only by calling the predicate outside the context that constrains it), and its producer-gate case drove a queue_persistent acquire carrying 'schedule' - a combination production cannot produce - now driven with 'agent', the shape that actually occurs." + "description": "PULL_MODE_PILOT_AGENTS can only route agent-to-agent work (#2048). _AUTONOMOUS_TRIGGERS declares seven triggers autonomous and pull_owns_dispatch consulted all seven, but capacity_manager only offers a row to the durable queue when its producer passed overflow_policy='queue_persistent', and exactly one of the three producers does: task_execution_service uses 'reject' (scheduler/ALL cron + fan-out), dispatch_admission_service uses 'queue_in_memory' (sequential chat, human chat), only chat_execution_service (POST /task) uses 'queue_persistent'. So on a cron-driven agent the pilot flag is INERT and #1766's soak measures nothing - silently, because the gate sits behind an overflow_policy short-circuit, so for a schedule row pull_owns_dispatch is never called at all and the row is pushed indistinguishably from the flag being unset, while PULL_MIGRATION_TESTING.md section 9 M1 told the operator to go check an env var that is present and correct. Issue decision was Option 1 (make the reach honest, do not extend it). Corrects the issue's own table: the reachable set is {agent, event}, not {agent} - #1578's task-completion loopback POSTs to the same /task route and is tagged 'event' via an internal-secret-gated X-Event-Trigger, with no capacity bypass. Tests pin: the reach set derived from what _derive_task_trigger can emit rather than hardcoded; every autonomous trigger classified as reachable or stranded so a new one cannot join the stranded set unreviewed; the narrowing is a runtime NO-OP over the only producer that consults the predicate (it stops a false claim, it does not change a verdict); and note_unreachable_pull_trigger reports a stranded row once per (agent, trigger), never for a non-pilot / reachable / interactive trigger, and never raises. A structural guard asserts the three producers' overflow_policy literals, so implementing the deferred Option 2 (wiring cron to queue_persistent) fails this test and forces PULL_REACHABLE_TRIGGERS and section 9's prose to be revisited in the same change - a silently-stale reach set being exactly the bug. Also updates test_1766: its 'pilot owns EVERY autonomous trigger' case asserted reach the system never had (it passed only by calling the predicate outside the context that constrains it), and its producer-gate case drove a queue_persistent acquire carrying 'schedule' - a combination production cannot produce - now driven with 'agent', the shape that actually occurs. SUPERSEDED IN PART BY #2391: the structural guard fired as designed when Option 2 landed and was re-derived, not relaxed - task_execution_service now carries BOTH 'reject' and 'queue_persistent', and a second assertion pins that the wider one is reachable only through pull_owns_dispatch, which is the property #2048 could not assert and the one that keeps the blast radius equal to the pilot allowlist. The reachable set is now {agent, event, schedule, webhook, reminder} and the stranded set {loop, fan_out, a2a, operator_response}; see unit/test_2391_scheduled_pull_reach.py." }, { "file": "unit/test_2052_scrubber_authority_parity.py", @@ -2633,6 +2633,19 @@ "added": "2026-09-01", "categories": ["backend", "agent-server", "prompt", "regression"], "description": "#2468 headless tool audit: PLATFORM_DENIED_TOOLS family (11 tools measured on claude 2.1.235) + PLATFORM_KEPT_TOOLS data + AUDIT_CLI_VERSION; every denied name has a recorded reason with verbatim description fragments; DENIED \u222a KEPT exactly covers the measured init list; CORE_TOOLS never-denied guard; guardrails merge/dedupe (AC 5); reverse vocabulary guard (prompt family paragraph \u2286 deny tuple); 'Nothing survives the end of your turn' counterweight registered always-tier, survives every tier/runtime, quotes the Bash promise, stays truthful about waited subagents and names the foreground wait idiom." + }, + { + "file": "unit/test_2391_scheduled_pull_reach.py", + "feature": "#2391", + "added": "2026-09-03", + "categories": [ + "backend", + "unit", + "reliability", + "pull-migration", + "scheduler" + ], + "description": "Scheduled work reaches the durable pull queue (#2391 - the Option 2 deferred by #2048). task_execution_service - the producer behind the scheduler, i.e. ALL cron plus webhooks, reminders, loops and fan-out - dispatched with overflow_policy='reject' unconditionally, and capacity_manager only offers a row to the durable queue when its producer passed 'queue_persistent', so PULL_MODE_PILOT_AGENTS was structurally inert for the fleet's dominant traffic class. It now passes 'queue_persistent' when - and only when - pull_owns_dispatch says a pilot owns the trigger, which widens PULL_REACHABLE_TRIGGERS from {agent, event} to {agent, event, schedule, webhook, reminder}. Tests pin the whole risk surface: with the flag OFF (the fleet default) the real acquire call still receives 'reject' with no payload for every trigger and a CapacityFull still writes the same 'Agent at capacity (N/N parallel tasks running)' FAILED terminal - verified against the call, not assumed; with the flag ON a schedule/webhook/reminder row returns TaskExecutionStatus.QUEUED without calling the agent, releasing a slot or writing the dispatch marker, and the payload carries the row's model/allowed_tools/timeout/resume_session_id plus the CALLER's uncomposed system_prompt (the claim path rebuilds the platform prompt around it, #1629); a full backlog names max_backlog_depth rather than parallel slots, because the pull branch never asks for a slot; loop/fan_out/a2a/operator_response stay pushed (each has a caller that reads the returned TaskExecutionResult - the first three structurally, operator_response by choice since it records result.status as the ent#329 dispatch receipt); interactive triggers are untouched (Open Question 7's scope cut); the payload builder refuses without an execution_id (the enqueue is a CAS RUNNING->QUEUED on an existing row) and when slot_already_held (queueing under a held slot leaks it for the lease TTL); #1083 fire-and-forget does not stack - pull wins by construction because a queued row is never dispatched, so no 202 can arrive, while a non-pilot's async_result dispatch is unchanged. An end-to-end class drives the REAL CapacityManager + BacklogService + _build_claim_response so 'it enqueued' is not mistaken for 'a worker can run it' (the #2317 failure mode). Scheduler side: 'queued' is NOT a terminal - _poll_execution_completion treated anything but 'running' as an outcome, so a queued row would have published schedule_execution_completed(status=queued), been classified a failure and handed to _maybe_schedule_retry, duplicating work that was queued and about to run; ExecutionStatus.QUEUED and _NON_TERMINAL_POLL_STATES fix that and are pinned here." } ] } diff --git a/tests/unit/test_1766_pull_pilot_exclusive.py b/tests/unit/test_1766_pull_pilot_exclusive.py index cef9018db..5142c90f8 100644 --- a/tests/unit/test_1766_pull_pilot_exclusive.py +++ b/tests/unit/test_1766_pull_pilot_exclusive.py @@ -179,29 +179,31 @@ def test_non_pilot_never_owns_dispatch(self, pilot): assert pull_owns_dispatch("bob", "schedule") is False - @pytest.mark.parametrize("trigger", ["agent", "event"]) + @pytest.mark.parametrize( + "trigger", ["agent", "event", "schedule", "webhook", "reminder"] + ) def test_pilot_owns_the_autonomous_triggers_dispatch_can_deliver(self, pilot, trigger): - """Narrowed from "every autonomous trigger" by #2048. - - This case used to parametrize all seven of ``_AUTONOMOUS_TRIGGERS`` and - assert True for each — encoding reach the system never had. Only - ``POST /task`` dispatches with ``overflow_policy="queue_persistent"``, - which is the sole path on which ``capacity_manager`` consults this - predicate at all, and that route can only emit ``agent`` / ``event`` from - the autonomous set. The old assertion passed only because it called the - predicate directly, outside the context that constrains it. See - ``test_2048_pull_pilot_reach.py``. + """Narrowed from "every autonomous trigger" by #2048, re-widened by #2391. + + This case originally parametrized all seven of ``_AUTONOMOUS_TRIGGERS`` + and asserted True for each — encoding reach the system did not have; it + passed only because it called the predicate directly, outside the context + that constrains it. #2048 cut it to what ``POST /task`` can emit. #2391 + then gave ``task_execution_service`` a pilot-gated ``queue_persistent`` + policy, so the scheduler's async-polled triggers genuinely reach the + queue now and belong here. See ``test_2048_pull_pilot_reach.py``. """ from services.agent_service.pull_mode import pull_owns_dispatch assert pull_owns_dispatch("alice", trigger) is True @pytest.mark.parametrize( - "trigger", ["schedule", "webhook", "loop", "fan_out", "reminder"] + "trigger", ["loop", "fan_out", "a2a", "operator_response"] ) def test_pilot_does_not_own_a_trigger_dispatch_cannot_deliver(self, pilot, trigger): - """The #2048 correction as a positive assertion: declaring a trigger - autonomous gives the durable queue no way to receive it.""" + """The #2048 correction as a positive assertion, on the four triggers + #2391 left stranded: each one's caller reads the ``TaskExecutionResult`` + synchronously, so a queued row returns nothing for it to consume.""" from services.agent_service.pull_mode import pull_owns_dispatch assert pull_owns_dispatch("alice", trigger) is False @@ -281,6 +283,8 @@ def test_pilot_interactive_work_still_pushes( def test_non_pilot_autonomous_work_unchanged( self, capacity, slot_service, backlog_service, pilot ): + """`bob` is not in the allowlist, so `schedule` admits normally even + though #2391 made that trigger reachable for a pilot.""" result = asyncio.run( capacity.acquire( agent_name="bob", diff --git a/tests/unit/test_2048_pull_pilot_reach.py b/tests/unit/test_2048_pull_pilot_reach.py index 4d3d85790..5d187aa2c 100644 --- a/tests/unit/test_2048_pull_pilot_reach.py +++ b/tests/unit/test_2048_pull_pilot_reach.py @@ -26,6 +26,18 @@ (`note_unreachable_pull_trigger`), which is the acceptance criterion the original code failed in both directions. +**#2391 fired this file's tripwire, deliberately.** Option 2 landed: the +`reject` producer now dispatches with `overflow_policy="queue_persistent"` when +— and only when — `pull_owns_dispatch` says a pilot owns the trigger, so +`schedule`, `webhook` and `reminder` joined the reachable set. The structural +test below was re-derived rather than relaxed: it now pins BOTH policies on that +producer *and* pins that the wider one is gated, which is the property that was +never asserted before and is the one that actually matters. Everything else in +this file still holds — a stranded trigger is still stranded (`loop`, `fan_out`, +`a2a`, `operator_response`, all four blocked by a synchronous result consumer +rather than by the producer's policy), and the diagnostic still tells the +operator so. + Pure unit test — no Redis, no DB, no agent. Path bootstrap and lazy imports follow `test_1766_pull_pilot_exclusive.py`. """ @@ -49,22 +61,30 @@ # The triggers that are declared autonomous but cannot reach the queue. # -# `a2a` (ent#157) is stranded, and the classification is deliberate rather than -# a default. `PULL_REACHABLE_TRIGGERS` is an explicit allow-list, so an -# unlisted trigger lands here automatically — this comment is the review the -# test below demands. A2A's `message/send` consumes `result.response` -# synchronously to build the JSON-RPC artifact it returns to the remote caller; -# a pull-claimed row is dispatched by the agent later and produces no -# synchronous response, so pull dispatch structurally cannot serve this -# trigger. Same reason `fan_out` and `loop` are here. -# ent#329: `operator_response` joins the stranded set. It is dispatched by a -# direct backend call from the respond endpoint, not by `POST /task`, so -# `_derive_task_trigger` can never emit it and the pilot flag is inert for it — -# the same shape as `schedule` and `reminder`. -_STRANDED = ["schedule", "webhook", "loop", "fan_out", "reminder", "a2a", - "operator_response"] -# The two that can. -_REACHABLE = ["agent", "event"] +# After #2391 none of these is stranded on the producer's overflow policy any +# more — that producer can queue. They are stranded on their CALLER, which reads +# the returned `TaskExecutionResult`; a row queued for a worker to claim later +# gives it nothing to read. Three structural, one a scope choice. +# +# `loop` — `loop_service` renders the next iteration from `result.response`. +# `fan_out` — `fan_out_service` builds each `FanOutTaskResult` from the result; +# the async fan-out join is #1081 Phase 4. +# `a2a` — `routers/a2a`'s `message/send` consumes `result.response` to +# build the JSON-RPC artifact it hands back to the remote caller +# (ent#157). +# `operator_response` — out by CHOICE, not structurally: it dispatches through +# the same producer #2391 widened, but the respond endpoint records +# `result.status` as the dispatch receipt (audit row + #525 +# idempotency completion) and ent#329 exists because this spends +# money on a person's answer. "queued" is not the outcome that +# contract reports. A deliberate follow-up, not an oversight. +# +# `PULL_REACHABLE_TRIGGERS` is an explicit allow-list, so an unlisted trigger +# lands here automatically — this comment is the review the test below demands. +_STRANDED = ["loop", "fan_out", "a2a", "operator_response"] +# The five that can. `agent` + `event` arrive via `POST /task`; `schedule`, +# `webhook` and `reminder` via the scheduler's async-poll dispatch (#2391). +_REACHABLE = ["agent", "event", "schedule", "webhook", "reminder"] @pytest.fixture @@ -83,21 +103,47 @@ def pilot(monkeypatch): # --------------------------------------------------------------------------- -def test_reachable_set_is_the_autonomous_triggers_post_task_can_actually_emit(): +def test_reachable_set_is_what_the_two_queueing_producers_can_actually_emit(): """The claim behind `PULL_REACHABLE_TRIGGERS`, stated as an assertion. - `POST /task` is the only `queue_persistent` producer, and its - `_derive_task_trigger` can only ever produce these five values. Intersect - with the autonomous set and exactly `agent` + `event` survive. Derived here - rather than hardcoded so the constant cannot quietly disagree with the - reasoning that justifies it. + Two producers can pass `overflow_policy="queue_persistent"` since #2391, and + the reachable set is exactly what they contribute from the autonomous set: + + * `POST /task` (`chat_execution_service`) — `_derive_task_trigger` can only + ever produce `{self_task, agent, mcp, manual, event}`, which contributes + `agent` + `event`. + * `task_execution_service` — contributes the autonomous triggers with no + synchronous result consumer, i.e. everything the scheduler dispatches + async-and-polls: `schedule`, `webhook`, `reminder`. + + Derived here rather than hardcoded so the constant cannot quietly disagree + with the reasoning that justifies it. """ from services.pull_pilot import PULL_REACHABLE_TRIGGERS from services.task_execution_service import _AUTONOMOUS_TRIGGERS task_route_can_emit = {"self_task", "agent", "mcp", "manual", "event"} - assert PULL_REACHABLE_TRIGGERS == (task_route_can_emit & _AUTONOMOUS_TRIGGERS) - assert PULL_REACHABLE_TRIGGERS == {"agent", "event"} + scheduler_async_polled = {"schedule", "webhook", "reminder"} + assert PULL_REACHABLE_TRIGGERS == ( + (task_route_can_emit | scheduler_async_polled) & _AUTONOMOUS_TRIGGERS + ) + assert PULL_REACHABLE_TRIGGERS == { + "agent", "event", "schedule", "webhook", "reminder" + } + + +def test_every_reachable_trigger_from_this_producer_is_also_1083_shaped(): + """Why `{schedule, webhook, reminder}` and not "the rest of the autonomous + set": all three reach `execute_task` from the scheduler, which dispatches + with `async_mode=True` and then polls the DB — nobody holds a coroutine on + the return value. That is the same property #1083 selects on for + fire-and-forget, so its eligible set must be a SUBSET of what this producer + can queue. If someone widens `ASYNC_DISPATCH_ELIGIBLE_TRIGGERS` to a trigger + with a synchronous consumer, this fails and names the contradiction.""" + from config import ASYNC_DISPATCH_ELIGIBLE_TRIGGERS + from services.pull_pilot import PULL_REACHABLE_TRIGGERS + + assert ASYNC_DISPATCH_ELIGIBLE_TRIGGERS <= PULL_REACHABLE_TRIGGERS def test_the_stranded_triggers_are_named_and_complete(): @@ -160,30 +206,35 @@ def test_a_stranded_trigger_on_a_pilot_is_reported(pilot, trigger, caplog): def test_the_message_contradicts_the_advice_that_used_to_misfire(pilot, caplog): """The operator-facing point of the line. §9 M1 said a pushed autonomous row means "check the backend actually has PULL_MODE_PILOT_AGENTS in its env" — - for a `schedule` row that sends them after a variable that is present and - correct. The log has to say so in words, not merely fire.""" + for a stranded row that sends them after a variable that is present and + correct. The log has to say so in words, not merely fire. + + Driven with `loop` since #2391: `schedule` is reachable now, and using it + here would assert the diagnostic over a case that no longer exists. + """ with caplog.at_level("WARNING"): - pilot.note_unreachable_pull_trigger("pilot-a", "schedule") + pilot.note_unreachable_pull_trigger("pilot-a", "loop") assert "flag is applied and correct" in caplog.text assert "topology" in caplog.text def test_it_reports_once_per_agent_and_trigger(pilot, caplog): - """This sits on the cron dispatch path, which fires on every scheduled run. - An un-deduped warning would emit thousands of identical lines a day and - train operators to filter it — a signal meant to be noticed becoming noise.""" + """This sits on the dispatch path of every loop iteration and fan-out + subtask. An un-deduped warning would emit thousands of identical lines a day + and train operators to filter it — a signal meant to be noticed becoming + noise.""" with caplog.at_level("WARNING"): - first = pilot.note_unreachable_pull_trigger("pilot-a", "schedule") - repeats = [pilot.note_unreachable_pull_trigger("pilot-a", "schedule") for _ in range(50)] + first = pilot.note_unreachable_pull_trigger("pilot-a", "loop") + repeats = [pilot.note_unreachable_pull_trigger("pilot-a", "loop") for _ in range(50)] assert first is True assert not any(repeats) assert caplog.text.count("#2048") == 1 def test_dedup_is_per_trigger_not_per_agent(pilot): - """A pilot firing cron AND webhooks has two distinct gaps to report.""" - assert pilot.note_unreachable_pull_trigger("pilot-a", "schedule") is True - assert pilot.note_unreachable_pull_trigger("pilot-a", "webhook") is True + """A pilot running loops AND fan-out has two distinct gaps to report.""" + assert pilot.note_unreachable_pull_trigger("pilot-a", "loop") is True + assert pilot.note_unreachable_pull_trigger("pilot-a", "fan_out") is True @pytest.mark.parametrize("trigger", _REACHABLE) @@ -247,21 +298,51 @@ def _overflow_policies(relative_path: str) -> set: def test_the_producer_topology_this_fix_rests_on_still_holds(): """The load-bearing fact, pinned as a test rather than a comment. - If someone later implements the issue's Option 2 — giving - `task_execution_service` a `queue_persistent` policy so cron can be claimed - — this fails, which is the intent: `PULL_REACHABLE_TRIGGERS` and §9's prose - both have to be revisited in that same change, and neither would otherwise - announce itself. A silently-stale reach set is exactly the bug #2048 is. + #2048 pinned `task_execution_service == {"reject"}` as a tripwire for its own + deferred Option 2. #2391 implemented Option 2, so the tripwire fired and is + re-derived here rather than deleted: the producer now carries BOTH literals, + which is the shape a *conditional* widening has and an unconditional one does + not. Pair it with `test_the_wider_policy_is_gated_on_the_pull_predicate` + below — that is the assertion #2048 could not make and the one that keeps the + blast radius equal to the pilot allowlist. """ - assert _overflow_policies("services/task_execution_service.py") == {"reject"}, ( - "task_execution_service no longer dispatches with overflow_policy='reject' — " - "cron may now be able to reach the durable queue. Re-derive " - "PULL_REACHABLE_TRIGGERS and update PULL_MIGRATION_TESTING.md §9 (#2048)." + assert _overflow_policies("services/task_execution_service.py") == { + "reject", "queue_persistent", + }, ( + "task_execution_service's overflow policies changed. Since #2391 it must " + "carry exactly two: 'reject' (the default, unchanged for every non-pilot " + "agent) and 'queue_persistent' (pilot-gated). Losing 'reject' would mean " + "scheduled work is queued for the whole fleet — the reliability-spine " + "change #2048 declined to make. Re-derive PULL_REACHABLE_TRIGGERS and " + "update PULL_MIGRATION_TESTING.md §9." ) assert "queue_persistent" in _overflow_policies("services/chat_execution_service.py") assert "queue_in_memory" in _overflow_policies("services/dispatch_admission_service.py") +def test_the_wider_policy_is_gated_on_the_pull_predicate(): + """The #2391 safety property, pinned structurally. + + `queue_persistent` on this producer is reachable ONLY through a non-None + `pull_overflow_payload`, and the only thing that builds one is + `build_pull_queue_payload`, whose sole widening condition is + `pull_owns_dispatch` (false for every agent outside the allowlist). Assert + the chain so nobody can later make the policy unconditional — which would + change capacity-pressure semantics for the entire fleet, flag or no flag — + without this failing. + """ + source = (_BACKEND / "services" / "task_execution_service.py").read_text(encoding="utf-8") + + # The policy choice is a branch on the payload, not a constant. + assert 'if pull_overflow_payload is not None:\n overflow_policy = "queue_persistent"' in source + assert 'overflow_policy = "reject"' in source + # The payload builder refuses unless the pull predicate says so. + builder = source[source.index("def build_pull_queue_payload("):] + builder = builder[: builder.index("\ndef ")] + assert "pull_owns_dispatch(agent_name, triggered_by)" in builder + assert "slot_already_held or not execution_id" in builder + + def test_the_reject_producer_reports_the_gap(): """The wiring, checked structurally: the diagnostic has to sit at the producer that actually knows `triggered_by`. `capacity_manager` cannot do it diff --git a/tests/unit/test_2391_scheduled_pull_reach.py b/tests/unit/test_2391_scheduled_pull_reach.py new file mode 100644 index 000000000..de94f13ef --- /dev/null +++ b/tests/unit/test_2391_scheduled_pull_reach.py @@ -0,0 +1,614 @@ +"""#2391 — scheduled work can reach the durable pull queue. + +Until this change `task_execution_service` — the producer behind the scheduler, +i.e. **all cron**, plus webhooks, reminders, loops and fan-out — dispatched with +`overflow_policy="reject"` unconditionally. `capacity_manager.acquire` only +offers a row to the durable queue when its producer passed +`"queue_persistent"`, so `PULL_MODE_PILOT_AGENTS` was structurally inert for the +fleet's dominant traffic class no matter how it was set. #2048 made that +legible (`PULL_REACHABLE_TRIGGERS`, `note_unreachable_pull_trigger`) and +deliberately deferred fixing it; this is the deferred half. + +The whole risk of the change is in *how* the policy widens, so that is what this +file pins: + + * **Gated, never unconditional.** The policy flips to `queue_persistent` only + when `pull_owns_dispatch` says a pilot owns the trigger. That keeps the + blast radius exactly equal to the pilot allowlist — the reason #2048 + declined to bundle Option 2 was that an *unconditional* persistent queue on + this producer would change what happens to scheduled work under capacity + pressure for every agent, flag or no flag. + * **Flag OFF is byte-for-byte.** Verified against the real acquire call, not + assumed: same policy, no payload, same `CapacityFull` → FAILED terminal. + * **Stranded stays stranded.** `loop` / `fan_out` / `a2a` / + `operator_response` have a caller that reads `TaskExecutionResult` + synchronously, so queueing them would silently return nothing to read. + * **#1083 does not stack.** A pull-queued row is never dispatched, so no 202 + ACK can arrive and `async_result` is never sent — one mechanism per turn. + * **The queued row is actually claimable.** The real producer's metadata is + fed to the real claim-response builder, so "it enqueued" is not mistaken for + "a worker can run it". + +Plus the scheduler side: `queued` is NOT a terminal, and the async-poll loop +must poll through it or it publishes a bogus failure and schedules a retry for +work that is queued and about to run. + +Pure unit test — mocked transport / capacity collaborators for the producer +assertions, the REAL `CapacityManager` + `BacklogService` + claim builder for +the end-to-end, and the scheduler service driven via `object.__new__` (the +`test_1823_scheduler_permanent_classification.py` precedent). +""" +from __future__ import annotations + +import asyncio +import json +import os +import sys +from pathlib import Path +from unittest.mock import AsyncMock, MagicMock, patch + +import pytest + +_REPO = Path(__file__).resolve().parents[2] +_BACKEND = _REPO / "src" / "backend" +if str(_BACKEND) not in sys.path: + sys.path.insert(0, str(_BACKEND)) +# Appended, never inserted at 0, so the repo root cannot shadow the +# conftest-managed `src/backend` entries (mirrors test_1823 / test_1994). +if str(_REPO) not in sys.path: + sys.path.append(str(_REPO)) + +# `src/scheduler/config.py` reads these at import time (#589). +os.environ.setdefault("REDIS_URL", "redis://test:test@redis:6379") +os.environ.setdefault("REDIS_PASSWORD", "test") +os.environ.setdefault("REDIS_BACKEND_PASSWORD", "test") + +pytestmark = pytest.mark.unit + + +PILOT = "pilot-a" +NON_PILOT = "bob" + + +def _await(coro): + loop = asyncio.new_event_loop() + try: + return loop.run_until_complete(coro) + finally: + loop.close() + + +def _resp(status_code, body=None): + r = MagicMock() + r.status_code = status_code + r.raise_for_status = MagicMock() + r.json.return_value = body if body is not None else {} + return r + + +_SUCCESS_BODY = { + "response": "done", + "session_id": "s1", + "metadata": {"cost_usd": 0.01, "context_window": 200000}, + "execution_log": [], +} + + +# --------------------------------------------------------------------------- +# Producer harness — real execute_task, mocked capacity + transport +# --------------------------------------------------------------------------- + + +def _run( + *, + agent_name=PILOT, + triggered_by="schedule", + execution_id="exec-2391", + acquire_result=None, + acquire_raises=None, + dispatch_async=False, + **execute_kwargs, +): + """Drive the real `execute_task` and return `(result, mocks)`. + + `capacity` is a mock so the test can read back the exact `overflow_policy` + and `overflow_payload` this producer chose — which is the decision under + test. The end-to-end class below uses the real one instead. + """ + import config + from services.task_execution_service import TaskExecutionService + + mock_db = MagicMock() + mock_db.get_max_parallel_tasks.return_value = 3 + mock_db.get_execution_timeout.return_value = 300 + mock_db.get_execution.return_value = MagicMock(status="cancelled") + mock_db.update_execution_status.return_value = True + + mock_capacity = MagicMock() + if acquire_raises is not None: + mock_capacity.acquire = AsyncMock(side_effect=acquire_raises) + else: + mock_capacity.acquire = AsyncMock( + return_value=acquire_result or MagicMock(state="admitted") + ) + mock_capacity.release = AsyncMock() + + mock_circuit = MagicMock() + mock_circuit.allow_request.return_value = True + mock_activity = MagicMock( + track_activity=AsyncMock(return_value="act-1"), + complete_activity=AsyncMock(), + ) + post_mock = AsyncMock(return_value=_resp(200, _SUCCESS_BODY)) + + with ( + patch.object(config, "DISPATCH_ASYNC", dispatch_async), + patch("services.task_execution_service.db", mock_db), + patch("services.task_execution_service.get_capacity_manager", return_value=mock_capacity), + patch("services.task_execution_service.activity_service", mock_activity), + patch("services.task_execution_service.CircuitState", return_value=mock_circuit), + patch("services.task_execution_service.agent_post_with_retry", post_mock), + patch("services.task_execution_service.dispatch_breaker_active", return_value=False), + patch("services.task_execution_service._record_dispatch_terminal", AsyncMock()), + ): + svc = TaskExecutionService() + result = _await( + svc.execute_task( + agent_name=agent_name, + message="do the thing", + triggered_by=triggered_by, + execution_id=execution_id, + timeout_seconds=300, + model="sonnet", + **execute_kwargs, + ) + ) + return result, {"db": mock_db, "capacity": mock_capacity, "post": post_mock} + + +def _acquire_kwargs(capacity_mock): + return capacity_mock.acquire.await_args.kwargs + + +@pytest.fixture +def pilot(monkeypatch): + monkeypatch.setenv("PULL_MODE_PILOT_AGENTS", PILOT) + + +@pytest.fixture +def no_pilots(monkeypatch): + """The default the whole fleet runs on.""" + monkeypatch.setenv("PULL_MODE_PILOT_AGENTS", "") + + +# --------------------------------------------------------------------------- +# 1. Flag OFF — scheduled dispatch is unchanged (AC: verified, not assumed) +# --------------------------------------------------------------------------- + + +class TestFlagOffIsUnchanged: + @pytest.mark.parametrize( + "trigger", ["schedule", "webhook", "reminder", "loop", "fan_out"] + ) + def test_policy_stays_reject_with_no_pilots(self, no_pilots, trigger): + """The property the entire risk assessment rests on: with an empty + allowlist this producer still asks for `reject` and offers no payload, + so `capacity_manager` cannot queue even if it wanted to.""" + _, m = _run(triggered_by=trigger) + kw = _acquire_kwargs(m["capacity"]) + assert kw["overflow_policy"] == "reject" + assert kw["overflow_payload"] is None + + def test_a_different_agent_being_a_pilot_changes_nothing(self, pilot): + """Blast radius is the allowlist, not the trigger.""" + _, m = _run(agent_name=NON_PILOT, triggered_by="schedule") + kw = _acquire_kwargs(m["capacity"]) + assert kw["overflow_policy"] == "reject" + assert kw["overflow_payload"] is None + + def test_at_capacity_still_fails_fast_with_the_same_terminal(self, no_pilots): + """The behaviour #2391 was warned not to change: a scheduled fire that + arrives at capacity is REJECTED, and the row is written FAILED with the + parallel-tasks wording.""" + from services.capacity_manager import CapacityFull + from services.task_execution_service import TaskExecutionStatus + + result, m = _run( + triggered_by="schedule", + acquire_raises=CapacityFull(NON_PILOT, 3, "rejected"), + ) + assert result.status == TaskExecutionStatus.FAILED + assert "at capacity (3/3 parallel tasks running)" in result.error + m["post"].assert_not_awaited() + m["db"].update_execution_status.assert_called_once() + + def test_the_agent_is_still_dispatched_normally(self, no_pilots): + from services.task_execution_service import TaskExecutionStatus + + result, m = _run(triggered_by="schedule") + assert result.status == TaskExecutionStatus.SUCCESS + m["post"].assert_awaited_once() + + +# --------------------------------------------------------------------------- +# 2. Flag ON — the scheduler's triggers reach the queue +# --------------------------------------------------------------------------- + + +class TestFlagOnQueuesScheduledWork: + @pytest.mark.parametrize("trigger", ["schedule", "webhook", "reminder"]) + def test_policy_widens_and_a_payload_is_offered(self, pilot, trigger): + _, m = _run(triggered_by=trigger) + kw = _acquire_kwargs(m["capacity"]) + assert kw["overflow_policy"] == "queue_persistent" + assert kw["overflow_payload"] is not None + assert kw["overflow_payload"].triggered_by == trigger + + def test_a_queued_row_returns_QUEUED_and_never_calls_the_agent(self, pilot): + """The handoff. No agent call, no slot to release — the worker owns the + turn from here and reports its terminal through the claim-token CAS.""" + from services.task_execution_service import TaskExecutionStatus + + result, m = _run( + triggered_by="schedule", + acquire_result=MagicMock(state="queued_persistent"), + ) + assert result.status == TaskExecutionStatus.QUEUED + assert result.execution_id == "exec-2391" + m["post"].assert_not_awaited() + m["capacity"].release.assert_not_awaited() + m["db"].mark_execution_dispatched.assert_not_called() + + def test_the_payload_carries_the_per_task_settings(self, pilot): + """`backlog_service.enqueue` reads these off the request; a pulled turn + that lost them would silently run on agent/global defaults (#2317).""" + _, m = _run( + triggered_by="schedule", + allowed_tools=["Bash", "Read"], + system_prompt="be brief", + resume_session_id="sess-9", + ) + req = _acquire_kwargs(m["capacity"])["overflow_payload"].request + assert req.message == "do the thing" + assert req.model == "sonnet" + assert req.allowed_tools == ["Bash", "Read"] + assert req.timeout_seconds == 300 + assert req.resume_session_id == "sess-9" + # The CALLER's prompt, uncomposed — the claim path rebuilds the platform + # prompt around it (#1629), so composing here would double it. + assert req.system_prompt == "be brief" + + def test_a_full_backlog_names_the_limit_that_actually_bound(self, pilot): + """On the pull path `CapacityFull` means max_backlog_depth, not busy + parallel slots — the pull branch never asks for a slot. Reporting + "3/3 parallel tasks running" would send the operator to the wrong knob.""" + from services.capacity_manager import CapacityFull + from services.task_execution_service import TaskExecutionStatus + + result, _ = _run( + triggered_by="schedule", + acquire_raises=CapacityFull(PILOT, 3, "persistent_full"), + ) + assert result.status == TaskExecutionStatus.FAILED + assert "backlog full" in result.error + assert "parallel tasks" not in result.error + + +# --------------------------------------------------------------------------- +# 3. The triggers that stay stranded, and why +# --------------------------------------------------------------------------- + + +class TestStrandedTriggersStayPushed: + @pytest.mark.parametrize( + "trigger", ["loop", "fan_out", "a2a", "operator_response"] + ) + def test_a_result_reading_caller_keeps_its_trigger_on_push( + self, pilot, trigger + ): + """Each of these callers reads the returned `TaskExecutionResult`, and a + queued row gives them nothing to read. + + `loop_service` renders the next iteration's template from it, + `fan_out_service` builds each `FanOutTaskResult` from it, and + `routers/a2a` turns it into the JSON-RPC artifact it hands a remote + caller — three structural blocks; widening those is #1081 Phase 4's + async join, not this change. `operator_response` is the fourth and is + NOT structural: it dispatches through this same producer, but the + respond endpoint records `result.status` as the dispatch receipt (audit + row + #525 idempotency completion) for a turn that spends money on a + person's answer (ent#329), and "queued" is not the outcome that contract + reports. Out by choice, and this test is what makes the choice explicit + rather than incidental.""" + _, m = _run(triggered_by=trigger) + kw = _acquire_kwargs(m["capacity"]) + assert kw["overflow_policy"] == "reject" + assert kw["overflow_payload"] is None + + @pytest.mark.parametrize("trigger", ["manual", "mcp", "chat", "public", "voice"]) + def test_interactive_triggers_are_untouched(self, pilot, trigger): + """Open Question 7's scope cut: a human turn keeps the synchronous push + path and today's Redis session lock.""" + _, m = _run(triggered_by=trigger) + assert _acquire_kwargs(m["capacity"])["overflow_policy"] == "reject" + + +# --------------------------------------------------------------------------- +# 4. Preconditions on the payload builder +# --------------------------------------------------------------------------- + + +class TestPayloadPreconditions: + def _build(self, **over): + from services.task_execution_service import build_pull_queue_payload + + kwargs = dict( + agent_name=PILOT, + triggered_by="schedule", + execution_id="exec-1", + message="m", + model="sonnet", + allowed_tools=None, + system_prompt=None, + timeout_seconds=300, + resume_session_id=None, + subscription_id=None, + source_user_id=None, + source_user_email=None, + source_agent_name=None, + slot_already_held=False, + ) + kwargs.update(over) + return build_pull_queue_payload(**kwargs) + + def test_refuses_without_an_execution_id(self, pilot): + """`BacklogService.enqueue` transitions an EXISTING row RUNNING→QUEUED + under a CAS; with no row there is nothing to transition.""" + assert self._build(execution_id=None) is None + + def test_refuses_when_a_slot_is_already_held(self, pilot): + """The caller (a drain, or /task's pre-flight) owns a real slot that + `execute_task`'s `finally` releases. Queueing under it would leak the + slot for the whole lease TTL.""" + assert self._build(slot_already_held=True) is None + + def test_builds_for_a_pilot(self, pilot): + assert self._build() is not None + + def test_refuses_for_a_non_pilot(self, pilot): + assert self._build(agent_name=NON_PILOT) is None + + +# --------------------------------------------------------------------------- +# 5. #1083 fire-and-forget does not stack on pull (AC #4) +# --------------------------------------------------------------------------- + + +class TestFireAndForgetInteraction: + def test_pull_wins_and_no_async_ack_is_ever_requested(self, pilot): + """`schedule` is eligible for BOTH mechanisms once `DISPATCH_ASYNC` is + on. They cannot both apply to one turn, and pull wins by construction + rather than by precedence: the row is queued before any dispatch, so + there is no request for the agent to answer 202 to and `async_result` is + never sent.""" + from services.task_execution_service import TaskExecutionStatus + + result, m = _run( + triggered_by="schedule", + dispatch_async=True, + acquire_result=MagicMock(state="queued_persistent"), + ) + assert result.status == TaskExecutionStatus.QUEUED + assert result.dispatched_async is False + m["post"].assert_not_awaited() + + def test_fire_and_forget_is_untouched_for_a_non_pilot(self, no_pilots): + """The other direction: #1083 keeps working exactly as it did on every + agent that is not a pull pilot.""" + from services.task_execution_service import TaskExecutionStatus + + result, m = _run( + triggered_by="schedule", + dispatch_async=True, + ) + m["post"].assert_awaited_once() + assert m["post"].await_args.args[2]["async_result"] is True + assert result.status == TaskExecutionStatus.SUCCESS + + +# --------------------------------------------------------------------------- +# 6. End-to-end: a scheduled row really lands in the queue, claimable +# --------------------------------------------------------------------------- + + +class _FakeQueueDb: + """The slice of the `database.db` singleton `BacklogService.enqueue` uses.""" + + def __init__(self): + self.queued: dict[str, str] = {} + + def get_queued_count(self, agent_name): + return 0 + + def get_max_backlog_depth(self, agent_name): + return 50 + + def update_execution_to_queued(self, execution_id, metadata, queued_at): + self.queued[execution_id] = metadata + return True + + +class TestEnqueuedScheduledRowIsClaimable: + """Drives the REAL `CapacityManager` + `BacklogService` so the assertion is + about what a worker would actually receive, not about a mock's call args. + + "It enqueued" and "a worker can run it" are different claims, and #2317 is + the precedent for them coming apart silently: the pull envelope read a + nested key no producer wrote, so every pulled turn ran on defaults while + every enqueue test stayed green. + """ + + def _enqueue_via_capacity(self, monkeypatch, trigger="schedule"): + from services import capacity_manager as cm_module + from services.backlog_service import BacklogService + from services.task_execution_service import build_pull_queue_payload + + fake_db = _FakeQueueDb() + monkeypatch.setattr(cm_module.redis, "from_url", lambda *_a, **_kw: MagicMock()) + + slots = AsyncMock() + slots.slots_prefix = "agent:slots:" + slots.acquire_slot = AsyncMock(return_value=True) # a slot IS free + slots.register_on_release = lambda cb: None + + backlog = BacklogService() + capacity = cm_module.CapacityManager( + redis_url="redis://test", slot_service=slots, backlog_service=backlog + ) + + payload = build_pull_queue_payload( + agent_name=PILOT, + triggered_by=trigger, + execution_id="exec-e2e", + message="run the nightly report", + model="opus", + allowed_tools=["Bash"], + system_prompt="be brief", + timeout_seconds=1800, + resume_session_id="sess-e2e", + subscription_id="sub-1", + source_user_id=None, + source_user_email=None, + source_agent_name=None, + slot_already_held=False, + ) + assert payload is not None + + # `acquire` and `enqueue` both late-import their collaborators, so patch + # the modules those names resolve from. The ephemeral-budget read is not + # stubbed deliberately: it is fail-open on any error, and letting the + # real guard run proves this path does not depend on stubbing it out. + with ( + patch("database.db", fake_db), + patch("services.settings_service.clamp_to_ceiling", lambda v: v), + ): + result = _await( + capacity.acquire( + agent_name=PILOT, + execution_id="exec-e2e", + max_concurrent=3, + overflow_policy="queue_persistent", + overflow_payload=payload, + ) + ) + return result, fake_db, slots + + def test_the_row_is_queued_not_admitted(self, pilot, monkeypatch): + result, fake_db, slots = self._enqueue_via_capacity(monkeypatch) + assert result.state == "queued_persistent" + # No ZADD even though a slot was free: capacity is the worker pool now. + slots.acquire_slot.assert_not_awaited() + assert "exec-e2e" in fake_db.queued + + def test_the_claim_envelope_a_worker_receives_is_complete(self, pilot, monkeypatch): + """Feed the REAL producer's metadata to the REAL claim-response builder, + through a row shaped like `claim_next_queued`'s RETURNING columns.""" + from services import pull_coordination_service as pcs + + _, fake_db, _ = self._enqueue_via_capacity(monkeypatch) + meta = fake_db.queued["exec-e2e"] + + claimed_row = { + "id": "exec-e2e", + "agent_name": PILOT, + "message": "run the nightly report", + "backlog_metadata": meta, + "source_user_id": None, + "source_user_email": None, + "source_agent_name": None, + "source_mcp_key_id": None, + "source_mcp_key_name": None, + "subscription_id": "sub-1", + "triggered_by": "schedule", + "claude_session_id": None, + "model_used": "opus", + "started_at": "2026-01-01T00:00:00.000000Z", + "claim_token": "tok-1", + "lease_expires_at": "2026-01-01T00:30:00.000000Z", + "claimed_by_worker": f"{PILOT}#w1", + "redelivery_count": 0, + } + with patch.object(pcs, "_compose_pull_system_prompt", lambda *a, **kw: "composed"): + claim = pcs._build_claim_response(claimed_row) + + assert claim["execution_id"] == "exec-e2e" + assert claim["claim_token"] == "tok-1" + env = claim["envelope"] + assert env["to"] == PILOT + assert env["payload"]["message"] == "run the nightly report" + # Session identity comes from the key the producer writes (#2317). + assert env["payload"]["session_id"] == "sess-e2e" + overrides = env["payload"]["task_overrides"] + assert overrides["model"] == "opus" + assert overrides["allowed_tools"] == ["Bash"] + assert overrides["timeout_seconds"] == 1800 + + def test_the_persisted_metadata_records_the_scheduled_trigger(self, pilot, monkeypatch): + """Provenance survives the round-trip: the queued row still says it came + from cron, which is what `_compose_pull_system_prompt` derives the + execution-context mode from.""" + _, fake_db, _ = self._enqueue_via_capacity(monkeypatch) + assert json.loads(fake_db.queued["exec-e2e"])["triggered_by"] == "schedule" + + +# --------------------------------------------------------------------------- +# 7. Scheduler side — `queued` is not a terminal +# --------------------------------------------------------------------------- + + +class _PollDb: + def __init__(self, statuses): + self._statuses = list(statuses) + self.calls = 0 + + def get_execution(self, execution_id): + self.calls += 1 + status = self._statuses[min(self.calls - 1, len(self._statuses) - 1)] + return MagicMock( + id=execution_id, status=status, response="ok", error=None, + cost=0.0, context_used=None, context_max=None, + ) + + +def _poller(monkeypatch, statuses): + import src.scheduler.service as scheduler_service + + monkeypatch.setattr(scheduler_service.config, "poll_interval", 0) + monkeypatch.setattr(scheduler_service.config, "poll_deadline_buffer", 5) + svc = object.__new__(scheduler_service.SchedulerService) + svc.db = _PollDb(statuses) + return scheduler_service, svc + + +class TestSchedulerPollsThroughQueued: + def test_queued_does_not_end_the_poll(self, monkeypatch): + """Without this the scheduler reads `queued` as "completed", publishes + `schedule_execution_completed(status=queued)` — a failure, since only + `success` is not — and hands it to `_maybe_schedule_retry`, duplicating + work that is queued and about to run.""" + _, svc = _poller(monkeypatch, ["queued", "queued", "running", "success"]) + result = _await(svc._poll_execution_completion("exec-1", 30)) + assert result["status"] == "success" + assert svc.db.calls == 4 + + def test_a_real_terminal_still_ends_the_poll_immediately(self, monkeypatch): + _, svc = _poller(monkeypatch, ["failed"]) + result = _await(svc._poll_execution_completion("exec-1", 30)) + assert result["status"] == "failed" + assert svc.db.calls == 1 + + def test_queued_is_declared_non_terminal(self, monkeypatch): + module, _ = _poller(monkeypatch, ["success"]) + assert module.ExecutionStatus.QUEUED in module._NON_TERMINAL_POLL_STATES + assert module.ExecutionStatus.RUNNING in module._NON_TERMINAL_POLL_STATES + assert module.ExecutionStatus.SUCCESS not in module._NON_TERMINAL_POLL_STATES