#!/usr/bin/env python3
# -*- coding: utf-8 -*-
"""
ccpc_zc_sync.py — ZC + CCPC + CONTACT 单文件入库链（物理合并版, 2026-09-10）
==========================================================================

合并来源（4 引擎 + 编排器 → 本文件；业务逻辑逐字保留，仅重命名冲突符号）：
  Phase ZC      ← sync_zc.py       WebDAV Downloads_nuc 的 PDF → /root/ZC.db
                + parse_zc.py      PDF 解析库(extract_basic/extract_page2/extract_contacts)
                                   ★采用旧版(live)启发式：长标题被 PDF 折行时回退文件名，标题更完整
  Phase CCPC    ← parse_ccpc_v2.py WebDAV Downloads|下载1|Downloads2 的 DOCX/XLS → /root/ccpc.db
  Phase CONTACT ← contact_lib.py   ccpc.db+ZC.db 联系人聚合去重(同公司同人合并) → /root/contact.db

为什么不统一 WebDAV 传输层：
  ZC 段用 urllib+正则列目录（既有可用路径），CCPC 段用 requests+ElementTree。
  两者都稳定运行，强行统一等于重写其中一条 → 违背"零行为变化"原则，故各自保留、前缀区分。

2026-09-10 加固（合并同批）：
  - parse_contacts_from_text: list(set(...)) → sorted(set(...))  多手机号拼接顺序跨进程随机(PYTHONHASHSEED)
  - load_ccpc / load_zc: 加 ORDER BY（稳定键）→ 聚合严格确定性，contact.db 不再每日无谓重写
  - WebDAV 列表/下载统一走 _retry(3 次)；列表最终失败 → 该阶段 rc=1，不再静默当作 0 个文件

用法（保留原三引擎的单跑能力）：
  python3 ccpc_zc_sync.py                    # 全链 daily 增量
  python3 ccpc_zc_sync.py force              # CCPC 全量重跑（耗时）；ZC/CONTACT 语义不变
  python3 ccpc_zc_sync.py test               # CCPC 只跑 3 文件
  python3 ccpc_zc_sync.py --only zc          # 只跑 ZC
  python3 ccpc_zc_sync.py --only ccpc        # 只跑 CCPC
  python3 ccpc_zc_sync.py --only contact     # 只跑 CONTACT
  python3 ccpc_zc_sync.py --zc-preview X.pdf # 单 PDF 解析预览（原 parse_zc.py 单文件用法）

凭据：从 /root/gov_crawler/.env（权限 600）读取 WEBDAV_USER / WEBDAV_PASS，
      或同名环境变量（环境变量优先）。

crontab：
  30 2 * * * cd /root/gov_crawler && python3 ccpc_zc_sync.py >> /var/log/ccpc_zc_sync.log 2>&1
"""

import os, sys, re, json, time, base64, sqlite3, hashlib, urllib.parse, urllib.request, urllib.error
from datetime import datetime
from collections import OrderedDict
from xml.etree import ElementTree as ET

import requests
from requests.auth import HTTPBasicAuth
import fitz                                    # PyMuPDF
from openpyxl import load_workbook             # XLS/XLSX
from docx import Document                      # DOCX


# ══════════════════════════════════════════════════════════════════════
# §0  CONFIG
# ══════════════════════════════════════════════════════════════════════

WORKDIR    = "/root/gov_crawler"
ENV_FILE   = os.path.join(WORKDIR, ".env")

# ── ZC ──────────────────────────────────────────────────────────────
ZC_WEBDAV_BASE = "https://dav.jianguoyun.com/dav/Downloads_nuc"   # 旧 sync_zc 语义(含目录)
ZC_PDF_DIR     = "/mnt/data/zc_pdfs"
ZC_DB          = "/root/ZC.db"
ZC_LOG         = "/var/log/zc_sync.log"

# ── CCPC ────────────────────────────────────────────────────────────
CCPC_WEBDAV_HOST  = "https://dav.jianguoyun.com"
CCPC_DOWNLOAD_DIRS = ['/dav/Downloads/', '/dav/%E4%B8%8B%E8%BD%BD1/', '/dav/Downloads2/']
CCPC_LOCAL_CACHE   = '/data/ccpc_source'
CCPC_DB            = '/root/ccpc.db'

# ── CONTACT ─────────────────────────────────────────────────────────
CONTACT_DB = "/root/contact.db"


def load_env(path=ENV_FILE):
    """极简 .env 读取（不引第三方依赖）；KEY=VALUE，# 注释。"""
    cfg = {}
    if os.path.exists(path):
        with open(path, encoding='utf-8') as fh:
            for line in fh:
                line = line.strip()
                if not line or line.startswith('#') or '=' not in line:
                    continue
                k, v = line.split('=', 1)
                cfg[k.strip()] = v.strip().strip('"').strip("'")
    return cfg


_ENV = load_env()
WEBDAV_USER = os.environ.get('WEBDAV_USER') or _ENV.get('WEBDAV_USER', '')
WEBDAV_PASS = os.environ.get('WEBDAV_PASS') or _ENV.get('WEBDAV_PASS', '')

_ccpc_auth = None


def _auth():
    """CCPC 段 requests 用的 HTTPBasicAuth（懒建，避免 import 期副作用）。"""
    global _ccpc_auth
    if _ccpc_auth is None:
        _ccpc_auth = HTTPBasicAuth(WEBDAV_USER, WEBDAV_PASS)
    return _ccpc_auth


def check_creds():
    if not WEBDAV_USER or not WEBDAV_PASS:
        sys.stderr.write(
            "[FATAL] 缺少 WebDAV 凭据。请在 %s 写入：\n"
            "        WEBDAV_USER=...\n        WEBDAV_PASS=...\n"
            "        （或导出同名环境变量）\n" % ENV_FILE)
        sys.exit(2)


# ══════════════════════════════════════════════════════════════════════
# §1  公共工具
# ══════════════════════════════════════════════════════════════════════

def log(msg, logfile=None):
    """logfile 给定 → 带时间戳落盘 + 裸 print（原 sync_zc.log 语义）；
    未给 → 带时间戳 print（原 parse_ccpc_v2.log 语义）。"""
    ts = datetime.now().strftime('%Y-%m-%d %H:%M:%S')
    if logfile:
        with open(logfile, 'a') as f:
            f.write(f"[{ts}] {msg}\n")
        print(msg)
    else:
        print(f"[{ts}] {msg}")


def sq(s):
    """SQL 单引号转义（原 sync_zc.sq / parse_zc.sq 同实现，合并为一处）。"""
    return "'" + str(s).replace("'", "''") + "'"


class PermanentError(Exception):
    """明确不该重试的错误（如 HTTP 4xx：资源不存在/名字非法，重试也白搭）"""


def _retry(fn, attempts=3, delay=3, what=''):
    """网络调用重试：最多 attempts 次，最后一次仍失败则抛出。
    （2026-09-10 加固：坚果云 WebDAV 偶发 ConnectTimeout，旧代码一崩就整阶段失败。
      PermanentError 直接抛出，不做无谓重试）"""
    last = None
    for i in range(1, attempts + 1):
        try:
            return fn()
        except PermanentError:
            raise
        except Exception as e:
            last = e
            if i < attempts:
                log(f"  {what} 第{i}/{attempts}次失败: {e} — {delay}s 后重试")
                time.sleep(delay)
    raise last


