
def test_a_chosen_setting_reaches_the_plugins_commands():
    from features.plugins.commands import Configure
    from features.plugins.declared import settings_of
    from features.plugins.environment import environment
    record = alone()
    row = installed(record, "linter", "exit 0", settings={"quiet": {"title": "Quiet", "default": "", "env": "QUIET"}})
    plugins = Plugins(record, actor=SYSTEM)
    assert environment(record.root, "linter", Manifest.of(row.manifest), row.token)["QUIET"] == "", "unchanged, a setting is its default"
    from controllers.types import Environments, Todos
    from engine.record import Record
    from features.plugins.queue import drain
    Environments(record, actor=SYSTEM).create("other")
    queued = environment(record.root, "linter", Manifest.of(row.manifest), row.token, env="other")["JOURNAL_QUEUE"]
    Path(queued).parent.mkdir(parents=True, exist_ok=True)
    Path(queued).write_text('todo create "From the event"\ntodo create "Somewhere else" --env ' + record.env + "\n")
    done, refusals = drain(record.root, "linter", record.env)
    assert (done, [(r.env, "names no --env" in r.why) for r in refusals]) == (2, [("other", True)]), "the drain hands back what it refused, and why"
    from features import FEATURES
    from features.plugins.host import Host
    from controllers.types import Notices
    host = Host(record.root, FEATURES["plugins"].journal)
    for _ in range(2):
        host.refusal("linter", refusals[0])
    told = [n.title for n in Notices(Record(record.root, "other"), actor=SYSTEM).all()]
    assert told == ["Plugin linter queued a line the journal refused"], "a refusal is told once, as a notice, in the environment of the line"
    titles = lambda env: [t.title for t in Todos(Record(record.root, env), actor=SYSTEM).all()]
    assert (titles("other"), "Somewhere else" in titles(record.env)) == (["From the event"], False), \
        "what a plugin queues answering an event runs in that event's environment, and a queued --env is refused"
    Configure().run(None, plugins, row.n, "quiet", "SourceReminder")
    chosen = settings_of(plugins.load(row.n)).chosen
    assert environment(record.root, "linter", Manifest.of(row.manifest), row.token, chosen=chosen)["QUIET"] == "SourceReminder", "and a chosen value reaches its env"
    assert "has no setting" in refused(lambda: Configure().run(None, plugins, row.n, "loud", "x"))
    from features.plugins.manifest import typed
    assert typed({"php": {"type": "flag"}, "strict": {"parent": "php"}})["strict"]["parent"] == "php", "a setting sits under the switch that turns it on"
    assert "no flag setting" in refused(lambda: typed({"php": {"type": "text"}, "strict": {"parent": "php"}})), "only under a switch"
    typed = installed(record, "typed", "exit 0", settings={"on": {"type": "flag", "default": "true"}, "level": {"type": "options", "options": ["low", "high"]}},
                      events={"sin-found": {"title": "Sin found", "tone": "warn", "card": {"icon": "warn"}}})
    assert "true or false" in refused(lambda: Configure().run(None, plugins, typed.n, "on", "yes")), "a switch takes true or false"
    assert "one of low, high" in refused(lambda: Configure().run(None, plugins, typed.n, "level", "mid")), "options take one of theirs"
    from features.plugins.answer import apply
    apply(record, None, "typed", "", {"settings": {"level": "high", "made-up": "x"}})
    assert settings_of(plugins.load(typed.n)).chosen == {"level": "high"}, "a plugin may fill in a setting it worked out, and only its own"
    from features import FEATURES
    from engine import bus
    heard = []
    Agents(record, actor=AGENT).create("s-1")
    off = bus.on("typed.sin-found", lambda event, record: heard.append(event.data["brief"]))
    apply(record, FEATURES["plugins"].journal, "typed", "", {"raise": {"event": "sin-found", "brief": "deep-nesting at src/A.php:12"}})
    raised = [e for e in record.event_log.events() if e.action == "raised"][-1]
    assert (raised.data["title"], raised.data["tone"], raised.data["brief"], heard) == ("Sin found", "warn", "deep-nesting at src/A.php:12", ["deep-nesting at src/A.php:12"]), \
        "a plugin raises an event it declared, styled from its manifest, and anything listening by its name hears it"
    card = Agents(record, actor=SYSTEM).primary().data["cards"][-1]
    assert (card["label"], card["tone"], card["icon"], card["detail"]) == ("Sin found", "warn", "warn", "deep-nesting at src/A.php:12"), \
        f"an event whose declaration carries a card puts it in the chat, looking as the manifest says: {card}"
    from features.format import VIEWER, shaped
    viewed = shaped(Agents(record, actor=SYSTEM).primary(), record, VIEWER)["data"]["cards"][-1]["detail"]
    assert "[[file src/A.php" in viewed, f"its words pass the formatters like any brief, so a file is a chip: {viewed}"
    from features.plugins.commands import Raise
    Raise().run(None, plugins, "typed", "sin-found", "again at src/B.php:3", open="sins/sin/deep-nesting/src/B.php")
    assert heard[-1] == "again at src/B.php:3", "journal plugin raise, from the queue, raises the same declared event"
    assert Agents(record, actor=SYSTEM).primary().data["cards"][-1]["page"] == "sins/sin/deep-nesting/src/B.php", \
        "and a raise that names a dashboard page gives its card that page to open"
    assert "declares no event" in refused(lambda: Raise().run(None, plugins, "typed", "made-up", "")), "and refuses one it does not declare"
    off()
    assert apply(record, FEATURES["plugins"].journal, "typed", "", {"raise": {"event": "made-up"}}) == [], "an event the manifest does not declare is refused"
    from features.plugins.manifest import typed as checked
    shown = checked({"php": {"type": "flag"}, "vue": {"type": "flag"}, "sin": {"type": "flag", "when": {"php": True}},
                     "either": {"type": "flag", "when": [{"php": True}, {"vue": True}]}})
    assert [shown["sin"]["when"], shown["either"]["when"]] == [[{"php": True}], [{"php": True}, {"vue": True}]], \
        "a setting may be shown only while another has a value, or while any of several do"
    assert "names settings" in refused(lambda: checked({"sin": {"type": "flag", "when": {"ruby": True}}})), "a condition names a setting that exists"


