import json
import os
import platform
import secrets
import shutil
import time
import urllib.request
from pathlib import Path

from engine.proc import ran as ran_command
from engine.stored import write_json
from typing import TypedDict
from engine.given import given
from engine.state import State
from engine.wording import slugged
from resources.base import Refused

TUNNEL_FILE = "sharing.json"
ADDRESS_REFUSED = "address_refused"
OWNED = "domain is owned by another user"
HELD = "409 Conflict"
LAST_LINES = 4
NAME_BYTES = 12
LOCAL_BIN = Path.home() / ".local" / "bin" / "tunler"
SERVER, TUNNEL = "sharing.server", "sharing.tunnel"
ARCHES = {"x86_64": "amd64", "aarch64": "arm64"}
DOWNLOAD_SECONDS = 60


def kept_address(root: Path) -> dict:
    path = Path(root) / TUNNEL_FILE
    if not path.exists():
        return {}
    unreadable = f"cannot read the tunnel address in {path}; the address has not changed"
    try:
        kept = json.loads(path.read_text())
    except (OSError, ValueError) as error:
        raise Refused(unreadable) from error
    if not isinstance(kept, dict):
        raise Refused(unreadable)
    return kept


def alerts(root: Path) -> State:
    return State(Path(root) / "runtime" / "sharing-tunnel.json")


def subdomain(root: Path) -> str:
    kept = kept_address(root)
    return kept.get("subdomain") or addressed(root, kept)


def new_address(root: Path) -> str:
    try:
        kept = kept_address(root)
    except Refused:
        kept = {}
    return addressed(root, {key: value for key, value in kept.items() if key != "subdomain"})


def addressed(root: Path, kept: dict) -> str:
    prefix = slugged(Path(root).resolve().parent.name, limit=20) or "journal"
    name = f"{prefix}-{secrets.token_hex(NAME_BYTES)}"
    write_json(Path(root) / TUNNEL_FILE, {**kept, "subdomain": name})
    return name


def last_lines(log: Path) -> list[str]:
    try:
        return log.read_text(errors="ignore").splitlines()[-LAST_LINES:]
    except OSError:
        return []


def refused_address(log: Path) -> bool:
    return any(OWNED in line for line in last_lines(log))


def held_by_server(log: Path) -> bool:
    lines = last_lines(log)
    return bool(lines) and HELD in lines[-1]


def tunler() -> str:
    found = shutil.which("tunler")
    if found:
        return found
    return str(LOCAL_BIN) if LOCAL_BIN.is_file() else ""

STATUS_SECONDS = 5
STATUS_KEPT = 60
KEPT_STATUS: dict = {}


def tunler_status() -> dict:
    if time.time() - KEPT_STATUS.get("at", 0) < STATUS_KEPT:
        return KEPT_STATUS["status"]
    KEPT_STATUS.update(at=time.time(), status=asked_status())
    return KEPT_STATUS["status"]


class TunnelStatus(TypedDict):
    installed: bool
    logged_in: bool
    account: str
    host: str


class Login(TypedDict):
    connected: bool
    needs_master: bool
    error: str


def asked_status() -> TunnelStatus:
    command = tunler()
    if not command:
        return TunnelStatus(installed=False, logged_in=False, account="", host="")
    done = ran_command([command, "status", "--json"], timeout=STATUS_SECONDS)
    try:
        told = json.loads(done.stdout if done else "{}") or {}
    except ValueError:
        told = {}
    return TunnelStatus(installed=True, logged_in=bool(told.get("logged_in") and told.get("auth_ok")), account=told.get("user") or told.get("email", ""),
                        host=told.get("host", ""))


LOGIN_SECONDS = 30


def ran(*words: str, hidden: dict | None = None) -> tuple[bool, str]:
    command = tunler()
    if not command:
        return False, "tunler isn't installed on this machine"
    done = ran_command([command, *words], timeout=LOGIN_SECONDS, stdin="", env={**os.environ, **(hidden or {})})
    if done is None:
        return False, f"tunler {words[0]} did not finish"
    KEPT_STATUS.clear()
    return done.returncode == 0, (done.stdout if done.returncode == 0 else done.stderr or done.stdout).strip()