# ══════════════════════════════════════════════════════════════════════
# §2  ZC 引擎   （源：sync_zc.py + parse_zc.py）
# ══════════════════════════════════════════════════════════════════════

def zc_webdav_list():
    """列 ZC 源目录 PDF；网络失败重试 3 次，最终失败抛异常（由 zc_run 捕获→rc=1 可见）"""
    url = ZC_WEBDAV_BASE + "/"
    creds = base64.b64encode(f"{WEBDAV_USER}:{WEBDAV_PASS}".encode()).decode()

    def _do():
        req = urllib.request.Request(url, method="PROPFIND", headers={"Depth": "1"})
        req.add_header("Authorization", f"Basic {creds}")
        resp = urllib.request.urlopen(req, timeout=30)
        return resp.read().decode()

    xml_data = _retry(_do, what="ZC WebDAV 列表")
    pdfs = []
    for m in re.finditer(r'<d:href>([^<]+\.pdf)</d:href>', xml_data):
        href = m.group(1)
        fname = href.rstrip('/').split('/')[-1]
        fname = urllib.parse.unquote(fname)
        pdfs.append(fname)
    return sorted(pdfs)


def zc_webdav_download(filename):
    url = ZC_WEBDAV_BASE + "/" + urllib.parse.quote(filename)
    creds = base64.b64encode(f"{WEBDAV_USER}:{WEBDAV_PASS}".encode()).decode()

    def _do():
        req = urllib.request.Request(url)
        req.add_header("Authorization", f"Basic {creds}")
        try:
            resp = urllib.request.urlopen(req, timeout=60)
        except urllib.error.HTTPError as he:
            if 400 <= he.code < 500:
                raise PermanentError(f"HTTP {he.code}")   # 不重试
            raise
        data = resp.read()
        path = os.path.join(ZC_PDF_DIR, filename)
        with open(path, 'wb') as f:
            f.write(data)
        return path

    try:
        return _retry(_do, what=f"下载 {filename[:30]}")
    except Exception as e:
        log(f"  下载失败 {filename}: {e}", ZC_LOG)
        return None


# ── Page 1 字段提取（标签: 值 格式）──────────────────────────────────

def extract_basic(text):
    """提取 page 1 基础字段"""
    p = {}
    lines = text.split('\n')

    field_map = {
        '项目编号': 'project_id',
        '版本类型': 'version_type',
        '发布时间': 'publish_date',
        '项目阶段': 'phase',
        '建设周期': 'construction_period',
        '总投资额': 'total_investment',
        '工程类型': 'project_type',
        '甲方性质': 'owner_nature',
        '所属行业': 'industry',
        '所属专题': 'topic',
        '项目规模': 'scale',
        '数量规模': 'quantity_scale',
        '行业级别': 'industry_level',
        '建筑面积': 'building_area',
        '占地面积': 'land_area',
        '建筑层数': 'floors',
        '外资参与': 'foreign_investment',
        '装修': 'decoration',
        '钢结构': 'steel_structure',
        '外墙材料': 'exterior_wall',
        '车库停车位': 'parking',
        '电梯': 'elevator',
        '空调': 'air_conditioning',
        '新风系统': 'fresh_air',
        '供暖方式': 'heating',
        '装配式建筑': 'prefab',
        '被动房': 'passive_house',
    }

    for i, l in enumerate(lines):
        for label, key in field_map.items():
            if l.startswith(label + '：') or l.startswith(label + ':'):
                val = l[len(label)+1:].strip()
                # 如果值为空且下一行不为空标签，取下一行（行业字段特殊处理）
                if not val and key != 'industry' and i + 1 < len(lines):
                    next_l = lines[i+1].strip()
                    if next_l and '：' not in next_l and ':' not in next_l:
                        val = next_l
                # 所属行业可能跨多行（如"其他仓储物流/石化/煤化工/基础\n化工"）
                if key == 'industry':
                    vals = []
                    n = 1
                    while i + n < len(lines):
                        next_l = lines[i + n].strip()
                        if '：' in next_l or ':' in next_l or len(next_l) < 2 or next_l.startswith('项目'):
                            break
                        vals.append(next_l)
                        n += 1
                    if vals:
                        val = ''.join(vals)
                p[key] = val
                break

    # 省/市（特殊格式：四川成都市 → 四川/成都市）
    for i, l in enumerate(lines):
        if l.startswith('省/市') or l.startswith('省／市'):
            val = l.split('：', 1)[-1].strip() if '：' in l else l.split(':', 1)[-1].strip()
            # 尝试找省份和城市：如"山东菏泽市"、"内蒙古包头市"
            m = re.match(r'^(黑龙江|内蒙古|广西|西藏|新疆|宁夏|青海|甘肃|四川|贵州|云南|陕西|山西|河北|山东|河南|湖北|湖南|江苏|浙江|安徽|江西|福建|广东|海南|辽宁|吉林|上海|北京|天津|重庆)(.*)$', val)
            if m:
                p['province'] = m.group(1).strip()
                p['city'] = m.group(2).strip()
            elif '/' in val:
                parts = val.split('/', 1)
                p['province'] = parts[0].strip()
                p['city'] = parts[1].strip()
            else:
                p['province'] = val
                p['city'] = ''

    # 详细地址
    for l in lines:
        if l.startswith('详细地址') and ('：' in l or ':' in l):
            p['detail_address'] = l.split('：', 1)[-1].strip() if '：' in l else l.split(':', 1)[-1].strip()

    # 建设内容描述（多行）
    desc_lines = []
    in_desc = False
    for l in lines:
        if l.startswith('建设内容描述') and ('：' in l or ':' in l):
            in_desc = True
            rest = l.split('：', 1)[-1].strip() if '：' in l else l.split(':', 1)[-1].strip()
            if rest:
                desc_lines.append(rest)
            continue
        if in_desc:
            if l.startswith('该项目可能') or l.startswith('温馨提示'):
                break
            if l.strip():
                desc_lines.append(l.strip())
    if desc_lines:
        p['construction_content'] = ''.join(desc_lines)

    # 设备（温馨提示之前的行）
    for l in lines:
        if l.startswith('该项目可能') or l.startswith('主要设备'):
            val = l.split('：', 1)[-1].strip() if '：' in l else l.split(':', 1)[-1].strip()
            # 去掉尾部温馨提示
            val = re.sub(r'温馨提示.*', '', val).strip()
            p['equipment_list'] = val
            break

    # 项目名称：从文件路径取或从页眉取
    # page header 通常是第一行
    header = lines[0].strip() if lines else ''
    if header and len(header) > 10 and '项目' not in header:
        p['project_name'] = header
    elif len(lines) > 1:
        # 有时标题在第二行
        header2 = lines[1].strip()
        if header2 and len(header2) > 10:
            p['project_name'] = header2

    return p


def extract_page2(text):
    """提取 page 2 字段"""
    p = {}
    lines = text.split('\n')

    for i, l in enumerate(lines):
        # 精准采购设备
        m = re.search(r'精准采购设备[：:]?\s*(.*)', l)
        if m:
            p['procurement_equipment'] = m.group(1).strip()
            continue

        # 工艺流程
        m = re.search(r'工艺流程[：:]?\s*(.*)', l)
        if m:
            p['process_flow'] = m.group(1).strip()
            continue

        # 工期概述
        m = re.search(r'工期概述[：：]?\s*(.*)', l)
        if m:
            p['schedule_overview'] = m.group(1).strip()
            continue

        # 立项审批/项目设计等阶段状态
        for phase_label, key in [
            ('立项审批', 'phase_approval'),
            ('项目设计', 'phase_design'),
            ('主设备材料采购', 'phase_procurement'),
            ('主体施工', 'phase_construction'),
            ('工程分包', 'phase_contractor'),
            ('暂存/取消/已完工', 'phase_status'),
        ]:
            if l.startswith(phase_label) and ('：' in l or ':' in l or len(l.strip()) == len(phase_label)):
                if len(l.strip()) == len(phase_label):
                    # 值在下一行
                    val = lines[i + 1].strip() if i + 1 < len(lines) else ''
                else:
                    val = l.split('：', 1)[-1].strip() if '：' in l else l.split(':', 1)[-1].strip()
                p[key] = val
                break

    return p


