"""Orchestrator graph: … → IMPLEMENTING → REVIEWING → TESTING → PR_CREATED | FAILED."""

from __future__ import annotations

import re
import uuid
from pathlib import Path
from typing import Any, Dict, Optional, TypedDict

import structlog
from langgraph.graph import END, StateGraph
from sqlalchemy import select
from sqlalchemy.ext.asyncio import AsyncSession

from packages.agents.code_reviewer import CodeReviewer
from packages.agents.qa_engineer import QAEngineer
from packages.agents.requirement_analyst import RequirementAnalyst
from packages.agents.salesforce_developer import SalesforceDeveloper
from packages.agents.solution_architect import SolutionArchitect
from packages.artifacts.design import DesignArtifact
from packages.artifacts.implementation import ImplementationArtifact
from packages.artifacts.requirements import RequirementsArtifact
from packages.artifacts.review import ReviewArtifact
from packages.config import get_settings
from packages.db.artifacts import ArtifactStore
from packages.db.models import ACTIVE_STATES, AuditEvent, Run, WorkflowState
from packages.integrations.confluence import (
    ConfluencePublisher,
    design_doc_storage_html,
    design_doc_title,
    get_confluence_publisher,
)
from packages.integrations.doc_validation import (
    DocReference,
    build_doc_check_prompt,
    doc_review_comment_body,
)
from packages.integrations.duplicates import (
    build_duplicate_check_prompt,
    candidate_search_jql,
    duplicate_review_comment_body,
)
from packages.integrations.git import GitClient, get_git_client
from packages.integrations.jira import JiraClient, JiraIssue, get_jira_client
from packages.integrations.salesforce import ApexTestRunner, SalesforceDeployer, get_apex_test_runner, get_salesforce_deployer
from packages.integrations.static_analysis import StaticAnalyzer, get_static_analyzer

logger = structlog.get_logger(__name__)

_REJECT_RE = re.compile(r"^(reject|rejected)\b", re.IGNORECASE)


class OrchestratorState(TypedDict, total=False):
    run_id: str
    jira_key: str
    resume: bool
    error: Optional[str]
    next_action: str
    design_decision: str  # approve | reject
    duplicate_decision: str  # continue | close
    doc_decision: str  # continue | close
    failure_phase: str
    review_cycle: int
    fix_mode: bool
    file_map: Dict[str, str]


async def _get_run(session: AsyncSession, run_id: uuid.UUID) -> Run:
    run = await session.get(Run, run_id)
    if run is None:
        raise ValueError(f"Run not found: {run_id}")
    return run


async def _audit(
    session: AsyncSession,
    run_id: uuid.UUID,
    event_type: str,
    message: str,
    payload: Optional[Dict[str, Any]] = None,
) -> None:
    session.add(
        AuditEvent(
            run_id=run_id,
            event_type=event_type,
            message=message,
            payload=payload or {},
        )
    )


def _is_design_reject_comment(body: Optional[str]) -> bool:
    if not body or not body.strip():
        return False
    return bool(_REJECT_RE.match(body.strip()))


def _reject_feedback_hint(body: Optional[str]) -> str:
    """Strip leading reject/rejected token; return remaining feedback if any."""
    if not body or not body.strip():
        return ""
    cleaned = _REJECT_RE.sub("", body.strip(), count=1).lstrip(" —-:").strip()
    return cleaned


def _design_reject_clarification_body(hint: str) -> str:
    lines = [
        "## Agent Nova needs clarification",
        "",
        "The previous design was rejected.",
    ]
    if hint:
        lines.extend(["", f'Your reject note: "{hint}"'])
    lines.extend(
        [
            "",
            "Please describe what should change in the design so Agent Nova can "
            "produce a new design document.",
            "",
            "Next step: reply in a comment with your clarification to resume.",
        ]
    )
    return "\n".join(lines)


def _apex_class_names(file_map: Dict[str, str]) -> list[str]:
    names: list[str] = []
    for path in file_map:
        if path.endswith(".cls") and "Test.cls" not in path:
            name = Path(path).stem
            if name and not name.endswith("Test"):
                names.append(name)
    return names


