#!/usr/bin/env python3
"""Manage nAvid batch queue coordination without duplicating item run state."""

from __future__ import annotations

import argparse
from datetime import datetime, timezone
import json
import os
from pathlib import Path
import subprocess
import sys
import tempfile
from typing import Any


SCHEMA_VERSION = 1
ITEM_STATUSES = ("pending", "running", "complete", "blocked", "failed", "skipped")
RUN_STATE_SCRIPT = Path(__file__).with_name("run_state.py")


def now_iso() -> str:
    return datetime.now(timezone.utc).replace(microsecond=0).isoformat()


def queue_path(queue_dir: str) -> Path:
    return Path(queue_dir).resolve() / "queue.json"


def output(value: dict[str, Any]) -> None:
    print(json.dumps(value, ensure_ascii=False, indent=2))


def emit_event(queue: dict[str, Any], event_type: str, **data: Any) -> None:
    event: dict[str, Any] = {"type": event_type, "timestamp": now_iso()}
    event.update(data)
    queue["events"].append(event)


def parse_item(raw: str) -> tuple[str, Path]:
    if "::" not in raw:
        raise ValueError("item must be SOURCE_URL::PROJECT_DIR")
    source_url, project_dir = raw.split("::", 1)
    source_url = source_url.strip()
    if not source_url or not project_dir.strip():
        raise ValueError("item source URL and project directory are required")
    return source_url, Path(project_dir).resolve()


def read_queue(queue_dir: str) -> dict[str, Any]:
    path = queue_path(queue_dir)
    if not path.is_file():
        raise FileNotFoundError(f"queue manifest not found: {path}")
    with path.open("r", encoding="utf-8") as handle:
        queue = json.load(handle)
    if queue.get("schema_version") != SCHEMA_VERSION:
        raise ValueError(f"unsupported schema_version: {queue.get('schema_version')}")
    return queue


def write_queue(queue_dir: str, queue: dict[str, Any]) -> None:
    path = queue_path(queue_dir)
    path.parent.mkdir(parents=True, exist_ok=True)
    queue["updated_at"] = now_iso()
    fd, temp_name = tempfile.mkstemp(prefix="queue-", suffix=".tmp", dir=str(path.parent))
    try:
        with os.fdopen(fd, "w", encoding="utf-8", newline="\n") as handle:
            json.dump(queue, handle, ensure_ascii=False, indent=2)
            handle.write("\n")
        os.replace(temp_name, path)
    finally:
        if os.path.exists(temp_name):
            os.unlink(temp_name)


def get_item(queue: dict[str, Any], item_id: str) -> dict[str, Any]:
    for item in queue["items"]:
        if item["id"] == item_id:
            return item
    raise ValueError(f"unknown item: {item_id}")


def run_resume_check(project_dir: str) -> dict[str, Any]:
    result = subprocess.run(
        [sys.executable, "-X", "utf8", str(RUN_STATE_SCRIPT), "resume-check",
         "--project-dir", project_dir],
        check=False,
        capture_output=True,
        text=True,
        encoding="utf-8",
    )
    payload = json.loads(result.stdout)
    if result.returncode != 0:
        raise FileNotFoundError(payload.get("error", f"unable to inspect {project_dir}"))
    state_file = Path(project_dir).resolve() / "run-state.json"
    with state_file.open("r", encoding="utf-8") as handle:
        item_state = json.load(handle)
    payload["template_selection"] = item_state.get("template_selection")
    payload["intake"] = item_state.get("intake", {})
    return payload


def artifact_path(project_dir: str, raw_path: str | None) -> str | None:
    if not raw_path:
        return None
    path = Path(raw_path)
    return str(path if path.is_absolute() else Path(project_dir).resolve() / path)


def fallback_from_stages(stages: dict[str, Any]) -> tuple[str | None, str | None]:
    for stage_name, record in stages.items():
        fallback = record.get("fallback")
        if fallback:
            state = fallback.get("status", "pending_approval")
            description = f"{stage_name}: {fallback.get('original', '-')} -> {fallback.get('proposal', '-')} ({state})"
            return description, fallback.get("reason") or fallback.get("impact")
    return None, None


