← Files Meetings (Beta)ARCHIVED FILE

scripts/companion_restart_budget.py

9.08 KB · Oct 8, 2026 · 12:02 UTC

↓ Download file

"""A small, process-shared limit for automatic native-owner replacement."""

from __future__ import annotations

import os
import stat
import tempfile
import time
from collections.abc import Callable, Generator
from contextlib import contextmanager
from pathlib import Path

from control_protocol import ControlUnavailable, canonical_json

from helpers import is_json, is_json_array, parse_bounded_json

MAXIMUM_AUTOMATIC_RESTARTS = 3
AUTOMATIC_RESTART_WINDOW_SECONDS = 300
_RECEIPT = "automatic-restart-budget.json"
_LOCK = "automatic-restart-budget.lock"
_RECOVERY_LOCK = "listener-recovery.lock"
_MAXIMUM_BYTES = 512
_MAXIMUM_CLOCK = (1 << 63) - 1


def _identity(metadata: os.stat_result) -> tuple[int, int]:
    return metadata.st_dev, metadata.st_ino


def _lock(descriptor: int) -> None:
    if os.name == "nt":
        import msvcrt

        os.lseek(descriptor, 0, os.SEEK_SET)
        msvcrt.locking(descriptor, msvcrt.LK_NBLCK, 1)
    else:
        import fcntl

        fcntl.flock(descriptor, fcntl.LOCK_EX | fcntl.LOCK_NB)


def _open_private(path: Path) -> tuple[int, bool]:
    import control_client

    flags = (
        os.O_RDWR
        | getattr(os, "O_BINARY", 0)
        | getattr(os, "O_CLOEXEC", 0)
        | getattr(os, "O_NOFOLLOW", 0)
        | getattr(os, "O_NONBLOCK", 0)
    )
    try:
        return os.open(path, flags | os.O_CREAT | os.O_EXCL, 0o600), True
    except FileExistsError:
        control_client.require_private(path, directory=False)
        return os.open(path, flags), False


def _validate_file(
    root: Path,
    root_identity: tuple[int, int],
    path: Path,
    descriptor: int,
    opened: os.stat_result,
) -> None:
    import control_client

    current = control_client.require_private(path, directory=False)
    held = os.fstat(descriptor)
    if (
        _identity(control_client.require_private(root, directory=True)) != root_identity
        or _identity(current) != _identity(opened)
        or _identity(held) != _identity(opened)
        or not stat.S_ISREG(held.st_mode)
        or held.st_nlink != 1
        or held.st_size > _MAXIMUM_BYTES
    ):
        raise ControlUnavailable("automatic recovery file changed")


@contextmanager
def _private_lock(
    root: Path, filename: str
) -> Generator[tuple[tuple[int, int], Callable[[], None]], None, None]:
    import control_client

    descriptor = -1
    try:
        root_identity = _identity(control_client.require_private(root, directory=True))
        path = root / filename
        descriptor, _ = _open_private(path)
        opened = os.fstat(descriptor)

        def revalidate() -> None:
            _validate_file(root, root_identity, path, descriptor, opened)

        revalidate()
        _lock(descriptor)
        revalidate()
        yield root_identity, revalidate
    except OSError as error:
        raise ControlUnavailable(
            "automatic recovery is unavailable or already in progress"
        ) from error
    finally:
        if descriptor >= 0:
            # Closing the sole descriptor releases flock / the Windows byte lock.
            os.close(descriptor)


@contextmanager
def automatic_recovery_guard(root: Path) -> Generator[Callable[[], None], None, None]:
    """Coalesce the whole health-check/termination attempt across MCP processes.

    Call the yielded validator at destructive boundaries. A concurrent attempt
    cannot overwrite the shared challenge while this guard remains held.
    """
    with _private_lock(root, _RECOVERY_LOCK) as (_, revalidate):
        yield revalidate


def _claims(payload: bytes) -> list[int]:
    value = parse_bounded_json(payload.decode("utf-8"))
    if (
        not is_json(value)
        or set(value) != {"version", "monotonicNanoseconds"}
        or type(value.get("version")) is not int
        or value["version"] != 1
    ):
        raise ValueError("invalid automatic restart budget")
    raw = value["monotonicNanoseconds"]
    if not is_json_array(raw) or not 1 <= len(raw) <= MAXIMUM_AUTOMATIC_RESTARTS:
        raise ValueError("invalid automatic restart budget")
    result: list[int] = []
    for item in raw:
        if type(item) is not int or not 0 <= item <= _MAXIMUM_CLOCK:
            raise ValueError("invalid automatic restart budget")
        if result and item < result[-1]:
            raise ValueError("invalid automatic restart budget")
        result.append(item)
    return result


