import fcntl
import threading
import time
from contextlib import contextmanager
from contextvars import ContextVar
from pathlib import Path
from typing import IO

from engine import waits

MIGRATIONS = ContextVar("migrations", default=())
MIGRATION_LOCK = ".migrations.lock"
LOCK_WAIT = 30.0
RUNTIME = "runtime"


def journal_roots(path: Path) -> tuple[Path, ...]:
    return tuple(parent for parent in path.parents if parent.name == ".journal")


class SharedWrites:
    def __init__(self, held: IO):
        self.held = held
        self.guard = waits.Lock("writes")
        self.writers = 0

    @contextmanager
    def joined(self):
        with self.guard:
            if not self.writers:
                with waits.waited("writes"):
                    acquire(self.held, fcntl.LOCK_SH)
            self.writers += 1
        try:
            yield
        finally:
            with self.guard:
                self.writers -= 1
                if not self.writers:
                    fcntl.flock(self.held, fcntl.LOCK_UN)


SHARED: dict[Path, SharedWrites] = {}
SHARING = threading.Lock()


def shared_writes(root: Path) -> SharedWrites:
    with SHARING:
        if root not in SHARED:
            root.mkdir(parents=True, exist_ok=True)
            SHARED[root] = SharedWrites((root / MIGRATION_LOCK).open("a"))
        return SHARED[root]


def acquire(held, operation: int) -> None:
    deadline = time.monotonic() + LOCK_WAIT
    while True:
        try:
            fcntl.flock(held, operation | fcntl.LOCK_NB)
            return
        except BlockingIOError as error:
            if time.monotonic() >= deadline:
                raise TimeoutError(held.name) from error
            time.sleep(0.05)


def claim(path: Path):
    path.parent.mkdir(parents=True, exist_ok=True)
    held = path.open("a")
    try:
        fcntl.flock(held, fcntl.LOCK_EX | fcntl.LOCK_NB)
    except BlockingIOError:
        held.close()
        return None
    return held


@contextmanager
def hold_record_writes(root: Path):
    root = Path(root)
    root.mkdir(parents=True, exist_ok=True)
    with (root / MIGRATION_LOCK).open("a") as held:
        acquire(held, fcntl.LOCK_EX)
        token = MIGRATIONS.set((*MIGRATIONS.get(), root))
        try:
            yield
        finally:
            MIGRATIONS.reset(token)
            fcntl.flock(held, fcntl.LOCK_UN)


def guarded(path: Path, roots: tuple[Path, ...]) -> bool:
    return bool(roots) and roots[0] not in MIGRATIONS.get() and path.relative_to(roots[0]).parts[0] != RUNTIME


@contextmanager
def writing(path: Path):
    roots = journal_roots(path)
    if not guarded(path, roots):
        yield
        return
    with shared_writes(roots[0]).joined():
        yield
import threading
from contextlib import contextmanager
from pathlib import Path

from engine.disk import replace


class Work(threading.local):
    def __init__(self):
        self.undo_snapshots: dict | None = None
        self.undo_releases: list | None = None
        self.bus_queue: list | None = None
        self.bus_after: dict | None = None


WORK = Work()


@contextmanager
def undoable():
    if WORK.undo_snapshots is not None:
        yield
        return
    WORK.undo_snapshots, WORK.undo_releases = {}, []
    try:
        yield
    except BaseException:
        for path, before in WORK.undo_snapshots.items():
            if before is None:
                Path(path).unlink(missing_ok=True)
            else:
                replace(Path(path), before)
        WORK.undo_snapshots, WORK.undo_releases = None, None
        raise
    events, WORK.undo_snapshots, WORK.undo_releases = WORK.undo_releases, None, None
    for release in events:
        release()


@contextmanager
def apart():
    held = WORK.undo_snapshots, WORK.undo_releases
    WORK.undo_snapshots, WORK.undo_releases = None, None
    try:
        yield
    finally:
        WORK.undo_snapshots, WORK.undo_releases = held


def held_back(release) -> bool:
    if WORK.undo_releases is None:
        return False
    WORK.undo_releases.append(release)
    return True


