#!/usr/bin/env python3
"""import_jsonl_auto.py - 调度器自动导入器 (根治"只抓不导入")

用途: 日跑调度器跑完脚本后自动调用, 把脚本输出的数据导入 search.db。
支持两类输出:
  1. --file <jsonl路径>      脚本写文件型 (直接导入)
  2. --stdout <out文件>      脚本 print 到 stdout 型 (智能解析)

stdout 解析支持三种格式:
  A. JSON_BEGIN:...:JSON_END 包裹 (crawl_hequ 等)
  B. 整体 JSON 数组 print(json.dumps(list, indent=2)) (crawl_shuozhou 等)
  C. 逐行 JSON 对象 (JSONL) (crawl_baixiang 等)

字段归一化: title, url/page_url/link, date/pub_date/publish_date, content, attachments, site_name
用法:
  python3 import_jsonl_auto.py --file /tmp/x.jsonl [--script crawl_x.py] [--site "站点名"]
  python3 import_jsonl_auto.py --stdout /tmp/x.out [--script crawl_x.py] [--site "站点名"]
"""
import json, sqlite3, sys, os, re, argparse, subprocess

BASE = "/root/gov_crawler"
DB_PATH = os.getenv("SEARCH_DB", "/root/search.db")


# ---------- stdout 解析 ----------
def parse_stdout_text(text):
    """从脚本 stdout 文本中提取 JSON 数据记录列表 (dict 列表)"""
    items = []
    # A. JSON_BEGIN:...:JSON_END
    m = re.search(r"JSON_BEGIN:(.*?):JSON_END", text, re.S)
    if m:
        try:
            data = json.loads(m.group(1))
            if isinstance(data, dict) and "items" in data:
                items = data["items"]
            elif isinstance(data, list):
                items = data
            else:
                items = [data]
            return items
        except Exception:
            pass
    # B. 整体 JSON (数组或对象)
    stripped = text.strip()
    # 去掉可能的日志前缀行: 找第一个 [ 或 { 开头到最后一个 ] 或 }
    try:
        data = json.loads(stripped)
        if isinstance(data, list):
            return data
        if isinstance(data, dict):
            if "items" in data and isinstance(data["items"], list):
                return data["items"]
            return [data]
    except Exception:
        pass
    # C. 逐行 JSONL (容忍日志行夹杂)
    for line in stripped.splitlines():
        line = line.strip()
        if not line:
            continue
        if line[0] in "[{":
            try:
                obj = json.loads(line)
                if isinstance(obj, dict):
                    items.append(obj)
                elif isinstance(obj, list):
                    items.extend(obj)
            except Exception:
                pass
    return items


# ---------- 字段归一化 ----------
def normalize_item(it):
    title = str(it.get("title") or it.get("name") or "").strip()[:500]
    page_url = str(it.get("page_url") or it.get("url") or it.get("link") or it.get("source_url") or "").strip()[:1000]
    content = str(it.get("content") or it.get("content_html") or it.get("body") or it.get("text") or "").strip()
    pub = str(it.get("publish_date") or it.get("pub_date") or it.get("date") or it.get("date_pub") or "")[:20]
    site = str(it.get("site_name") or "").strip()[:100]
    att = it.get("attachments") or []
    if isinstance(att, str):
        try:
            att = json.loads(att)
        except Exception:
            att = []
    if not isinstance(att, list):
        att = []
    return {"title": title, "page_url": page_url, "content": content,
            "publish_date": pub, "site_name": site, "attachments": att}