def test_a_service_no_plugin_declares_is_stopped_and_forgotten():
    record = fresh()
    left = subprocess.Popen(["sleep", "30"], start_new_session=True)
    status_file(record.root, "gone.web").parent.mkdir(parents=True, exist_ok=True)
    status_file(record.root, "gone.web").write_text(json.dumps({"state": "running", "keeper": left.pid, "pgid": left.pid}))
    keeping = Manager(record.root, sources=(plugin_services,))
    keeping.tick()
    assert left.wait(timeout=5) is not None, "its process is stopped"
    assert not status_file(record.root, "gone.web").exists(), "and it is no longer listed"
    other = subprocess.Popen(["sleep", "30"], start_new_session=True)
    status_file(record.root, "gone.web").write_text(json.dumps({"state": "running", "keeper": other.pid, "pgid": other.pid}))
    Manager(record.root, sources=(plugin_services,)).tick()
    assert other.poll() is None, "a second keeper of the same project keeps nothing while the first holds the services"
    keeping.tick()
    assert other.wait(timeout=5) is not None, "the one that holds them does"


def test_stopping_a_service_stops_every_process_it_forked():
    from engine.keeper import gone, teardown
    service = subprocess.Popen(["/bin/sh", "-c", "sleep 30 & sleep 30 & wait"], start_new_session=True)
    time.sleep(0.2)
    teardown(service.pid, 1.0)
    assert (service.wait(timeout=5) is not None, gone(service.pid)) == (True, True), \
        "the service and the workers it forked go together, as one process group"
    from engine.services import lock_file, log_file, want
    record = fresh()
    restarted = subprocess.Popen(["/bin/sh", "-c", "sleep 30 & wait"], start_new_session=True)
    status_file(record.root, "fresh.web").parent.mkdir(parents=True, exist_ok=True)
    status_file(record.root, "fresh.web").write_text(json.dumps({"state": "ready", "keeper": restarted.pid, "pgid": restarted.pid}))
    for place in (lock_file, log_file):
        place(record.root, "fresh.web").write_text("old")
    want(record.root, "fresh.web", "up", nonce=time.time())
    Manager(record.root).remove("fresh.web")
    assert restarted.wait(timeout=5) is not None, "a restart first stops the old run"
    assert [place(record.root, "fresh.web").exists() for place in (status_file, lock_file, log_file)] == [False, False, False], \
        "and removes what it left, so the new run starts from nothing of the old one's"
    stray = subprocess.Popen(["sleep", "30"], start_new_session=True)
    lock_file(record.root, "stray.web").write_text(str(stray.pid))
    Manager(record.root).remove("stray.web")
    assert stray.wait(timeout=5) is not None, "a keeper the status no longer names but that still holds the lock is stopped before its lock is removed"
    from engine.keeper import ServiceSpec
    from engine.services import files_for
    kept = subprocess.Popen(["/bin/sh", "-c", "sleep 30 & wait"], start_new_session=True)
    asked = time.time()
    want(record.root, "kept.web", "up", nonce=asked)
    status_file(record.root, "kept.web").write_text(json.dumps({"state": "ready", "keeper": kept.pid, "pgid": kept.pid, "nonce": asked}))
    Manager(record.root).one(ServiceSpec(id="kept.web", plugin="kept", service="web", run=["true"], **files_for(record.root, "kept.web")))
    assert kept.poll() is None, "a restart already carried out is never carried out again by the next agent's manager"
    kept.kill()
    started = []
    manager = Manager(record.root, start=lambda spec, lifeline: started.append(spec.id) or 0)
    for sid, when in (("idle.web", "echo no C# here; exit 1"), ("busy.web", "exit 0")):
        manager.one(ServiceSpec(id=sid, plugin=sid.split(".")[0], service="web", run=["true"], when=when, **files_for(record.root, sid)))
    idle = json.loads(status_file(record.root, "idle.web").read_text())
    assert (started, idle["state"], "no C# here" in idle["why"]) == (["busy.web"], "not needed", True), \
        "a service whose when-command fails is left unstarted as not needed, with the command's own words; one that answers 0 starts"
    want(record.root, "busy.web", UP, nonce=time.time())
    manager.one(ServiceSpec(id="busy.web", plugin="busy", service="web", run=["true"], when="echo none here; exit 1", **files_for(record.root, "busy.web")))
    assert json.loads(status_file(record.root, "busy.web").read_text())["state"] == "not needed", \
        "a restart, as after a plugin upgrade, asks the when-command again instead of keeping the old answer"
    stuck = subprocess.Popen(["/bin/sh", "-c", "sleep 30 & wait"], start_new_session=True)
    status_file(record.root, "stuck.web").write_text(json.dumps({"state": "starting", "keeper": stuck.pid, "pgid": stuck.pid}))
    manager.one(ServiceSpec(id="stuck.web", plugin="stuck", service="web", run=["true"], when="exit 1", **files_for(record.root, "stuck.web")))
    assert (stuck.wait(5) is not None, json.loads(status_file(record.root, "stuck.web").read_text())["state"]) == (True, "not needed"), \
        "a service already running is stopped once its when-command says it is not needed"
    with socket.socket() as probe:
        probe.bind(("127.0.0.1", 0))
        port = probe.getsockname()[1]
    Manager(record.root).one(ServiceSpec(id="real.web", plugin="real", service="web", run=[sys.executable, "-m", "http.server", str(port), "--bind", "127.0.0.1"],
                                         port=port, url=f"http://127.0.0.1:{port}", **files_for(record.root, "real.web")))
    deadline = time.time() + 15
    while time.time() < deadline and json.loads(status_file(record.root, "real.web").read_text()).get("state") != "ready":
        time.sleep(0.2)
    assert json.loads(status_file(record.root, "real.web").read_text()).get("state") == "ready", \
        "the keeper it ships reads its spec from disk and brings a real service up, as it does for the phone's server and tunnel"
    Manager(record.root).remove("real.web")
    from engine.services import allocate
    with socket.socket() as busy:
        busy.bind(("127.0.0.1", 0))
        busy.listen()
        port = busy.getsockname()[1]
        status_file(record.root, "own.web").write_text(json.dumps({"state": "stopped", "port": port}))
        assert allocate(record.root, "own.web", None, set())[0] != port, \
            "a held port cannot be assigned to a service without proving which process owns it"
        status_file(record.root, "own.web").write_text(json.dumps({"state": "exited", "port": port}))
        assert allocate(record.root, "own.web", None, set())[0] != port, "one whose run broke with its port taken gets another"
    holding = subprocess.Popen(["sleep", "30"], start_new_session=True)
    lock_file(record.root, "held.web").write_text(str(holding.pid))
    status_file(record.root, "held.web").write_text(json.dumps({"state": "starting", "keeper": 999999}))
    spawned = []
    Manager(record.root, start=lambda spec, lifeline: spawned.append(spec.id) or 0).one(
        ServiceSpec(id="held.web", plugin="held", service="web", run=["true"], **files_for(record.root, "held.web")))
    assert spawned == [] and json.loads(status_file(record.root, "held.web").read_text())["keeper"] == holding.pid, \
        "a live keeper that holds the lock is adopted, never started again beside itself"
    holding.kill()
    running = subprocess.Popen(["sleep", "30"], start_new_session=True)
    rebuilt = ServiceSpec(id="moved.web", plugin="moved", service="web", run=["true"], env={"JOURNAL_BUILD": "port 8442"}, **files_for(record.root, "moved.web"))
