"""Bringing decisions back from the CMS (ADR-0008).

A workflow that submitted a proposal is stopped, holding a proposal id and
nothing else. Somebody in the CMS eventually approves, rejects, or sends it back,
and this is what notices and starts the run moving again.

Polling rather than a callback. The CMS knows how to reach nothing; this process
knows how to reach the CMS. A push would need the installation to hold an address
for this platform and to be able to open connections to it, and neither is true
on shared hosting.
"""

from __future__ import annotations

import logging
from collections.abc import Iterable
from dataclasses import dataclass
from typing import Any

from app.api.client import AgentApiError, CmsClient
from app.workflow.engine import RunStore, Workflow, WorkflowEngine
from app.workflow.state import RunState, StepState, WorkflowRun

logger = logging.getLogger(__name__)

#: Which step a reviewer's directive sends the workflow back to.
#:
#: "Read the document again" means exactly that — re-run the OCR and everything
#: downstream of it, because a different reading changes the classification, the
#: extracted fields and possibly the client.
REDO_TARGETS = {
    "rerun_ocr": "read",
    "rerun_classification": "classify",
    "rerun_client_search": "identify",
}


@dataclass(frozen=True, slots=True)
class DecisionOutcome:
    """What came back for one run."""

    run_id: str
    proposal_id: int
    status: str
    applied: bool
    detail: str = ""


class DecisionPoller:
    """Asks the CMS what happened, and moves the affected runs on.

    Holds every workflow this deployment runs, not one. A decision comes back
    with a proposal id and nothing else — which workflow it belongs to is a
    property of the paused run, and resuming a run against the wrong definition
    would step it through somebody else's sequence.
    """

    def __init__(
        self,
        cms: CmsClient,
        engine: WorkflowEngine,
        workflows: Workflow | Iterable[Workflow],
    ) -> None:
        self._cms = cms
        self._engine = engine
        self._workflows = {
            w.name: w
            for w in ([workflows] if isinstance(workflows, Workflow) else workflows)
        }
        self._since: str | None = None

    @property
    def store(self) -> RunStore:
        return self._engine.store

    def poll(self) -> list[DecisionOutcome]:
        """One pass. Returns what it did, which is what a caller logs."""
        try:
            decided = self._cms.decided_proposals(since=self._since)
        except AgentApiError as exc:
            # The CMS being unreachable is an operational fact, not a reason to
            # crash a long-running poller. The next pass asks again, and `since`
            # is deliberately not advanced so nothing is skipped.
            logger.warning("Could not reach the CMS for decisions: %s", exc)

            return []

        waiting = self._waiting_runs()
        outcomes: list[DecisionOutcome] = []

        for proposal in decided:
            run = waiting.get(proposal.get("id"))

            if run is None:
                # A decision for a run this process does not know about — most
                # often one lost in a restart before there was a durable store.
                # Skipped rather than treated as an error; the document is filed
                # either way, and the CMS holds the record.
                continue

            outcomes.append(self._apply(run, proposal))

        # Advanced only after everything was applied, so a failure mid-batch
        # leaves the batch to be seen again rather than silently dropped.
        if decided:
            self._since = max(
                (p.get("updated_at") for p in decided if p.get("updated_at")),
                default=self._since,
            )

        return outcomes

    # ── Internals ─────────────────────────────────────────────────────────

    def _waiting_runs(self) -> dict[int, WorkflowRun]:
        """Runs stopped on a proposal, keyed by the proposal they are waiting on."""
        waiting: dict[int, WorkflowRun] = {}

        for run in self.store.resumable():
            if run.state is not RunState.AWAITING_APPROVAL:
                continue

            paused = next((s for s in run.steps if s.state is StepState.AWAITING_APPROVAL), None)

            if paused and (proposal_id := paused.output.get("proposal_id")):
                waiting[int(proposal_id)] = run

        return waiting

    def _definition(self, run: WorkflowRun) -> Workflow | None:
        """The definition this run was started from."""
        return self._workflows.get(run.workflow)

    def _redo_target(self, workflow: Workflow, directive: str | None) -> str:
        """Which step "do it again" means, in this workflow.

        The directives are named for document intake's steps, and two of them —
        re-read, re-classify — describe work the bank statement workflow does not
        do and never will. Sending it back to a step it has no name for fails the
        run with "no such step", which turns a reviewer asking a reasonable
        question into a lost document.

        So a directive that does not name a step in this workflow rewinds to its
        first one. For the forwarding workflow that is the client lookup, which
        is the only thing there is to do again — and is what all three directives
        amount to asking for.
        """
        names = [s.name for s in workflow.steps]
        target = REDO_TARGETS.get(str(directive)) if directive else None

        if target in names:
            return target

        return names[0]

    def _apply(self, run: WorkflowRun, proposal: dict[str, Any]) -> DecisionOutcome:
        status = str(proposal.get("status", ""))
        proposal_id = int(proposal.get("id"))
        workflow = self._definition(run)

        if workflow is None:
            # A run of a workflow this deployment no longer defines — a release
            # that removed one, most likely. Left paused rather than failed: the
            # CMS holds the decision and an upgrade that restores the definition
            # picks it up, whereas failing it here throws that away.
            logger.warning(
                "Run %s is waiting on proposal %s but '%s' is not a workflow this "
                "deployment defines.", run.id, proposal_id, run.workflow,
            )

            return DecisionOutcome(run.id, proposal_id, status, False, "Unknown workflow.")

        if status == "executed":
            self._engine.apply_decision(
                workflow, run,
                accepted=True,
                data={
                    "document_id": proposal.get("result_id"),
                    "client_id": (proposal.get("amendments") or {}).get("client_id"),
                    "corrected": bool(proposal.get("amendments")),
                },
            )

            return DecisionOutcome(run.id, proposal_id, status, True, "Filed.")

        if status == "rejected":
            self._engine.apply_decision(
                workflow, run,
                accepted=False,
                reason=proposal.get("decision_note") or "Rejected by a reviewer.",
            )

            return DecisionOutcome(run.id, proposal_id, status, True, "Rejected.")

        if status == "changes_requested":
            target = self._redo_target(workflow, proposal.get("directive"))

            self._engine.apply_decision(
                workflow, run,
                accepted=False,
                reason=proposal.get("decision_note"),
                redo_from=target,
            )

            return DecisionOutcome(run.id, proposal_id, status, True, f"Redoing from '{target}'.")

        # approved-but-not-yet-filed, or failed and awaiting a retry over there.
        # Both are still in progress: leaving the run paused is correct, and
        # advancing it would claim an outcome the CMS has not reached.
        return DecisionOutcome(run.id, proposal_id, status, False, "Still in progress.")
