import fcntl
import os
import signal
import subprocess
import sys
import threading
import time
from pathlib import Path

from engine import runtime, typist
from engine.record import Record
from engine.seats import Seat, seat_file
from engine.stored import read_json
from engine.sessions import Sessions
from controllers.faults import threw
from controllers.types import Messages, Notices
from resources.base import AGENT, SYSTEM, USER
from runner.engine import TICK, Engine
from engine.package import CODE, ZIPPED, build_file
from engine.locks import claim

ENDING = 5.0
UNHEARD_AFTER = 120.0
LOOKED_EVERY = 30.0
CHILD = "import commands.cli; from runner.engines import child"

SUPERVISING = "engines-supervisor.lock"


def always() -> bool:
    return True


def keep_ticking(stopping, tick, fault, going=always) -> None:
    while not stopping.is_set() and going():
        try:
            tick()
        except Exception:
            fault()
        stopping.wait(TICK)


class Engines:
    def __init__(self, root: Path, env: str):
        self.root = Path(root)
        self.env = env
        self.parent = os.getppid()
        self.held: dict[str, Engine] = {}

    def seated(self, session: str) -> Engine | None:
        from providers import DRIVERS
        provider = Sessions(self.root).read(session).provider
        if provider not in DRIVERS:
            return None
        record = Record(self.root, self.env)
        engine = Engine(record, DRIVERS[provider](record, session))
        engine.start()
        return engine

    def mine(self) -> list[str]:
        sessions = Sessions(self.root)
        return [session for session in typist.live(self.root) if sessions.environment(session) == self.env]

    def tick(self) -> None:
        mine = self.mine()
        self.held = {session: engine for session, engine in self.held.items() if session in mine}
        for session in mine:
            engine = self.held.get(session) or self.seated(session)
            if engine:
                self.held[session] = engine
                engine.step()

    def run(self, stopping) -> None:
        with (runtime.folder(self.root) / f"engines-{self.env}.lock").open("a") as held:
            while not stopping.is_set() and self.going() and not self.owned(held):
                stopping.wait(TICK)
            keep_ticking(stopping, self.tick, lambda: threw(self.root, self.env, f"the engines of {self.env}"), self.going)

    def going(self) -> bool:
        return os.getppid() == self.parent and current(self.root)

    def owned(self, held) -> bool:
        try:
            fcntl.flock(held, fcntl.LOCK_EX | fcntl.LOCK_NB)
            return True
        except OSError:
            return False


def current(root: Path) -> bool:
    return not ZIPPED or build_file(root) == CODE


def leftovers(root: Path) -> list[int]:
    listed = subprocess.run(["ps", "-eo", "pid=,command="], capture_output=True, text=True, timeout=5).stdout
    return [int(line.split(None, 1)[0]) for line in listed.splitlines()
            if CHILD in line and str(root) in line and repr(str(CODE)) not in line]


def child(root: str, env: str) -> None:
    Engines(Path(root), env).run(threading.Event())