import time
from dataclasses import replace
from pathlib import Path

from controllers.faults import threw
from controllers.types import Notices, Plugins
from engine.runtime import default_env
from engine.keeper import ServiceSpec
from engine.record import Record
from engine.services import BLOCKED, FAILED, claimed, local_url, log_file, service_spec, states
from engine.wording import fill
from features.plugins.declared import declared, settings_of
from features.plugins.environment import environment, port_values
from features.plugins.paths import folder
from resources.base import SYSTEM

WATCH = 5.0
TOLD = "service"


def keep(root: Path, feature) -> None:
    while True:
        try:
            notice_stopped(root, feature)
        except Exception:
            threw(root, default_env(root), "watching plugin services")
        time.sleep(WATCH)


def here(root: Path) -> Record:
    return Record(Path(root), default_env(Path(root)))


def notice_stopped(root: Path, feature) -> list[str]:
    record = here(root)
    if not feature.enabled(record):
        return []
    open_ = {n.data.get(TOLD): n for n in feature.standing(record, Notices) if n.data.get(TOLD)}
    stopped = []
    for sid, state in states(root).items():
        failing = state.state in (FAILED, BLOCKED)
        if failing and sid not in open_:
            feature.journal.notice(record, "stopped", name=sid, why=state.why if state.why else "it stopped", log=log_file(root, sid), tone="warn", **{TOLD: sid})
            stopped.append(sid)
        if not failing and sid in open_:
            feature.journal.clear(record, open_[sid], "it is running again")
    return stopped


