diff --git a/erd_queue.py b/erd_queue.py index 3a058d83..3d4777fd 100644 --- a/erd_queue.py +++ b/erd_queue.py @@ -3552,17 +3552,134 @@ def finalize_bundle_stats(self, branch_key): return (row["n_bundles"], row["max_bundle_nodes"], row["total_bundle_wall_millis"], row["censored_units"]) - def complete_candidate(self, branch_key, idx): - """Mark a candidate claim authoritatively complete (done=1).""" + def claim_is_current(self, branch_key, idx, claimed_by=None, + bundle_id=None, budget=None): + """Does this caller still hold an unfinished claim on this branch? + + One question for the whole set of writes an evaluation produces. A + result carries branch state in several places -- the taint flag, the + running best, the cut flag, the nodes spent, and the completion itself + -- and every one of them is only meaningful for the branch incarnation + the candidate was evaluated against. Guarding them one at a time + cannot be made safe: refusing one while accepting the others leaves the + branch describing a mixture of two incarnations, which is how a stale + OVER_ERD_LIMIT sets cut_occurred on a replacement that has no ceiling, + and finalize then reaches add_cut_result with a NULL bound. + + Answered in one indexed read: the claim must still exist unfinished and + belong to this caller, and the branch must still be open at the budget + the caller evaluated at. A NULL stored budget predates the column and + is admitted, as everywhere else. + """ + branch_id = self._intern_branch(branch_key) + if branch_id is None: + return False + row = self._conn.execute(""" + SELECT 1 + FROM candidate_claims c + JOIN active_branches a ON a.branch_id = c.branch_id + WHERE c.branch_id = ? AND c.idx = ? AND c.done = 0 + AND (? IS NULL OR c.claimed_by = ?) + AND (? IS NULL OR c.bundle_id = ?) + AND a.status = 'open' + AND (? IS NULL OR a.budget IS NULL OR a.budget = ?) + LIMIT 1 + """, (branch_id, idx, claimed_by, claimed_by, bundle_id, bundle_id, + budget, budget)).fetchone() + return row is not None + + def apply_candidate_result(self, branch_key, idx, *, claimed_by=None, + bundle_id=None, budget=None, nodes_spent=0, + infeasible=False, tainted=False, best=None, + cut=False): + """Apply every write one candidate evaluation produces, or none. + + Returns True when the result was applied. + + The writes are the branch's nodes spent, its taint flag, its running + best, its cut flag, and the candidate's completion. All five describe + the branch incarnation the candidate was evaluated against, and a + reclaimed claim can be reissued -- or the branch finalized, deleted and + re-created at another budget -- while the evaluation runs. + + Checking first and writing after cannot close that: the check and each + write are separate statements, so the branch can change between them + and leave a mixture of two incarnations behind. That is how a stale + OVER_ERD_LIMIT sets cut_occurred on a replacement with no ceiling, and + finalize then reaches add_cut_result with a NULL bound against a NOT + NULL column. So the validation and the writes share one transaction, + and the claim is re-read inside it. + + `best` is (best_guess, best_erd, max_depth) or None. Each write keeps + its own guard as well: they cost nothing here and they still hold for + the callers that use them directly. + """ + opened_transaction = not self._conn.in_transaction + if opened_transaction: + self._conn.execute("BEGIN IMMEDIATE") + try: + if not self.claim_is_current(branch_key, idx, + claimed_by=claimed_by, + bundle_id=bundle_id, budget=budget): + applied = False + else: + if nodes_spent or infeasible: + self.add_nodes_spent(branch_key, nodes_spent, + infeasible=infeasible) + if tainted: + self.mark_branch_tainted(branch_key) + if best is not None: + best_guess, best_erd, max_remaining_depth = best + self.update_branch_best(branch_key, best_guess, best_erd, + max_remaining_depth, budget=budget) + if cut: + self.mark_branch_cut(branch_key) + self.complete_candidate(branch_key, idx, + claimed_by=claimed_by, + bundle_id=bundle_id) + applied = True + except Exception: + if opened_transaction: + self._conn.execute("ROLLBACK") + raise + if opened_transaction: + self._conn.execute("COMMIT") + return applied + + def complete_candidate(self, branch_key, idx, claimed_by=None, + bundle_id=None): + """Mark a candidate claim authoritatively complete (done=1). + + Returns True when the row completed was the caller's own claim. + + Scoped to that claim, because branch_key and idx alone do not identify + one. A worker whose claim was reclaimed while it was still evaluating + goes on to finish; by then the index may have been reissued, to another + worker on this branch or to a different incarnation of it after a + finalize and re-creation. Completing by key and index alone marks that + live claim done while contributing nothing to it, and the branch can + then finalize without the candidate ever having been evaluated at the + budget it now holds -- cacheing an optimum some candidate beats, or a + loss that is not one. + + claimed_by and bundle_id are the pair `claim_next_bundle` stamps: the + bundle id is unique to one claim call and settles the case where the + same worker re-claimed the same index, and claimed_by carries a bare + claim that has no bundle. Passing neither asks for no check. + """ now = int(time.time()) branch_id = self._intern_branch(branch_key, create=True) self._conn.execute(""" UPDATE candidate_claims SET done = 1, done_at = ? WHERE branch_id = ? AND idx = ? - """, (now, branch_id, idx)) + AND (? IS NULL OR claimed_by = ?) + AND (? IS NULL OR bundle_id = ?) + """, (now, branch_id, idx, claimed_by, claimed_by, + bundle_id, bundle_id)) n = self._conn.execute("SELECT changes()").fetchone()[0] self._tally_wal_traffic( 'candidate_claims/complete', n, n * _CLAIM_ROW_WAL_BYTES) + return n > 0 def complete_bundle_two_level_erd_prunes(self, branch_key, bundle_id, candidate_indices, nodes_spent=0, @@ -3650,13 +3767,27 @@ def complete_bundle_two_level_erd_prunes(self, branch_key, bundle_id, updated_branch_count * _CLAIM_ROW_WAL_BYTES) return completed_candidate_count - def update_branch_best(self, branch_key, best_guess, best_erd, max_depth=None): + def update_branch_best(self, branch_key, best_guess, best_erd, + max_depth=None, budget=None): """Lower the branch's running best (monotone — never raises it). max_depth is the winning candidate's worst-case line length; it is stored atomically with the best it belongs to, so best_max_depth always describes the current best_guess. + budget is the budget the caller evaluated at, and the update applies + only to a branch still open at that budget. A branch can finalize and + be re-created at another budget under the same branch_key — the same + answer set reached by a second spine of a different length — while a + worker holding a claim on the old branch is still evaluating. Its cost + belongs to the budget it was computed at, and a cost from a larger + budget is below what a smaller one can achieve, so the monotone test + below would accept it and drive the new branch's best under its own + optimum. Ownership and priority both survive the re-creation and so + catch nothing. A stored budget of NULL predates the column and is + admitted, matching how callers derive a budget from the spine for + those; a caller passing no budget asks for no check. + The same statement stamps first_best_at/nodes_at_first_best on the update that creates the branch's first incumbent, and leaves them alone on every later improvement — COALESCE keeps the first value, so "how @@ -3675,8 +3806,9 @@ def update_branch_best(self, branch_key, best_guess, best_erd, max_depth=None): nodes_at_first_best = COALESCE(nodes_at_first_best, nodes_spent) WHERE branch_id = ? AND (best_erd IS NULL OR ? < best_erd) + AND (? IS NULL OR budget IS NULL OR budget = ?) """, (best_erd, best_guess, max_depth, now, now, branch_id, - best_erd)) + best_erd, budget, budget)) def read_branch_best(self, branch_key): """Return (best_guess, best_erd, ceiling) or (None, None, None). diff --git a/erd_swarm.py b/erd_swarm.py index 84dcecd1..743e1968 100644 --- a/erd_swarm.py +++ b/erd_swarm.py @@ -549,7 +549,8 @@ def check(self, token, candidate_list, last_index, # ceiling, if any, rides on the branch's ceiling column instead). if best_guess is not None: self._worker.queue.update_branch_best( - branch_key, best_guess, best_erd, best_max_remaining_depth) + branch_key, best_guess, best_erd, best_max_remaining_depth, + budget=budget) result = self._worker.cooperative_solve( branch_words, budget, @@ -1258,6 +1259,22 @@ def _heartbeat(self, branch_key, n_words, claim_idx, claim_started_at, # Count every invocation (one per node) BEFORE the throttle, so the # node counter is exact even though we only write every HB_SECONDS. self._nodes += 1 + self._liveness_tick(branch_key, n_words, claim_idx, claim_started_at, + best_guess, best_erd, force=force, + bound_erd=bound_erd) + + def _liveness_tick(self, branch_key, n_words, claim_idx, claim_started_at, + best_guess, best_erd, force=False, + bound_erd=None): + """Prove the worker is alive without counting a node. + + `_nodes` means candidate evaluations — the cost model, add_nodes_spent + and the accuracy rows all read it as one — so a signal that fires per + response group, or anywhere else below a candidate, must come through + here instead of `_heartbeat`. A worker that has not reached this + within HB_TIMEOUT_SECONDS has its in-flight claims reclaimed and + handed to another worker. + """ now = time.time() if not force and now - self._last_hb < HB_SECONDS: return @@ -1601,6 +1618,9 @@ def _metric_observer(group_sizes, has_self, candidate_cost_lower_bound, branch_floor_table=self.branch_floor_table, hint_cache=self.hint_cache, heartbeat=lambda: self._heartbeat( + branch_key, n_words, idx, claim_started, + local_candidate, local_best, bound_erd=_eff_bound()), + liveness_tick=lambda: self._liveness_tick( branch_key, n_words, idx, claim_started, local_candidate, local_best, bound_erd=_eff_bound())) cand_elapsed = time.time() - cand_t0 @@ -1611,15 +1631,11 @@ def _metric_observer(group_sizes, has_self, candidate_cost_lower_bound, idx, cand_elapsed, status, self._cand_max_depth) nodes_delta = self._nodes - nodes_before - if self._adaptive and (nodes_delta > 0 or status == OVER_DEPTH_BUDGET): - # These counters cover candidates proven infeasible at this level, - # not candidates whose taint arrived from a deeper branch. They - # are therefore a lower bound on local infeasibility. Every such - # proof also carries budget_tainted, so infeasible_candidates > 0 - # implies that the branch is marked tainted below. - self.queue.add_nodes_spent( - branch_key, nodes_delta, - infeasible=status == OVER_DEPTH_BUDGET) + # Counted with the rest of the result rather than on its own. An + # aborted candidate keeps its claim open for another worker to redo, so + # charging its nodes here as well would count the same candidate twice. + record_nodes = (self._adaptive + and (nodes_delta > 0 or status == OVER_DEPTH_BUDGET)) candidate_outcome = { SOLVED: 'exact', @@ -1655,32 +1671,48 @@ def _record_candidate_accuracy(): _record_candidate_accuracy() return False - # A candidate excluded by the depth cap (anywhere in its subtree) - # taints the branch: its ERD is only valid at this budget. Marked - # for any candidate, winner or not — see the taint rule. - if budget_tainted: - self.queue.mark_branch_tainted(branch_key) + # What this evaluation has to say about the branch, decided before any + # of it is written. A candidate excluded by the depth cap (anywhere in + # its subtree) taints the branch: its ERD is only valid at this budget, + # and that holds for any candidate, winner or not — see the taint rule. + improved_best = None + mark_cut = False if status == SOLVED: self.n_ok += 1 if local_best is None or cost < local_best: local_best, local_candidate, local_md = cost, candidate, cand_md - self.queue.update_branch_best(branch_key, local_candidate, - local_best, local_md) + improved_best = (local_candidate, local_best, local_md) shared_best = local_best elif status == OVER_ERD_LIMIT: self.n_cutoff += 1 - if branch_ceiling is not None: - # Priced out on a ceilinged branch. Only consulted at finalize - # when best_guess is NULL — where no real best ever existed, so - # every price-out was against the ceiling and the branch is a - # cut, not a proven loss. - self.queue.mark_branch_cut(branch_key) + # Priced out on a ceilinged branch. Only consulted at finalize + # when best_guess is NULL — where no real best ever existed, so + # every price-out was against the ceiling and the branch is a cut, + # not a proven loss. + mark_cut = branch_ceiling is not None elif status == OVER_DEPTH_BUDGET: self.n_pruned += 1 else: # pragma: no cover self.n_useless += 1 - self.queue.complete_candidate(branch_key, idx) + # One transaction: the claim is re-read inside it, so either every one + # of these lands on the incarnation this candidate was evaluated + # against or none of them lands at all. + if not self.queue.apply_candidate_result( + branch_key, idx, claimed_by=self.name, bundle_id=bundle_id, + budget=budget, + nodes_spent=nodes_delta if record_nodes else 0, + infeasible=record_nodes and status == OVER_DEPTH_BUDGET, + tainted=budget_tainted, best=improved_best, cut=mark_cut): + # The claim was reissued, or the branch was re-created, while this + # candidate ran. The bundle is not abandoned with it: a one-level + # prune sweep replaces a single claim row, so the siblings may + # still be this worker's to finish, and a worker that is alive and + # heartbeating never has them reclaimed for it. + logger.warning( + '%s lost candidate %s (idx=%d) mid-evaluation; its result ' + 'describes a branch incarnation this worker no longer holds ' + 'and was discarded', self.name, candidate, idx) # The outbound claim telemetry is required for branch ETA reporting, # regardless of whether this worker uses adaptive decomposition. now_complete = time.time() @@ -1794,7 +1826,10 @@ def _complete_bundle_two_level_erd_prunes( words, candidate, self.rcache, guesses=self.all_words, pattern_matrix=self.pattern_matrix, branch_indices=branch_indices, - branch_floor_table=self.branch_floor_table) + branch_floor_table=self.branch_floor_table, + liveness_tick=lambda: self._liveness_tick( + branch_key, n_words, candidate_index, claim_started_at, + best_guess, best_erd, bound_erd=bound_erd)) if candidate_cost_lower_bound >= bound_erd: pruned_candidate_indices.append(candidate_index) diff --git a/tests/test_branch_cost_lower_bound.py b/tests/test_branch_cost_lower_bound.py index 0a6a4505..82865ab1 100644 --- a/tests/test_branch_cost_lower_bound.py +++ b/tests/test_branch_cost_lower_bound.py @@ -14,6 +14,7 @@ import unittest from unittest import mock +import wordle_engine from cache_sqlite import ScoreCache from erd_queue import encode_subset from pattern_matrix import PatternMatrix @@ -21,7 +22,7 @@ from wordle_engine import ( ERD_ALL, BranchFloorTable, ResponseCache, all_singletons_floor, candidate_two_level_cost_lower_bound, evaluate_candidate, - _candidate_cost_lower_bound, min_expected_guesses, + _ALL_GREEN_PATTERN, _candidate_cost_lower_bound, min_expected_guesses, sub_branch_cost_lower_bound, ) @@ -694,3 +695,147 @@ def test_two_level_erd_prune_bound_is_the_engine_entry_gate(self): if __name__ == "__main__": unittest.main() + + +class TestPricingGroupsProvesLiveness(_VocabularyMixin, unittest.TestCase): + """Pricing a candidate's response groups must signal liveness as it goes. + + `evaluate_candidate` ticks once on entry and then relies on recursion to + reach the next tick. A candidate that prunes on its bound never recurses, + so that single tick covers its entire evaluation. + + That gap is closed because it is cheap to close, not because the time is + spent here: `_remaining_groups_cost_lower_bounds` measures 1.4 ms to 44 ms + per group and 0.52 s to 0.91 s for a candidate priced over the whole answer + list, which is two orders of magnitude inside HB_TIMEOUT_SECONDS. This + loop is not where a worker falls silent. + + The tick is observation only: it can never change a bound, so these assert + on the signal alone and on the bound being unchanged by its presence. + """ + + def _ticks_for_evaluate(self, branch_words, candidate, best_erd, + liveness_tick): + return evaluate_candidate( + branch_words, candidate, self._response_cache(), None, + best_erd=best_erd, guesses=self.guess_words, policy=ERD_ALL, + budget=5, pattern_matrix=self.pattern_matrix, + branch_floor_table=self._table(), + liveness_tick=liveness_tick) + + def _two_level_pruning_bound(self, branch_words, candidate): + """A bound that only the group-pricing gate can prove. + + Below the closed-form bound the candidate is priced out by the + vectorized one-level check and the group loop never runs at all -- the + cheap prune, which needs no tick. The expensive prune is the band + above it, where every group must be priced before the candidate can be + rejected -- so it is the band where a tick in that loop is the only + signal a candidate emits after its entry tick. + """ + cache = self._response_cache() + groups = cache.group_words( + candidate, branch_words, pattern_matrix=self.pattern_matrix, + branch_indices=self.pattern_matrix.answer_indices(branch_words)) + closed_form = _candidate_cost_lower_bound( + groups.values(), _ALL_GREEN_PATTERN in groups, len(branch_words)) + two_level = candidate_two_level_cost_lower_bound( + branch_words, candidate, cache, guesses=self.guess_words, + pattern_matrix=self.pattern_matrix, + branch_floor_table=self._table()) + self.assertGreater(two_level, closed_form, + "fixture has no band only group pricing can decide") + return len(groups), (closed_form + two_level) / 2 + + def test_a_candidate_that_prunes_without_recursing_ticks_per_group(self): + branch_words = self._branch(40, seed=77) + candidate = self.guess_words[0] + group_count, bound = self._two_level_pruning_bound( + branch_words, candidate) + ticks = [] + status, _cost, max_remaining_depth, _floor = self._ticks_for_evaluate( + branch_words, candidate, bound, lambda: ticks.append(1)) + self.assertIsNone(max_remaining_depth, + "fixture recursed; it must prune on the bound") + self.assertEqual( + len(ticks), group_count, + "a non-recursing candidate produced no signal beyond entry") + + def test_the_tick_cannot_change_the_answer(self): + branch_words = self._branch(40, seed=77) + candidate = self.guess_words[0] + _group_count, bound = self._two_level_pruning_bound( + branch_words, candidate) + without = self._ticks_for_evaluate(branch_words, candidate, bound, None) + with_tick = self._ticks_for_evaluate( + branch_words, candidate, bound, lambda: None) + self.assertEqual(without, with_tick) + + def test_the_two_level_bound_ticks_per_group_too(self): + # The swarm prices a whole bundle through this entry before evaluating + # any of it, and that pass has no recursion to fall back on at all. + branch_words = self._branch(40, seed=78) + candidate = self.guess_words[0] + ticks = [] + bound = candidate_two_level_cost_lower_bound( + branch_words, candidate, self._response_cache(), + guesses=self.guess_words, pattern_matrix=self.pattern_matrix, + branch_floor_table=self._table(), + liveness_tick=lambda: ticks.append(1)) + unticked = candidate_two_level_cost_lower_bound( + branch_words, candidate, self._response_cache(), + guesses=self.guess_words, pattern_matrix=self.pattern_matrix, + branch_floor_table=self._table()) + self.assertGreater(len(ticks), 1) + self.assertEqual(bound, unticked) + + +class TestDescendantFramesPriceGroupsWithATick(_VocabularyMixin, + unittest.TestCase): + """The tick must reach every frame, not just the entry one. + + `evaluate_candidate` recurses through `_solve_subset`, which evaluates each + descendant candidate in turn. A descendant priced its response groups with + no tick would go silent exactly as the entry frame did, and every result + assertion would still pass -- so this asserts on the frames actually + reached rather than on the answer. + """ + + def _frames_that_priced_groups(self, branch_words, candidate, best_erd, + liveness_tick): + """Each (branch_size, got_a_tick) the floor loop was entered with.""" + frames = [] + real = wordle_engine._remaining_groups_cost_lower_bounds + + def recording(ordered_groups, group_candidate, branch_size, + branch_floor_table, liveness_tick=None): + frames.append((branch_size, liveness_tick is not None)) + return real(ordered_groups, group_candidate, branch_size, + branch_floor_table, liveness_tick=liveness_tick) + + with mock.patch.object(wordle_engine, + "_remaining_groups_cost_lower_bounds", + recording): + evaluate_candidate( + branch_words, candidate, self._response_cache(), + None, best_erd=best_erd, guesses=self.guess_words, + policy=ERD_ALL, budget=5, + pattern_matrix=self.pattern_matrix, + branch_floor_table=self._table(), + liveness_tick=liveness_tick) + return frames + + def test_every_frame_that_prices_groups_is_given_the_tick(self): + branch_words = self._branch(40, seed=77) + frames = self._frames_that_priced_groups( + branch_words, self.guess_words[0], float("inf"), lambda: None) + + descendant_sizes = {size for size, _ in frames + if size != len(branch_words)} + self.assertTrue( + descendant_sizes, + "fixture never recursed; it cannot cover descendant frames") + unticked = sorted({size for size, ticked in frames if not ticked}) + self.assertEqual( + unticked, [], + f"frames priced groups with no liveness tick: sizes {unticked}") diff --git a/tests/test_erd_queue_unit.py b/tests/test_erd_queue_unit.py index d2e57666..0b1a4788 100644 --- a/tests/test_erd_queue_unit.py +++ b/tests/test_erd_queue_unit.py @@ -778,6 +778,177 @@ def test_update_branch_best_is_monotone(self): self.assertEqual(guess, "crane") self.assertAlmostEqual(erd, 2.0) + def test_update_branch_best_at_a_stale_budget_cannot_lower_a_recreated_branch(self): + # A branch finalizes, is deleted, and is re-created at a smaller budget + # -- the same answer set reached by a longer spine. A worker whose + # claim on the old incarnation was reclaimed is still evaluating, and + # folds in a cost computed at the larger budget. It is below anything + # the smaller budget can achieve, so the monotone test alone accepts + # it and the branch finalizes under its own optimum. + self.q.create_branch(self.key, len(WORDS), N_CANDIDATES, budget=5) + self.q.delete_branch(self.key) + self.q.create_branch(self.key, len(WORDS), N_CANDIDATES, budget=3) + self.q.update_branch_best(self.key, "crane", 3.0, max_depth=3, budget=3) + + self.q.update_branch_best(self.key, "slate", 1.5, max_depth=5, budget=5) + + guess, erd, _ceiling = self.q.read_branch_best(self.key) + self.assertEqual(guess, "crane") + self.assertAlmostEqual(erd, 3.0) + + def test_update_branch_best_at_the_branch_budget_still_lowers(self): + # The guard rejects a stale budget, never a legitimate improvement. + self.q.create_branch(self.key, len(WORDS), N_CANDIDATES, budget=3) + self.q.update_branch_best(self.key, "crane", 3.0, max_depth=3, budget=3) + self.q.update_branch_best(self.key, "slate", 2.0, max_depth=3, budget=3) + guess, erd, _ceiling = self.q.read_branch_best(self.key) + self.assertEqual(guess, "slate") + self.assertAlmostEqual(erd, 2.0) + + def test_update_branch_best_admits_a_branch_whose_budget_predates_the_column(self): + # A NULL stored budget carries no assertion to contradict, so it is + # admitted -- the same rule the claim transaction applies. + self.q.create_branch(self.key, len(WORDS), N_CANDIDATES) + self.q.update_branch_best(self.key, "crane", 2.0, max_depth=3, budget=5) + guess, erd, _ceiling = self.q.read_branch_best(self.key) + self.assertEqual(guess, "crane") + self.assertAlmostEqual(erd, 2.0) + + def test_update_branch_best_without_a_budget_asks_for_no_check(self): + self.q.create_branch(self.key, len(WORDS), N_CANDIDATES, budget=3) + self.q.update_branch_best(self.key, "crane", 2.0, max_depth=3) + guess, erd, _ceiling = self.q.read_branch_best(self.key) + self.assertEqual(guess, "crane") + self.assertAlmostEqual(erd, 2.0) + + def test_completing_a_candidate_reissued_to_another_worker_is_refused(self): + # The stale worker's whole hazard in one case: its claim was reclaimed, + # the index reissued, and it now finishes. Completing by key and index + # alone would mark the new holder's live claim done with no result + # behind it, and the branch could finalize a candidate nobody evaluated + # at the budget it now holds. + self.q.create_branch(self.key, len(WORDS), N_CANDIDATES) + idx = self._claim_one_idx(self.key, worker_id="worker-0") + self.q.reclaim_claims_of_worker("worker-0") + reissued = self._claim_one_idx(self.key, worker_id="worker-1") + self.assertEqual(reissued, idx, "fixture did not reissue the index") + + self.assertFalse( + self.q.complete_candidate(self.key, idx, claimed_by="worker-0")) + + self.assertEqual(self.q.branch_done_candidates(self.key), 0, + "a live claim was marked done by a stale worker") + + def test_completing_a_candidate_this_worker_still_holds_succeeds(self): + self.q.create_branch(self.key, len(WORDS), N_CANDIDATES) + idx = self._claim_one_idx(self.key, worker_id="worker-0") + self.assertTrue( + self.q.complete_candidate(self.key, idx, claimed_by="worker-0")) + self.assertEqual(self.q.branch_done_candidates(self.key), 1) + + def test_completing_without_an_owner_still_asks_for_no_check(self): + self.q.create_branch(self.key, len(WORDS), N_CANDIDATES) + idx = self._claim_one_idx(self.key, worker_id="worker-0") + self.assertTrue(self.q.complete_candidate(self.key, idx)) + self.assertEqual(self.q.branch_done_candidates(self.key), 1) + + def test_claim_is_current_is_false_for_a_branch_never_registered(self): + # No branch_id means no claim can exist against it, so the answer is + # no -- and asking must not intern the key, which would register a + # branch as a side effect of a read. + self.assertFalse(self.q.claim_is_current( + b"notakey", 0, claimed_by="worker-0")) + + def test_claim_is_current_is_true_for_the_worker_that_holds_it(self): + self.q.create_branch(self.key, len(WORDS), N_CANDIDATES, budget=5) + idx = self._claim_one_idx(self.key, worker_id="worker-0") + self.assertTrue(self.q.claim_is_current( + self.key, idx, claimed_by="worker-0", budget=5)) + + def test_claim_is_current_is_false_once_the_claim_is_reissued(self): + self.q.create_branch(self.key, len(WORDS), N_CANDIDATES, budget=5) + idx = self._claim_one_idx(self.key, worker_id="worker-0") + self.q.reclaim_claims_of_worker("worker-0") + self._claim_one_idx(self.key, worker_id="worker-1") + self.assertFalse(self.q.claim_is_current( + self.key, idx, claimed_by="worker-0", budget=5)) + + def test_claim_is_current_is_false_at_a_budget_the_branch_no_longer_holds(self): + # The budget clause has to be what decides this. delete_branch also + # deletes the claim rows, so re-creating the branch and asking straight + # away answers False because the JOIN finds no claim at all -- true + # with the budget clause deleted as well. A live claim under the new + # incarnation is what isolates it. + self.q.create_branch(self.key, len(WORDS), N_CANDIDATES, budget=5) + self._claim_one_idx(self.key, worker_id="worker-0") + self.q.delete_branch(self.key) + self.q.create_branch(self.key, len(WORDS), N_CANDIDATES, budget=3) + idx = self._claim_one_idx(self.key, worker_id="worker-0") + + self.assertTrue( + self.q.claim_is_current(self.key, idx, claimed_by="worker-0", + budget=3), + "fixture has no live claim; the budget clause decides nothing") + self.assertFalse(self.q.claim_is_current( + self.key, idx, claimed_by="worker-0", budget=5)) + + def test_claim_is_current_is_false_for_a_branch_that_finalized(self): + self.q.create_branch(self.key, len(WORDS), N_CANDIDATES, budget=5) + idx = self._claim_one_idx(self.key, worker_id="worker-0") + self.q.delete_branch(self.key) + self.assertFalse(self.q.claim_is_current( + self.key, idx, claimed_by="worker-0", budget=5)) + + def test_apply_candidate_result_writes_everything_or_nothing(self): + self.q.create_branch(self.key, len(WORDS), N_CANDIDATES, budget=5) + idx = self._claim_one_idx(self.key, worker_id="worker-0") + + applied = self.q.apply_candidate_result( + self.key, idx, claimed_by="worker-0", budget=5, + nodes_spent=9, infeasible=False, tainted=True, + best=("crane", 2.5, 3), cut=False) + + self.assertTrue(applied) + self.assertEqual(self.q.branch_done_candidates(self.key), 1) + guess, erd, _ceiling = self.q.read_branch_best(self.key) + self.assertEqual(guess, "crane") + self.assertAlmostEqual(erd, 2.5) + row = self.q.get_branch(self.key) + self.assertEqual(row["nodes_spent"], 9) + self.assertTrue(row["tainted"]) + + def test_apply_candidate_result_writes_nothing_once_the_claim_is_gone(self): + # The whole point: a refusal must leave no trace of any of the five, + # not just of the completion. + self.q.create_branch(self.key, len(WORDS), N_CANDIDATES, budget=5) + idx = self._claim_one_idx(self.key, worker_id="worker-0") + self.q.reclaim_claims_of_worker("worker-0") + self._claim_one_idx(self.key, worker_id="worker-1") + + applied = self.q.apply_candidate_result( + self.key, idx, claimed_by="worker-0", budget=5, + nodes_spent=9, infeasible=True, tainted=True, + best=("crane", 0.5, 3), cut=True) + + self.assertFalse(applied) + self.assertEqual(self.q.branch_done_candidates(self.key), 0, + "a live claim was completed by a stale worker") + self.assertEqual(self.q.read_branch_best(self.key)[0], None, + "a stale best was published") + row = self.q.get_branch(self.key) + self.assertEqual(row["nodes_spent"], 0, "stale nodes were charged") + self.assertFalse(row["tainted"], "a stale taint was set") + self.assertFalse(row["cut_occurred"], "a stale cut was set") + + def test_apply_candidate_result_leaves_no_transaction_open(self): + # It opens BEGIN IMMEDIATE; a leaked transaction would block every + # other writer for as long as this worker lives. + self.q.create_branch(self.key, len(WORDS), N_CANDIDATES, budget=5) + idx = self._claim_one_idx(self.key, worker_id="worker-0") + self.q.apply_candidate_result( + self.key, idx, claimed_by="worker-0", budget=5, nodes_spent=1) + self.assertFalse(self.q._conn.in_transaction) + def test_read_branch_best_returns_none_none_for_missing_key(self): self.assertEqual(self.q.read_branch_best(b"notakey"), (None, None, None)) diff --git a/tests/test_erd_swarm_unit.py b/tests/test_erd_swarm_unit.py index 8abf0220..2cbb96e4 100644 --- a/tests/test_erd_swarm_unit.py +++ b/tests/test_erd_swarm_unit.py @@ -171,6 +171,40 @@ def test_node_counter_increments_even_when_db_write_is_throttled(self): self.assertEqual(w._nodes, 2) # counter still incremented self.assertEqual(w.queue.heartbeat.call_count, 1) # still only one DB write + def test_liveness_tick_writes_a_heartbeat_without_counting_a_node(self): + """A signal that fires below one candidate must not move `_nodes`. + + `_nodes` means candidate evaluations, and the cost model, + add_nodes_spent and the accuracy rows all read it as one. Routing a + per-response-group tick through `_heartbeat` would prove liveness and + inflate every one of them, which is the tempting simplification this + pins against. + """ + w = _bare_worker() + branch_key = ScoreCache.encode_subset(BRANCH) + + w._liveness_tick(branch_key, len(BRANCH), 0, 0, None, None, force=True) + self.assertEqual(w._nodes, 0, "a liveness tick counted a node") + self.assertEqual(w.queue.heartbeat.call_count, 1, + "a liveness tick did not prove liveness") + + def test_liveness_tick_is_throttled_on_the_same_clock_as_a_heartbeat(self): + # Pricing a branch's groups fires this per group, so an unthrottled + # tick would write a heartbeat row per group. + w = _bare_worker() + branch_key = ScoreCache.encode_subset(BRANCH) + w._liveness_tick(branch_key, len(BRANCH), 0, 0, None, None, force=True) + for _ in range(50): + w._liveness_tick(branch_key, len(BRANCH), 0, 0, None, None) + self.assertEqual(w.queue.heartbeat.call_count, 1) + self.assertEqual(w._nodes, 0) + + def test_heartbeat_still_counts_its_node(self): + w = _bare_worker() + branch_key = ScoreCache.encode_subset(BRANCH) + w._heartbeat(branch_key, len(BRANCH), 0, 0, None, None, force=True) + self.assertEqual(w._nodes, 1) + def test_hb_max_spine_reset_after_each_db_write(self): """_hb_max_spine is cleared after each DB write so the 2-second window starts fresh — the next heartbeat builds a new spine from scratch.""" @@ -3202,9 +3236,10 @@ def evaluate(*_args, **_kwargs): self.assertTrue(worker.evaluate_claim( b"branch", BRANCH, len(BRANCH), 0, budget=4)) - worker.queue.add_nodes_spent.assert_called_once_with( - b"branch", 7, infeasible=True) - worker.queue.mark_branch_tainted.assert_called_once_with(b"branch") + kwargs = worker.queue.apply_candidate_result.call_args.kwargs + self.assertEqual(kwargs["nodes_spent"], 7) + self.assertTrue(kwargs["infeasible"]) + self.assertTrue(kwargs["tainted"]) class TestTwoLevelERDPruneBundles(unittest.TestCase): @@ -4086,9 +4121,11 @@ def test_check_calls_update_branch_best_when_best_guess_known(self): self.assertIsNotNone(result) # The seed carries the winner's worst-case line, not just its cost: a # branch seeded with an unknown depth finalizes into a cache row no - # budget can ever reuse, so it reads as unsolved forever. + # budget can ever reuse, so it reads as unsolved forever. It carries + # the budget it was achieved at too, so a branch re-created at another + # budget under the same key does not take this seed as its own. w.queue.update_branch_best.assert_called_once_with( - ScoreCache.encode_subset(BRANCH[:6]), "crane", 1.5, 3) + ScoreCache.encode_subset(BRANCH[:6]), "crane", 1.5, 3, budget=5) def test_check_skips_update_branch_best_when_no_best_guess(self): result, w, _ = self._pub_overrun(best_guess=None) @@ -4599,7 +4636,7 @@ def test_achieved_best_seeds_and_publishes_exact(self): pub.check(token, CANDIDATES, 1, "crane", 1.8, 4, 5) self.assertIsNone(w.queue.create_branch.call_args.kwargs["ceiling"]) w.queue.update_branch_best.assert_called_once_with( - ScoreCache.encode_subset(BRANCH[:6]), "crane", 1.8, 4) + ScoreCache.encode_subset(BRANCH[:6]), "crane", 1.8, 4, budget=5) w.queue.mark_claims_done.assert_called_once() w.cooperative_solve.assert_called_once_with( BRANCH[:6], 5, ceiling=float('inf')) @@ -6244,3 +6281,64 @@ def _lose_then_stop(*_a, **_kw): self.assertEqual(len(rows), 1) self.assertEqual(rows[0]["blocked_millis"], 0) self._assert_blocks(rows[0]) + +class TestAStaleResultIsNotAppliedPiecemeal(unittest.TestCase): + """The worker hands its whole result to one call, and survives a refusal. + + Whether those writes land atomically is the queue's property, tested + against a real database in test_erd_queue_unit. What belongs here is that + the worker states the whole result in one place, issues no branch write + outside it, and does the right thing when it is refused. + """ + + def _worker(self, applied): + w = _bare_worker() + w._adaptive = True + w.queue.apply_candidate_result.return_value = applied + # A ceiling in scope is what makes a price-out a cut. + w.queue.read_branch_best.return_value = (None, None, 4.5) + return w + + def _evaluate(self, worker): + branch_key = ScoreCache.encode_subset(BRANCH) + + def _evaluated(*args, **kwargs): + worker._nodes += 5 + return (erd_swarm.OVER_ERD_LIMIT, 4.0, 2, True) + + with mock.patch.object(erd_swarm, "evaluate_candidate", _evaluated): + return worker.evaluate_claim(branch_key, BRANCH, len(BRANCH), + idx=0, budget=5) + + def test_the_whole_result_is_handed_over_in_one_call(self): + w = self._worker(applied=True) + self._evaluate(w) + w.queue.apply_candidate_result.assert_called_once() + kwargs = w.queue.apply_candidate_result.call_args.kwargs + self.assertEqual(kwargs["claimed_by"], w.name) + self.assertEqual(kwargs["budget"], 5) + self.assertTrue(kwargs["tainted"], "the taint was not carried") + self.assertTrue(kwargs["cut"], "the cut was not carried") + self.assertEqual(kwargs["nodes_spent"], 5) + + def test_no_branch_write_is_issued_outside_that_call(self): + # Every one of these used to be its own statement on this path. + w = self._worker(applied=True) + self._evaluate(w) + for name in ("mark_branch_tainted", "mark_branch_cut", + "add_nodes_spent", "update_branch_best", + "complete_candidate"): + with self.subTest(write=name): + getattr(w.queue, name).assert_not_called() + + def test_a_refused_result_does_not_abandon_the_rest_of_the_bundle(self): + """A lost claim is not a cancellation. + + claim_next_bundle's one-level sweep replaces a single claim row, so the + siblings may still be this worker's to finish -- and a worker that is + alive and heartbeating never has them reclaimed for it, so abandoning + them leaves the branch unable to finalize at all. + """ + w = self._worker(applied=False) + self.assertTrue(self._evaluate(w), + "a refused result abandoned the bundle") diff --git a/wordle_engine.py b/wordle_engine.py index 1ece533d..6a9ecbdc 100644 --- a/wordle_engine.py +++ b/wordle_engine.py @@ -1318,10 +1318,30 @@ def _candidate_response_groups(branch_words, candidate, cache, def _remaining_groups_cost_lower_bounds(ordered_groups, candidate, - branch_size, branch_floor_table): - """Suffix sums used by both the two-level entry gate and sub-ceilings.""" + branch_size, branch_floor_table, + liveness_tick=None): + """Suffix sums used by both the two-level entry gate and sub-ceilings. + + liveness_tick fires once per response group, and is observation only: it + can never change a bound. + + Measured at production vocabulary, this loop is cheap. One group's floor + is a single pass over (guess vocabulary x group): 1.4 ms on a two-word + group, 44 ms on the largest group that can exist -- the whole answer list. + A candidate's groups partition its branch, so the loop is not 243 + full-branch scans; priced over the entire answer list it runs in 0.52 s to + 0.91 s. + + So this is not a place a worker can fall silent for HB_TIMEOUT_SECONDS. + The tick is here because it costs one call per group and closes a gap + wherever uninterrupted work happens, not because the time is spent here. + Do not cite this loop as the cause of a stale-claim reclamation without + measuring again. + """ remaining_groups_cost_lower_bound = [0.0] * (len(ordered_groups) + 1) for index in range(len(ordered_groups) - 1, -1, -1): + if liveness_tick is not None: + liveness_tick() sub_branch = ordered_groups[index][1] remaining_groups_cost_lower_bound[index] = ( remaining_groups_cost_lower_bound[index + 1] @@ -1333,7 +1353,8 @@ def _remaining_groups_cost_lower_bounds(ordered_groups, candidate, def candidate_two_level_cost_lower_bound( branch_words, candidate, cache, guesses=None, - pattern_matrix=None, branch_indices=None, branch_floor_table=None): + pattern_matrix=None, branch_indices=None, branch_floor_table=None, + liveness_tick=None): """Admissible two-level ERD lower bound for one candidate. This is evaluate_candidate's entry proof without recursion. It performs @@ -1364,7 +1385,8 @@ def candidate_two_level_cost_lower_bound( groups.values(), has_self, branch_size) ordered_groups = sorted(groups.items(), key=_by_group_size, reverse=True) remaining_groups_cost_lower_bound = _remaining_groups_cost_lower_bounds( - ordered_groups, candidate, branch_size, branch_floor_table) + ordered_groups, candidate, branch_size, branch_floor_table, + liveness_tick=liveness_tick) two_level_cost_lower_bound = ( 1.0 + remaining_groups_cost_lower_bound[0] ) @@ -1464,7 +1486,8 @@ def evaluate_candidate(branch_words, candidate, cache, score_cache, *, subbranch_solver=None, bound_provider=None, mid_loop_publisher=None, metric_observer=None, pattern_matrix=None, branch_indices=None, - branch_floor_table=None, hint_cache=None): + branch_floor_table=None, hint_cache=None, + liveness_tick=None): """Evaluate one `candidate`'s exact ERD for solving `branch_words`. This is the body of the top-level candidate loop, extracted so a parallel @@ -1575,7 +1598,8 @@ def _observe(pruned): # *after* position i (each sub-branch of size k costs >= lb(k)). The self # singleton contributes 0. remaining_groups_cost_lower_bound = _remaining_groups_cost_lower_bounds( - ordered, candidate, n, branch_floor_table) + ordered, candidate, n, branch_floor_table, + liveness_tick=liveness_tick) def _sub_lb(sub_branch): return sub_branch_cost_lower_bound( @@ -1624,7 +1648,8 @@ def _sub_lb(sub_branch): subbranch_solver, ceiling=sub_ceiling, entry_guess=candidate, entry_pattern=pattern_code, mid_loop_publisher=mid_loop_publisher, pattern_matrix=pattern_matrix, - branch_floor_table=branch_floor_table, hint_cache=hint_cache) + branch_floor_table=branch_floor_table, hint_cache=hint_cache, + liveness_tick=liveness_tick) if sub in _ABORT_STATUSES: return (sub, None, None, False) sub_status, sub_cost, sub_max_remaining_depth, sub_budget_tainted = sub @@ -1695,7 +1720,7 @@ def _solve_subset(branch_words, cache, score_cache, budget, deadline, guesses, branch_floor_table=None, ceiling=float('inf'), entry_guess=None, entry_pattern=None, mid_loop_publisher=None, pattern_matrix=None, - hint_cache=None): + hint_cache=None, liveness_tick=None): """Budget-aware core of min_expected_guesses. Returns (cost, max_depth, floor_hit, cutoff), or None on deadline/cancel @@ -1886,6 +1911,7 @@ def _solve_subset(branch_words, cache, score_cache, budget, deadline, guesses, branch_indices=branch_indices, branch_floor_table=branch_floor_table, hint_cache=hint_cache, + liveness_tick=liveness_tick, ) if status in _ABORT_STATUSES: return status