# ── 联系人提取（可选文字，无需OCR）──────────────────────────────────

def extract_contacts(text):
    """从 '业主方' / '设计院' 等多个区域提取联系人"""
    contacts = []

    CONTACT_SECTIONS = ["业主方", "设计院", "施工单位", "施工方", "业主单位", "设计单位", "承包商", "承包方"]

    def parse_one_section(blk, default_role=""):
        """解析一个联系人区域块"""
        lines = blk.split('\n')
        # 首行为区域名
        section_title = lines[0].strip() if lines else ""
        # 角色映射
        role_map = {"业主方": "业主", "设计院": "设计院", "施工单位": "施工单位", "施工方": "施工单位", "业主单位": "业主", "设计单位": "设计院", "承包商": "施工单位", "承包方": "施工单位"}
        role = role_map.get(section_title, default_role)

        # 提取单位名称
        company = ''
        for line in lines:
            if '单位名称' in line and ('：' in line or ':' in line):
                # 格式: 单位名称：[主体承建商]XXX公司(私营)
                val = line.split('：', 1)[-1].strip() if '：' in line else line.split(':', 1)[-1].strip()
                # 去掉尾部 (私营) (外资) (国有/集体所有) 等
                val = re.sub(r'\s*[（(].*?[）)]$', '', val).strip()
                # 去掉 [业主] 前缀（role已标记），保留 [施工图设计] [主体承建商] 等
                val = re.sub(r'^\[业主\]', '', val).strip()
                company = val
                break

        # 每个联系人以 '姓名：' 开头
        person_blocks = re.split(r'\n(?=姓名[：:])', blk)

        results = []
        for pblk in person_blocks:
            plines = pblk.split('\n')
            c = {'company': company, 'contact_name': '', 'department': '',
                 'position': '', 'phone': '', 'remarks': '', 'address': '', 'role': role}
            for l in plines:
                l = l.strip()
                if not l:
                    continue
                if l.startswith('姓名') and ('：' in l):
                    name = l.split('：', 1)[-1].strip()
                    c['contact_name'] = name
                elif l.startswith('部门') and ('：' in l):
                    c['department'] = l.split('：', 1)[-1].strip()
                elif l.startswith('职务') and ('：' in l):
                    c['position'] = l.split('：', 1)[-1].strip()
                elif l.startswith('手机') and ('：' in l):
                    c['phone'] = l.split('：', 1)[-1].strip()
                elif l.startswith('备注') and ('：' in l):
                    c['remarks'] = l.split('：', 1)[-1].strip()
                elif l.startswith('单位注册地址') and ('：' in l):
                    c['address'] = l.split('：', 1)[-1].strip()
            # 跳过区域标题块本身（没有姓名的）
            if c['contact_name'] or c['phone']:
                results.append(c)
        return results

    # 找所有区域位置
    positions = {}
    for sec in CONTACT_SECTIONS:
        idx = text.find('\n' + sec + '\n')
        if idx < 0:
            idx = text.find(sec)
        if idx >= 0:
            positions[sec] = idx

    if not positions:
        return contacts

    # 按文本顺序排序
    sections = sorted(positions.items(), key=lambda x: x[1])

    # 解析每个区域，边界到下一个区域或文本结尾
    for i, (sec_name, sec_pos) in enumerate(sections):
        # 边界：从 sec_pos 到 next section 或 "Powered by TCPDF"
        next_pos = len(text)
        if i + 1 < len(sections):
            next_pos = sections[i + 1][1]
        else:
            # 最后一个区域，到 Powered by 或文本结尾
            pbr = text.find('Powered by TCPDF', sec_pos)
            if pbr >= 0:
                next_pos = pbr

        blk = text[sec_pos:next_pos].strip()
        if blk:
            contacts.extend(parse_one_section(blk))

    return contacts


# ── SQL 生成（供 --zc-preview / 批量导出用）────────────────────────

def gen_sql(p, contacts, filename):
    pid = p.get('project_id', '')
    name = p.get('project_name', '') or os.path.splitext(filename)[0]

    fields = [
        'project_id', 'project_name', 'version_type', 'publish_date', 'phase',
        'construction_period', 'total_investment', 'project_type', 'owner_nature',
        'industry', 'topic', 'scale', 'quantity_scale', 'industry_level',
        'province', 'city', 'detail_address', 'building_area', 'land_area',
        'floors', 'foreign_investment', 'decoration', 'steel_structure',
        'exterior_wall', 'parking', 'elevator', 'air_conditioning', 'fresh_air',
        'heating', 'prefab', 'passive_house', 'construction_content',
        'equipment_list', 'procurement_equipment', 'process_flow',
        'schedule_overview', 'phase_approval', 'phase_design',
        'phase_procurement', 'phase_construction', 'phase_contractor',
        'phase_status',
    ]

    vals = [sq(p.get(f, '')) for f in fields]
    sql = f"INSERT INTO zc_projects ({', '.join(fields)}) VALUES ({', '.join(vals)});\n"

    for c in contacts:
        sql += (
            f"INSERT INTO zc_contacts "
            f"(project_id, company, contact_name, department, position, phone, remarks, address, role) "
            f"VALUES ({sq(pid)}, {sq(c['company'])}, {sq(c['contact_name'])}, "
            f"{sq(c['department'])}, {sq(c['position'])}, {sq(c['phone'])}, "
            f"{sq(c['remarks'])}, {sq(c['address'])}, {sq(c.get('role', ''))});\n"
        )

    return sql


# ── 解析单个 PDF ───────────────────────────────────────────────────

def zc_parse_one(pdf_path):
    doc = fitz.open(pdf_path)
    text_p1 = doc[0].get_text() if doc.page_count > 0 else ''
    text_p2 = doc[1].get_text() if doc.page_count > 1 else ''
    text_contacts = ''
    for pi in range(1, doc.page_count):
        text_contacts += doc[pi].get_text() + '\n'
    doc.close()

    p = extract_basic(text_p1)
    p.update(extract_page2(text_p2))

    if not p.get('project_name'):
        fname = os.path.basename(pdf_path)
        p['project_name'] = os.path.splitext(fname)[0]

    contacts = extract_contacts(text_contacts)
    return p, contacts


# ── ZC 入库 ────────────────────────────────────────────────────────

def zc_check_exists(c, pid, name):
    """查 DB：project_id 或 project_name 是否已存在"""
    if pid:
        r = c.execute("SELECT 1 FROM zc_projects WHERE project_id=?", (pid,)).fetchone()
        if r:
            return True
    if name:
        r = c.execute("SELECT 1 FROM zc_projects WHERE project_name=?", (name,)).fetchone()
        if r:
            return True
    return False