def plugins(root: Path) -> list:
    return Plugins(here(root), actor=SYSTEM)._installed()


def planned(root: Path, name: str, service, port: int, blocked: str, env: dict, where: Path) -> ServiceSpec:
    return service_spec(root, f"{name}.{service.name}", plugin=name, service=service.name, port=port, blocked=blocked, run=service.run,
                        cwd=str(where / service.cwd), env={**env, **service.env}, path=service.ready.path, restart=service.restart,
                        grace=service.grace, show=service.show, when=service.when)


def plugin_services(root: Path, taken: set[int]) -> list[ServiceSpec]:
    out: list[ServiceSpec] = []
    for row in plugins(root):
        manifest = declared(row)
        name = manifest.name
        where = folder(root, name)
        settings = settings_of(row)
        kept = settings.ports
        env = environment(root, name, manifest, row.token, kept, settings.chosen)
        ports = port_values(kept)
        made = []
        for service in manifest.services:
            asked = kept[service.name] if kept.get(service.name) else service.port
            port, blocked = claimed(root, f"{name}.{service.name}", asked, taken) if service.port is not None else (0, "")
            if port:
                ports[f"ports.{service.name}"] = port
            made.append(planned(root, name, service, port, blocked, env, where))
        for spec in made:
            places = {**ports, "port": spec.port, "dir": str(where)}
            out.append(replace(spec, run=fill(spec.run, places), when=str(fill(spec.when, places)), env={key: str(fill(value, places)) for key, value in spec.env.items()},
                               url=local_url(spec.port)))
    return out
import fcntl
import os
import select
import signal
import socket
import subprocess
import sys
import time
import urllib.error
import urllib.request
from dataclasses import asdict, dataclass, field, fields, replace
from pathlib import Path

sys.path.insert(0, str(Path(__file__).resolve().parents[1]))
from engine.stored import read_json, write_json  # noqa: E402

TAKEN = 3
WATCH = 0.5
GRACE = 5.0
KILL_AFTER = 2.0
PROBE = 1.0
STARTING, READY, STOPPED, EXITED = "starting", "ready", "stopped", "exited"


def known(cls, raw: dict) -> dict:
    return {f.name: raw[f.name] for f in fields(cls) if f.name in raw}


BUILD = "JOURNAL_BUILD"


@dataclass(frozen=True)
class ServiceSpec:
    id: str
    plugin: str
    service: str
    run: object
    lock: str
    log: str
    status: str
    spec: str
    port: int = 0
    blocked: str = ""
    cwd: str = ""
    env: dict = field(default_factory=dict)
    path: str = ""
    restart: str = "on-failure"
    grace: float = GRACE
    show: dict = field(default_factory=dict)
    url: str = ""
    owner: int = 0
    when: str = ""

    @classmethod
    def from_json(cls, raw: dict) -> "ServiceSpec":
        return cls(**known(cls, raw))

    @property
    def build(self) -> str:
        return self.env.get(BUILD, "")


