|
6 | 6 | one bounded question. LoopX therefore admits exactly one *executing* Turn per |
7 | 7 | lane and refuses the second with a typed, retryable refusal naming the holder. |
8 | 8 |
|
9 | | -The fence is a kernel lock held by the executing process, so a crashed or killed |
10 | | -Turn releases the lane instead of leaving a stale claim no later Turn can enter. |
11 | | -Previews and other non-executing decisions never take it. |
| 9 | +The fence is two facts, because a kernel lock is only an authority inside one |
| 10 | +machine. The executing process holds a kernel lock, so a crashed or killed Turn |
| 11 | +releases the lane instead of leaving a stale claim no later Turn can enter. When |
| 12 | +the runtime root is shared -- an SSH-driven executor beside a local one, for |
| 13 | +example -- a lock the other host cannot see is not a fence, so the same Turn |
| 14 | +also holds a durable lane lease under the runtime root: an exclusive create |
| 15 | +naming this host and this Turn, refused while it is live, and taken over once it |
| 16 | +expires. Previews and other non-executing decisions never take either. |
12 | 17 | """ |
13 | 18 |
|
14 | 19 | from __future__ import annotations |
15 | 20 |
|
16 | 21 | from collections.abc import Callable, Iterator, Mapping |
17 | 22 | from contextlib import contextmanager |
| 23 | +from datetime import timedelta |
18 | 24 | from functools import wraps |
19 | 25 | import hashlib |
20 | 26 | import json |
| 27 | +import os |
21 | 28 | from pathlib import Path |
22 | 29 | import re |
| 30 | +import socket |
23 | 31 | from typing import Any |
24 | 32 |
|
25 | 33 | from ...file_lock import lock_holder_path, try_exclusive_file_lock |
| 34 | +from ..runtime.time import now_utc as runtime_now_utc |
| 35 | +from ..runtime.time import parse_timestamp, utc_isoformat |
26 | 36 |
|
27 | 37 | # Typed refusal for a lane whose single executor is already busy. The reason is |
28 | 38 | # a fact about this lane, so a caller can retry it unchanged once it clears. |
|
32 | 42 | TURN_LANE_OPERATION = "loopx_turn_lane" |
33 | 43 | TURN_LANE_DIR_NAME = ".lanes" |
34 | 44 | TURN_LANE_UNATTRIBUTED_AGENT = "unattributed" |
| 45 | +# A durable lane lease exists because a kernel lock is not an authority across |
| 46 | +# hosts. Its holder names the host it was taken on, so a second host sharing the |
| 47 | +# runtime root refuses it, and it expires so a host that died cannot hold a lane |
| 48 | +# forever. |
| 49 | +TURN_LANE_LEASE_SCHEMA_VERSION = "turn_lane_lease_v0" |
| 50 | +TURN_LANE_LEASE_SUFFIX = ".lease.json" |
| 51 | +# Long enough to cover a bounded Turn, short enough that a host which died |
| 52 | +# without releasing the lane does not block the lane for the rest of the day. |
| 53 | +TURN_LANE_LEASE_TTL_SECONDS = 30 * 60 |
| 54 | +# Public-safe holder facts. The host is kept as a fingerprint: a refusal has to |
| 55 | +# say "another host is running this lane" without publishing which machine it is. |
| 56 | +TURN_LANE_LEASE_HOLDER_TEXT_FIELDS = ("agent_id", "operation", "acquired_at") |
35 | 57 | # Public-safe holder fields only: the lock record also carries a lock id, a |
36 | 58 | # policy name, and the private lock path, which never leave this process. |
37 | 59 | TURN_LANE_HOLDER_TEXT_FIELDS = ("agent_id", "operation", "acquired_at") |
@@ -93,12 +115,212 @@ def turn_lane_singleflight( |
93 | 115 | """ |
94 | 116 |
|
95 | 117 | target = turn_lane_target(runtime_root=runtime_root, goal_id=goal_id, plan=plan) |
| 118 | + lease_target = turn_lane_lease_target(target) |
| 119 | + agent_id = turn_lane_agent_id(plan) |
| 120 | + # The durable lease is read first: a lane another host is executing is |
| 121 | + # refused before this process even contends for its own kernel lock. |
| 122 | + if turn_lane_lease_holder(lease_target) is not None: |
| 123 | + yield None |
| 124 | + return |
96 | 125 | with try_exclusive_file_lock( |
97 | 126 | target, |
98 | | - agent_id=turn_lane_agent_id(plan), |
| 127 | + agent_id=agent_id, |
99 | 128 | operation=TURN_LANE_OPERATION, |
100 | 129 | ) as lock_path: |
101 | | - yield lock_path |
| 130 | + if lock_path is None: |
| 131 | + yield None |
| 132 | + return |
| 133 | + record = turn_lane_lease_record( |
| 134 | + agent_id=agent_id, turn_instance_id=turn_lane_turn_instance_id(plan) |
| 135 | + ) |
| 136 | + if not _claim_turn_lane_lease(lease_target, record): |
| 137 | + # Another host claimed the lane between the read above and this |
| 138 | + # process's kernel lock. The lease is the authority, so this Turn is |
| 139 | + # refused and its own kernel lock is released by the context manager. |
| 140 | + yield None |
| 141 | + return |
| 142 | + try: |
| 143 | + yield lock_path |
| 144 | + finally: |
| 145 | + _release_turn_lane_lease(lease_target, record) |
| 146 | + |
| 147 | + |
| 148 | +def turn_lane_lease_target(target: Path) -> Path: |
| 149 | + """Return the durable lease path that sits beside one lane's kernel lock.""" |
| 150 | + |
| 151 | + return target.with_name(f"{target.name}{TURN_LANE_LEASE_SUFFIX}") |
| 152 | + |
| 153 | + |
| 154 | +def turn_lane_host_fingerprint() -> str: |
| 155 | + """Return a public-safe fingerprint of the machine holding a lane lease.""" |
| 156 | + |
| 157 | + try: |
| 158 | + hostname = socket.gethostname() |
| 159 | + except OSError: # pragma: no cover - a host without a name is still a host |
| 160 | + hostname = "" |
| 161 | + return hashlib.sha256(hostname.encode("utf-8")).hexdigest()[:12] |
| 162 | + |
| 163 | + |
| 164 | +def turn_lane_turn_instance_id(plan: Mapping[str, Any]) -> str: |
| 165 | + """Return the Turn identity a lease names, taken from the Turn envelope.""" |
| 166 | + |
| 167 | + envelope = plan.get("turn_envelope") |
| 168 | + if isinstance(envelope, Mapping): |
| 169 | + for field in ("turn_instance_id", "turn_id", "run_id"): |
| 170 | + value = str(envelope.get(field) or "").strip() |
| 171 | + if value: |
| 172 | + return value |
| 173 | + return "" |
| 174 | + |
| 175 | + |
| 176 | +def turn_lane_lease_record( |
| 177 | + *, agent_id: str, turn_instance_id: str = "" |
| 178 | +) -> dict[str, Any]: |
| 179 | + """Build this process's durable claim on one lane.""" |
| 180 | + |
| 181 | + acquired = runtime_now_utc() |
| 182 | + return { |
| 183 | + "schema_version": TURN_LANE_LEASE_SCHEMA_VERSION, |
| 184 | + "agent_id": agent_id or TURN_LANE_UNATTRIBUTED_AGENT, |
| 185 | + "operation": TURN_LANE_OPERATION, |
| 186 | + "turn_instance_id": turn_instance_id, |
| 187 | + "acquired_at": utc_isoformat(acquired), |
| 188 | + "expires_at": utc_isoformat( |
| 189 | + acquired + timedelta(seconds=TURN_LANE_LEASE_TTL_SECONDS) |
| 190 | + ), |
| 191 | + "host_fingerprint": turn_lane_host_fingerprint(), |
| 192 | + "pid": os.getpid(), |
| 193 | + } |
| 194 | + |
| 195 | + |
| 196 | +def _pid_is_alive(pid: int) -> bool | None: |
| 197 | + """Report whether a recorded holder process still exists, when knowable.""" |
| 198 | + |
| 199 | + if pid <= 0: |
| 200 | + return None |
| 201 | + try: |
| 202 | + os.kill(pid, 0) |
| 203 | + except ProcessLookupError: |
| 204 | + return False |
| 205 | + except PermissionError: |
| 206 | + # It exists; this process is simply not allowed to signal it. |
| 207 | + return True |
| 208 | + except (OSError, AttributeError, ValueError, NotImplementedError): |
| 209 | + # The platform cannot answer. The lease's own expiry is then the only |
| 210 | + # thing that clears it, which the caller already applies. |
| 211 | + return None |
| 212 | + return True |
| 213 | + |
| 214 | + |
| 215 | +def _turn_lane_lease_is_live(record: Mapping[str, Any]) -> bool: |
| 216 | + """Report whether one lease record still holds its lane.""" |
| 217 | + |
| 218 | + expires_at = record.get("expires_at") |
| 219 | + if not isinstance(expires_at, str) or not expires_at: |
| 220 | + # A record without an expiry is not a claim this module can honor: it |
| 221 | + # would hold the lane forever, so it is treated as already released. |
| 222 | + return False |
| 223 | + try: |
| 224 | + deadline = parse_timestamp(expires_at) |
| 225 | + except (TypeError, ValueError): |
| 226 | + return False |
| 227 | + if deadline is None or runtime_now_utc() >= deadline: |
| 228 | + return False |
| 229 | + if str(record.get("host_fingerprint") or "") != turn_lane_host_fingerprint(): |
| 230 | + # A different machine is executing this lane. Only its own expiry can |
| 231 | + # clear the claim, because this process cannot observe its process. |
| 232 | + return True |
| 233 | + pid = record.get("pid") |
| 234 | + if isinstance(pid, bool) or not isinstance(pid, int): |
| 235 | + return True |
| 236 | + # The same machine: a holder that no longer exists is a crashed Turn, and |
| 237 | + # the lane is free. A live holder is still running, though the kernel lock |
| 238 | + # is what actually keeps two local processes apart. |
| 239 | + return _pid_is_alive(pid) is not False |
| 240 | + |
| 241 | + |
| 242 | +def turn_lane_lease_readback(record: Mapping[str, Any]) -> dict[str, Any]: |
| 243 | + """Return the public-safe identity of a durable lease holder.""" |
| 244 | + |
| 245 | + projection: dict[str, Any] = {} |
| 246 | + for field in TURN_LANE_LEASE_HOLDER_TEXT_FIELDS: |
| 247 | + value = record.get(field) |
| 248 | + if isinstance(value, str) and value: |
| 249 | + projection[field] = value |
| 250 | + pid = record.get("pid") |
| 251 | + if isinstance(pid, int) and not isinstance(pid, bool): |
| 252 | + projection["pid"] = pid |
| 253 | + return projection |
| 254 | + |
| 255 | + |
| 256 | +def turn_lane_lease_holder(target: Path) -> dict[str, Any] | None: |
| 257 | + """Return the readback of the holder of one lane, or ``None`` if it is free.""" |
| 258 | + |
| 259 | + record = _read_turn_lane_lease(target) |
| 260 | + if record is None or not _turn_lane_lease_is_live(record): |
| 261 | + return None |
| 262 | + return turn_lane_lease_readback(record) |
| 263 | + |
| 264 | + |
| 265 | +def _read_turn_lane_lease(target: Path) -> dict[str, Any] | None: |
| 266 | + try: |
| 267 | + payload = json.loads(target.read_text(encoding="utf-8")) |
| 268 | + except (OSError, ValueError): |
| 269 | + return None |
| 270 | + return dict(payload) if isinstance(payload, Mapping) else None |
| 271 | + |
| 272 | + |
| 273 | +def _claim_turn_lane_lease(target: Path, record: Mapping[str, Any]) -> bool: |
| 274 | + """Take one lane's durable lease, replacing a stale record. |
| 275 | +
|
| 276 | + The exclusive create is the claim: two hosts that both find the lane free |
| 277 | + cannot both succeed, so exactly one of them executes. A record that already |
| 278 | + exists is replaced only when it no longer holds the lane (expired, or left |
| 279 | + by a crashed process on this host). |
| 280 | + """ |
| 281 | + |
| 282 | + try: |
| 283 | + handle = os.open(target, os.O_CREAT | os.O_EXCL | os.O_WRONLY, 0o600) |
| 284 | + except FileExistsError: |
| 285 | + existing = _read_turn_lane_lease(target) |
| 286 | + if existing is not None and _turn_lane_lease_is_live(existing): |
| 287 | + return False |
| 288 | + _write_turn_lane_lease(target, record) |
| 289 | + return True |
| 290 | + except OSError: |
| 291 | + return False |
| 292 | + with os.fdopen(handle, "w", encoding="utf-8") as stream: |
| 293 | + json.dump(dict(record), stream, ensure_ascii=False, indent=2) |
| 294 | + stream.write("\n") |
| 295 | + return True |
| 296 | + |
| 297 | + |
| 298 | +def _write_turn_lane_lease(target: Path, record: Mapping[str, Any]) -> None: |
| 299 | + """Replace one lane's lease atomically.""" |
| 300 | + |
| 301 | + target.parent.mkdir(parents=True, exist_ok=True) |
| 302 | + temp_path = target.with_name(f".{target.name}.{os.getpid()}.tmp") |
| 303 | + temp_path.write_text( |
| 304 | + json.dumps(dict(record), ensure_ascii=False, indent=2) + "\n", |
| 305 | + encoding="utf-8", |
| 306 | + ) |
| 307 | + temp_path.replace(target) |
| 308 | + |
| 309 | + |
| 310 | +def _release_turn_lane_lease(target: Path, record: Mapping[str, Any]) -> None: |
| 311 | + """Release this process's lease, leaving any other holder's record alone.""" |
| 312 | + |
| 313 | + existing = _read_turn_lane_lease(target) |
| 314 | + if existing is None: |
| 315 | + return |
| 316 | + if str(existing.get("turn_instance_id") or "") != str( |
| 317 | + record.get("turn_instance_id") or "" |
| 318 | + ) or existing.get("pid") != record.get("pid"): |
| 319 | + return |
| 320 | + try: |
| 321 | + target.unlink() |
| 322 | + except OSError: # pragma: no cover - already released by a peer |
| 323 | + return |
102 | 324 |
|
103 | 325 |
|
104 | 326 | def turn_lane_holder_readback(target: Path) -> dict[str, Any]: |
@@ -188,11 +410,15 @@ def single_lane_turn( |
188 | 410 | ) as held: |
189 | 411 | if held is not None: |
190 | 412 | return execute_turn(plan, *args, **kwargs) |
| 413 | + # Either this machine's kernel lock or another host's durable |
| 414 | + # lease refused the Turn, and the durable lease is the only one |
| 415 | + # a caller on a shared runtime root can actually wait for. |
| 416 | + holder = turn_lane_lease_holder( |
| 417 | + turn_lane_lease_target(target) |
| 418 | + ) or turn_lane_holder_readback(target) |
191 | 419 | return execution_payload( |
192 | 420 | plan, |
193 | | - turn_lane_in_flight_record( |
194 | | - plan, holder=turn_lane_holder_readback(target) |
195 | | - ), |
| 421 | + turn_lane_in_flight_record(plan, holder=holder), |
196 | 422 | execute=True, |
197 | 423 | replayed=False, |
198 | 424 | effects=TURN_LANE_NO_EFFECTS, |
|
0 commit comments