def snapshot(path: Path) -> None:
    if WORK.undo_snapshots is not None and str(path) not in WORK.undo_snapshots:
        WORK.undo_snapshots[str(path)] = path.read_bytes() if path.is_file() else None
import os
import signal
import time
from contextlib import suppress
from pathlib import Path
from engine import runtime
from engine.proc import ran
from engine.sessions import alive
from engine.stored import write_text
from engine.viewer import last, running

WAIT = 15.0
EVERY = 0.2
ESCALATE = 5.0


def flag(root: Path) -> Path:
    return runtime.folder(root) / "stop"


def at(root: Path) -> float:
    try:
        text = flag(root).read_text().strip()
        return float(text) if text else 0.0
    except (OSError, ValueError):
        return 0.0


def asked(root: Path, since: float = 0.0) -> bool:
    return at(root) > since


def session_flag(root: Path, terminal: str) -> Path:
    return runtime.session_file(Path(root), terminal, "stop")


def ask_session(root: Path, terminal: str) -> None:
    write_text(session_flag(root, terminal), f"{time.time()}\n")


def ask(root: Path) -> None:
    where = flag(root)
    where.parent.mkdir(parents=True, exist_ok=True)
    write_text(where, f"{time.time()}\n")


def clear(root: Path) -> None:
    flag(root).unlink(missing_ok=True)


def gone(root: Path, seconds: float = WAIT) -> bool:
    server = serving(Path(root))
    end = time.monotonic() + seconds
    while time.monotonic() < end:
        if not (alive(server) if server else running(Path(root))):
            return True
        time.sleep(EVERY)
    return False


def ended(root: Path) -> bool:
    if gone(root):
        return True
    for signal_, seconds in ((signal.SIGTERM, ESCALATE), (signal.SIGKILL, ESCALATE)):
        server = serving(Path(root))
        if not server:
            return not running(Path(root))
        with suppress(OSError):
            os.kill(server, signal_)
        if gone(root, seconds):
            return True
    return False


def serving(root: Path) -> int:
    pid = last(root).pid
    listed = ran(["ps", "-o", "command=", "-p", str(pid)]) if pid else None
    command = listed.stdout if listed else ""
    return pid if " serve" in command and any(form in command for form in (str(root), str(root.resolve()))) else 0
import threading
import time
from pathlib import Path

from engine import runtime
from engine.stored import write_text
from install import fetch, released, version_key

__all__ = ["fetch"]

UPSTREAM_FOR = 900
FETCHING = threading.Lock()


def stale(root: Path) -> bool:
    try:
        return time.time() - runtime.upstream_cache(root).stat().st_mtime >= UPSTREAM_FOR
    except OSError:
        return True


def upstream(root: Path) -> str:
    cache = runtime.upstream_cache(root)
    try:
        held = cache.read_text().strip()
    except OSError:
        held = ""
    if stale(root) and not FETCHING.locked():
        threading.Thread(target=fetched, args=(cache,), daemon=True).start()
    return held


def fetched(cache: Path) -> None:
    with FETCHING:
        kept_newest(cache)


def kept_newest(cache: Path) -> None:
    cache.parent.mkdir(parents=True, exist_ok=True)
    newest = released()
    if not newest:
        cache.touch(exist_ok=True)
        return
    write_text(cache, newest)


def check_now(root: Path) -> None:
    if not FETCHING.acquire(blocking=False):
        return

    def check() -> None:
        try:
            kept_newest(runtime.upstream_cache(root))
        finally:
            FETCHING.release()

    threading.Thread(target=check, daemon=True).start()


def newer(version: str, than: str) -> bool:
    return bool(version) and (not than or than == "0" or version_key(version) > version_key(than))
import json
import os
import threading
from pathlib import Path
from typing import Any, Callable, TypeVar

from engine.memo import Memo

T = TypeVar("T")
LOG_BYTES = 262144


class Growth:
    def __init__(self):
        self.size = -1

    def grew(self, path: Path) -> bool:
        try:
            size = path.stat().st_size
        except OSError:
            return False
        if size == self.size:
            return False
        self.size = size
        return True