def log_in(host: str, username: str, password: str, master: str | None = None) -> Login:
    ok, said = ran("login", username, f"--host={host}", hidden=given(TUNLER_PASSWORD=password, TUNLER_MASTER_PASSWORD=master))
    if ok:
        return Login(connected=True, needs_master=False, error="")
    lines = said.splitlines() or ["tunler refused the login"]
    needs = "master password" in said.lower() and master is None
    return Login(connected=False, needs_master=needs, error=lines[0] if needs else lines[-1])


class TunlerVersion(TypedDict):
    current: str
    latest: str
    update_available: bool


def installed() -> str:
    ok, said = ran("version")
    return said.split()[-1] if ok and said else ""


def latest(host: str) -> str:
    try:
        with urllib.request.urlopen(f"https://{host}/_tunler/version", timeout=STATUS_SECONDS) as answer:
            return str(json.loads(answer.read()).get("version", ""))
    except (OSError, ValueError):
        return ""


def versions(host: str) -> TunlerVersion:
    ok, said = ran("update", "--check", "--json", f"--host={host}")
    try:
        told = json.loads(said) if ok else {}
    except ValueError:
        told = {}
    if "update_available" in told:
        return TunlerVersion(current=told.get("current", ""), latest=told.get("latest", ""), update_available=bool(told["update_available"]))
    current, newest = installed(), latest(host)
    return TunlerVersion(current=current, latest=newest, update_available=bool(current and newest and current != newest))


def server_name(address: str) -> str:
    server = address.strip().removeprefix("https://").removeprefix("http://").split("/")[0]
    if not server:
        raise Refused("give the address of the tunler server to install from, such as tunler.example.com")
    return server


def install(host: str) -> str:
    machine = platform.machine().lower()
    build = f"tunler-{platform.system().lower()}-{ARCHES.get(machine, machine)}"
    part = LOCAL_BIN.with_name("tunler.part")
    LOCAL_BIN.parent.mkdir(parents=True, exist_ok=True)
    try:
        with urllib.request.urlopen(f"https://{host}/dl/{build}", timeout=DOWNLOAD_SECONDS) as answer:
            part.write_bytes(answer.read())
    except OSError as error:
        part.unlink(missing_ok=True)
        raise Refused(f"tunler could not be downloaded from {host}: {error}") from error
    part.chmod(0o700)
    part.replace(LOCAL_BIN)
    return f"tunler {installed()} is installed"


def updated() -> str:
    ok, said = ran("update")
    if not ok:
        return said or "tunler did not update"
    return said.splitlines()[-1] if said else "tunler is up to date"


def log_out() -> str:
    ok, said = ran("logout")
    return "" if ok else said or "tunler did not log out"


def owned() -> list[str]:
    ok, said = ran("domains")
    return [line.strip() for line in said.splitlines() if line.strip()] if ok else []


def unclaim(domain: str, host: str) -> str:
    ok, said = ran("release", domain.removesuffix(f".{host}"))
    return "" if ok else said or f"tunler did not release {domain}"
import os
import signal
import subprocess
import time
from dataclasses import asdict, dataclass, replace
from pathlib import Path

from engine.fields import Loaded
from engine.keeper import BUILD, ServiceSpec, ServiceState, gone, teardown
from engine.stored import read_json, write_json
from engine.locks import claim
from engine.package import ARCHIVE, entry
from engine.ports import free
from engine.runtime import env, folder
from engine.sessions import alive
from typing import TypedDict
from engine.extension import Extension
from engine.record import Record

