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
204 changes: 202 additions & 2 deletions docs/testing/PULL_MIGRATION_TESTING.md
Original file line number Diff line number Diff line change
Expand Up @@ -329,8 +329,208 @@ Ran on `gtm-synthesizer` (rebuilt image, `AGENT_AUTH_SECRET` fixed; chosen becau

## 8. Net remaining before a default-ON decision

Canary lease-awareness (E-05/E-01, §3 T3.6/T3.7) · Tier-6 `effect_guard` `execution_id` injection · ~~G3
canary-on-PG (#1540)~~ ✅ closed · B6 runtime-verify on the rebuilt image · the ≥2-week soak (#856).
~~Canary lease-awareness E-05 (§3 T3.6)~~ ✅ **closed by #1766** — E-05 now excludes leased rows, mirroring
S-01 and the `mark_no_session_executions_failed` sweep it watches (which already carried
`lease_expires_at IS NULL`); without it E-05 fired on every pull turn past 60s · canary lease-awareness
E-01 (§3 T3.7) · Tier-6 `effect_guard` `execution_id` injection · ~~G3 canary-on-PG (#1540)~~ ✅ closed ·
B6 runtime-verify on the rebuilt image · the ≥2-week soak (#856 / #1766, measurement set in §9).

**Also closed by #1766:** the pilot flag was purely additive, so a pilot ran push AND pull concurrently —
the producer never force-queued (a free slot still meant a push, so rows only queued on overflow) and the
backend's own `drain_next` raced the agent's worker for whatever did queue. Two independent capacity
counters meant up to 2x `max_parallel_tasks`, invisible to S-02. The flag is now a true either/or for
autonomous triggers; interactive turns keep the synchronous path (Open Question 7 scope cut).

## 9. Soak measurement set (#1766)

The queries that turn "we ran it for a few days" into a verdict. **Nothing else
currently covers this**: `instance-monitor` has zero references to leases,
`claim_token`, `redelivery_count`, or queue depth — its probes report generic
health and will stay green whether pull carries the fleet or does nothing at all.
Run these instead (or fold M1/M3/M5 into the monitor's deep-probe).

PostgreSQL flavour. All timestamp columns are `Text` holding ISO-Z strings
(Invariant #16), hence the explicit `::timestamptz` casts. Substitute the pilot
name for `'<agent>'`.

### Pre-flight — capture a baseline BEFORE the flip

Run M4 + M6 + M7 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
unattributable and the write-up degrades to "felt fine".

### The set

| ID | Answers | AC |
|----|---------|-----|
| M1 | Is pull actually carrying load? | coverage (the #1766 premise) |
| M2 | Push/pull split over time | coverage |
| M3 | Is the agent over its configured concurrency? | no slot/backlog drift |
| M4 | Lost or phantom executions | no lost/phantom |
| M5 | Re-delivery + poison-park behaviour | leases behave as designed |
| M6 | Outcome regression vs baseline | no regression |
| M7 | Canary state | canary green throughout |
| M8 | Queue starvation | no drift |

**M1 — is the pull path carrying anything?** The single most important query: if
`pulled` is ~0 the soak is measuring nothing and everything downstream is void.

```sql
SELECT COUNT(*) FILTER (WHERE claimed_by_worker IS NOT NULL) AS pulled,
COUNT(*) FILTER (WHERE claimed_by_worker IS NULL) AS pushed
FROM schedule_executions
WHERE agent_name = '<agent>'
AND started_at::timestamptz > now() - interval '24 hours';
```

Expected after the #1766 gate: autonomous triggers ~100% `pulled`; anything still
`pushed` should be interactive (`manual`/`user`/`chat`) — confirm with:

```sql
SELECT triggered_by,
COUNT(*) FILTER (WHERE claimed_by_worker IS NOT NULL) AS pulled,
COUNT(*) FILTER (WHERE claimed_by_worker IS NULL) AS pushed
FROM schedule_executions
WHERE agent_name = '<agent>'
AND started_at::timestamptz > now() - interval '24 hours'
GROUP BY 1 ORDER BY 1;
```

A pushed **autonomous** row means the producer gate did not engage — check the
backend actually has `PULL_MODE_PILOT_AGENTS` in its env.

**M2 — split over time.** Confirms the flip took effect at the moment you think
it did, and shows drift.

```sql
SELECT date_trunc('hour', started_at::timestamptz) AS hr,
COUNT(*) FILTER (WHERE claimed_by_worker IS NOT NULL) AS pulled,
COUNT(*) FILTER (WHERE claimed_by_worker IS NULL) AS pushed
FROM schedule_executions
WHERE agent_name = '<agent>'
AND started_at::timestamptz > now() - interval '7 days'
GROUP BY 1 ORDER BY 1;
```

**M3 — overbooking.** Concurrent leased rows must never exceed the agent's
`max_parallel_tasks`. This is the 2x-concurrency failure the #1766 gate closes;
canary S-02 cannot see it (it counts `ZCARD` only), so it needs its own check.
Sample on a short interval during peak, not once.

```sql
SELECT o.agent_name, o.max_parallel_tasks AS cap,
COUNT(*) FILTER (WHERE e.lease_expires_at IS NOT NULL) AS leased_running,
COUNT(*) FILTER (WHERE e.lease_expires_at IS NULL) AS pushed_running
FROM agent_ownership o
LEFT JOIN schedule_executions e
ON e.agent_name = o.agent_name AND e.status = 'running'
WHERE o.agent_name = '<agent>'
GROUP BY 1, 2;
```

Verdict: `leased_running <= cap` always. `leased_running > 0 AND pushed_running > 0`
for an autonomous workload means coexistence is still live.

**M4 — lost / phantom executions.** Any non-terminal row older than the lease
window is the headline AC failing.

```sql
SELECT id, status, triggered_by, started_at, lease_expires_at,
claimed_by_worker, redelivery_count,
EXTRACT(epoch FROM (now() - started_at::timestamptz))::int AS age_s
FROM schedule_executions
WHERE agent_name = '<agent>'
AND status IN ('queued', 'running')
AND started_at::timestamptz < now() - interval '2 hours'
ORDER BY started_at;
```

Expect **zero rows**. A `running` row past its `lease_expires_at` that the reaper
has not touched is a reaper failure; a `queued` row aging with idle workers is a
claim failure.

**M5 — re-delivery + poison-park.** Proves the lease machinery behaves under real
load rather than in the 2026-07-08 synthetic pilot.

```sql
SELECT redelivery_count, COUNT(*)
FROM schedule_executions
WHERE agent_name = '<agent>'
AND started_at::timestamptz > now() - interval '7 days'
GROUP BY 1 ORDER BY 1;
```

Healthy: overwhelmingly `0`, a thin tail at 1–2, nothing at `MAX_REDELIVERY`
without a matching operator alert. Cross-check the parks:

```sql
SELECT id, agent_name, title, created_at, status
FROM operator_queue
WHERE id LIKE 'poison-%' AND created_at::timestamptz > now() - interval '7 days'
ORDER BY created_at DESC;
```

Every row at the cap in the first query must have a `poison-` item here. A park
with no alert is a silent drop — treat as a blocker.

**M6 — outcome regression.** Compare against the pre-flight baseline.

```sql
SELECT status, COUNT(*),
ROUND(AVG(duration_ms)::numeric, 0) AS avg_ms,
ROUND(SUM(cost)::numeric, 4) AS cost
FROM schedule_executions
WHERE agent_name = '<agent>'
AND started_at::timestamptz > now() - interval '7 days'
GROUP BY 1 ORDER BY 2 DESC;
```

Watch for a `failed` share above baseline and for `duration_ms` inflation (pull
adds up to one poll interval of latency by design — the worker idles up to 15s
between claims, so a small p50 rise is expected, a p95 blow-up is not).

**M7 — canary.** With E-05 lease-excluded (this branch), violations should be
genuinely rare.

```sql
SELECT invariant_id, severity, COUNT(*), MAX(snapshot_time) AS last_seen
FROM canary_violations
WHERE snapshot_time::timestamptz > now() - interval '7 days'
GROUP BY 1, 2 ORDER BY 3 DESC;
```

Any **S-01 / S-02 / E-01 / E-02 / L-03** hit on the pilot is an abort signal.
E-05 hits mean the lease exclusion is not deployed — verify the build.

**M8 — starvation.** Queue depth should oscillate, not climb monotonically.

```sql
SELECT COUNT(*) AS queued,
EXTRACT(epoch FROM (now() - MIN(queued_at)::timestamptz))::int AS oldest_s
FROM schedule_executions
WHERE agent_name = '<agent>' AND status = 'queued';
```

A steadily rising `oldest_s` with healthy workers means claims are failing;
with busy workers it just means the agent is undersized — check M3 before
concluding.

### Abort criteria

Stop the soak and revert if any of these hold: a non-empty **M4**; **M3**
`leased_running > cap`; a critical canary violation on the pilot (**M7**); a
park with no operator alert (**M5**); or a `failed` share materially above the
pre-flight baseline (**M6**).

### Rollback

Remove the agent from `PULL_MODE_PILOT_AGENTS` → restart backend → **recreate the
agent** (the flag only clears at create/recreate; `PULL_MODE_ENV_KEYS` is popped
on recreate precisely so de-piloting takes effect). Queued rows are picked up by
the backend drain again as soon as the pilot guard stops matching, so in-flight
work is not stranded.

---

## Appendix — key file references

Expand Down
23 changes: 23 additions & 0 deletions src/backend/canary/invariants/e05_dispatched_rows_have_session.py
Original file line number Diff line number Diff line change
Expand Up @@ -31,6 +31,17 @@
and the observability link to the JSONL file in the container. Real
operational problem, but not the lights-out kind S-01/S-02/E-01 flag.

## Pull-claimed rows are excluded (#1766)

A `#1081` pull-CLAIMED row (`lease_expires_at IS NOT NULL`) is `running`
with a NULL `claude_session_id` by design — the claim is a pure SQL
UPDATE and the worker reports its session id back only with the terminal.
`mark_no_session_executions_failed`, the very sweep this invariant
watches, already carries `lease_expires_at IS NULL` for that reason, so
without the same exclusion here E-05 fires on every pull turn older than
60s: flagging rows the sweep is deliberately leaving alone, and burying a
real #106 regression under noise for a whole soak window.

Tier B, severity major.
"""

Expand Down Expand Up @@ -62,6 +73,18 @@ def check(snapshot: Snapshot) -> List[ViolationReport]:

for agent in snapshot.agents:
for eid in sorted(agent.running_exec_ids):
# #1766: exclude pull-CLAIMED rows (mirrors S-01's exclusion and,
# decisively, the very sweep this invariant watches —
# `mark_no_session_executions_failed` already carries
# `lease_expires_at IS NULL` for exactly this reason). A leased row
# is `running` with a NULL claude_session_id BY DESIGN: the claim is
# a pure SQL UPDATE and the worker reports its session id back only
# with the terminal. Without this, E-05 fires on every pull turn
# older than 60s — flagging the rows the sweep is deliberately
# leaving alone, and drowning a real #106 regression in noise for
# the whole soak window (#1081 Phase 3 / T3.6).
if agent.running_lease_expires_at.get(eid) is not None:
continue
session_id = agent.running_claude_session_ids.get(eid)
if session_id:
continue
Expand Down
7 changes: 5 additions & 2 deletions src/backend/canary/snapshot.py
Original file line number Diff line number Diff line change
Expand Up @@ -222,8 +222,11 @@ class AgentSnapshot:
# owned EXCLUSIVELY by the lease-reaper and NEVER enters the slot ZSET (a
# claim is a pure SQL UPDATE with no ZADD). S-01 uses this to exclude leased
# rows from the SQL side of its slot–row bijection, so a legitimately-
# unslotted pull row is not flagged `in_sql_only`. Only S-01 reads this;
# E-01/E-02/E-05 keep seeing the full `running_exec_ids` unchanged.
# unslotted pull row is not flagged `in_sql_only`. Read by S-01 and, since
# #1766, by E-05 — a leased row is `running` with a NULL claude_session_id
# by design, and `mark_no_session_executions_failed` (the sweep E-05 watches)
# already excludes leased rows for that reason. E-01/E-02 keep seeing the
# full `running_exec_ids` unchanged.
running_lease_expires_at: Dict[str, Optional[str]] = field(default_factory=dict)
# `claude_session_id` per running id (str or None); used by E-05 to detect
# dispatched rows that never acquired a backing session.
Expand Down
2 changes: 1 addition & 1 deletion src/backend/routers/internal.py
Original file line number Diff line number Diff line change
Expand Up @@ -113,7 +113,7 @@ def _pull_authorized(request: Request, agent_name: str) -> bool:
# with a valid scoped key, so rollback takes effect on backend restart rather
# than only after a container recreate. The trusted-backend (internal-secret)
# path is unchanged.
from services.agent_service.pull_mode import is_pull_pilot_agent
from services.pull_pilot import is_pull_pilot_agent
if not is_pull_pilot_agent(agent_name):
return False
return heartbeat_service.authorize_heartbeat(_validated_agent_key(request), agent_name)
Expand Down
28 changes: 17 additions & 11 deletions src/backend/services/agent_service/pull_mode.py
Original file line number Diff line number Diff line change
Expand Up @@ -27,17 +27,23 @@
from __future__ import annotations

import os
from typing import Dict, Set


def _pilot_allowlist() -> Set[str]:
raw = os.getenv("PULL_MODE_PILOT_AGENTS", "")
return {name.strip() for name in raw.split(",") if name.strip()}


def is_pull_pilot_agent(agent_name: str) -> bool:
"""True when ``agent_name`` is in the ``PULL_MODE_PILOT_AGENTS`` allowlist."""
return agent_name in _pilot_allowlist()
from typing import Dict

# The predicates live in the LEAF module `services/pull_pilot.py` and are
# re-exported here so every existing importer keeps working. They CANNOT live in
# this file: `services/agent_service/__init__.py` eagerly imports the whole
# agent-lifecycle stack (helpers/lifecycle/crud/deploy/terminal → models), so a
# dispatch-path module doing `from services.agent_service.pull_mode import ...`
# drags all of it in. See the docstring in `services/pull_pilot.py`.
#
# Dispatch-path callers (`capacity_manager`, `backlog_service`, `routers/internal`)
# import from `services.pull_pilot` DIRECTLY — importing them from here would
# reintroduce exactly the dependency this split removes.
from services.pull_pilot import ( # noqa: F401 — re-exported for back-compat
_pilot_allowlist,
is_pull_pilot_agent,
pull_owns_dispatch,
)


# The container env keys this module manages. Recreate (``lifecycle.py``) pops
Expand Down
20 changes: 20 additions & 0 deletions src/backend/services/backlog_service.py
Original file line number Diff line number Diff line change
Expand Up @@ -152,6 +152,26 @@ async def drain_next(self, agent_name: str) -> bool:

Returns True if a row was drained, False otherwise.
"""
# ---- #1766: a pull pilot's queue has exactly one consumer ----------
# The backend and the agent's worker pool both claim through
# `db.claim_next_queued`, so before this guard they raced for every
# queued row — the atomic UPDATE meant no double-run, but the winner was
# whoever polled first, and the backend was structurally favoured (it
# drains on slot release plus the 60s `drain_orphans_all` sweep, while a
# worker idles up to 15s between polls). That made pull coverage
# nondeterministic and drain-biased: exactly the thing a soak is supposed
# to measure. One guard here covers every drain path, since the release
# callback, the orphan sweep, and `drain_on_release` all funnel through
# this method.
from services.pull_pilot import is_pull_pilot_agent

if is_pull_pilot_agent(agent_name):
logger.debug(
f"[Backlog] Drain skipped for pull pilot '{agent_name}': "
"its worker pool is the sole consumer (#1766)"
)
return False

from database import db

if db.get_queued_count(agent_name) == 0:
Expand Down
43 changes: 33 additions & 10 deletions src/backend/services/capacity_manager.py
Original file line number Diff line number Diff line change
Expand Up @@ -345,15 +345,38 @@ async def acquire(
raise CircuitOpen(agent_name, breaker.retry_after_seconds())
return AcquireResult(state="admitted", execution_id=execution_id)

# First try to acquire a slot directly. SlotService already does the
# atomic ZADD + count check.
admitted = await self._slots.acquire_slot(
agent_name=agent_name,
execution_id=execution_id,
max_parallel_tasks=max_concurrent,
message_preview=message_preview,
timeout_seconds=timeout_seconds,
# ---- #1766: a pull pilot's autonomous work is queue-ONLY -----------
# The pilot flag used to be purely additive: the agent started pulling,
# but the backend kept admitting-and-pushing whenever a slot was free, so
# the two paths ran in parallel over one queue with two independent
# capacity counters. Skipping admission here makes the durable queue the
# single entry point for this agent, so its worker pool IS its capacity
# (#1081 Phase 5, pilot-scoped). `pull_owns_dispatch` excludes
# interactive triggers (Open Question 7 scope cut) and fails safe to
# push, so a non-pilot's path is byte-for-byte unchanged.
from services.pull_pilot import pull_owns_dispatch

pull_exclusive = (
overflow_policy == "queue_persistent"
and overflow_payload is not None
and pull_owns_dispatch(agent_name, overflow_payload.triggered_by)
)

if pull_exclusive:
# No ZADD: the row falls through to the persistent enqueue below and
# is claimed by the agent's own worker. Capacity is enforced
# physically by the pool, not by this counter.
admitted = False
else:
# First try to acquire a slot directly. SlotService already does the
# atomic ZADD + count check.
admitted = await self._slots.acquire_slot(
agent_name=agent_name,
execution_id=execution_id,
max_parallel_tasks=max_concurrent,
message_preview=message_preview,
timeout_seconds=timeout_seconds,
)
if admitted:
return AcquireResult(state="admitted", execution_id=execution_id)

Expand Down Expand Up @@ -517,7 +540,7 @@ async def get_all_states(
# Gated on the EXISTING pilot allowlist; empty allowlist (the default)
# short-circuits to the unchanged ZSET-only path — inert, zero added
# cost/delta. Metering ONLY: acquire/release are not touched.
from services.agent_service.pull_mode import _pilot_allowlist
from services.pull_pilot import _pilot_allowlist

pilots = _pilot_allowlist() & set(clamped)
if pilots:
Expand Down Expand Up @@ -556,7 +579,7 @@ async def get_slot_state(self, agent_name: str, max_concurrent: int):
# short-circuits per-agent, so a non-pilot pays nothing and its output
# is byte-for-byte identical whether the allowlist is set or not.
# Metering ONLY: acquire/release are not touched.
from services.agent_service.pull_mode import is_pull_pilot_agent
from services.pull_pilot import is_pull_pilot_agent

if is_pull_pilot_agent(agent_name):
from database import db
Expand Down
Loading
Loading