#!/usr/bin/env python3
import os
"""达州市生态环境局-公示公告 爬虫
CMS: 自定义PHP (AJAX动态加载)
列表: POST api-ajax_list-{page}.html -> JSON (20条/页)
详情: news-show-{id}.html -> div.xxq_essay
"""
import sys, os, json, re, urllib.request, time
from datetime import datetime
from concurrent.futures import ThreadPoolExecutor, as_completed

SEARCH_DB = os.getenv("SEARCH_DB", "/root/search.db")
SITE_NAME = "达州市生态环境局-公示公告"
CUTOFF_DATE = "2023-06-18"
BASE_URL = "https://sthjj.dazhou.gov.cn"
MAX_WORKERS = 10
LIST_TIMEOUT = 20
DETAIL_TIMEOUT = 20

def log(msg):
    print(f"[{datetime.now().strftime('%H:%M:%S')}] {msg}")

def fetch_url(url, post_body=None, retries=3, timeout=LIST_TIMEOUT):
    headers = {
        "User-Agent": "Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36",
        "Referer": BASE_URL + "/news-list-ggtz.html",
        "X-Requested-With": "XMLHttpRequest",
        "Content-Type": "application/x-www-form-urlencoded; charset=UTF-8",
    }
    for attempt in range(retries):
        try:
            data = post_body.encode() if post_body else None
            req = urllib.request.Request(url, data=data, headers=headers)
            with urllib.request.urlopen(req, timeout=timeout) as resp:
                return resp.read().decode('utf-8', errors='replace')
        except Exception as e:
            if attempt < retries - 1:
                time.sleep(2)
    return None

def fetch_list(page):
    """获取列表页, 返回 (items, total)"""
    url = f"{BASE_URL}/api-ajax_list-{page}.html"
    body = "ajax_type%5B%5D=13_news&ajax_type%5B%5D=20&ajax_type%5B%5D=13&ajax_type%5B%5D=news&ajax_type%5B%5D=Y-m-d&ajax_type%5B%5D=40&ajax_type%5B%5D=20&ajax_type%5B%5D=0&ajax_type%5B%5D=&is_ds=1"
    resp = fetch_url(url, body)
    if not resp:
        return [], 0
    try:
        j = json.loads(resp)
        return j.get("data", []), j.get("total", 0)
    except:
        return [], 0

def fetch_detail(news_id):
    """获取详情页正文 (div.xxq_essay)"""
    url = f"{BASE_URL}/news-show-{news_id}.html"
    html = fetch_url(url, timeout=DETAIL_TIMEOUT)
    if not html:
        return ""
    m = re.search(r'<div class="xxq_essay">(.*?)</div>\s*</div>', html, re.DOTALL)
    if m:
        return m.group(1).strip()
    m = re.search(r'class="xxq_essay">(.*?)</div>\s*</div>', html, re.DOTALL)
    if m:
        return m.group(1).strip()
    return ""

def main():
    import sqlite3
    
    conn = sqlite3.connect(SEARCH_DB, timeout=60)
    conn.execute("PRAGMA busy_timeout=10000")
    c = conn.cursor()
    
    # 1. 获取总数
    items, total_items = fetch_list(1)
    if not items:
        log("ERROR: 无法获取列表")
        return
    
    total_pages = (total_items // 20) + (1 if total_items % 20 else 0)
    log(f"共 {total_items} 条, {total_pages} 页")
    
    # 2. 收集所有列表中近3年的条目
    all_items = []
    for page in range(1, min(total_pages + 1, 3)):  # max 2pg
        items, _ = fetch_list(page)
        if not items:
            break
        all_items.extend(items)
        if page % 10 == 0:
            log(f"  列表第{page}/{total_pages}页 ({len(all_items)}条)")
        # 检查最旧条目是否已过截止日期
        oldest = items[-1].get("inputtime", "")[:10]
        if oldest and oldest < CUTOFF_DATE:
            log(f"  第{page}页已达截止日期{oldest}，停止")
            break
        time.sleep(0.3)
    
    log(f"列表共 {len(all_items)} 条")
    
    # 3. 过滤日期 + 去重
    to_fetch = []
    skipped_exist = 0
    skipped_date = 0
    for item in all_items:
        inputtime = item.get("inputtime", "").strip()[:10]
        if inputtime and inputtime < CUTOFF_DATE:
            skipped_date += 1
            continue
        url = item.get("url", "")
        c.execute("SELECT COUNT(*) FROM gov_raw WHERE page_url = ?", (url,))
        if c.fetchone()[0] > 0:
            skipped_exist += 1
            continue
        to_fetch.append(item)
    
    log(f"新增: {len(to_fetch)}, 已存在: {skipped_exist}, 日期过滤: {skipped_date}")
    
    if not to_fetch:
        log("无新增条目")
        conn.close()
        return
    
    # 4. 并发获取详情
    def process_item(item):
        content = fetch_detail(item["id"])
        if content:
            return (item.get("title",""), content, item.get("url",""), item.get("inputtime","")[:10])
        return None
    
    log(f"并发获取{len(to_fetch)}条详情...")
    results = []
    with ThreadPoolExecutor(max_workers=MAX_WORKERS) as executor:
        futures = {executor.submit(process_item, item): item for item in to_fetch}
        done = 0
        for f in as_completed(futures):
            done += 1
            results.append(f.result())
            if done % 100 == 0:
                log(f"  {done}/{len(to_fetch)} 完成")
    
    # 5. 分批插入
    insert_data = [(SITE_NAME, r[0], r[1], r[2], r[3]) for r in results if r]
    empty_count = len(results) - len(insert_data)
    
    log(f"入库 {len(insert_data)} 条 (空正文 {empty_count} 条)...")
    BATCH_SIZE = 50
    inserted = 0
    for i in range(0, len(insert_data), BATCH_SIZE):
        batch = insert_data[i:i+BATCH_SIZE]
        c.executemany(
            "INSERT OR IGNORE INTO gov_raw(site_name, title, content, page_url, publish_date) VALUES(?,?,?,?,?)",
            batch
        )
        conn.commit()
        inserted += len(batch)
    
    log(f"入库完成: {inserted} 条")
    
    # 6. 重建FTS
    try:
        c.execute("SELECT 1 /* noop: gov_search 由触发器维护, 无需 rebuild */")
        conn.commit()
        log("FTS重建完成")
    except Exception as e:
        log(f"FTS重建错误: {e}")
    
    conn.close()
    log(f"[DONE] 达州爬取完成: +{inserted} 新增")

if __name__ == "__main__":
    main()