PORTS = range(8440, 8500)
UP, DOWN = "up", "down"
BLOCKED, FAILED, NOT_NEEDED = "blocked", "failed", "not needed"
RESTING = (BLOCKED, FAILED, NOT_NEEDED, "stopped", "exited")
NEEDED_FOR = 600.0
ASKED_WITHIN = 10.0
BACKOFF = (1.0, 2.0, 4.0, 8.0, 16.0, 30.0)
KEEPER_EXIT = 3.0
KEEPING = "services-keeper.lock"
CRASHES, WITHIN = 5, 60.0


def status_file(root: Path, sid: str) -> Path:
    return folder(root) / f"service-{sid}.json"


def holder(root: Path, sid: str) -> int:
    try:
        said = lock_file(root, sid).read_text().strip()
    except OSError:
        return 0
    return int(said) if said.isdigit() else 0


def lock_file(root: Path, sid: str) -> Path:
    return folder(root) / f"service-{sid}.lock"


def current_build(root: Path) -> str:
    return (Path(root) / ARCHIVE).resolve().name


def spec_file(root: Path, sid: str) -> Path:
    return folder(root) / f"spec-{sid}.json"


def want_file(root: Path, sid: str) -> Path:
    return folder(root) / f"service-{sid}.want"


def log_file(root: Path, sid: str) -> Path:
    return folder(root) / f"service-{sid}.log"


@dataclass(frozen=True)
class Wanted(Loaded):
    want: str = UP
    nonce: float = 0.0

    @classmethod
    def read(cls, root: Path, sid: str) -> "Wanted":
        return read_json(want_file(root, sid), cls.from_json, cls.from_json({}))


def status(root: Path, sid: str) -> ServiceState:
    return ServiceState.read(status_file(root, sid))


def states(root: Path) -> dict[str, ServiceState]:
    home = folder(root)
    found = sorted(home.glob("service-*.json")) if home.is_dir() else []
    return {p.stem.removeprefix("service-"): ServiceState.read(p) for p in found}


def wanted(root: Path, sid: str) -> str:
    return Wanted.read(root, sid).want


def want(root: Path, sid: str, state: str, nonce: float = 0.0) -> dict:
    asked = {"want": state if state in (UP, DOWN) else UP, "nonce": nonce}
    write_json(want_file(root, sid), asked)
    return asked


def allocate(root: Path, sid: str, wants, taken: set[int]) -> tuple[int, str]:
    if isinstance(wants, int):
        return (wants, "") if free(wants) else (wants, f"port {wants} is in use")
    held = status(root, sid)
    if held.port and held.port not in taken and (free(held.port) or alive(held.keeper)):
        return held.port, ""
    for port in PORTS:
        if port not in taken and free(port):
            return port, ""
    return 0, f"no port free from {PORTS.start} through {PORTS.stop - 1}"


def claimed(root: Path, sid: str, wants, taken: set[int]) -> tuple[int, str]:
    port, blocked = allocate(root, sid, wants, taken)
    if port:
        taken.add(port)
    return port, blocked


def local_url(port: int) -> str:
    return f"http://127.0.0.1:{port}" if port else ""


SOURCES = Extension()


class ServiceFiles(TypedDict):
    lock: str
    log: str
    status: str
    spec: str


def files_for(root: Path, sid: str) -> ServiceFiles:
    return {"lock": str(lock_file(root, sid)), "log": str(log_file(root, sid)), "status": str(status_file(root, sid)), "spec": str(spec_file(root, sid))}


def service_spec(root: Path, sid: str, **fields) -> ServiceSpec:
    return ServiceSpec(id=sid, **fields, **files_for(root, sid))


def specs(root: Path, sources) -> list[ServiceSpec]:
    taken: set[int] = set()
    return [spec for source in (*sources, *SOURCES.each(Record(root, env(root)))) for spec in source(root, taken)]


