import errno
import socket
import time
from pathlib import Path
from engine.wording import digest

PACKET = 2048
FULL_FOR = 2.0


def folder(root: Path) -> Path:
    return Path("/tmp") / f"journal-{digest(str(Path(root).resolve()), 16)}"


def path(root: Path, session: str) -> Path:
    return folder(root) / f"typist-{session}.sock"


def listen(root: Path, session: str) -> socket.socket:
    where = path(root, session)
    where.parent.mkdir(parents=True, exist_ok=True)
    where.unlink(missing_ok=True)
    inbox = socket.socket(socket.AF_UNIX, socket.SOCK_DGRAM)
    inbox.bind(str(where))
    inbox.setblocking(False)
    return inbox


def receive(inbox: socket.socket) -> list[bytes]:
    packets = []
    while True:
        try:
            packets.append(inbox.recv(PACKET))
        except (BlockingIOError, InterruptedError):
            return packets


def close(inbox: socket.socket, root: Path, session: str) -> None:
    inbox.close()
    path(root, session).unlink(missing_ok=True)


def send(root: Path, session: str, raw: bytes) -> bool:
    mouth = socket.socket(socket.AF_UNIX, socket.SOCK_DGRAM)
    try:
        return all(sent(mouth, raw[at:at + PACKET], str(path(root, session))) for at in range(0, len(raw), PACKET))
    finally:
        mouth.close()


def sent(mouth: socket.socket, packet: bytes, where: str) -> bool:
    until = time.time() + FULL_FOR
    while True:
        try:
            mouth.sendto(packet, where)
            return True
        except (BlockingIOError, InterruptedError):
            pass
        except OSError as e:
            if e.errno != errno.ENOBUFS or time.time() > until:
                return False
        if time.time() > until:
            return False
        time.sleep(0.01)


def live(root: Path) -> list[str]:
    return sorted(p.name.removeprefix("typist-").removesuffix(".sock") for p in folder(root).glob("typist-*.sock") if reachable(p))


def reachable(where: Path) -> bool:
    mouth = socket.socket(socket.AF_UNIX, socket.SOCK_DGRAM)
    try:
        mouth.connect(str(where))
        return True
    except OSError:
        return False
    finally:
        mouth.close()
import json
import os
import subprocess
from dataclasses import dataclass, replace
from pathlib import Path

from engine.sessions import ACTIVE_ENV, Sessions, hold_build
from engine import runtime
from engine.stored import read_json, write_json
from engine.package import CODE, code, entry_in
from engine.fields import Loaded
from engine.worktree import checkout, environment, share_journal
from typing import TypedDict

from supervisor import LAUNCHED

LAUNCH = 2
CARRIED = "AGENT_JOURNAL_CARRIED"
LAUNCH_ARGS: list = []
OUTPUT_LINES: list = []


@dataclass(frozen=True)
class Launched(Loaded):
    pid: int = 0
    command: tuple = ()
    args: tuple = ()
    cwd: str = ""
    launch: int = 0

    @classmethod
    def read(cls, root: Path, session: str) -> "Launched":
        return read_json(runtime.session_file(root, session, LAUNCHED), cls.from_json, cls.from_json({}))


def agent_environment(base: dict | None = None, env: str | None = None, capped: dict | None = None) -> dict:
    made = {**(base if base is not None else os.environ), ACTIVE_ENV: "1"}
    if env is not None:
        made["JOURNAL_ENV"] = env
    return {**made, **(capped or {})}


def session_named(provider) -> dict:
    return {"JOURNAL_SESSION_VARIABLE": provider.session_variable} if provider.session_variable else {}


def output_lines(record) -> int:
    return max((int(lines(record)) for lines in OUTPUT_LINES), default=0)


def shaped_args(record, agent: str, args: list[str]) -> list[str]:
    for shape in LAUNCH_ARGS:
        args = shape(record, agent, args)
    return args


def output_cap(root: Path, env: str, provider) -> dict:
    from engine.record import Record
    lines = output_lines(Record(root, env))
    wrapper = provider.shell_wrapper(code(root) / "output_cap.sh") if lines > 0 else {}
    return {**wrapper, "JOURNAL_OUTPUT_LINES": str(lines), "JOURNAL_PROVIDER": provider.name, "JOURNAL_OUTPUT_DIR": str(runtime.folder(root) / "outputs")} if wrapper else {}


