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:
Ouroboros 2026-09-02 20:28:41 +00:00
parent 72bb494928
commit 626b48b73a
5 changed files with 182 additions and 6 deletions

View file

@ -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",

View file

@ -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(

View file

@ -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

View file

@ -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
# ===========================================================================

View file

@ -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)