def spawn(spec: ServiceSpec, lifeline: int) -> int:
    replace(ServiceState.read(spec.status), port=spec.port, url=spec.url, owner=os.getpid()).write(spec.status)
    Path(spec.log).parent.mkdir(parents=True, exist_ok=True)
    with open(spec.log, "ab", buffering=0) as log:
        kept = subprocess.Popen([*entry("engine.keeper"), str(lifeline), spec.spec],
                                pass_fds=(lifeline,) if lifeline >= 0 else (), stdin=subprocess.DEVNULL,
                                stdout=log, stderr=log, start_new_session=True)
    return kept.pid


def excerpt(output: str) -> str:
    words = output.strip()[:160]
    return f": {words}" if words else ""


class Manager:
    def __init__(self, root: Path, lifeline: int = -1, start=spawn, clock=time.time, living=alive, sources=()):
        self.root = Path(root)
        self.sources = sources
        self.lifeline = lifeline
        self.start = start
        self.clock = clock
        self.living = living
        self.crashes: dict = {}
        self.seen: dict = {}
        self.waiting: dict = {}
        self.needed: dict = {}
        self.owned = None

    def tick(self) -> list[str]:
        self.owned = self.owned or claim(folder(self.root) / KEEPING)
        if self.owned is None:
            return []
        started = []
        declared = specs(self.root, self.sources)
        for spec in declared:
            if self.one(spec):
                started.append(spec.id)
        self.retire({spec.id for spec in declared})
        self.sweep()
        return started

    def retire(self, declared: set[str]) -> list[str]:
        gone_now = [sid for sid in states(self.root) if sid not in declared]
        for sid in gone_now:
            self.remove(sid)
            want_file(self.root, sid).unlink(missing_ok=True)
        return gone_now

    def remove(self, sid: str) -> None:
        current = status(self.root, sid)
        self.stop(sid, current)
        stray = holder(self.root, sid)
        if stray != current.keeper:
            self.stop(sid, replace(current, keeper=stray))
        if current.pgid and not gone(current.pgid):
            teardown(current.pgid, 1.0)
        until = time.monotonic() + KEEPER_EXIT
        while (self.living(current.keeper) or self.living(stray)) and time.monotonic() < until:
            time.sleep(0.05)
        for place in (status_file, spec_file, lock_file, log_file):
            place(self.root, sid).unlink(missing_ok=True)

    def stored_build(self, sid: str) -> str:
        stored = read_json(spec_file(self.root, sid), ServiceSpec.from_json, None)
        return stored.build if stored is not None else ""

    def one(self, spec: ServiceSpec) -> bool:
        sid = spec.id
        current = status(self.root, sid)
        stray = holder(self.root, sid)
        if not self.living(current.keeper) and self.living(stray):
            current = replace(current, keeper=stray)
            current.write(spec.status)
        if self.living(current.keeper) and (current.build or self.stored_build(sid)) != spec.build:
            self.remove(sid)
            current = status(self.root, sid)
        asked = Wanted.read(self.root, sid)
        now = self.clock()
        if asked.want == DOWN:
            self.stop(sid, current)
            return False
        if asked.nonce > current.nonce:
            self.remove(sid)
            self.crashes.pop(sid, None)
            self.waiting.pop(sid, None)
            self.needed.pop(sid, None)
            current = ServiceState(nonce=asked.nonce)
            current.write(spec.status)
        unneeded = self.unneeded(spec, now)
        if unneeded:
            self.stop(sid, current)
            replace(current, state=NOT_NEEDED, why=unneeded, at=now).write(spec.status)
            return False
        if self.living(current.keeper) and current.state not in RESTING:
            return False
        if spec.blocked:
            replace(current, state=BLOCKED, why=spec.blocked, at=now).write(spec.status)
            return False
        if current.state == "exited" and self.seen.get(sid) != current.at:
            self.seen[sid] = current.at
            self.crashed(sid, now)
        if len(self.crashes.get(sid, [])) >= CRASHES:
            replace(current, state=FAILED, why=f"it stopped {CRASHES} times within {WITHIN:g} seconds", at=now).write(spec.status)
            return False
        if self.waiting.get(sid, 0) > now:
            return False
        if spec.restart == "never" and current.state in ("exited", "stopped"):
            return False
        write_json(spec_file(self.root, sid), asdict(replace(spec, owner=os.getpid())))
        replace(current, state="starting", keeper=0, owner=os.getpid(), port=spec.port, url=spec.url, at=now).write(spec.status)
        self.start(spec, self.lifeline)
        return True

    def unneeded(self, spec: ServiceSpec, now: float) -> str:
        if not spec.when:
            return ""
        asked, why = self.needed.get(spec.id, (0.0, ""))
        if now - asked < NEEDED_FOR:
            return why
        try:
            ran = subprocess.run(spec.when, shell=True, cwd=spec.cwd or None, env={**os.environ, **spec.env}, capture_output=True, text=True, timeout=ASKED_WITHIN)
            why = "" if ran.returncode == 0 else f"not needed here: {spec.when} answered {ran.returncode}{excerpt(ran.stdout)}"
        except (OSError, subprocess.SubprocessError) as e:
            why = f"not needed here: {spec.when} could not be asked ({e})"
        self.needed[spec.id] = (now, why)
        return why

    def crashed(self, sid: str, now: float) -> None:
        seen = [at for at in self.crashes.get(sid, []) if now - at < WITHIN] + [now]
        self.crashes[sid] = seen
        self.waiting[sid] = now + BACKOFF[min(len(seen), len(BACKOFF)) - 1]

    def stop(self, sid: str, current: ServiceState) -> None:
        if self.living(current.keeper):
            try:
                os.kill(current.keeper, signal.SIGTERM)
            except (ProcessLookupError, PermissionError):
                pass

    def sweep(self) -> list[int]:
        killed = []
        for sid, current in states(self.root).items():
            group = current.pgid
            if not group or gone(group) or self.living(current.keeper):
                continue
            teardown(group, 1.0)
            killed.append(group)
            replace(current, state="stopped", why="its keeper is gone", at=self.clock()).write(status_file(self.root, sid))
        return killed