def zc_insert_to_db(p, contacts, filename):
    conn = sqlite3.connect(ZC_DB)
    c = conn.cursor()
    pid = p.get('project_id', '')
    name = p.get('project_name', '') or os.path.splitext(filename)[0]
    if zc_check_exists(c, pid, name):
        conn.close()
        return False
    fields = ['project_id', 'project_name', 'version_type', 'publish_date', 'phase',
              'construction_period', 'total_investment', 'project_type', 'owner_nature',
              'industry', 'topic', 'scale', 'quantity_scale', 'industry_level',
              'province', 'city', 'detail_address', 'building_area', 'land_area',
              'floors', 'foreign_investment', 'decoration', 'steel_structure',
              'exterior_wall', 'parking', 'elevator', 'air_conditioning', 'fresh_air',
              'heating', 'prefab', 'passive_house', 'construction_content',
              'equipment_list', 'procurement_equipment', 'process_flow',
              'schedule_overview', 'phase_approval', 'phase_design',
              'phase_procurement', 'phase_construction', 'phase_contractor', 'phase_status',
              'raw_json']
    raw_data = dict(p)
    raw_data['contacts'] = contacts
    p['raw_json'] = json.dumps(raw_data, ensure_ascii=False)
    vals = [sq(p.get(f, '')) for f in fields]
    sql = f"INSERT INTO zc_projects ({', '.join(fields)}) VALUES ({', '.join(vals)})"
    c.execute(sql)
    for ct in contacts:
        c.execute(
            "INSERT INTO zc_contacts (project_id, company, contact_name, department, position, phone, remarks, address, role) VALUES (?,?,?,?,?,?,?,?,?)",
            (pid, ct.get('company',''), ct.get('contact_name',''), ct.get('department',''),
             ct.get('position',''), ct.get('phone',''), ct.get('remarks',''),
             ct.get('address',''), ct.get('role',''))
        )
    conn.commit()
    conn.close()
    return True


def zc_run():
    """返回 0 成功 / 1 失败（列表拿不到=失败，不再静默当成 0 个文件）"""
    os.makedirs(ZC_PDF_DIR, exist_ok=True)
    log("=== 开始 ZC 同步 ===", ZC_LOG)

    try:
        pdfs = zc_webdav_list()
    except Exception as e:
        log(f"!!! ZC 源目录列表失败(已重试 3 次): {e}", ZC_LOG)
        log("=== 中止 ZC ===", ZC_LOG)
        return 1
    log(f"WebDAV 共有 {len(pdfs)} 个 PDF", ZC_LOG)

    new_count = 0
    skip_count = 0
    for fname in pdfs:
        # 已下载则跳过下载，直接解析
        local_path = os.path.join(ZC_PDF_DIR, fname)
        if os.path.exists(local_path):
            pdf_path = local_path
        else:
            pdf_path = zc_webdav_download(fname)
            if not pdf_path:
                continue

        try:
            p, contacts = zc_parse_one(pdf_path)
            inserted = zc_insert_to_db(p, contacts, fname)
            if inserted:
                new_count += 1
                log(f"  ✅ {fname[:50]} (联系人: {len(contacts)})", ZC_LOG)
            else:
                skip_count += 1
        except Exception as e:
            log(f"  ✗ {fname[:50]}: {e}", ZC_LOG)

    count = sqlite3.connect(ZC_DB).execute("SELECT COUNT(*) FROM zc_projects").fetchone()[0]
    log(f"新增: {new_count}, 跳过: {skip_count}, ZC.db 总计: {count} 条", ZC_LOG)
    log("=== 完成 ===", ZC_LOG)
    return 0


def zc_preview(pdf_path):
    """单 PDF 解析预览（原 parse_zc.py 单文件用法）"""
    print(f"📄 {os.path.basename(pdf_path)}")
    p, contacts = zc_parse_one(pdf_path)

    print("\n── 项目信息 ──")
    for k, v in p.items():
        if v:
            output = str(v)[:100] + '...' if len(str(v)) > 100 else str(v)
            print(f"  {k}: {output}")

    print("\n── 联系人 ──")
    if contacts:
        for i, c in enumerate(contacts, 1):
            print(f"  [{i}]")
            for k, v in c.items():
                if v:
                    print(f"    {k}: {v}")
            print()
    else:
        print("  (无)")

    print("\n── SQL ──")
    print(gen_sql(p, contacts, os.path.basename(pdf_path)))


# ══════════════════════════════════════════════════════════════════════
# §3  CCPC 引擎   （源：parse_ccpc_v2.py）
# ══════════════════════════════════════════════════════════════════════

# ── 字段映射表（中 → 英）──────────────────────────────────────────

FIELD_MAP = OrderedDict([
    # 基础字段
    ('项目名称',    'project_name'),
    ('项目类型',    'project_type'),
    ('项目ID',      'project_id'),
    ('项目所属行业', 'industry'),
    ('所属领域类型',  'field_type'),
    ('所属省份',    'province'),
    ('所属地级市',   'city'),
    ('进展阶段',    'phase'),
    ('发布时间',    'publish_date'),
    ('跟踪版本号',   'version'),
    ('项目性质',    'nature'),
    ('预算投资总额(万元)', 'budget'),
    ('投资性质',    'invest_nature'),
    ('资金到位情况',  'funding'),
    ('建设等级',    'grade'),
    ('预计开建时间',  'start_time'),
    ('预计截止时间',  'end_time'),
    ('设备来源',    'equipment_source'),
    ('建筑面积',    'building_area'),
    ('占地面积',    'land_area'),
    ('有无钢结构',   'steel_structure'),
    ('建筑层数',    'building_floors'),
    ('供暖方式',    'heating_method'),
    ('外墙材料',    'wall_material'),
    ('装修',       'decoration'),
    ('有无空调',    'has_ac'),
    ('有无立体停车位', 'has_parking'),
    ('有无电梯',    'has_elevator'),
    ('项目所在地',   'project_address'),
    ('专题标签',    'tags'),
    ('工艺流程',    'process_flow'),
    ('主体施工阶段',  'phase_construction'),
    ('主体设计阶段',  'phase_design'),
    ('商机描述',    'opportunity'),
    ('关键工程',    'key_works'),
    ('具体需求情况',  'procurement_needs'),
    ('业主单位',    'owner_company'),
    # DOCX 跟踪版本字段
    ('项目标题',    'project_name'),
    ('更新版本',    'version'),
    ('更新时间',    'publish_date'),
    ('所属栏目',    'project_type'),
])

ALL_PROJECT_FIELDS = [
    'project_name', 'project_type', 'project_id', 'industry', 'field_type',
    'province', 'city', 'phase', 'publish_date', 'version',
    'nature', 'budget', 'invest_nature', 'funding', 'grade',
    'start_time', 'end_time', 'equipment_source', 'building_area', 'land_area',
    'steel_structure', 'equipment_list', 'progress', 'opportunity', 'key_works',
    'overview', 'detail', 'owner_company', 'contacts', 'raw_json',
    'building_floors', 'heating_method', 'wall_material', 'decoration',
    'has_ac', 'has_parking', 'has_elevator', 'project_address', 'tags',
    'process_flow', 'phase_construction', 'phase_design', 'procurement_equipment',
    'procurement_needs',
]

# ── 工具 ───────────────────────────────────────────────────────────