class Children:
    def __init__(self, root: Path):
        self.root = Path(root)
        self.running: dict[str, subprocess.Popen] = {}
        self.alarmed: set[str] = set()
        self.looked = 0.0
        for pid in leftovers(self.root):
            try:
                os.kill(pid, signal.SIGTERM)
            except OSError:
                continue

    def wanted(self) -> set[str]:
        sessions = Sessions(self.root)
        return {env for session in typist.live(self.root) if (env := sessions.environment(session) or self.healed(sessions, session))}

    def healed(self, sessions: Sessions, session: str) -> str | None:
        seat = Seat.of(read_json(seat_file(self.root, session), dict, {}), session)
        if not seat.env or sessions.read(session).evicted_since_start:
            return None
        sessions.write(session, environment=seat.env)
        return seat.env

    def tick(self) -> None:
        wanted = self.wanted()
        for env, running in list(self.running.items()):
            if env not in wanted or running.poll() is not None:
                self.end(env)
        for env in sorted(wanted - set(self.running)):
            self.running[env] = self.spawn(env)
        if time.time() - self.looked >= LOOKED_EVERY:
            self.looked = time.time()
            for env in sorted(wanted & set(self.running)):
                self.unheard(env)

    def unheard(self, env: str) -> None:
        record = Record(self.root, env)
        messages = Messages(record, actor=SYSTEM)
        waiting = [row["n"] for row in messages.rows.summaries() if not row["deleted"] and not row["completed"] and AGENT not in row["seen"]
                   and f"{env}:{row['n']}" not in self.alarmed]
        for row in [messages.load(n) for n in waiting]:
            if row.seen[:1] != [USER] or row.data.get("delivered") or time.time() - row.created < UNHEARD_AFTER:
                continue
            self.alarmed.add(f"{env}:{row.n}")
            Notices(record, actor=SYSTEM).create(f"Message {row.n} has not reached the agent", tone="warn",
                                                 brief=f"It waited {int((time.time() - row.created) // 60)} minutes; the agent's engine is restarted to deliver it.")
            self.end(env)
            return

    def spawn(self, env: str) -> subprocess.Popen:
        log = runtime.folder(self.root) / f"engine-{env}.log"
        log.parent.mkdir(parents=True, exist_ok=True)
        with log.open("a") as output:
            return subprocess.Popen([sys.executable, "-c", f"import sys; sys.path.insert(0, {str(CODE)!r}); {CHILD}; child(sys.argv[1], sys.argv[2])",
                                     str(self.root), env], cwd=self.root.parent, stdin=subprocess.DEVNULL, stdout=output, stderr=output)

    def end(self, env: str) -> None:
        running = self.running.pop(env)
        if running.poll() is None:
            running.terminate()
            try:
                running.wait(timeout=ENDING)
            except subprocess.TimeoutExpired:
                running.kill()

    def stop(self) -> None:
        for running in self.running.values():
            if running.poll() is None:
                running.terminate()
        for env in list(self.running):
            self.end(env)

    def run(self, stopping) -> None:
        keep_ticking(stopping, self.tick, lambda: threw(self.root, runtime.env(self.root), "starting the engines"))


def supervise(root: Path, stopping) -> None:
    while not stopping.is_set():
        held = claim(runtime.folder(root) / SUPERVISING)
        if held is None:
            stopping.wait(TICK)
            continue
        with held:
            children = Children(root)
            try:
                children.run(stopping)
            finally:
                children.stop()
import re
from dataclasses import asdict
from pathlib import Path

from controllers.types import Agents, Notices
from engine.inputs import BACKGROUND, FORCE, PAUSE, PERMIT, QueuedCommand, RESUME, SHELL, STALE, queue, waiting_commands
from engine.record import Record
from engine.seats import live_session, offline
from providers.drivers import AGENT_COMMAND
from agents.terminal import relaunch as relaunch_session
from providers import DRIVERS, PROVIDERS
from providers.base import Provider
from resources.base import SYSTEM, Refused
from typing import TypedDict

RELOAD_GRACE = 60.0


def current(group: str, value: str, model: str, effort: str) -> bool:
    if group == "model":
        return value == model or value in re.split(r"[-\[\]]", model)
    return group == "effort" and value == effort


class ControlOptions(TypedDict):
    provider: str
    groups: list[dict]
    note: str


def options(provider: str, current_model: str, current_effort: str) -> ControlOptions:
    controls = PROVIDERS.get(provider, Provider).control_options(current_model)
    return {
        "provider": provider,
        "groups": [
            {**group, "choices": [{**{k: v for k, v in choice.items() if k not in ("command", "commands")},
                                   "current": current(group["key"], choice["value"], current_model, current_effort)} for choice in group["choices"]]}
            for group in controls["groups"]
        ],
        "note": controls["note"],
    }


def choice(provider: str, action: str, value: str, current_model: str) -> dict:
    cls = PROVIDERS.get(provider)
    if not cls:
        raise Refused(f"{provider or 'this agent'} does not support {action} {value!r}")
    return cls.control_choice(action, value, current_model)


def online(root: Path, env: str, session: str) -> dict:
    pair = live_session(Path(root), session, within=RELOAD_GRACE)
    if not pair:
        raise Refused(offline(Path(root), session, within=RELOAD_GRACE))
    found = pair[1]
    if found.environment != env:
        raise Refused(f"session {session!r} belongs to environment {found.environment!r}")
    return found


def pressed(root: Path, env: str, session: str, label: str, action: str, value: str = "") -> dict:
    found = online(root, env, session)
    queued = queue(Path(root), session, (), label, provider=found.provider, action=action, value=value)
    return queued.for_viewer


def force(root: Path, env: str, session: str) -> dict:
    return pressed(root, env, session, "Force through", FORCE)


