#!/usr/bin/env python3
# """社媒营销数据读写-抖小B 工作流"""
import json, os, sys, re, urllib.request
from datetime import datetime

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

def call_workflow():
    """调用Dify workflow（streaming），返回完整SSE输出"""
    payload = json.dumps({
        "inputs": {"method": "GET", "key": "all"},
        "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'
    )

    full_data = b""
    with urllib.request.urlopen(req, timeout=120) as resp:
        while True:
            chunk = resp.read(4096)
            if not chunk:
                break
            full_data += chunk
    return full_data.decode('utf-8')

def parse_sse_output(sse_text: str) -> dict | None:
    """从SSE流中提取最后一个 data: 块的 outputs.output"""
    outputs_text = None
    for line in sse_text.splitlines():
        line = line.strip()
        if line.startswith("data:"):
            data_str = line[5:].strip()
            try:
                obj = json.loads(data_str)
                if "data" in obj and "outputs" in obj["data"]:
                    raw = obj["data"]["outputs"].get("output", "")
                    if raw:
                        outputs_text = raw
            except json.JSONDecodeError:
                continue
    if outputs_text:
        return outputs_text if isinstance(outputs_text, dict) else json.loads(outputs_text)
    return None

def main():
    if len(sys.argv) < 2:
        print("用法: python3 get_buzz_from_dify.py <写入文件路径>")
        sys.exit(1)

    out_path = sys.argv[1]
    ts = datetime.now().strftime('%H:%M:%S')

    print(f"[{ts}] 调用社媒营销工作流...")
    sse = call_workflow()
    print(f"  收到 {len(sse)} bytes")

    result = parse_sse_output(sse)
    if result is None:
        print("  ❌ 未能解析 outputs.output")
        sys.exit(1)

    # 读取目标文件（如果存在），追加数据
    existing = []
    if os.path.exists(out_path):
        with open(out_path) as f:
            existing = json.load(f)

    if isinstance(result, dict) and isinstance(existing, dict):
        # 目标文件是 dict → 合并 key
        existing.update(result)
        merged = existing
        count = len(result)
    elif isinstance(result, list):
        existing.extend(result)
        merged = existing
        count = len(result)
    else:
        existing.append(result)
        merged = existing
        count = 1

    with open(out_path, 'w') as f:
        json.dump(merged, f, ensure_ascii=False, indent=2)

    print(f"  ✅ 追加 {count} 个 key，共 {len(merged)} 个 key" if isinstance(merged, dict) else f"  ✅ 追加 {count} 条，共 {len(merged)} 条")
    print(f"  已保存到: {out_path}")

if __name__ == '__main__':
    main()
