#!/usr/bin/env python3
"""boss_byom.py — 用你自己的模型跑完 boss 判断流水线 (BYOM 三连), 一条命令。

把「prepare → 跑 synthesis → 逐评委打分 → submit → render」压成一次调用。
**零第三方依赖** (只用标准库), 拷进你自己的项目随便改。

用法
----
    export BOSS_SVC_TOKEN=...          # 向服务管理员申请
    export ANTHROPIC_API_KEY=...       # 或 OPENAI_API_KEY (见下)
    python3 boss_byom.py "评估 X 产品线明年是否扩张"

    # 评议一份现成文档
    python3 boss_byom.py --doc plan.md

    # 指定场景 / 输出文件 / 并发
    python3 boss_byom.py "议题" --scene default -o report.md --concurrency 3

环境变量
--------
    BOSS_SVC_URL       服务地址 (默认 https://svc.apex.fan)
    BOSS_SVC_TOKEN     服务 token (必填)

    模型二选一:
    ANTHROPIC_API_KEY  + 可选 ANTHROPIC_MODEL (默认 claude-sonnet-5)
    OPENAI_API_KEY     + 可选 OPENAI_BASE_URL (默认 https://api.openai.com/v1)
                       + 可选 OPENAI_MODEL

流程说明
--------
boss 出题与纪律 (context / prompt 包 / 打分契约 / 确定性聚合),
**你的模型出判断** —— 模型选型与费用都在你手里, prompt 包不会离开你的进程。

注意: judges[].user_template 与 merge_task.user_template 里的哨兵
<<SYNTHESIS_MD>> / <<REVIEWS_MD>> 必须用**字符串替换**, 不能用 str.format —
正文里可能含花括号。
"""

from __future__ import annotations

import argparse
import concurrent.futures
import json
import os
import sys
import urllib.error
import urllib.request

SYNTHESIS_SENTINEL = "<<SYNTHESIS_MD>>"
REVIEWS_SENTINEL = "<<REVIEWS_MD>>"

DEFAULT_BOSS_URL = "https://svc.apex.fan"
DEFAULT_ANTHROPIC_MODEL = "claude-sonnet-5"
DEFAULT_OPENAI_MODEL = "gpt-4o"


# ── 通用 HTTP ────────────────────────────────────────────────────────────

def _post_json(url: str, payload: dict, headers: dict, timeout: int = 300) -> dict:
    body = json.dumps(payload, ensure_ascii=False).encode("utf-8")
    req = urllib.request.Request(url, data=body, method="POST",
                                 headers={"Content-Type": "application/json", **headers})
    try:
        with urllib.request.urlopen(req, timeout=timeout) as r:
            return json.loads(r.read().decode("utf-8"))
    except urllib.error.HTTPError as e:
        detail = e.read().decode("utf-8", "replace")[:800]
        raise SystemExit(f"✗ HTTP {e.code} ← {url}\n  {detail}\n{_hint(e.code)}") from None
    except urllib.error.URLError as e:
        raise SystemExit(f"✗ 连不上 {url}: {e.reason}") from None


def _hint(code: int) -> str:
    return {
        401: "  提示: token 不对。确认用的是 BOSS_SVC_TOKEN (与 BOSS_API_TOKEN 不通用)。",
        404: "  提示: run 可能已过期 (TTL 24 小时), 或路径写错。",
        422: "  提示: review 校验未过 —— 全批原子, 看上面 errors 逐条明细。"
             "维度评委必填 adversarial_view 三字段。",
        502: "  提示: 网关连不上源站, 多半是服务重启的几秒空窗, 等 10 秒重试。",
        503: "  提示: 服务端未配 token (fail-close) —— 不是你的问题, 找部署方。",
    }.get(code, "")


# ── 模型调用 (Anthropic / OpenAI 兼容) ───────────────────────────────────

def call_model(system: str, user: str, *, max_tokens: int = 4096) -> str:
    """调你自己的模型。想换别的供应商, 改这一个函数即可。"""
    if os.environ.get("ANTHROPIC_API_KEY"):
        out = _post_json(
            "https://api.anthropic.com/v1/messages",
            {"model": os.environ.get("ANTHROPIC_MODEL", DEFAULT_ANTHROPIC_MODEL),
             "max_tokens": max_tokens,
             "system": system,
             "messages": [{"role": "user", "content": user}]},
            {"x-api-key": os.environ["ANTHROPIC_API_KEY"],
             "anthropic-version": "2023-06-01"})
        return "".join(b.get("text", "") for b in out.get("content", []))

    if os.environ.get("OPENAI_API_KEY"):
        base = os.environ.get("OPENAI_BASE_URL", "https://api.openai.com/v1").rstrip("/")
        out = _post_json(
            f"{base}/chat/completions",
            {"model": os.environ.get("OPENAI_MODEL", DEFAULT_OPENAI_MODEL),
             "max_tokens": max_tokens,
             "messages": [{"role": "system", "content": system},
                          {"role": "user", "content": user}]},
            {"Authorization": f"Bearer {os.environ['OPENAI_API_KEY']}"})
        return out["choices"][0]["message"]["content"]

    raise SystemExit("✗ 没配模型: 需要 ANTHROPIC_API_KEY 或 OPENAI_API_KEY")