def sync_item(queue: dict[str, Any], item: dict[str, Any]) -> dict[str, Any]:
    state = run_resume_check(item["project_dir"])
    resume = state["resume"]
    stages = state["stages"]
    terminal = next(
        ((name, record) for name, record in stages.items()
         if record.get("status") in ("blocked", "failed")),
        None,
    )
    template = state.get("template_selection")
    fallback_disclosure, blocked_reason = fallback_from_stages(stages)
    item["resume_stage"] = terminal[0] if terminal else resume.get("stage")
    item["qa_status"] = stages["qa"].get("status")
    item["template_override"] = template
    item["voice_disclosure"] = stages["voice"].get("metadata") or None
    item["fallback_disclosure"] = fallback_disclosure
    item["output_path"] = artifact_path(
        item["project_dir"],
        (stages["render"].get("artifacts") or [{}])[0].get("path"),
    )
    publish_path = Path(item["project_dir"]).resolve() / "PUBLISH.md"
    item["publish_path"] = str(publish_path) if publish_path.is_file() else None
    item["error"] = None
    item["blocked_reason"] = None

    if template and template.get("tier") != "approved":
        item["status"] = "blocked"
        item["blocked_reason"] = "selected template is not approved"
    elif terminal and terminal[1]["status"] == "blocked":
        item["status"] = "blocked"
        item["blocked_reason"] = blocked_reason or f"{terminal[0]} awaits approval"
    elif terminal and terminal[1]["status"] == "failed":
        item["status"] = "failed"
        attempts = terminal[1].get("attempts") or []
        item["error"] = attempts[-1].get("error") if attempts else f"{terminal[0]} failed"
    elif resume["status"] == "complete":
        item["status"] = "complete"
    elif resume["status"] == "blocked":
        item["status"] = "blocked"
        item["blocked_reason"] = blocked_reason or f"{resume['stage']} awaits approval"
    elif resume["status"] == "failed":
        item["status"] = "failed"
        attempts = stages[resume["stage"]].get("attempts") or []
        item["error"] = attempts[-1].get("error") if attempts else f"{resume['stage']} failed"
    else:
        item["status"] = "pending"
    emit_event(queue, "item_synchronized", item_id=item["id"], status=item["status"],
               resume_stage=item["resume_stage"])
    return item


def item_is_ready(item: dict[str, Any]) -> bool:
    selection = item.get("template_override")
    return (
        item["status"] == "complete"
        and item.get("qa_status") == "complete"
        and bool(item.get("output_path"))
        and bool(item.get("publish_path"))
        and (selection is None or selection.get("tier") == "approved")
        and not item.get("fallback_disclosure")
    )


def compact_value(value: Any) -> str:
    if value in (None, "", {}):
        return "-"
    if isinstance(value, dict):
        if "id" in value:
            suffix = f" ({value.get('tier', '-')})"
            return f"{value['id']}{suffix}"
        return ", ".join(f"{key}={text}" for key, text in sorted(value.items()))
    return str(value).replace("|", "/").replace("\n", " ")


def publish_sections(path: str | None) -> dict[str, str]:
    if not path or not Path(path).is_file():
        return {}
    sections: dict[str, str] = {}
    current: str | None = None
    values: list[str] = []
    for line in Path(path).read_text(encoding="utf-8").splitlines():
        if line.startswith("## "):
            if current:
                sections[current] = " ".join(part.strip() for part in values if part.strip())
            current = line[3:].strip()
            values = []
        elif current:
            values.append(line)
    if current:
        sections[current] = " ".join(part.strip() for part in values if part.strip())
    return sections


def summary_text(queue: dict[str, Any]) -> str:
    statuses = {status: 0 for status in ITEM_STATUSES}
    for item in queue["items"]:
        statuses[item["status"]] += 1
    intake = queue["intake"]
    mode = queue["mode"]
    lines = [
        f"# Batch Summary: {queue['batch']['id']}",
        "",
        "## Shared Intake",
        "",
        f"- Profile: `{intake['profile']}`",
        f"- Language: `{intake['language']}`",
        f"- Template strategy: `{intake['template_strategy']}`",
        f"- Template set: `{intake.get('template_set') or '-'}`",
        f"- Mode: `{mode['execution']}` / `{mode['reporting']}`",
        "",
        "## Counts",
        "",
        "| pending | running | complete | blocked | failed | skipped |",
        "| ---: | ---: | ---: | ---: | ---: | ---: |",
        "| " + " | ".join(str(statuses[status]) for status in ITEM_STATUSES) + " |",
        "",
        "## Items",
        "",
        "| # | Item | Status | Output | QA | Template | Voice | Fallback | Detail |",
        "| ---: | --- | --- | --- | --- | --- | --- | --- | --- |",
    ]
    for item in queue["items"]:
        detail = item.get("blocked_reason") or item.get("error") or "-"
        lines.append(
            "| {order} | `{item_id}` `{state}` | `{status}` | `{output}` | `{qa}` | `{template}` | `{voice}` | `{fallback}` | {detail} |".format(
                order=item["order"],
                item_id=item["id"],
                state=compact_value(item["state_path"]),
                status=item["status"],
                output=compact_value(item.get("output_path")),
                qa=compact_value(item.get("qa_status")),
                template=compact_value(item.get("template_override") or queue["intake"].get("template_set")),
                voice=compact_value(item.get("voice_disclosure")),
                fallback=compact_value(item.get("fallback_disclosure")),
                detail=compact_value(detail),
            )
        )
    ready = [item for item in queue["items"] if item_is_ready(item)]
    lines += ["", "## Ready to publish", ""]
    if ready:
        for item in ready:
            publish = publish_sections(item.get("publish_path"))
            lines.append(f"- `{item['id']}`: `{item['output_path']}`")
            lines.append(f"  - Publish Title: {publish.get('Publish Title', '-')}")
            lines.append(f"  - Hashtags: {publish.get('Hashtags', '-')}")
    else:
        lines.append("- None.")
    for heading, status in (("Blocked", "blocked"), ("Failed", "failed"), ("Skipped", "skipped")):
        entries = [item for item in queue["items"] if item["status"] == status]
        lines += ["", f"## {heading}", ""]
        if entries:
            for item in entries:
                detail = item.get("blocked_reason") or item.get("error") or "recorded"
                lines.append(f"- `{item['id']}`: {detail}")
        else:
            lines.append("- None.")
    return "\n".join(lines) + "\n"


