Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
11 changes: 8 additions & 3 deletions docs/memory/architecture.md
Original file line number Diff line number Diff line change
Expand Up @@ -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`.

Expand Down
2 changes: 1 addition & 1 deletion docs/memory/feature-flows/capacity-management.md
Original file line number Diff line number Diff line change
Expand Up @@ -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) |
Expand Down
2 changes: 2 additions & 0 deletions docs/memory/feature-flows/task-execution-service.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
140 changes: 102 additions & 38 deletions docs/testing/PULL_MIGRATION_TESTING.md
Original file line number Diff line number Diff line change
Expand Up @@ -363,11 +363,20 @@ name for `'<agent>'`.

### 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
Expand Down Expand Up @@ -435,8 +444,9 @@ WHERE agent_name = '<agent>'
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:

Expand All @@ -450,65 +460,119 @@ WHERE agent_name = '<agent>'
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.
Expand Down
Loading
Loading