def listed(root: Path, sources) -> list[dict]:
    known = {spec.id: spec for spec in specs(root, sources)}
    current = states(root)
    out = []
    for sid in sorted({*known, *current}):
        spec, state = known.get(sid), current.get(sid, ServiceState())
        plugin, _, service = sid.partition(".")
        out.append({"id": sid, "plugin": spec.plugin if spec else plugin, "service": spec.service if spec else service,
                    "state": state.state if state.state else "not running", "why": state.why, "url": spec.url if spec and spec.url else state.url,
                    "port": spec.port if spec and spec.port else state.port, "since": state.started, "declared": spec is not None})
    return out
from pathlib import Path

from engine import runtime
from engine.keeper import ServiceSpec
from engine.package import entry
from engine.record import Record
from engine.services import BUILD, allocate, current_build, files_for
from features.sharing.controller import Shares
from features.sharing.resource import ended
from features.sharing.tunnel import SERVER, TUNNEL, subdomain, tunler
from resources.base import Refused, SYSTEM
from engine.extension import Extension

KEEP_UP = Extension()


def open_shares(root: Path) -> list:
    shares = Shares(Record(root, runtime.env(root)), actor=SYSTEM)
    return [row for row in shares.rows.summaries() if row.get("token") and row.get("approved") and not row["deleted"]
            and not ended(row["completed"], row.get("expires"))]


def wanted(root: Path) -> bool:
    from features import running
    from features.sharing.feature import SharingFeature
    record = Record(root, runtime.env(root))
    feature = running(SharingFeature)
    sharing = bool(feature) and feature.enabled(record)
    return (sharing and bool(open_shares(root))) or any(keep(root) for keep in KEEP_UP.each(record))


