mirror of
https://github.com/razzant/ouroboros.git
synced 2026-10-03 20:27:56 +00:00
fix: root the two CI races — copy-back source-handle promotion and the /proc env marker
O3 (CI 33579445704, windows-latest attempt 1). Two concurrent copy-backs of the same task promote the SAME content-addressed source handle. Both miss the destination, both write; on Windows the loser's os.replace over a destination the winner (or a verifying reader) holds open is a sharing violation, CPython opening files without FILE_SHARE_DELETE. The loser then published an INCOMPLETE promotion with a pending ref and promoted_source_handle_count=0 while the winner published complete/1 — two different custody projections for one settled fact, which is the `child_ref_promotion` mismatch the Windows leg read. Fixed as a class, at both seams, with no sleeps and no test-side retry: * `store_actor_source_bytes` is WRITE-ONCE. The digest is in the file name, so a destination already holding exactly these bytes IS this handle; re-storing it no longer replaces the file, which removes the contended replace outright for the ordinary second copier. * `_promote_task_source_ref` judges by the POSTCONDITION, not by authorship of the write: when the store raises, it asks the destination once more through `read_actor_source_bytes`, and a verified handle there counts as promoted. Only a still-unreadable destination stays a pending ref. The mechanism is Windows-only (POSIX rename(2) never refuses an open destination), but the fix is platform-neutral: the postcondition is the same everywhere and the pins run on every OS. Env-marker race (CI 33671108287, system-e2e-mock attempt 1, `assert 3898 in []`). `Popen` returns once the exec SUCCEEDED — the CLOEXEC error pipe closes inside execve — but the kernel publishes the new image's env_start/env_end later in that same path, so a /proc read landing in that window sees an EMPTY environ for a live, correctly marked child. The harness now splits the two oracles: `wait_pid_env_value` polls THE ONE pid for a bounded window (the positive claim), while `pids_with_env_value` keeps its single scan (the no-orphans postcondition, where an orphan was execed long before the scan and the window cannot hide it). Both read through one seam so the window can be pinned deterministically.
This commit is contained in:
parent
72bb494928
commit
626b48b73a
5 changed files with 182 additions and 6 deletions
|
|
@ -573,7 +573,17 @@ def store_actor_source_bytes(
|
|||
f"{safe_id}-{digest}.{safe_extension}",
|
||||
)
|
||||
target = artifact_dir.joinpath(*relative.parts)
|
||||
write_bytes_atomic(target, bytes(data))
|
||||
# WRITE-ONCE. The name carries the digest, so an existing file with exactly
|
||||
# these bytes IS this handle: rewriting it would only republish identical
|
||||
# content while racing another writer for the same path (on Windows an
|
||||
# os.replace over a destination a concurrent reader holds open is a sharing
|
||||
# violation, which is how a second copy-back of the same handle used to fail).
|
||||
try:
|
||||
already_stored = not target.is_symlink() and target.read_bytes() == bytes(data)
|
||||
except OSError:
|
||||
already_stored = False
|
||||
if not already_stored:
|
||||
write_bytes_atomic(target, bytes(data))
|
||||
return {
|
||||
"kind": "task_source",
|
||||
"root": "artifact_store",
|
||||
|
|
|
|||
|
|
@ -639,6 +639,21 @@ def _promote_task_source_ref(
|
|||
state["promoted_source_handle_count"] += 1
|
||||
return dict(ref)
|
||||
except Exception as exc:
|
||||
# A CONCURRENT copy-back of the same task may have claimed this exact
|
||||
# content-addressed destination between our miss above and our write
|
||||
# (its os.replace can refuse ours on Windows, where a destination another
|
||||
# thread holds open cannot be replaced). The promotion's postcondition is
|
||||
# the verified handle at the destination, not authorship of the write, so
|
||||
# ask the destination once more: if it verifies, the copy DID happen and
|
||||
# this caller must publish the same complete custody projection as the
|
||||
# winner. Only a still-unreadable destination is a pending ref.
|
||||
try:
|
||||
read_actor_source_bytes(parent_root, task_id, ref)
|
||||
except Exception:
|
||||
pass
|
||||
else:
|
||||
state["promoted_source_handle_count"] += 1
|
||||
return dict(ref)
|
||||
_append_promotion_fact(
|
||||
state["pending_refs"],
|
||||
_promotion_fact(
|
||||
|
|
|
|||
|
|
@ -731,6 +731,46 @@ def process_tree_pids(root_pid: int) -> list:
|
|||
return pids
|
||||
|
||||
|
||||
def _read_proc_environ_bytes(pid) -> bytes:
|
||||
"""Raw ``/proc/<pid>/environ`` bytes, or ``b""`` when it cannot be read.
|
||||
|
||||
The single read seam of both /proc environ oracles below, so the empty-window
|
||||
behaviour they must survive can be pinned deterministically.
|
||||
"""
|
||||
try:
|
||||
return pathlib.Path(f"/proc/{pid}/environ").read_bytes()
|
||||
except OSError:
|
||||
return b""
|
||||
|
||||
|
||||
def wait_pid_env_value(pid: int, value: str, timeout: float = 10.0) -> bool:
|
||||
"""True once ``/proc/<pid>/environ`` carries *value*, within a bounded window.
|
||||
|
||||
THE POSITIVE ORACLE. ``Popen`` returns as soon as the exec SUCCEEDED — the
|
||||
CLOEXEC error pipe closes inside ``execve`` — but the kernel publishes the NEW
|
||||
image's ``env_start``/``env_end`` later in that same exec path, so a read that
|
||||
lands in that window sees an EMPTY environ for a live, correctly marked child.
|
||||
A positive claim ("this child carries the marker") must therefore poll THE ONE
|
||||
pid until its environ becomes readable or a deadline passes; scanning all of
|
||||
/proc again would only re-roll the same window against a moving target.
|
||||
|
||||
The negative oracle (``pids_with_env_value`` as the no-orphans postcondition)
|
||||
deliberately keeps its SINGLE scan: an orphan that survived a teardown was
|
||||
execed long before the scan, so the post-exec window cannot hide it, while a
|
||||
wait there would only slow every clean teardown down.
|
||||
"""
|
||||
needle = str(value).encode()
|
||||
deadline = time.monotonic() + float(timeout)
|
||||
while True:
|
||||
if needle in _read_proc_environ_bytes(pid):
|
||||
return True
|
||||
if not pathlib.Path(f"/proc/{pid}").exists():
|
||||
return False
|
||||
if time.monotonic() >= deadline:
|
||||
return False
|
||||
time.sleep(0.01)
|
||||
|
||||
|
||||
def pids_with_env_value(value: str) -> list:
|
||||
"""Every live pid whose /proc environ carries *value* (readable procs only).
|
||||
|
||||
|
|
@ -741,17 +781,17 @@ def pids_with_env_value(value: str) -> list:
|
|||
teardown would show up here reparented outside it. Same-uid procs only by
|
||||
construction (/proc environ of other users is unreadable), which covers the
|
||||
whole tree an isolated server can have spawned.
|
||||
|
||||
NEGATIVE-USE ORACLE: one scan, no waiting (see ``wait_pid_env_value`` for why
|
||||
a positive claim about a just-spawned child needs a bounded wait instead).
|
||||
"""
|
||||
needle = str(value).encode()
|
||||
found = []
|
||||
for pid_dir in pathlib.Path("/proc").iterdir():
|
||||
if not pid_dir.name.isdigit():
|
||||
continue
|
||||
try:
|
||||
if needle in (pid_dir / "environ").read_bytes():
|
||||
found.append(int(pid_dir.name))
|
||||
except OSError:
|
||||
continue
|
||||
if needle in _read_proc_environ_bytes(pid_dir.name):
|
||||
found.append(int(pid_dir.name))
|
||||
return found
|
||||
|
||||
|
||||
|
|
|
|||
|
|
@ -62,6 +62,7 @@ import uuid
|
|||
|
||||
import pytest
|
||||
|
||||
from tests.system_e2e import harness
|
||||
from tests.system_e2e.harness import (
|
||||
LANE_MOCK,
|
||||
MOCK_SLUG,
|
||||
|
|
@ -80,6 +81,7 @@ from tests.system_e2e.harness import (
|
|||
start_server,
|
||||
submit_running,
|
||||
wait_durable_result,
|
||||
wait_pid_env_value,
|
||||
wait_until,
|
||||
)
|
||||
|
||||
|
|
@ -134,9 +136,13 @@ def test_replay_model_callable_steps_and_explicit_model_ids():
|
|||
|
||||
@pytest.mark.skipif(sys.platform != "linux", reason="/proc environ scan is Linux-only")
|
||||
def test_pids_with_env_value_sees_and_loses_a_marked_child():
|
||||
"""Both directions of the /proc environ oracle, each through its own contract:
|
||||
the positive claim waits a bounded window for the just-execed child's environ
|
||||
(``wait_pid_env_value``), the no-orphans postcondition keeps its single scan."""
|
||||
marker = f"e2e-w2-marker-{uuid.uuid4().hex}"
|
||||
child = subprocess.Popen(["sleep", "30"], env={**os.environ, "E2E_W2_MARK": marker})
|
||||
try:
|
||||
assert wait_pid_env_value(child.pid, marker)
|
||||
assert child.pid in pids_with_env_value(marker)
|
||||
finally:
|
||||
child.kill()
|
||||
|
|
@ -144,6 +150,28 @@ def test_pids_with_env_value_sees_and_loses_a_marked_child():
|
|||
assert child.pid not in pids_with_env_value(marker)
|
||||
|
||||
|
||||
def test_pid_env_wait_rides_out_the_post_exec_empty_environ_window(monkeypatch):
|
||||
"""The exact CI interleaving (33671108287): ``Popen`` has returned — the exec
|
||||
succeeded — but the kernel has not published the new image's env_start/env_end
|
||||
yet, so the environ reads EMPTY for a live, correctly marked pid. A single scan
|
||||
misses it (and that is what the positive assertion used to be); the bounded
|
||||
per-pid wait rides the window out and still returns True."""
|
||||
marker = "e2e-w2-window-marker"
|
||||
reads = {"n": 0}
|
||||
real = harness._read_proc_environ_bytes
|
||||
|
||||
def _windowed(pid):
|
||||
if str(pid) != str(os.getpid()):
|
||||
return real(pid)
|
||||
reads["n"] += 1
|
||||
return b"" if reads["n"] <= 3 else f"E2E_W2_MARK={marker}\x00".encode()
|
||||
|
||||
monkeypatch.setattr(harness, "_read_proc_environ_bytes", _windowed)
|
||||
assert os.getpid() not in pids_with_env_value(marker) # the single scan misses it
|
||||
assert wait_pid_env_value(os.getpid(), marker, timeout=5) # the bounded wait does not
|
||||
assert reads["n"] > 3
|
||||
|
||||
|
||||
# ===========================================================================
|
||||
# Shared helpers of the mock-lane scenarios
|
||||
# ===========================================================================
|
||||
|
|
|
|||
|
|
@ -520,6 +520,89 @@ def test_concurrent_copyback_is_idempotent_and_copies_only_referenced_source_han
|
|||
assert not (parent / "observability" / "blobs" / pathlib.Path(unreferenced_blob["path"]).name).exists()
|
||||
|
||||
|
||||
def test_copyback_source_handle_promotion_survives_a_lost_write_race(tmp_path, monkeypatch):
|
||||
"""The Windows interleaving of the concurrent copy-back (CI 33579445704): two
|
||||
copiers both miss the destination handle, the winner's identical bytes land,
|
||||
and the loser's own os.replace over that destination is refused as a sharing
|
||||
violation (CPython opens files without FILE_SHARE_DELETE). The promotion's
|
||||
postcondition is the VERIFIED handle at the destination, not authorship of the
|
||||
write, so the loser must publish the same complete custody projection — not an
|
||||
incomplete promotion with a pending ref."""
|
||||
from ouroboros import artifacts as artifacts_module
|
||||
from ouroboros.artifacts import store_actor_source_bytes
|
||||
|
||||
task_id = "phase3c-lost-write-race"
|
||||
parent, child = _child(tmp_path, task_id)
|
||||
source_ref = store_actor_source_bytes(
|
||||
child, task_id, category="tool_results", source_id="tool",
|
||||
data=b"actor promised source", extension="txt",
|
||||
)
|
||||
write_task_result(
|
||||
child, task_id, STATUS_COMPLETED, result="done", artifact_status="ready",
|
||||
review_evidence={"exact_source_ref": source_ref},
|
||||
)
|
||||
real_store = artifacts_module.store_actor_source_bytes
|
||||
|
||||
def _losing_store(*args, **kwargs):
|
||||
real_store(*args, **kwargs) # the concurrent winner's identical copy lands
|
||||
raise PermissionError("[WinError 32] the file is in use by another process")
|
||||
|
||||
monkeypatch.setattr(artifacts_module, "store_actor_source_bytes", _losing_store)
|
||||
|
||||
copied = copy_child_task_result(parent, {"id": task_id, "drive_root": str(child)})
|
||||
|
||||
promotion = copied["child_ref_promotion"]
|
||||
assert promotion["status"] == "complete"
|
||||
assert promotion["promoted_source_handle_count"] == 1
|
||||
assert promotion["pending_refs"] == []
|
||||
assert copied["review_evidence"]["exact_source_ref"] == source_ref
|
||||
promoted = (
|
||||
parent / "task_results" / "artifacts" / task_id / source_ref["path"]
|
||||
)
|
||||
assert promoted.read_bytes() == b"actor promised source"
|
||||
|
||||
|
||||
def test_store_actor_source_bytes_does_not_rewrite_an_identical_handle(tmp_path, monkeypatch):
|
||||
"""Content-addressed write-once: the digest is in the name, so re-storing the
|
||||
same bytes must not replace the file at all — that replace is the operation a
|
||||
concurrent copier can lose (and on Windows must lose, when a reader holds the
|
||||
destination open)."""
|
||||
from ouroboros import artifacts as artifacts_module
|
||||
from ouroboros.artifacts import store_actor_source_bytes
|
||||
|
||||
task_id = "phase3c-write-once"
|
||||
drive = tmp_path / "data"
|
||||
drive.mkdir()
|
||||
first = store_actor_source_bytes(
|
||||
drive, task_id, category="tool_results", source_id="tool",
|
||||
data=b"exact bytes", extension="txt",
|
||||
)
|
||||
target = drive / "task_results" / "artifacts" / task_id / first["path"]
|
||||
before = target.stat()
|
||||
writes = []
|
||||
real_write = artifacts_module.write_bytes_atomic
|
||||
monkeypatch.setattr(
|
||||
artifacts_module, "write_bytes_atomic",
|
||||
lambda path, content, **kw: (writes.append(str(path)), real_write(path, content, **kw))[1],
|
||||
)
|
||||
|
||||
again = store_actor_source_bytes(
|
||||
drive, task_id, category="tool_results", source_id="tool",
|
||||
data=b"exact bytes", extension="txt",
|
||||
)
|
||||
|
||||
assert again == first
|
||||
assert writes == []
|
||||
after = target.stat()
|
||||
assert (after.st_ino, after.st_mtime_ns) == (before.st_ino, before.st_mtime_ns)
|
||||
# Different bytes are a DIFFERENT handle and are still written.
|
||||
other = store_actor_source_bytes(
|
||||
drive, task_id, category="tool_results", source_id="tool",
|
||||
data=b"other bytes", extension="txt",
|
||||
)
|
||||
assert other["path"] != first["path"] and writes
|
||||
|
||||
|
||||
def test_legacy_missing_child_ref_is_typed_gap_without_permanent_retention(tmp_path):
|
||||
task_id = "phase3c-legacy-gap"
|
||||
parent, child = _child(tmp_path, task_id)
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue