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 json
import os
import subprocess
import sys
import threading
import time
from dataclasses import dataclass
from pathlib import Path

from engine import runtime
from engine.sessions import ACTIVE_ENV, agent_pid
from engine.fields import Loaded
from engine.package import build_file

PROTOCOL = "2025-06-18"
NAME = "journal"
WAIT = 0.3
READ_AT = "JOURNAL_CHANNEL_READ_AT"
CHECKING = "JOURNAL_CHANNEL_CHECK"
CHECK_FOR = 20

queue, alive = runtime.channel_queue, runtime.channel_alive


def say(message: dict) -> None:
    sys.stdout.write(json.dumps(message) + "\n")
    sys.stdout.flush()


def start(f: Path) -> int:
    carried = os.environ.pop(READ_AT, "")
    if carried:
        return int(carried)
    return f.stat().st_size if f.exists() else 0


def fresh_lines(f: Path, at: int) -> tuple[list[str], int]:
    try:
        size = f.stat().st_size
    except OSError:
        return [], at
    if size < at:
        return [], size
    if size == at:
        return [], at
    with f.open() as lines:
        lines.seek(at)
        read = lines.read()
        return read.splitlines(), lines.tell()


@dataclass(frozen=True)
class Queued(Loaded):
    content: str = ""


@dataclass(frozen=True)
class Asked(Loaded):
    method: str = ""
    id: object = None


def contents(lines: list[str]) -> list[str]:
    texts = []
    for line in lines:
        try:
            texts.append(Queued.from_json(json.loads(line)).content)
        except ValueError:
            continue
    return [s for s in texts if s]


def renewed(root: Path, began: Path) -> bool:
    return build_file(root) != began and build_file(root).is_file()


def starts() -> bool:
    try:
        return subprocess.run(sys.orig_argv, env={**os.environ, CHECKING: "1"}, stdin=subprocess.DEVNULL, capture_output=True,
                              timeout=CHECK_FOR).returncode == 0
    except (OSError, subprocess.SubprocessError):
        return False


def push(root: Path, pid: int) -> None:
    f = queue(root, pid)
    at = start(f)
    began = build_file(root)
    while True:
        time.sleep(WAIT)
        if renewed(root, began):
            if not starts():
                began = build_file(root)
                continue
            sys.stdout.flush()
            os.environ[READ_AT] = str(at)
            os.execv(sys.executable, sys.orig_argv)
        try:
            alive(root, pid).touch()
            lines, at = fresh_lines(f, at)
            texts = contents(lines)
            if texts:
                say({"jsonrpc": "2.0", "method": "notifications/claude/channel", "params": {"content": "; ".join(texts), "meta": {"from": "journal"}}})
        except Exception:
            from controllers.faults import threw
            threw(root, os.environ.get("JOURNAL_ENV") or runtime.env(root), "the channel that carries lines to the agent")


def launched() -> bool:
    return os.environ.get(ACTIVE_ENV) == "1"


def answer(asked: Asked) -> dict | None:
    method = asked.method
    if method == "initialize" and not launched():
        return {"protocolVersion": PROTOCOL, "serverInfo": {"name": NAME, "version": "1"}, "capabilities": {}}
    if method == "initialize":
        return {"protocolVersion": PROTOCOL, "serverInfo": {"name": NAME, "version": "1"},
                "capabilities": {"experimental": {"claude/channel": {}}},
                "instructions": "The journal pushes its own notices here: unread messages, answered questions, work it wants you to close. Each is an instruction to you, not a message from the user: act on it or note it, and never answer or mention it in the chat."}
    if method in ("tools/list", "prompts/list", "resources/list"):
        return {method.split("/")[0]: []}
    return {} if method and not method.startswith("notifications/") else None


def main(argv: list[str]) -> int:
    if os.environ.get(CHECKING):
        return 0
    root = Path(argv[0]) if argv else Path.cwd() / ".journal"
    pid = agent_pid(os.getppid())
    queue(root, pid).parent.mkdir(parents=True, exist_ok=True)
    if launched():
        threading.Thread(target=push, args=(root, pid), daemon=True).start()
    for line in sys.stdin:
        try:
            asked = Asked.from_json(json.loads(line))
        except ValueError:
            continue
        reply = answer(asked)
        if reply is not None and asked.id is not None:
            say({"jsonrpc": "2.0", "id": asked.id, "result": reply})
    return 0
import re
import time
from pathlib import Path