class Launching(TypedDict):
    command: list[str]
    args: list[str]
    launch: int
    exit: str
    environ: dict[str, str]


def launching(root: Path, cwd: Path, env: str, agent: str, args: list[str], conversation: str = "") -> Launching:
    from providers import DRIVERS, PROVIDERS
    from engine.record import Record
    driver = DRIVERS[agent]
    named = driver.command(driver.resumed(shaped_args(Record(root, env), agent, args), conversation), cwd)
    command = [driver.binary(os.environ.get("PATH", "")), *named[1:]]
    provider = PROVIDERS[agent]()
    inherited = {name: value for name, value in os.environ.items() if name not in provider.session_markers}
    return {"command": command, "args": args, "launch": LAUNCH, "exit": "" if driver.worktree(args) else driver.EXIT,
            "environ": agent_environment(inherited, env=env, capped={**output_cap(root, env, provider), **session_named(provider)})}


@dataclass(frozen=True)
class TerminalSession:
    root: Path
    env: str
    agent: str
    session: str


def seat_session(sessions: Sessions, env: str, session: str, **bound) -> None:
    from controllers.types import Environments
    from engine.record import Record
    from resources.base import SYSTEM
    sessions.bind(session, env, **bound)
    Environments(Record(sessions.root, env), actor=SYSTEM)._seat(env, session)


def seated(seat: TerminalSession) -> TerminalSession:
    launched = Launched.read(seat.root, seat.session)
    sessions = Sessions(seat.root)
    if not sessions.known(seat.session):
        seat_session(sessions, seat.env, seat.session)
    sessions.write(seat.session, pid=launched.pid, provider=seat.agent, args=list(launched.command), launch=launched.launch)
    return replace(seat, env=sessions.environment(seat.session) or seat.env)


def relaunch(root: Path, env: str, session: str, conversation: str) -> Path:
    launched = Launched.read(root, session)
    provider = Sessions(root).read(session).provider
    agent = provider if provider else session.split("-", 1)[0]
    cwd = Path(launched.cwd) if launched.cwd else root.parent
    asked = runtime.relaunch_file(root, session)
    write_json(asked, launching(root, cwd, env, agent, list(launched.args), conversation))
    return asked


def lifeline() -> tuple[int, int]:
    read, write = os.pipe()
    os.set_inheritable(read, True)
    os.set_inheritable(write, False)
    return read, write


def carried() -> dict | None:
    given = os.environ.pop(CARRIED, "")
    return json.loads(given) if given else None


def launch_spec(root: Path, cwd: Path, env: str, agent: str, args: list[str], taken: dict | None = None, conversation: str = "") -> dict:
    from providers import DRIVERS, workspace_folders
    if not taken:
        cwd, args = DRIVERS[agent].placed(cwd, args)
        top = checkout(cwd, workspace_folders())
        if top:
            share_journal(top, root, workspace_folders())
        worked = environment(top)
        env = Sessions(root).free(worked) if worked else env
    journal = [*entry_in(root, "journal"), "--root", str(root)]
    return {"root": str(root), "cwd": str(cwd), "env": env, "agent": agent,
            "worker": entry_in(root, "worker"), "heal": [*journal, "heal"], "ended": [*journal, "--env", env, "ended"],
            **({"adopt": {"pid": taken["pid"], "fd": taken["fd"], "session": taken["session"], "saved": taken["saved"]}, "args": args}
               if taken else launching(root, cwd, env, agent, args, conversation))}


def supervise(root: Path, cwd: Path, env: str, agent: str, args: list[str], taken: dict | None = None) -> None:
    hold_build(root, CODE)
    spec = launch_spec(root, cwd, env, agent, args, taken)
    if taken:
        os.set_inheritable(taken["fd"], True)
    else:
        print(f"journal: environment {spec['env']}")
    started = entry_in(root, "supervisor")
    os.execv(started[0], [*started, json.dumps(spec)])


def launch_log(root: Path, env: str) -> Path:
    return runtime.folder(root) / "launches" / f"{env}.log"


def detached(root: Path, cwd: Path, env: str, agent: str, args: list[str], conversation: str = "") -> int:
    started = entry_in(root, "supervisor")
    log = launch_log(root, env)
    log.parent.mkdir(parents=True, exist_ok=True)
    with log.open("ab") as kept:
        child = subprocess.Popen([*started, json.dumps({**launch_spec(root, cwd, env, agent, args, conversation=conversation), "headless": True})],
                                 stdin=subprocess.DEVNULL, stdout=kept, stderr=kept, start_new_session=True)
    hold_build(root, CODE, child.pid)
    return child.pid