@dataclass(frozen=True)
class ServiceState:
    state: str = ""
    at: float = 0.0
    keeper: int = 0
    pgid: int = 0
    owner: int = 0
    port: int = 0
    url: str = ""
    why: str = ""
    started: float = 0.0
    ready_at: float = 0.0
    last_exit: int = 0
    nonce: float = 0.0
    build: str = ""

    @classmethod
    def read(cls, path: Path) -> "ServiceState":
        raw = read_json(Path(path), dict, {})
        return cls(**known(cls, raw if isinstance(raw, dict) else {}))

    def write(self, path: Path) -> None:
        write_json(Path(path), asdict(self))


def answers(port: int, path: str) -> bool:
    if not port:
        return True
    if not path:
        try:
            with socket.create_connection(("127.0.0.1", port), timeout=PROBE):
                return True
        except OSError:
            return False
    try:
        with urllib.request.urlopen(f"http://127.0.0.1:{port}{path}", timeout=PROBE) as answered:
            return answered.status < 500
    except urllib.error.HTTPError as answered:
        return answered.code < 500
    except (urllib.error.URLError, OSError, ValueError):
        return False


def gone(group: int) -> bool:
    try:
        os.killpg(group, 0)
    except ProcessLookupError:
        return True
    except PermissionError:
        return False
    return False


def teardown(group: int, grace: float) -> None:
    for sign, waited in ((signal.SIGTERM, grace), (signal.SIGKILL, KILL_AFTER)):
        try:
            os.killpg(group, sign)
        except (ProcessLookupError, PermissionError):
            return
        until = time.monotonic() + waited
        while time.monotonic() < until:
            if gone(group):
                return
            time.sleep(0.05)


def state(spec: ServiceSpec, name: str, child=None, **more) -> None:
    replace(ServiceState.read(spec.status), state=name, at=time.time(), keeper=os.getpid(), pgid=child.pid if child else 0, owner=spec.owner,
            port=spec.port, url=spec.url, build=spec.build, **more).write(spec.status)


def lease(spec: ServiceSpec):
    held = open(spec.lock, "a")
    try:
        fcntl.flock(held, fcntl.LOCK_EX | fcntl.LOCK_NB)
    except BlockingIOError:
        held.close()
        return None
    os.set_inheritable(held.fileno(), False)
    os.ftruncate(held.fileno(), 0)
    held.write(str(os.getpid()))
    held.flush()
    return held


def watch(spec: ServiceSpec, child, lifeline: int, stopping) -> str:
    ready = False
    while child.poll() is None and not stopping[0]:
        seen, _, _ = select.select([lifeline], [], [], WATCH) if lifeline >= 0 else ([], [], [])
        if seen and not os.read(lifeline, 1):
            return STOPPED
        if not ready and answers(spec.port, spec.path):
            ready = True
            state(spec, READY, child, ready_at=time.time())
    return STOPPED if stopping[0] else EXITED


def main(argv: list[str]) -> int:
    lifeline, spec_path = int(argv[0]), Path(argv[1])
    spec = read_json(spec_path, ServiceSpec.from_json, None)
    if spec is None:
        sys.exit(f"no service spec at {spec_path}")
    held = lease(spec)
    if not held:
        return TAKEN
    stopping = [False]
    signal.signal(signal.SIGTERM, lambda *_: stopping.__setitem__(0, True))
    signal.signal(signal.SIGHUP, signal.SIG_IGN)
    Path(spec.log).parent.mkdir(parents=True, exist_ok=True)
    with open(spec.log, "ab", buffering=0) as log:
        command = spec.run
        child = subprocess.Popen(["/bin/sh", "-c", command] if isinstance(command, str) else list(command),
                                 cwd=spec.cwd if spec.cwd else None, env={**os.environ, **spec.env},
                                 stdin=subprocess.DEVNULL, stdout=log, stderr=subprocess.STDOUT, start_new_session=True)
        state(spec, STARTING, child, started=time.time())
        ended = watch(spec, child, lifeline, stopping)
        teardown(child.pid, spec.grace)
        state(spec, ended, child, last_exit=child.wait())
    held.close()
    return 0


if __name__ == "__main__":
    sys.exit(main(sys.argv[1:]))
