"""The process that keeps running.

A small scheduler: named jobs, each on its own interval, run in one thread until
something asks the process to stop.

Two details carry the weight.

**Waiting is interruptible.** `time.sleep(60)` would leave a SIGTERM unanswered
for up to a minute — past Docker's ten-second grace period, at which point it
sends SIGKILL and whatever was running dies mid-write. Waiting on an Event means
a stop request is answered in milliseconds.

**A failing job never stops the loop.** The CMS restarts, a network blips, a
document is malformed. Each of those is an ordinary Tuesday, and a supervisor
that exits on the first one turns a thirty-second outage into one lasting until
somebody notices.
"""

from __future__ import annotations

import logging
import signal
import threading
import time
from collections.abc import Callable
from dataclasses import dataclass, field
from datetime import UTC, datetime

logger = logging.getLogger(__name__)


@dataclass(slots=True)
class Job:
    """Something to do, repeatedly."""

    name: str
    run: Callable[[], object]
    interval_seconds: float

    #: Monotonic, not wall clock. A clock correction — NTP, a daylight-saving
    #: jump — must not make a job run twice or stall for an hour.
    next_due: float = 0.0

    runs: int = 0
    failures: int = 0
    last_error: str | None = None
    last_run_at: datetime | None = None

    def due(self, now: float) -> bool:
        return now >= self.next_due

    def schedule_after(self, now: float) -> None:
        self.next_due = now + self.interval_seconds


