#!/usr/bin/env python3
"""客户数据补拉 — 带重试的 OFFSET 分页，修复 _full_refresh_offset.py 第26页超时导致的缺漏"""
import json
import os
import re
import time
import urllib.request

MCP_URL = "https://project.feishu.cn/mcp_server/v1"
MCP_TOKEN = "m-abfb29e8-3104-434f-9944-8d0bb592f8cd"
DATA_DIR = "/Users/liuxinyuan/Desktop/Hermes输出-工作类/数据"
SALES_PK = "6593cd71471290e3cc6be6e6"
TYPE_CUSTOMER = "65ae1e403c87b152f3365ca6"


def mcp_call(method, args, timeout=120, retries=3):
    payload = {
        "jsonrpc": "2.0", "id": 1,
        "method": "tools/call",
        "params": {"name": method, "arguments": args}
    }
    for attempt in range(1, retries + 1):
        try:
            req = urllib.request.Request(
                MCP_URL,
                data=json.dumps(payload).encode('utf-8'),
                headers={"X-Mcp-Token": MCP_TOKEN, "Content-Type": "application/json"}
            )
            with urllib.request.urlopen(req, timeout=timeout) as r:
                raw = r.read().decode('utf-8')
            break
        except Exception as e:
            if attempt == retries:
                print(f"  ❌ HTTP error (重试{retries}次后): {e}")
                return None
            print(f"  ⚠️ 第{attempt}次失败: {e}, 2秒后重试...")
            time.sleep(2)

    if raw.startswith("data:"):
        raw = raw.split("\n")[0].replace("data: ", "", 1)
    try:
        d = json.loads(raw)
    except Exception as e:
        print(f"  ❌ JSON parse error: {e}")
        return None
    if "error" in d:
        print(f"  ❌ MCP error: {json.dumps(d['error'], ensure_ascii=False)[:200]}")
        return None
    for c in d.get("result", {}).get("content", []):
        t = c.get("text", "")
        if not t or "log_id" in t:
            continue
        try:
            return json.loads(re.sub(r'\nlog_id:.*$', '', t.strip()))
        except Exception:
            pass
    return None


def parse_items(result):
    items = []
    if not result:
        return items
    for gid, gitems in result.get("data", {}).items():
        if not isinstance(gitems, list):
            continue
        for item in gitems:
            fields = {}
            for f in (item.get("moql_field_list") or []):
                k = f["key"]
                v = f.get("value")
                if isinstance(v, list) and len(v) > 0:
                    if isinstance(v[0], dict) and "label" in v[0]:
                        fields[k] = v[0]["label"]
                        continue
                    v = v[0]
                if v is None:
                    fields[k] = ""
                elif isinstance(v, dict):
                    if "string_value" in v:
                        fields[k] = v["string_value"]
                    elif "double_value" in v:
                        fields[k] = v["double_value"]
                    elif "long_value" in v:
                        fields[k] = v["long_value"]
                    elif "key_label_value" in v:
                        fields[k] = v["key_label_value"]
                    elif "user_value" in v:
                        fields[k] = v["user_value"]
                    elif "key_label_value_list" in v:
                        fields[k] = v["key_label_value_list"]
                    else:
                        fields[k] = str(v)
                else:
                    fields[k] = str(v)
            items.append(fields)
    return items


def main():
    mql = (
        "SELECT name, field_c3224b, field_5ed7ab, field_c8e80d, "
        "field_17186c, field_6415cf, work_item_id "
        "FROM `销售管理`.`65ae1e403c87b152f3365ca6` "
        "ORDER BY name"
    )
    all_items = []
    seen_ids = set()
    offset = 0
    empty_streak = 0
    max_pages = 100

    for page in range(1, max_pages + 1):
        paginated_mql = f"{mql} LIMIT 50 OFFSET {offset}"
        result = mcp_call("search_by_mql", {
            "project_key": SALES_PK, "mql": paginated_mql, "session_id": ""
        })
        if not result:
            print(f"  ⚠️ 第{page}页(OFFSET={offset}) 失败，中断")
            break
        items = parse_items(result)
        if len(items) == 0:
            empty_streak += 1
            if empty_streak >= 2:
                print(f"  ✅ 连续2页空，停止。累计 {len(all_items)}")
                break
        else:
            empty_streak = 0
            for item in items:
                wid = item.get("work_item_id", "")
                if wid and wid in seen_ids:
                    continue
                if wid:
                    seen_ids.add(wid)
                all_items.append(item)
        print(f"  第{page}页(OFFSET={offset}): {len(items)}条, 累计: {len(all_items)}")
        offset += 50

    fp = os.path.join(DATA_DIR, "mcp_customers.json")
    with open(fp, "w", encoding="utf-8") as f:
        json.dump({"updated_at": time.strftime("%Y-%m-%d %H:%M:%S"), "customers": all_items},
                  f, ensure_ascii=False, indent=2)
    print(f"  ✅ 写入 mcp_customers.json: {len(all_items)} 条, {os.path.getsize(fp)/1024:.1f}KB")


if __name__ == "__main__":
    main()