def _clean_company_name(name):
    """清理公司名中的所有权标记 (私营)/(国有)/(集体) 及尾部附属备注。"""
    if not name:
        return name
    cleaned = re.sub(r"\s*[（(](?:私营|国有(?:/集体所有)?|集体(?:所有)?)[)）]\s*(?:[（(][^)）]*[)）])?\s*$", "", name)
    cleaned = re.sub(r"\s*[（(][^)）]*$", "", cleaned)
    cleaned = re.sub(r"\s*[(（][^)）]*[)）]\s*$", "", cleaned)
    return cleaned.strip()


def _apply_field(record, cn_key, val):
    """通过 FIELD_MAP 将中文键映射为英文字段写入 record。"""
    en = FIELD_MAP.get(cn_key)
    if en and not record.get(en):
        val = val.strip()
        if en == "owner_company":
            val = _clean_company_name(val)
        record[en] = val


# ── WebDAV 工具（requests 版，原 parse_ccpc_v2 语义）───────────────

def ccpc_webdav_list(remote_dir):
    def _do():
        return requests.request('PROPFIND', CCPC_WEBDAV_HOST + remote_dir, auth=_auth(),
                                headers={'Depth': '1'}, timeout=30)
    r = _retry(_do, what=f"CCPC 列表 {remote_dir}")
    if r.status_code != 207:
        return []
    root = ET.fromstring(r.content)
    ns = {'d': 'DAV:'}
    files = []
    for resp in root.findall('.//d:response', ns):
        href_el = resp.find('d:href', ns)
        if href_el is None:
            continue
        fn = (href_el.text or '').rstrip('/').rsplit('/', 1)[-1]
        fn = urllib.parse.unquote(fn)
        if fn:
            files.append(fn)
    return files


def ccpc_webdav_download(remote_base, filename, local_path):
    # 文件名可能含 # { } 等特殊字符（如 httpsvip.ccpc360.com#trackcompany-…docx）；
    # 不编码时 '#' 会被当作 URL fragment 截断 → 永久 404（历史每日 err1 的根因）
    url = CCPC_WEBDAV_HOST + remote_base + urllib.parse.quote(filename)

    def _do():
        r = requests.get(url, auth=_auth(), timeout=60)
        if r.status_code == 200:
            with open(local_path, 'wb') as f:
                f.write(r.content)
            return True
        if 400 <= r.status_code < 500:
            raise PermanentError(f"HTTP {r.status_code}")   # 不重试
        raise RuntimeError(f"HTTP {r.status_code}")

    try:
        return _retry(_do, what=f"CCPC 下载 {filename[:30]}")
    except Exception:
        return False


# ── 联系人解析 ─────────────────────────────────────────────────────

PHONE_RE = re.compile(r'1[3-9]\d{9}')
EMAIL_RE = re.compile(r'[a-zA-Z0-9._%+-]+@[a-zA-Z0-9.-]+\.[a-zA-Z]{2,}')
NAME_RE  = re.compile(r'联系人[：:]\s*([\u4e00-\u9fff·]{2,6})')


def parse_contacts_from_text(text_blob, role_label=''):
    """从联系人文本块解析为 [dict] 列表。"""
    contacts = []
    if not text_blob or not text_blob.strip():
        return contacts
    lines = text_blob.strip().split('\n')
    blocks = []
    cur = []
    for line in lines:
        line = line.strip()
        if not line:
            continue
        is_company_line = (len(line) > 4 and
                          re.search(r'[（(](业主|设计|施工|法人|项目|股东|总集团|技术|采购|施工图纸)[^)）]*[)）]', line) or
                          '有限公司' in line or '公司' in line or '厂' in line or '院' in line) and \
                          not line.startswith(('联系人', '联系方式', '手机', '地址', '电话', '姓名', '部门', '邮箱', 'Email', 'email'))
        if is_company_line and cur:
            blocks.append('\n'.join(cur))
            cur = [line]
        else:
            cur.append(line)
    if cur:
        blocks.append('\n'.join(cur))
    for block in blocks:
        block = block.strip()
        if not block or '仅供参考' in block or '版权所有' in block:
            continue
        lines_b = block.split('\n')
        first_line = lines_b[0].strip()
        company = re.sub(r'\s*[（(][^)）]*[)）]\s*$', '', first_line).strip()
        if not company:
            continue
        rest = '\n'.join(lines_b[1:])
        names = NAME_RE.findall(block)
        phones = sorted(set(PHONE_RE.findall(block)))
        emails = sorted(set(EMAIL_RE.findall(block)))   # sorted: set 迭代顺序跨进程随机(PYTHONHASHSEED)
        landline_num = ''
        m_ll = re.search(r'(?:固定电话|座机|landline)[：:]\s*(\d[\d\- ]{5,15})', block)
        if m_ll:
            landline_num = m_ll.group(1).strip()
        addr = ''
        m = re.search(r'地址[：:]\s*(.+)', block)
        if m:
            addr = m.group(1).strip()
        if names:
            for name in names:
                contacts.append({
                    'company': company,
                    'contact_name': name,
                    'phone': ' / '.join(phones) if phones else '',
                    'landline': landline_num,
                    'email': ' / '.join(emails) if emails else '',
                    'address': addr,
                    'role': role_label,
                })
        elif phones:
            contacts.append({
                'company': company,
                'contact_name': '',
                'phone': ' / '.join(phones),
                'landline': landline_num,
                'email': ' / '.join(emails) if emails else '',
                'address': addr,
                'role': role_label,
            })
    return contacts


CONTACT_SECTION_MAP = {
    '业主单位联系人以及联系方式': '业主单位',
    '设计院联系人以及联系方式': '设计院',
    '施工单位联系人以及联系方式': '施工单位',
    '参考单位（仅供参考）': '参考联系人',
}


# ── DOCX 解析 ──────────────────────────────────────────────────────

SECTION_NAMES = {
    '目录', '项目跟踪版本', '项目基础信息',
    '所需材料设备', '设备采购情况', '项目概况',
    '项目进展', '项目详情',
    '业主单位联系人以及联系方式', '设计院联系人以及联系方式',
    '施工单位联系人以及联系方式', '参考单位（仅供参考）',
}