import base64
import os
import select
import sys
import termios
import tty
from dataclasses import dataclass
from pathlib import Path

from engine import runtime, typist
from engine.stored import read_json
from supervisor import SCREEN, SCREEN_SHAPE


DETACH = b"\x1d"
SHOWN_BACK = 65536


@dataclass(frozen=True)
class ScreenPart:
    data: str
    at: int
    rows: int
    cols: int

    @classmethod
    def blank(cls, rows: int, cols: int) -> "ScreenPart":
        return cls("", 0, rows, cols)


def screen_since(root: Path, terminal: str, since: int) -> ScreenPart:
    screen = runtime.session_file(root, terminal, SCREEN)
    shape = read_json(runtime.session_file(root, terminal, SCREEN_SHAPE), dict, {"rows": 40, "cols": 120})
    if not screen.is_file():
        return ScreenPart.blank(int(shape["rows"]), int(shape["cols"]))
    size = screen.stat().st_size
    at = max(0, size - SHOWN_BACK) if since < 0 or since > size else since
    with screen.open("rb") as shown:
        shown.seek(at)
        fresh = shown.read(size - at)
    return ScreenPart(base64.b64encode(fresh).decode(), at + len(fresh), int(shape["rows"]), int(shape["cols"]))


def attach(root: Path, session: str) -> str:
    screen = runtime.session_file(root, session, SCREEN)
    if not screen.is_file():
        return f"journal: no session {session} to attach to"
    at = max(0, screen.stat().st_size - SHOWN_BACK)
    saved = termios.tcgetattr(sys.stdin.fileno())
    tty.setraw(sys.stdin.fileno())
    try:
        while True:
            with screen.open("rb") as shown:
                shown.seek(at)
                fresh = shown.read()
            at += len(fresh)
            os.write(sys.stdout.fileno(), fresh)
            ready, _, _ = select.select([sys.stdin.fileno()], [], [], 0.2)
            if ready:
                keys = os.read(sys.stdin.fileno(), 4096)
                if not keys or DETACH in keys:
                    break
                typist.send(root, session, keys)
    finally:
        termios.tcsetattr(sys.stdin.fileno(), termios.TCSADRAIN, saved)
    return f"\njournal: left session {session}; it runs on"

import fcntl
import os
import time
from dataclasses import dataclass, field
from pathlib import Path

from engine.fields import Loaded

from resources.types import TYPES
from engine.stored import read_json, write_json, write_text
from engine.proc import run
from engine import runtime, waits


RECENT = 600.0
ACTIVE_ENV = "AGENT_JOURNAL_ACTIVE"
SHELLS = {"sh", "bash", "zsh", "dash", "fish"}


def alive(pid) -> bool:
    try:
        pid = int(pid)
    except (ValueError, TypeError):
        return False
    if pid <= 0:
        return False
    try:
        if os.waitpid(pid, os.WNOHANG)[0]:
            return False
    except ChildProcessError:
        pass
    try:
        os.kill(pid, 0)
    except PermissionError:
        return True
    except OSError:
        return False
    return True


def hold_build(root: Path, build: Path, pid: int | None = None) -> None:
    if build.suffix == ".pyz":
        write_text(runtime.builds(root) / str(pid or os.getpid()), build.name)


def held_builds(root: Path) -> set[str]:
    held = set()
    for marker in runtime.builds(root).glob("*") if runtime.builds(root).is_dir() else ():
        if marker.name.isdigit() and alive(int(marker.name)):
            held.add(marker.read_text().strip())
        else:
            marker.unlink(missing_ok=True)
    return held


@dataclass(frozen=True)
class SessionRecord(Loaded):
    environment: str = ""
    provider: str = ""
    pid: int = 0
    since: float = 0.0
    seen: float = 0.0
    args: tuple = ()
    launch: int = 0
    evicted: dict = field(default_factory=dict)
    grants: tuple = ()

    @property
    def evicted_since_start(self) -> bool:
        return "at" in self.evicted and self.evicted["at"] > self.since

    @property
    def last_heard(self) -> float:
        return self.seen if self.seen else self.since


def live(session: SessionRecord) -> bool:
    if session.pid:
        return alive(session.pid)
    return time.time() - session.last_heard < RECENT