def write_summary(queue_dir: str, text: str) -> Path:
    path = Path(queue_dir).resolve() / "SUMMARY.md"
    path.parent.mkdir(parents=True, exist_ok=True)
    fd, temp_name = tempfile.mkstemp(prefix="summary-", suffix=".tmp", dir=str(path.parent))
    try:
        with os.fdopen(fd, "w", encoding="utf-8", newline="\n") as handle:
            handle.write(text)
        os.replace(temp_name, path)
    finally:
        if os.path.exists(temp_name):
            os.unlink(temp_name)
    return path


def new_queue(args: argparse.Namespace) -> dict[str, Any]:
    created = now_iso()
    queue_dir = Path(args.queue_dir).resolve()
    queue = {
        "schema_version": SCHEMA_VERSION,
        "batch": {"id": args.batch_id or queue_dir.name, "path": str(queue_dir)},
        "created_at": created,
        "updated_at": created,
        "intake": {
            "profile": args.profile,
            "language": args.language,
            "template_strategy": args.template_strategy,
            "template_set": args.template_set,
        },
        "mode": {"execution": args.mode, "reporting": args.reporting},
        "items": [],
        "events": [],
    }
    seen: set[str] = set()
    duplicate_sources: list[str] = []
    for raw_item in args.item:
        source_url, project_dir = parse_item(raw_item)
        if source_url in seen:
            duplicate_sources.append(source_url)
            continue
        seen.add(source_url)
        order = len(queue["items"]) + 1
        queue["items"].append(
            {
                "id": f"item-{order:03d}",
                "order": order,
                "source_url": source_url,
                "project_dir": str(project_dir),
                "state_path": str(project_dir / "run-state.json"),
                "status": "pending",
                "template_override": None,
                "output_path": None,
                "qa_status": None,
                "voice_disclosure": None,
                "fallback_disclosure": None,
                "error": None,
                "blocked_reason": None,
                "resume_stage": "source",
            }
        )
    emit_event(queue, "queue_initialized", retained_items=len(queue["items"]))
    if duplicate_sources:
        emit_event(queue, "exact_duplicates_suppressed", sources=duplicate_sources)
    return queue


def cmd_init(args: argparse.Namespace) -> dict[str, Any]:
    path = queue_path(args.queue_dir)
    if path.exists() and not args.force:
        raise FileExistsError(f"queue manifest already exists: {path}")
    queue = new_queue(args)
    write_queue(args.queue_dir, queue)
    return queue


def cmd_inspect(args: argparse.Namespace) -> dict[str, Any]:
    return read_queue(args.queue_dir)


def cmd_set_status(args: argparse.Namespace) -> dict[str, Any]:
    queue = read_queue(args.queue_dir)
    item = get_item(queue, args.item_id)
    item["status"] = args.status
    if args.error is not None:
        item["error"] = args.error
    if args.blocked_reason is not None:
        item["blocked_reason"] = args.blocked_reason
    emit_event(queue, "item_status_recorded", item_id=item["id"], status=args.status)
    write_queue(args.queue_dir, queue)
    return {"batch": queue["batch"], "item": item, "events": queue["events"]}


def cmd_sync_item(args: argparse.Namespace) -> dict[str, Any]:
    queue = read_queue(args.queue_dir)
    item = sync_item(queue, get_item(queue, args.item_id))
    write_queue(args.queue_dir, queue)
    return {"batch": queue["batch"], "item": item, "events": queue["events"]}


def cmd_next_item(args: argparse.Namespace) -> dict[str, Any]:
    queue = read_queue(args.queue_dir)
    selected = None
    for item in queue["items"]:
        if item["status"] == "skipped":
            continue
        if Path(item["state_path"]).is_file():
            sync_item(queue, item)
        if item["status"] in ("complete", "blocked", "failed", "skipped"):
            continue
        selected = item
        break
    emit_event(queue, "next_item_selected", item_id=selected["id"] if selected else None)
    write_queue(args.queue_dir, queue)
    return {"batch": queue["batch"], "item": selected, "items": queue["items"]}


