diff --git a/loopx/control_plane/turn_driver/lane_fence.py b/loopx/control_plane/turn_driver/lane_fence.py index 24ba3c1d2..b9a9b986c 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 8bdd88102..870ede5f6 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