def agent_pid(pid: int) -> int:
    for _ in range(4):
        try:
            parent, name = run(["ps", "-o", "ppid=,comm=", "-p", str(pid)], timeout=2).split(None, 1)
        except ValueError:
            return pid
        if Path(name.strip()).name.lstrip("-") not in SHELLS:
            return pid
        pid = int(parent)
    return pid


class Sessions:
    def __init__(self, root: Path):
        self.root = Path(root)

    def path(self, session: str) -> Path:
        return runtime.session_file(self.root, session, "session.json")

    def read(self, session: str) -> SessionRecord:
        return read_json(self.path(session), SessionRecord.from_json, SessionRecord.from_json({}))

    def known(self, session: str) -> bool:
        return self.path(session).is_file()

    def write(self, session: str, **fields) -> SessionRecord:
        self.path(session).parent.mkdir(parents=True, exist_ok=True)
        with (self.path(session).parent / "session.lock").open("w") as lock:
            with waits.waited("sessions"):
                fcntl.flock(lock, fcntl.LOCK_EX)
            raw = read_json(self.path(session), dict, {})
            got = {**(raw if isinstance(raw, dict) else {}), **fields}
            write_json(self.path(session), got)
        return SessionRecord.from_json(got)

    def bind(self, session: str, env: str, pid: int = 0, provider: str = "") -> dict:
        given = {"pid": pid, "provider": provider}
        return self.write(session, environment=env, since=time.time(), **{key: value for key, value in given.items() if value})

    def choose(self, session: str, provider: str, prefer: str, owned: set[str]) -> str:
        own = self.environment(session)
        if own and self.holder(own) in ("", session):
            return own
        others = self.all()
        ended = sorted((s for name, s in others.items() if name != session and s.provider == provider and s.environment and not live(s)
                        and s.environment not in owned), key=lambda s: s.last_heard)
        for s in reversed(ended):
            if not self.holder(s.environment):
                return s.environment
        return prefer

    def last(self, env: str, provider: str) -> str:
        ended = [(s.last_heard, name) for name, s in self.all().items()
                 if s.environment == env and s.provider == provider and s.pid and not live(s)
                 and not name.startswith(f"{provider}-")]
        return max(ended)[1] if ended else ""

    def touch(self, session: str) -> None:
        self.write(session, seen=time.time())

    def unbind(self, session: str) -> None:
        self.write(session, environment="")

    def rebind(self, old: str, new: str) -> None:
        for session, s in self.all().items():
            if s.environment == old:
                self.write(session, environment=new)
        if runtime.env_file(self.root).is_file() and runtime.env(self.root) == old:
            runtime.set_env(self.root, new)
        runtime.note_rename(self.root, old, new)

    def terminal(self, provider: str, pid: int) -> str:
        return next((name for name, s in self.all().items() if name.startswith(f"{provider}-") and s.pid == pid), "")

    def environment(self, session: str) -> str:
        return self.read(session).environment

    def all(self) -> dict[str, SessionRecord]:
        return {p.parent.name: read_json(p, SessionRecord.from_json, SessionRecord.from_json({})) for p in sorted(runtime.sessions(self.root).glob("*/session.json"))}

    def running(self) -> list[str]:
        return [name for name, s in self.all().items() if s.pid and alive(s.pid)]

    def holder(self, env: str) -> str:
        return next(iter(self.holders(env)), "")

    def free(self, env: str) -> str:
        return next(name for name in (env, *(f"{env}-{i}" for i in range(2, 100))) if not self.holder(name))

    def holders(self, env: str) -> list[str]:
        held = [(s.last_heard, session) for session, s in self.all().items() if s.environment == env and live(s)]
        return [session for _, session in sorted(held, reverse=True)]

    def evict(self, session: str, by: str, env: str, why: str) -> None:
        self.write(session, environment="", evicted={"by": by, "environment": env, "why": why, "at": time.time()})

    def grant(self, session: str, env: str, on: bool = True) -> list[str]:
        lent = set(self.read(session).grants)
        if on:
            lent.add(env)
        else:
            lent.discard(env)
        return list(self.write(session, grants=sorted(lent)).grants)

    def granted(self, session: str, env: str) -> bool:
        return env in self.read(session).grants


class SessionsSnapshot(Sessions):
    def __init__(self, root: Path):
        super().__init__(root)
        self.held: dict[str, SessionRecord] | None = None

    def all(self) -> dict[str, SessionRecord]:
        if self.held is None:
            self.held = super().all()
        return self.held

