#!/usr/bin/env python3
"""客诉数据管道 — MQL全量拉取（替代旧 list_todo 方式）
用法: python3 pull_complaints.py
输出: mcp_complaints.json
"""
import subprocess, json, os, sys
from datetime import datetime, date
from collections import defaultdict

TK = "m-abfb29e8-3104-434f-9944-8d0bb592f8cd"
PK = "658288abfb8bd616b17025f1"  # 售后管理
BASE = os.path.dirname(os.path.abspath(__file__))
OUTPUT = os.path.join(BASE, "mcp_complaints.json")
TODAY = date.today()

# 离职人员客诉归属映射：黄明月/张立娅 的客诉归并到接替人刘子研
RESIGN_REASSIGN = {'黄明月': '刘子研', '张立娅': '刘子研'}

def mcp(method, args):
    r = subprocess.run(['curl','-s','-X','POST','https://project.feishu.cn/mcp_server/v1',
        '-H',f'X-Mcp-Token: {TK}','-H','Content-Type: application/json',
        '-d', json.dumps({"jsonrpc":"2.0","method":"tools/call",
        "params":{"name":method,"arguments":args},"id":1})],
        capture_output=True, text=True, timeout=30)
    try:
        d = json.loads(r.stdout)
        if 'error' in d: return None
        for c in d['result']['content']:
            t = c.get('text','')
            if 'log_id' in t: continue
            return json.loads(t)
    except: return None

def parse_items(result):
    items = []
    if not result: return items
    for gid, gitems in result.get('data', {}).items():
        for item in gitems:
            fields = {}
            for f in item.get('moql_field_list', []):
                k = f['key']
                v = f.get('value')
                if v is None:
                    fields[k] = ''
                elif 'string_value' in v:
                    fields[k] = v['string_value']
                elif 'long_value' in v:
                    fields[k] = v['long_value']
                elif 'user_value' in v:
                    fields[k] = v['user_value']
                elif 'user_value_list' in v:
                    fields[k] = v['user_value_list']
                elif 'key_label_value' in v:
                    fields[k] = v['key_label_value']
                elif 'key_label_value_list' in v:
                    fields[k] = v['key_label_value_list']
                else:
                    fields[k] = ''
            items.append(fields)
    return items

def extract_label(v):
    if isinstance(v, dict):
        return v.get('label', '')
    if isinstance(v, list) and v:
        return ', '.join(x.get('label', '') for x in v)
    return ''

def extract_name(v):
    if isinstance(v, dict):
        return v.get('name_cn', v.get('name_en', ''))
    if isinstance(v, list) and v:
        return ', '.join(x.get('name_cn', x.get('name_en', '')) for x in v)
    return ''

def query_all():
    all_items = []
    seen = set()
    page = 0
    COLUMNS = 'name, start_time, current_status_operator, work_item_id, work_item_status, field_9e2144, owner, field_3ce9fe, field_ac5caf'
    while page < 50:
        offset = page * 50
        mql = f"SELECT {COLUMNS} FROM `售后管理`.`客户反馈单` ORDER BY start_time DESC LIMIT 50 OFFSET {offset}"
        result = mcp("search_by_mql", {"project_key": PK, "mql": mql})
        if not result: break
        items = parse_items(result)
        if not items: break
        new = []
        for i in items:
            wid = str(i.get('work_item_id', ''))
            if wid and wid not in seen:
                seen.add(wid)
                new.append(i)
        all_items.extend(new)
        page += 1
        if len(items) < 50: break
        print(f"  page {page}: +{len(new)} (total {len(all_items)})")
    return all_items

print(f"[{datetime.now().strftime('%H:%M:%S')}] 拉取客诉数据...")
raw = query_all()
print(f"  共 {len(raw)} 条")

if not raw:
    print("❌ 拉取失败或返回空数据，保留旧文件，不覆盖（exit 1）")
    sys.exit(1)

output = []
by_urgency = defaultdict(int)
for r in raw:
    created = r.get('start_time', '')
    if created:
        try:
            days = (TODAY - datetime.strptime(created[:10], '%Y-%m-%d').date()).days
        except:
            days = 0
    else:
        days = 0

    if days >= 365: urgency = 'critical'
    elif days >= 200: urgency = 'high'
    elif days >= 90: urgency = 'medium'
    else: urgency = 'normal'
    by_urgency[urgency] += 1

    ops = r.get('current_status_operator', [])
    owner = r.get('owner', {})
    creator_name = extract_name(owner)
    creator_name = RESIGN_REASSIGN.get(creator_name, creator_name)  # 离职归并

    output.append({
        'id': str(r.get('work_item_id', '')),
        'name': (r.get('name', '') or '')[:100],
        'days': days,
        'created': created,
        'urgency': urgency,
        'area': extract_label(r.get('field_9e2144', '')),
        'status': extract_label(r.get('work_item_status', '')),
        'type': extract_label(r.get('field_3ce9fe', '')),
        'operator': [x.get('name_cn', x.get('name_en', '')) for x in ops] if isinstance(ops, list) else ([ops.get('name_cn', ops.get('name_en', ''))] if isinstance(ops, dict) and ops else []),
        'customer': extract_label(r.get('field_ac5caf', '')),
        'creator': creator_name,
        'node': extract_label(r.get('work_item_status', '')),
    })

output.sort(key=lambda x: -x['days'])

std = {
    'updated': datetime.now().strftime('%Y-%m-%d %H:%M:%S'),
    'total': len(output),
    'by_urgency': dict(by_urgency),
    'items': output,
}

with open(OUTPUT, 'w', encoding='utf-8') as f:
    json.dump(std, f, ensure_ascii=False, indent=2)

print(f"✅ 写入 {OUTPUT} ({len(output)} 条)")
print(f"   紧急度: {dict(by_urgency)}")