class Orchestrator:
    def __init__(
        self,
        session: AsyncSession,
        jira: Optional[JiraClient] = None,
        analyst: Optional[RequirementAnalyst] = None,
        architect: Optional[SolutionArchitect] = None,
        developer: Optional[SalesforceDeveloper] = None,
        reviewer: Optional[CodeReviewer] = None,
        qa: Optional[QAEngineer] = None,
        git: Optional[GitClient] = None,
        deployer: Optional[SalesforceDeployer] = None,
        static_analysis: Optional[StaticAnalyzer] = None,
        apex_runner: Optional[ApexTestRunner] = None,
        artifacts: Optional[ArtifactStore] = None,
        confluence: Optional[ConfluencePublisher] = None,
    ) -> None:
        self.session = session
        self.jira = jira or get_jira_client()
        self.analyst = analyst or RequirementAnalyst()
        self.architect = architect or SolutionArchitect()
        self.developer = developer or SalesforceDeveloper()
        self.reviewer = reviewer or CodeReviewer()
        self.qa = qa or QAEngineer()
        self.git = git or get_git_client()
        self.deployer = deployer or get_salesforce_deployer()
        self.static_analysis = static_analysis or get_static_analyzer()
        self.apex_runner = apex_runner or get_apex_test_runner()
        self.artifacts = artifacts or ArtifactStore()
        self.settings = get_settings()
        self.confluence = confluence or get_confluence_publisher(self.settings)
        self.graph = self._build_graph()

    def _build_graph(self):
        graph = StateGraph(OrchestratorState)
        graph.add_node("route", self._route_node)
        graph.add_node("analyze", self._analyze_node)
        graph.add_node("clarification_gate", self._clarification_gate_node)
        graph.add_node("duplicate_review_gate", self._duplicate_review_gate_node)
        graph.add_node("doc_review_gate", self._doc_review_gate_node)
        graph.add_node("architect", self._architect_node)
        graph.add_node("design_review_gate", self._design_review_gate_node)
        graph.add_node("design_approve", self._design_approve_node)
        graph.add_node("design_reject_clarify", self._design_reject_clarify_node)
        graph.add_node("implement", self._implement_node)
        graph.add_node("review", self._review_node)
        graph.add_node("qa", self._qa_node)
        graph.add_node("create_pr", self._create_pr_node)
        graph.add_node("fail", self._fail_node)
        graph.add_node("noop", self._noop_node)

        graph.set_entry_point("route")
        graph.add_conditional_edges(
            "route",
            self._after_route,
            {
                "analyze": "analyze",
                "design_approve": "design_approve",
                "design_reject_clarify": "design_reject_clarify",
                "fail": "fail",
                "noop": "noop",
            },
        )
        graph.add_edge("noop", END)
        graph.add_conditional_edges(
            "analyze",
            self._after_analyze,
            {
                "clarification_gate": "clarification_gate",
                "duplicate_review_gate": "duplicate_review_gate",
                "doc_review_gate": "doc_review_gate",
                "architect": "architect",
                "fail": "fail",
            },
        )
        graph.add_conditional_edges(
            "architect",
            self._after_architect,
            {"design_review_gate": "design_review_gate", "fail": "fail"},
        )
        graph.add_conditional_edges(
            "implement",
            self._after_implement,
            {"review": "review", "fail": "fail"},
        )
        graph.add_conditional_edges(
            "review",
            self._after_review,
            {"implement": "implement", "qa": "qa", "fail": "fail"},
        )
        graph.add_conditional_edges(
            "qa",
            self._after_qa,
            {"create_pr": "create_pr", "fail": "fail"},
        )
        graph.add_conditional_edges(
            "create_pr",
            self._after_create_pr,
            {"end": END, "fail": "fail"},
        )
        graph.add_edge("clarification_gate", END)
        graph.add_edge("duplicate_review_gate", END)
        graph.add_edge("doc_review_gate", END)
        graph.add_edge("design_review_gate", END)
        graph.add_edge("design_reject_clarify", END)
        graph.add_edge("design_approve", "implement")
        graph.add_edge("fail", END)
        return graph.compile()

    async def run(self, run_id: uuid.UUID, resume: bool = False) -> Run:
        initial: OrchestratorState = {
            "run_id": str(run_id),
            "resume": resume,
            "next_action": "analyze",
            "review_cycle": 0,
            "fix_mode": False,
        }
        await self.graph.ainvoke(initial)
        await self.session.commit()
        return await _get_run(self.session, run_id)

    async def _route_node(self, state: OrchestratorState) -> OrchestratorState:
        run_id = uuid.UUID(state["run_id"])
        run = await _get_run(self.session, run_id)
        if state.get("resume") and run.state == WorkflowState.DUPLICATE_REVIEW:
            decision = (run.artifacts or {}).get("duplicate_decision") or ""
            original = (run.artifacts or {}).get("duplicate_of") or ""
            if decision == "close":
                state["error"] = (
                    f"Closed as duplicate of {original}" if original else "Closed as duplicate"
                )
                state["failure_phase"] = "duplicate"
                state["next_action"] = "fail"
                state["duplicate_decision"] = "close"
            else:
                # continue (default when resume decision is continue)
                arts = dict(run.artifacts or {})
                arts["duplicate_check_skipped"] = True
                arts.pop("duplicate_decision", None)
                run.artifacts = arts
                await self.session.flush()
                state["next_action"] = "analyze"
                state["duplicate_decision"] = "continue"
            return state
        if state.get("resume") and run.state == WorkflowState.DOC_REVIEW:
            decision = (run.artifacts or {}).get("doc_decision") or ""
            if decision == "close":
                state["error"] = (
                    "Closed: functionality already documented/implemented in Confluence"
                )
                state["failure_phase"] = "doc"
                state["next_action"] = "fail"
                state["doc_decision"] = "close"
            else:
                arts = dict(run.artifacts or {})
                arts["doc_check_skipped"] = True
                arts["duplicate_check_skipped"] = True
                arts.pop("doc_decision", None)
                run.artifacts = arts
                await self.session.flush()
                state["next_action"] = "analyze"
                state["doc_decision"] = "continue"
            return state
        if state.get("resume") and run.state == WorkflowState.DESIGN_REVIEW:
            decision = (run.artifacts or {}).get("design_decision") or "approve"
            if decision == "reject" or _is_design_reject_comment(
                (run.artifacts or {}).get("design_resume_comment")
            ):
                state["next_action"] = "design_reject_clarify"
                state["design_decision"] = "reject"
            else:
                state["next_action"] = "design_approve"
                state["design_decision"] = "approve"
            return state
        # Duplicate Jira webhooks often enqueue a second resume while the first is
        # already past DESIGN_REVIEW. Never re-analyze in that case — it overwrites success.
        if state.get("resume") and run.state in {
            WorkflowState.IMPLEMENTING,
            WorkflowState.REVIEWING,
            WorkflowState.TESTING,
            WorkflowState.PR_CREATED,
            WorkflowState.COMPLETED,
            WorkflowState.FAILED,
            WorkflowState.HUMAN_REVIEW,
        }:
            state["next_action"] = "noop"
            return state
        state["next_action"] = "analyze"
        return state

    async def _noop_node(self, state: OrchestratorState) -> OrchestratorState:
        return state

    def _after_route(self, state: OrchestratorState) -> str:
        return state.get("next_action") or "analyze"

    async def _analyze_node(self, state: OrchestratorState) -> OrchestratorState:
        run_id = uuid.UUID(state["run_id"])
        run = await _get_run(self.session, run_id)
        run.state = WorkflowState.ANALYZING
        await _audit(self.session, run_id, "analyze.start", f"Analyzing {run.jira_key}")
        try:
            issue = await self.jira.get_issue(run.jira_key)
            arts = dict(run.artifacts or {})
            if not arts.get("duplicate_check_skipped"):
                dup_action = await self._maybe_flag_duplicate(run, issue, arts)
                if dup_action:
                    run.artifacts = arts
                    state["next_action"] = dup_action
                    await self.session.flush()
                    return state

            if not arts.get("doc_check_skipped"):
                doc_action = await self._maybe_flag_existing_docs(run, issue, arts)
                if doc_action:
                    run.artifacts = arts
                    state["next_action"] = doc_action
                    await self.session.flush()
                    return state

            artifact = await self.analyst.analyze(issue)
            path = self.artifacts.write_json(run_id, "requirements.json", artifact)
            run.artifact_dir = str(self.artifacts.run_dir(run_id))
            arts["requirements.json"] = path
            run.artifacts = arts
            await _audit(
                self.session,
                run_id,
                "analyze.complete",
                "Requirement analysis complete",
                {"is_clear": artifact.is_clear, "confidence": artifact.confidence},
            )
            state["jira_key"] = run.jira_key
            state["next_action"] = "clarification_gate" if not artifact.is_clear else "architect"
            run.error = None
            await self.session.flush()
            arts["is_clear"] = artifact.is_clear
            arts["clarification_body"] = (
                artifact.clarification_comment_body() if not artifact.is_clear else None
            )
            run.artifacts = arts
            await self.session.flush()
            return state
        except Exception as exc:  # noqa: BLE001
            logger.exception("analyze_failed", run_id=str(run_id))
            state["error"] = str(exc)
            state["next_action"] = "fail"
            return state

    async def _maybe_flag_existing_docs(
        self,
        run: Run,
        issue: JiraIssue,
        arts: Dict[str, Any],
    ) -> Optional[str]:
        """Return 'doc_review_gate' when Confluence suggests already implemented."""
        try:
            query = f"{issue.summary} {issue.description}".strip()
            hits = await self.analyst.knowledge.confluence.search(query, limit=8)
            conf_hits = [h for h in hits if (h.source or "").lower() == "confluence"]
            if not conf_hits:
                return None
            prompt = build_doc_check_prompt(issue, conf_hits)
            result = await self.analyst.llm.check_already_implemented(prompt)
            await _audit(
                self.session,
                run.id,
                "doc.check",
                "Confluence already-implemented check completed",
                {
                    "already_implemented": result.already_implemented,
                    "confidence": result.confidence,
                    "hits": len(conf_hits),
                    "references": [{"title": r.title, "url": r.url} for r in result.references],
                },
            )
            if not result.already_implemented or not result.references:
                return None
            arts["doc_already_implemented"] = True
            arts["doc_rationale"] = result.rationale
            arts["doc_confidence"] = result.confidence
            arts["doc_references"] = [
                {"title": r.title, "url": r.url, "snippet": r.snippet} for r in result.references
            ]
            arts["doc_review_body"] = doc_review_comment_body(result.references, result.rationale)
            return "doc_review_gate"
        except Exception as exc:  # noqa: BLE001
            logger.warning("doc_check_failed", jira_key=run.jira_key, error=str(exc))
            await _audit(
                self.session,
                run.id,
                "doc.check_failed",
                f"Confluence doc check failed (continuing): {exc}",
            )
            return None

    async def _doc_review_gate_node(self, state: OrchestratorState) -> OrchestratorState:
        run_id = uuid.UUID(state["run_id"])
        run = await _get_run(self.session, run_id)
        arts = dict(run.artifacts or {})
        refs = [
            DocReference(
                title=r.get("title") or "",
                url=r.get("url") or "",
                snippet=r.get("snippet") or "",
            )
            for r in (arts.get("doc_references") or [])
            if isinstance(r, dict)
        ]
        body = arts.get("doc_review_body") or doc_review_comment_body(
            refs, arts.get("doc_rationale") or ""
        )
        await self.jira.add_comment(run.jira_key, body)
        run.state = WorkflowState.DOC_REVIEW
        await _audit(
            self.session,
            run_id,
            "doc.review_posted",
            "Posted Confluence already-implemented review",
            payload={"references": arts.get("doc_references") or []},
        )
        await self.session.flush()
        return state

    async def _maybe_flag_duplicate(
        self,
        run: Run,
        issue: JiraIssue,
        arts: Dict[str, Any],
    ) -> Optional[str]:
        """Return 'duplicate_review_gate' when a likely duplicate is found; else None."""
        try:
            project_key = issue.project_key or run.jira_key.split("-")[0]
            jql = candidate_search_jql(project_key, run.jira_key)
            candidates = await self.jira.search_issues(jql, max_results=20)
            if not candidates:
                return None
            prompt = build_duplicate_check_prompt(issue, candidates)
            result = await self.analyst.llm.check_duplicate(prompt)
            await _audit(
                self.session,
                run.id,
                "duplicate.check",
                "Duplicate check completed",
                {
                    "is_duplicate": result.is_duplicate,
                    "original_key": result.original_key,
                    "confidence": result.confidence,
                    "candidates": len(candidates),
                },
            )
            if not result.is_duplicate or not result.original_key:
                return None
            arts["duplicate_of"] = result.original_key
            arts["duplicate_rationale"] = result.rationale
            arts["duplicate_confidence"] = result.confidence
            arts["duplicate_review_body"] = duplicate_review_comment_body(
                result.original_key, result.rationale
            )
            return "duplicate_review_gate"
        except Exception as exc:  # noqa: BLE001
            logger.warning(
                "duplicate_check_failed",
                jira_key=run.jira_key,
                error=str(exc),
            )
            await _audit(
                self.session,
                run.id,
                "duplicate.check_failed",
                f"Duplicate check failed (continuing): {exc}",
            )
            return None

    def _after_analyze(self, state: OrchestratorState) -> str:
        return state.get("next_action") or "fail"

    async def _duplicate_review_gate_node(self, state: OrchestratorState) -> OrchestratorState:
        run_id = uuid.UUID(state["run_id"])
        run = await _get_run(self.session, run_id)
        arts = dict(run.artifacts or {})
        original = arts.get("duplicate_of") or "UNKNOWN"
        body = arts.get("duplicate_review_body") or duplicate_review_comment_body(original)
        await self.jira.add_comment(run.jira_key, body)
        run.state = WorkflowState.DUPLICATE_REVIEW
        await _audit(
            self.session,
            run_id,
            "duplicate.review_posted",
            f"Posted possible duplicate of {original}",
            payload={"duplicate_of": original},
        )
        await self.session.flush()
        return state

    async def _clarification_gate_node(self, state: OrchestratorState) -> OrchestratorState:
        run_id = uuid.UUID(state["run_id"])
        run = await _get_run(self.session, run_id)
        body = (run.artifacts or {}).get("clarification_body") or "Agent Nova needs clarification."
        await self.jira.add_comment(run.jira_key, body)
        run.state = WorkflowState.NEEDS_CLARIFICATION
        await _audit(
            self.session,
            run_id,
            "clarification.posted",
            "Posted clarification questions to Jira",
        )
        await self.session.flush()
        return state

    async def _architect_node(self, state: OrchestratorState) -> OrchestratorState:
        run_id = uuid.UUID(state["run_id"])
        run = await _get_run(self.session, run_id)
        run.state = WorkflowState.ARCHITECTING
        await _audit(self.session, run_id, "architect.start", "Solution Architect designing")
        try:
            req_path = (run.artifacts or {}).get("requirements.json")
            if not req_path:
                raise ValueError("requirements.json missing before architect stage")
            raw = self.artifacts.read_json(run_id, "requirements.json")
            requirements = RequirementsArtifact.model_validate(raw)
            design = await self.architect.design(requirements)
            path = self.artifacts.write_json(run_id, "design.json", design)
            arts = dict(run.artifacts or {})
            arts["design.json"] = path
            arts["design_review_body"] = design.design_review_comment_body()
            run.artifacts = arts
            await _audit(
                self.session,
                run_id,
                "architect.complete",
                "Technical design and metadata plan written",
                {
                    "out_of_allowlist": design.out_of_allowlist,
                    "plan_items": len(design.metadata_plan),
                    "confidence": design.confidence,
                },
            )
            if not design.metadata_plan:
                state["error"] = (
                    "Design metadata plan empty after allowlist enforcement — human takeover required"
                )
                state["next_action"] = "fail"
                await self.session.flush()
                return state
            state["next_action"] = "design_review_gate"
            await self.session.flush()
            return state
        except Exception as exc:  # noqa: BLE001
            logger.exception("architect_failed", run_id=str(run_id))
            state["error"] = str(exc)
            state["next_action"] = "fail"
            await self.session.flush()
            return state

    def _after_architect(self, state: OrchestratorState) -> str:
        return state.get("next_action") or "fail"

    async def _design_review_gate_node(self, state: OrchestratorState) -> OrchestratorState:
        run_id = uuid.UUID(state["run_id"])
        run = await _get_run(self.session, run_id)
        body = (run.artifacts or {}).get("design_review_body")
        if not body:
            raw = self.artifacts.read_json(run_id, "design.json")
            design = DesignArtifact.model_validate(raw)
            body = design.design_review_comment_body()
            arts = dict(run.artifacts or {})
            arts["design_review_body"] = body
            run.artifacts = arts

        arts = dict(run.artifacts or {})
        confluence_note = await self._publish_design_confluence_doc(run_id, run.jira_key, arts)
        if confluence_note:
            body = f"{confluence_note}\n\n{body}"
        run.artifacts = arts

        await self.jira.add_comment(run.jira_key, body)
        run.state = WorkflowState.DESIGN_REVIEW
        await _audit(
            self.session,
            run_id,
            "design_review.posted",
            "Posted design for human review",
            payload={"confluence_url": arts.get("confluence_url")},
        )
        await self.session.flush()
        return state

    async def _publish_design_confluence_doc(
        self,
        run_id: uuid.UUID,
        jira_key: str,
        arts: Dict[str, Any],
    ) -> Optional[str]:
        """Create Confluence design doc; soft-fail so design review still proceeds."""
        space_key = self.settings.confluence_write_space
        if not space_key:
            await _audit(
                self.session,
                run_id,
                "confluence.create_failed",
                "No Confluence write space configured",
            )
            return (
                "## Confluence design doc\n"
                "Not created (set CONFLUENCE_WRITE_SPACE_KEY or CONFLUENCE_SPACE_KEYS)."
            )
        try:
            req_raw = self.artifacts.read_json(run_id, "requirements.json")
            design_raw = self.artifacts.read_json(run_id, "design.json")
            requirements = RequirementsArtifact.model_validate(req_raw)
            design = DesignArtifact.model_validate(design_raw)
            title = design_doc_title(jira_key, design.summary)
            storage = design_doc_storage_html(requirements, design)
            page = await self.confluence.create_page(space_key, title, storage)
            arts["confluence_url"] = page.url
            arts["confluence_page_id"] = page.id
            await _audit(
                self.session,
                run_id,
                "confluence.created",
                "Created Confluence design doc",
                payload={"url": page.url, "page_id": page.id, "space_key": space_key},
            )
            return f"## Confluence design doc\n{page.url}"
        except Exception as exc:  # noqa: BLE001
            logger.warning(
                "confluence_create_failed",
                jira_key=jira_key,
                run_id=str(run_id),
                error=str(exc),
            )
            await _audit(
                self.session,
                run_id,
                "confluence.create_failed",
                f"Failed to create Confluence design doc: {exc}",
            )
            return f"## Confluence design doc\nNot created ({exc})."

    async def _design_approve_node(self, state: OrchestratorState) -> OrchestratorState:
        run_id = uuid.UUID(state["run_id"])
        run = await _get_run(self.session, run_id)
        arts = dict(run.artifacts or {})
        arts["design_approved"] = True
        arts["review_cycle"] = 0
        run.artifacts = arts
        run.state = WorkflowState.IMPLEMENTING
        run.error = None
        state["review_cycle"] = 0
        state["fix_mode"] = False
        await _audit(
            self.session,
            run_id,
            "design.approved",
            "Human approved design; starting implementation",
        )
        await self.session.flush()
        return state

    async def _design_reject_clarify_node(self, state: OrchestratorState) -> OrchestratorState:
        """After design reject: ask for clarification and pause (do not fail)."""
        run_id = uuid.UUID(state["run_id"])
        run = await _get_run(self.session, run_id)
        reject_comment = (run.artifacts or {}).get("design_resume_comment") or ""
        hint = _reject_feedback_hint(reject_comment)
        body = _design_reject_clarification_body(hint)

        arts = dict(run.artifacts or {})
        arts.pop("design_decision", None)
        arts.pop("design_approved", None)
        arts["design_reject_hint"] = hint
        arts["clarification_body"] = body
        arts["design_revision"] = int(arts.get("design_revision") or 0) + 1
        run.artifacts = arts
        run.error = None

        await self.jira.add_comment(run.jira_key, body)
        run.state = WorkflowState.NEEDS_CLARIFICATION
        await _audit(
            self.session,
            run_id,
            "design.rejected_needs_clarification",
            "Design rejected; asked for clarification before redesign",
            payload={"hint": hint, "design_revision": arts["design_revision"]},
        )
        await self.session.flush()
        return state

    async def _load_context(
        self, run_id: uuid.UUID, run: Run
    ) -> tuple[DesignArtifact, Optional[RequirementsArtifact], Dict[str, Any]]:
        arts = dict(run.artifacts or {})
        if not arts.get("design_approved"):
            raise ValueError("design_approved missing before implement stage")
        if not arts.get("design.json"):
            raise ValueError("design.json missing before implement stage")
        design = DesignArtifact.model_validate(self.artifacts.read_json(run_id, "design.json"))
        requirements: Optional[RequirementsArtifact] = None
        if arts.get("requirements.json"):
            requirements = RequirementsArtifact.model_validate(
                self.artifacts.read_json(run_id, "requirements.json")
            )
        return design, requirements, arts

    async def _implement_node(self, state: OrchestratorState) -> OrchestratorState:
        run_id = uuid.UUID(state["run_id"])
        run = await _get_run(self.session, run_id)
        run.state = WorkflowState.IMPLEMENTING
        fix_mode = bool(state.get("fix_mode"))
        await _audit(
            self.session,
            run_id,
            "implement.start",
            "Applying fixes" if fix_mode else "Salesforce Developer implementing",
        )
        try:
            design, requirements, arts = await self._load_context(run_id, run)
            file_map = dict(state.get("file_map") or {})

            if fix_mode:
                review_raw = self.artifacts.read_json(run_id, "review.json")
                review = ReviewArtifact.model_validate(review_raw)
                impl_artifact, file_map = await self.developer.apply_fixes(
                    review=review,
                    design=design,
                    requirements=requirements,
                    run_id=run_id,
                    file_map=file_map,
                )
            else:
                impl_artifact, file_map = await self.developer.implement(
                    design, requirements, run_id
                )
                impl_path = self.artifacts.write_json(run_id, "implementation.json", impl_artifact)
                arts["implementation.json"] = impl_path

            source_dir = self.artifacts.run_dir(run_id) / "force-app"
            deploy_result = await self.deployer.deploy_validate(
                source_dir=source_dir,
                jira_key=run.jira_key,
            )
            deploy_path = self.artifacts.write_json(run_id, "deploy.json", deploy_result)
            arts["deploy.json"] = deploy_path
            run.artifacts = arts
            await _audit(
                self.session,
                run_id,
                "deploy.complete",
                deploy_result.message,
                {
                    "success": deploy_result.success,
                    "status": deploy_result.status,
                    "check_only": deploy_result.check_only,
                },
            )
            if not deploy_result.success:
                raise RuntimeError(
                    f"Salesforce deploy validate failed ({deploy_result.status}): {deploy_result.message}"
                )

            state["file_map"] = file_map
            state["next_action"] = "review"
            state["fix_mode"] = False
            await self.session.flush()
            return state
        except Exception as exc:  # noqa: BLE001
            logger.exception("implement_failed", run_id=str(run_id))
            state["error"] = str(exc)
            state["failure_phase"] = "implement"
            state["next_action"] = "fail"
            await self.session.flush()
            return state

    def _after_implement(self, state: OrchestratorState) -> str:
        if state.get("next_action") == "fail":
            return "fail"
        return "review"

    async def _review_node(self, state: OrchestratorState) -> OrchestratorState:
        run_id = uuid.UUID(state["run_id"])
        run = await _get_run(self.session, run_id)
        run.state = WorkflowState.REVIEWING
        arts = dict(run.artifacts or {})
        cycle = int(arts.get("review_cycle") or state.get("review_cycle") or 0)
        await _audit(self.session, run_id, "review.start", f"Code review cycle {cycle}")
        try:
            design, requirements, _ = await self._load_context(run_id, run)
            if requirements is None:
                raise ValueError("requirements.json missing before review stage")

            impl_raw = self.artifacts.read_json(run_id, "implementation.json")
            implementation = ImplementationArtifact.model_validate(impl_raw)
            file_map = dict(state.get("file_map") or {})

            source_dir = self.artifacts.run_dir(run_id) / "force-app"
            static_result = await self.static_analysis.analyze(source_dir=source_dir)
            pmd_path = self.artifacts.write_json(run_id, "pmd.json", static_result)
            arts["pmd.json"] = pmd_path

            review = await self.reviewer.review(
                requirements=requirements,
                design=design,
                implementation=implementation,
                file_map=file_map,
                static_analysis=static_result,
                cycle=cycle,
                run_id=run_id,
            )
            arts["review.json"] = str(self.artifacts.run_dir(run_id) / "review.json")
            cycle += 1
            arts["review_cycle"] = cycle
            run.retry_count = cycle
            run.artifacts = arts

            await _audit(
                self.session,
                run_id,
                "review.complete",
                review.summary,
                {"passed": review.passed, "cycle": cycle, "findings": len(review.findings)},
            )

            max_cycles = self.settings.max_review_cycles
            if review.passed:
                state["next_action"] = "qa"
                state["review_cycle"] = cycle
                state["file_map"] = file_map
            elif review.has_fixable_blocking() and cycle < max_cycles:
                state["next_action"] = "implement"
                state["fix_mode"] = True
                state["review_cycle"] = cycle
                state["file_map"] = file_map
                await self.jira.add_comment(
                    run.jira_key,
                    (
                        "## Agent Nova code review\n"
                        f"Status: failed (cycle {cycle}/{max_cycles})\n"
                        f"Action: applying fixes for {len(review.blocking_findings())} "
                        "blocking finding(s)."
                    ),
                )
            else:
                if cycle >= max_cycles:
                    state["error"] = (
                        f"Code review failed after {max_cycles} cycles — human escalation required"
                    )
                else:
                    state["error"] = (
                        f"Code review failed with non-fixable findings: {review.summary}"
                    )
                state["failure_phase"] = "review"
                state["next_action"] = "fail"

            await self.session.flush()
            return state
        except Exception as exc:  # noqa: BLE001
            logger.exception("review_failed", run_id=str(run_id))
            state["error"] = str(exc)
            state["failure_phase"] = "review"
            state["next_action"] = "fail"
            await self.session.flush()
            return state

    def _after_review(self, state: OrchestratorState) -> str:
        action = state.get("next_action") or "fail"
        if action == "implement":
            return "implement"
        if action == "qa":
            return "qa"
        return "fail"

    async def _qa_node(self, state: OrchestratorState) -> OrchestratorState:
        run_id = uuid.UUID(state["run_id"])
        run = await _get_run(self.session, run_id)
        run.state = WorkflowState.TESTING
        await _audit(self.session, run_id, "qa.start", "QA Engineer validating")
        try:
            design, requirements, arts = await self._load_context(run_id, run)
            if requirements is None:
                raise ValueError("requirements.json missing before QA stage")

            file_map = dict(state.get("file_map") or {})
            source_dir = self.artifacts.run_dir(run_id) / "force-app"
            class_names = _apex_class_names(file_map)
            apex_result = await self.apex_runner.run_tests(
                source_dir=source_dir,
                jira_key=run.jira_key,
                class_names=class_names or None,
            )

            qa_report = await self.qa.run_qa(
                requirements=requirements,
                design=design,
                file_map=file_map,
                apex_result=apex_result,
                run_id=run_id,
            )
            arts["qa-report.json"] = str(self.artifacts.run_dir(run_id) / "qa-report.json")
            run.artifacts = arts

            await _audit(
                self.session,
                run_id,
                "qa.complete",
                qa_report.summary,
                {"passed": qa_report.passed, "scenarios": len(qa_report.test_matrix)},
            )

            if qa_report.passed:
                state["next_action"] = "create_pr"
                state["file_map"] = file_map
            else:
                state["error"] = f"QA failed: {qa_report.summary}"
                state["failure_phase"] = "qa"
                state["next_action"] = "fail"

            await self.session.flush()
            return state
        except Exception as exc:  # noqa: BLE001
            logger.exception("qa_failed", run_id=str(run_id))
            state["error"] = str(exc)
            state["failure_phase"] = "qa"
            state["next_action"] = "fail"
            await self.session.flush()
            return state

    def _after_qa(self, state: OrchestratorState) -> str:
        if state.get("next_action") == "create_pr":
            return "create_pr"
        return "fail"

    async def _create_pr_node(self, state: OrchestratorState) -> OrchestratorState:
        run_id = uuid.UUID(state["run_id"])
        run = await _get_run(self.session, run_id)
        await _audit(self.session, run_id, "create_pr.start", "Creating draft PR")
        try:
            impl_raw = self.artifacts.read_json(run_id, "implementation.json")
            impl_artifact = ImplementationArtifact.model_validate(impl_raw)
            file_map = dict(state.get("file_map") or {})

            git_files = {
                path: content
                for path, content in file_map.items()
                if path.startswith("force-app/") or path == "sfdx-project.json"
            }
            branch_name = f"agent-nova/{run.jira_key.lower()}"
            git_result = await self.git.create_draft_pr(
                jira_key=run.jira_key,
                branch_name=branch_name,
                commit_message=impl_artifact.commit_message,
                files=git_files,
            )
            run.branch_name = git_result.branch_name
            run.pr_url = git_result.pr_url
            run.error = None
            run.state = WorkflowState.PR_CREATED

            deploy_raw = self.artifacts.read_json(run_id, "deploy.json")
            deploy_status = deploy_raw.get("status", "unknown")

            await _audit(
                self.session,
                run_id,
                "implement.complete",
                "Draft PR created",
                {
                    "branch_name": git_result.branch_name,
                    "pr_url": git_result.pr_url,
                    "files": len(impl_artifact.files_written),
                    "deploy_status": deploy_status,
                },
            )
            await self.jira.add_comment(
                run.jira_key,
                (
                    "## Agent Nova draft PR created\n"
                    f"PR: {git_result.pr_url}\n"
                    f"Branch: {git_result.branch_name}\n"
                    f"Deploy validate: {deploy_status}\n"
                    "Code review: passed\n"
                    "QA: passed"
                ),
            )
            state["next_action"] = "end"
            await self.session.flush()
            return state
        except Exception as exc:  # noqa: BLE001
            logger.exception("create_pr_failed", run_id=str(run_id))
            state["error"] = str(exc)
            state["failure_phase"] = "implement"
            state["next_action"] = "fail"
            await self.session.flush()
            return state

    def _after_create_pr(self, state: OrchestratorState) -> str:
        if state.get("next_action") == "fail":
            return "fail"
        return "end"

    async def _fail_node(self, state: OrchestratorState) -> OrchestratorState:
        run_id = uuid.UUID(state["run_id"])
        run = await _get_run(self.session, run_id)
        run.state = WorkflowState.FAILED
        run.error = state.get("error") or "Unknown failure"
        await _audit(self.session, run_id, "run.failed", run.error)
        try:
            if state.get("design_decision") == "reject":
                prefix = "## Agent Nova design rejected\n"
            elif state.get("failure_phase") == "duplicate":
                original = (run.artifacts or {}).get("duplicate_of") or ""
                prefix = (
                    f"## Agent Nova closed as duplicate\n"
                    f"Original / main ticket: {original}\n"
                    if original
                    else "## Agent Nova closed as duplicate\n"
                )
            elif state.get("failure_phase") == "doc":
                refs = (run.artifacts or {}).get("doc_references") or []
                ref_lines = []
                for r in refs:
                    if isinstance(r, dict) and r.get("title"):
                        link = f" — {r['url']}" if r.get("url") else ""
                        ref_lines.append(f"- {r['title']}{link}")
                ref_block = ("\n".join(ref_lines) + "\n") if ref_lines else ""
                prefix = (
                    "## Agent Nova closed — already documented/implemented\n"
                    f"{ref_block}"
                )
            elif state.get("failure_phase") == "review":
                prefix = "## Agent Nova failed during code review\n"
            elif state.get("failure_phase") == "qa":
                prefix = "## Agent Nova failed during QA\n"
            elif state.get("failure_phase") == "implement":
                prefix = "## Agent Nova failed during implementation\n"
            else:
                prefix = "## Agent Nova failed during analysis\n"
            await self.jira.add_comment(
                run.jira_key,
                f"{prefix}Error: {run.error}\nNext step: a human should take over.",
            )
        except Exception:  # noqa: BLE001
            logger.exception("failure_comment_failed", jira_key=run.jira_key)
        await self.session.flush()
        return state


async def find_active_run(session: AsyncSession, jira_key: str) -> Optional[Run]:
    result = await session.execute(
        select(Run)
        .where(Run.jira_key == jira_key, Run.state.in_(list(ACTIVE_STATES)))
        .order_by(Run.created_at.desc())
        .limit(1)
    )
    return result.scalar_one_or_none()


async def create_run(
    session: AsyncSession,
    jira_key: str,
    idempotency_key: Optional[str] = None,
) -> Run:
    run = Run(
        jira_key=jira_key,
        state=WorkflowState.RECEIVED,
        idempotency_key=idempotency_key,
        artifacts={},
    )
    session.add(run)
    await session.flush()
    store = ArtifactStore()
    run.artifact_dir = str(store.run_dir(run.id))
    await _audit(session, run.id, "run.created", f"Run created for {jira_key}")
    await session.commit()
    await session.refresh(run)
    return run