def _read_receipt(
    root: Path, root_identity: tuple[int, int]
) -> tuple[bytes, tuple[int, int]] | None:
    import control_client

    path = root / _RECEIPT
    try:
        path.lstat()
    except FileNotFoundError:
        return None
    expected = control_client.require_private(path, directory=False)
    descriptor = os.open(
        path,
        os.O_RDONLY
        | getattr(os, "O_BINARY", 0)
        | getattr(os, "O_CLOEXEC", 0)
        | getattr(os, "O_NOFOLLOW", 0)
        | getattr(os, "O_NONBLOCK", 0),
    )
    try:
        _validate_file(root, root_identity, path, descriptor, expected)
        payload = os.read(descriptor, _MAXIMUM_BYTES + 1)
        _validate_file(root, root_identity, path, descriptor, expected)
        if len(payload) > _MAXIMUM_BYTES:
            raise ControlUnavailable("automatic recovery file is too large")
        return payload, _identity(expected)
    finally:
        os.close(descriptor)


def claim_automatic_restart(root: Path, *, monotonic_ns: int | None = None) -> bool:
    """Durably reserve one attempt before termination; never poll on healthy I/O.

    The system monotonic clock is shared across Python 3.10+ processes. A lower
    clock origin clears the previous boot's claims, whose owners cannot survive
    reboot. Wall-clock corrections, owner epochs and successful connections do
    not replenish the rolling budget. Contention and unsafe receipts fail closed.
    """
    import control_client

    try:
        # The independent empty lock remains stable while the receipt is first
        # initialized. Competing creators cannot mistake partial initialization
        # for a malformed receipt or strand the creator behind another reader.
        with _private_lock(root, _LOCK) as (root_identity, revalidate_lock):
            path = root / _RECEIPT
            existing = _read_receipt(root, root_identity)
            previous = [] if existing is None else _claims(existing[0])
            now = time.monotonic_ns() if monotonic_ns is None else monotonic_ns
            if type(now) is not int or not 0 <= now <= _MAXIMUM_CLOCK:
                return False
            if previous and now < previous[-1]:
                previous = []
            window = AUTOMATIC_RESTART_WINDOW_SECONDS * 1_000_000_000
            retained = [claim for claim in previous if now - claim < window]
            if len(retained) >= MAXIMUM_AUTOMATIC_RESTARTS:
                return False
            retained.append(now)
            updated = canonical_json({"version": 1, "monotonicNanoseconds": retained}).encode(
                "utf-8"
            )
            revalidate_lock()
            descriptor, temporary_name = tempfile.mkstemp(
                prefix=".automatic-restart-budget-", dir=root
            )
            temporary = Path(temporary_name)
            created = os.fstat(descriptor)
            try:
                # An incomplete write stays unpublished. Validate inherited
                # Windows ACLs before bytes, then close before atomic rename.
                _validate_file(root, root_identity, temporary, descriptor, created)
                if os.write(descriptor, updated) != len(updated):
                    return False
                os.fsync(descriptor)
                _validate_file(root, root_identity, temporary, descriptor, created)
                os.close(descriptor)
                descriptor = -1
                revalidate_lock()
                current = control_client.require_private(temporary, directory=False)
                if _identity(current) != _identity(created) or current.st_nlink != 1:
                    return False
                if _read_receipt(root, root_identity) != existing:
                    return False
                revalidate_lock()
                os.replace(temporary, path)
                if _read_receipt(root, root_identity) != (updated, _identity(created)):
                    return False
            finally:
                if descriptor >= 0:
                    os.close(descriptor)
                try:
                    current = control_client.require_private(temporary, directory=False)
                    revalidate_lock()
                    if _identity(current) == _identity(created) and current.st_nlink == 1:
                        temporary.unlink()
                except (OSError, ControlUnavailable):
                    pass
            if os.name != "nt":
                directory = os.open(root, os.O_RDONLY | os.O_DIRECTORY | os.O_NOFOLLOW)
                try:
                    if _identity(os.fstat(directory)) != root_identity:
                        return False
                    os.fsync(directory)
                finally:
                    os.close(directory)
            return True
    except (OSError, ValueError, RecursionError, ControlUnavailable):
        return False

SHA-256: f94f07313c42cf9abaf5b27ae3ec8fb06d9625e3390de7bf1a733b8c978929c3