from providers.payload import Asking
from providers.drivers import Driver, plain, squeezed


class CodexDriver(Driver):
    PRODUCT = "the Codex CLI"
    AUTO_ARGS = ("--approve-for-me",)
    SKIP_ARGS = ("--dangerously-bypass-approvals-and-sandbox",)
    APPROVAL_FLAGS = frozenset({"-a", "--ask-for-approval", "--approve-for-me", "--full-auto", "--dangerously-bypass-approvals-and-sandbox"})
    TRUSTS_HOOKS = "--dangerously-bypass-hook-trust"
    READY = b"AskCodextodoanything"
    BUSY = b"esctointerrupt"
    ASKING = (b"Wouldyouliketorun", b"Yes,proceed", b"Allowcommand", b"Approve")
    ASKED_COMMAND = re.compile(r"Would you like to run the following command\?.*\$ (.+?)\s*›\s*1\.\s*Yes, proceed", re.S)
    ALLOW = b"y"
    ASKS_ON_SCREEN = True
    QUEUED = b"Messagestobesubmittedafternexttoolcall"
    RUNNING = b"backgroundterminalrunning"
    SEND_NOW = b"\x1b"
    TRUSTING = re.compile(rb"(?:Doyoutrustthecontentsofthisdirectory|Trustthisfolder\?).*?(\d)\.(?:Yes,continue|Trustandcontinue)", re.S)
    UPDATING = re.compile(rb"Updateavailable.*?(\d)\.Skip(?!until)", re.S)
    PROMPT_TAIL = 8192
    OPENING = "The journal started this session."
    CONFIRM_AFTER = 3.0
    RESUME = "resume"
    CONTINUING = ("continue", "--continue")
    name = "codex"

    @classmethod
    def opening(cls, printed: bytes) -> str:
        return cls.OPENING if cls.READY in squeezed(printed) and not cls.consent(printed) else ""

    @classmethod
    def consent(cls, printed: bytes) -> bytes:
        plain = squeezed(printed)
        asked = max((*cls.TRUSTING.finditer(plain), *cls.UPDATING.finditer(plain)), key=lambda match: match.start(), default=None)
        return asked.group(1) + b"\r" if asked and asked.start() > plain.rfind(cls.READY) else b""

    def at_prompt(self) -> bool:
        plain = self._screen()
        return plain.rfind(self.READY) > plain.rfind(self.BUSY) and self.quiet_for() >= self.QUIET

    def asking(self) -> bool:
        plain = self._screen()
        return max(plain.rfind(phrase) for phrase in self.ASKING) > max(plain.rfind(self.READY), plain.rfind(self.BUSY))

    def asked(self) -> Asking | None:
        if not self.asking():
            return None
        found = list(self.ASKED_COMMAND.finditer(self._screen_text()))
        return Asking("exec_command", found[-1][1].strip()[:300] if found else "a command", time.time())

    def _screen_text(self) -> str:
        return self._printed_tail().decode(errors="replace")

    def _screen(self) -> bytes:
        return b"".join(self._printed_tail().split())

    def _printed_tail(self) -> bytes:
        return plain(self.printed_tail(self.PROMPT_TAIL))

    @classmethod
    def command(cls, args: list[str], cwd: Path | None = None) -> list[str]:
        trusted = ["-c", f'projects."{Path(cwd).resolve()}".trust_level="trusted"'] if cwd else []
        trust = [] if cls.TRUSTS_HOOKS in args else [cls.TRUSTS_HOOKS]
        return ["codex", *trust, *trusted, *cls.carried_on(args)]

    @classmethod
    def carried_on(cls, args: list[str]) -> list[str]:
        rest = [arg for arg in args if arg not in cls.CONTINUING]
        if len(rest) < len(args):
            return [cls.RESUME, "--last", *rest]
        if "--resume" in args[:-1]:
            at = args.index("--resume")
            return [cls.RESUME, args[at + 1], *args[:at], *args[at + 2:]]
        return args

    @classmethod
    def resuming(cls, args: list[str]) -> bool:
        return cls.carried_on(args)[:1] == [cls.RESUME]

    @classmethod
    def resumed(cls, args: list[str], conversation: str) -> list[str]:
        if not conversation:
            return args
        rest = cls.carried_on(args)
        if rest[:1] == [cls.RESUME]:
            rest = rest[2:] if len(rest) > 1 and (rest[1] == "--last" or not rest[1].startswith("-")) else rest[1:]
        return [cls.RESUME, conversation, *rest]