def pause(root: Path, env: str, session: str) -> dict:
    return pressed(root, env, session, "Pause", PAUSE)


def resume(root: Path, env: str, session: str) -> dict:
    return pressed(root, env, session, "Resume", RESUME)


def move_to_background(root: Path, env: str, session: str) -> dict:
    return pressed(root, env, session, "Move to the background", BACKGROUND)


def shell(root: Path, env: str, session: str, command: str, now: bool = False) -> dict:
    found = online(root, env, session)
    if not command.strip():
        raise Refused("type a command to run")
    for_agent = command.strip().startswith(AGENT_COMMAND)
    if not for_agent and not DRIVERS[found.provider].SHELL:
        raise Refused(f"{found.provider} has no shell command to type")
    queued = queue(Path(root), session, (), f"Run {command.strip()}", provider=found.provider, action=SHELL, value=command.strip())
    if for_agent:
        return queued.for_viewer
    agents = Agents(Record(Path(root), env), actor=SYSTEM)
    row = agents.by_session(session)
    waiting = [c for c in waiting_commands(row) if queued.at - c.at < STALE]
    agents.update(row.n, queued_commands=[*(asdict(c) for c in waiting), asdict(QueuedCommand(queued.at, queued.value))])
    if now:
        queue(Path(root), session, (), "Run now", provider=found.provider, action=FORCE)
    return queued.for_viewer


def permit(root: Path, env: str, session: str, allow: bool) -> dict:
    answer = "allow" if allow else "deny"
    return pressed(root, env, session, answer.capitalize(), PERMIT, answer)


def relaunch(root: Path, env: str, session: str) -> dict:
    found = online(root, env, session)
    relaunch_session(Path(root), env, found.terminal, session)
    return {"relaunching": True}


def request(root: Path, env: str, session: str, action: str, value: str) -> dict:
    root = Path(root)
    found = online(root, env, session)
    selected = choice(found.provider, action, value, found.model)
    queued = queue(root, session, selected.get("commands") or [selected["command"]], selected["label"], provider=found.provider, action=action, value=value)
    record = Record(root, env)
    agents = Agents(record, actor=SYSTEM)
    row = agents.by_session(session)
    agents.update(row.n, pending={**row.pending, action: {"value": value, "at": queued.at}})
    waits = "" if action in PROVIDERS[found.provider].applies_at_once else " — waiting for the agent"
    Notices(record, actor=SYSTEM).create(f"Setting {action} to {selected['label'].lower()}{waits}", tone="note", session=session, action=action)
    return queued.for_viewer
import signal
import sys
import threading
import time
from pathlib import Path

sys.path.insert(0, str(Path(__file__).resolve().parents[1]))
from providers import DRIVERS  # noqa: E402
from engine import viewer  # noqa: E402
from engine.services import Manager  # noqa: E402
from runner.engines import ENDING, supervise  # noqa: E402
from features.plugins.services import plugin_services  # noqa: E402
from engine import runtime  # noqa: E402
from controllers.faults import threw  # noqa: E402
from engine.stop import asked, session_flag  # noqa: E402
from supervisor import HEAL, RELAUNCH, RELOAD, STOP  # noqa: E402
from agents.terminal import TerminalSession, seated  # noqa: E402
from agents.actors import Agent  # noqa: E402
from engine.record import Record  # noqa: E402
from engine.package import CODE, installed_stamp, own_build  # noqa: E402
from engine.sessions import Sessions, hold_build  # noqa: E402
import features  # noqa: E402
from features.switches import watch_change_log  # noqa: E402
from features.auto_update.check import UpdateCheck  # noqa: E402
from features.auto_update.relaunch import Relaunch  # noqa: E402
from features.work_tracking.auto import CheckIn  # noqa: E402

TICK = 0.25
RELOAD_EVERY = 1.0
VIEWER_EVERY = 10.0
SERVICES_EVERY = 1.0
CHECKS_EVERY = 1.0
SERVER_CRASHES = 3
RETRY_AFTER = 1.0
STARTUP, EARLY = 30.0, 16384
CONSENT_EVERY = 3.0
TERMINATED = threading.Event()


def keep_viewer(root: Path, cwd: Path, watching, exits: list) -> object:
    if watching and watching.is_alive():
        return watching
    thread = threading.Thread(target=lambda: exits.append(viewer.launch(root, cwd)[1]), daemon=True)
    thread.start()
    return thread


