Skip to content
50 changes: 49 additions & 1 deletion AGENTS.md
Original file line number Diff line number Diff line change
Expand Up @@ -79,6 +79,53 @@ missing in CI. The `scaling` job deliberately keeps its explicit
`coverage numpy`: every assertion it makes is a wall-clock ratio, and a
first-call JIT compile would land inside the measurement.

### idle_millis is a residual; dependency_wait is its attribution

`idle_millis` is not measured. It is whatever remains of a coordination window
after the four measured phases, so it totals waiting without saying what was
waited on — and it is the swarm's largest single cost (61.5% of all worker time
on epoch 21, against 3.4% for the scan and 2.3% for lock waits).

`telemetry.dependency_wait` carries one row per `cooperative_solve` episode that
reached the wait loop: the dependency's spine, size and budget, how the episode
divided between claiming that branch, helping elsewhere, and being stuck, and
how it ended.

**The `blocks_*` counters say why each sleep happened.** A worker in that loop
tries the branch alone, then anywhere else, then a pair onto the branch, and
only then sleeps:

| counter | cause | reached by |
|---|---|---|
| `blocks_worker_cap` | the branch is at `MAX_WORKERS_PER_BRANCH` | raising the cap |
| `blocks_no_candidates` | nothing left to hand out | nothing the cap can do |
| `blocks_awaiting_finalize` | every candidate done, a rival finalizing | nothing the cap can do |
| `blocks_help_capped` | `MAX_HELP_RECURSION_DEPTH` forbade scanning | raising that cap |
| `blocks_other` | the dependency changed identity, or a retry ran out | nothing; it keeps the sum honest |

Those are different problems with different fixes, and `idle_millis` cannot
tell them apart.

**Every sleep increments exactly one counter, so the five sum to the episode's
sleep count.** That is what `blocks_other` is for — a claim transaction can
decline for reasons with no column of their own, and dropping those would leave
`blocked_millis` holding time no counter accounts for, which is precisely the
defect `idle_millis` has. A counter that is a partition can be audited; a
counter that is a selection cannot.

**Every reason is reported by the code that decided it, never sampled
afterwards.** The two claim outcomes come from `ERDQueue.last_claim_decline()`,
set inside the claim transaction at each of its `return None` sites; the other
two are the branch condition the worker is standing in. Nothing on this path
reads live state, so no counter can disagree with the moment it describes.

That property is the design, not an optimization. **Occupancy is counted
against the claims the claim transaction is about to create, so only that
transaction can say whether the cap refused a claim** — anything reconstructing
it afterwards is reading a system that has already moved, and a holder count
from one instant beside an availability count from another does not blur the
answer, it inverts it. An episode that never reached the loop writes no row.

### Priority ladders, and the fan-out they prevent

**Openers tied at one priority all become eligible at once, and the swarm
Expand Down Expand Up @@ -1166,7 +1213,8 @@ successful merge unless `--keep-source` is given.
Swarm telemetry lives in a **separate** Linux-only file,
`runtime/erd_queue_telemetry.sqlite3` (`<stem>_telemetry<ext>`, computed by
`derive_telemetry_path`), which `ERDQueue` opens as an attached schema named
`telemetry`. The `claim_telemetry` and `branch_finalize_log` tables are there, not
`telemetry`. The `claim_telemetry`, `branch_finalize_log` and `dependency_wait`
tables are there, not
in the main queue file — `add_claim_telemetry` and `add_branch_finalize_log` write
`telemetry.<table>`, and reads join through the `telemetry.` prefix. Because
attached-schema tables do not appear in the main file's `sqlite_master`, running
Expand Down
128 changes: 127 additions & 1 deletion erd_queue.py
Original file line number Diff line number Diff line change
Expand Up @@ -38,6 +38,16 @@
# excluded from live-worker totals.
WORKER_LIVENESS_SECONDS = 30

# Why a claim transaction handed out no bundle. These are decided inside the
# transaction, which is the only place that can decide them: the cap counts
# occupancy against the very claims the transaction is about to create, so a
# caller asking afterwards is asking a system that has already moved on.
CLAIM_DECLINE_BRANCH_GONE = "branch_gone"
CLAIM_DECLINE_BUDGET_MISMATCH = "budget_mismatch"
CLAIM_DECLINE_OWNER_MISMATCH = "owner_mismatch"
CLAIM_DECLINE_WORKER_CAP = "worker_cap" # branch at MAX_WORKERS_PER_BRANCH
CLAIM_DECLINE_NO_CANDIDATES = "no_candidates" # nothing left to hand out

