#!/usr/bin/python3.12
"""调用Dify全量采集工作流，输出保存到 origin_data/alldata/"""
import json, os, urllib.request, http.client
from datetime import datetime

WORKFLOW_URL = "http://43.164.190.156/v1/workflows/run"
WORKFLOW_TOKEN = "app-G3tJ1EbxiLWvHS1tgbcnJEcK"

OUTPUT_DIR = "/data/news/json/origin_data/alldata"


def call_workflow():
    """调用Dify workflow（streaming），返回outputs字典"""
    payload = json.dumps({
        "inputs": {"method": "GET"},
        "response_mode": "streaming",
        "user": "root"
    }).encode('utf-8')

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

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

    # 读取全部SSE数据
    import signal

    def timeout_handler(signum, frame):
        raise TimeoutError("Stream read timeout")

    signal.signal(signal.SIGALRM, timeout_handler)
    signal.alarm(300)  # 5分钟超时

    buffer = b""
    while True:
        chunk = resp.read(65536)
        if not chunk:
            break
        buffer += chunk

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

    # 解析SSE，找到 workflow_finished 事件中的 outputs
    for line in full_text.split('\n'):
        line = line.strip()
        if line.startswith('data: '):
            try:
                evt = json.loads(line[6:])
                if evt.get('event') == 'workflow_finished':
                    return evt.get('data', {}).get('outputs', {})
            except json.JSONDecodeError:
                continue

    # 如果没找到 workflow_finished，尝试直接解析整个 body
    try:
        return json.loads(full_text)
    except json.JSONDecodeError:
        return None


def extract_output(outputs):
    """从outputs中提取实际数据（output字段可能是JSON字符串或dict）"""
    if not outputs:
        return None
    val = outputs.get('output')
    if val is None:
        # 尝试 text 字段
        val = outputs.get('text')
    if val is None and len(outputs) > 0:
        # 取第一个非空值
        for v in outputs.values():
            if v:
                val = v
                break
    if val is None:
        return None

    # 如果是字符串，尝试解析为JSON
    if isinstance(val, str):
        try:
            val = json.loads(val)
        except json.JSONDecodeError:
            return val  # 非JSON字符串原样返回

    return val


def main():
    mmdd = datetime.now().strftime("%m%d")
    os.makedirs(OUTPUT_DIR, exist_ok=True)
    out_path = os.path.join(OUTPUT_DIR, f"{mmdd}.json")
    timestamp = datetime.now().strftime("%H:%M:%S")

    print(f"[{timestamp}] 调用Dify全量采集工作流...")

    outputs = call_workflow()
    if not outputs:
        print("❌ 工作流无返回")
        print(f"   可用keys: {list(outputs.keys()) if outputs else 'N/A'}")
        return

    print(f"  workflow outputs keys: {list(outputs.keys())}")

    data = extract_output(outputs)
    if data is None:
        print("❌ 未能提取数据")
        return

    # 写入文件
    with open(out_path, 'w', encoding='utf-8') as f:
        if isinstance(data, str):
            f.write(data)
        else:
            json.dump(data, f, ensure_ascii=False, indent=2)

    # 报告类型和大小
    dtype = type(data).__name__
    if isinstance(data, list):
        sz = len(data)
    elif isinstance(data, dict):
        sz = len(data)
    else:
        sz = len(str(data))
    file_size = os.path.getsize(out_path)

    print(f"\n✅ 已保存: {out_path}")
    print(f"   类型: {dtype}, 条目: {sz}, 大小: {file_size//1024}KB")


if __name__ == '__main__':
    main()