# ── 三连 ─────────────────────────────────────────────────────────────────

def run(topic: str = "", *, doc: tuple[str, str] | None = None,
        scene: str = "default", concurrency: int = 3,
        model=call_model, log=lambda m: print(m, file=sys.stderr)) -> str:
    """跑完整闭环, 返回 report_md。model/log 可注入, 便于测试与改造。"""
    boss = os.environ.get("BOSS_SVC_URL", DEFAULT_BOSS_URL).rstrip("/")
    token = os.environ.get("BOSS_SVC_TOKEN")
    if not token:
        raise SystemExit("✗ 缺 BOSS_SVC_TOKEN")
    auth = {"Authorization": f"Bearer {token}"}

    # 1) prepare — 拿冻结 context + 全套 prompt 包
    payload: dict = {"scene": scene, "caller": "boss_byom.py"}
    if topic:
        payload["topic"] = topic
    if doc:
        payload["review_doc"] = {"name": doc[0], "text": doc[1]}
    log("① prepare …")
    pack = _post_json(f"{boss}/v1/prepare", payload, auth)
    run_id, judges = pack["run_id"], pack["judges"]
    log(f"   run_id={run_id} · panel={pack['panel_name']} · {len(judges)} 位评委")

    # 2) synthesis — user_template 可直接用 (没有哨兵)
    log("② synthesis …")
    synthesis_md = model(pack["synthesis_task"]["system"],
                         pack["synthesis_task"]["user_template"], max_tokens=6144)

    # 3) 逐评委打分 (互不可见 —— 独立打分是 panel 设计的核心, 别把别人的 review 喂进去)
    log(f"③ {len(judges)} 位评委打分 (并发 {concurrency}) …")

    def one(j: dict) -> tuple[str, str]:
        user = j["user_template"].replace(SYNTHESIS_SENTINEL, synthesis_md)
        md = model(j["system"], user, max_tokens=4096)
        log(f"   ✓ {j['slug']} ({j['category']})")
        return j["slug"], md

    with concurrent.futures.ThreadPoolExecutor(max_workers=max(1, concurrency)) as ex:
        reviews = list(ex.map(one, judges))

    # 4) submit — 全批原子: 任一份不合格, 整批 422
    log("④ submit …")
    res = _post_json(f"{boss}/v1/submit",
                     {"run_id": run_id,
                      "reviews": [{"judge": s, "review_md": m} for s, m in reviews]},
                     auth)
    log(f"   accepted: {', '.join(res['accepted'])}")

    # 5) merge — 可选; 不跑也能 render (那就只有确定性聚合, 没有行文)
    log("⑤ merge …")
    merge_user = (pack["merge_task"]["user_template"]
                  .replace(SYNTHESIS_SENTINEL, synthesis_md)
                  .replace(REVIEWS_SENTINEL,
                           "\n\n---\n\n".join(f"## {s}\n\n{m}" for s, m in reviews)))
    body_prose = model(pack["merge_task"]["system"], merge_user, max_tokens=6144)

    # 6) render — 确定性聚合出富报告
    log("⑥ render …")
    out = _post_json(f"{boss}/v1/render",
                     {"run_id": run_id, "body_prose": body_prose}, auth)
    summary = out.get("panel_summary") or {}
    if summary:
        log(f"   panel_summary: {json.dumps(summary, ensure_ascii=False)[:200]}")
    return out["report_md"]


def main(argv: list[str] | None = None) -> int:
    ap = argparse.ArgumentParser(description="用你自己的模型跑完 boss 判断流水线")
    ap.add_argument("topic", nargs="?", default="", help="议题 (与 --doc 至少给一个)")
    ap.add_argument("--doc", help="评议一份现成文档 (文件路径)")
    ap.add_argument("--scene", default="default", help="场景 slug (默认 default)")
    ap.add_argument("-o", "--out", help="报告写到文件 (默认打到 stdout)")
    ap.add_argument("--concurrency", type=int, default=3, help="评委并发数 (默认 3)")
    a = ap.parse_args(argv)

    if not a.topic and not a.doc:
        ap.error("议题与 --doc 至少给一个")

    doc = None
    if a.doc:
        with open(a.doc, encoding="utf-8") as f:
            doc = (os.path.basename(a.doc), f.read())

    report = run(a.topic, doc=doc, scene=a.scene, concurrency=a.concurrency)

    if a.out:
        with open(a.out, "w", encoding="utf-8") as f:
            f.write(report)
        print(f"✓ 报告已写入 {a.out}", file=sys.stderr)
    else:
        print(report)
    return 0


if __name__ == "__main__":
    raise SystemExit(main())
