From 98413fefedc23bca4b2f70e78389e1345ba3085f Mon Sep 17 00:00:00 2001 From: Sean Ahern Date: Thu, 17 Sep 2026 16:31:55 -0400 Subject: [PATCH 1/7] Attribute dependency waiting, the swarm's largest unexplained cost claim_telemetry.idle_millis is computed as the remainder of a coordination window after the four measured phases, so it totals waiting without naming what was waited on. It is also the largest single line item in the swarm: 61.5% of all worker time on epoch 21, against 3.4% for the scheduling scan and 2.3% for lock waits. telemetry.dependency_wait records 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. holders_at_first_block and unclaimed_at_first_block are the columns the worker-cap question turns on. A worker in that loop tries the branch alone, then anywhere else, then a pair onto the branch, and only then sleeps; holders at MAX_WORKERS_PER_BRANCH means the cap refused the pair, while zero unclaimed candidates means the branch had nothing left to hand out and the wait is for its finalize. Those are different problems with different fixes and idle_millis cannot separate them. They are sampled once, at the first blocked iteration: that path already runs every 50 ms on a starving worker, and re-reading them per poll would put two queries on the branch everyone is waiting for. branch_claim_holders is the single-branch form of claim_holders_by_branch, so the sample costs one indexed query rather than a map of every branch. The table is created by _TELEMETRY_SCHEMA_SQL with CREATE TABLE IF NOT EXISTS, which runs on every open, so an existing telemetry file gains it without a migration step to sequence. Workers fork from the supervisor, so deploying still needs a swarm stop. Co-Authored-By: Claude Opus 5 Claude-Session: https://claude.ai/code/session_019rawcKNbdKoVpUW1Avr4sS --- AGENTS.md | 29 ++++++- erd_queue.py | 101 +++++++++++++++++++++- erd_swarm.py | 131 ++++++++++++++++++++++++++++- tests/test_claim_packing.py | 2 +- tests/test_erd_swarm_unit.py | 158 +++++++++++++++++++++++++++++++++++ 5 files changed, 417 insertions(+), 4 deletions(-) diff --git a/AGENTS.md b/AGENTS.md index 3a1604d5..3c7dbabf 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -79,6 +79,32 @@ 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 two columns that matter are `holders_at_first_block` and +`unclaimed_at_first_block`.** A worker in that loop tries the branch alone, +then anywhere else, then a pair onto the branch, and only then sleeps. Holders +at `MAX_WORKERS_PER_BRANCH` means the cap refused the pair — raising it would +admit this worker. Zero unclaimed candidates means the branch had nothing left +to hand out and the wait is for its finalize, which no cap change reaches. +Those are different problems with different fixes, and `idle_millis` cannot +tell them apart. + +They are sampled once, at the first blocked iteration, because that path +already runs every 50 ms on a starving worker and two more queries per turn +would be paid by the branch everyone is waiting for. 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 @@ -1166,7 +1192,8 @@ successful merge unless `--keep-source` is given. Swarm telemetry lives in a **separate** Linux-only file, `runtime/erd_queue_telemetry.sqlite3` (`_telemetry`, 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.`, and reads join through the `telemetry.` prefix. Because attached-schema tables do not appear in the main file's `sqlite_master`, running diff --git a/erd_queue.py b/erd_queue.py index 990b4ab6..1b9f1756 100644 --- a/erd_queue.py +++ b/erd_queue.py @@ -846,6 +846,53 @@ 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 holders/unclaimed at first block say whether the pair was refused +-- because the branch was already at MAX_WORKERS_PER_BRANCH or because there +-- was nothing left on it to claim. Those are different problems with +-- different fixes, and idle_millis cannot tell them apart. +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, + -- Sampled once, at the first blocked iteration: re-reading them every + -- 50 ms poll would put a query on the one path that is already starving. + holders_at_first_block INTEGER, + unclaimed_at_first_block INTEGER, + 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 @@ -853,7 +900,7 @@ def guess_depth_from_spine(spine) -> int: _TELEMETRY_TABLES = ( "bundle_stats", "cost_samples", "claim_telemetry", "branch_finalize_log", "candidate_accuracy", "backstop_telemetry", - "cut_reuse_misses", + "cut_reuse_misses", "dependency_wait", ) @@ -3774,6 +3821,28 @@ def branch_done_candidates(self, branch_key) -> int: "SELECT COUNT(*) FROM candidate_claims WHERE branch_id = ? AND done = 1", (branch_id,)).fetchone()[0] + def branch_claim_holders(self, branch_key, exclude_worker_id=None) -> int: + """Workers other than exclude_worker_id holding unfinished claims on + one branch. + + The single-branch form of claim_holders_by_branch, for a caller that + wants occupancy for the branch it is already looking at rather than a + map of every branch. Occupancy is unfinished claims, never heartbeats: + the claim row is written by the transaction that hands out the bundle, + so a branch reads as taken the instant it is taken. + """ + branch_id = self._intern_branch(branch_key) + if branch_id is None: + return 0 + if exclude_worker_id is None: + return self._conn.execute( + "SELECT COUNT(DISTINCT claimed_by) FROM candidate_claims " + "WHERE branch_id = ? AND done = 0", (branch_id,)).fetchone()[0] + return self._conn.execute( + "SELECT COUNT(DISTINCT claimed_by) FROM candidate_claims " + "WHERE branch_id = ? AND done = 0 AND claimed_by IS NOT ?", + (branch_id, exclude_worker_id)).fetchone()[0] + 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) @@ -6798,6 +6867,36 @@ 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, + holders_at_first_block=None, + unclaimed_at_first_block=None, + 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 what its alternatives were when it + first stalled. 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, + holders_at_first_block, unclaimed_at_first_block, + 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, + holders_at_first_block, unclaimed_at_first_block, + 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 diff --git a/erd_swarm.py b/erd_swarm.py index 0a2c98bb..165cbd52 100644 --- a/erd_swarm.py +++ b/erd_swarm.py @@ -585,6 +585,75 @@ def record_inline(self, token): buf[key] = (log_n, log_n * log_n, 1) +class _DependencyWait: + """One cooperative_solve wait episode, accumulated for telemetry. + + `claim_telemetry.idle_millis` is a residual — whatever is left of a + coordination window after the four measured phases — so it totals waiting + without saying what was waited on. This carries the attribution: the + dependency, how the episode divided between working on it, helping + elsewhere, and being stuck, and what the worker's alternatives were the + first time it had none. + + `blocked_millis` is the stuck state alone, entered only after the worker + has failed to claim the branch by itself, failed to find work anywhere + else, and failed to pair onto the branch. `holders_at_first_block` and + `unclaimed_at_first_block` are sampled once at that moment, which is what + separates "the cap refused the pair" from "the branch had nothing left to + claim". + """ + + __slots__ = ("spine", "n_words", "budget", "help_depth", "outcome", + "_started", "iterations", "blocked_millis", "empty_scans", + "helped_scans", "bundles_claimed", "pair_attempts", + "pair_successes", "holders_at_first_block", + "unclaimed_at_first_block") + + def __init__(self, spine, n_words, budget, help_depth): + self.spine = spine + self.n_words = n_words + self.budget = budget + self.help_depth = help_depth + self.outcome = None + self._started = time.perf_counter() + self.iterations = 0 + self.blocked_millis = 0 + self.empty_scans = 0 + self.helped_scans = 0 + self.bundles_claimed = 0 + self.pair_attempts = 0 + self.pair_successes = 0 + self.holders_at_first_block = None + self.unclaimed_at_first_block = None + + @property + def episode_millis(self): + return int((time.perf_counter() - self._started) * 1000) + + def note_first_block(self, holders, unclaimed): + """Record the alternatives once, on the first blocked iteration. + + Sampled once rather than per poll: the blocked path already runs every + 50 ms on a starving worker, and two more queries per turn of it would + be paid by the branch everyone is waiting for. + """ + if self.holders_at_first_block is None: + self.holders_at_first_block = holders + self.unclaimed_at_first_block = unclaimed + + def note_blocked(self, millis): + self.blocked_millis += millis + + def should_record(self): + """True when the episode reached the wait loop at all. + + An episode answered from cache before the loop has nothing to + attribute, and writing a row for it would put queue traffic on the + path that is already free. + """ + return self.iterations > 0 + + class _BranchWorker: """One worker process's state and operations on branches and candidates.""" @@ -2327,6 +2396,35 @@ def _read_satisfying_cut(self, branch_key, budget, ceiling, n_words): return (OVER_ERD_LIMIT, cut_bound, None, cut_tainted) return None + def _record_first_block(self, wait, branch_key): + """Sample the worker's alternatives the first time it has none. + + Reached only after the sole-worker claim, the scan for work elsewhere, + and (on the uncapped path) the pair attempt have all failed, so the two + counts answer why: holders at MAX_WORKERS_PER_BRANCH means the cap + refused the pair, while zero unclaimed candidates means the branch had + nothing left to hand out and the wait is for its finalize. + """ + if wait.holders_at_first_block is not None: + return + holders = self.queue.branch_claim_holders( + branch_key, exclude_worker_id=self.name) + done = self.queue.branch_done_candidates(branch_key) + wait.note_first_block(holders, max(0, self.n_candidates - done)) + + def _record_dependency_wait(self, wait): + """Persist one wait episode, unless it never reached the wait loop.""" + if not wait.should_record(): + return + self.queue.add_dependency_wait( + self.name, wait.spine, wait.n_words, wait.budget, + wait.episode_millis, wait.blocked_millis, wait.iterations, + wait.empty_scans, wait.helped_scans, wait.bundles_claimed, + wait.pair_attempts, wait.pair_successes, + holders_at_first_block=wait.holders_at_first_block, + unclaimed_at_first_block=wait.unclaimed_at_first_block, + help_depth=wait.help_depth, outcome=wait.outcome) + def cooperative_solve(self, words, budget, ceiling=float('inf')): """Solve sub-branch `words` at `budget` cooperatively, returning the engine's (status, cost, max_depth, floor) tuple. @@ -2454,20 +2552,27 @@ def cooperative_solve(self, words, budget, ceiling=float('inf')): context_stack.enter_context(self._entered( parent_context.descend(branch_key, child_spine))) + wait = _DependencyWait(child_spine, n_words, budget, + self._help_recursion_depth) + context_stack.callback(self._record_dependency_wait, wait) while not self.cancel(): + wait.iterations += 1 # Finished? Check before claiming so we never touch a branch that # another worker just finalized and deleted. reuse = _cache_reuse( self.score_cache.read_for_budget(branch_key, ERD_ALL, budget), budget) if reuse is not None: + wait.outcome = 'solved' return (SOLVED, *reuse) loss_budget = self.score_cache.read_loss( branch_key, ERD_ALL, refresh=True) if loss_budget is not None and budget <= loss_budget: + wait.outcome = 'loss' return (OVER_DEPTH_BUDGET, float('inf'), None, True) cut = self._read_satisfying_cut(branch_key, budget, ceiling, n_words) if cut is not None: + wait.outcome = 'cut' return cut if self.queue.get_branch(branch_key) is None: # Finalized + deleted: a cut lands in cut_results and a @@ -2476,10 +2581,13 @@ def cooperative_solve(self, words, budget, ceiling=float('inf')): loss_budget = self.score_cache.read_loss( branch_key, ERD_ALL, refresh=True) if loss_budget is not None and budget <= loss_budget: + wait.outcome = 'loss' return (OVER_DEPTH_BUDGET, float('inf'), None, True) cut = self._read_satisfying_cut(branch_key, budget, ceiling, n_words) if cut is not None: + wait.outcome = 'cut' return cut + wait.outcome = 'deleted' break # finalized as a loss + deleted # Sole worker only: a dependency someone else already holds is # being worked, and this worker is worth far more opening a @@ -2489,6 +2597,7 @@ def cooperative_solve(self, words, budget, ceiling=float('inf')): expected_opener_work_id=self._work_context.opener_work_id, max_other_workers=0) if claim is not None: + wait.bundles_claimed += 1 self._evaluate_dependency_bundle( branch_key, words, n_words, claim, budget) elif self.queue.branch_done_candidates(branch_key) >= self.n_candidates: @@ -2504,7 +2613,11 @@ def cooperative_solve(self, words, budget, ceiling=float('inf')): # _help_other_branch's capped-depth contract already # promises its callers. self._cur_candidate = None + blocked_at = time.perf_counter() self._idle_wait(0.05) + wait.note_blocked( + int((time.perf_counter() - blocked_at) * 1000)) + self._record_first_block(wait, branch_key) else: # No bundle: every candidate is claimed, or another worker # holds the branch. Try free or promotable work first — @@ -2519,7 +2632,10 @@ def cooperative_solve(self, words, budget, ceiling=float('inf')): # contends against the very lock every other worker's # writes need too. self._cur_candidate = None # coordinating, no candidate in flight - if not self._help_other_branch(branch_key): + if self._help_other_branch(branch_key): + wait.helped_scans += 1 + else: + wait.empty_scans += 1 # A completed scan found nothing free or promotable # anywhere (the recursion cap is ruled out by the # branch above, so this really is an empty scan, not a @@ -2535,20 +2651,33 @@ def cooperative_solve(self, words, budget, ceiling=float('inf')): self._heartbeat(branch_key, n_words, None, None, None, None, force=True) self.queue.reclaim_stale_claims(HB_TIMEOUT_SECONDS) + wait.pair_attempts += 1 paired = self._claim_bundle( branch_key, self.n_candidates, words, budget, expected_opener_work_id=( self._work_context.opener_work_id), max_other_workers=MAX_WORKERS_PER_BRANCH - 1) if paired is not None: + wait.pair_successes += 1 + wait.bundles_claimed += 1 self._evaluate_dependency_bundle( branch_key, words, n_words, paired, budget) else: + # Nothing claimable anywhere, no pair available on + # the dependency: the stuck state idle_millis + # totals without naming. + self._record_first_block(wait, branch_key) + blocked_at = time.perf_counter() self._idle_wait(0.05) # let claims land + wait.note_blocked( + int((time.perf_counter() - blocked_at) * 1000)) if self.cancel(): # pragma: no cover + wait.outcome = 'cancelled' return CANCEL_RECVD # Finalized as a loss: proven unsolvable within budget (not a cutoff). + if wait.outcome is None: # pragma: no cover + wait.outcome = 'loss' return (OVER_DEPTH_BUDGET, float('inf'), None, True) # pragma: no cover # -- scheduling: claim one candidate from the best available branch ------ diff --git a/tests/test_claim_packing.py b/tests/test_claim_packing.py index 3e79c0db..c8e6ba73 100644 --- a/tests/test_claim_packing.py +++ b/tests/test_claim_packing.py @@ -221,7 +221,7 @@ def test_telemetry_file_is_created_alongside_queue(self): "bundle_stats", "cost_samples", "claim_telemetry", "branch_finalize_log", "candidate_accuracy", "backstop_telemetry", "cut_reuse_misses", - "two_level_prune_telemetry"}) + "two_level_prune_telemetry", "dependency_wait"}) N_CANDIDATES = 40 diff --git a/tests/test_erd_swarm_unit.py b/tests/test_erd_swarm_unit.py index 68b447c7..b0269547 100644 --- a/tests/test_erd_swarm_unit.py +++ b/tests/test_erd_swarm_unit.py @@ -5918,3 +5918,161 @@ def test_pairing_still_joins_a_branch_from_the_reused_list(self): work = self._worker(90)._claim_one_uninstrumented() self.assertEqual(self._claimed_key(work), joinable_key) + + +class TestDependencyWaitAttribution(unittest.TestCase): + """cooperative_solve records what it waited on, so idle time has a subject. + + `claim_telemetry.idle_millis` is a residual and totals waiting without + naming it. These pin the attribution: which dependency, how much of the + episode was spent stuck rather than working or helping, and — the field the + worker-cap question turns on — whether anything was claimable when the + worker first had nothing to do. + """ + + def setUp(self): + self._tmp = tempfile.TemporaryDirectory() + self.addCleanup(self._tmp.cleanup) + self.answer_file = self._write("answers.txt", BRANCH) + self.words_file = self._write("words.txt", CANDIDATES) + for attr, path in [("ANSWER_FILE", self.answer_file), + ("WORDS_FILE", self.words_file)]: + p = mock.patch.object(erd_swarm, attr, path) + p.start() + self.addCleanup(p.stop) + self.cache_path = os.path.join(self._tmp.name, "cache.sqlite3") + self.queue_path = os.path.join(self._tmp.name, "queue.sqlite3") + + def _write(self, name, words): + p = os.path.join(self._tmp.name, name) + with open(p, "w") as f: + f.write("\n".join(words) + "\n") + return p + + def _worker(self): + w = _BranchWorker(0, self.cache_path, self.queue_path, None) + self.addCleanup(w.close) + return w + + def _rival_claims(self, w, key): + """A rival worker takes a bundle, the way production leaves occupancy: + an unfinished claim row written by the transaction that handed it out.""" + order = list(range(w.n_candidates)) + return w.queue.claim_next_bundle( + key, "rival", w.n_candidates, order, [0.0] * w.n_candidates, + small_count=1, count_cap=1) + + def _wait_rows(self): + path = erd_queue.derive_telemetry_path(self.queue_path) + conn = sqlite3.connect(path) + conn.row_factory = sqlite3.Row + try: + return [dict(r) for r in + conn.execute("SELECT * FROM dependency_wait ORDER BY id")] + finally: + conn.close() + + def test_cache_hit_writes_no_wait_row(self): + """An episode answered before the wait loop never opens an episode.""" + words = BRANCH[:3] + sc = ScoreCache(self.cache_path, BRANCH) + sc.write(ScoreCache.encode_subset(words), ERD_ALL, "crane", 1.5, + max_depth=2, solve_budget=None) + sc.close() + w = self._worker() + self.assertEqual(w.cooperative_solve(words, ROOT_BUDGET)[0], SOLVED) + self.assertEqual(self._wait_rows(), []) + + def test_an_episode_that_never_iterates_writes_no_row(self): + """A worker asked to stop before the first iteration attributed nothing. + + The wait object exists by then — it is built immediately before the + loop — so only should_record keeps this row out. Distinct from the + cache-hit path above, which returns before the episode is opened at + all. + """ + w = self._worker() + w.request_stop() + self.assertEqual(w.cooperative_solve(BRANCH[:3], ROOT_BUDGET), + erd_swarm.CANCEL_RECVD) + self.assertEqual(self._wait_rows(), []) + + def test_a_solved_dependency_names_the_branch_it_waited_on(self): + words = BRANCH[:3] + w = self._worker() + status, _cost, _md, _taint = w.cooperative_solve(words, ROOT_BUDGET) + self.assertEqual(status, SOLVED) + rows = self._wait_rows() + self.assertEqual(len(rows), 1) + row = rows[0] + self.assertEqual(row["n_words"], len(words)) + self.assertEqual(row["budget"], ROOT_BUDGET) + self.assertEqual(row["outcome"], "solved") + self.assertGreaterEqual(row["iterations"], 1) + self.assertGreaterEqual(row["bundles_claimed"], 1) + + def test_blocked_time_never_exceeds_the_episode(self): + """blocked_millis is a part of the episode, not a separate clock.""" + w = self._worker() + w.cooperative_solve(BRANCH[:3], ROOT_BUDGET) + for row in self._wait_rows(): + self.assertLessEqual(row["blocked_millis"], row["episode_millis"]) + self.assertGreaterEqual(row["blocked_millis"], 0) + + def test_first_block_is_sampled_once(self): + """Re-sampling would put two queries on every 50 ms poll of a stall.""" + w = self._worker() + wait = erd_swarm._DependencyWait("SPINE -----", 3, 5, 0) + key = ScoreCache.encode_subset(BRANCH[:3]) + w.queue.create_branch(key, 3, w.n_candidates, budget=5) + w._record_first_block(wait, key) + first = (wait.holders_at_first_block, wait.unclaimed_at_first_block) + w.queue.mark_claims_done(key, list(range(w.n_candidates))) + w._record_first_block(wait, key) + self.assertEqual( + (wait.holders_at_first_block, wait.unclaimed_at_first_block), first) + + def test_cap_refusal_is_distinguishable_from_an_exhausted_branch(self): + """The two reasons a pair fails must not read alike. + + A branch another worker holds with candidates left is a cap refusal — + raising MAX_WORKERS_PER_BRANCH would admit this worker. A branch with + nothing left to claim is waiting on its own finalize, which no cap + change reaches. idle_millis cannot tell these apart; these columns + must. + """ + w = self._worker() + key = ScoreCache.encode_subset(BRANCH[:3]) + w.queue.create_branch(key, 3, w.n_candidates, budget=5) + self.assertIsNotNone(self._rival_claims(w, key)) + + occupied = erd_swarm._DependencyWait("SPINE -----", 3, 5, 0) + w._record_first_block(occupied, key) + self.assertGreaterEqual(occupied.holders_at_first_block, 1) + self.assertGreater(occupied.unclaimed_at_first_block, 0) + + w.queue.mark_claims_done(key, list(range(w.n_candidates))) + exhausted = erd_swarm._DependencyWait("SPINE -----", 3, 5, 0) + w._record_first_block(exhausted, key) + self.assertEqual(exhausted.unclaimed_at_first_block, 0) + + def test_branch_claim_holders_of_an_unknown_branch_is_zero(self): + """A branch the queue has never interned holds nobody.""" + w = self._worker() + self.assertEqual( + w.queue.branch_claim_holders(ScoreCache.encode_subset(BRANCH)), 0) + + def test_branch_claim_holders_agrees_with_the_map(self): + """The single-branch count must not drift from the map scheduling uses.""" + w = self._worker() + key = ScoreCache.encode_subset(BRANCH[:3]) + w.queue.create_branch(key, 3, w.n_candidates, budget=5) + self.assertEqual(w.queue.branch_claim_holders(key), 0) + self.assertIsNotNone(self._rival_claims(w, key)) + self.assertEqual( + w.queue.branch_claim_holders(key, exclude_worker_id="rival"), 0) + self.assertEqual(w.queue.branch_claim_holders(key), 1) + self.assertEqual( + w.queue.branch_claim_holders(key, exclude_worker_id=w.name), + w.queue.claim_holders_by_branch( + exclude_worker_id=w.name).get(bytes(key), 0)) From 14b30b459a655891c8329d43379a906e440e073e Mon Sep 17 00:00:00 2001 From: Sean Ahern Date: Thu, 17 Sep 2026 16:48:10 -0400 Subject: [PATCH 2/7] Count only slots no claim row covers, and snapshot before the poll MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Two review findings on the first-block sample, both of which defeated the distinction the columns exist to draw. unclaimed_at_first_block used n_candidates - done, which counts a candidate held in an unfinished claim as available. That is the ordinary finalize-wait shape — one rival holding every remaining slot — so the pair attempt would find no bundle while telemetry reported unclaimed work, reading as a cap refusal on a branch with nothing left to hand out. branch_unclaimed_candidates counts slots no claim row covers, so a freed position still counts as claimable and an in-flight one does not. The recursion-capped path sampled after its 50 ms poll rather than before it. A holder that finished mid-sleep read as zero holders, so both columns described a moment that blocked nothing. It now samples first, as the uncapped path already did. test_candidates_held_in_flight_do_not_read_as_unclaimed and test_the_capped_path_snapshots_before_it_sleeps each fail with their fix reverted. The second drives the capped path directly, with a patched _idle_wait that finishes the rival mid-sleep, because no existing test reaches that branch and the ordering is invisible without it. Co-Authored-By: Claude Opus 5 Claude-Session: https://claude.ai/code/session_019rawcKNbdKoVpUW1Avr4sS --- erd_queue.py | 20 +++++++++ erd_swarm.py | 11 +++-- tests/test_erd_swarm_unit.py | 79 +++++++++++++++++++++++++++++++++++- 3 files changed, 105 insertions(+), 5 deletions(-) diff --git a/erd_queue.py b/erd_queue.py index 1b9f1756..f9d65a9d 100644 --- a/erd_queue.py +++ b/erd_queue.py @@ -3843,6 +3843,26 @@ def branch_claim_holders(self, branch_key, exclude_worker_id=None) -> int: "WHERE branch_id = ? AND done = 0 AND claimed_by IS NOT ?", (branch_id, exclude_worker_id)).fetchone()[0] + def branch_unclaimed_candidates(self, branch_key, n_candidates) -> int: + """Candidate slots on this branch that no claim row covers. + + NOT `n_candidates - done`: a candidate held in an unfinished claim is + taken, not available, so counting it as unclaimed would report work the + packer cannot hand out — which is the normal finalize-wait shape, with + one rival holding every remaining slot. A freed position (reclaim or + republish) has no row and is correctly counted as available again. + + Zero for a branch the registry does not know: there is nothing left to + claim on a branch that no longer exists. + """ + branch_id = self._intern_branch(branch_key) + if branch_id is None: + return 0 + taken = self._conn.execute( + "SELECT COUNT(*) FROM candidate_claims WHERE branch_id = ?", + (branch_id,)).fetchone()[0] + return max(0, n_candidates - taken) + 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) diff --git a/erd_swarm.py b/erd_swarm.py index 165cbd52..867c436a 100644 --- a/erd_swarm.py +++ b/erd_swarm.py @@ -2404,13 +2404,18 @@ def _record_first_block(self, wait, branch_key): counts answer why: holders at MAX_WORKERS_PER_BRANCH means the cap refused the pair, while zero unclaimed candidates means the branch had nothing left to hand out and the wait is for its finalize. + + Taken before the poll that follows it on every path, so the snapshot + describes the state that caused the block rather than whatever the + branch became while this worker slept. """ if wait.holders_at_first_block is not None: return holders = self.queue.branch_claim_holders( branch_key, exclude_worker_id=self.name) - done = self.queue.branch_done_candidates(branch_key) - wait.note_first_block(holders, max(0, self.n_candidates - done)) + unclaimed = self.queue.branch_unclaimed_candidates( + branch_key, self.n_candidates) + wait.note_first_block(holders, unclaimed) def _record_dependency_wait(self, wait): """Persist one wait episode, unless it never reached the wait loop.""" @@ -2613,11 +2618,11 @@ def cooperative_solve(self, words, budget, ceiling=float('inf')): # _help_other_branch's capped-depth contract already # promises its callers. self._cur_candidate = None + self._record_first_block(wait, branch_key) blocked_at = time.perf_counter() self._idle_wait(0.05) wait.note_blocked( int((time.perf_counter() - blocked_at) * 1000)) - self._record_first_block(wait, branch_key) else: # No bundle: every candidate is claimed, or another worker # holds the branch. Try free or promotable work first — diff --git a/tests/test_erd_swarm_unit.py b/tests/test_erd_swarm_unit.py index b0269547..aad8234a 100644 --- a/tests/test_erd_swarm_unit.py +++ b/tests/test_erd_swarm_unit.py @@ -5954,13 +5954,13 @@ def _worker(self): self.addCleanup(w.close) return w - def _rival_claims(self, w, key): + def _rival_claims(self, w, key, count_cap=1): """A rival worker takes a bundle, the way production leaves occupancy: an unfinished claim row written by the transaction that handed it out.""" order = list(range(w.n_candidates)) return w.queue.claim_next_bundle( key, "rival", w.n_candidates, order, [0.0] * w.n_candidates, - small_count=1, count_cap=1) + small_count=count_cap, count_cap=count_cap) def _wait_rows(self): path = erd_queue.derive_telemetry_path(self.queue_path) @@ -6056,6 +6056,81 @@ def test_cap_refusal_is_distinguishable_from_an_exhausted_branch(self): w._record_first_block(exhausted, key) self.assertEqual(exhausted.unclaimed_at_first_block, 0) + def test_the_capped_path_snapshots_before_it_sleeps(self): + """The snapshot must describe the state that caused the block. + + On the recursion-capped path the worker polls instead of scanning. If + the sample were taken after the 50 ms sleep it would observe whatever + the branch became while this worker slept — a holder that finished + reads as zero holders — and both diagnostic columns would describe a + moment that never blocked anything. + """ + w = self._worker() + words = BRANCH[:3] + key = ScoreCache.encode_subset(words) + w.queue.create_branch(key, len(words), w.n_candidates, + budget=ROOT_BUDGET) + self.assertIsNotNone( + self._rival_claims(w, key, count_cap=w.n_candidates)) + + def _sleep_that_changes_the_branch(_seconds): + # The rival finishes mid-sleep: holders drop to zero, so a sample + # taken afterwards would report nobody was on the branch. + w.queue.mark_claims_done(key, list(range(w.n_candidates))) + w.request_stop() + + with mock.patch.object(erd_swarm, "MAX_HELP_RECURSION_DEPTH", 0), \ + mock.patch.object(w, "_idle_wait", + side_effect=_sleep_that_changes_the_branch): + w.cooperative_solve(words, ROOT_BUDGET) + + rows = self._wait_rows() + self.assertEqual(len(rows), 1) + self.assertEqual(rows[0]["holders_at_first_block"], 1) + self.assertEqual(rows[0]["unclaimed_at_first_block"], 0) + + def test_candidates_held_in_flight_do_not_read_as_unclaimed(self): + """A slot inside an unfinished claim is taken, not available. + + This is the ordinary finalize-wait shape: one rival holds every + remaining candidate, so the pair attempt finds no bundle. Counting + those slots as unclaimed would report a cap refusal on a branch that + simply has nothing left to hand out, which is precisely the distinction + these columns exist to draw. + """ + w = self._worker() + key = ScoreCache.encode_subset(BRANCH[:3]) + w.queue.create_branch(key, 3, w.n_candidates, budget=5) + self.assertIsNotNone( + self._rival_claims(w, key, count_cap=w.n_candidates)) + self.assertEqual(w.queue.branch_done_candidates(key), 0) + + wait = erd_swarm._DependencyWait("SPINE -----", 3, 5, 0) + w._record_first_block(wait, key) + self.assertEqual(wait.unclaimed_at_first_block, 0) + self.assertGreaterEqual(wait.holders_at_first_block, 1) + + def test_a_freed_position_counts_as_claimable_again(self): + """A reclaimed or republished slot has no row and is available.""" + w = self._worker() + key = ScoreCache.encode_subset(BRANCH[:3]) + w.queue.create_branch(key, 3, w.n_candidates, budget=5) + self.assertIsNotNone( + self._rival_claims(w, key, count_cap=w.n_candidates)) + self.assertEqual( + w.queue.branch_unclaimed_candidates(key, w.n_candidates), 0) + w.queue.reclaim_claims_of_worker("rival") + self.assertEqual( + w.queue.branch_unclaimed_candidates(key, w.n_candidates), + w.n_candidates) + + def test_unclaimed_of_an_unknown_branch_is_zero(self): + """A branch the registry never saw has nothing left to claim.""" + w = self._worker() + self.assertEqual( + w.queue.branch_unclaimed_candidates( + ScoreCache.encode_subset(BRANCH), w.n_candidates), 0) + def test_branch_claim_holders_of_an_unknown_branch_is_zero(self): """A branch the queue has never interned holds nobody.""" w = self._worker() From 4f51b53bc130965dc0d131ddb1b3b2d0546502af Mon Sep 17 00:00:00 2001 From: Sean Ahern Date: Thu, 17 Sep 2026 16:56:21 -0400 Subject: [PATCH 3/7] Attribute the finalize wait, and read a deleted branch as finished MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Two more review findings, both on cases the columns claim to name. _await_rival_finalize ends in the same 50 ms poll but was called outside the accounting, so an episode that lost the finalize race put that sleep in episode_millis while blocked_millis read zero and both first-block counters stayed NULL. That is the wait-for-finalize case these columns exist to identify, and it is the common one. It now samples before the call and returns whether it polled, so the takeover path — which reopens a dead finalizer's row and completes the finalize — is counted as the work it is rather than as waiting. branch_unclaimed_candidates trusted a registry id as evidence the branch exists. delete_branch deliberately keeps the append-only branches row so branch_id stays stable across a re-promotion, while dropping every candidate_claims row, so a finished branch counted zero claims and reported all n_candidates as claimable: completed work described as untouched, and a cap refusal where nothing remains to claim. It now requires a live active_branches row. Each fix has a test verified to fail with it reverted, including one that pins the takeover path as unblocked so the poll and the takeover cannot be conflated. Co-Authored-By: Claude Opus 5 Claude-Session: https://claude.ai/code/session_019rawcKNbdKoVpUW1Avr4sS --- erd_queue.py | 13 +++++- erd_swarm.py | 23 +++++++++-- tests/test_erd_swarm_unit.py | 78 ++++++++++++++++++++++++++++++++++++ 3 files changed, 108 insertions(+), 6 deletions(-) diff --git a/erd_queue.py b/erd_queue.py index b2b570fc..24886374 100644 --- a/erd_queue.py +++ b/erd_queue.py @@ -3857,12 +3857,21 @@ def branch_unclaimed_candidates(self, branch_key, n_candidates) -> int: one rival holding every remaining slot. A freed position (reclaim or republish) has no row and is correctly counted as available again. - Zero for a branch the registry does not know: there is nothing left to - claim on a branch that no longer exists. + Zero for a branch that is not open. A registry id is not evidence the + branch exists: delete_branch deliberately keeps the append-only + branches row (branch_id must stay stable across a re-promotion) while + dropping every candidate_claims row, so a finished branch would + otherwise count zero claims and report all n_candidates as claimable — + describing completed work as untouched. """ branch_id = self._intern_branch(branch_key) if branch_id is None: return 0 + live = self._conn.execute( + "SELECT 1 FROM active_branches WHERE branch_id = ?", + (branch_id,)).fetchone() + if live is None: + return 0 taken = self._conn.execute( "SELECT COUNT(*) FROM candidate_claims WHERE branch_id = ?", (branch_id,)).fetchone()[0] diff --git a/erd_swarm.py b/erd_swarm.py index 867c436a..54ed21a3 100644 --- a/erd_swarm.py +++ b/erd_swarm.py @@ -2166,7 +2166,12 @@ def _await_rival_finalize(self, branch_key, words, n_words, n_candidates): silently at full CPU. Heartbeat first (this wait must stay visible and must not get this worker's parent claims reclaimed), then reopen the row once it has been 'finalized' past FINALIZE_TAKEOVER_SECONDS - and complete the finalize from its intact claims and meta.""" + and complete the finalize from its intact claims and meta. + + Returns True when it polled and False when it took the finalize over: + only the poll is time this worker spent waiting, and a caller + attributing blocked time must not charge the takeover, which is work. + """ self._cur_candidate = None # coordinating, no candidate in flight self._heartbeat(branch_key, n_words, None, None, None, None, force=True) @@ -2175,8 +2180,9 @@ def _await_rival_finalize(self, branch_key, words, n_words, n_candidates): logger.warning('%s reopened branch (%d words): finalizer died ' 'mid-finalize', self.name, n_words) self.maybe_finalize(branch_key, words, n_candidates) - return + return False self._idle_wait(0.05) + return True # -- recursive cooperative solving -------------------------------------- @@ -2608,8 +2614,17 @@ def cooperative_solve(self, words, budget, ceiling=float('inf')): elif self.queue.branch_done_candidates(branch_key) >= self.n_candidates: if not self.maybe_finalize(branch_key, words, self.n_candidates): - self._await_rival_finalize(branch_key, words, n_words, - self.n_candidates) + # Every candidate is done and a rival holds the + # finalize: the wait-for-finalize case these columns + # exist to name, so it is sampled here rather than left + # NULL with its poll unattributed. + self._record_first_block(wait, branch_key) + blocked_at = time.perf_counter() + if self._await_rival_finalize(branch_key, words, + n_words, + self.n_candidates): + wait.note_blocked( + int((time.perf_counter() - blocked_at) * 1000)) elif self._help_recursion_depth >= MAX_HELP_RECURSION_DEPTH: # _help_other_branch would refuse to scan at all here (see # its own docstring): False from it below would mean diff --git a/tests/test_erd_swarm_unit.py b/tests/test_erd_swarm_unit.py index aad8234a..1236611a 100644 --- a/tests/test_erd_swarm_unit.py +++ b/tests/test_erd_swarm_unit.py @@ -6124,6 +6124,84 @@ def test_a_freed_position_counts_as_claimable_again(self): w.queue.branch_unclaimed_candidates(key, w.n_candidates), w.n_candidates) + def test_a_deleted_branch_has_no_unclaimed_slots(self): + """delete_branch keeps the registry row; that is not existence. + + branch_id stays stable across a re-promotion, so the append-only + branches row survives while every candidate_claims row is dropped. + Counting from the id alone would see zero claims on a finished branch + and report all n_candidates as claimable — completed work described as + untouched, and a cap refusal where there is nothing left to claim. + """ + w = self._worker() + key = ScoreCache.encode_subset(BRANCH[:3]) + w.queue.create_branch(key, 3, w.n_candidates, budget=5) + self.assertEqual( + w.queue.branch_unclaimed_candidates(key, w.n_candidates), + w.n_candidates) + w.queue.delete_branch(key) + self.assertIsNone(w.queue.get_branch(key)) + self.assertEqual( + w.queue.branch_unclaimed_candidates(key, w.n_candidates), 0) + + def test_losing_the_finalize_race_is_attributed_as_blocked(self): + """The wait-for-finalize case must not report zero blocked time. + + Every candidate is done and a rival holds the finalize, so the worker + polls in _await_rival_finalize. That poll is the commonest wait these + columns exist to name; leaving it outside the accounting would put the + sleep in episode_millis while blocked_millis read zero and both + first-block counters stayed NULL. + """ + w = self._worker() + words = BRANCH[:3] + key = ScoreCache.encode_subset(words) + w.queue.create_branch(key, len(words), w.n_candidates, + budget=ROOT_BUDGET) + w.queue.mark_claims_done(key, list(range(w.n_candidates))) + + def _lose_then_stop(*_a, **_kw): + w.request_stop() + return False + + with mock.patch.object(w, "maybe_finalize", side_effect=_lose_then_stop), \ + mock.patch.object(w, "_await_rival_finalize", + return_value=True) as await_rival: + w.cooperative_solve(words, ROOT_BUDGET) + + await_rival.assert_called() + rows = self._wait_rows() + self.assertEqual(len(rows), 1) + self.assertEqual(rows[0]["unclaimed_at_first_block"], 0) + self.assertIsNotNone(rows[0]["holders_at_first_block"]) + + def test_a_finalize_takeover_is_not_charged_as_blocked_time(self): + """Taking the finalize over is work, not waiting. + + _await_rival_finalize returns False when it reopened a dead + finalizer's row and completed the finalize itself; charging that span + to blocked_millis would inflate the stuck figure with real work. + """ + w = self._worker() + words = BRANCH[:3] + key = ScoreCache.encode_subset(words) + w.queue.create_branch(key, len(words), w.n_candidates, + budget=ROOT_BUDGET) + w.queue.mark_claims_done(key, list(range(w.n_candidates))) + + def _lose_then_stop(*_a, **_kw): + w.request_stop() + return False + + with mock.patch.object(w, "maybe_finalize", side_effect=_lose_then_stop), \ + mock.patch.object(w, "_await_rival_finalize", + return_value=False): + w.cooperative_solve(words, ROOT_BUDGET) + + rows = self._wait_rows() + self.assertEqual(len(rows), 1) + self.assertEqual(rows[0]["blocked_millis"], 0) + def test_unclaimed_of_an_unknown_branch_is_zero(self): """A branch the registry never saw has nothing left to claim.""" w = self._worker() From 988bd150f099c6932633e7bd489757018a21a301 Mon Sep 17 00:00:00 2001 From: Sean Ahern Date: Thu, 17 Sep 2026 17:05:29 -0400 Subject: [PATCH 4/7] Read branch liveness and its claim count in one snapshot MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The liveness check added for the deleted-branch case was a separate autocommit select from the claim count, so the two described different database states. A rival finalizing between them leaves the branch present and its claim rows already deleted, which reports a completed branch as fully claimable — the same telemetry error, reached by timing rather than by order. Both now come from one statement, which sees one snapshot. The race is not reproducible from a sequential fixture, so the guard is structural: test_unclaimed_reads_liveness_and_claims_in_one_snapshot traces the connection and requires a single SELECT naming both tables. It fails when the query is split back in two. Co-Authored-By: Claude Opus 5 Claude-Session: https://claude.ai/code/session_019rawcKNbdKoVpUW1Avr4sS --- erd_queue.py | 17 ++++++++++------- tests/test_erd_swarm_unit.py | 25 +++++++++++++++++++++++++ 2 files changed, 35 insertions(+), 7 deletions(-) diff --git a/erd_queue.py b/erd_queue.py index 24886374..216bb87a 100644 --- a/erd_queue.py +++ b/erd_queue.py @@ -3867,14 +3867,17 @@ def branch_unclaimed_candidates(self, branch_key, n_candidates) -> int: branch_id = self._intern_branch(branch_key) if branch_id is None: return 0 - live = self._conn.execute( - "SELECT 1 FROM active_branches WHERE branch_id = ?", - (branch_id,)).fetchone() - if live is None: + # Liveness and the claim count come from ONE statement, so they + # describe one snapshot. Read separately they do not: these are + # autocommit selects, and a rival finalizing between them leaves the + # branch present and its claim rows already deleted — which is the + # deleted-branch error again, reached by a race instead of by order. + live, taken = self._conn.execute( + "SELECT EXISTS(SELECT 1 FROM active_branches WHERE branch_id = ?)," + " (SELECT COUNT(*) FROM candidate_claims WHERE branch_id = ?)", + (branch_id, branch_id)).fetchone() + if not live: return 0 - taken = self._conn.execute( - "SELECT COUNT(*) FROM candidate_claims WHERE branch_id = ?", - (branch_id,)).fetchone()[0] return max(0, n_candidates - taken) def branch_bulk_done_candidates(self, branch_key) -> int: diff --git a/tests/test_erd_swarm_unit.py b/tests/test_erd_swarm_unit.py index 1236611a..96600291 100644 --- a/tests/test_erd_swarm_unit.py +++ b/tests/test_erd_swarm_unit.py @@ -6144,6 +6144,31 @@ def test_a_deleted_branch_has_no_unclaimed_slots(self): self.assertEqual( w.queue.branch_unclaimed_candidates(key, w.n_candidates), 0) + def test_unclaimed_reads_liveness_and_claims_in_one_snapshot(self): + """Two autocommit selects are two snapshots, and the gap is a race. + + A rival finalizing between them leaves the branch present and its + claim rows already deleted, which reports a completed branch as fully + claimable — the deleted-branch error reached by timing rather than by + order, and not reproducible from a sequential fixture. Pinned + structurally instead: the count must come from one statement. + """ + w = self._worker() + key = ScoreCache.encode_subset(BRANCH[:3]) + w.queue.create_branch(key, 3, w.n_candidates, budget=5) + w.queue.branch_unclaimed_candidates(key, w.n_candidates) # warm the id + + statements = [] + w.queue._conn.set_trace_callback(statements.append) + try: + w.queue.branch_unclaimed_candidates(key, w.n_candidates) + finally: + w.queue._conn.set_trace_callback(None) + selects = [q for q in statements if q.lstrip().upper().startswith("SELECT")] + self.assertEqual(len(selects), 1, selects) + self.assertIn("active_branches", selects[0]) + self.assertIn("candidate_claims", selects[0]) + def test_losing_the_finalize_race_is_attributed_as_blocked(self): """The wait-for-finalize case must not report zero blocked time. From d62ba308f672a2acbe2f07924477ceb162571d5e Mon Sep 17 00:00:00 2001 From: Sean Ahern Date: Thu, 17 Sep 2026 17:10:58 -0400 Subject: [PATCH 5/7] Read both first-block counters from one snapshot MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Making liveness and the claim count atomic inside one method did not make the two diagnostic columns atomic with each other. _record_first_block still called branch_claim_holders and branch_unclaimed_candidates separately, so a rival could finish, reclaim or finalize between them and leave a holder count from one state beside an unclaimed count from another. That pair is read against itself — holders at MAX_WORKERS_PER_BRANCH versus nothing left to claim — so a mismatched pair does not blur the answer, it inverts it. branch_block_snapshot replaces both methods and returns liveness, the holder count and the unclaimed count from a single statement, which sees a single snapshot. The race is not reachable from a sequential fixture, so the guard stays structural: the trace test now covers all three reads and fails when the query is split. Co-Authored-By: Claude Opus 5 Claude-Session: https://claude.ai/code/session_019rawcKNbdKoVpUW1Avr4sS --- AGENTS.md | 13 ++++-- erd_queue.py | 78 +++++++++++++++--------------------- erd_swarm.py | 6 +-- tests/test_erd_swarm_unit.py | 64 ++++++++++++++--------------- 4 files changed, 76 insertions(+), 85 deletions(-) diff --git a/AGENTS.md b/AGENTS.md index 3c7dbabf..ef96afd2 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -100,10 +100,15 @@ to hand out and the wait is for its finalize, which no cap change reaches. Those are different problems with different fixes, and `idle_millis` cannot tell them apart. -They are sampled once, at the first blocked iteration, because that path -already runs every 50 ms on a starving worker and two more queries per turn -would be paid by the branch everyone is waiting for. An episode that never -reached the loop writes no row. +They are sampled once, at the first blocked iteration and before the poll that +follows it, because that path already runs every 50 ms on a starving worker and +another query per turn would be paid by the branch everyone is waiting for. +`branch_block_snapshot` returns **both counters and branch liveness from one +statement**: they are read against each other, so a holder count from one state +beside an unclaimed count from another misclassifies the block, and a branch +present at a liveness lookup whose claim rows are gone by the count reports +finished work as fully claimable. An episode that never reached the loop writes +no row. ### Priority ladders, and the fan-out they prevent diff --git a/erd_queue.py b/erd_queue.py index 216bb87a..880af699 100644 --- a/erd_queue.py +++ b/erd_queue.py @@ -3826,59 +3826,47 @@ def branch_done_candidates(self, branch_key) -> int: "SELECT COUNT(*) FROM candidate_claims WHERE branch_id = ? AND done = 1", (branch_id,)).fetchone()[0] - def branch_claim_holders(self, branch_key, exclude_worker_id=None) -> int: - """Workers other than exclude_worker_id holding unfinished claims on - one branch. - - The single-branch form of claim_holders_by_branch, for a caller that - wants occupancy for the branch it is already looking at rather than a - map of every branch. Occupancy is unfinished claims, never heartbeats: - the claim row is written by the transaction that hands out the bundle, - so a branch reads as taken the instant it is taken. - """ - branch_id = self._intern_branch(branch_key) - if branch_id is None: - return 0 - if exclude_worker_id is None: - return self._conn.execute( - "SELECT COUNT(DISTINCT claimed_by) FROM candidate_claims " - "WHERE branch_id = ? AND done = 0", (branch_id,)).fetchone()[0] - return self._conn.execute( - "SELECT COUNT(DISTINCT claimed_by) FROM candidate_claims " - "WHERE branch_id = ? AND done = 0 AND claimed_by IS NOT ?", - (branch_id, exclude_worker_id)).fetchone()[0] - - def branch_unclaimed_candidates(self, branch_key, n_candidates) -> int: - """Candidate slots on this branch that no claim row covers. - - NOT `n_candidates - done`: a candidate held in an unfinished claim is - taken, not available, so counting it as unclaimed would report work the - packer cannot hand out — which is the normal finalize-wait shape, with - one rival holding every remaining slot. A freed position (reclaim or - republish) has no row and is correctly counted as available again. - - Zero for a branch that is not open. A registry id is not evidence the - branch exists: delete_branch deliberately keeps the append-only + def branch_block_snapshot(self, branch_key, n_candidates, + exclude_worker_id=None): + """(holders, unclaimed) for one branch, from ONE statement. + + holders is the workers other than exclude_worker_id holding unfinished + claims; unclaimed is the candidate slots no claim row covers. A caller + that reads them separately gets a holder count from one state and an + unclaimed count from another — a rival can finish, reclaim or finalize + in between — and these two are compared against each other to decide + whether the cap or a pending finalize caused a block, so a mismatched + pair misclassifies it. One statement sees one snapshot. + + Occupancy is unfinished claims, never heartbeats: the claim row is + written by the transaction that hands out the bundle, so a branch reads + as taken the instant it is taken. + + unclaimed is NOT `n_candidates - done`: a candidate held in an + unfinished claim is taken, not available, and one rival holding every + remaining slot is the ordinary finalize-wait shape. A position freed + by a reclaim or republish has no row and counts as claimable again. + + (0, 0) for a branch that is not open. A registry id is not evidence + the branch exists: delete_branch deliberately keeps the append-only branches row (branch_id must stay stable across a re-promotion) while dropping every candidate_claims row, so a finished branch would - otherwise count zero claims and report all n_candidates as claimable — - describing completed work as untouched. + otherwise report no holders and every slot claimable. """ branch_id = self._intern_branch(branch_key) if branch_id is None: - return 0 - # Liveness and the claim count come from ONE statement, so they - # describe one snapshot. Read separately they do not: these are - # autocommit selects, and a rival finalizing between them leaves the - # branch present and its claim rows already deleted — which is the - # deleted-branch error again, reached by a race instead of by order. - live, taken = self._conn.execute( + return (0, 0) + live, holders, taken = self._conn.execute( "SELECT EXISTS(SELECT 1 FROM active_branches WHERE branch_id = ?)," + " (SELECT COUNT(DISTINCT claimed_by) FROM candidate_claims" + " WHERE branch_id = ? AND done = 0" + " AND (? IS NULL OR claimed_by IS NOT ?))," " (SELECT COUNT(*) FROM candidate_claims WHERE branch_id = ?)", - (branch_id, branch_id)).fetchone() + (branch_id, branch_id, exclude_worker_id, exclude_worker_id, + branch_id)).fetchone() if not live: - return 0 - return max(0, n_candidates - taken) + return (0, 0) + return (holders, max(0, n_candidates - taken)) def branch_bulk_done_candidates(self, branch_key) -> int: """Return the legacy combined count completed by ERD pruning.""" diff --git a/erd_swarm.py b/erd_swarm.py index 54ed21a3..ab9876b5 100644 --- a/erd_swarm.py +++ b/erd_swarm.py @@ -2417,10 +2417,8 @@ def _record_first_block(self, wait, branch_key): """ if wait.holders_at_first_block is not None: return - holders = self.queue.branch_claim_holders( - branch_key, exclude_worker_id=self.name) - unclaimed = self.queue.branch_unclaimed_candidates( - branch_key, self.n_candidates) + holders, unclaimed = self.queue.branch_block_snapshot( + branch_key, self.n_candidates, exclude_worker_id=self.name) wait.note_first_block(holders, unclaimed) def _record_dependency_wait(self, wait): diff --git a/tests/test_erd_swarm_unit.py b/tests/test_erd_swarm_unit.py index 96600291..7ebf8bcc 100644 --- a/tests/test_erd_swarm_unit.py +++ b/tests/test_erd_swarm_unit.py @@ -6118,10 +6118,10 @@ def test_a_freed_position_counts_as_claimable_again(self): self.assertIsNotNone( self._rival_claims(w, key, count_cap=w.n_candidates)) self.assertEqual( - w.queue.branch_unclaimed_candidates(key, w.n_candidates), 0) + w.queue.branch_block_snapshot(key, w.n_candidates)[1], 0) w.queue.reclaim_claims_of_worker("rival") self.assertEqual( - w.queue.branch_unclaimed_candidates(key, w.n_candidates), + w.queue.branch_block_snapshot(key, w.n_candidates)[1], w.n_candidates) def test_a_deleted_branch_has_no_unclaimed_slots(self): @@ -6137,31 +6137,35 @@ def test_a_deleted_branch_has_no_unclaimed_slots(self): key = ScoreCache.encode_subset(BRANCH[:3]) w.queue.create_branch(key, 3, w.n_candidates, budget=5) self.assertEqual( - w.queue.branch_unclaimed_candidates(key, w.n_candidates), - w.n_candidates) + w.queue.branch_block_snapshot(key, w.n_candidates), + (0, w.n_candidates)) w.queue.delete_branch(key) self.assertIsNone(w.queue.get_branch(key)) self.assertEqual( - w.queue.branch_unclaimed_candidates(key, w.n_candidates), 0) - - def test_unclaimed_reads_liveness_and_claims_in_one_snapshot(self): - """Two autocommit selects are two snapshots, and the gap is a race. - - A rival finalizing between them leaves the branch present and its - claim rows already deleted, which reports a completed branch as fully - claimable — the deleted-branch error reached by timing rather than by - order, and not reproducible from a sequential fixture. Pinned - structurally instead: the count must come from one statement. + w.queue.branch_block_snapshot(key, w.n_candidates), (0, 0)) + + def test_both_counters_and_liveness_come_from_one_snapshot(self): + """Separate autocommit selects are separate snapshots. + + The two counters are compared against each other — holders at the cap + versus nothing left to claim — so reading them apart lets a rival + finish, reclaim or finalize in between and yields a holder count from + one state beside an unclaimed count from another, misclassifying the + block. Liveness has the same problem: a branch present at the lookup + whose claim rows are gone by the count reports a finished branch as + fully claimable. None of that is reproducible from a sequential + fixture, so the guard is structural: one statement, one snapshot. """ w = self._worker() key = ScoreCache.encode_subset(BRANCH[:3]) w.queue.create_branch(key, 3, w.n_candidates, budget=5) - w.queue.branch_unclaimed_candidates(key, w.n_candidates) # warm the id + w.queue.branch_block_snapshot(key, w.n_candidates) # warm the id statements = [] w.queue._conn.set_trace_callback(statements.append) try: - w.queue.branch_unclaimed_candidates(key, w.n_candidates) + w.queue.branch_block_snapshot(key, w.n_candidates, + exclude_worker_id="worker-0") finally: w.queue._conn.set_trace_callback(None) selects = [q for q in statements if q.lstrip().upper().startswith("SELECT")] @@ -6227,30 +6231,26 @@ def _lose_then_stop(*_a, **_kw): self.assertEqual(len(rows), 1) self.assertEqual(rows[0]["blocked_millis"], 0) - def test_unclaimed_of_an_unknown_branch_is_zero(self): - """A branch the registry never saw has nothing left to claim.""" - w = self._worker() - self.assertEqual( - w.queue.branch_unclaimed_candidates( - ScoreCache.encode_subset(BRANCH), w.n_candidates), 0) - - def test_branch_claim_holders_of_an_unknown_branch_is_zero(self): - """A branch the queue has never interned holds nobody.""" + def test_an_unknown_branch_snapshots_as_empty(self): + """A branch the registry never saw holds nobody and owes nothing.""" w = self._worker() self.assertEqual( - w.queue.branch_claim_holders(ScoreCache.encode_subset(BRANCH)), 0) + w.queue.branch_block_snapshot( + ScoreCache.encode_subset(BRANCH), w.n_candidates), (0, 0)) - def test_branch_claim_holders_agrees_with_the_map(self): - """The single-branch count must not drift from the map scheduling uses.""" + def test_snapshot_holders_agree_with_the_map(self): + """The snapshot's holder count must not drift from the map scheduling uses.""" w = self._worker() key = ScoreCache.encode_subset(BRANCH[:3]) w.queue.create_branch(key, 3, w.n_candidates, budget=5) - self.assertEqual(w.queue.branch_claim_holders(key), 0) + self.assertEqual(w.queue.branch_block_snapshot(key, w.n_candidates)[0], 0) self.assertIsNotNone(self._rival_claims(w, key)) self.assertEqual( - w.queue.branch_claim_holders(key, exclude_worker_id="rival"), 0) - self.assertEqual(w.queue.branch_claim_holders(key), 1) + w.queue.branch_block_snapshot( + key, w.n_candidates, exclude_worker_id="rival")[0], 0) + self.assertEqual(w.queue.branch_block_snapshot(key, w.n_candidates)[0], 1) self.assertEqual( - w.queue.branch_claim_holders(key, exclude_worker_id=w.name), + w.queue.branch_block_snapshot( + key, w.n_candidates, exclude_worker_id=w.name)[0], w.queue.claim_holders_by_branch( exclude_worker_id=w.name).get(bytes(key), 0)) From 738c196ef5938aa1fc7f5563723cf62ffe188577 Mon Sep 17 00:00:00 2001 From: Sean Ahern Date: Thu, 17 Sep 2026 23:22:05 -0400 Subject: [PATCH 6/7] Report why a claim declined instead of sampling for it afterwards MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Both first-block columns reconstructed, from outside and after the fact, a decision claim_next_bundle already makes atomically inside its own BEGIN IMMEDIATE. That reconstruction was the source of four separate defects: a sample taken after the sleep it described, a deleted branch read as fully claimable, liveness and a claim count taken in two statements, and two counters taken in two calls. Each fix closed one instance of the same mistake. 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. It now records which of its five return-None sites it took, and last_claim_decline reports it. dependency_wait's holders_at_first_block and unclaimed_at_first_block are replaced by four counters — blocks_worker_cap, blocks_no_candidates, blocks_awaiting_finalize, blocks_help_capped — one per cause, counted over the whole episode rather than sampled once. Two come from the claim transaction, two from 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, and the race class is gone by construction rather than by vigilance. branch_block_snapshot and _record_first_block are deleted along with their query: the attribution is now free. test_a_served_claim_clears_the_previous_decline asserts the clear after a real decline; on a fresh connection the field is already None and removing the clear proved nothing. Co-Authored-By: Claude Opus 5 Claude-Session: https://claude.ai/code/session_019rawcKNbdKoVpUW1Avr4sS --- AGENTS.md | 38 ++-- erd_queue.py | 110 ++++++------ erd_swarm.py | 75 +++----- tests/test_erd_swarm_unit.py | 329 ++++++++++++++--------------------- 4 files changed, 236 insertions(+), 316 deletions(-) diff --git a/AGENTS.md b/AGENTS.md index ef96afd2..02ea8787 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -91,24 +91,32 @@ 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 two columns that matter are `holders_at_first_block` and -`unclaimed_at_first_block`.** A worker in that loop tries the branch alone, -then anywhere else, then a pair onto the branch, and only then sleeps. Holders -at `MAX_WORKERS_PER_BRANCH` means the cap refused the pair — raising it would -admit this worker. Zero unclaimed candidates means the branch had nothing left -to hand out and the wait is for its finalize, which no cap change reaches. +**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 | + Those are different problems with different fixes, and `idle_millis` cannot tell them apart. -They are sampled once, at the first blocked iteration and before the poll that -follows it, because that path already runs every 50 ms on a starving worker and -another query per turn would be paid by the branch everyone is waiting for. -`branch_block_snapshot` returns **both counters and branch liveness from one -statement**: they are read against each other, so a holder count from one state -beside an unclaimed count from another misclassifies the block, and a branch -present at a liveness lookup whose claim rows are gone by the count reports -finished work as fully claimable. An episode that never reached the loop writes -no row. +**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 diff --git a/erd_queue.py b/erd_queue.py index 880af699..84f35564 100644 --- a/erd_queue.py +++ b/erd_queue.py @@ -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). @@ -865,10 +875,14 @@ def guess_depth_from_spine(spine) -> int: -- 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 holders/unclaimed at first block say whether the pair was refused --- because the branch was already at MAX_WORKERS_PER_BRANCH or because there --- was nothing left on it to claim. Those are different problems with --- different fixes, and idle_millis cannot tell them apart. +-- 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, @@ -888,10 +902,12 @@ def guess_depth_from_spine(spine) -> int: bundles_claimed INTEGER NOT NULL, pair_attempts INTEGER NOT NULL, pair_successes INTEGER NOT NULL, - -- Sampled once, at the first blocked iteration: re-reading them every - -- 50 ms poll would put a query on the one path that is already starving. - holders_at_first_block INTEGER, - unclaimed_at_first_block INTEGER, + -- Why this episode's sleeps happened, one count per cause. Raising + -- MAX_WORKERS_PER_BRANCH reaches blocks_worker_cap and nothing else. + 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, help_depth INTEGER, -- How the wait ended: 'solved', 'loss', 'cut', 'deleted', 'cancelled'. outcome TEXT, @@ -984,6 +1000,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 @@ -3193,10 +3215,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: @@ -3218,11 +3242,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 @@ -3325,8 +3351,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 @@ -3826,47 +3854,16 @@ def branch_done_candidates(self, branch_key) -> int: "SELECT COUNT(*) FROM candidate_claims WHERE branch_id = ? AND done = 1", (branch_id,)).fetchone()[0] - def branch_block_snapshot(self, branch_key, n_candidates, - exclude_worker_id=None): - """(holders, unclaimed) for one branch, from ONE statement. - - holders is the workers other than exclude_worker_id holding unfinished - claims; unclaimed is the candidate slots no claim row covers. A caller - that reads them separately gets a holder count from one state and an - unclaimed count from another — a rival can finish, reclaim or finalize - in between — and these two are compared against each other to decide - whether the cap or a pending finalize caused a block, so a mismatched - pair misclassifies it. One statement sees one snapshot. - - Occupancy is unfinished claims, never heartbeats: the claim row is - written by the transaction that hands out the bundle, so a branch reads - as taken the instant it is taken. - - unclaimed is NOT `n_candidates - done`: a candidate held in an - unfinished claim is taken, not available, and one rival holding every - remaining slot is the ordinary finalize-wait shape. A position freed - by a reclaim or republish has no row and counts as claimable again. - - (0, 0) for a branch that is not open. A registry id is not evidence - the branch exists: delete_branch deliberately keeps the append-only - branches row (branch_id must stay stable across a re-promotion) while - dropping every candidate_claims row, so a finished branch would - otherwise report no holders and every slot claimable. + 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. """ - branch_id = self._intern_branch(branch_key) - if branch_id is None: - return (0, 0) - live, holders, taken = self._conn.execute( - "SELECT EXISTS(SELECT 1 FROM active_branches WHERE branch_id = ?)," - " (SELECT COUNT(DISTINCT claimed_by) FROM candidate_claims" - " WHERE branch_id = ? AND done = 0" - " AND (? IS NULL OR claimed_by IS NOT ?))," - " (SELECT COUNT(*) FROM candidate_claims WHERE branch_id = ?)", - (branch_id, branch_id, exclude_worker_id, exclude_worker_id, - branch_id)).fetchone() - if not live: - return (0, 0) - return (holders, max(0, n_candidates - taken)) + return self._last_claim_decline def branch_bulk_done_candidates(self, branch_key) -> int: """Return the legacy combined count completed by ERD pruning.""" @@ -6896,16 +6893,15 @@ 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, - holders_at_first_block=None, - unclaimed_at_first_block=None, + blocks_worker_cap=0, blocks_no_candidates=0, + blocks_awaiting_finalize=0, blocks_help_capped=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 what its alternatives were when it - first stalled. A caller that never blocked has nothing to attribute - and should not write a row. + 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(""" @@ -6913,13 +6909,15 @@ def add_dependency_wait(self, worker_id, spine, n_words, budget, (worker_id, spine, n_words, budget, episode_millis, blocked_millis, iterations, empty_scans, helped_scans, bundles_claimed, pair_attempts, pair_successes, - holders_at_first_block, unclaimed_at_first_block, + blocks_worker_cap, blocks_no_candidates, + blocks_awaiting_finalize, blocks_help_capped, help_depth, outcome, epoch, recorded_at) - VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?) + VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?) """, (worker_id, spine, n_words, budget, episode_millis, blocked_millis, iterations, empty_scans, helped_scans, bundles_claimed, pair_attempts, pair_successes, - holders_at_first_block, unclaimed_at_first_block, + blocks_worker_cap, blocks_no_candidates, + blocks_awaiting_finalize, blocks_help_capped, help_depth, outcome, self.epoch, now)) def add_cut_reuse_miss(self, branch_key, n_words, budget, wanted_ceiling, diff --git a/erd_swarm.py b/erd_swarm.py index ab9876b5..5e12918e 100644 --- a/erd_swarm.py +++ b/erd_swarm.py @@ -58,7 +58,8 @@ DEFAULT_SMALL_COUNT, DEFAULT_COUNT_CAP, DEFAULT_REPUBLISH_LIMIT, CLAIM_RETRY, SCHEDULING_ROLE_PREFERRED, SCHEDULING_ROLE_FALLBACK, - SCHEDULING_ROLE_DIRECT) + SCHEDULING_ROLE_DIRECT, + CLAIM_DECLINE_WORKER_CAP, CLAIM_DECLINE_NO_CANDIDATES) from wordle_ui import fmt_pattern from runtime_paths import ( @@ -266,6 +267,12 @@ def descend(self, branch_key, spine): # when it has no unoccupied branch to take instead. MAX_HELP_RECURSION_DEPTH = 4 +# Why a dependency wait slept, for the two sleeps no claim transaction decides. +# The other two reasons come from ERDQueue's CLAIM_DECLINE_* constants, reported +# by the transaction that refused the pair. +BLOCK_AWAITING_FINALIZE = "awaiting_finalize" # every candidate done, rival finalizing +BLOCK_HELP_CAPPED = "help_capped" # recursion cap forbade scanning + # _promote_opener_work sentinel: a pending branch was promoted (or found # already solved and marked done) but no bundle is claimable from it right # now — distinct from None, which means the opener has no pending branch @@ -597,17 +604,17 @@ class _DependencyWait: `blocked_millis` is the stuck state alone, entered only after the worker has failed to claim the branch by itself, failed to find work anywhere - else, and failed to pair onto the branch. `holders_at_first_block` and - `unclaimed_at_first_block` are sampled once at that moment, which is what - separates "the cap refused the pair" from "the branch had nothing left to - claim". + else, and failed to pair onto the branch. The `blocks_*` counters say why + each of those sleeps happened, and every one of them is reported by the + code that made the decision — the claim transaction for the two claim + outcomes, the branch condition for the other two — so none of them can + disagree with the moment it describes. """ __slots__ = ("spine", "n_words", "budget", "help_depth", "outcome", "_started", "iterations", "blocked_millis", "empty_scans", "helped_scans", "bundles_claimed", "pair_attempts", - "pair_successes", "holders_at_first_block", - "unclaimed_at_first_block") + "pair_successes", "blocks") def __init__(self, spine, n_words, budget, help_depth): self.spine = spine @@ -623,26 +630,16 @@ def __init__(self, spine, n_words, budget, help_depth): self.bundles_claimed = 0 self.pair_attempts = 0 self.pair_successes = 0 - self.holders_at_first_block = None - self.unclaimed_at_first_block = None + self.blocks = collections.Counter() @property def episode_millis(self): return int((time.perf_counter() - self._started) * 1000) - def note_first_block(self, holders, unclaimed): - """Record the alternatives once, on the first blocked iteration. - - Sampled once rather than per poll: the blocked path already runs every - 50 ms on a starving worker, and two more queries per turn of it would - be paid by the branch everyone is waiting for. - """ - if self.holders_at_first_block is None: - self.holders_at_first_block = holders - self.unclaimed_at_first_block = unclaimed - - def note_blocked(self, millis): + def note_blocked(self, reason, millis): + """Charge one sleep, against the reason the deciding code gave for it.""" self.blocked_millis += millis + self.blocks[reason] += 1 def should_record(self): """True when the episode reached the wait loop at all. @@ -2402,25 +2399,6 @@ def _read_satisfying_cut(self, branch_key, budget, ceiling, n_words): return (OVER_ERD_LIMIT, cut_bound, None, cut_tainted) return None - def _record_first_block(self, wait, branch_key): - """Sample the worker's alternatives the first time it has none. - - Reached only after the sole-worker claim, the scan for work elsewhere, - and (on the uncapped path) the pair attempt have all failed, so the two - counts answer why: holders at MAX_WORKERS_PER_BRANCH means the cap - refused the pair, while zero unclaimed candidates means the branch had - nothing left to hand out and the wait is for its finalize. - - Taken before the poll that follows it on every path, so the snapshot - describes the state that caused the block rather than whatever the - branch became while this worker slept. - """ - if wait.holders_at_first_block is not None: - return - holders, unclaimed = self.queue.branch_block_snapshot( - branch_key, self.n_candidates, exclude_worker_id=self.name) - wait.note_first_block(holders, unclaimed) - def _record_dependency_wait(self, wait): """Persist one wait episode, unless it never reached the wait loop.""" if not wait.should_record(): @@ -2430,8 +2408,10 @@ def _record_dependency_wait(self, wait): wait.episode_millis, wait.blocked_millis, wait.iterations, wait.empty_scans, wait.helped_scans, wait.bundles_claimed, wait.pair_attempts, wait.pair_successes, - holders_at_first_block=wait.holders_at_first_block, - unclaimed_at_first_block=wait.unclaimed_at_first_block, + blocks_worker_cap=wait.blocks[CLAIM_DECLINE_WORKER_CAP], + blocks_no_candidates=wait.blocks[CLAIM_DECLINE_NO_CANDIDATES], + blocks_awaiting_finalize=wait.blocks[BLOCK_AWAITING_FINALIZE], + blocks_help_capped=wait.blocks[BLOCK_HELP_CAPPED], help_depth=wait.help_depth, outcome=wait.outcome) def cooperative_solve(self, words, budget, ceiling=float('inf')): @@ -2616,12 +2596,12 @@ def cooperative_solve(self, words, budget, ceiling=float('inf')): # finalize: the wait-for-finalize case these columns # exist to name, so it is sampled here rather than left # NULL with its poll unattributed. - self._record_first_block(wait, branch_key) blocked_at = time.perf_counter() if self._await_rival_finalize(branch_key, words, n_words, self.n_candidates): wait.note_blocked( + BLOCK_AWAITING_FINALIZE, int((time.perf_counter() - blocked_at) * 1000)) elif self._help_recursion_depth >= MAX_HELP_RECURSION_DEPTH: # _help_other_branch would refuse to scan at all here (see @@ -2631,10 +2611,10 @@ def cooperative_solve(self, words, budget, ceiling=float('inf')): # _help_other_branch's capped-depth contract already # promises its callers. self._cur_candidate = None - self._record_first_block(wait, branch_key) blocked_at = time.perf_counter() self._idle_wait(0.05) wait.note_blocked( + BLOCK_HELP_CAPPED, int((time.perf_counter() - blocked_at) * 1000)) else: # No bundle: every candidate is claimed, or another worker @@ -2683,11 +2663,14 @@ def cooperative_solve(self, words, budget, ceiling=float('inf')): else: # Nothing claimable anywhere, no pair available on # the dependency: the stuck state idle_millis - # totals without naming. - self._record_first_block(wait, branch_key) + # totals without naming. The claim transaction + # that just refused the pair reports why, so this + # is the decision itself rather than a re-read of + # the state it was made against. blocked_at = time.perf_counter() self._idle_wait(0.05) # let claims land wait.note_blocked( + self.queue.last_claim_decline(), int((time.perf_counter() - blocked_at) * 1000)) if self.cancel(): # pragma: no cover diff --git a/tests/test_erd_swarm_unit.py b/tests/test_erd_swarm_unit.py index 7ebf8bcc..058ee4e7 100644 --- a/tests/test_erd_swarm_unit.py +++ b/tests/test_erd_swarm_unit.py @@ -5920,14 +5920,84 @@ def test_pairing_still_joins_a_branch_from_the_reused_list(self): self.assertEqual(self._claimed_key(work), joinable_key) +class TestClaimDeclineReason(unittest.TestCase): + """The claim transaction reports why it handed out no bundle. + + The cap counts occupancy against the very claims the transaction is about + to create, so only that transaction can say whether the cap refused a + claim. A caller asking afterwards is asking a system that has moved. + """ + + def setUp(self): + self._tmp = tempfile.TemporaryDirectory() + self.addCleanup(self._tmp.cleanup) + self.q = erd_queue.ERDQueue(os.path.join(self._tmp.name, "q.sqlite3")) + self.addCleanup(self.q.close) + self.key = ScoreCache.encode_subset(BRANCH) + self.n = len(CANDIDATES) + self.order = list(range(self.n)) + self.bounds = [0.0] * self.n + + def _claim(self, worker, **kw): + return self.q.claim_next_bundle( + self.key, worker, self.n, self.order, self.bounds, **kw) + + def test_a_served_claim_clears_the_previous_decline(self): + """A stale reason would be charged to whatever blocks next. + + The clear has to be observed after a decline, not on a fresh + connection where the field is already None and nothing is proven. + """ + self.q.create_branch(self.key, len(BRANCH), self.n, budget=5) + self.assertIsNone(self._claim("w0", small_count=1, count_cap=1, + expected_budget=4)) + self.assertEqual(self.q.last_claim_decline(), + erd_queue.CLAIM_DECLINE_BUDGET_MISMATCH) + self.assertIsNotNone(self._claim("w0", small_count=1, count_cap=1)) + self.assertIsNone(self.q.last_claim_decline()) + + def test_the_worker_cap_is_named_as_the_reason(self): + """The signal the cap question turns on.""" + self.q.create_branch(self.key, len(BRANCH), self.n, budget=5) + self.assertIsNotNone(self._claim("rival", small_count=1, count_cap=1)) + self.assertIsNone(self._claim("w0", small_count=1, count_cap=1, + max_other_workers=0)) + self.assertEqual(self.q.last_claim_decline(), + erd_queue.CLAIM_DECLINE_WORKER_CAP) + + def test_an_exhausted_branch_is_not_reported_as_capped(self): + """Nothing left to hand out is a different problem from the cap.""" + self.q.create_branch(self.key, len(BRANCH), self.n, budget=5) + self.assertIsNotNone( + self._claim("rival", small_count=self.n, count_cap=self.n)) + self.assertIsNone(self._claim("w0", small_count=1, count_cap=1, + max_other_workers=9)) + self.assertEqual(self.q.last_claim_decline(), + erd_queue.CLAIM_DECLINE_NO_CANDIDATES) + + def test_a_missing_branch_is_named(self): + self.q.create_branch(self.key, len(BRANCH), self.n, budget=5) + self.q.delete_branch(self.key) + self.assertIsNone(self._claim("w0", small_count=1, count_cap=1)) + self.assertEqual(self.q.last_claim_decline(), + erd_queue.CLAIM_DECLINE_BRANCH_GONE) + + def test_a_budget_mismatch_is_named(self): + self.q.create_branch(self.key, len(BRANCH), self.n, budget=5) + self.assertIsNone(self._claim("w0", small_count=1, count_cap=1, + expected_budget=4)) + self.assertEqual(self.q.last_claim_decline(), + erd_queue.CLAIM_DECLINE_BUDGET_MISMATCH) + + class TestDependencyWaitAttribution(unittest.TestCase): """cooperative_solve records what it waited on, so idle time has a subject. `claim_telemetry.idle_millis` is a residual and totals waiting without naming it. These pin the attribution: which dependency, how much of the - episode was spent stuck rather than working or helping, and — the field the - worker-cap question turns on — whether anything was claimable when the - worker first had nothing to do. + episode was spent stuck rather than working or helping, and why each sleep + happened. Every reason is reported by the code that decided it, so no + counter can disagree with the moment it describes. """ def setUp(self): @@ -5972,6 +6042,21 @@ def _wait_rows(self): finally: conn.close() + def _stall_once(self, w, key, words, cap): + """Drive one blocked iteration of the wait loop, then stop the worker. + + _help_other_branch is forced to report an empty scan so the loop + reaches its pair attempt, which is the decision under test. + """ + with mock.patch.object(erd_swarm, "MAX_WORKERS_PER_BRANCH", cap), \ + mock.patch.object(w, "_help_other_branch", return_value=False), \ + mock.patch.object(w, "_idle_wait", + side_effect=lambda *_a: w.request_stop()): + w.cooperative_solve(words, ROOT_BUDGET) + rows = self._wait_rows() + self.assertEqual(len(rows), 1) + return rows[0] + def test_cache_hit_writes_no_wait_row(self): """An episode answered before the wait loop never opens an episode.""" words = BRANCH[:3] @@ -5984,13 +6069,7 @@ def test_cache_hit_writes_no_wait_row(self): self.assertEqual(self._wait_rows(), []) def test_an_episode_that_never_iterates_writes_no_row(self): - """A worker asked to stop before the first iteration attributed nothing. - - The wait object exists by then — it is built immediately before the - loop — so only should_record keeps this row out. Distinct from the - cache-hit path above, which returns before the episode is opened at - all. - """ + """A worker asked to stop before the first iteration attributed nothing.""" w = self._worker() w.request_stop() self.assertEqual(w.cooperative_solve(BRANCH[:3], ROOT_BUDGET), @@ -6004,12 +6083,10 @@ def test_a_solved_dependency_names_the_branch_it_waited_on(self): self.assertEqual(status, SOLVED) rows = self._wait_rows() self.assertEqual(len(rows), 1) - row = rows[0] - self.assertEqual(row["n_words"], len(words)) - self.assertEqual(row["budget"], ROOT_BUDGET) - self.assertEqual(row["outcome"], "solved") - self.assertGreaterEqual(row["iterations"], 1) - self.assertGreaterEqual(row["bundles_claimed"], 1) + self.assertEqual(rows[0]["n_words"], len(words)) + self.assertEqual(rows[0]["budget"], ROOT_BUDGET) + self.assertEqual(rows[0]["outcome"], "solved") + self.assertGreaterEqual(rows[0]["bundles_claimed"], 1) def test_blocked_time_never_exceeds_the_episode(self): """blocked_millis is a part of the episode, not a separate clock.""" @@ -6019,52 +6096,26 @@ def test_blocked_time_never_exceeds_the_episode(self): self.assertLessEqual(row["blocked_millis"], row["episode_millis"]) self.assertGreaterEqual(row["blocked_millis"], 0) - def test_first_block_is_sampled_once(self): - """Re-sampling would put two queries on every 50 ms poll of a stall.""" - w = self._worker() - wait = erd_swarm._DependencyWait("SPINE -----", 3, 5, 0) - key = ScoreCache.encode_subset(BRANCH[:3]) - w.queue.create_branch(key, 3, w.n_candidates, budget=5) - w._record_first_block(wait, key) - first = (wait.holders_at_first_block, wait.unclaimed_at_first_block) - w.queue.mark_claims_done(key, list(range(w.n_candidates))) - w._record_first_block(wait, key) - self.assertEqual( - (wait.holders_at_first_block, wait.unclaimed_at_first_block), first) - - def test_cap_refusal_is_distinguishable_from_an_exhausted_branch(self): - """The two reasons a pair fails must not read alike. + def test_a_pair_refused_by_the_cap_is_recorded_as_worker_cap(self): + """The measurement the worker-cap decision rests on. - A branch another worker holds with candidates left is a cap refusal — - raising MAX_WORKERS_PER_BRANCH would admit this worker. A branch with - nothing left to claim is waiting on its own finalize, which no cap - change reaches. idle_millis cannot tell these apart; these columns - must. + A rival holds the branch and candidates remain, so the pair attempt is + refused by the cap alone. Raising MAX_WORKERS_PER_BRANCH reaches this + block and no other, which is why it must not be conflated with a + branch that has simply run out of work. """ w = self._worker() - key = ScoreCache.encode_subset(BRANCH[:3]) - w.queue.create_branch(key, 3, w.n_candidates, budget=5) - self.assertIsNotNone(self._rival_claims(w, key)) - - occupied = erd_swarm._DependencyWait("SPINE -----", 3, 5, 0) - w._record_first_block(occupied, key) - self.assertGreaterEqual(occupied.holders_at_first_block, 1) - self.assertGreater(occupied.unclaimed_at_first_block, 0) + words = BRANCH[:3] + key = ScoreCache.encode_subset(words) + w.queue.create_branch(key, len(words), w.n_candidates, + budget=ROOT_BUDGET) + self.assertIsNotNone(self._rival_claims(w, key, count_cap=1)) + row = self._stall_once(w, key, words, cap=1) + self.assertEqual(row["blocks_worker_cap"], 1) + self.assertEqual(row["blocks_no_candidates"], 0) - w.queue.mark_claims_done(key, list(range(w.n_candidates))) - exhausted = erd_swarm._DependencyWait("SPINE -----", 3, 5, 0) - w._record_first_block(exhausted, key) - self.assertEqual(exhausted.unclaimed_at_first_block, 0) - - def test_the_capped_path_snapshots_before_it_sleeps(self): - """The snapshot must describe the state that caused the block. - - On the recursion-capped path the worker polls instead of scanning. If - the sample were taken after the 50 ms sleep it would observe whatever - the branch became while this worker slept — a holder that finished - reads as zero holders — and both diagnostic columns would describe a - moment that never blocked anything. - """ + def test_a_pair_refused_for_want_of_candidates_is_not_worker_cap(self): + """An exhausted branch must not read as a cap refusal.""" w = self._worker() words = BRANCH[:3] key = ScoreCache.encode_subset(words) @@ -6072,116 +6123,29 @@ def test_the_capped_path_snapshots_before_it_sleeps(self): budget=ROOT_BUDGET) self.assertIsNotNone( self._rival_claims(w, key, count_cap=w.n_candidates)) + row = self._stall_once(w, key, words, cap=9) + self.assertEqual(row["blocks_no_candidates"], 1) + self.assertEqual(row["blocks_worker_cap"], 0) - def _sleep_that_changes_the_branch(_seconds): - # The rival finishes mid-sleep: holders drop to zero, so a sample - # taken afterwards would report nobody was on the branch. - w.queue.mark_claims_done(key, list(range(w.n_candidates))) - w.request_stop() - + def test_the_recursion_cap_names_its_own_block(self): + """The capped path polls without ever asking the claim transaction.""" + w = self._worker() + words = BRANCH[:3] + key = ScoreCache.encode_subset(words) + w.queue.create_branch(key, len(words), w.n_candidates, + budget=ROOT_BUDGET) + self.assertIsNotNone(self._rival_claims(w, key, count_cap=1)) with mock.patch.object(erd_swarm, "MAX_HELP_RECURSION_DEPTH", 0), \ mock.patch.object(w, "_idle_wait", - side_effect=_sleep_that_changes_the_branch): + side_effect=lambda *_a: w.request_stop()): w.cooperative_solve(words, ROOT_BUDGET) - rows = self._wait_rows() self.assertEqual(len(rows), 1) - self.assertEqual(rows[0]["holders_at_first_block"], 1) - self.assertEqual(rows[0]["unclaimed_at_first_block"], 0) - - def test_candidates_held_in_flight_do_not_read_as_unclaimed(self): - """A slot inside an unfinished claim is taken, not available. - - This is the ordinary finalize-wait shape: one rival holds every - remaining candidate, so the pair attempt finds no bundle. Counting - those slots as unclaimed would report a cap refusal on a branch that - simply has nothing left to hand out, which is precisely the distinction - these columns exist to draw. - """ - w = self._worker() - key = ScoreCache.encode_subset(BRANCH[:3]) - w.queue.create_branch(key, 3, w.n_candidates, budget=5) - self.assertIsNotNone( - self._rival_claims(w, key, count_cap=w.n_candidates)) - self.assertEqual(w.queue.branch_done_candidates(key), 0) - - wait = erd_swarm._DependencyWait("SPINE -----", 3, 5, 0) - w._record_first_block(wait, key) - self.assertEqual(wait.unclaimed_at_first_block, 0) - self.assertGreaterEqual(wait.holders_at_first_block, 1) - - def test_a_freed_position_counts_as_claimable_again(self): - """A reclaimed or republished slot has no row and is available.""" - w = self._worker() - key = ScoreCache.encode_subset(BRANCH[:3]) - w.queue.create_branch(key, 3, w.n_candidates, budget=5) - self.assertIsNotNone( - self._rival_claims(w, key, count_cap=w.n_candidates)) - self.assertEqual( - w.queue.branch_block_snapshot(key, w.n_candidates)[1], 0) - w.queue.reclaim_claims_of_worker("rival") - self.assertEqual( - w.queue.branch_block_snapshot(key, w.n_candidates)[1], - w.n_candidates) - - def test_a_deleted_branch_has_no_unclaimed_slots(self): - """delete_branch keeps the registry row; that is not existence. - - branch_id stays stable across a re-promotion, so the append-only - branches row survives while every candidate_claims row is dropped. - Counting from the id alone would see zero claims on a finished branch - and report all n_candidates as claimable — completed work described as - untouched, and a cap refusal where there is nothing left to claim. - """ - w = self._worker() - key = ScoreCache.encode_subset(BRANCH[:3]) - w.queue.create_branch(key, 3, w.n_candidates, budget=5) - self.assertEqual( - w.queue.branch_block_snapshot(key, w.n_candidates), - (0, w.n_candidates)) - w.queue.delete_branch(key) - self.assertIsNone(w.queue.get_branch(key)) - self.assertEqual( - w.queue.branch_block_snapshot(key, w.n_candidates), (0, 0)) - - def test_both_counters_and_liveness_come_from_one_snapshot(self): - """Separate autocommit selects are separate snapshots. - - The two counters are compared against each other — holders at the cap - versus nothing left to claim — so reading them apart lets a rival - finish, reclaim or finalize in between and yields a holder count from - one state beside an unclaimed count from another, misclassifying the - block. Liveness has the same problem: a branch present at the lookup - whose claim rows are gone by the count reports a finished branch as - fully claimable. None of that is reproducible from a sequential - fixture, so the guard is structural: one statement, one snapshot. - """ - w = self._worker() - key = ScoreCache.encode_subset(BRANCH[:3]) - w.queue.create_branch(key, 3, w.n_candidates, budget=5) - w.queue.branch_block_snapshot(key, w.n_candidates) # warm the id - - statements = [] - w.queue._conn.set_trace_callback(statements.append) - try: - w.queue.branch_block_snapshot(key, w.n_candidates, - exclude_worker_id="worker-0") - finally: - w.queue._conn.set_trace_callback(None) - selects = [q for q in statements if q.lstrip().upper().startswith("SELECT")] - self.assertEqual(len(selects), 1, selects) - self.assertIn("active_branches", selects[0]) - self.assertIn("candidate_claims", selects[0]) + self.assertEqual(rows[0]["blocks_help_capped"], 1) + self.assertEqual(rows[0]["blocks_worker_cap"], 0) def test_losing_the_finalize_race_is_attributed_as_blocked(self): - """The wait-for-finalize case must not report zero blocked time. - - Every candidate is done and a rival holds the finalize, so the worker - polls in _await_rival_finalize. That poll is the commonest wait these - columns exist to name; leaving it outside the accounting would put the - sleep in episode_millis while blocked_millis read zero and both - first-block counters stayed NULL. - """ + """The commonest wait: every candidate done, a rival finalizing.""" w = self._worker() words = BRANCH[:3] key = ScoreCache.encode_subset(words) @@ -6194,23 +6158,15 @@ def _lose_then_stop(*_a, **_kw): return False with mock.patch.object(w, "maybe_finalize", side_effect=_lose_then_stop), \ - mock.patch.object(w, "_await_rival_finalize", - return_value=True) as await_rival: + mock.patch.object(w, "_await_rival_finalize", return_value=True): w.cooperative_solve(words, ROOT_BUDGET) - - await_rival.assert_called() rows = self._wait_rows() self.assertEqual(len(rows), 1) - self.assertEqual(rows[0]["unclaimed_at_first_block"], 0) - self.assertIsNotNone(rows[0]["holders_at_first_block"]) + self.assertEqual(rows[0]["blocks_awaiting_finalize"], 1) + self.assertEqual(rows[0]["blocks_worker_cap"], 0) def test_a_finalize_takeover_is_not_charged_as_blocked_time(self): - """Taking the finalize over is work, not waiting. - - _await_rival_finalize returns False when it reopened a dead - finalizer's row and completed the finalize itself; charging that span - to blocked_millis would inflate the stuck figure with real work. - """ + """Taking the finalize over is work, not waiting.""" w = self._worker() words = BRANCH[:3] key = ScoreCache.encode_subset(words) @@ -6223,34 +6179,9 @@ def _lose_then_stop(*_a, **_kw): return False with mock.patch.object(w, "maybe_finalize", side_effect=_lose_then_stop), \ - mock.patch.object(w, "_await_rival_finalize", - return_value=False): + mock.patch.object(w, "_await_rival_finalize", return_value=False): w.cooperative_solve(words, ROOT_BUDGET) - rows = self._wait_rows() self.assertEqual(len(rows), 1) self.assertEqual(rows[0]["blocked_millis"], 0) - - def test_an_unknown_branch_snapshots_as_empty(self): - """A branch the registry never saw holds nobody and owes nothing.""" - w = self._worker() - self.assertEqual( - w.queue.branch_block_snapshot( - ScoreCache.encode_subset(BRANCH), w.n_candidates), (0, 0)) - - def test_snapshot_holders_agree_with_the_map(self): - """The snapshot's holder count must not drift from the map scheduling uses.""" - w = self._worker() - key = ScoreCache.encode_subset(BRANCH[:3]) - w.queue.create_branch(key, 3, w.n_candidates, budget=5) - self.assertEqual(w.queue.branch_block_snapshot(key, w.n_candidates)[0], 0) - self.assertIsNotNone(self._rival_claims(w, key)) - self.assertEqual( - w.queue.branch_block_snapshot( - key, w.n_candidates, exclude_worker_id="rival")[0], 0) - self.assertEqual(w.queue.branch_block_snapshot(key, w.n_candidates)[0], 1) - self.assertEqual( - w.queue.branch_block_snapshot( - key, w.n_candidates, exclude_worker_id=w.name)[0], - w.queue.claim_holders_by_branch( - exclude_worker_id=w.name).get(bytes(key), 0)) + self.assertEqual(rows[0]["blocks_awaiting_finalize"], 0) From 19ae050ae62673a8335957e2e1c70fbe3e7c8f8b Mon Sep 17 00:00:00 2001 From: Sean Ahern Date: Thu, 17 Sep 2026 23:29:19 -0400 Subject: [PATCH 7/7] Make the block counters a partition of the episode's sleeps MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit wait.blocks counted whatever reason came back, but only four of them had a column. A claim transaction can also report branch_gone, budget_mismatch or owner_mismatch when the dependency changes identity under the waiter, and an exhausted retry loop reports nothing at all. Those sleeps were counted and then dropped, leaving blocked_millis holding time no counter accounted for — which is the defect idle_millis has and this table exists to avoid repeating. blocks_other carries them, so every sleep increments exactly one counter and the five sum to the episode's sleep count. A counter that is a partition can be audited; a counter that is a selection cannot. The re-loop-without-sleeping alternative is deliberately not taken here: it changes what a blocked worker does, which is a scheduling change and wants its own justification rather than arriving inside instrumentation. _assert_blocks checks every column on every test rather than the one under test, because the first version of these tests passed with other_blocks summing named reasons too — the only reasons present were unnamed ones, so a double-count was indistinguishable. Co-Authored-By: Claude Opus 5 Claude-Session: https://claude.ai/code/session_019rawcKNbdKoVpUW1Avr4sS --- AGENTS.md | 8 ++++ erd_queue.py | 17 ++++++-- erd_swarm.py | 19 +++++++++ tests/test_erd_swarm_unit.py | 81 +++++++++++++++++++++++++++++++----- 4 files changed, 110 insertions(+), 15 deletions(-) diff --git a/AGENTS.md b/AGENTS.md index 02ea8787..f6e7582d 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -101,10 +101,18 @@ only then sleeps: | `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 diff --git a/erd_queue.py b/erd_queue.py index 84f35564..92c2c93d 100644 --- a/erd_queue.py +++ b/erd_queue.py @@ -904,10 +904,18 @@ def guess_depth_from_spine(spine) -> int: 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, @@ -6895,7 +6903,7 @@ def add_dependency_wait(self, worker_id, spine, n_words, budget, pair_attempts, pair_successes, blocks_worker_cap=0, blocks_no_candidates=0, blocks_awaiting_finalize=0, blocks_help_capped=0, - help_depth=None, outcome=None): + 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 @@ -6910,14 +6918,15 @@ def add_dependency_wait(self, worker_id, spine, n_words, budget, 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_awaiting_finalize, blocks_help_capped, blocks_other, help_depth, outcome, epoch, recorded_at) - VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?) + 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_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, diff --git a/erd_swarm.py b/erd_swarm.py index 5e12918e..84dcecd1 100644 --- a/erd_swarm.py +++ b/erd_swarm.py @@ -636,11 +636,29 @@ def __init__(self, spine, n_words, budget, help_depth): def episode_millis(self): return int((time.perf_counter() - self._started) * 1000) + #: The reasons with a column of their own. Anything else a claim + #: transaction can report — a dependency whose identity changed under the + #: waiter, or a retry loop that ran out — is summed into blocks_other. + NAMED_BLOCK_REASONS = frozenset(( + CLAIM_DECLINE_WORKER_CAP, CLAIM_DECLINE_NO_CANDIDATES, + BLOCK_AWAITING_FINALIZE, BLOCK_HELP_CAPPED)) + def note_blocked(self, reason, millis): """Charge one sleep, against the reason the deciding code gave for it.""" self.blocked_millis += millis self.blocks[reason] += 1 + def other_blocks(self): + """Sleeps whose reason has no column of its own. + + Keeps the counters a partition of the episode's sleeps: without it a + reason nobody anticipated would leave blocked_millis holding time no + counter accounts for, which is the shape of defect idle_millis already + has and this table exists to avoid repeating. + """ + return sum(count for reason, count in self.blocks.items() + if reason not in self.NAMED_BLOCK_REASONS) + def should_record(self): """True when the episode reached the wait loop at all. @@ -2412,6 +2430,7 @@ def _record_dependency_wait(self, wait): blocks_no_candidates=wait.blocks[CLAIM_DECLINE_NO_CANDIDATES], blocks_awaiting_finalize=wait.blocks[BLOCK_AWAITING_FINALIZE], blocks_help_capped=wait.blocks[BLOCK_HELP_CAPPED], + blocks_other=wait.other_blocks(), help_depth=wait.help_depth, outcome=wait.outcome) def cooperative_solve(self, words, budget, ceiling=float('inf')): diff --git a/tests/test_erd_swarm_unit.py b/tests/test_erd_swarm_unit.py index 058ee4e7..8abf0220 100644 --- a/tests/test_erd_swarm_unit.py +++ b/tests/test_erd_swarm_unit.py @@ -6042,6 +6042,20 @@ def _wait_rows(self): finally: conn.close() + _BLOCK_COLUMNS = ("blocks_worker_cap", "blocks_no_candidates", + "blocks_awaiting_finalize", "blocks_help_capped", + "blocks_other") + + def _assert_blocks(self, row, **expected): + """Assert the whole partition, not just the counter under test. + + Every column is checked, unnamed ones defaulting to zero, so a change + that double-counts a sleep into two columns fails here rather than + passing because the test only looked at the one it cared about. + """ + for column in self._BLOCK_COLUMNS: + self.assertEqual(row[column], expected.get(column, 0), column) + def _stall_once(self, w, key, words, cap): """Drive one blocked iteration of the wait loop, then stop the worker. @@ -6110,9 +6124,8 @@ def test_a_pair_refused_by_the_cap_is_recorded_as_worker_cap(self): w.queue.create_branch(key, len(words), w.n_candidates, budget=ROOT_BUDGET) self.assertIsNotNone(self._rival_claims(w, key, count_cap=1)) - row = self._stall_once(w, key, words, cap=1) - self.assertEqual(row["blocks_worker_cap"], 1) - self.assertEqual(row["blocks_no_candidates"], 0) + self._assert_blocks(self._stall_once(w, key, words, cap=1), + blocks_worker_cap=1) def test_a_pair_refused_for_want_of_candidates_is_not_worker_cap(self): """An exhausted branch must not read as a cap refusal.""" @@ -6123,9 +6136,8 @@ def test_a_pair_refused_for_want_of_candidates_is_not_worker_cap(self): budget=ROOT_BUDGET) self.assertIsNotNone( self._rival_claims(w, key, count_cap=w.n_candidates)) - row = self._stall_once(w, key, words, cap=9) - self.assertEqual(row["blocks_no_candidates"], 1) - self.assertEqual(row["blocks_worker_cap"], 0) + self._assert_blocks(self._stall_once(w, key, words, cap=9), + blocks_no_candidates=1) def test_the_recursion_cap_names_its_own_block(self): """The capped path polls without ever asking the claim transaction.""" @@ -6141,8 +6153,7 @@ def test_the_recursion_cap_names_its_own_block(self): w.cooperative_solve(words, ROOT_BUDGET) rows = self._wait_rows() self.assertEqual(len(rows), 1) - self.assertEqual(rows[0]["blocks_help_capped"], 1) - self.assertEqual(rows[0]["blocks_worker_cap"], 0) + self._assert_blocks(rows[0], blocks_help_capped=1) def test_losing_the_finalize_race_is_attributed_as_blocked(self): """The commonest wait: every candidate done, a rival finalizing.""" @@ -6162,8 +6173,56 @@ def _lose_then_stop(*_a, **_kw): w.cooperative_solve(words, ROOT_BUDGET) rows = self._wait_rows() self.assertEqual(len(rows), 1) - self.assertEqual(rows[0]["blocks_awaiting_finalize"], 1) - self.assertEqual(rows[0]["blocks_worker_cap"], 0) + self._assert_blocks(rows[0], blocks_awaiting_finalize=1) + + def test_an_unnamed_reason_still_lands_in_a_counter(self): + """The counters must partition the episode's sleeps. + + A claim transaction can decline for reasons with no column of their own + — a dependency whose identity changed under the waiter, or a retry loop + that ran out and reported nothing. Dropping those would leave + blocked_millis holding time no counter accounts for, which is exactly + the defect idle_millis has and this table exists to avoid repeating. + """ + w = self._worker() + words = BRANCH[:3] + key = ScoreCache.encode_subset(words) + w.queue.create_branch(key, len(words), w.n_candidates, + budget=ROOT_BUDGET) + self.assertIsNotNone(self._rival_claims(w, key, count_cap=1)) + + with mock.patch.object(erd_swarm, "MAX_WORKERS_PER_BRANCH", 1), \ + mock.patch.object(w, "_help_other_branch", return_value=False), \ + mock.patch.object(w.queue, "last_claim_decline", + return_value=erd_queue.CLAIM_DECLINE_BRANCH_GONE), \ + mock.patch.object(w, "_idle_wait", + side_effect=lambda *_a: w.request_stop()): + w.cooperative_solve(words, ROOT_BUDGET) + + rows = self._wait_rows() + self.assertEqual(len(rows), 1) + self._assert_blocks(rows[0], blocks_other=1) + + def test_a_retry_exhausted_claim_is_still_counted(self): + """A claim that reports no reason at all must not vanish from the sum.""" + w = self._worker() + words = BRANCH[:3] + key = ScoreCache.encode_subset(words) + w.queue.create_branch(key, len(words), w.n_candidates, + budget=ROOT_BUDGET) + self.assertIsNotNone(self._rival_claims(w, key, count_cap=1)) + + with mock.patch.object(erd_swarm, "MAX_WORKERS_PER_BRANCH", 1), \ + mock.patch.object(w, "_help_other_branch", return_value=False), \ + mock.patch.object(w.queue, "last_claim_decline", + return_value=None), \ + mock.patch.object(w, "_idle_wait", + side_effect=lambda *_a: w.request_stop()): + w.cooperative_solve(words, ROOT_BUDGET) + + rows = self._wait_rows() + self.assertEqual(len(rows), 1) + self._assert_blocks(rows[0], blocks_other=1) def test_a_finalize_takeover_is_not_charged_as_blocked_time(self): """Taking the finalize over is work, not waiting.""" @@ -6184,4 +6243,4 @@ def _lose_then_stop(*_a, **_kw): rows = self._wait_rows() self.assertEqual(len(rows), 1) self.assertEqual(rows[0]["blocked_millis"], 0) - self.assertEqual(rows[0]["blocks_awaiting_finalize"], 0) + self._assert_blocks(rows[0])