def share_services(root: Path, taken: set) -> list:
    if not wanted(root):
        return []
    port, blocked = allocate(root, SERVER, None, taken)
    taken.add(port)
    specs = [ServiceSpec(id=SERVER, plugin="sharing", service="server", run=[*entry("features.sharing.server"), str(root), str(port)],
                         cwd=str(Path(root).parent), port=port, blocked=blocked, url=f"http://127.0.0.1:{port}", env={BUILD: current_build(root)},
                         **files_for(root, SERVER))]
    command = tunler()
    if not command:
        return specs
    try:
        domain = subdomain(root)
    except Refused:
        return specs
    inspector, _ = allocate(root, TUNNEL, None, taken)
    taken.add(inspector)
    specs.append(ServiceSpec(id=TUNNEL, plugin="sharing", service="tunnel", cwd=str(Path(root).parent), port=inspector, url=f"http://127.0.0.1:{inspector}",
                             run=[command, str(port), f"--domain={domain}", f"--inspect={inspector}"],
                             env={BUILD: f"{current_build(root)}:{port}:{inspector}"}, **files_for(root, TUNNEL)))
    return specs
158:    def _unusable(self, standing: TunnelStatus) -> str:
159-        if not standing["installed"]:
160-            return NOT_INSTALLED
161-        return "" if standing["logged_in"] else LOGGED_OUT
162-
163-    def _problems(self, standing: TunnelStatus) -> list[str]:
164-        if unusable := self._unusable(standing):
165-            return [unusable]
166-        if refused_address(log_file(self.record.root, TUNNEL)):
167-            return [ADDRESS_TAKEN]
168-        if status(self.record.root, TUNNEL).state == FAILED:
169-            return [TUNNEL_STOPPED]
170-        return []
171-
172-    def _address(self) -> str:
173-        host = self._host()
174-        return f"{subdomain(self.record.root)}.{host}" if host else ""
175-
176-    @action
177-    def login(self, username: str, password: str, endpoint: str | None = None, master_password: str | None = None) -> dict:
178-        self._user_only("log tunler in")
179-        host = endpoint.strip() if endpoint else self._host()
180-        made = log_in(host, username.strip(), password, master_password or None)
--
210:    def _host(self) -> str:
211-        return SharingDetails.values(self.record).host or tunler_status()["host"]
212-
213-    @action
214-    def domains(self) -> list[str]:
215-        return owned()
216-
217-    @action
218-    def release(self, domain: str) -> list[str]:
219-        self._user_only("release a tunler domain")
220-        failed = unclaim(domain, self._host())
221-        if failed:
222-            raise Refused(failed)
223-        return owned()
224-
225-    @action
226-    def readdress(self) -> str:
227-        self._user_only("choose a new tunnel address")
228-        from engine.services import UP, want
229-        from features.sharing.services import TUNNEL
230-        name = new_address(self.record.root)
231-        alerts(self.record.root).set(ADDRESS_REFUSED, 0)
232-        want(self.record.root, TUNNEL, UP, nonce=time.time())
--
245:        return {"reachable": self._answering()}
246-
247:    def _answering(self, wait: float = REACH_SECONDS) -> bool:
248-        return answers(f"https://{self._address()}/{HEALTH}", wait)
249-
250-    def _link(self, token: str) -> str:
251-        return f"https://{self._address()}/s/{token}"
252-
253-    def _target(self, ref: str):
254-        kind, _, number = str(ref).strip().replace(" ", ":", 1).partition(":")
255-        if kind not in SHARED_TYPES or not number.isdigit():
256-            raise Refused(f"only a document, a report, a collection or a plan can be shared: write it as doc:12, report:4, collection:3 or plan:2, not {ref!r}")
257-        row = CONTROLLERS[kind](self.record, actor=SYSTEM).load(number)
258-        if row.deleted:
259-            raise Refused(f"{kind} {number} is deleted")
260-        return row
261-
262-    def _by_token(self, token: str):
263-        if not TOKEN.match(token):
264-            return None
265-        found = next((row["n"] for row in self.rows.summaries() if row.get("token") == token), None)
266-        return self.load(found) if found else None
267-
268-    def _home(self, share) -> Record:
269-        return Record(self.record.root, share.data.get("environment") or self.record.env)