# ---------- 写库 (与 import_jsonl_raw 同逻辑) ----------
def write_records(records, force_site, force_script):
    if not records:
        print("  [AUTO-IMP] 无有效记录")
        return 0
    db = sqlite3.connect(DB_PATH)
    db.execute("PRAGMA journal_mode=WAL")
    db.execute("PRAGMA busy_timeout=300000")
    db.execute("PRAGMA synchronous=NORMAL")
    db.execute("""CREATE VIRTUAL TABLE IF NOT EXISTS gov_search USING fts5(
        title, site_name, summary, tokenize=trigram
    )""")
    db.execute("""CREATE TRIGGER IF NOT EXISTS trg_gov_raw_fts_ins AFTER INSERT ON gov_raw BEGIN
      INSERT INTO gov_search(rowid, title, site_name, summary)
      VALUES (new.id, new.title, new.site_name, new.summary);
    END""")
    db.execute("""CREATE TRIGGER IF NOT EXISTS trg_gov_raw_fts_del AFTER DELETE ON gov_raw BEGIN
      DELETE FROM gov_search WHERE rowid = old.id;
    END""")
    db.execute("""CREATE TRIGGER IF NOT EXISTS trg_gov_raw_fts_upd AFTER UPDATE ON gov_raw BEGIN
      DELETE FROM gov_search WHERE rowid = old.id;
      INSERT INTO gov_search(rowid, title, site_name, summary)
      VALUES (new.id, new.title, new.site_name, new.summary);
    END""")
    db.commit()

    added, skipped = 0, 0
    for it in records:
        n = normalize_item(it)
        if not n["title"] or not n["page_url"]:
            skipped += 1
            continue
        site_name = force_site or n["site_name"] or "unknown"
        script_name = force_script or "import_jsonl_auto"
        att_str = json.dumps(n["attachments"], ensure_ascii=False) if n["attachments"] else ""
        # summary 供 FTS 搜索: 纯文本截断 (与 import_jsonl.py 一致 content[:500])
        summary = n["content"][:500] if n["content"] else n["title"][:500]
        db.execute("DELETE FROM gov_raw WHERE page_url = ?", (n["page_url"],))
        db.execute(
            "INSERT INTO gov_raw (title, page_url, content, publish_date, site_name, source_url, status, attachments, script_name, summary) VALUES (?, ?, ?, ?, ?, ?, 'synced', ?, ?, ?)",
            (n["title"], n["page_url"], n["content"], n["publish_date"], site_name, n["page_url"], att_str, script_name, summary)
        )
        added += 1
    db.commit()
    db.close()
    print(f"  [AUTO-IMP] ✅ 导入 +{added} (跳过 {skipped})")
    return added


def main():
    ap = argparse.ArgumentParser()
    ap.add_argument("--file", help="JSONL 文件路径")
    ap.add_argument("--stdout", help="脚本 stdout 输出文件路径")
    ap.add_argument("--script", default="", help="脚本名 (写入 script_name)")
    ap.add_argument("--site", default="", help="强制 site_name")
    args = ap.parse_args()

    if args.file:
        # 文件型: 直接按 JSONL 读 (也可能整体 JSON 数组, 容错)
        if not os.path.exists(args.file) or os.path.getsize(args.file) == 0:
            print(f"  [AUTO-IMP] 文件不存在或为空: {args.file}")
            sys.exit(0)
        records = []
        with open(args.file, encoding="utf-8", errors="ignore") as f:
            content = f.read()
        records = parse_stdout_text(content) if content.strip().startswith(("[", "{")) else []
        if not records:
            # 回退逐行 JSONL
            with open(args.file, encoding="utf-8", errors="ignore") as f:
                for line in f:
                    line = line.strip()
                    if not line:
                        continue
                    try:
                        records.append(json.loads(line))
                    except Exception:
                        pass
        added = write_records(records, args.site, args.script)
        sys.exit(0)

    if args.stdout:
        if not os.path.exists(args.stdout):
            print(f"  [AUTO-IMP] stdout 文件不存在: {args.stdout}")
            sys.exit(0)
        with open(args.stdout, encoding="utf-8", errors="ignore") as f:
            text = f.read()
        records = parse_stdout_text(text)
        added = write_records(records, args.site, args.script)
        sys.exit(0)

    print("Usage: import_jsonl_auto.py --file <jsonl> | --stdout <outfile> [--script X] [--site Y]")
    sys.exit(1)


if __name__ == "__main__":
    main()