def cmd_retry_item(args: argparse.Namespace) -> dict[str, Any]:
    queue = read_queue(args.queue_dir)
    item = get_item(queue, args.item_id)
    if item["status"] != "failed":
        raise ValueError("selected retry is allowed only for a failed item")
    item["status"] = "pending"
    item["error"] = None
    emit_event(queue, "item_retry_selected", item_id=item["id"])
    write_queue(args.queue_dir, queue)
    return {"batch": queue["batch"], "item": item, "events": queue["events"]}


def cmd_skip_item(args: argparse.Namespace) -> dict[str, Any]:
    queue = read_queue(args.queue_dir)
    item = get_item(queue, args.item_id)
    item["status"] = "skipped"
    item["blocked_reason"] = args.reason
    emit_event(queue, "item_skipped", item_id=item["id"], reason=args.reason)
    write_queue(args.queue_dir, queue)
    return {"batch": queue["batch"], "item": item, "events": queue["events"]}


def cmd_summary(args: argparse.Namespace) -> dict[str, Any]:
    queue = read_queue(args.queue_dir)
    if args.refresh:
        for item in queue["items"]:
            if item["status"] != "skipped" and Path(item["state_path"]).is_file():
                sync_item(queue, item)
    text = summary_text(queue)
    path = write_summary(args.queue_dir, text)
    emit_event(queue, "summary_written", path=str(path))
    write_queue(args.queue_dir, queue)
    return {
        "batch": queue["batch"],
        "summary_path": str(path),
        "ready_to_publish": [item["id"] for item in queue["items"] if item_is_ready(item)],
        "items": queue["items"],
    }


def build_parser() -> argparse.ArgumentParser:
    parser = argparse.ArgumentParser(description=__doc__)
    commands = parser.add_subparsers(dest="command", required=True)

    init = commands.add_parser("init", help="Create queue.json from resolved batch intake")
    init.add_argument("--queue-dir", required=True)
    init.add_argument("--batch-id")
    init.add_argument("--item", action="append", required=True, help="SOURCE_URL::PROJECT_DIR")
    init.add_argument("--profile", required=True)
    init.add_argument("--language", required=True)
    init.add_argument("--template-strategy", required=True)
    init.add_argument("--template-set")
    init.add_argument("--mode", choices=("auto-run", "review-first"), default="auto-run")
    init.add_argument("--reporting", choices=("default", "quiet", "verbose"), default="default")
    init.add_argument("--force", action="store_true")
    init.set_defaults(handler=cmd_init)

    inspect = commands.add_parser("inspect", help="Inspect queue.json")
    inspect.add_argument("--queue-dir", required=True)
    inspect.set_defaults(handler=cmd_inspect)

    status = commands.add_parser("set-status", help="Record aggregate coordination status")
    status.add_argument("--queue-dir", required=True)
    status.add_argument("--item-id", required=True)
    status.add_argument("--status", choices=ITEM_STATUSES, required=True)
    status.add_argument("--error")
    status.add_argument("--blocked-reason")
    status.set_defaults(handler=cmd_set_status)

    sync = commands.add_parser("sync-item", help="Synchronize one item from linked run-state.json")
    sync.add_argument("--queue-dir", required=True)
    sync.add_argument("--item-id", required=True)
    sync.set_defaults(handler=cmd_sync_item)

    next_item = commands.add_parser("next-item", help="Return the next actionable sequential item")
    next_item.add_argument("--queue-dir", required=True)
    next_item.set_defaults(handler=cmd_next_item)

    retry = commands.add_parser("retry-item", help="Select one failed item for retry")
    retry.add_argument("--queue-dir", required=True)
    retry.add_argument("--item-id", required=True)
    retry.set_defaults(handler=cmd_retry_item)

    skip = commands.add_parser("skip-item", help="Explicitly skip an item")
    skip.add_argument("--queue-dir", required=True)
    skip.add_argument("--item-id", required=True)
    skip.add_argument("--reason", required=True)
    skip.set_defaults(handler=cmd_skip_item)

    summary = commands.add_parser("summary", help="Write operator-facing SUMMARY.md")
    summary.add_argument("--queue-dir", required=True)
    summary.add_argument("--refresh", action="store_true")
    summary.set_defaults(handler=cmd_summary)
    return parser


def main() -> int:
    parser = build_parser()
    args = parser.parse_args()
    try:
        result = args.handler(args)
        output(result)
        return 0
    except (ValueError, FileNotFoundError, FileExistsError, KeyError, json.JSONDecodeError) as exc:
        output({"status": "error", "error": str(exc)})
        return 1


if __name__ == "__main__":
    sys.exit(main())