# Time-weighted geometric mean EMA: half-life for the cost model.
_COST_MODEL_TAU = 86400.0 # seconds (≈ 1 day)
# Effective-weight below which a cost-model bucket reads cold (no prediction).
Expand Down Expand Up @@ -851,14 +861,75 @@ def guess_depth_from_spine(spine) -> int:
epoch INTEGER NOT NULL DEFAULT 0,
recorded_at INTEGER NOT NULL
);

-- One row per cooperative_solve that had to wait on a dependency: the worker
-- needed a sub-branch it could not simply claim, and spent time in the wait
-- loop before the branch produced an answer.
--
-- claim_telemetry's idle_millis is a RESIDUAL — the part of a coordination
-- window left after the four measured phases — so it says how much waiting
-- happened but nothing about what was waited on. This table is the
-- attribution: which branch, for how long, and what the worker's alternatives
-- were at the moment it stalled.
--
-- The distinction the columns exist to draw: a worker in this loop first tries
-- to claim the branch alone, then to help anywhere else, then to pair onto the
-- branch, and only then sleeps. blocked_millis covers that last state alone,
-- and the blocks_* counters say WHY each sleep happened. Those are different
-- problems with different fixes, and idle_millis cannot tell them apart.
--
-- Every reason is reported by the code that decided it rather than sampled
-- afterwards: blocks_worker_cap and blocks_no_candidates come from the claim
-- transaction that enforces the cap, the other two from the branch condition
-- the worker is standing in. Nothing here reads live state, so no counter can
-- disagree with the moment it describes.
CREATE TABLE IF NOT EXISTS telemetry.dependency_wait (
id INTEGER PRIMARY KEY AUTOINCREMENT,
worker_id TEXT,
spine TEXT, -- the dependency, not the waiter
n_words INTEGER,
budget INTEGER,
-- Whole time inside the wait loop, and the part of it spent asleep with no
-- claim held and nothing claimable: episode >= blocked always.
episode_millis INTEGER NOT NULL,
blocked_millis INTEGER NOT NULL,
iterations INTEGER NOT NULL,
-- _help_other_branch outcomes: a completed scan that found nothing free or
-- promotable anywhere, versus one that sent this worker off to real work.
empty_scans INTEGER NOT NULL,
helped_scans INTEGER NOT NULL,
-- Bundles this worker claimed on the dependency itself, sole or paired.
bundles_claimed INTEGER NOT NULL,
pair_attempts INTEGER NOT NULL,
pair_successes INTEGER NOT NULL,
-- Why this episode's sleeps happened, one count per cause. Raising
-- MAX_WORKERS_PER_BRANCH reaches blocks_worker_cap and nothing else.
-- Every sleep increments exactly one of these, so they sum to the
-- episode's sleep count and blocked_millis is never time no counter
-- accounts for. blocks_other carries the reasons that are neither
-- actionable nor common -- a dependency that changed identity under the
-- waiter (branch_gone, budget_mismatch, owner_mismatch) and a claim whose
-- retry loop ran out -- which exist to keep that sum honest rather than to
-- be read on their own.
blocks_worker_cap INTEGER NOT NULL DEFAULT 0,
blocks_no_candidates INTEGER NOT NULL DEFAULT 0,
blocks_awaiting_finalize INTEGER NOT NULL DEFAULT 0,
blocks_help_capped INTEGER NOT NULL DEFAULT 0,
blocks_other INTEGER NOT NULL DEFAULT 0,
help_depth INTEGER,
-- How the wait ended: 'solved', 'loss', 'cut', 'deleted', 'cancelled'.
outcome TEXT,
epoch INTEGER NOT NULL DEFAULT 0,
recorded_at INTEGER NOT NULL
);
"""

# Every table _TELEMETRY_SCHEMA_SQL creates, used to detect a queue file that
# predates the telemetry split (see _absorb_legacy_telemetry_tables).
_TELEMETRY_TABLES = (
"bundle_stats", "cost_samples", "claim_telemetry",
"branch_finalize_log", "candidate_accuracy", "backstop_telemetry",
"cut_reuse_misses",
"cut_reuse_misses", "dependency_wait",
)


Expand Down Expand Up @@ -937,6 +1008,12 @@ def __init__(self, db_path: str, timeout: float = 30.0,
self._last_claim_retries = 0
self._last_claim_transaction_millis = 0
self._last_claim_commit_millis = 0
# Why the last claim_next_bundle declined, or None when it handed out a
# bundle. Decided inside the claim transaction, which is the only
# place that can decide it: occupancy is counted against claims the
# same transaction is about to create. A caller reconstructing this
# afterwards is reading a system that has already moved.
self._last_claim_decline = None
# Monotonic per-connection counter for bundle_id generation: paired
# with worker_id and this process's pid, it is unique without a
# timestamp-collision risk (two bundles claimed by the same worker
Expand Down Expand Up @@ -3146,10 +3223,12 @@ def claim_next_bundle(self, branch_key, worker_id, n_candidates,
"FROM active_branches "
"WHERE branch_id = ?", (branch_id,)).fetchone()
if br is None or br["status"] != "open":
self._last_claim_decline = CLAIM_DECLINE_BRANCH_GONE
self._commit_claim_transaction(_txn_t0)
return None
if (expected_budget is not None and br["budget"] is not None
and br["budget"] != expected_budget):
self._last_claim_decline = CLAIM_DECLINE_BUDGET_MISMATCH
self._commit_claim_transaction(_txn_t0)
return None
if expected_opener_work_id is None:
Expand All @@ -3171,11 +3250,13 @@ def claim_next_bundle(self, branch_key, worker_id, n_candidates,
""", (branch_id, expected_opener_work_id,
expected_opener_priority)).fetchone() is not None
if not owner_matches:
self._last_claim_decline = CLAIM_DECLINE_OWNER_MISMATCH
self._commit_claim_transaction(_txn_t0)
return None
if (max_other_workers is not None
and self._other_claim_holders(branch_id, worker_id)
> max_other_workers):
self._last_claim_decline = CLAIM_DECLINE_WORKER_CAP
self._commit_claim_transaction(_txn_t0)
return None
# The branch ceiling is a bound like any achieved best: candidates
Expand Down Expand Up @@ -3278,8 +3359,10 @@ def claim_next_bundle(self, branch_key, worker_id, n_candidates,
packed = self._pack_recorded_holes(
branch_id, bound, cost_lower_bound, survivor_limit)
if not packed:
self._last_claim_decline = CLAIM_DECLINE_NO_CANDIDATES
self._commit_claim_transaction(_txn_t0)
return None
self._last_claim_decline = None
bundle = [idx for idx, _position in packed]
bundle_id = f"{worker_id}:{self._pid}:{self._bundle_seq}"
self._bundle_seq += 1
Expand Down Expand Up @@ -3779,6 +3862,17 @@ def branch_done_candidates(self, branch_key) -> int:
"SELECT COUNT(*) FROM candidate_claims WHERE branch_id = ? AND done = 1",
(branch_id,)).fetchone()[0]

