From 50f6bdd3966ec4349ba9fae191de21a204d26df6 Mon Sep 17 00:00:00 2001 From: YingzuoLiu <156630675+YingzuoLiu@users.noreply.github.com> Date: Wed, 26 Aug 2026 18:24:21 +0800 Subject: [PATCH 1/8] fix: stabilize P5 proof synchronization --- docs/p5-multi-worker-recovery-proof.md | 6 ++ examples/p5_multi_worker_proof.py | 5 +- examples/p5_proof_worker.py | 37 ++++--- tests/test_p5_multi_worker_proof.py | 130 +++++++++++++++++++++++++ 4 files changed, 165 insertions(+), 13 deletions(-) diff --git a/docs/p5-multi-worker-recovery-proof.md b/docs/p5-multi-worker-recovery-proof.md index a62fb53..167cc7c 100644 --- a/docs/p5-multi-worker-recovery-proof.md +++ b/docs/p5-multi-worker-recovery-proof.md @@ -7,6 +7,12 @@ controller process submits Runs, controls named one-shot barriers, injects exact store-time lease expiry, and rereads durable evidence. The synthetic HTTP provider is a fourth, independent process with its own SQLite effect ledger. +Each reusable barrier has an explicit consumer-completed acknowledgement. The +controller cannot clear and re-arm a barrier until the worker has observed the prior +release, preventing one generation from consuming the next generation's signal. The +post-run PostgreSQL session check runs only after both polling workers have stopped; +an in-flight polling transaction is not misclassified as a leaked session. + The proof adds no production scheduler, distributed queue, connection pool, or production control endpoint. `RuntimeManager._wake` remains a process-local latency optimization; bounded PostgreSQL polling is the durable progress mechanism. Fault diff --git a/examples/p5_multi_worker_proof.py b/examples/p5_multi_worker_proof.py index 3fb28dc..061af7a 100644 --- a/examples/p5_multi_worker_proof.py +++ b/examples/p5_multi_worker_proof.py @@ -1328,6 +1328,7 @@ def run_proof( port=_free_loopback_port(), ) controller: ProofController | None = None + workers_stopped = False scenarios: list[dict[str, Any]] = [] cleanup_error: Exception | None = None try: @@ -1352,6 +1353,8 @@ def run_proof( result = scenario(controller) scenarios.append(result) _log_scenario(result) + controller.stop_workers() + workers_stopped = True idle_sessions = idle_in_transaction_count(proof_dsn) _require(idle_sessions == 0, "P5 left an idle-in-transaction PostgreSQL session") schedule_result = next( @@ -1396,7 +1399,7 @@ def run_proof( encoding="utf-8", ) finally: - if controller is not None: + if controller is not None and not workers_stopped: controller.stop_workers() provider.stop() provider_temp.cleanup() diff --git a/examples/p5_proof_worker.py b/examples/p5_proof_worker.py index 81c1be2..c8de326 100644 --- a/examples/p5_proof_worker.py +++ b/examples/p5_proof_worker.py @@ -83,18 +83,25 @@ class ProcessHook: enabled: Any reached: Any release: Any + completed: Any metadata: Any @classmethod def create(cls, context: Any) -> ProcessHook: + completed = context.Event() + completed.set() return cls( enabled=context.Event(), reached=context.Event(), release=context.Event(), + completed=completed, metadata=context.Queue(maxsize=1), ) def arm(self) -> None: + if not self.completed.wait(P5_HOOK_TIMEOUT_SECONDS): + raise TimeoutError("P5 previous proof hook consumer did not finish") + self.completed.clear() self.reached.clear() self.release.clear() while True: @@ -111,10 +118,13 @@ def consume(self) -> bool: return True def block_consumed(self, payload: dict[str, Any]) -> None: - self.metadata.put(payload) - self.reached.set() - if not self.release.wait(P5_HOOK_TIMEOUT_SECONDS): - raise TimeoutError("P5 controller did not release a reached proof hook") + try: + self.metadata.put(payload) + self.reached.set() + if not self.release.wait(P5_HOOK_TIMEOUT_SECONDS): + raise TimeoutError("P5 controller did not release a reached proof hook") + finally: + self.completed.set() def hit(self, payload: dict[str, Any]) -> bool: if not self.consume(): @@ -127,14 +137,17 @@ def pause_process(self, payload: dict[str, Any]) -> bool: if not self.consume(): return False - self.metadata.put(payload) - # multiprocessing.Queue publishes through a feeder thread. Flush this - # one-shot channel before SIGSTOP so the controller never observes the - # reached event before its matching payload is readable. - self.metadata.close() - self.metadata.join_thread() - self.reached.set() - os.kill(os.getpid(), signal.SIGSTOP) + try: + self.metadata.put(payload) + # multiprocessing.Queue publishes through a feeder thread. Flush this + # one-shot channel before SIGSTOP so the controller never observes the + # reached event before its matching payload is readable. + self.metadata.close() + self.metadata.join_thread() + self.reached.set() + os.kill(os.getpid(), signal.SIGSTOP) + finally: + self.completed.set() return True diff --git a/tests/test_p5_multi_worker_proof.py b/tests/test_p5_multi_worker_proof.py index 3911ceb..092389d 100644 --- a/tests/test_p5_multi_worker_proof.py +++ b/tests/test_p5_multi_worker_proof.py @@ -3,12 +3,15 @@ import multiprocessing import json import os +import queue import signal import threading from pathlib import Path +from types import SimpleNamespace import pytest +import examples.p5_multi_worker_proof as p5_proof from examples.p5_multi_worker_proof import ( REPORT_VERSION, P5ProofFailure, @@ -19,6 +22,7 @@ from examples.p5_proof_worker import ( P5ProofHooks, P5WorkerConfig, + ProcessHook, _TerminatedConnection, ) from tests.conformance.p5_mutation_proof import p5_mutants @@ -143,6 +147,132 @@ def hit() -> None: assert completed.is_set() +def test_process_hook_rearm_waits_for_previous_consumer_acknowledgement() -> None: + class DelayedRelease: + def __init__(self) -> None: + self._set = threading.Event() + self.allow_wait_to_return = threading.Event() + + def clear(self) -> None: + self._set.clear() + + def set(self) -> None: + self._set.set() + + def wait(self, timeout: float) -> bool: + if not self._set.wait(timeout): + return False + return self.allow_wait_to_return.wait(timeout) + + completed = threading.Event() + completed.set() + release = DelayedRelease() + hook = ProcessHook( + enabled=threading.Event(), + reached=threading.Event(), + release=release, + completed=completed, + metadata=queue.Queue(maxsize=1), + ) + hook.arm() + consumer = threading.Thread(target=hook.hit, args=({"generation": 1},)) + consumer.start() + assert hook.reached.wait(2) + assert hook.metadata.get(timeout=1) == {"generation": 1} + + release.set() + rearmed = threading.Event() + + def rearm() -> None: + hook.arm() + rearmed.set() + + controller = threading.Thread(target=rearm) + controller.start() + try: + assert not rearmed.wait(0.05) + finally: + release.allow_wait_to_return.set() + consumer.join(timeout=2) + controller.join(timeout=2) + assert not consumer.is_alive() + assert rearmed.is_set() + + +def test_run_proof_measures_session_hygiene_after_worker_polling_stops( + monkeypatch: pytest.MonkeyPatch, + tmp_path: Path, +) -> None: + events: list[str] = [] + + class Worker: + def __init__(self, worker_id: str) -> None: + self.worker_id = worker_id + + def start_claiming(self) -> None: + events.append(f"{self.worker_id}-started") + + class Controller: + def __init__(self, **_kwargs: object) -> None: + self.workers = (Worker("worker-a"), Worker("worker-b")) + self.bundle = SimpleNamespace( + metadata=SimpleNamespace(schema_versions={"run": "test"}) + ) + + @property + def worker_a(self) -> Worker: + return self.workers[0] + + @property + def worker_b(self) -> Worker: + return self.workers[1] + + def start_workers(self) -> None: + events.append("workers-started") + + def stop_workers(self) -> None: + events.append("workers-stopped") + + class Provider: + url = "http://127.0.0.1:1" + + def __init__(self, **_kwargs: object) -> None: + pass + + def start(self) -> None: + events.append("provider-started") + + def stop(self) -> None: + events.append("provider-stopped") + + def count_idle_sessions(dsn: str) -> int: + assert dsn == "sanitized-dsn" + assert events[-1] == "workers-stopped" + events.append("sessions-measured") + return 0 + + monkeypatch.setattr(p5_proof, "ARTIFACT_PATH", tmp_path / "proof.json") + monkeypatch.setattr(p5_proof, "make_conninfo", lambda dsn, **_kwargs: dsn) + monkeypatch.setattr(p5_proof, "ProviderProcess", Provider) + monkeypatch.setattr(p5_proof, "ProofController", Controller) + monkeypatch.setattr( + p5_proof, + "bootstrap_postgres_application_schema", + lambda *_args, **_kwargs: None, + ) + monkeypatch.setattr(p5_proof, "idle_in_transaction_count", count_idle_sessions) + monkeypatch.setattr(p5_proof, "postgres_version", lambda _dsn: "test") + monkeypatch.setattr(p5_proof, "drop_schema", lambda _dsn, _schema: None) + monkeypatch.setattr(p5_proof, "_git_value", lambda *_args: "test") + monkeypatch.setattr(p5_proof, "assert_secret_safe", lambda *_args, **_kwargs: None) + + artifact = p5_proof.run_proof("sanitized-dsn", scenario_ids=()) + + assert artifact == tmp_path / "proof.json" + assert events.index("workers-stopped") < events.index("sessions-measured") + assert events.count("workers-stopped") == 1 + + def test_process_pause_hook_is_inert_until_explicitly_armed() -> None: context = multiprocessing.get_context("spawn") hooks = P5ProofHooks.create(context, ("paused",)) From 38a28ab95412c318a51a399b4af225eae4e0268a Mon Sep 17 00:00:00 2001 From: YingzuoLiu <156630675+YingzuoLiu@users.noreply.github.com> Date: Wed, 26 Aug 2026 18:40:59 +0800 Subject: [PATCH 2/8] test: add bounded P5 worker stack diagnostics --- examples/p5_multi_worker_proof.py | 21 ++++++++++++++++++ examples/p5_proof_worker.py | 4 ++++ tests/test_p5_multi_worker_proof.py | 34 +++++++++++++++++++++++++++++ 3 files changed, 59 insertions(+) diff --git a/examples/p5_multi_worker_proof.py b/examples/p5_multi_worker_proof.py index 061af7a..6ba6747 100644 --- a/examples/p5_multi_worker_proof.py +++ b/examples/p5_multi_worker_proof.py @@ -258,6 +258,7 @@ def wait_hook(self, name: str, *, timeout: float = WAIT_SECONDS) -> dict[str, An while not hook.reached.wait(0.05): self.raise_if_failed() if time.monotonic() >= deadline: + self.dump_thread_stacks(reason=f"timeout waiting for {name}") raise P5ProofFailure(f"{self.worker_id} did not reach {name}") self.raise_if_failed() try: @@ -267,6 +268,26 @@ def wait_hook(self, name: str, *, timeout: float = WAIT_SECONDS) -> dict[str, An _require(isinstance(payload, dict), f"{self.worker_id} {name} metadata is invalid") return payload + def dump_thread_stacks(self, *, reason: str) -> None: + """Ask a live proof worker for bounded, locals-free stack diagnostics.""" + + if self.process is None or not self.process.is_alive(): + return + print( + json.dumps( + { + "diagnostic": "p5-worker-thread-stacks", + "reason": reason, + "worker": self.worker_id, + }, + sort_keys=True, + ), + file=sys.stderr, + flush=True, + ) + os.kill(self.pid, signal.SIGUSR1) + time.sleep(0.1) + def release(self, name: str) -> None: self.hooks.hooks[name].release.set() diff --git a/examples/p5_proof_worker.py b/examples/p5_proof_worker.py index c8de326..2287db8 100644 --- a/examples/p5_proof_worker.py +++ b/examples/p5_proof_worker.py @@ -1,8 +1,10 @@ from __future__ import annotations +import faulthandler import os import queue import signal +import sys import threading import time from dataclasses import dataclass, field @@ -695,6 +697,7 @@ def run_p5_worker( manager: RuntimeManager | None = None previous_thread_excepthook = threading.excepthook + faulthandler.register(signal.SIGUSR1, file=sys.stderr, all_threads=True) def report_thread_failure(args: threading.ExceptHookArgs) -> None: failures.put( @@ -765,5 +768,6 @@ def report_thread_failure(args: threading.ExceptHookArgs) -> None: finally: if manager is not None: manager.stop() + faulthandler.unregister(signal.SIGUSR1) threading.excepthook = previous_thread_excepthook time.sleep(0.01) diff --git a/tests/test_p5_multi_worker_proof.py b/tests/test_p5_multi_worker_proof.py index 092389d..448b799 100644 --- a/tests/test_p5_multi_worker_proof.py +++ b/tests/test_p5_multi_worker_proof.py @@ -126,6 +126,40 @@ def test_worker_config_repr_redacts_postgres_dsn() -> None: assert "dsn=" not in repr(config) +def test_worker_timeout_diagnostic_requests_locals_free_stack_dump( + monkeypatch: pytest.MonkeyPatch, + capsys: pytest.CaptureFixture[str], +) -> None: + class LiveProcess: + pid = 4242 + + @staticmethod + def is_alive() -> bool: + return True + + worker = p5_proof.WorkerProcess( + context=None, + config=P5WorkerConfig( + worker_id="p5-worker-a", + schema="p5_test", + dsn="postgresql://proof:canary-password@localhost/proof", + ), + hooks=P5ProofHooks({}), + process=LiveProcess(), + ) + signals: list[tuple[int, int]] = [] + monkeypatch.setattr(p5_proof.os, "kill", lambda pid, sig: signals.append((pid, sig))) + monkeypatch.setattr(p5_proof.time, "sleep", lambda _seconds: None) + + worker.dump_thread_stacks(reason="timeout waiting for claim.before") + + diagnostic = capsys.readouterr().err + assert signals == [(4242, signal.SIGUSR1)] + assert '"worker": "p5-worker-a"' in diagnostic + assert '"reason": "timeout waiting for claim.before"' in diagnostic + assert "canary-password" not in diagnostic + + def test_process_hook_is_one_shot_and_carries_only_controller_payload() -> None: context = multiprocessing.get_context("spawn") hooks = P5ProofHooks.create(context, ("point",)) From 9ca262c494acd91c18c593286216a92ad010798b Mon Sep 17 00:00:00 2001 From: YingzuoLiu <156630675+YingzuoLiu@users.noreply.github.com> Date: Wed, 26 Aug 2026 18:51:10 +0800 Subject: [PATCH 3/8] fix: make P5 proof barriers generation-safe --- docs/p5-multi-worker-recovery-proof.md | 13 +-- examples/p5_multi_worker_proof.py | 14 ++- examples/p5_proof_worker.py | 114 ++++++++++++++++++++---- tests/test_p5_multi_worker_proof.py | 118 +++++++++++++++---------- 4 files changed, 187 insertions(+), 72 deletions(-) diff --git a/docs/p5-multi-worker-recovery-proof.md b/docs/p5-multi-worker-recovery-proof.md index 167cc7c..489914d 100644 --- a/docs/p5-multi-worker-recovery-proof.md +++ b/docs/p5-multi-worker-recovery-proof.md @@ -7,11 +7,14 @@ controller process submits Runs, controls named one-shot barriers, injects exact store-time lease expiry, and rereads durable evidence. The synthetic HTTP provider is a fourth, independent process with its own SQLite effect ledger. -Each reusable barrier has an explicit consumer-completed acknowledgement. The -controller cannot clear and re-arm a barrier until the worker has observed the prior -release, preventing one generation from consuming the next generation's signal. The -post-run PostgreSQL session check runs only after both polling workers have stopped; -an in-flight polling transaction is not misclassified as a leaked session. +Each reusable barrier has monotonic armed, consumed, reached, released, and completed +generation counters. Events are wake-up hints only: a stale event cannot satisfy a +different generation, and the controller cannot re-arm until the prior generation's +consumer has acknowledged completion. Timeout diagnostics request locals-free worker +thread stacks so a failed barrier identifies the exact blocking boundary without +publishing credentials or payloads. The post-run PostgreSQL session check runs only +after both polling workers have stopped; an in-flight polling transaction is not +misclassified as a leaked session. The proof adds no production scheduler, distributed queue, connection pool, or production control endpoint. `RuntimeManager._wake` remains a process-local latency diff --git a/examples/p5_multi_worker_proof.py b/examples/p5_multi_worker_proof.py index 6ba6747..3f6756c 100644 --- a/examples/p5_multi_worker_proof.py +++ b/examples/p5_multi_worker_proof.py @@ -254,17 +254,23 @@ def arm(self, *names: str) -> None: def wait_hook(self, name: str, *, timeout: float = WAIT_SECONDS) -> dict[str, Any]: hook = self.hooks.hooks[name] + generation = hook.current_generation() deadline = time.monotonic() + timeout - while not hook.reached.wait(0.05): + while hook.reached_generation() < generation: + hook.reached.wait(0.05) self.raise_if_failed() if time.monotonic() >= deadline: self.dump_thread_stacks(reason=f"timeout waiting for {name}") raise P5ProofFailure(f"{self.worker_id} did not reach {name}") self.raise_if_failed() try: - payload = hook.metadata.get(timeout=1) + payload_generation, payload = hook.metadata.get(timeout=1) except queue.Empty: raise P5ProofFailure(f"{self.worker_id} {name} metadata is missing") from None + _require( + payload_generation == generation, + f"{self.worker_id} {name} metadata generation is invalid", + ) _require(isinstance(payload, dict), f"{self.worker_id} {name} metadata is invalid") return payload @@ -289,11 +295,11 @@ def dump_thread_stacks(self, *, reason: str) -> None: time.sleep(0.1) def release(self, name: str) -> None: - self.hooks.hooks[name].release.set() + self.hooks.hooks[name].release_current() def release_all(self) -> None: for hook in self.hooks.hooks.values(): - hook.release.set() + hook.release_current() def suspend(self) -> None: os.kill(self.pid, signal.SIGSTOP) diff --git a/examples/p5_proof_worker.py b/examples/p5_proof_worker.py index 2287db8..22bc2f6 100644 --- a/examples/p5_proof_worker.py +++ b/examples/p5_proof_worker.py @@ -87,6 +87,13 @@ class ProcessHook: release: Any completed: Any metadata: Any + generations: Any + + _ARMED = 0 + _CONSUMED = 1 + _REACHED = 2 + _RELEASED = 3 + _COMPLETED = 4 @classmethod def create(cls, context: Any) -> ProcessHook: @@ -98,11 +105,30 @@ def create(cls, context: Any) -> ProcessHook: release=context.Event(), completed=completed, metadata=context.Queue(maxsize=1), + generations=context.Array("Q", 5, lock=True), ) - def arm(self) -> None: - if not self.completed.wait(P5_HOOK_TIMEOUT_SECONDS): - raise TimeoutError("P5 previous proof hook consumer did not finish") + def _generation(self, index: int) -> int: + with self.generations.get_lock(): + return int(self.generations[index]) + + def current_generation(self) -> int: + return self._generation(self._ARMED) + + def reached_generation(self) -> int: + return self._generation(self._REACHED) + + def arm(self) -> int: + deadline = time.monotonic() + P5_HOOK_TIMEOUT_SECONDS + while True: + with self.generations.get_lock(): + previous = int(self.generations[self._ARMED]) + previous_completed = int(self.generations[self._COMPLETED]) + if previous_completed >= previous: + break + remaining = deadline - time.monotonic() + if remaining <= 0 or not self.completed.wait(remaining): + raise TimeoutError("P5 previous proof hook consumer did not finish") self.completed.clear() self.reached.clear() self.release.clear() @@ -111,45 +137,92 @@ def arm(self) -> None: self.metadata.get_nowait() except queue.Empty: break + with self.generations.get_lock(): + generation = int(self.generations[self._ARMED]) + 1 + self.generations[self._ARMED] = generation self.enabled.set() + return generation - def consume(self) -> bool: + def consume(self) -> int | None: if not self.enabled.is_set(): - return False + return None + with self.generations.get_lock(): + generation = int(self.generations[self._ARMED]) + if int(self.generations[self._CONSUMED]) >= generation: + return None + self.generations[self._CONSUMED] = generation self.enabled.clear() - return True + return generation + + def _mark_reached(self, generation: int, payload: dict[str, Any]) -> None: + self.metadata.put((generation, payload)) + with self.generations.get_lock(): + self.generations[self._REACHED] = max( + int(self.generations[self._REACHED]), + generation, + ) + self.reached.set() - def block_consumed(self, payload: dict[str, Any]) -> None: + def _mark_completed(self, generation: int) -> None: + with self.generations.get_lock(): + self.generations[self._COMPLETED] = max( + int(self.generations[self._COMPLETED]), + generation, + ) + self.completed.set() + + def release_current(self) -> None: + with self.generations.get_lock(): + generation = int(self.generations[self._ARMED]) + self.generations[self._RELEASED] = max( + int(self.generations[self._RELEASED]), + generation, + ) + self.release.set() + + def block_consumed(self, generation: int, payload: dict[str, Any]) -> None: try: - self.metadata.put(payload) - self.reached.set() - if not self.release.wait(P5_HOOK_TIMEOUT_SECONDS): - raise TimeoutError("P5 controller did not release a reached proof hook") + self._mark_reached(generation, payload) + deadline = time.monotonic() + P5_HOOK_TIMEOUT_SECONDS + while self._generation(self._RELEASED) < generation: + remaining = deadline - time.monotonic() + if remaining <= 0: + raise TimeoutError("P5 controller did not release a reached proof hook") + self.release.wait(remaining) + if self._generation(self._RELEASED) < generation: + self.release.clear() finally: - self.completed.set() + self._mark_completed(generation) def hit(self, payload: dict[str, Any]) -> bool: - if not self.consume(): + generation = self.consume() + if generation is None: return False - self.block_consumed(payload) + self.block_consumed(generation, payload) return True def pause_process(self, payload: dict[str, Any]) -> bool: """Stop this proof process after publishing one bounded observation.""" - if not self.consume(): + generation = self.consume() + if generation is None: return False try: - self.metadata.put(payload) + self.metadata.put((generation, payload)) # multiprocessing.Queue publishes through a feeder thread. Flush this # one-shot channel before SIGSTOP so the controller never observes the # reached event before its matching payload is readable. self.metadata.close() self.metadata.join_thread() + with self.generations.get_lock(): + self.generations[self._REACHED] = max( + int(self.generations[self._REACHED]), + generation, + ) self.reached.set() os.kill(os.getpid(), signal.SIGSTOP) finally: - self.completed.set() + self._mark_completed(generation) return True @@ -192,6 +265,7 @@ def __init__( self, connection: Any, hook: ProcessHook, + hook_generation: int, failure_hook: ProcessHook | None, *, worker_id: str, @@ -199,6 +273,7 @@ def __init__( ) -> None: self._connection = connection self._hook = hook + self._hook_generation = hook_generation self._failure_hook = failure_hook self._worker_id = worker_id self._backend_pid = backend_pid @@ -208,6 +283,7 @@ def execute(self, query: Any, params: Any = None) -> Any: if self._first_execute: self._first_execute = False self._hook.block_consumed( + self._hook_generation, { "point": "db.connection.open", "worker": self._worker_id, @@ -260,7 +336,8 @@ def _install_connection_fault_hook(self) -> None: def connect_with_optional_fault() -> Any: connection = original() hook = self.hooks.hooks.get("db.connection.open") - if hook is None or not hook.consume(): + generation = hook.consume() if hook is not None else None + if hook is None or generation is None: return connection row = connection.execute("SELECT pg_backend_pid() AS pid").fetchone() if row is None: @@ -269,6 +346,7 @@ def connect_with_optional_fault() -> Any: return _TerminatedConnection( connection, hook, + generation, self.hooks.hooks.get("db.connection.failed"), worker_id=self.worker_id, backend_pid=int(row["pid"]), diff --git a/tests/test_p5_multi_worker_proof.py b/tests/test_p5_multi_worker_proof.py index 448b799..7a2ddd2 100644 --- a/tests/test_p5_multi_worker_proof.py +++ b/tests/test_p5_multi_worker_proof.py @@ -3,7 +3,6 @@ import multiprocessing import json import os -import queue import signal import threading from pathlib import Path @@ -163,7 +162,7 @@ def is_alive() -> bool: def test_process_hook_is_one_shot_and_carries_only_controller_payload() -> None: context = multiprocessing.get_context("spawn") hooks = P5ProofHooks.create(context, ("point",)) - hooks.arm("point") + generation = hooks.hooks["point"].arm() completed = threading.Event() def hit() -> None: @@ -175,46 +174,38 @@ def hit() -> None: thread.start() hook = hooks.hooks["point"] assert hook.reached.wait(2) - assert hook.metadata.get(timeout=1) == {"point": "point", "attempt": 1} - hook.release.set() + assert hook.metadata.get(timeout=1) == ( + generation, + {"point": "point", "attempt": 1}, + ) + hook.release_current() thread.join(timeout=2) assert completed.is_set() -def test_process_hook_rearm_waits_for_previous_consumer_acknowledgement() -> None: - class DelayedRelease: - def __init__(self) -> None: - self._set = threading.Event() - self.allow_wait_to_return = threading.Event() - - def clear(self) -> None: - self._set.clear() - - def set(self) -> None: - self._set.set() - - def wait(self, timeout: float) -> bool: - if not self._set.wait(timeout): - return False - return self.allow_wait_to_return.wait(timeout) - - completed = threading.Event() - completed.set() - release = DelayedRelease() - hook = ProcessHook( - enabled=threading.Event(), - reached=threading.Event(), - release=release, - completed=completed, - metadata=queue.Queue(maxsize=1), - ) - hook.arm() +def test_process_hook_rearm_waits_for_previous_consumer_acknowledgement( + monkeypatch: pytest.MonkeyPatch, +) -> None: + context = multiprocessing.get_context("spawn") + hook = ProcessHook.create(context) + generation = hook.arm() + completion_entered = threading.Event() + allow_completion = threading.Event() + original_mark_completed = hook._mark_completed + + def delayed_completion(completed_generation: int) -> None: + completion_entered.set() + assert allow_completion.wait(2) + original_mark_completed(completed_generation) + + monkeypatch.setattr(hook, "_mark_completed", delayed_completion) consumer = threading.Thread(target=hook.hit, args=({"generation": 1},)) consumer.start() assert hook.reached.wait(2) - assert hook.metadata.get(timeout=1) == {"generation": 1} + assert hook.metadata.get(timeout=1) == (generation, {"generation": 1}) - release.set() + hook.release_current() + assert completion_entered.wait(2) rearmed = threading.Event() def rearm() -> None: @@ -226,13 +217,41 @@ def rearm() -> None: try: assert not rearmed.wait(0.05) finally: - release.allow_wait_to_return.set() + allow_completion.set() consumer.join(timeout=2) controller.join(timeout=2) assert not consumer.is_alive() assert rearmed.is_set() +def test_process_hook_stale_release_signal_cannot_release_new_generation() -> None: + context = multiprocessing.get_context("spawn") + hook = ProcessHook.create(context) + + first_generation = hook.arm() + first = threading.Thread(target=hook.hit, args=({"generation": 1},)) + first.start() + assert hook.reached.wait(2) + assert hook.metadata.get(timeout=1) == (first_generation, {"generation": 1}) + hook.release_current() + first.join(timeout=2) + assert not first.is_alive() + + second_generation = hook.arm() + second = threading.Thread(target=hook.hit, args=({"generation": 2},)) + second.start() + assert hook.reached.wait(2) + assert hook.metadata.get(timeout=1) == (second_generation, {"generation": 2}) + + hook.release.set() + second.join(timeout=0.05) + assert second.is_alive() + + hook.release_current() + second.join(timeout=2) + assert not second.is_alive() + + def test_run_proof_measures_session_hygiene_after_worker_polling_stops( monkeypatch: pytest.MonkeyPatch, tmp_path: Path, @@ -325,7 +344,10 @@ def test_process_pause_hook_flushes_metadata_before_sigstop() -> None: try: hook = hooks.hooks["paused"] assert hook.reached.wait(5) - assert hook.metadata.get(timeout=1) == {"point": "paused", "attempt": 1} + assert hook.metadata.get(timeout=1) == ( + hook.current_generation(), + {"point": "paused", "attempt": 1}, + ) finally: if process.is_alive() and process.pid is not None: os.kill(process.pid, signal.SIGCONT) @@ -352,11 +374,14 @@ def execute(self, _query, _params=None): ) open_hook = hooks.hooks["db.connection.open"] failure_hook = hooks.hooks["db.connection.failed"] - open_hook.release.set() - failure_hook.enabled.set() + open_generation = open_hook.arm() + assert open_hook.consume() == open_generation + open_hook.release_current() + failure_generation = failure_hook.arm() proxy = _TerminatedConnection( BrokenConnection(), open_hook, + open_generation, failure_hook, worker_id="p5-worker-a", backend_pid=123, @@ -371,12 +396,15 @@ def execute() -> None: thread = threading.Thread(target=execute) thread.start() assert failure_hook.reached.wait(2) - assert failure_hook.metadata.get(timeout=1) == { - "point": "db.connection.failed", - "worker": "p5-worker-a", - "error_type": "AdminShutdown", - "sqlstate": "57P01", - } - failure_hook.release.set() + assert failure_hook.metadata.get(timeout=1) == ( + failure_generation, + { + "point": "db.connection.failed", + "worker": "p5-worker-a", + "error_type": "AdminShutdown", + "sqlstate": "57P01", + }, + ) + failure_hook.release_current() thread.join(timeout=2) assert observed.is_set() From 389ad352a430ebb706964adbbc60a72d1b6bad29 Mon Sep 17 00:00:00 2001 From: YingzuoLiu <156630675+YingzuoLiu@users.noreply.github.com> Date: Wed, 26 Aug 2026 18:56:29 +0800 Subject: [PATCH 4/8] test: report P5 barrier generations on timeout --- examples/p5_multi_worker_proof.py | 16 ++++++++++++++-- examples/p5_proof_worker.py | 10 ++++++++++ 2 files changed, 24 insertions(+), 2 deletions(-) diff --git a/examples/p5_multi_worker_proof.py b/examples/p5_multi_worker_proof.py index 3f6756c..9b39252 100644 --- a/examples/p5_multi_worker_proof.py +++ b/examples/p5_multi_worker_proof.py @@ -260,7 +260,13 @@ def wait_hook(self, name: str, *, timeout: float = WAIT_SECONDS) -> dict[str, An hook.reached.wait(0.05) self.raise_if_failed() if time.monotonic() >= deadline: - self.dump_thread_stacks(reason=f"timeout waiting for {name}") + self.dump_thread_stacks( + reason=f"timeout waiting for {name}", + hook_states={ + hook_name: candidate.generation_state() + for hook_name, candidate in self.hooks.hooks.items() + }, + ) raise P5ProofFailure(f"{self.worker_id} did not reach {name}") self.raise_if_failed() try: @@ -274,7 +280,12 @@ def wait_hook(self, name: str, *, timeout: float = WAIT_SECONDS) -> dict[str, An _require(isinstance(payload, dict), f"{self.worker_id} {name} metadata is invalid") return payload - def dump_thread_stacks(self, *, reason: str) -> None: + def dump_thread_stacks( + self, + *, + reason: str, + hook_states: dict[str, dict[str, int]] | None = None, + ) -> None: """Ask a live proof worker for bounded, locals-free stack diagnostics.""" if self.process is None or not self.process.is_alive(): @@ -283,6 +294,7 @@ def dump_thread_stacks(self, *, reason: str) -> None: json.dumps( { "diagnostic": "p5-worker-thread-stacks", + "hook_states": hook_states or {}, "reason": reason, "worker": self.worker_id, }, diff --git a/examples/p5_proof_worker.py b/examples/p5_proof_worker.py index 22bc2f6..e922437 100644 --- a/examples/p5_proof_worker.py +++ b/examples/p5_proof_worker.py @@ -118,6 +118,16 @@ def current_generation(self) -> int: def reached_generation(self) -> int: return self._generation(self._REACHED) + def generation_state(self) -> dict[str, int]: + with self.generations.get_lock(): + return { + "armed": int(self.generations[self._ARMED]), + "consumed": int(self.generations[self._CONSUMED]), + "reached": int(self.generations[self._REACHED]), + "released": int(self.generations[self._RELEASED]), + "completed": int(self.generations[self._COMPLETED]), + } + def arm(self) -> int: deadline = time.monotonic() + P5_HOOK_TIMEOUT_SECONDS while True: From d685863e6938d499af2f893b387697c021e763d5 Mon Sep 17 00:00:00 2001 From: YingzuoLiu <156630675+YingzuoLiu@users.noreply.github.com> Date: Wed, 26 Aug 2026 19:14:19 +0800 Subject: [PATCH 5/8] fix: drain P5 barriers between schedules --- docs/p5-multi-worker-recovery-proof.md | 13 +++-- examples/p5_multi_worker_proof.py | 11 ++++ examples/p5_proof_worker.py | 5 ++ tests/test_p5_multi_worker_proof.py | 80 ++++++++++++++++++++++++++ 4 files changed, 104 insertions(+), 5 deletions(-) diff --git a/docs/p5-multi-worker-recovery-proof.md b/docs/p5-multi-worker-recovery-proof.md index 489914d..3b75161 100644 --- a/docs/p5-multi-worker-recovery-proof.md +++ b/docs/p5-multi-worker-recovery-proof.md @@ -10,11 +10,14 @@ a fourth, independent process with its own SQLite effect ledger. Each reusable barrier has monotonic armed, consumed, reached, released, and completed generation counters. Events are wake-up hints only: a stale event cannot satisfy a different generation, and the controller cannot re-arm until the prior generation's -consumer has acknowledged completion. Timeout diagnostics request locals-free worker -thread stacks so a failed barrier identifies the exact blocking boundary without -publishing credentials or payloads. The post-run PostgreSQL session check runs only -after both polling workers have stopped; an in-flight polling transaction is not -misclassified as a leaked session. +consumer has acknowledged completion. Starting a new schedule first releases every +prior named barrier, so a worker parked on one hook cannot prevent it from reaching a +different hook in the next schedule. Timeout diagnostics request locals-free worker +thread stacks and integer-only generation state so a failed barrier identifies the +exact blocking boundary without publishing credentials or payloads. Failure cleanup +resumes a stopped proof worker and has a bounded SIGKILL backstop. The post-run +PostgreSQL session check runs only after both polling workers have stopped; an +in-flight polling transaction is not misclassified as a leaked session. The proof adds no production scheduler, distributed queue, connection pool, or production control endpoint. `RuntimeManager._wake` remains a process-local latency diff --git a/examples/p5_multi_worker_proof.py b/examples/p5_multi_worker_proof.py index 9b39252..7710194 100644 --- a/examples/p5_multi_worker_proof.py +++ b/examples/p5_multi_worker_proof.py @@ -334,8 +334,19 @@ def stop(self) -> None: self.shutdown_event.set() self.process.join(timeout=5) if self.process.is_alive(): + # SIGTERM remains pending for a SIGSTOPed proof process. Resume it + # before bounded termination, then retain SIGKILL as the final + # proof-owned cleanup backstop. + try: + self.resume() + except ProcessLookupError: + pass self.process.terminate() self.process.join(timeout=5) + if self.process.is_alive(): + self.process.kill() + self.process.join(timeout=5) + _require(not self.process.is_alive(), f"{self.worker_id} did not stop") def raise_if_failed(self) -> None: if self.failures is not None: diff --git a/examples/p5_proof_worker.py b/examples/p5_proof_worker.py index e922437..146fd71 100644 --- a/examples/p5_proof_worker.py +++ b/examples/p5_proof_worker.py @@ -245,6 +245,11 @@ def create(cls, context: Any, names: tuple[str, ...]) -> P5ProofHooks: return cls({name: ProcessHook.create(context) for name in names}) def arm(self, *names: str) -> None: + # Arming starts a new deterministic schedule. Drain any barrier left by + # the previous schedule first so a worker cannot remain parked on a + # different hook while the controller waits for the newly armed one. + for hook in self.hooks.values(): + hook.release_current() for name in names: self.hooks[name].arm() diff --git a/tests/test_p5_multi_worker_proof.py b/tests/test_p5_multi_worker_proof.py index 7a2ddd2..552e075 100644 --- a/tests/test_p5_multi_worker_proof.py +++ b/tests/test_p5_multi_worker_proof.py @@ -5,6 +5,7 @@ import os import signal import threading +import time from pathlib import Path from types import SimpleNamespace @@ -34,6 +35,24 @@ def _pause_hook_in_child(hooks: P5ProofHooks) -> None: hooks.pause_process("paused", {"point": "paused", "attempt": 1}) +def _repeat_hook_in_child(hook: ProcessHook, count: int) -> None: + for expected in range(1, count + 1): + deadline = time.monotonic() + 5 + while not hook.hit({"generation": expected}): + if time.monotonic() >= deadline: + raise TimeoutError(f"generation {expected} was never armed") + time.sleep(0.001) + + +def _hit_hook_sequence_in_child(hooks: P5ProofHooks, names: tuple[str, ...]) -> None: + for name in names: + deadline = time.monotonic() + 5 + while not hooks.hit(name, {"point": name}): + if time.monotonic() >= deadline: + raise TimeoutError(f"hook {name} was never armed") + time.sleep(0.001) + + def _valid_report() -> dict: return { "proof": REPORT_VERSION, @@ -252,6 +271,67 @@ def test_process_hook_stale_release_signal_cannot_release_new_generation() -> No assert not second.is_alive() +def test_process_hook_rearms_across_spawned_process_generations() -> None: + context = multiprocessing.get_context("spawn") + hook = ProcessHook.create(context) + count = 200 + process = context.Process(target=_repeat_hook_in_child, args=(hook, count)) + process.start() + + try: + for expected in range(1, count + 1): + generation = hook.arm() + assert hook.reached.wait(5) + assert hook.reached_generation() == generation + assert hook.metadata.get(timeout=1) == ( + generation, + {"generation": expected}, + ) + hook.release_current() + process.join(timeout=5) + finally: + if process.is_alive(): + process.kill() + process.join(timeout=2) + + assert process.exitcode == 0 + + +def test_arming_new_schedule_drains_worker_blocked_on_different_hook() -> None: + context = multiprocessing.get_context("spawn") + hooks = P5ProofHooks.create(context, ("previous", "next")) + hooks.arm("previous") + process = context.Process( + target=_hit_hook_sequence_in_child, + args=(hooks, ("previous", "next")), + ) + process.start() + + try: + previous = hooks.hooks["previous"] + assert previous.reached.wait(5) + previous_generation, previous_payload = previous.metadata.get(timeout=1) + assert previous_generation == previous.current_generation() + assert previous_payload == {"point": "previous"} + + hooks.arm("next") + following = hooks.hooks["next"] + assert following.reached.wait(5) + following_generation, following_payload = following.metadata.get(timeout=1) + assert following_generation == following.current_generation() + assert following_payload == {"point": "next"} + following.release_current() + process.join(timeout=5) + finally: + for hook in hooks.hooks.values(): + hook.release_current() + if process.is_alive(): + process.kill() + process.join(timeout=2) + + assert process.exitcode == 0 + + def test_run_proof_measures_session_hygiene_after_worker_polling_stops( monkeypatch: pytest.MonkeyPatch, tmp_path: Path, From 376b41e2ebbfe97fcdbe6800acbdc4ac7dbf535d Mon Sep 17 00:00:00 2001 From: YingzuoLiu <156630675+YingzuoLiu@users.noreply.github.com> Date: Wed, 26 Aug 2026 19:27:41 +0800 Subject: [PATCH 6/8] test: bound and diagnose P5 proof hangs --- .github/workflows/ci.yml | 4 +++- examples/p5_multi_worker_proof.py | 3 +++ tests/test_p5_multi_worker_proof.py | 12 ++++++++++++ 3 files changed, 18 insertions(+), 1 deletion(-) diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index 0ab3d4e..d245fc5 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -177,7 +177,9 @@ jobs: - run: python -m pip install --upgrade pip - run: pip install -r requirements-dev.txt - name: Two-process PostgreSQL recovery proof (required; never skipped) - run: python examples/p5_multi_worker_proof.py + run: >- + timeout --signal=TERM --kill-after=15s 300s + python -u examples/p5_multi_worker_proof.py - name: Recheck sanitized proof artifact run: >- python -c "import json, os, pathlib; diff --git a/examples/p5_multi_worker_proof.py b/examples/p5_multi_worker_proof.py index 7710194..dd5092f 100644 --- a/examples/p5_multi_worker_proof.py +++ b/examples/p5_multi_worker_proof.py @@ -1,6 +1,7 @@ from __future__ import annotations import argparse +import faulthandler import json import multiprocessing import os @@ -1369,6 +1370,7 @@ def run_proof( *, scenario_ids: tuple[str, ...] = ALL_SCENARIO_IDS, ) -> Path: + faulthandler.dump_traceback_later(60, repeat=True, file=sys.stderr) ARTIFACT_PATH.unlink(missing_ok=True) schema = f"p5_{uuid.uuid4().hex[:20]}" proof_dsn = make_conninfo(dsn, application_name="p5_multi_worker_proof") @@ -1449,6 +1451,7 @@ def run_proof( encoding="utf-8", ) finally: + faulthandler.cancel_dump_traceback_later() if controller is not None and not workers_stopped: controller.stop_workers() provider.stop() diff --git a/tests/test_p5_multi_worker_proof.py b/tests/test_p5_multi_worker_proof.py index 552e075..4bc5e6b 100644 --- a/tests/test_p5_multi_worker_proof.py +++ b/tests/test_p5_multi_worker_proof.py @@ -398,12 +398,24 @@ def count_idle_sessions(dsn: str) -> int: monkeypatch.setattr(p5_proof, "drop_schema", lambda _dsn, _schema: None) monkeypatch.setattr(p5_proof, "_git_value", lambda *_args: "test") monkeypatch.setattr(p5_proof, "assert_secret_safe", lambda *_args, **_kwargs: None) + monkeypatch.setattr( + p5_proof.faulthandler, + "dump_traceback_later", + lambda *_args, **_kwargs: events.append("watchdog-armed"), + ) + monkeypatch.setattr( + p5_proof.faulthandler, + "cancel_dump_traceback_later", + lambda: events.append("watchdog-cancelled"), + ) artifact = p5_proof.run_proof("sanitized-dsn", scenario_ids=()) assert artifact == tmp_path / "proof.json" assert events.index("workers-stopped") < events.index("sessions-measured") assert events.count("workers-stopped") == 1 + assert events.index("watchdog-armed") < events.index("provider-started") + assert events.index("sessions-measured") < events.index("watchdog-cancelled") def test_process_pause_hook_is_inert_until_explicitly_armed() -> None: From ea488be424008cb2f3ae1ee515abd30ea5749354 Mon Sep 17 00:00:00 2001 From: YingzuoLiu <156630675+YingzuoLiu@users.noreply.github.com> Date: Wed, 26 Aug 2026 19:35:22 +0800 Subject: [PATCH 7/8] fix: serialize P5 hook arming with claim cycles --- docs/p5-multi-worker-recovery-proof.md | 16 +++--- examples/p5_proof_worker.py | 28 ++++++++-- tests/test_p5_multi_worker_proof.py | 75 ++++++++++++++++++++++++++ 3 files changed, 110 insertions(+), 9 deletions(-) diff --git a/docs/p5-multi-worker-recovery-proof.md b/docs/p5-multi-worker-recovery-proof.md index 3b75161..4970f71 100644 --- a/docs/p5-multi-worker-recovery-proof.md +++ b/docs/p5-multi-worker-recovery-proof.md @@ -12,12 +12,16 @@ generation counters. Events are wake-up hints only: a stale event cannot satisfy different generation, and the controller cannot re-arm until the prior generation's consumer has acknowledged completion. Starting a new schedule first releases every prior named barrier, so a worker parked on one hook cannot prevent it from reaching a -different hook in the next schedule. Timeout diagnostics request locals-free worker -thread stacks and integer-only generation state so a failed barrier identifies the -exact blocking boundary without publishing credentials or payloads. Failure cleanup -resumes a stopped proof worker and has a bounded SIGKILL backstop. The post-run -PostgreSQL session check runs only after both polling workers have stopped; an -in-flight polling transaction is not misclassified as a leaked session. +different hook in the next schedule. Arming the new schedule is also serialized with +the proof adapter's whole `claim_next_run` cycle. A poll already between +`claim.before` and `claim.result` must finish before either new hook becomes visible, +so one call cannot miss the new before generation and consume its result generation. +Timeout diagnostics request locals-free worker thread stacks and integer-only +generation state so a failed barrier identifies the exact blocking boundary without +publishing credentials or payloads. Failure cleanup resumes a stopped proof worker +and has a bounded SIGKILL backstop. The post-run PostgreSQL session check runs only +after both polling workers have stopped; an in-flight polling transaction is not +misclassified as a leaked session. The proof adds no production scheduler, distributed queue, connection pool, or production control endpoint. `RuntimeManager._wake` remains a process-local latency diff --git a/examples/p5_proof_worker.py b/examples/p5_proof_worker.py index 146fd71..1ef726f 100644 --- a/examples/p5_proof_worker.py +++ b/examples/p5_proof_worker.py @@ -239,10 +239,14 @@ def pause_process(self, payload: dict[str, Any]) -> bool: @dataclass class P5ProofHooks: hooks: dict[str, ProcessHook] + claim_cycle: Any = field(default_factory=threading.Lock) @classmethod def create(cls, context: Any, names: tuple[str, ...]) -> P5ProofHooks: - return cls({name: ProcessHook.create(context) for name in names}) + return cls( + {name: ProcessHook.create(context) for name in names}, + claim_cycle=context.Lock(), + ) def arm(self, *names: str) -> None: # Arming starts a new deterministic schedule. Drain any barrier left by @@ -250,8 +254,12 @@ def arm(self, *names: str) -> None: # different hook while the controller waits for the newly armed one. for hook in self.hooks.values(): hook.release_current() - for name in names: - self.hooks[name].arm() + # A worker already inside claim_next_run must finish that whole cycle + # before new before/result hooks become visible. Otherwise it can miss + # the new before hook and consume the new result hook from the same call. + with self.claim_cycle: + for name in names: + self.hooks[name].arm() def hit(self, name: str, payload: dict[str, Any]) -> bool: hook = self.hooks.get(name) @@ -379,6 +387,20 @@ def claim_next_run( owner_id: str, lease_duration_seconds: int, reconciliation_pending_code: str | None = None, + ) -> RunLeaseClaim | None: + with self.hooks.claim_cycle: + return self._claim_next_run_in_cycle( + owner_id=owner_id, + lease_duration_seconds=lease_duration_seconds, + reconciliation_pending_code=reconciliation_pending_code, + ) + + def _claim_next_run_in_cycle( + self, + *, + owner_id: str, + lease_duration_seconds: int, + reconciliation_pending_code: str | None = None, ) -> RunLeaseClaim | None: self.hooks.hit( "claim.before", diff --git a/tests/test_p5_multi_worker_proof.py b/tests/test_p5_multi_worker_proof.py index 4bc5e6b..fbbdbfd 100644 --- a/tests/test_p5_multi_worker_proof.py +++ b/tests/test_p5_multi_worker_proof.py @@ -8,6 +8,7 @@ import time from pathlib import Path from types import SimpleNamespace +from typing import Any import pytest @@ -53,6 +54,26 @@ def _hit_hook_sequence_in_child(hooks: P5ProofHooks, names: tuple[str, ...]) -> time.sleep(0.001) +def _straddle_claim_cycle_in_child( + hooks: P5ProofHooks, + between_hooks: Any, + continue_cycle: Any, + start_next_cycle: Any, +) -> None: + with hooks.claim_cycle: + assert hooks.hit("claim.before", {"cycle": 1}) is False + between_hooks.set() + if not continue_cycle.wait(5): + raise TimeoutError("controller did not release the in-flight claim cycle") + assert hooks.hit("claim.result", {"cycle": 1}) is False + + if not start_next_cycle.wait(5): + raise TimeoutError("controller did not arm the next claim cycle") + with hooks.claim_cycle: + assert hooks.hit("claim.before", {"cycle": 2}) is True + assert hooks.hit("claim.result", {"cycle": 2}) is True + + def _valid_report() -> dict: return { "proof": REPORT_VERSION, @@ -332,6 +353,60 @@ def test_arming_new_schedule_drains_worker_blocked_on_different_hook() -> None: assert process.exitcode == 0 +def test_arming_claim_schedule_cannot_straddle_inflight_claim_cycle() -> None: + context = multiprocessing.get_context("spawn") + hooks = P5ProofHooks.create(context, ("claim.before", "claim.result")) + between_hooks = context.Event() + continue_cycle = context.Event() + start_next_cycle = context.Event() + process = context.Process( + target=_straddle_claim_cycle_in_child, + args=(hooks, between_hooks, continue_cycle, start_next_cycle), + ) + process.start() + armed = threading.Event() + + def arm_next_cycle() -> None: + hooks.arm("claim.before", "claim.result") + armed.set() + + controller = threading.Thread(target=arm_next_cycle) + try: + assert between_hooks.wait(5) + controller.start() + assert not armed.wait(0.05) + continue_cycle.set() + controller.join(timeout=5) + assert armed.is_set() + start_next_cycle.set() + + before = hooks.hooks["claim.before"] + assert before.reached.wait(5) + before_generation, before_payload = before.metadata.get(timeout=1) + assert before_generation == before.current_generation() + assert before_payload == {"cycle": 2} + before.release_current() + + result = hooks.hooks["claim.result"] + assert result.reached.wait(5) + result_generation, result_payload = result.metadata.get(timeout=1) + assert result_generation == result.current_generation() + assert result_payload == {"cycle": 2} + result.release_current() + process.join(timeout=5) + finally: + continue_cycle.set() + start_next_cycle.set() + for hook in hooks.hooks.values(): + hook.release_current() + controller.join(timeout=1) + if process.is_alive(): + process.kill() + process.join(timeout=2) + + assert process.exitcode == 0 + + def test_run_proof_measures_session_hygiene_after_worker_polling_stops( monkeypatch: pytest.MonkeyPatch, tmp_path: Path, From aa2887f1ab260115da2a619d9cb48226cc6edd66 Mon Sep 17 00:00:00 2001 From: YingzuoLiu <156630675+YingzuoLiu@users.noreply.github.com> Date: Wed, 26 Aug 2026 19:46:23 +0800 Subject: [PATCH 8/8] fix: keep P5 generations readable across SIGSTOP --- docs/p5-multi-worker-recovery-proof.md | 10 ++-- examples/p5_proof_worker.py | 69 ++++++++++---------------- tests/test_p5_multi_worker_proof.py | 2 + 3 files changed, 34 insertions(+), 47 deletions(-) diff --git a/docs/p5-multi-worker-recovery-proof.md b/docs/p5-multi-worker-recovery-proof.md index 4970f71..3a82b42 100644 --- a/docs/p5-multi-worker-recovery-proof.md +++ b/docs/p5-multi-worker-recovery-proof.md @@ -7,10 +7,12 @@ controller process submits Runs, controls named one-shot barriers, injects exact store-time lease expiry, and rereads durable evidence. The synthetic HTTP provider is a fourth, independent process with its own SQLite effect ledger. -Each reusable barrier has monotonic armed, consumed, reached, released, and completed -generation counters. Events are wake-up hints only: a stale event cannot satisfy a -different generation, and the controller cannot re-arm until the prior generation's -consumer has acknowledged completion. Starting a new schedule first releases every +Each reusable barrier has monotonic, lock-free armed, consumed, reached, released, +and completed generation counters. Every counter has one writer, so the deliberate +SIGSTOP schedule cannot freeze a process-shared generation mutex. Events are wake-up +hints only: a stale event cannot satisfy a different generation, and the controller +cannot re-arm until the prior generation's consumer has acknowledged completion. +Starting a new schedule first releases every prior named barrier, so a worker parked on one hook cannot prevent it from reaching a different hook in the next schedule. Arming the new schedule is also serialized with the proof adapter's whole `claim_next_run` cycle. A poll already between diff --git a/examples/p5_proof_worker.py b/examples/p5_proof_worker.py index 1ef726f..83f236f 100644 --- a/examples/p5_proof_worker.py +++ b/examples/p5_proof_worker.py @@ -105,12 +105,15 @@ def create(cls, context: Any) -> ProcessHook: release=context.Event(), completed=completed, metadata=context.Queue(maxsize=1), - generations=context.Array("Q", 5, lock=True), + # Every slot has exactly one writer: the controller writes armed / + # released, and the worker writes consumed / reached / completed. + # A mutex is both unnecessary and unsafe here because SIGSTOP is a + # deliberate proof action and would freeze any lock held in-process. + generations=context.Array("Q", 5, lock=False), ) def _generation(self, index: int) -> int: - with self.generations.get_lock(): - return int(self.generations[index]) + return int(self.generations[index]) def current_generation(self) -> int: return self._generation(self._ARMED) @@ -119,21 +122,19 @@ def reached_generation(self) -> int: return self._generation(self._REACHED) def generation_state(self) -> dict[str, int]: - with self.generations.get_lock(): - return { - "armed": int(self.generations[self._ARMED]), - "consumed": int(self.generations[self._CONSUMED]), - "reached": int(self.generations[self._REACHED]), - "released": int(self.generations[self._RELEASED]), - "completed": int(self.generations[self._COMPLETED]), - } + return { + "armed": int(self.generations[self._ARMED]), + "consumed": int(self.generations[self._CONSUMED]), + "reached": int(self.generations[self._REACHED]), + "released": int(self.generations[self._RELEASED]), + "completed": int(self.generations[self._COMPLETED]), + } def arm(self) -> int: deadline = time.monotonic() + P5_HOOK_TIMEOUT_SECONDS while True: - with self.generations.get_lock(): - previous = int(self.generations[self._ARMED]) - previous_completed = int(self.generations[self._COMPLETED]) + previous = int(self.generations[self._ARMED]) + previous_completed = int(self.generations[self._COMPLETED]) if previous_completed >= previous: break remaining = deadline - time.monotonic() @@ -147,47 +148,33 @@ def arm(self) -> int: self.metadata.get_nowait() except queue.Empty: break - with self.generations.get_lock(): - generation = int(self.generations[self._ARMED]) + 1 - self.generations[self._ARMED] = generation + generation = int(self.generations[self._ARMED]) + 1 + self.generations[self._ARMED] = generation self.enabled.set() return generation def consume(self) -> int | None: if not self.enabled.is_set(): return None - with self.generations.get_lock(): - generation = int(self.generations[self._ARMED]) - if int(self.generations[self._CONSUMED]) >= generation: - return None - self.generations[self._CONSUMED] = generation + generation = int(self.generations[self._ARMED]) + if int(self.generations[self._CONSUMED]) >= generation: + return None + self.generations[self._CONSUMED] = generation self.enabled.clear() return generation def _mark_reached(self, generation: int, payload: dict[str, Any]) -> None: self.metadata.put((generation, payload)) - with self.generations.get_lock(): - self.generations[self._REACHED] = max( - int(self.generations[self._REACHED]), - generation, - ) + self.generations[self._REACHED] = generation self.reached.set() def _mark_completed(self, generation: int) -> None: - with self.generations.get_lock(): - self.generations[self._COMPLETED] = max( - int(self.generations[self._COMPLETED]), - generation, - ) + self.generations[self._COMPLETED] = generation self.completed.set() def release_current(self) -> None: - with self.generations.get_lock(): - generation = int(self.generations[self._ARMED]) - self.generations[self._RELEASED] = max( - int(self.generations[self._RELEASED]), - generation, - ) + generation = int(self.generations[self._ARMED]) + self.generations[self._RELEASED] = generation self.release.set() def block_consumed(self, generation: int, payload: dict[str, Any]) -> None: @@ -224,11 +211,7 @@ def pause_process(self, payload: dict[str, Any]) -> bool: # reached event before its matching payload is readable. self.metadata.close() self.metadata.join_thread() - with self.generations.get_lock(): - self.generations[self._REACHED] = max( - int(self.generations[self._REACHED]), - generation, - ) + self.generations[self._REACHED] = generation self.reached.set() os.kill(os.getpid(), signal.SIGSTOP) finally: diff --git a/tests/test_p5_multi_worker_proof.py b/tests/test_p5_multi_worker_proof.py index fbbdbfd..d2f3791 100644 --- a/tests/test_p5_multi_worker_proof.py +++ b/tests/test_p5_multi_worker_proof.py @@ -515,6 +515,8 @@ def test_process_pause_hook_flushes_metadata_before_sigstop() -> None: hook.current_generation(), {"point": "paused", "attempt": 1}, ) + assert not hasattr(hook.generations, "get_lock") + assert hook.reached_generation() == hook.current_generation() finally: if process.is_alive() and process.pid is not None: os.kill(process.pid, signal.SIGCONT)