def parse_docx(filepath):
    """解析 DOCX → (record dict, contacts list)"""
    doc = Document(filepath)
    record = {}
    detail_parts = []
    all_contacts = []
    contact_sections = {}
    all_sections_raw = {}

    for table in doc.tables:
        rows_data = []
        section_name = None

        def flush_section():
            nonlocal section_name, rows_data
            if not section_name or not rows_data:
                return
            if section_name == '项目跟踪版本':
                if len(rows_data) >= 2:
                    hdr = rows_data[0]
                    lines = []
                    for row in rows_data[1:]:
                        for i in range(0, min(len(hdr), len(row))):
                            k = hdr[i].strip()
                            v = row[i].strip()
                            if k and v:
                                lines.append(f'{k}：{v}')
                                _apply_field(record, k, v)
                    if lines:
                        all_sections_raw[section_name] = '\n'.join(lines)

            elif section_name == '项目基础信息':
                kv_lines = []
                for rcells in rows_data:
                    deduped = []
                    for c in rcells:
                        if not deduped or c != deduped[-1]:
                            deduped.append(c)
                    for i in range(0, len(deduped), 2):
                        if i + 1 < len(deduped):
                            k, v = deduped[i].strip(), deduped[i+1].strip()
                            if k and v:
                                _apply_field(record, k, v)
                                if k in FIELD_MAP:
                                    kv_lines.append(f'{k}：{v}')
                if kv_lines:
                    txt = '\n'.join(kv_lines)
                    all_sections_raw[section_name] = txt

            elif section_name == '项目详情':
                v = rows_data[0][0].strip() if rows_data and rows_data[0] else ''
                if v and '版权所有' not in v:
                    all_sections_raw[section_name] = v
                    if v not in (record.get('detail') or ''):
                        detail_parts.append(v)
                    if not record.get('detail'):
                        record['detail'] = v

            elif section_name == '项目概况':
                v = rows_data[0][0].strip() if rows_data and rows_data[0] else ''
                if v and '版权所有' not in v:
                    all_sections_raw[section_name] = v
                    if not record.get('overview'):
                        record['overview'] = v

            elif section_name in ('所需材料设备', '设备采购情况'):
                v = rows_data[0][0].strip() if rows_data and rows_data[0] else ''
                if v and '版权所有' not in v:
                    all_sections_raw[section_name] = v
                    if not record.get('equipment_list'):
                        record['equipment_list'] = v

            elif section_name == '项目进展':
                v = rows_data[0][0].strip() if rows_data and rows_data[0] else ''
                if v and '版权所有' not in v:
                    all_sections_raw[section_name] = v
                    if not record.get('progress'):
                        record['progress'] = v

            elif section_name in CONTACT_SECTION_MAP:
                v = rows_data[0][0].strip() if rows_data and rows_data[0] else ''
                if v and '版权所有' not in v and '仅供参考' not in v:
                    role = CONTACT_SECTION_MAP[section_name]
                    contact_sections[section_name] = v
                    all_sections_raw[section_name] = v
                    all_contacts.extend(parse_contacts_from_text(v, role))

        for ri, row in enumerate(table.rows):
            cells = [c.text.strip() for c in row.cells]
            all_same = cells and all(c == cells[0] for c in cells)
            ct = cells[0] if cells else ''
            is_hdr = all_same and (ct in SECTION_NAMES or (ri > 0 and ct and len(ct) < 25))

            if is_hdr and ri > 0:
                flush_section()
                rows_data = []
                section_name = ct
            elif not is_hdr:
                rows_data.append(cells)

        # Last section
        flush_section()

    # Set detail
    if not record.get('detail'):
        record['detail'] = '\n\n'.join(detail_parts)
    elif detail_parts:
        # Combine existing detail with any extra sections not already captured
        # （确定性：用有序去重替代 set —— set 迭代顺序跨进程随机，会让 detail 拼接顺序漂移）
        existing = record['detail']
        extra = []
        for part in detail_parts:
            if part not in existing and part not in extra:
                extra.append(part)
        if extra:
            record['detail'] = existing + '\n\n' + '\n\n'.join(extra)

    # 项目名称
    if not record.get('project_name'):
        names = [p.text.strip() for p in doc.paragraphs if p.text.strip()]
        for n in names:
            if len(n) > 5 and not n.startswith('目录'):
                record['project_name'] = n
                break
    if not record.get('project_name'):
        fname = os.path.basename(filepath)
        record['project_name'] = re.sub(r'-\d+\.docx$', '', fname)

    # Build raw_json from ALL parsed data
    raw_dict = {}
    for f in ALL_PROJECT_FIELDS:
        if record.get(f):
            raw_dict[f] = record[f]
    for sn, sr in all_sections_raw.items():
        raw_dict[f'_section_{sn}'] = sr
    paragraphs = [p.text.strip() for p in doc.paragraphs if p.text.strip()]
    if paragraphs:
        raw_dict['_paragraphs'] = paragraphs
    record['raw_json'] = json.dumps(raw_dict, ensure_ascii=False, default=str)
    return record, all_contacts


# ── XLS 解析 ───────────────────────────────────────────────────────

def parse_xls(filepath):
    import shutil
    tmp = '/tmp/ccpc_xls_temp.xlsx'
    shutil.copy2(filepath, tmp)
    try:
        wb = load_workbook(tmp, read_only=True)
        ws = wb[wb.sheetnames[0]]
        headers, data = None, None
        for i, row in enumerate(ws.iter_rows(values_only=True)):
            if i == 0:
                headers = [str(v) if v else '' for v in row]
            elif i == 1:
                data = [str(v) if v else '' for v in row]
                break
        wb.close()
        if not headers or not data:
            return None, []
        raw = dict(zip(headers, data))
        record = {}
        contacts_raw = {}
        for cn, en in FIELD_MAP.items():
            val = raw.get(cn, '') or ''
            if cn in ('业主单位联系人以及联系方式', '设计院联系人以及联系方式',
                       '施工单位联系人以及联系方式', '参考联系人', '参考单位（仅供参考）'):
                contacts_raw[cn] = str(val)
            elif en and not en.startswith('_'):
                record[en] = str(val)
        if not record.get('project_name'):
            fname = os.path.basename(filepath)
            record['project_name'] = re.sub(r'-\d+\.(xls|xlsx)$', '', fname)
        if not record.get('detail'):
            raw_detail = raw.get('项目详情', '').strip()
            if raw_detail:
                record['detail'] = raw_detail
        record['raw_json'] = json.dumps(raw, ensure_ascii=False)
        all_contacts = []
        for cn, text in contacts_raw.items():
            role = CONTACT_SECTION_MAP.get(cn, '')
            if text and text.strip():
                all_contacts.extend(parse_contacts_from_text(text, role))
        return record, all_contacts
    finally:
        if os.path.exists(tmp):
            os.remove(tmp)


# ── CCPC 入库 ──────────────────────────────────────────────────────

def ccpc_upsert_project(record):
    db = sqlite3.connect(CCPC_DB)
    pid, pname = record.get('project_id', ''), record.get('project_name', '')
    existing = None
    if pid:
        existing = db.execute('SELECT id FROM cceup_projects WHERE project_id = ?', (pid,)).fetchone()
    if existing is None and pname:
        existing = db.execute('SELECT id FROM cceup_projects WHERE project_name = ?', (pname,)).fetchone()
    vals = [record.get(f, '') for f in ALL_PROJECT_FIELDS]
    if existing:
        set_c = ', '.join(f'{f}=?' for f in ALL_PROJECT_FIELDS)
        db.execute(f'UPDATE cceup_projects SET {set_c} WHERE id=?', vals + [existing[0]])
        proj_id = existing[0]
        action = 'UPDATE'
    else:
        ph = ', '.join(['?'] * len(ALL_PROJECT_FIELDS))
        db.execute(f'INSERT INTO cceup_projects ({",".join(ALL_PROJECT_FIELDS)}) VALUES ({ph})', vals)
        proj_id = db.execute('SELECT last_insert_rowid()').fetchone()[0]
        action = 'INSERT'
    db.commit()
    db.close()
    return action, proj_id


def ccpc_save_contacts(project_id, contacts):
    if not contacts:
        return 0
    db = sqlite3.connect(CCPC_DB)
    db.execute('DELETE FROM cceup_contacts WHERE project_id = ?', (project_id,))
    saved = 0
    for c in contacts:
        db.execute(
            'INSERT INTO cceup_contacts (project_id, company, contact_name, phone, landline, address, role, email) VALUES (?,?,?,?,?,?,?,?)',
            (project_id, c.get('company', ''), c.get('contact_name', ''),
             c.get('phone', ''), c.get('landline', '') or c.get('phone', ''),
             c.get('address', ''), c.get('role', ''), c.get('email', ''))
        )
        saved += 1
    db.commit()
    db.close()
    return saved