class JsonFiles:
    def __init__(self):
        self.held = Memo()

    def read(self, path: Path, into: Callable[[Any], T], default: T) -> T:
        try:
            found = path.stat()
        except OSError:
            return default
        return self.held.get(path, (found.st_mtime_ns, found.st_size), lambda: read_json(path, into, default))


def read_json(path: Path, into: Callable[[Any], T], default: T) -> T:
    try:
        return into(json.loads(Path(path).read_text()))
    except (OSError, ValueError, TypeError, KeyError, AttributeError):
        return default


def replace(path: Path, raw: bytes) -> None:
    path.parent.mkdir(parents=True, exist_ok=True)
    spare = path.with_name(f".{path.name}.{os.getpid()}.{threading.get_ident()}")
    spare.write_bytes(raw)
    os.replace(spare, path)


def last_lines(path, lines: int) -> str:
    try:
        with Path(path).open("rb") as source:
            start = max(0, source.seek(0, 2) - LOG_BYTES)
            source.seek(start)
            raw = source.read()
    except (OSError, TypeError):
        return ""
    if start:
        raw = raw.split(b"\n", 1)[-1]
    return "\n".join(raw.decode(errors="replace").splitlines()[-lines:])
from pathlib import Path

from resources.base import Refused

ENVIRONMENTS = "environments"
ROUTED = frozenset({"agent-controls", "agent-hooks", "hook", "journals", "plugins", "services", "update"})


def contained(folder: Path, name: str, nested: bool = False) -> Path:
    if not isinstance(name, str) or not name or "\\" in name or "\x00" in name:
        raise Refused(f"invalid path {name!r}")
    part = Path(name)
    if part.is_absolute() or any(piece in ("", ".", "..") for piece in name.split("/")) or (not nested and len(part.parts) != 1):
        raise Refused(f"invalid path {name!r}")
    target = folder / part
    if not target.resolve().is_relative_to(folder.resolve()):
        raise Refused(f"path {name!r} leaves its folder")
    return target


def environments(root: Path) -> Path:
    return Path(root) / ENVIRONMENTS


def environment_home(root: Path, name: str) -> Path:
    return contained(environments(root), name)