@dataclass(slots=True)
class Daemon:
    """Runs jobs until asked to stop."""

    jobs: list[Job] = field(default_factory=list)
    _stopping: threading.Event = field(default_factory=threading.Event)
    started_at: datetime | None = None

    running_job: str | None = None
    """The job currently inside the loop, if any."""

    running_since: datetime | None = None
    """When that job started.

    Together these are what separates "busy" from "hung". Without them a long
    job is indistinguishable from a deadlock, and on 1 August 2026 that had
    production reporting itself unhealthy while it read a client's document —
    which is what the process is for. Anything acting on that signal, an
    orchestrator most of all, would have restarted a working agent mid-read.
    """

    #: How long a single job may run before the loop is called stalled anyway.
    #:
    #: Generous, because OCR is measured in tens of seconds a page and a long
    #: PDF is minutes of legitimate work. Bounded, because "a job is running"
    #: must not become an excuse that never expires: a genuinely hung job has to
    #: surface eventually, and this is when.
    max_job_seconds: float = 15 * 60

    last_tick_at: datetime | None = None
    """When the loop last came round.

    The liveness signal. If a job hangs — a socket with no timeout, a deadlock —
    this stops advancing while the process still exists, answers signals and
    looks perfectly healthy to anything checking that it is running. A stalled
    tick is the one condition a restart genuinely fixes.
    """

    #: Longest the loop will wait between wake-ups. Bounded so a long interval
    #: does not delay the *first* run of a job added later, and so the loop stays
    #: responsive without spinning.
    max_wait_seconds: float = 5.0

    def add(self, name: str, run: Callable[[], object], interval_seconds: float) -> Job:
        job = Job(name=name, run=run, interval_seconds=interval_seconds)
        self.jobs.append(job)

        return job

    def install_signal_handlers(self) -> None:
        """Answer SIGTERM and SIGINT by stopping cleanly.

        SIGTERM is what `docker stop` sends; SIGINT is Ctrl-C. Installed
        separately from `run` so a caller embedding this — a test, or a process
        that owns its own signals — is not forced to take ours.
        """
        for received in (signal.SIGTERM, signal.SIGINT):
            try:
                signal.signal(received, self._on_signal)
            except (ValueError, OSError):
                # Not the main thread, or a platform without that signal.
                # Worth continuing: the daemon still stops through stop().
                logger.debug("Could not install a handler for %s.", received)

    def stop(self) -> None:
        self._stopping.set()

    @property
    def stopping(self) -> bool:
        return self._stopping.is_set()

    def is_ticking(self, tolerance_seconds: float | None = None) -> bool:
        """Whether the loop is still coming round.

        Tolerant by default of several missed passes. A restart throws away
        in-flight work, so the bar for declaring the process dead should be
        higher than one slow iteration — a job that takes a moment longer than
        usual is not a reason to kill it.

        A daemon that has been asked to stop is still 'ticking': it is shutting
        down on purpose, and reporting that as a failure would have an
        orchestrator restart something that was told to exit.
        """
        if self.started_at is None or self.stopping:
            return True

        # A job in flight is the loop working, not the loop stalled. It cannot
        # come round until the job returns, so judging it by the tick would
        # report every long read as a stall. Still bounded: a job that has run
        # past `max_job_seconds` is no longer plausibly working.
        if self.running_since is not None:
            return (datetime.now(UTC) - self.running_since).total_seconds() <= self.max_job_seconds

        if self.last_tick_at is None:
            return False

        tolerance = tolerance_seconds or (self.max_wait_seconds * 6)

        return (datetime.now(UTC) - self.last_tick_at).total_seconds() <= tolerance

    def run(self, max_iterations: int | None = None) -> None:
        """The loop.

        ``max_iterations`` exists for tests, which need this to end. Production
        passes nothing and it runs until signalled.
        """
        self.started_at = datetime.now(UTC)
        iterations = 0

        logger.info("Started with %d job(s): %s", len(self.jobs),
                    ", ".join(j.name for j in self.jobs) or "none")

        while not self.stopping:
            if max_iterations is not None and iterations >= max_iterations:
                break

            iterations += 1
            self.last_tick_at = datetime.now(UTC)
            now = time.monotonic()

            for job in self.jobs:
                if self.stopping:
                    # Checked between jobs as well as between iterations, so a
                    # stop arriving during a slow pass is not held up by every
                    # remaining job.
                    break

                if job.due(now):
                    self._run_job(job)
                    job.schedule_after(time.monotonic())

            self._wait(self._seconds_until_next_due())

        logger.info("Stopped cleanly.")

    # ── Internals ─────────────────────────────────────────────────────────

    def _run_job(self, job: Job) -> None:
        job.runs += 1
        job.last_run_at = datetime.now(UTC)

        from app.runtime.metrics import METRICS

        # Published so liveness can tell "working" from "hung". The loop cannot
        # come round while a job is inside it, and reading a document takes tens
        # of seconds a page — long enough that the tick goes stale during
        # entirely ordinary work.
        self.running_job = job.name
        self.running_since = datetime.now(UTC)

        try:
            with METRICS.time("taxpilot_job_seconds", "Time spent in a scheduled job.",
                              job=job.name):
                job.run()

            job.last_error = None
            METRICS.counter("taxpilot_jobs_total", "Scheduled job runs.",
                            job=job.name, outcome="ok")
        except Exception as exc:  # noqa: BLE001 - see the module docstring
            job.failures += 1
            job.last_error = str(exc)
            METRICS.counter("taxpilot_jobs_total", "Scheduled job runs.",
                            job=job.name, outcome="failed")
            logger.exception("Job '%s' failed: %s", job.name, exc)
        finally:
            # Cleared in `finally` so a job that raised does not leave the loop
            # looking permanently busy, which would suppress the very signal
            # this exists to preserve.
            self.running_job = None
            self.running_since = None

    def _seconds_until_next_due(self) -> float:
        if not self.jobs:
            return self.max_wait_seconds

        now = time.monotonic()
        soonest = min(job.next_due for job in self.jobs)

        return max(0.0, min(soonest - now, self.max_wait_seconds))

    def _wait(self, seconds: float) -> None:
        """Sleep, unless asked to stop.

        Event.wait rather than time.sleep: this is the difference between
        answering a SIGTERM immediately and being killed mid-write when Docker
        loses patience.
        """
        if seconds > 0:
            self._stopping.wait(timeout=seconds)

    def _on_signal(self, received: int, _frame) -> None:
        logger.info("Received %s; finishing the current job and stopping.",
                    signal.Signals(received).name)
        self.stop()