def ccpc_run(mode='daily'):
    """返回 0 成功 / 1 有目录列表失败（部分源缺失→rc=1 可见，不再静默）"""
    os.makedirs(CCPC_LOCAL_CACHE, exist_ok=True)
    log(f'=== CCPC v2 parse [mode={mode}] ===')

    # 1) List files from WebDAV（逐目录容错：一个目录挂了不影响其余）
    all_remote = []
    list_fail = False
    for d in CCPC_DOWNLOAD_DIRS:
        try:
            files = ccpc_webdav_list(d)
        except Exception as e:
            list_fail = True
            log(f'  !!! {d} 列表失败(已重试 3 次): {e}')
            continue
        supported = [f for f in files if f.endswith(('.docx', '.xls', '.xlsx'))]
        log(f'  {d}: {len(supported)} files')
        if supported:
            all_remote.extend((d, f) for f in supported)
    log(f'Total: {len(all_remote)} files')

    if mode == 'test':
        seen = {}
        for d, f in all_remote:
            if f not in seen:
                seen[f] = d
        to_process = [(seen[f], f) for f in list(seen.keys())[:3]]
    elif mode == 'force':
        to_process = all_remote
    else:
        to_process = all_remote

    if not to_process:
        db = sqlite3.connect(CCPC_DB)
        n = db.execute('SELECT COUNT(*) FROM cceup_projects').fetchone()[0]
        nc = db.execute('SELECT COUNT(*) FROM cceup_contacts').fetchone()[0]
        db.close()
        log(f'Nothing to process. Projects: {n}, Contacts: {nc}')
        return 1 if list_fail else 0

    ins = upd = err = 0
    for base_dir, fname in to_process:
        lp = os.path.join(CCPC_LOCAL_CACHE, fname)
        if not os.path.exists(lp) and not ccpc_webdav_download(base_dir, fname, lp):
            err += 1
            continue
        try:
            if fname.endswith(('.xls', '.xlsx')):
                rec, contacts = parse_xls(lp)
            else:
                rec, contacts = parse_docx(lp)
            if rec and rec.get('project_name'):
                # Populate contacts text field
                if contacts:
                    clines = []
                    for c in contacts:
                        parts = []
                        if c.get('company'): parts.append(c['company'])
                        if c.get('contact_name'): parts.append(f"联系人：{c['contact_name']}")
                        if c.get('phone'): parts.append(f"联系方式：{c['phone']}")
                        if c.get('landline'): parts.append(f"固定电话：{c['landline']}")
                        if c.get('email'): parts.append(f"邮箱：{c['email']}")
                        if c.get('address'): parts.append(f"地址：{c['address']}")
                        if c.get('role'): parts.append(f"角色：{c['role']}")
                        clines.append('\n'.join(parts))
                    rec['contacts'] = '\n\n'.join(clines)
                act, proj_id = ccpc_upsert_project(rec)
                c_saved = ccpc_save_contacts(proj_id, contacts)
                if act == 'INSERT':
                    ins += 1
                else:
                    upd += 1
            else:
                err += 1
        except Exception as e:
            err += 1
            log(f'  ERROR {fname[:40]}: {e}')

    db = sqlite3.connect(CCPC_DB)
    n = db.execute('SELECT COUNT(*) FROM cceup_projects').fetchone()[0]
    nc = db.execute('SELECT COUNT(*) FROM cceup_contacts').fetchone()[0]
    db.close()
    log(f'Done: +{ins} U{upd} err{err}')
    log(f'ccpc.db: {n} projects, {nc} contacts')
    return 1 if list_fail else 0


# ══════════════════════════════════════════════════════════════════════
# §4  CONTACT 段   （源：contact_lib.py）
# ══════════════════════════════════════════════════════════════════════

PHONE_NORM = re.compile(r'[\s\-—–_()（）]+')
MOBILE_RE = re.compile(r'^1[3-9]\d{9}$')


def norm_phone(v):
    if not v:
        return ''
    return PHONE_NORM.sub('', str(v)).strip()


def clean_landline(v, phone_val=''):
    """座机清洗(2026-09-10): 源数据 99% 的 'landline' 是手机号重复存 →
    手机格式(1[3-9]+9位) 或 与手机同号 的值一律剔除; 多值逐项过滤; 真座机(区号开头)保留"""
    if not v:
        return ''
    keep = []
    for part in str(v).split('/'):
        p = norm_phone(part)
        if not p:
            continue
        if MOBILE_RE.match(p):
            continue           # 手机格式 → 不是座机
        if norm_phone(phone_val) and p == norm_phone(phone_val):
            continue           # 与手机同号 → 重复
        if p not in keep:
            keep.append(p)
    return ' / '.join(keep)


def merge_val(a, b):
    """字段并集: 两侧都按 ' / ' 拆成原子值再去重合并。
    （2026-09-10 修正：原实现只把 a 拆分、b 整体比较，当 b 本身是多值时
      会产生重复项，如 '07958726678 / 07958726678 / 07957179566'）"""
    vals = [x.strip() for x in (a or '').split(' / ') if x.strip()]
    for part in (b or '').split(' / '):
        p = part.strip()
        if p and p not in vals:
            vals.append(p)
    return ' / '.join(vals)


class ContactMerger:
    def __init__(self):
        self.by_key = {}  # (company, contact_name) -> {fields..., projects:[]}

    def add(self, company, name, project, **fields):
        company = (company or '').strip()
        name = (name or '').strip()
        if not company or not name:
            return False
        key = (company, name)
        rec = self.by_key.get(key)
        if rec is None:
            rec = {'company': company, 'contact_name': name,
                   'department': '', 'position': '', 'role': '', 'phone': '',
                   'landline': '', 'email': '', 'address': '', 'remarks': '',
                   'projects': []}
            self.by_key[key] = rec
        for f in ('department', 'position', 'role', 'address', 'remarks'):
            if fields.get(f):
                rec[f] = merge_val(rec[f], fields[f])
        # 电话归一: 手机 phone 与 座机 landline 分开并集(landline 先过真座机清洗)
        if fields.get('phone'):
            rec['phone'] = merge_val(rec['phone'], norm_phone(fields['phone']))
        if fields.get('landline'):
            _ln = clean_landline(fields['landline'], fields.get('phone') or '')
            if _ln:
                rec['landline'] = merge_val(rec['landline'], _ln)
        if fields.get('email'):
            rec['email'] = merge_val(rec['email'], (fields['email'] or '').strip())
        if project and project not in rec['projects']:
            rec['projects'].append(project)
        return True


def load_ccpc(m):
    """CCPC 源: cceup_contacts + cceup_projects(id→project_id/project_name)"""
    con = sqlite3.connect(CCPC_DB)
    con.row_factory = sqlite3.Row
    rows = con.execute("""
        SELECT ct.project_id AS fk, ct.company, ct.contact_name, ct.phone,
               ct.landline, ct.email, ct.role, ct.address
        FROM cceup_contacts ct
        ORDER BY ct.project_id, ct.company, ct.contact_name, ct.phone
    """).fetchall()
    pmap = {r['id']: (r['project_id'], r['project_name']) for r in con.execute(
        "SELECT id, project_id, project_name FROM cceup_projects")}
    con.close()
    n = 0
    for r in rows:
        pid, pname = pmap.get(r['fk'], ('', ''))
        if m.add(r['company'], r['contact_name'], ('ccpc', pid, pname),
                 department='', position='', role=r['role'], phone=r['phone'],
                 landline=r['landline'], email=r['email'], address=r['address'], remarks=''):
            n += 1
    return n


