diff --git a/AGENTS.md b/AGENTS.md index 3a1604d5..f6e7582d 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -79,6 +79,53 @@ missing in CI. The `scaling` job deliberately keeps its explicit `coverage numpy`: every assertion it makes is a wall-clock ratio, and a first-call JIT compile would land inside the measurement. +### idle_millis is a residual; dependency_wait is its attribution + +`idle_millis` is not measured. It is whatever remains of a coordination window +after the four measured phases, so it totals waiting without saying what was +waited on — and it is the swarm's largest single cost (61.5% of all worker time +on epoch 21, against 3.4% for the scan and 2.3% for lock waits). + +`telemetry.dependency_wait` carries one row per `cooperative_solve` episode that +reached the wait loop: the dependency's spine, size and budget, how the episode +divided between claiming that branch, helping elsewhere, and being stuck, and +how it ended. + +**The `blocks_*` counters say why each sleep happened.** A worker in that loop +tries the branch alone, then anywhere else, then a pair onto the branch, and +only then sleeps: + +| counter | cause | reached by | +|---|---|---| +| `blocks_worker_cap` | the branch is at `MAX_WORKERS_PER_BRANCH` | raising the cap | +| `blocks_no_candidates` | nothing left to hand out | nothing the cap can do | +| `blocks_awaiting_finalize` | every candidate done, a rival finalizing | nothing the cap can do | +| `blocks_help_capped` | `MAX_HELP_RECURSION_DEPTH` forbade scanning | raising that cap | +| `blocks_other` | the dependency changed identity, or a retry ran out | nothing; it keeps the sum honest | + +Those are different problems with different fixes, and `idle_millis` cannot +tell them apart. + +**Every sleep increments exactly one counter, so the five sum to the episode's +sleep count.** That is what `blocks_other` is for — a claim transaction can +decline for reasons with no column of their own, and dropping those would leave +`blocked_millis` holding time no counter accounts for, which is precisely the +defect `idle_millis` has. A counter that is a partition can be audited; a +counter that is a selection cannot. + +**Every reason is reported by the code that decided it, never sampled +afterwards.** The two claim outcomes come from `ERDQueue.last_claim_decline()`, +set inside the claim transaction at each of its `return None` sites; the other +two are the branch condition the worker is standing in. Nothing on this path +reads live state, so no counter can disagree with the moment it describes. + +That property is the design, not an optimization. **Occupancy is counted +against the claims the claim transaction is about to create, so only that +transaction can say whether the cap refused a claim** — anything reconstructing +it afterwards is reading a system that has already moved, and a holder count +from one instant beside an availability count from another does not blur the +answer, it inverts it. An episode that never reached the loop writes no row. + ### Priority ladders, and the fan-out they prevent **Openers tied at one priority all become eligible at once, and the swarm @@ -1166,7 +1213,8 @@ successful merge unless `--keep-source` is given. Swarm telemetry lives in a **separate** Linux-only file, `runtime/erd_queue_telemetry.sqlite3` (`_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 b943a314..92c2c93d 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). @@ -851,6 +861,67 @@ def guess_depth_from_spine(spine) -> int: epoch INTEGER NOT NULL DEFAULT 0, recorded_at INTEGER NOT NULL ); + +-- One row per cooperative_solve that had to wait on a dependency: the worker +-- needed a sub-branch it could not simply claim, and spent time in the wait +-- loop before the branch produced an answer. +-- +-- claim_telemetry's idle_millis is a RESIDUAL — the part of a coordination +-- window left after the four measured phases — so it says how much waiting +-- happened but nothing about what was waited on. This table is the +-- attribution: which branch, for how long, and what the worker's alternatives +-- were at the moment it stalled. +-- +-- The distinction the columns exist to draw: a worker in this loop first tries +-- to claim the branch alone, then to help anywhere else, then to pair onto the +-- branch, and only then sleeps. blocked_millis covers that last state alone, +-- and the blocks_* counters say WHY each sleep happened. Those are different +-- problems with different fixes, and idle_millis cannot tell them apart. +-- +-- Every reason is reported by the code that decided it rather than sampled +-- afterwards: blocks_worker_cap and blocks_no_candidates come from the claim +-- transaction that enforces the cap, the other two from the branch condition +-- the worker is standing in. Nothing here reads live state, so no counter can +-- disagree with the moment it describes. +CREATE TABLE IF NOT EXISTS telemetry.dependency_wait ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + worker_id TEXT, + spine TEXT, -- the dependency, not the waiter + n_words INTEGER, + budget INTEGER, + -- Whole time inside the wait loop, and the part of it spent asleep with no + -- claim held and nothing claimable: episode >= blocked always. + episode_millis INTEGER NOT NULL, + blocked_millis INTEGER NOT NULL, + iterations INTEGER NOT NULL, + -- _help_other_branch outcomes: a completed scan that found nothing free or + -- promotable anywhere, versus one that sent this worker off to real work. + empty_scans INTEGER NOT NULL, + helped_scans INTEGER NOT NULL, + -- Bundles this worker claimed on the dependency itself, sole or paired. + bundles_claimed INTEGER NOT NULL, + pair_attempts INTEGER NOT NULL, + pair_successes INTEGER NOT NULL, + -- Why this episode's sleeps happened, one count per cause. Raising + -- MAX_WORKERS_PER_BRANCH reaches blocks_worker_cap and nothing else. + -- Every sleep increments exactly one of these, so they sum to the + -- episode's sleep count and blocked_millis is never time no counter + -- accounts for. blocks_other carries the reasons that are neither + -- actionable nor common -- a dependency that changed identity under the + -- waiter (branch_gone, budget_mismatch, owner_mismatch) and a claim whose + -- retry loop ran out -- which exist to keep that sum honest rather than to + -- be read on their own. + blocks_worker_cap INTEGER NOT NULL DEFAULT 0, + blocks_no_candidates INTEGER NOT NULL DEFAULT 0, + blocks_awaiting_finalize INTEGER NOT NULL DEFAULT 0, + blocks_help_capped INTEGER NOT NULL DEFAULT 0, + blocks_other INTEGER NOT NULL DEFAULT 0, + help_depth INTEGER, + -- How the wait ended: 'solved', 'loss', 'cut', 'deleted', 'cancelled'. + outcome TEXT, + epoch INTEGER NOT NULL DEFAULT 0, + recorded_at INTEGER NOT NULL +); """ # Every table _TELEMETRY_SCHEMA_SQL creates, used to detect a queue file that @@ -858,7 +929,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", ) @@ -937,6 +1008,12 @@ def __init__(self, db_path: str, timeout: float = 30.0, self._last_claim_retries = 0 self._last_claim_transaction_millis = 0 self._last_claim_commit_millis = 0 + # Why the last claim_next_bundle declined, or None when it handed out a + # bundle. Decided inside the claim transaction, which is the only + # place that can decide it: occupancy is counted against claims the + # same transaction is about to create. A caller reconstructing this + # afterwards is reading a system that has already moved. + self._last_claim_decline = None # Monotonic per-connection counter for bundle_id generation: paired # with worker_id and this process's pid, it is unique without a # timestamp-collision risk (two bundles claimed by the same worker @@ -3146,10 +3223,12 @@ def claim_next_bundle(self, branch_key, worker_id, n_candidates, "FROM active_branches " "WHERE branch_id = ?", (branch_id,)).fetchone() if br is None or br["status"] != "open": + self._last_claim_decline = CLAIM_DECLINE_BRANCH_GONE self._commit_claim_transaction(_txn_t0) return None if (expected_budget is not None and br["budget"] is not None and br["budget"] != expected_budget): + self._last_claim_decline = CLAIM_DECLINE_BUDGET_MISMATCH self._commit_claim_transaction(_txn_t0) return None if expected_opener_work_id is None: @@ -3171,11 +3250,13 @@ def claim_next_bundle(self, branch_key, worker_id, n_candidates, """, (branch_id, expected_opener_work_id, expected_opener_priority)).fetchone() is not None if not owner_matches: + self._last_claim_decline = CLAIM_DECLINE_OWNER_MISMATCH self._commit_claim_transaction(_txn_t0) return None if (max_other_workers is not None and self._other_claim_holders(branch_id, worker_id) > max_other_workers): + self._last_claim_decline = CLAIM_DECLINE_WORKER_CAP self._commit_claim_transaction(_txn_t0) return None # The branch ceiling is a bound like any achieved best: candidates @@ -3278,8 +3359,10 @@ def claim_next_bundle(self, branch_key, worker_id, n_candidates, packed = self._pack_recorded_holes( branch_id, bound, cost_lower_bound, survivor_limit) if not packed: + self._last_claim_decline = CLAIM_DECLINE_NO_CANDIDATES self._commit_claim_transaction(_txn_t0) return None + self._last_claim_decline = None bundle = [idx for idx, _position in packed] bundle_id = f"{worker_id}:{self._pid}:{self._bundle_seq}" self._bundle_seq += 1 @@ -3779,6 +3862,17 @@ def branch_done_candidates(self, branch_key) -> int: "SELECT COUNT(*) FROM candidate_claims WHERE branch_id = ? AND done = 1", (branch_id,)).fetchone()[0] + def last_claim_decline(self): + """Why the last claim_next_bundle on this connection declined, or None + when it handed out a bundle. + + Read straight after a claim that returned None. This is an observation + of the decision, not a reconstruction of it: the transaction that + enforces the cap is the same one that reports it, so the answer cannot + disagree with the state the cap was enforced against. + """ + return self._last_claim_decline + def branch_bulk_done_candidates(self, branch_key) -> int: """Return the legacy combined count completed by ERD pruning.""" branch_id = self._intern_branch(branch_key) @@ -6803,6 +6897,38 @@ def add_branch_finalize_log(self, branch_key, spine, n_words, budget, first_best_at, nodes_at_first_best, now)) + def add_dependency_wait(self, worker_id, spine, n_words, budget, + episode_millis, blocked_millis, iterations, + empty_scans, helped_scans, bundles_claimed, + pair_attempts, pair_successes, + blocks_worker_cap=0, blocks_no_candidates=0, + blocks_awaiting_finalize=0, blocks_help_capped=0, + blocks_other=0, help_depth=None, outcome=None): + """Record one cooperative_solve wait episode. + + Attributes what claim_telemetry's idle_millis can only total: which + dependency the worker waited on, how much of the episode it was asleep + rather than helping or claiming, and why each sleep happened. A caller + that never blocked has nothing to attribute and should not write a row. + """ + now = int(time.time()) + self._conn.execute(""" + INSERT INTO telemetry.dependency_wait + (worker_id, spine, n_words, budget, episode_millis, + blocked_millis, iterations, empty_scans, helped_scans, + bundles_claimed, pair_attempts, pair_successes, + blocks_worker_cap, blocks_no_candidates, + blocks_awaiting_finalize, blocks_help_capped, blocks_other, + help_depth, outcome, epoch, recorded_at) + VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, + ?) + """, (worker_id, spine, n_words, budget, episode_millis, + blocked_millis, iterations, empty_scans, helped_scans, + bundles_claimed, pair_attempts, pair_successes, + blocks_worker_cap, blocks_no_candidates, + blocks_awaiting_finalize, blocks_help_capped, blocks_other, + help_depth, outcome, self.epoch, now)) + def add_cut_reuse_miss(self, branch_key, n_words, budget, wanted_ceiling, available_bound, available_budget): """Log a re-solve forced by a cut: a consumer needed this branch but diff --git a/erd_swarm.py b/erd_swarm.py index 0a2c98bb..84dcecd1 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 @@ -585,6 +592,83 @@ 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. 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", "blocks") + + 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.blocks = collections.Counter() + + @property + 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. + + 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.""" @@ -2097,7 +2181,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) @@ -2106,8 +2195,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 -------------------------------------- @@ -2327,6 +2417,22 @@ def _read_satisfying_cut(self, branch_key, budget, ceiling, n_words): return (OVER_ERD_LIMIT, cut_bound, None, cut_tainted) return None + 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, + 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], + blocks_other=wait.other_blocks(), + 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 +2560,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 +2589,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,13 +2605,23 @@ 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: 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. + 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 # its own docstring): False from it below would mean @@ -2504,7 +2630,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( + BLOCK_HELP_CAPPED, + int((time.perf_counter() - blocked_at) * 1000)) else: # No bundle: every candidate is claimed, or another worker # holds the branch. Try free or promotable work first — @@ -2519,7 +2649,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 +2668,36 @@ 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. 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 + 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..8abf0220 100644 --- a/tests/test_erd_swarm_unit.py +++ b/tests/test_erd_swarm_unit.py @@ -5918,3 +5918,329 @@ 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 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 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): + 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, 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=count_cap, count_cap=count_cap) + + 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() + + _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. + + _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] + 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.""" + 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) + 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.""" + 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_a_pair_refused_by_the_cap_is_recorded_as_worker_cap(self): + """The measurement the worker-cap decision rests on. + + 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() + 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)) + 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.""" + 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)) + 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.""" + 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=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_help_capped=1) + + def test_losing_the_finalize_race_is_attributed_as_blocked(self): + """The commonest wait: every candidate done, a rival finalizing.""" + 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): + w.cooperative_solve(words, ROOT_BUDGET) + rows = self._wait_rows() + self.assertEqual(len(rows), 1) + 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.""" + 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) + self._assert_blocks(rows[0])