#!/usr/bin/env python3
"""
parse_ccpc_v2.py — 统一解析 CCPC DOCX/XLS 源文件，逐字段映射 + 联系人提取
==============================================================

数据来源：坚果云 WebDAV（/dav/Downloads/ + /dav/下载1/）
缓存目录：/data/ccpc_source/
目标数据库：ccpc.db (cceup_projects + cceup_contacts)

用法：
  python3 parse_ccpc_v2.py             # 增量（按 project_name 去重，已有则 UPDATE）
  python3 parse_ccpc_v2.py force       # 全量重跑
  python3 parse_ccpc_v2.py test        # 测试 3 个文件
"""

import sys, json, os, re, urllib.parse, sqlite3
from datetime import datetime
from collections import OrderedDict

import requests
from requests.auth import HTTPBasicAuth
from xml.etree import ElementTree as ET

# ─── 第三方依赖：openpyxl（XLS）, python-docx（DOCX）────────
from openpyxl import load_workbook
from docx import Document

# ════════════════════════════════════════════════════════════
#  CONFIG
# ════════════════════════════════════════════════════════════

WEBDAV_USER = 'chniir0000@outlook.com'
WEBDAV_PASS = 'REDACTED_SEE_DOTENV'
WEBDAV_BASE = 'https://dav.jianguoyun.com'
DOWNLOAD_DIRS = ['/dav/Downloads/', '/dav/%E4%B8%8B%E8%BD%BD1/', '/dav/Downloads2/']
LOCAL_CACHE = '/data/ccpc_source'
DB_PATH     = '/root/ccpc.db'

os.makedirs(LOCAL_CACHE, exist_ok=True)
auth = HTTPBasicAuth(WEBDAV_USER, WEBDAV_PASS)


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

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 log(msg):
    print(f'[{datetime.now().strftime("%Y-%m-%d %H:%M:%S")}] {msg}')

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 工具
# ════════════════════════════════════════════════════════════

def webdav_list(remote_dir):
    url = WEBDAV_BASE + remote_dir
    r = requests.request('PROPFIND', url, auth=auth, headers={'Depth': '1'}, timeout=30)
    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 webdav_download(remote_base, filename, local_path):
    url = WEBDAV_BASE + remote_base + filename
    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
    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] 列表。
    文本格式：
      公司名(角色)
      联系人：XXX（备注）
      联系方式：1XXXXXXXXXX
      邮箱：xxx@xx.com
      地址：XXX
    返回 [{company, contact_name, phone, landline, address, role, email}]"""
    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
        # 新公司块判断：行长度>4，包含有限公司/公司/厂/院 或 (业主)/(设计)/(施工)
        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 = list(set(PHONE_RE.findall(block)))
        emails = list(set(EMAIL_RE.findall(block)))
        # 固定电话（座机）
        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 = {}  # section_name -> raw text
    all_sections_raw = {}  # ALL sections' raw text for raw_json

    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 == '项目跟踪版本':
                # 多行数据表：Row0=header, Row1+=data
                # 仅映射字段到record，不追加到detail_parts（避免重复）
                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)
                        # 不加入 detail_parts — 已有独立字段

            elif section_name == '项目基础信息':
                # 6列表：k1,v1,k2,v2,k3,v3
                # 仅映射字段到record，不追加到detail_parts（避免重复）
                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)
                                # Only keep kv_lines for keys in FIELD_MAP (skip 预测说明 etc.)
                                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
                    # 不加入 detail_parts — 已有独立字段

            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
                    # 不加【项目详情】标签，仅当与已有detail不同时才追加
                    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
                    # 已有 overview 字段，不追加到 detail_parts
                    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
                    # 已有 equipment_list 字段，不追加到 detail_parts
                    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
                    # 已有 progress 字段，不追加到 detail_parts
                    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
        detail_set = set()
        existing = record['detail']
        for part in detail_parts:
            if part not in existing:
                detail_set.add(part)
        if detail_set:
            record['detail'] = existing + '\n\n' + '\n\n'.join(detail_set)

    # 项目名称
    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 = {}
    # All project fields
    for f in ALL_PROJECT_FIELDS:
        if record.get(f):
            raw_dict[f] = record[f]
    # All section raw texts (including contact sections)
    for sn, sr in all_sections_raw.items():
        raw_dict[f'_section_{sn}'] = sr
    # Also store full document paragraphs
    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)
        # Build detail — only include 项目详情 if not already set
        if not record.get('detail'):
            raw_detail = raw.get('项目详情', '').strip()
            if raw_detail:
                record['detail'] = raw_detail
        # 项目概况/项目进展/所需材料设备/项目所在地 已有独立字段，不塞入detail
        record['raw_json'] = json.dumps(raw, ensure_ascii=False)
        # Contacts
        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)


# ════════════════════════════════════════════════════════════
#  DB 操作
# ════════════════════════════════════════════════════════════

def upsert_project(record):
    db = sqlite3.connect(DB_PATH)
    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 save_contacts(project_id, contacts):
    if not contacts:
        return 0
    db = sqlite3.connect(DB_PATH)
    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


# ════════════════════════════════════════════════════════════
#  MAIN
# ════════════════════════════════════════════════════════════

def main():
    mode = sys.argv[1] if len(sys.argv) > 1 else 'daily'
    log(f'=== CCPC v2 parse [mode={mode}] ===')

    # 1) List files from WebDAV
    all_remote = []
    for d in DOWNLOAD_DIRS:
        files = webdav_list(d)
        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(DB_PATH)
        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

    ins = upd = err = 0
    for base_dir, fname in to_process:
        lp = os.path.join(LOCAL_CACHE, fname)
        if not os.path.exists(lp) and not 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 = upsert_project(rec)
                c_saved = 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(DB_PATH)
    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')


if __name__ == '__main__':
    main()