def load_zc(m):
    con = sqlite3.connect(ZC_DB)
    con.row_factory = sqlite3.Row
    rows = con.execute("""
        SELECT ct.project_id AS fk, ct.company, ct.contact_name, ct.department,
               ct.position, ct.phone, ct.role, ct.address, ct.remarks
        FROM zc_contacts ct
        ORDER BY ct.project_id, ct.company, ct.contact_name, ct.phone
    """).fetchall()
    pmap = {r['project_id']: r['project_name'] for r in con.execute(
        "SELECT project_id, project_name FROM zc_projects")}
    con.close()
    n = 0
    for r in rows:
        pname = pmap.get(r['fk'], '')
        if m.add(r['company'], r['contact_name'], ('zc', r['fk'], pname),
                 department=r['department'], position=r['position'], role=r['role'],
                 phone=r['phone'], landline='', email='', address=r['address'],
                 remarks=r['remarks']):
            n += 1
    return n


def snapshot(merger=None):
    """统一快照: [[10字段..., [[source,key,name],...]], ...] 排序。merger=None 时读现库"""
    if merger is not None:
        out = []
        for key in sorted(merger.by_key):
            r = merger.by_key[key]
            out.append([r['company'], r['contact_name'], r['department'], r['position'],
                        r['role'], r['phone'], r['landline'], r['email'], r['address'],
                        r['remarks'], sorted([list(p) for p in r['projects']])])
        return out
    if not os.path.exists(CONTACT_DB):
        return None
    try:
        con = sqlite3.connect(CONTACT_DB)
        rows = con.execute("""SELECT company, contact_name, department, position, role,
            phone, landline, email, address, remarks FROM contacts
            ORDER BY company, contact_name""").fetchall()
        cps = con.execute("""SELECT contact_id, source, project_key, project_name
            FROM contact_projects ORDER BY contact_id""").fetchall()
        proj = {}
        for cid, s, pk, pn in cps:
            proj.setdefault(cid, []).append([s, pk, pn])
        ids = con.execute("SELECT id FROM contacts ORDER BY company, contact_name").fetchall()
        out = []
        for i, row in enumerate(rows):
            cid = ids[i][0]
            out.append([list(row) + [sorted(proj.get(cid, []))]])
        con.close()
        return out
    except Exception:
        return None


def content_hash(items):
    if items is None:
        return None
    return hashlib.sha256(json.dumps(items, ensure_ascii=False).encode()).hexdigest()


def write_db(merger):
    new_items = snapshot(merger)
    new_h = content_hash(new_items)
    cur_h = content_hash(snapshot(None))
    if cur_h == new_h:
        return False  # 无变化, 不落盘

    # 全量重建
    con = sqlite3.connect(CONTACT_DB)
    c = con.cursor()
    c.execute("DROP TABLE IF EXISTS contact_projects")
    c.execute("DROP TABLE IF EXISTS contacts")
    c.execute("""CREATE TABLE contacts (
        id INTEGER PRIMARY KEY AUTOINCREMENT,
        company TEXT NOT NULL,
        contact_name TEXT NOT NULL DEFAULT '',
        department TEXT DEFAULT '', position TEXT DEFAULT '', role TEXT DEFAULT '',
        phone TEXT DEFAULT '', landline TEXT DEFAULT '', email TEXT DEFAULT '',
        address TEXT DEFAULT '', remarks TEXT DEFAULT '')""")
    c.execute("CREATE UNIQUE INDEX idx_contacts_company_name ON contacts(company, contact_name)")
    c.execute("""CREATE TABLE contact_projects (
        contact_id INTEGER NOT NULL, source TEXT NOT NULL,
        project_key TEXT NOT NULL, project_name TEXT DEFAULT '')""")
    c.execute("CREATE INDEX idx_cp_contact ON contact_projects(contact_id)")
    for key in sorted(merger.by_key):
        r = merger.by_key[key]
        c.execute(
            "INSERT INTO contacts (company, contact_name, department, position, role, phone,"
            " landline, email, address, remarks) VALUES (?,?,?,?,?,?,?,?,?,?)",
            (r['company'], r['contact_name'], r['department'], r['position'], r['role'],
             r['phone'], r['landline'], r['email'], r['address'], r['remarks']))
        cid = c.lastrowid
        c.executemany(
            "INSERT INTO contact_projects (contact_id, source, project_key, project_name) VALUES (?,?,?,?)",
            [(cid, s, pk, pn) for s, pk, pn in r['projects']])
    con.commit()
    n = c.execute("SELECT COUNT(*) FROM contacts").fetchone()[0]
    nc = c.execute("SELECT COUNT(*) FROM contact_projects").fetchone()[0]
    con.close()
    return (n, nc)


def contact_run():
    m = ContactMerger()
    n_ccpc = load_ccpc(m)
    n_zc = load_zc(m)
    src_rows = n_ccpc + n_zc
    result = write_db(m)
    if result is False:
        return 0  # 静默
    n_contacts, n_proj = result
    n_company = len({k[0] for k in m.by_key})
    print(f"联系人库: 源行 {src_rows} (ccpc {n_ccpc} / zc {n_zc}) → 去重后 {n_contacts} 人 / "
          f"{n_company} 公司 / {n_proj} 项目关联")
    return 0


# ══════════════════════════════════════════════════════════════════════
# §5  MAIN — 编排
# ══════════════════════════════════════════════════════════════════════

PHASES = ('zc', 'ccpc', 'contact')


def main():
    argv = sys.argv[1:]

    # --zc-preview <pdf>：单 PDF 解析预览
    if '--zc-preview' in argv:
        i = argv.index('--zc-preview')
        if i + 1 >= len(argv):
            sys.stderr.write("用法: ccpc_zc_sync.py --zc-preview <file.pdf>\n")
            return 2
        zc_preview(argv[i + 1])
        return 0

    # --only <phase>
    only = None
    if '--only' in argv:
        i = argv.index('--only')
        only = argv[i + 1].lower() if i + 1 < len(argv) else ''
        if only not in PHASES:
            sys.stderr.write(f"--only 只接受 {'|'.join(PHASES)}，收到 {only!r}\n")
            return 2
    phases = [only] if only else list(PHASES)

    mode = 'force' if 'force' in argv else ('test' if 'test' in argv else 'daily')

    check_creds()

    print(f"[{time.strftime('%Y-%m-%d %H:%M:%S')}] ccpc+zc 入库编排开始 "
          f"(phases={'+'.join(phases)}, mode={mode})", flush=True)

    rc_all = 0
    for name in phases:
        print(f"\n===== {name.upper()} [{time.strftime('%H:%M:%S')}] =====", flush=True)
        t0 = time.time()
        try:
            if name == 'zc':
                rc = zc_run()
            elif name == 'ccpc':
                rc = ccpc_run(mode)
            else:
                rc = contact_run()
        except Exception as e:
            import traceback
            traceback.print_exc()
            print(f"!!! {name.upper()} FAILED: {e} ({time.time()-t0:.0f}s)", flush=True)
            rc_all = 1
            continue
        if rc:
            rc_all = 1
            print(f"!!! {name.upper()} 非零返回 rc={rc} ({time.time()-t0:.0f}s)", flush=True)
        else:
            print(f"--- {name.upper()} OK ({time.time()-t0:.0f}s) ---", flush=True)

    print(f"\n[{time.strftime('%Y-%m-%d %H:%M:%S')}] 编排结束 rc={rc_all}", flush=True)
    return rc_all


if __name__ == '__main__':
    sys.exit(main())