def last_claim_decline(self):
"""Why the last claim_next_bundle on this connection declined, or None
when it handed out a bundle.

Read straight after a claim that returned None. This is an observation
of the decision, not a reconstruction of it: the transaction that
enforces the cap is the same one that reports it, so the answer cannot
disagree with the state the cap was enforced against.
"""
return self._last_claim_decline

def branch_bulk_done_candidates(self, branch_key) -> int:
"""Return the legacy combined count completed by ERD pruning."""
branch_id = self._intern_branch(branch_key)
Expand Down Expand Up @@ -6803,6 +6897,38 @@ def add_branch_finalize_log(self, branch_key, spine, n_words, budget,
first_best_at, nodes_at_first_best,
now))

def add_dependency_wait(self, worker_id, spine, n_words, budget,
episode_millis, blocked_millis, iterations,
empty_scans, helped_scans, bundles_claimed,
pair_attempts, pair_successes,
blocks_worker_cap=0, blocks_no_candidates=0,
blocks_awaiting_finalize=0, blocks_help_capped=0,
blocks_other=0, help_depth=None, outcome=None):
"""Record one cooperative_solve wait episode.

Attributes what claim_telemetry's idle_millis can only total: which
dependency the worker waited on, how much of the episode it was asleep
rather than helping or claiming, and why each sleep happened. A caller
that never blocked has nothing to attribute and should not write a row.
"""
now = int(time.time())
self._conn.execute("""
INSERT INTO telemetry.dependency_wait
(worker_id, spine, n_words, budget, episode_millis,
blocked_millis, iterations, empty_scans, helped_scans,
bundles_claimed, pair_attempts, pair_successes,
blocks_worker_cap, blocks_no_candidates,
blocks_awaiting_finalize, blocks_help_capped, blocks_other,
help_depth, outcome, epoch, recorded_at)
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?,
?)
""", (worker_id, spine, n_words, budget, episode_millis,
blocked_millis, iterations, empty_scans, helped_scans,
bundles_claimed, pair_attempts, pair_successes,
blocks_worker_cap, blocks_no_candidates,
blocks_awaiting_finalize, blocks_help_capped, blocks_other,
help_depth, outcome, self.epoch, now))

def add_cut_reuse_miss(self, branch_key, n_words, budget, wanted_ceiling,
available_bound, available_budget):
"""Log a re-solve forced by a cut: a consumer needed this branch but
Expand Down
Loading
Loading