From 2938a678445bcd9010f2e0ef8875acf9dcebfafe Mon Sep 17 00:00:00 2001 From: huangruiteng <14976749+huangruiteng@users.noreply.github.com> Date: Thu, 17 Sep 2026 04:48:24 +0800 Subject: [PATCH] fix(turn-driver): fence one Turn lane across hosts, not only one machine A Turn lane admits exactly one executing Turn, and the fence was a kernel lock held by the executing process. That is an authority inside one machine only: when the runtime root is shared -- an SSH-driven executor beside a local one -- the other host cannot see this lock, so two hosts could execute the same lane at once, each invoking its own host and spending its own slot. The same Turn now also holds a durable lane lease under the runtime root: an exclusive create naming this host (as a fingerprint), this process and this Turn, read before the kernel lock and refused while it is live. A lease whose holder is on another machine clears only by its own expiry, so a host that died mid-Turn cannot hold the lane forever; a lease left by a process on this machine clears as soon as that process is gone, because the kernel lock is what actually keeps two local Turns apart. A settled Turn releases its lease, and the typed turn_lane_in_flight refusal now names whichever holder refused -- this machine's lock or the other host's lease -- with the same all-false effect shape and the same pre-journal timing. Validation: tests/test_turn_lane_fence.py (8 tests) plus the 217-test turn driver set. The new cases cover a live lease from another host being refused before the journal, an expired lease and a crashed same-host holder not holding the lane, the holder's claim being readable by the next Turn, and the release leaving no claim behind. Signed-off-by: huangruiteng <14976749+huangruiteng@users.noreply.github.com> --- loopx/control_plane/turn_driver/lane_fence.py | 242 +++++++++++++++++- tests/test_turn_lane_fence.py | 147 ++++++++++- 2 files changed, 377 insertions(+), 12 deletions(-) diff --git a/loopx/control_plane/turn_driver/lane_fence.py b/loopx/control_plane/turn_driver/lane_fence.py index 24ba3c1d29..b9a9b986c4 100644 --- a/loopx/control_plane/turn_driver/lane_fence.py +++ b/loopx/control_plane/turn_driver/lane_fence.py @@ -6,23 +6,33 @@ one bounded question. LoopX therefore admits exactly one *executing* Turn per lane and refuses the second with a typed, retryable refusal naming the holder. -The fence is a kernel lock held by the executing process, so a crashed or killed -Turn releases the lane instead of leaving a stale claim no later Turn can enter. -Previews and other non-executing decisions never take it. +The fence is two facts, because a kernel lock is only an authority inside one +machine. The executing process holds a kernel lock, so a crashed or killed Turn +releases the lane instead of leaving a stale claim no later Turn can enter. When +the runtime root is shared -- an SSH-driven executor beside a local one, for +example -- a lock the other host cannot see is not a fence, so the same Turn +also holds a durable lane lease under the runtime root: an exclusive create +naming this host and this Turn, refused while it is live, and taken over once it +expires. Previews and other non-executing decisions never take either. """ from __future__ import annotations from collections.abc import Callable, Iterator, Mapping from contextlib import contextmanager +from datetime import timedelta from functools import wraps import hashlib import json +import os from pathlib import Path import re +import socket from typing import Any from ...file_lock import lock_holder_path, try_exclusive_file_lock +from ..runtime.time import now_utc as runtime_now_utc +from ..runtime.time import parse_timestamp, utc_isoformat # Typed refusal for a lane whose single executor is already busy. The reason is # a fact about this lane, so a caller can retry it unchanged once it clears. @@ -32,6 +42,18 @@ TURN_LANE_OPERATION = "loopx_turn_lane" TURN_LANE_DIR_NAME = ".lanes" TURN_LANE_UNATTRIBUTED_AGENT = "unattributed" +# A durable lane lease exists because a kernel lock is not an authority across +# hosts. Its holder names the host it was taken on, so a second host sharing the +# runtime root refuses it, and it expires so a host that died cannot hold a lane +# forever. +TURN_LANE_LEASE_SCHEMA_VERSION = "turn_lane_lease_v0" +TURN_LANE_LEASE_SUFFIX = ".lease.json" +# Long enough to cover a bounded Turn, short enough that a host which died +# without releasing the lane does not block the lane for the rest of the day. +TURN_LANE_LEASE_TTL_SECONDS = 30 * 60 +# Public-safe holder facts. The host is kept as a fingerprint: a refusal has to +# say "another host is running this lane" without publishing which machine it is. +TURN_LANE_LEASE_HOLDER_TEXT_FIELDS = ("agent_id", "operation", "acquired_at") # Public-safe holder fields only: the lock record also carries a lock id, a # policy name, and the private lock path, which never leave this process. TURN_LANE_HOLDER_TEXT_FIELDS = ("agent_id", "operation", "acquired_at") @@ -93,12 +115,212 @@ def turn_lane_singleflight( """ target = turn_lane_target(runtime_root=runtime_root, goal_id=goal_id, plan=plan) + lease_target = turn_lane_lease_target(target) + agent_id = turn_lane_agent_id(plan) + # The durable lease is read first: a lane another host is executing is + # refused before this process even contends for its own kernel lock. + if turn_lane_lease_holder(lease_target) is not None: + yield None + return with try_exclusive_file_lock( target, - agent_id=turn_lane_agent_id(plan), + agent_id=agent_id, operation=TURN_LANE_OPERATION, ) as lock_path: - yield lock_path + if lock_path is None: + yield None + return + record = turn_lane_lease_record( + agent_id=agent_id, turn_instance_id=turn_lane_turn_instance_id(plan) + ) + if not _claim_turn_lane_lease(lease_target, record): + # Another host claimed the lane between the read above and this + # process's kernel lock. The lease is the authority, so this Turn is + # refused and its own kernel lock is released by the context manager. + yield None + return + try: + yield lock_path + finally: + _release_turn_lane_lease(lease_target, record) + + +def turn_lane_lease_target(target: Path) -> Path: + """Return the durable lease path that sits beside one lane's kernel lock.""" + + return target.with_name(f"{target.name}{TURN_LANE_LEASE_SUFFIX}") + + +def turn_lane_host_fingerprint() -> str: + """Return a public-safe fingerprint of the machine holding a lane lease.""" + + try: + hostname = socket.gethostname() + except OSError: # pragma: no cover - a host without a name is still a host + hostname = "" + return hashlib.sha256(hostname.encode("utf-8")).hexdigest()[:12] + + +def turn_lane_turn_instance_id(plan: Mapping[str, Any]) -> str: + """Return the Turn identity a lease names, taken from the Turn envelope.""" + + envelope = plan.get("turn_envelope") + if isinstance(envelope, Mapping): + for field in ("turn_instance_id", "turn_id", "run_id"): + value = str(envelope.get(field) or "").strip() + if value: + return value + return "" + + +def turn_lane_lease_record( + *, agent_id: str, turn_instance_id: str = "" +) -> dict[str, Any]: + """Build this process's durable claim on one lane.""" + + acquired = runtime_now_utc() + return { + "schema_version": TURN_LANE_LEASE_SCHEMA_VERSION, + "agent_id": agent_id or TURN_LANE_UNATTRIBUTED_AGENT, + "operation": TURN_LANE_OPERATION, + "turn_instance_id": turn_instance_id, + "acquired_at": utc_isoformat(acquired), + "expires_at": utc_isoformat( + acquired + timedelta(seconds=TURN_LANE_LEASE_TTL_SECONDS) + ), + "host_fingerprint": turn_lane_host_fingerprint(), + "pid": os.getpid(), + } + + +def _pid_is_alive(pid: int) -> bool | None: + """Report whether a recorded holder process still exists, when knowable.""" + + if pid <= 0: + return None + try: + os.kill(pid, 0) + except ProcessLookupError: + return False + except PermissionError: + # It exists; this process is simply not allowed to signal it. + return True + except (OSError, AttributeError, ValueError, NotImplementedError): + # The platform cannot answer. The lease's own expiry is then the only + # thing that clears it, which the caller already applies. + return None + return True + + +def _turn_lane_lease_is_live(record: Mapping[str, Any]) -> bool: + """Report whether one lease record still holds its lane.""" + + expires_at = record.get("expires_at") + if not isinstance(expires_at, str) or not expires_at: + # A record without an expiry is not a claim this module can honor: it + # would hold the lane forever, so it is treated as already released. + return False + try: + deadline = parse_timestamp(expires_at) + except (TypeError, ValueError): + return False + if deadline is None or runtime_now_utc() >= deadline: + return False + if str(record.get("host_fingerprint") or "") != turn_lane_host_fingerprint(): + # A different machine is executing this lane. Only its own expiry can + # clear the claim, because this process cannot observe its process. + return True + pid = record.get("pid") + if isinstance(pid, bool) or not isinstance(pid, int): + return True + # The same machine: a holder that no longer exists is a crashed Turn, and + # the lane is free. A live holder is still running, though the kernel lock + # is what actually keeps two local processes apart. + return _pid_is_alive(pid) is not False + + +def turn_lane_lease_readback(record: Mapping[str, Any]) -> dict[str, Any]: + """Return the public-safe identity of a durable lease holder.""" + + projection: dict[str, Any] = {} + for field in TURN_LANE_LEASE_HOLDER_TEXT_FIELDS: + value = record.get(field) + if isinstance(value, str) and value: + projection[field] = value + pid = record.get("pid") + if isinstance(pid, int) and not isinstance(pid, bool): + projection["pid"] = pid + return projection + + +def turn_lane_lease_holder(target: Path) -> dict[str, Any] | None: + """Return the readback of the holder of one lane, or ``None`` if it is free.""" + + record = _read_turn_lane_lease(target) + if record is None or not _turn_lane_lease_is_live(record): + return None + return turn_lane_lease_readback(record) + + +def _read_turn_lane_lease(target: Path) -> dict[str, Any] | None: + try: + payload = json.loads(target.read_text(encoding="utf-8")) + except (OSError, ValueError): + return None + return dict(payload) if isinstance(payload, Mapping) else None + + +def _claim_turn_lane_lease(target: Path, record: Mapping[str, Any]) -> bool: + """Take one lane's durable lease, replacing a stale record. + + The exclusive create is the claim: two hosts that both find the lane free + cannot both succeed, so exactly one of them executes. A record that already + exists is replaced only when it no longer holds the lane (expired, or left + by a crashed process on this host). + """ + + try: + handle = os.open(target, os.O_CREAT | os.O_EXCL | os.O_WRONLY, 0o600) + except FileExistsError: + existing = _read_turn_lane_lease(target) + if existing is not None and _turn_lane_lease_is_live(existing): + return False + _write_turn_lane_lease(target, record) + return True + except OSError: + return False + with os.fdopen(handle, "w", encoding="utf-8") as stream: + json.dump(dict(record), stream, ensure_ascii=False, indent=2) + stream.write("\n") + return True + + +def _write_turn_lane_lease(target: Path, record: Mapping[str, Any]) -> None: + """Replace one lane's lease atomically.""" + + target.parent.mkdir(parents=True, exist_ok=True) + temp_path = target.with_name(f".{target.name}.{os.getpid()}.tmp") + temp_path.write_text( + json.dumps(dict(record), ensure_ascii=False, indent=2) + "\n", + encoding="utf-8", + ) + temp_path.replace(target) + + +def _release_turn_lane_lease(target: Path, record: Mapping[str, Any]) -> None: + """Release this process's lease, leaving any other holder's record alone.""" + + existing = _read_turn_lane_lease(target) + if existing is None: + return + if str(existing.get("turn_instance_id") or "") != str( + record.get("turn_instance_id") or "" + ) or existing.get("pid") != record.get("pid"): + return + try: + target.unlink() + except OSError: # pragma: no cover - already released by a peer + return def turn_lane_holder_readback(target: Path) -> dict[str, Any]: @@ -188,11 +410,15 @@ def single_lane_turn( ) as held: if held is not None: return execute_turn(plan, *args, **kwargs) + # Either this machine's kernel lock or another host's durable + # lease refused the Turn, and the durable lease is the only one + # a caller on a shared runtime root can actually wait for. + holder = turn_lane_lease_holder( + turn_lane_lease_target(target) + ) or turn_lane_holder_readback(target) return execution_payload( plan, - turn_lane_in_flight_record( - plan, holder=turn_lane_holder_readback(target) - ), + turn_lane_in_flight_record(plan, holder=holder), execute=True, replayed=False, effects=TURN_LANE_NO_EFFECTS, diff --git a/tests/test_turn_lane_fence.py b/tests/test_turn_lane_fence.py index 8bdd881021..870ede5f6a 100644 --- a/tests/test_turn_lane_fence.py +++ b/tests/test_turn_lane_fence.py @@ -1,24 +1,34 @@ """One executing bounded Turn per Turn lane. -The fence is a process-level lock, so these tests hold the lane the same way a -running Turn does and assert what the second Turn is told. Admission and release -after a settled Turn are covered end to end by the public dsh smokes, which run -one Turn and then its replay through the same entry. +The fence is a process-level lock plus a durable lane lease, so these tests hold +the lane the same way a running Turn does and assert what the second Turn is +told -- locally, and from another host that shares the runtime root and cannot +see this machine's lock. Admission and release after a settled Turn are covered +end to end by the public dsh smokes, which run one Turn and then its replay +through the same entry. """ from __future__ import annotations +import json +import os from pathlib import Path +from datetime import timedelta from loopx.control_plane.turn_driver.executor import run_loopx_turn_once from loopx.control_plane.turn_driver.lane_fence import ( REMEDY_WAIT_FOR_IN_FLIGHT_TURN, TURN_LANE_IN_FLIGHT, + TURN_LANE_LEASE_SCHEMA_VERSION, TURN_LANE_OPERATION, turn_lane_holder_readback, + turn_lane_host_fingerprint, + turn_lane_lease_record, + turn_lane_lease_target, turn_lane_singleflight, turn_lane_target, ) +from loopx.control_plane.runtime.time import now_utc, utc_isoformat GOAL_ID = "lane-fence-goal" AGENT_ID = "lane-fence-agent" @@ -123,3 +133,132 @@ def test_the_holder_readback_stays_public_safe(tmp_path: Path) -> None: # The private lock identity and the runtime path never leave the process. assert str(tmp_path) not in str(holder) assert turn_lane_holder_readback(tmp_path / "absent.lane") == {} + + +def _write_lease( + tmp_path: Path, + *, + host_fingerprint: str, + expires_in_seconds: int, + pid: int | None = None, +) -> Path: + """Write a durable lane lease the way another host would have left it.""" + + target = turn_lane_lease_target( + turn_lane_target( + runtime_root=tmp_path / "runtime", goal_id=GOAL_ID, plan=_plan() + ) + ) + acquired = now_utc() + record = { + "schema_version": TURN_LANE_LEASE_SCHEMA_VERSION, + "agent_id": AGENT_ID, + "operation": TURN_LANE_OPERATION, + "turn_instance_id": "another-host-turn", + "acquired_at": utc_isoformat(acquired), + "expires_at": utc_isoformat(acquired + timedelta(seconds=expires_in_seconds)), + "host_fingerprint": host_fingerprint, + "pid": os.getpid() if pid is None else pid, + } + target.parent.mkdir(parents=True, exist_ok=True) + target.write_text(json.dumps(record, ensure_ascii=False, indent=2), encoding="utf-8") + return target + + +def test_a_second_host_sharing_the_runtime_root_is_refused(tmp_path: Path) -> None: + """A kernel lock cannot fence a host that cannot see it. + + The runtime root is shared -- an SSH-driven executor beside this machine is + the case that motivated the durable lease -- so the lane's holder is read + from the runtime root and refused before the journal, the host and quota. + """ + + _write_lease(tmp_path, host_fingerprint="a-different-machine", expires_in_seconds=600) + + payload = _execute(tmp_path) + + assert payload["ok"] is False + assert payload["status"] == "unavailable" + assert payload["reason"] == TURN_LANE_IN_FLIGHT + assert payload["remediation"] == [REMEDY_WAIT_FOR_IN_FLIGHT_TURN] + # The refusal is a readback of the lease, not an invocation. + assert payload["host"] == {"executable": "not_invoked", "kind": "dsh"} + assert payload["effects"] == { + "host_invoked": False, + "state_written": False, + "quota_spent": False, + "scheduler_acknowledged": False, + } + assert payload["quota_slot_spend_count"] == 0 + assert payload["in_flight"]["agent_id"] == AGENT_ID + assert payload["in_flight"]["operation"] == TURN_LANE_OPERATION + assert payload["in_flight"]["acquired_at"] + # Nothing about the other machine leaves the process: no host name, no path. + assert "a-different-machine" not in json.dumps(payload) + assert str(tmp_path) not in json.dumps(payload) + + +def test_an_expired_lease_from_another_host_does_not_hold_the_lane( + tmp_path: Path, +) -> None: + """A host that died mid-Turn must not hold its lane forever.""" + + lease_target = _write_lease( + tmp_path, host_fingerprint="a-different-machine", expires_in_seconds=-60 + ) + + with turn_lane_singleflight( + runtime_root=tmp_path / "runtime", goal_id=GOAL_ID, plan=_plan() + ) as held: + assert held is not None + # The stale record was replaced by this Turn's own claim, so the lane is + # held by a live lease rather than by the host that never released it. + stored = json.loads(lease_target.read_text(encoding="utf-8")) + assert stored["host_fingerprint"] == turn_lane_host_fingerprint() + assert stored["expires_at"] > stored["acquired_at"] + + assert not lease_target.exists() + + +def test_a_crashed_holder_on_this_host_does_not_hold_the_lane(tmp_path: Path) -> None: + """A holder on this machine that no longer exists is a finished fence.""" + + _write_lease( + tmp_path, + host_fingerprint=turn_lane_host_fingerprint(), + expires_in_seconds=600, + pid=999_999, + ) + + with turn_lane_singleflight( + runtime_root=tmp_path / "runtime", goal_id=GOAL_ID, plan=_plan() + ) as held: + assert held is not None + + +def test_a_held_lane_writes_the_lease_the_next_turn_reads(tmp_path: Path) -> None: + """The holder writes the durable fact the second host decides from.""" + + with turn_lane_singleflight( + runtime_root=tmp_path / "runtime", goal_id=GOAL_ID, plan=_plan() + ) as held: + assert held is not None + record = turn_lane_lease_record(agent_id=AGENT_ID, turn_instance_id="turn-1") + assert record["schema_version"] == TURN_LANE_LEASE_SCHEMA_VERSION + lease_target = turn_lane_lease_target( + turn_lane_target( + runtime_root=tmp_path / "runtime", goal_id=GOAL_ID, plan=_plan() + ) + ) + assert lease_target.exists() + stored = json.loads(lease_target.read_text(encoding="utf-8")) + assert stored["host_fingerprint"] == turn_lane_host_fingerprint() + assert stored["pid"] == os.getpid() + assert stored["agent_id"] == AGENT_ID + + # A settled Turn leaves no claim behind for the next Turn or host. + assert not lease_target.exists() + with turn_lane_singleflight( + runtime_root=tmp_path / "runtime", goal_id=GOAL_ID, plan=_plan() + ) as again: + assert again is not None