#!/usr/bin/python3.12
"""从Dify选车网工作流获取车型热点并写入日报数据"""
import json, re, os, sys
import urllib.request
from datetime import datetime
# 每日营销日报-车型热点-多渠道-Api
DIFY_URL = "http://43.164.190.156/v1/workflows/run"
DIFY_TOKEN = "app-UH9ps9DtWQpMKqPVOzjsx5rI"


# 车型名关键词（有这些才算有具体车型）
_MODEL_KEYWORDS = [
    'GT', 'GTi', 'EV', 'SUV', 'MPV', 'PHEV', 'DM-i', 'L', 'Pro', 'Plus', 'MAX',
    '纯电', '插混', '混动', '增程', '电动',
    '款', '换代', '改款', '新款', '全新',
]
_MODEL_PATTERN = __import__('re').compile(r'[A-Z0-9]{2,}')


def _has_model(title):
    """检查标题是否有具体车型名"""
    # 包含数字+字母组合（如 Q7、X5、Z20、K50、i60、G700 等）
    if _MODEL_PATTERN.search(title):
        return True
    # 包含车型关键词
    for kw in _MODEL_KEYWORDS:
        if kw in title:
            return True
    # 包含常见车型品牌后缀
    for suffix in ['版', '型', '代', '系']:
        for num in ['0', '1', '2', '3', '4', '5', '6', '7', '8', '9']:
            if num in title:
                return True
    return False

def call_dify_workflow(page_url, user="root"):
    """调用Dify工作流（streaming），返回完整输出"""
    payload = json.dumps({
        "inputs": {"pageUrl": page_url},
        "response_mode": "streaming",
        "user": user
    }).encode('utf-8')

    req = urllib.request.Request(
        DIFY_URL, data=payload,
        headers={'Authorization': f'Bearer {DIFY_TOKEN}', 'Content-Type': 'application/json'},
        method='POST'
    )

    import http.client
    try:
        resp = urllib.request.urlopen(req, timeout=300)
    except Exception as e:
        raise ConnectionError(f"Dify请求失败: {e}")

    buffer = b""
    import signal
    def timeout_handler(signum, frame):
        raise TimeoutError("Stream read timeout")
    signal.signal(signal.SIGALRM, timeout_handler)
    signal.alarm(300)
    while True:
        chunk = resp.read(65536)
        if not chunk:
            break
        buffer += chunk

    full_text = buffer.decode('utf-8', errors='replace')

    lines = full_text.strip().split('\n')
    data_line = ''
    for line in lines:
        if line.startswith('data: '):
            data_line = line[6:]
    if data_line:
        result = json.loads(data_line)
        outputs = result.get('data', {}).get('outputs', {})
        return outputs
    else:
        raise ValueError(f"未找到有效的SSE数据行: {full_text[:200]}")
def merge_to_origin(mmd_data, vehicle_list):
    """追加到已有列表，按title去重"""
    for item in mmd_data:
        if item.get('name') == 'vehicle_hotposts' or item.get('name') == 'vehicle_hotspots':
            existing = item.get('list', [])
            seen_titles = {e.get('title', '') for e in existing}
            for ni in vehicle_list:
                t = ni.get('title', '')
                if t and t not in seen_titles:
                    existing.append(ni)
                    seen_titles.add(t)
            item['list'] = existing
            return mmd_data
    
    mmd_data.append({
        "name": "vehicle_hotspots",
        "opinion": "",
        "list": vehicle_list
    })
    
    return mmd_data

PAGE_URLS = [
    "https://www.autohome.com.cn/all/#pvareaid=3311229",
    "https://news.yiche.com/xinche/"
]

def main():
    mmdd = datetime.now().strftime("%m%d")
    origin_path = f"/data/news/json/origin_data/{mmdd}data.json"
    
    # 读取或初始化origin文件
    if os.path.exists(origin_path):
        with open(origin_path, 'r', encoding='utf-8') as f:
            mmd_data = json.load(f)
        import sys; print(f"📄 已读取: {origin_path}")
    else:
        mmd_data = []
        import sys; print(f"📄 新建文件: {origin_path}")
    
    all_vehicle = []
    seen_titles = set()
    
    for idx, page_url in enumerate(PAGE_URLS):
        import sys; print(f"\n[{idx+1}/{len(PAGE_URLS)}] 请求: {page_url}")
        
        outputs = call_dify_workflow(page_url)
        if not outputs:
            import sys; print(f"  ⚠️ 第{idx+1}次无返回，跳过")
            continue
        
        vehicle_data = None
        for key in ['vehicle_hotspots', 'vehicle_hotposts', 'vehicle_data']:
            if key in outputs:
                vehicle_data = outputs[key]
                import sys; print(f"  ✅ 找到数据: {key}")
                break
        
        if vehicle_data is None:
            import sys; print(f"  ⚠️ 未找到车型热点，可用keys: {list(outputs.keys())}")
            continue
        
        if isinstance(vehicle_data, str):
            try:
                vehicle_data = json.loads(vehicle_data)
            except:
                import sys; print(f"  ⚠️ vehicle数据非JSON")
                continue
        
        if not isinstance(vehicle_data, list):
            import sys; print(f"  ⚠️ 数据类型错误: {type(vehicle_data)}")
            continue
        
        # 去重合并
        added = 0
        for item in vehicle_data:
            title = item.get('title', '')
            model = item.get('model', '')
            # 跳过无model或model为空的条目
            if not model or model == "未提供":
                continue
            item["yesorno"] = ""
            if title and title not in seen_titles:
                seen_titles.add(title)
                all_vehicle.append(item)
                added += 1
        
        import sys; print(f"  本次新增 {added} 条（去重后），累计 {len(all_vehicle)} 条")
        # 每轮写入，防止后续卡死丢数据
        if all_vehicle:
            merged = merge_to_origin(mmd_data, all_vehicle)
            os.makedirs(os.path.dirname(origin_path), exist_ok=True)
            with open(origin_path, 'w', encoding='utf-8') as f:
                json.dump(merged, f, ensure_ascii=False, indent=2)
    
    all_vehicle = [v for v in all_vehicle if _has_model(v.get("title", ""))]
    if not all_vehicle:
        print("❌ 所有请求均未获取到数据")
        sys.exit(1)
    
    # 写入origin文件
    merged = merge_to_origin(mmd_data, all_vehicle)
    os.makedirs(os.path.dirname(origin_path), exist_ok=True)
    with open(origin_path, 'w', encoding='utf-8') as f:
        json.dump(merged, f, ensure_ascii=False, indent=2)
    
    import sys; print(f"\n✅ 共获取 {len(all_vehicle)} 条去重车型热点，已写入 {origin_path}")

if __name__ == '__main__':
    main()