def moved(seat: TerminalSession) -> bool:
    return Sessions(seat.root).environment(seat.session) not in ("", seat.env)


def crashing(exits: list) -> bool:
    return len(exits) >= SERVER_CRASHES and all(exits[-SERVER_CRASHES:])


def checks(seat: TerminalSession, driver) -> list:
    try:
        features.load()
        watcher = Agent(driver.record, driver)
        return [CheckIn(watcher), UpdateCheck(watcher), Relaunch(watcher)]
    except Exception:
        threw(seat.root, seat.env, "starting the worker's checks")
        return []


def run_checks(seat: TerminalSession, driver, kept: list) -> None:
    for step in [driver.pump, *(check.tick for check in kept)] if kept else []:
        try:
            step()
        except Exception:
            threw(seat.root, seat.env, f"a worker check: {type(getattr(step, '__self__', step)).__name__}")


class Confirm:
    def __init__(self, driver):
        self.driver = driver
        self.at = self.driver.printed.stat().st_size if self.driver.printed.is_file() else 0
        self.started = self.consented = time.time()
        self.ready = 0.0
        self.answered = False

    def tick(self) -> None:
        if self.answered or time.time() - self.started >= STARTUP or not self.driver.printed.is_file():
            return
        fresh = self.driver.printed.stat().st_size - self.at
        early = self.driver.printed_tail(min(fresh, EARLY)) if fresh > 0 else b""
        if self.driver.consent(early):
            self.consent(early)
            return
        opening = self.driver.opening(early)
        self.ready = (self.ready or time.time()) if opening else 0.0
        if opening and time.time() - self.ready >= self.driver.CONFIRM_AFTER:
            self.answered = True
            self.driver.send(opening, now=True)

    def consent(self, early: bytes) -> None:
        if time.time() - self.consented < CONSENT_EVERY:
            return
        self.consented = time.time()
        self.driver.press_raw(self.driver.consent(early))


def run(root: Path, cwd: Path, env: str, agent: str, session: str, lifeline: int = -1) -> int:
    hold_build(root, CODE)
    watch_change_log()
    seat = seated(TerminalSession(root, env, agent, session))
    relaunching = runtime.relaunch_file(root, session)
    stopping = session_flag(root, session)
    stamps = installed_stamp(root)
    began = time.time()
    driver = DRIVERS[agent](Record(root, env), session)
    confirm = Confirm(driver)
    kept = checks(seat, driver)
    services = Manager(root, lifeline, sources=(plugin_services,))
    last_check = last_viewer = last_services = last_checks = 0.0
    watching = None
    exits: list = []
    keeps = own_build(root)
    engines = threading.Event()
    supervising = threading.Thread(target=supervise, args=(root, engines), daemon=True)
    if keeps:
        supervising.start()
    try:
        while True:
            time.sleep(TICK)
            confirm.tick()
            if asked(root, began) or stopping.is_file() or TERMINATED.is_set():
                stopping.unlink(missing_ok=True)
                return STOP
            if relaunching.is_file():
                return RELAUNCH
            now = time.time()
            if keeps and now - last_services >= SERVICES_EVERY:
                last_services = now
                services.tick()
            if keeps and now - last_viewer >= (RETRY_AFTER if exits and exits[-1] else VIEWER_EVERY):
                last_viewer = now
                watching = keep_viewer(root, cwd, watching, exits)
                if crashing(exits):
                    return HEAL
            if now - last_checks >= CHECKS_EVERY:
                last_checks = now
                if moved(seat):
                    return RELOAD
                if not kept:
                    kept = checks(seat, driver)
                run_checks(seat, driver, kept)
            if now - last_check >= RELOAD_EVERY:
                last_check = now
                if installed_stamp(root) != stamps and not runtime.upgrading(root):
                    return RELOAD
    finally:
        engines.set()
        if keeps:
            supervising.join(timeout=ENDING)


def ended(signum, frame) -> None:
    TERMINATED.set()


if __name__ == "__main__":
    signal.signal(signal.SIGTERM, ended)
    root, cwd, env, agent, session = sys.argv[1:6]
    raise SystemExit(run(Path(root), Path(cwd), env, agent, session, int(sys.argv[6]) if len(sys.argv) > 6 else -1))