event_log.py:15:class Recent:
event_log.py:24:def parsed(lines):
event_log.py:32:class EventLog:
event_log.py:33:    def __init__(self, home: Path, locked: Callable[[], AbstractContextManager]):
event_log.py:38:    def append(self, event: Event) -> None:
event_log.py:41:    def events(self, since: int = 0, last: int = 0) -> list[Event]:
event_log.py:48:    def recent(self) -> list[Event]:
event_log.py:70:    def back(self, since: int = 0, last: int = 0) -> list[Event]:
event_log.py:78:    def lines_back(self, block: int = 65536):
event_log.py:94:    def last_id(self) -> int:
event_log.py:98:    def trim(self, keep: int, readers_since: float) -> int:
event_log.py:113:    def cursor_file(self, name: str) -> Path:
event_log.py:116:    def cursor_text(self, name: str) -> str:
event_log.py:122:    def set_cursor_text(self, name: str, text: str) -> None:
event_log.py:127:    def cursor(self, name: str) -> int:
event_log.py:134:    def set_cursor(self, name: str, n: int) -> None:
keeper.py:25:def known(cls, raw: dict) -> dict:
keeper.py:33:class ServiceSpec:
keeper.py:55:    def from_json(cls, raw: dict) -> "ServiceSpec":
keeper.py:59:    def build(self) -> str:
keeper.py:64:class ServiceState:
keeper.py:80:    def read(cls, path: Path) -> "ServiceState":
keeper.py:84:    def write(self, path: Path) -> None:
keeper.py:88:def answers(port: int, path: str) -> bool:
keeper.py:106:def gone(group: int) -> bool:
keeper.py:116:def teardown(group: int, grace: float) -> None:
keeper.py:129:def state(spec: ServiceSpec, name: str, child=None, **more) -> None:
keeper.py:134:def lease(spec: ServiceSpec):
keeper.py:148:def watch(spec: ServiceSpec, child, lifeline: int, stopping) -> str:
keeper.py:160:def main(argv: list[str]) -> int:
services.py:32:def status_file(root: Path, sid: str) -> Path:
services.py:36:def holder(root: Path, sid: str) -> int:
services.py:44:def lock_file(root: Path, sid: str) -> Path:
services.py:48:def current_build(root: Path) -> str:
services.py:52:def spec_file(root: Path, sid: str) -> Path:
services.py:56:def want_file(root: Path, sid: str) -> Path:
services.py:60:def log_file(root: Path, sid: str) -> Path:
services.py:65:class Wanted(Loaded):
services.py:70:    def read(cls, root: Path, sid: str) -> "Wanted":
services.py:74:def status(root: Path, sid: str) -> ServiceState:
services.py:78:def states(root: Path) -> dict[str, ServiceState]:
services.py:84:def wanted(root: Path, sid: str) -> str:
services.py:88:def want(root: Path, sid: str, state: str, nonce: float = 0.0) -> dict:
services.py:94:def allocate(root: Path, sid: str, wants, taken: set[int]) -> tuple[int, str]:
services.py:106:def claimed(root: Path, sid: str, wants, taken: set[int]) -> tuple[int, str]:
services.py:113:def local_url(port: int) -> str:
services.py:120:class ServiceFiles(TypedDict):
services.py:127:def files_for(root: Path, sid: str) -> ServiceFiles:
services.py:131:def service_spec(root: Path, sid: str, **fields) -> ServiceSpec:
services.py:135:def specs(root: Path, sources) -> list[ServiceSpec]:
services.py:140:def spawn(spec: ServiceSpec, lifeline: int) -> int:
services.py:150:def excerpt(output: str) -> str:
services.py:155:class Manager:
services.py:156:    def __init__(self, root: Path, lifeline: int = -1, start=spawn, clock=time.time, living=alive, sources=()):
services.py:169:    def tick(self) -> list[str]:
services.py:182:    def retire(self, declared: set[str]) -> list[str]:
services.py:189:    def remove(self, sid: str) -> None:
services.py:203:    def stored_build(self, sid: str) -> str:
services.py:207:    def one(self, spec: ServiceSpec) -> bool:
services.py:253:    def unneeded(self, spec: ServiceSpec, now: float) -> str:
services.py:267:    def crashed(self, sid: str, now: float) -> None:
services.py:272:    def stop(self, sid: str, current: ServiceState) -> None:
services.py:279:    def sweep(self) -> list[int]:
services.py:291:def listed(root: Path, sources) -> list[dict]:
worktree.py:23:class WorkspaceFolders:
worktree.py:31:    def managed(self) -> tuple[str, ...]:
worktree.py:35:    def worktree_home(self) -> tuple[str, ...]:
worktree.py:39:def checkout(start: Path, folders: WorkspaceFolders) -> Path | None:
worktree.py:52:def environment(top: Path | None) -> str:
worktree.py:64:def unused_name(project: Path) -> str:
worktree.py:71:def linked(project: Path) -> dict[str, Path]:
worktree.py:77:def main_checkout(start: Path) -> Path:
worktree.py:82:def opened(project: Path, folder: Path, branch: str, name: str = "") -> Path:
worktree.py:100:def repositories(folder: Path) -> list[Path]:
worktree.py:105:def spread(project: Path) -> bool:
worktree.py:109:def roots(project: Path) -> dict[str, Path]:
worktree.py:113:def changed(project: Path, branch: str, base: str) -> bool:
worktree.py:118:def workspace(project: Path, folder: Path, folders: WorkspaceFolders) -> Path:
worktree.py:138:def keep(project: Path, name: str, branch: str) -> None:
worktree.py:143:def merged(project: Path, branch: str, base: str, into: str = "HEAD") -> bool:
worktree.py:151:def current_branch(project: Path) -> str:
worktree.py:155:def checked_out(project: Path, branch: str) -> Path | None:
worktree.py:161:def merged_into(project: Path, branch: str, into: str) -> str:
worktree.py:176:def contains(project: Path, commit: str, branch: str) -> bool:
worktree.py:180:def branched(project: Path, branch: str, start: str, fresh: bool = False) -> str:
worktree.py:202:def tip(project: Path, ref: str = "HEAD") -> str:
worktree.py:210:def present(project: Path, ref: str) -> bool:
worktree.py:214:def matched(project: Path, pattern: str) -> list[Path]:
worktree.py:222:def included(project: Path, folder: Path) -> None:
worktree.py:232:def git(project: Path, *args: str) -> subprocess.CompletedProcess:
worktree.py:237:def lines(project: Path, *args: str) -> list[str]:
worktree.py:245:def share_journal(top: Path, root: Path, folders: WorkspaceFolders) -> None:
worktree.py:263:def cleared(top: Path, path: Path) -> None:
worktree.py:273:def belongs(top: Path, project: Path) -> bool:
worktree.py:283:def ignored(project: Path, paths: list[Path], folders: WorkspaceFolders) -> list[Path]:
worktree.py:294:def untracked(project: Path, paths: list[Path]) -> list[Path]:
worktree.py:301:def linked_to(place: Path, target: Path) -> None:
worktree.py:315:def is_linked(place: Path, target: Path) -> bool:
worktree.py:322:def unshared(top: Path, project: Path, wanted: set[Path], folders: WorkspaceFolders) -> None:
worktree.py:333:def excluded(top: Path, patterns: list[str], folders: WorkspaceFolders) -> None:
worktree.py:339:def ignore(exclude: Path, patterns: list[str], managed: tuple = ()) -> None:
project_files.py:21:class ProjectSource:
project_files.py:29:def project_path(project: Path, asked: str) -> Path:
project_files.py:37:def readable_path(project: Path, target: Path) -> bool:
project_files.py:45:def read_source(project: Path, asked: str) -> ProjectSource:
project_files.py:59:def walked(project: Path) -> tuple[list[Path], dict[str, list[str]]]:
project_files.py:69:def walk(project: Path) -> tuple[list[Path], dict[str, list[str]]]:
project_files.py:89:def project_paths(project: Path) -> list[Path]:
project_files.py:93:def matching(project: Path, asked: str) -> list[str]:
project_files.py:101:def list_folder(project: Path, asked: str) -> list[dict]:
files.py:19:def change_kind(path: str, last: dict, now: dict) -> str:
files.py:25:def internal(record, project: Path, homes: tuple[str, ...]) -> tuple[str, ...]:
files.py:38:def journals_own(path: str, marks: tuple[str, ...]) -> bool:
files.py:43:class LineCount:
files.py:48:    def between(cls, old: str, new: str) -> "LineCount":
files.py:54:class Indexed:
files.py:60:class Hashed:
files.py:67:class FoundFile:
files.py:84:def project_repositories(project: Path) -> tuple[Path, ...]:
files.py:93:def nested_repositories(folder: Path, depth: int) -> list[Path]:
files.py:103:def prefixed(project: Path, repository: Path, paths: dict) -> dict:
files.py:108:def index_file(project: Path) -> Path:
files.py:112:def tracked_in(project: Path) -> dict:
files.py:129:def unchanged(held: Hashed | None, stat) -> bool:
files.py:133:def hashed(project: Path, paths: list[str]) -> dict:
files.py:142:def blobs(record, project: Path, homes: tuple[str, ...]) -> dict:
files.py:151:def blobs_in(project: Path) -> dict:
files.py:162:def blob_texts(project: Path, shas: list[str]) -> dict[str, str]:
files.py:172:def project_paths(project: Path) -> set[str]:
files.py:176:def paths_in(project: Path) -> set[str]:
files.py:182:def found_files(project: Path, needle: str) -> list[FoundFile]:
files.py:189:def match_rank(path: str, wanted: str) -> int:
files.py:196:def line_counts(project: Path, pairs: list[tuple[str, str]]) -> dict[tuple[str, str], LineCount]:
files.py:201:class Coalesced:
files.py:202:    def __init__(self):
files.py:207:    def run(self, key, job) -> None:
files.py:230:def announce_writes(record, agent: int, homes: tuple[str, ...]) -> None:
files.py:234:def announce(record, agent: int, homes: tuple[str, ...]) -> None:
heal.py:12:def ledger(root: Path) -> Path:
heal.py:16:def broken(root: Path) -> list[str]:
heal.py:20:def refused(root: Path, version: str) -> bool:
heal.py:25:def heal(root: Path) -> str:
