# -*- coding: utf-8 -*-
"""正式巡检编排：读基础表 → 按站点调用独立浏览器抓取 → 写抓取表(挂关联/基准/差异/简报) → 写汇总表 → 更新基础表最新简报。
用法: /home/user/pwenv/bin/python run_round.py [--dry]  (--dry 只抓取与计算，不写飞书)
"""
import json, os, re, subprocess, sys, time, random, hashlib
from datetime import datetime, timezone, timedelta

sys.path.insert(0, "/home/user/monitor")
from browser import SiteBrowsers, ensure_location, glow_text
from extract import extract_product
from sites import SITE_CONFIGS, CLEANERS

BASE = "MSW3bNqVxalMXJsNwkXcjTfwnid"
T_BASE, T_SCRAPE, T_SUM = "tbl0GbALr9ULYo94", "tblAZV708RHQ8ZPb", "tblc1FU2ivPJU7m9"
ENV = dict(os.environ, LARKSUITE_CLI_NO_UPDATE_NOTIFIER="1")
CST = timezone(timedelta(hours=8))
VERSION = "v2026-09-14.1"
DRY = "--dry" in sys.argv
SLOTS = [(9,0),(11,0),(13,30),(17,30),(21,0)]

def lark(*args):
    r = subprocess.run(["lark-cli", *args], capture_output=True, text=True, env=ENV, cwd="/home/user")
    try:
        return json.loads(r.stdout)
    except Exception:
        return {"ok": False, "error": {"message": r.stdout[:200] + r.stderr[:200]}}

def read_table(tid):
    d = lark("base","+record-list","--base-token",BASE,"--table-id",tid,"--as","user","--limit","200","--format","json")
    assert d.get("ok"), d
    names = d["data"]["fields"]; rows = d["data"]["data"]; ids = d["data"]["record_id_list"]
    out = []
    for rid, row in zip(ids, rows):
        rec = {"_id": rid}
        for n, v in zip(names, row):
            if isinstance(v, list) and v and isinstance(v[0], dict) and set(v[0]) <= {"id","text","link"}:
                v = [x.get("id") or x.get("text") or x.get("link") for x in v]
            if isinstance(v, str):
                v = re.sub(r"\[([^\]]*)\]\(([^)]*)\)", r"\2", v)   # 剥 markdown 链接
            rec[n] = v
        out.append(rec)
    return out

def field_names(tid):
    d = lark("base","+field-list","--base-token",BASE,"--table-id",tid,"--as","user","--format","json")
    assert d.get("ok"), d
    return {f["name"] for f in d["data"]["fields"]}

def now_cst(): return datetime.now(CST)

def slot_of(dt):
    mins = dt.hour*60+dt.minute
    best, bd = None, 9999
    for h,m in SLOTS:
        d_ = abs(mins-(h*60+m))
        if d_ < bd: best, bd = f"{h}:{m:02d}", d_
    return best if bd <= 45 else "其他"

def gen_id(site):
    ms = int(time.time()*1000)
    h = hashlib.md5(f"{site}-{ms}".encode()).hexdigest()[:8]
    return f"{site}-9331-{ms}-{h}"

VAT = {"UK": 1.20, "DE": 1.19}

def tax_split(price, site):
    incl = round(price, 2) if price is not None else None
    excl = round(price / VAT[site], 2) if price is not None else None
    return incl, excl

# ---------- 评论星级（点击直方图进入 product-reviews，登录墙则放弃） ----------
def try_review_counts(page, site):
    """返回 (评论总数, {1..5星}, 状态, 范围说明)"""
    try:
        links = page.locator("#histogramTable a[href*='filterByStar']")
        if links.count() == 0:
            return None, None, "未采集", "商品页无直方图链接"
        href = links.first.get_attribute("href")
        m = re.search(r"(/product-reviews/[A-Z0-9]{10})", href)
        if not m:
            return None, None, "未采集", f"非product-reviews链接:{href[:60]}"
        path = m.group(1)
        base_url = f"https://www.amazon.{'co.uk' if site=='UK' else 'de'}{path}"
        star_names = {"five_star":"5星评论数","four_star":"4星评论数","three_star":"3星评论数","two_star":"2星评论数","one_star":"1星评论数"}
        counts, total = {}, None
        for star in ["five_star","four_star","three_star","two_star","one_star"]:
            page.goto(f"{base_url}?filterByStar={star}&reviewerType=all_reviews", wait_until="domcontentloaded", timeout=45000)
            time.sleep(random.uniform(1.8, 2.8))
            title = (page.title() or "")
            if re.search(r"Sign in|Anmelden", title):
                return None, None, "未采集", "product-reviews 登录墙，匿名不可得"
            fi = page.locator('[data-hook="cr-filter-info-review-rating-count"]')
            if fi.count():
                txt = " ".join(fi.first.inner_text().split())
                nums = re.findall(r"[\d.,]+", txt)
                # 第一个数=总ratings，第二个数=本筛选评论数（英文）；德文只有一个数
                n = CLEANERS[site]["count"](nums[1] if len(nums) > 1 else (nums[0] if nums else None))
                counts[star_names[star]] = n
        page.goto(f"{base_url}?reviewerType=all_reviews", wait_until="domcontentloaded", timeout=45000)
        time.sleep(2)
        fi = page.locator('[data-hook="cr-filter-info-review-rating-count"]')
        if fi.count():
            txt = " ".join(fi.first.inner_text().split())
            nums = re.findall(r"[\d.,]+", txt)
            total = CLEANERS[site]["count"](nums[1] if len(nums) > 1 else (nums[0] if nums else None))
        if len(counts) == 5:
            return total, counts, "成功", f"{path} filterByStar one..five DOM matching; sum={sum(v for v in counts.values() if v)}"
        return total, counts, "部分覆盖", f"{path} filterByStar 部分星级未取到"
    except Exception as e:
        return None, None, "失败", f"{type(e).__name__}: {e}"

# ---------- 差异计算（脚本，不猜测） ----------
DIFF_FIELDS = ["标题","当前含税价格","当前不含税价格","历史参考价","建议零售价","单位价格","综合评分","评分人数",
               "评论总数","1星评论数","2星评论数","3星评论数","4星评论数","5星评论数","可售状态","Buy Box状态",
               "当前销售店铺","购买按钮状态","限时优惠标识","优惠券金额","优惠券折扣百分比","Coupon及促销原文",
               "主图URL","副图URL列表","A+图片URL列表","视频展示数量","变体清单","变体数量","所属类目路径",
               "大类排名","小类排名","卖点列表","父ASIN","实际ASIN"]
NUM_FIELDS = {"当前含税价格","当前不含税价格","历史参考价","建议零售价","单位价格","综合评分","评分人数","评论总数",
              "1星评论数","2星评论数","3星评论数","4星评论数","5星评论数","优惠券金额","优惠券折扣百分比","变体数量","视频展示数量"}

def diff_records(base_rec, new_rec):
    changed, unchanged, incomparable = [], [], []
    for f in DIFF_FIELDS:
        b, n = base_rec.get(f), new_rec.get(f)
        if b is None and n is None:
            continue
        if b is None or n is None:
            incomparable.append(f"{f}:{'基准缺失' if b is None else '本次缺失'}")
            continue
        if f in NUM_FIELDS:
            try:
                if abs(float(b) - float(n)) < 1e-9:
                    unchanged.append(f); continue
                delta = float(n) - float(b)
                changed.append((f, b, n, f"{fmt_v(b)} → {fmt_v(n)}（Δ{delta:+g}）"))
            except (TypeError, ValueError):
                incomparable.append(f"{f}:非数值")
        elif isinstance(b, list) or isinstance(n, list):
            bl, nl = b if isinstance(b, list) else [b], n if isinstance(n, list) else [n]
            if bl == nl:
                unchanged.append(f)
            else:
                add = [x for x in nl if x not in bl]; rem = [x for x in bl if x not in nl]
                seq = "顺序变化" if sorted(map(str,bl))==sorted(map(str,nl)) else ""
                changed.append((f, len(bl), len(n), f"{len(bl)}项 → {len(n)}项（新增{len(add)}/消失{len(rem)}{('，'+seq) if seq else ''}）"))
        else:
            if str(b) == str(n):
                unchanged.append(f)
            else:
                changed.append((f, b, n, f"{str(b)[:40]} → {str(n)[:40]}"))
    return changed, unchanged, incomparable

def fmt_v(v):
    return str(v)

def eval_status(changed, incomparable, base_exists, new_fields):
    if not base_exists:
        return "建立基线"
    if new_fields.get("页面访问状态") != "正常商品页":
        return "异常"
    key = [c[0] for c in changed]
    if "可售状态" in key or "Buy Box状态" in key or "当前销售店铺" in key or "购买按钮状态" in key:
        return "需关注"
    for f, b, n, _ in changed:
        if f in ("当前含税价格","当前不含税价格","历史参考价") and b is not None:
            try:
                if abs(float(n)-float(b))/float(b) >= 0.03: return "需关注"
            except Exception: pass
    if "综合评分" in key:
        for f,b,n,_ in changed:
            if f=="综合评分":
                try:
                    if abs(float(n)-float(b)) >= 0.1: return "需关注"
                except Exception: pass
    if len(incomparable) >= 8:
        return "数据不足"
    return "有变化" if changed else "无变化"

def build_diff_summary(base_id, base_time, changed, unchanged, incomparable, result):
    parts = [f"对比基准：{base_id or '（本轮建立基线）'}/{base_time or '-'}；结果：{result}；"]
    for f,b,n,desc in changed:
        parts.append(f"【{f}】上次：{fmt_v(b)[:60]}｜本次：{fmt_v(n)[:60]}｜差异：{desc}；")
    parts.append(f"未变化字段：{'、'.join(unchanged) or '无'}；")
    parts.append(f"无法比较字段及原因：{'；'.join(incomparable) or '无'}")
    return "".join(parts)

def build_brief(name, site, asin, fields, changed, unchanged, base_id, base_time, result, rev_scope, cfg):
    t = now_cst().strftime("%Y-%m-%d %H:%M")
    title = (fields.get("标题") or "")[:50]
    head = f"======= 监控简报 =======\n一句话结论：{name}｜{site}｜{asin}｜{t}｜【{result}】｜" + ("；".join(c[0] for c in changed) if changed else "无实质监控字段变化")
    lines = [head, f"内部名称：{name}", f"站点 / ASIN：{site} / {asin}", f"页面标题：{title}", f"抓取时间：{t}（东八区）", f"总体状态：【{result}】", ""]
    lines.append("【重点变化】")
    if changed:
        for f,b,n,desc in changed[:12]:
            lines.append(f"· 【变】{f}：{desc}")
    else:
        lines.append("· 无实质监控字段变化")
    lines += ["", "【评论变化】"]
    rev = [c for c in changed if c[0] in ("评论总数","1星评论数","2星评论数","3星评论数","4星评论数","5星评论数","评分人数","综合评分")]
    if rev:
        for f,b,n,desc in rev:
            label = "净变化" if "评论" in f or f=="评分人数" else "变化"
            lines.append(f"· {f}：{desc}（{label}）")
    else:
        lines.append("· 评论字段无净变化" if fields.get("评分人数") else "· 评论数据未采集")
    lines += ["", "【图片 / 视频】", f"· 主图：{'有变化' if any(c[0]=='主图URL' for c in changed) else '无变化'}",
              f"· 副图：{'有变化' if any(c[0]=='副图URL列表' for c in changed) else '无变化'}",
              f"· A+：{'有变化' if any(c[0]=='A+图片URL列表' for c in changed) else '无变化'}",
              "· 说明：只比 URL，不评价画面或播放", ""]
    lines.append("【建议核查】")
    lines.append("· 请按下方「有变化字段」逐项确认是否为真实页面变化（排除采错）" if changed else "· 无需核查")
    lines += ["", "======= 差异对比 =======",
              f"对比基准：{base_id or '本轮建立基线'}/{base_time or '-'}",
              "采集环境：", f"· 站点：{site}", f"· ASIN：{fields.get('实际ASIN') or asin}",
              f"· 配送：{cfg['postcode']}", f"· 币种：{cfg['currency']}",
              f"· 评论范围：{rev_scope or '未采集'}", f"对比结果：【{result}】", "", "—— 有变化字段 ——"]
    if changed:
        for f,b,n,desc in changed:
            lines += [f"【{f}】", f"· 上次：{fmt_v(b)[:80]}", f"· 本次：{fmt_v(n)[:80]}", f"· 差异：{desc}", "· 补充：无", ""]
    else:
        lines.append("· （无）")
    lines += ["", "—— 未变化字段 ——", "· " + "、".join(unchanged) if unchanged else "· （无）"]
    return "\n".join(lines)

# ---------- 主流程 ----------
def main():
    t_start = time.time()
    products = read_table(T_BASE)
    print(f"[i] 基础表商品 {len(products)} 条")
    old = read_table(T_SCRAPE)
    FIELDS = field_names(T_SCRAPE)
    global DIFF_FIELDS
    DIFF_FIELDS = [f for f in DIFF_FIELDS if f in FIELDS]
    print(f"[i] 抓取表存量 {len(old)} 条")
    sb = SiteBrowsers()
    results = []
    try:
        for p in products:
            site = (p.get("站点") or ["其他"])[0] if isinstance(p.get("站点"), list) else p.get("站点")
            if site not in SITE_CONFIGS:
                results.append({"商品": p.get("内部产品名称"), "状态": "跳过：站点未配置"}); continue
            cfg = SITE_CONFIGS[site]
            url = (p.get("监控主链接") or [""])[0] if isinstance(p.get("监控主链接"), list) else p.get("监控主链接")
            name = p.get("内部产品名称") or ""
            t0 = time.time()
            rec = {"商品": name, "站点": site, "_基础ID": p["_id"]}
            page = None
            try:
                ctx = sb.get(site)
                page = ctx.new_page()
                page.goto(url, wait_until="domcontentloaded", timeout=45000)
                time.sleep(random.uniform(2, 3))
                ok, g = ensure_location(page, site, log=print)
                rec["邮编"] = g
                if not ok:
                    page.reload(wait_until="domcontentloaded", timeout=45000); time.sleep(2)
                f = extract_product(page, site)
                # 含税换算
                incl, excl = tax_split(f.get("当前展示价格"), site)
                f["当前含税价格"] = incl; f["当前不含税价格"] = excl
                f["实际配送地区"] = glow_text(page).replace("‌","")
                # 评论星级
                total, stars, rev_state, rev_scope = try_review_counts(page, site)
                f["评论总数"] = total
                for k, v in (stars or {}).items(): f[k] = v
                rec_fields = dict(f)
                rec_fields.update({
                    "抓取记录标识": gen_id(site),
                    "抓取时间": now_cst().strftime("%Y-%m-%dT%H:%M:%S+08:00"),
                    "监控时间段": [slot_of(now_cst())],
                    "评估版本": VERSION,
                    "评论采集状态": [rev_state],
                    "评论覆盖范围": rev_scope,
                    "实际配送地区": cfg["postcode"],
                })
                # 找基准（固定：同商品同站点同邮编同币种的最早成功记录）
                pid = p["_id"]
                cands = [r for r in old if (r.get("关联监控商品") or [None])[0] == pid
                         and str(r.get("实际站点","")) == site
                         and str(r.get("实际配送地区","")) == cfg["postcode"]
                         and (r.get("抓取状态") or [""])[0] in ("成功","部分成功")
                         and r.get("当前含税价格") is not None]
                base_rec = cands[0] if cands else None
                base_id = (base_rec or {}).get("对比基准记录ID") or (base_rec or {}).get("_id") if base_rec else None
                if base_rec and not (base_rec.get("对比基准记录ID") or "").strip():
                    base_id = base_rec["_id"]
                changed, unchanged, incompar = diff_records(base_rec, rec_fields) if base_rec else ([],[],[])
                result = eval_status(changed, incompar, base_rec is not None, rec_fields)
                rec_fields["评估状态"] = [result]
                rec_fields["对比基准记录ID"] = base_id or ""
                rec_fields["差异汇总"] = build_diff_summary(base_id, (base_rec or {}).get("抓取时间"), changed, unchanged, incompar, result)
                rec_fields["简报说明"] = build_brief(name, site, (p.get("ASIN") or [""])[0], rec_fields, changed, unchanged, base_id, (base_rec or {}).get("抓取时间"), result, rev_scope, cfg)
                miss = [f_ for f_ in ["当前含税价格","综合评分","评分人数","标题"] if rec_fields.get(f_) is None]
                rec_fields["缺失字段及原因"] = ("；".join(f"{x}:页面未展示或解析失败" for x in miss)) or ""
                rec_fields["抓取状态"] = ["成功" if not miss else "部分成功"]
                rec_fields["关联监控商品"] = [pid]
                rec["_fields"] = rec_fields
                rec["评估"] = result
                rec["变化数"] = len(changed)
                rec["变化字段"] = "、".join(c[0] for c in changed)[:80]
            except Exception as e:
                rec["评估"] = "异常"
                rec["_fields"] = {"抓取记录标识": gen_id(site), "抓取时间": now_cst().strftime("%Y-%m-%dT%H:%M:%S+08:00"),
                                  "监控时间段": [slot_of(now_cst())], "评估版本": VERSION, "抓取状态": ["失败"],
                                  "页面访问状态": "加载失败", "评估状态": ["异常"], "关联监控商品": [p["_id"]],
                                  "差异汇总": "对比基准：-/-；结果：无法比较；原因：页面加载异常",
                                  "简报说明": f"======= 监控简报 =======\n一句话结论：{name}｜{site}｜-｜{now_cst().strftime('%Y-%m-%d %H:%M')}｜【异常】｜抓取失败，不推断商品状态\n抓取异常：{type(e).__name__}",
                                  "缺失字段及原因": "全部字段:页面加载异常"}
                rec["错误"] = f"{type(e).__name__}: {e}"
            finally:
                if page:
                    try: page.close()
                    except Exception: pass
            results.append(rec)
            print(f"[i] {name}: {rec.get('评估')} 变化{rec.get('变化数',0)}项 耗时{round(time.time()-t0,1)}s")
            time.sleep(random.uniform(3, 6))
    finally:
        sb.close()

    # ---- 写入飞书 ----
    if DRY:
        print("[dry] 不写飞书")
        json.dump(results, open("/home/user/dry_result.json","w"), ensure_ascii=False, indent=1)
        return results
    new_ids = []
    LIST_TEXT = {"副图URL列表", "A+图片URL列表", "卖点列表", "变体清单"}
    for rec in results:
        fields = {}
        for k, v in rec["_fields"].items():
            if k.startswith("_") or v in (None, "", []) or k in ("配送信息原文",) or k not in FIELDS:
                continue
            if k in LIST_TEXT and isinstance(v, list):
                v = "\n".join(str(x) for x in v)
            fields[k] = v
        payload = {"create_records": [fields]}
        d = lark("base","+record-batch-create","--base-token",BASE,"--table-id",T_SCRAPE,"--as","user","--json",json.dumps(payload,ensure_ascii=False),"--format","json")
        rid = None
        if d.get("ok"):
            rid = d["data"]["record_id_list"][0]; rec["新记录"] = rid; new_ids.append(rid)
        else:
            rec["写入失败"] = json.dumps(d.get("error",{}),ensure_ascii=False)[:200]
        print(f"[w] {rec['商品']}: {rid or rec.get('写入失败')}")
    # ---- 汇总表 ----
    n_att = sum(1 for r in results if r.get("评估")=="异常")
    n_focus = sum(1 for r in results if r.get("评估")=="需关注")
    n_base = sum(1 for r in results if r.get("评估")=="建立基线")
    n_same = sum(1 for r in results if r.get("评估")=="无变化")
    n_chg = sum(1 for r in results if r.get("评估") in ("有变化","需关注"))
    overall = "有异常" if n_att else ("有关注项" if n_focus else ("含基线" if n_base else "全无变化"))
    uk = sum(1 for r in results if r.get("站点")=="UK"); de = sum(1 for r in results if r.get("站点")=="DE")
    day = now_cst().strftime("%Y-%m-%d"); t = now_cst().strftime("%Y-%m-%d %H:%M"); slot = slot_of(now_cst())
    tldr = f"{day}｜{slot}｜共{len(results)}条｜【{'有变化' if n_chg else '没变化'}】｜有变化{n_chg}｜没变化{n_same}"
    focus_ids = [r["新记录"] for r in results if r.get("评估") in ("需关注","异常") and r.get("新记录")]
    ssum = [f"======= 时段汇总简报 =======", f"一句话结论：{tldr}", f"生成时间：{t}（东八区）", "",
            "【本批概览】", f"· 链接总数：{len(results)}（UK {uk} / DE {de} / FR 0）", f"· 有变化：{n_chg}", f"· 没变化：{n_same}", "",
            "【有变化】"]
    for r in results:
        if r.get("评估") in ("有变化","需关注","异常","建立基线"):
            ssum.append(f"· {r['商品']}｜{r.get('站点')}｜{r.get('变化字段') or r.get('评估')}　记录：{r.get('新记录','-')}")
    if not any(r.get("评估") in ("有变化","需关注","异常","建立基线") for r in results): ssum.append("· 无")
    ssum += ["", "【没变化】"]
    ssum += [f"· {r['商品']}" for r in results if r.get("评估")=="无变化"] or ["· 无"]
    ssum += ["", "【建议】", "· 优先核查【有变化】清单" if n_chg else "· 本批无需特别核查", "", "【备注】",
             f"· 建立基线 {n_base} 条" if n_base else "· 无基线项", "· 说明：正文只汇总「有变化」与「没变化」；基线与异常仅作备注，不占主清单"]
    payload = {"create_records": [{
        "汇报日期": f"{day}T00:00:00+08:00", "监控时间段": [slot], "生成状态": ["成功"],
        "链接总数": len(results), "UK条数": uk, "DE条数": de, "建立基线数": n_base, "无变化数": n_same,
        "需关注数": n_focus + n_att, "异常失败数": n_att, "总体状态": [overall],
        "一句话结论": tldr, "汇总简报": "\n".join(ssum),
        "关联抓取记录": new_ids, "需要关注的记录": focus_ids,
    }]}
    d = lark("base","+record-batch-create","--base-token",BASE,"--table-id",T_SUM,"--as","user","--json",json.dumps(payload,ensure_ascii=False),"--format","json")
    sum_id = d["data"]["record_id_list"][0] if d.get("ok") else None
    print(f"[w] 汇总记录: {sum_id or d}")
    # ---- 基础表最新简报 ----
    upd = {}
    for r in results:
        if r.get("新记录"):
            upd[r["_基础ID"]] = {"最新简报": [r["新记录"]]}
    if upd:
        d = lark("base","+record-batch-update","--base-token",BASE,"--table-id",T_BASE,"--as","user","--json",
                 json.dumps({"update_records": upd},ensure_ascii=False),"--format","json")
        if not d.get("ok"): print("[w] 最新简报更新失败:", json.dumps(d.get("error",{}),ensure_ascii=False)[:150])
    # ---- PNG 简报卡 ----
    try:
        from PIL import Image, ImageDraw, ImageFont
        W, H = 750, 260 + 46*len(results)
        img = Image.new("RGB", (W, H), "#1a2233"); dr = ImageDraw.Draw(img)
        try: fnt = ImageFont.truetype("/usr/share/fonts/truetype/dejavu/DejaVuSans-Bold.ttf", 22)
        except Exception: fnt = ImageFont.load_default()
        try: fnt_s = ImageFont.truetype("/usr/share/fonts/truetype/dejavu/DejaVuSans.ttf", 18)
        except Exception: fnt_s = ImageFont.load_default()
        dr.text((28, 24), f"亚马逊监控简报 {day} {slot}", font=fnt, fill="#ffffff")
        dr.text((28, 60), tldr, font=fnt_s, fill="#9fb3d1")
        y = 110
        for r in results:
            color = {"需关注":"#ffb84d","异常":"#ff6b6b","建立基线":"#4dc3ff","无变化":"#7ee08b"}.get(r.get("评估"), "#dddddd")
            dr.text((28, y), f"· {r['商品']}｜{r.get('站点')}｜{r.get('评估')}", font=fnt_s, fill=color)
            dr.text((48, y+26), (r.get("变化字段") or "无变化/基线")[:64], font=fnt_s, fill="#8899aa")
            y += 46
        dr.text((28, y+8), f"生成 {t}（东八区）", font=fnt_s, fill="#667788")
        png = "slot_brief.png"
        img.save("/home/user/" + png)
        if sum_id:
            d = lark("base","+record-upload-attachment","--base-token",BASE,"--table-id",T_SUM,"--record-id",sum_id,"--field-id","简报图片","--file",png,"--as","user","--format","json")
            print(f"[w] PNG附件: {'ok' if d.get('ok') else json.dumps(d)[:150]}")
    except Exception as e:
        print(f"[w] PNG生成失败: {e}")
    print(f"[done] 总耗时 {round(time.time()-t_start,1)}s")
    return results

if __name__ == "__main__":
    main()
