import json
import os
import re
import ssl
import time
import urllib.parse
import urllib.request
import uuid
import pymssql
import requests
from datetime import datetime
from qcloud_cos import CosConfig
from qcloud_cos import CosS3Client

# ---- 数据库超时配置（单位：秒；与 rpa_get.py 保持一致）----
DB_LOGIN_TIMEOUT = 10      # 建立数据库连接的超时时间
DB_QUERY_TIMEOUT = 30      # 单次查询/执行的超时时间

# ---- 异步结果回调配置（RPA 进程不读 .env，此处为部署常量）----
# ⚠️ 本文件与 rpa_get.py 各自持有一份同名实现（RPA 可能按路径加载单文件，不依赖同目录 import），
#    两份必须同步修改；DEFAULT_CALLBACK_URL / API_FLOW_OSS_PREFIX 须与同机 PHP 侧 .env 一致
DEFAULT_CALLBACK_URL = 'http://192.168.1.211:8080/delivery/rpa/callback'
CALLBACK_TIMEOUT = 10
CALLBACK_MAX_RETRIES = 3
CALLBACK_RETRY_DELAY = 5

# 与 PHP 侧 EPR_RESULT_CALLBACK_VERIFY_SSL=false 对齐：回调域名证书链含内网自签 CA
CALLBACK_SSL_CONTEXT = ssl.create_default_context()
CALLBACK_SSL_CONTEXT.check_hostname = False
CALLBACK_SSL_CONTEXT.verify_mode = ssl.CERT_NONE

# ---- API 流程 OSS 相对路径常量 ----
# API_FLOW_OSS_PREFIX 必须与同机 PHP 侧 EPR_API_OSS_PREFIX 一致，否则两份文件落在不同前缀
API_FLOW_OSS_PREFIX = 'common-test/generatefile/'
API_FLOW_OSS_MODULE_DIR = 'epr_italia'
# API 行专用桶：必须与 PHP 侧 EPR_API_COS_BUCKET 一致（source 行仍用 RPA 传入的 bucket_name）
API_FLOW_BUCKET_DEFAULT = 'usaeu-1259285998'

# ---- 回调 files[].type 枚举（UNIFIED_API_DESIGN §5.2）----
FILE_TYPE_DESC = 'EPR注册文件'
FILE_TYPE_AUTH = '授权书'

NOTIFY_COLUMNS = ['notify_status', 'notify_request', 'notify_response', 'notify_attempts', 'notify_time']

# 迁移未执行时的提示只打一次，避免每次回调刷屏
_notify_columns_probe_warned = False


def post_json(url, payload, timeout=CALLBACK_TIMEOUT):
    """POST JSON 到回调地址，2xx 视为成功（不解析响应体业务语义）。任何异常不抛出。"""
    try:
        data = json.dumps(payload, ensure_ascii=False).encode('utf-8')
        request = urllib.request.Request(
            url,
            data=data,
            headers={'Content-Type': 'application/json; charset=utf-8'},
            method='POST'
        )
        with urllib.request.urlopen(request, timeout=timeout, context=CALLBACK_SSL_CONTEXT) as response:
            body = response.read().decode('utf-8', errors='replace')
            return {
                'success': 200 <= response.getcode() < 300,
                'http_code': response.getcode(),
                'response': body[:500],
                'error': None,
            }
    except Exception as e:
        return {'success': False, 'http_code': None, 'response': None, 'error': str(e)}


def notify_api_result(task_id, callback_url, success, error='', files=None, biz_param=None,
                      persist=None, file_urls=None, max_retries=None):
    """
    发送异步结果回调（成功/失败），回调地址为空时回退 DEFAULT_CALLBACK_URL。绝不抛出异常。
    信封对齐 UNIFIED_API_DESIGN §7.2：
      成功 data={task_id, status:'success', files:[{url,name,type}]}、msg='success'
      失败 data=null、错误原因放 msg
    persist: 可选回调，每次投递尝试后调用（用于写 notify_* 列）
    max_retries: 覆盖默认重试次数（取数/保存热路径可传 0）
    file_urls: 旧参数名兼容（等价于 files）
    """
    if file_urls is not None and files is None:
        files = [{'url': url, 'name': url.rsplit('/', 1)[-1] if url else url, 'type': '授权书'}
                 for url in ([file_urls] if isinstance(file_urls, str) else file_urls) if url]

    if not callback_url:
        callback_url = DEFAULT_CALLBACK_URL
    if not callback_url:
        print('[回调] callback_url 为空且无部署常量，跳过结果通知')
        return False

    if success:
        payload = {
            'code': 200,
            'msg': 'success',
            'ProcessMode': 'async',
            'data': {'task_id': task_id, 'status': 'success', 'files': files or []},
        }
    else:
        payload = {
            'code': 500,
            'msg': error or '未知错误',
            'ProcessMode': 'async',
            'data': None,
        }

    payload['bizParam'] = biz_param if biz_param else None

    retries = CALLBACK_MAX_RETRIES if max_retries is None else max(0, int(max_retries))

    attempts = 0
    result = {'success': False, 'http_code': None, 'response': None, 'error': '未尝试'}
    while attempts <= retries:
        attempts += 1
        result = post_json(callback_url, payload)

        if persist:
            try:
                persist(result, attempts, payload)
            except Exception as e:
                print('[WARN] 回调结果落库失败（已忽略）: {}'.format(e))

        if result['success']:
            break
        if attempts <= retries and CALLBACK_RETRY_DELAY > 0:
            time.sleep(CALLBACK_RETRY_DELAY)

    print(
        "[回调] 结果通知{}: task_id={}, status={}, 尝试={}, http={}".format(
            '成功' if result['success'] else '失败',
            task_id,
            'success' if success else 'failed',
            attempts,
            result['http_code']
        )
    )
    return result['success']


def strip_cos_domain(url, bucket_name, region):
    """
    把完整 COS URL 裁成 OSS 相对路径（去掉域名与签名查询串）；非本桶 URL 原样返回。
    files[].url 契约为「OSS 相对路径，域名由 SaaS 侧拼接」。
    """
    if not url:
        return ''

    url = str(url).strip()

    if not url.startswith(('http://', 'https://')):
        return url

    prefix = 'https://{}.cos.{}.myqcloud.com/'.format(bucket_name, region)
    if not url.startswith(prefix):
        # 非本桶域名（如 file.usaeu.com 的完整 URL）原样透传，交由调用方处理
        return url

    relative = url[len(prefix):].split('?')[0]

    return urllib.parse.unquote(relative)


def resolve_api_flow_upload_bucket(api_bucket_name=None):
    """API 行的上传桶：RPA 显式传参优先，否则用部署常量（与 PHP 侧 EPR_API_COS_BUCKET 一致）。"""
    bucket = str(api_bucket_name or '').strip()

    return bucket or API_FLOW_BUCKET_DEFAULT


def sanitize_business_code(value):
    """
    清洗业务流水号作为 OSS 路径段（与 PHP 侧 sanitizeBusinessCode 同规则）。
    首尾空白按 PHP trim 的 ASCII 集合处理（不用 str.strip()，它会吃掉 NBSP/全角空格，
    导致同一流水号在 PHP 与 RPA 两侧落到不同子目录）。
    """
    cleaned = str(value or '').strip(' \t\n\r\0\x0b')
    cleaned = re.sub(r'[\x00-\x1f\x7f/\\:*?"<>|]', '_', cleaned)
    cleaned = cleaned.strip('.')

    return cleaned or 'unknown'


def build_api_cos_key(file_name, business_serial_number, fallback_identifier=''):
    """
    API 行授权书 COS 对象键：{前缀}{年}/epr_italia/{业务流水号}/{文件名}
    与 PHP 侧 FileUploadService::buildApiCosKey 拼法必须一致。
    """
    sub_dir = sanitize_business_code(business_serial_number)

    if sub_dir == 'unknown' and fallback_identifier:
        sub_dir = sanitize_business_code(fallback_identifier)

    year = datetime.now().strftime('%Y')

    return '{}/{}/{}/{}/{}'.format(
        API_FLOW_OSS_PREFIX.rstrip('/'),
        year,
        API_FLOW_OSS_MODULE_DIR,
        sub_dir,
        file_name
    )


def parse_biz_param(raw):
    """解析行内 biz_param（dict 或 JSON 字符串），失败返回 None。"""
    try:
        if isinstance(raw, dict):
            return raw
        if not raw:
            return None
        if isinstance(raw, str):
            return json.loads(raw)
        return raw
    except Exception:
        return None


def _read_row_api_attributes(cursor_biz, italia_epr_file_id, data_source, callback_url):
    """
    读取行内 data_source / callback_url / biz_param / desc_file_url。
    老库缺少这些列时回退到入参（保证与历史行为兼容）。
    """
    row_info = {}
    try:
        cursor_biz.execute(
            "SELECT data_source, callback_url, biz_param, desc_file_url FROM italia_epr_file WHERE id = %s",
            (italia_epr_file_id,)
        )
        row = cursor_biz.fetchone()
        if row:
            row_info = {
                'data_source': row[0],
                'callback_url': row[1],
                'biz_param': row[2],
                'desc_file_url': row[3],
            }
    except Exception as e:
        # 读不到行内属性时 data_source 只能按入参判定；若调用方也没透传（默认 source），
        # API 行会走 source 分支且不发回调，因此这里必须留痕便于排查
        print('[WARN] 读取行内 API 属性失败（按入参 {} 处理，若为 API 行将不发回调）: {}'.format(data_source, e))

    return {
        'data_source': str(row_info.get('data_source') or data_source or 'source').strip().lower(),
        'callback_url': (row_info.get('callback_url') or callback_url or '').strip(),
        'biz_param': row_info.get('biz_param'),
        'desc_file_url': row_info.get('desc_file_url'),
    }


def _has_notify_columns(cursor_biz):
    """探测 5 个 notify_* 列是否全部存在（迁移未执行时跳过落库，提示只打一次）。"""
    global _notify_columns_probe_warned
    try:
        placeholders = ', '.join(['%s'] * len(NOTIFY_COLUMNS))
        cursor_biz.execute(
            "SELECT COUNT(*) FROM INFORMATION_SCHEMA.COLUMNS "
            "WHERE TABLE_NAME = 'italia_epr_file' AND TABLE_SCHEMA = SCHEMA_NAME() "
            "AND COLUMN_NAME IN ({})".format(placeholders),
            tuple(NOTIFY_COLUMNS)
        )
        row = cursor_biz.fetchone()
        count = row[0] if row else 0

        if count < len(NOTIFY_COLUMNS) and not _notify_columns_probe_warned:
            _notify_columns_probe_warned = True
            print('[提示] italia_epr_file 缺少 notify_* 列，回调投递明细将不落库（请执行 api_columns.sql）')

        return count >= len(NOTIFY_COLUMNS)
    except Exception as e:
        if not _notify_columns_probe_warned:
            _notify_columns_probe_warned = True
            print('[提示] 探测 notify_* 列失败，回调投递明细将不落库: {}'.format(e))
        return False


def _persist_notify_result(db_ip, db_port, db_username, db_password, dbname, italia_epr_file_id,
                           result, attempts, payload):
    """
    独立连接写入回调投递明细（notify_* 列），失败只告警不抛出。
    """
    conn = None
    cursor = None
    try:
        conn = pymssql.connect(
            server=db_ip, port=db_port,
            user=db_username, password=db_password,
            database=dbname, charset='utf8',
            login_timeout=DB_LOGIN_TIMEOUT,
            timeout=DB_QUERY_TIMEOUT
        )
        cursor = conn.cursor()

        if not _has_notify_columns(cursor):
            return

        cursor.execute(
            """
            UPDATE italia_epr_file
            SET notify_status   = %s,
                notify_attempts = %s,
                notify_request  = %s,
                notify_response = %s,
                notify_time     = %s
            WHERE id = %s
            """,
            (
                'success' if result.get('success') else 'failed',
                attempts,
                json.dumps(payload, ensure_ascii=False),
                json.dumps({
                    'http_code': result.get('http_code'),
                    'response': result.get('response'),
                    'error': result.get('error'),
                }, ensure_ascii=False),
                datetime.now().strftime('%Y-%m-%d %H:%M:%S'),
                italia_epr_file_id,
            )
        )
        conn.commit()
    except Exception as e:
        print('[WARN] 回调结果落库失败（已忽略）: {}'.format(e))
        try:
            conn.rollback()
        except Exception:
            pass
    finally:
        try:
            if cursor:
                cursor.close()
        except Exception:
            pass
        try:
            if conn:
                conn.close()
        except Exception:
            pass


def _build_callback_files(desc_file_url, auth_file_url, bucket_name, region):
    """
    组装成功回调的 files[]：说明书（EPR注册文件）在前、授权书在后，空项跳过。
    url 为 OSS 相对路径，name 取文件名，type 为 §5.2 枚举。
    """
    files = []

    for url, file_type in ((desc_file_url, FILE_TYPE_DESC), (auth_file_url, FILE_TYPE_AUTH)):
        if not url:
            continue

        relative = strip_cos_domain(url, bucket_name, region)

        if relative.startswith(('http://', 'https://')):
            # 非本桶域名无法裁成相对路径，原样透传但留痕（通常是 bucket/region 与 PHP 侧不一致）
            print('[WARN] 文件 URL 非本桶域名，无法转为 OSS 相对路径: {}'.format(relative))

        name = os.path.basename(relative) or relative

        files.append({'url': relative, 'name': name, 'type': file_type})

    return files


def process_italia_epr_file(
    db_ip, db_port, db_username, db_password,
    database,
    dbname,
    secret_id, secret_key, region, bucket_name, cos_path,
    local_file_path,
    desc_file_oss_url,
    italia_epr_file_id,
    saas_id,
    wx_webhook_url,
    serial_number,
    company_name,
    notify_phone,
    data_source='source',
    callback_url=None,
    api_bucket_name=None
):
    """
    处理 italia_epr_file 表的授权书和说明书附件上传及数据库更新。

    source 数据源处理步骤：
      1. 将授权书上传到腾讯云COS，获得 auth_file_url
      2. 用 saas_id 查询 italia_epr_file 表中旧附件ID，若存在则从附件表删除
      3. 将授权书COS URL 插入附件表，获得 auth_file_id；
         将说明书OSS URL 插入附件表，获得 desc_file_id
      4. 用 italia_epr_file_id 将结果回写到 italia_epr_file 表，
         记录 auth_status / auth_process_time / auth_error_msg
      5. 发送企业微信通知：流水号 + 空格 + 公司名称 + @手机号

    api 数据源处理步骤（对接新系统）：
      1. 上传到 API 专用桶（API_FLOW_BUCKET_DEFAULT / api_bucket_name），键
         {API_FLOW_OSS_PREFIX}{年}/epr_italia/{业务流水号}/{文件名}
      2. 不访问附件表、不更新 EPRRegInfo，只回写 auth_status=2 / auth_file_url
      3. 发送统一异步结果回调：data.files 含两份文件（说明书 + 授权书）；失败时 data=null、原因放 msg
      4. 发送企业微信通知

    参数说明：
        db_ip                : 数据库服务器IP
        db_port              : 数据库端口
        db_username          : 数据库用户名
        db_password          : 数据库密码
        database             : 附件表所在数据库名
        dbname               : 业务表(italia_epr_file)所在数据库名
        secret_id            : 腾讯云 SecretId
        secret_key           : 腾讯云 SecretKey
        region               : 腾讯云地域，如 'ap-guangzhou'
        bucket_name          : COS 存储桶名称
        cos_path             : COS 目录前缀，如 'epr/italia/'（仅 source 数据源使用）
        local_file_path      : 授权书本地文件路径
        desc_file_oss_url    : 说明书已有 OSS URL（直接写入附件表，不上传）
        italia_epr_file_id   : italia_epr_file 表主键，用于精确 UPDATE
        saas_id              : 业务ID，用于查询同业务下的旧附件
        wx_webhook_url       : 企业微信机器人 Webhook 地址
        serial_number        : 流水号（用于通知文本）
        company_name         : 公司名称（用于通知文本）
        notify_phone         : 被@成员的手机号
        data_source          : 数据来源（'source' 默认；'api' 走 API 分支；行内可读时以行内值为准）
        callback_url         : 异步结果回调地址（api 数据源使用；行内可读时以行内值为准）
        api_bucket_name      : API 行上传桶（缺省用部署常量 API_FLOW_BUCKET_DEFAULT，须与 PHP 侧
                               EPR_API_COS_BUCKET 一致）

    返回:
        dict: 包含操作结果的字典
    """
    conn = None
    cursor = None
    conn_biz = None
    cursor_biz = None
    process_time = datetime.now().strftime('%Y-%m-%d %H:%M:%S')

    # 数据源与回调上下文（步骤 0 读行后覆盖）
    is_api_task = str(data_source or 'source').strip().lower() == 'api'
    resolved_callback_url = (callback_url or '').strip()
    resolved_biz_param = None
    desc_file_url_used = desc_file_oss_url
    auth_file_id = None
    desc_file_id = None
    upload_bucket = bucket_name  # API 行在步骤 1 改写为 API 专用桶

    def _close_all():
        for cur in [cursor, cursor_biz]:
            try:
                if cur:
                    cur.close()
            except Exception:
                pass
        for con in [conn, conn_biz]:
            try:
                if con:
                    con.close()
            except Exception:
                pass

    def _persist(result, attempts, payload):
        _persist_notify_result(
            db_ip, db_port, db_username, db_password, dbname,
            italia_epr_file_id, result, attempts, payload
        )

    def _notify_failure(error_msg):
        """API 行失败回调（source 行跳过）。"""
        if not is_api_task:
            return

        notify_api_result(
            italia_epr_file_id,
            resolved_callback_url,
            False,
            error_msg,
            biz_param=resolved_biz_param,
            persist=_persist
        )

    def _update_status_fail(error_msg):
        """失败时单独更新 italia_epr_file 状态，独立连接保证能写入（API 行同时发失败回调）。"""
        _conn = None
        _cur = None
        try:
            _conn = pymssql.connect(
                server=db_ip, port=db_port,
                user=db_username, password=db_password,
                database=dbname, charset='utf8',
                login_timeout=DB_LOGIN_TIMEOUT,
                timeout=DB_QUERY_TIMEOUT
            )
            _cur = _conn.cursor()
            _cur.execute(
                """
                UPDATE italia_epr_file
                SET auth_status       = 3,
                    auth_process_time = %s,
                    auth_error_msg    = %s
                WHERE id = %s
                """,
                (process_time, error_msg[:500], italia_epr_file_id)
            )
            _conn.commit()
        except Exception as ex:
            print('[WARN] 回写失败状态时出错: {}'.format(ex))
            try:
                if _conn:
                    _conn.rollback()
            except Exception:
                pass
        finally:
            try:
                if _cur:
                    _cur.close()
            except Exception:
                pass
            try:
                if _conn:
                    _conn.close()
            except Exception:
                pass

        _notify_failure(error_msg)

    # ──────────────────────────────────────────────
    # 步骤 0：读取行内属性，判定数据源
    # ──────────────────────────────────────────────
    try:
        conn_biz = pymssql.connect(
            server=db_ip, port=db_port,
            user=db_username, password=db_password,
            database=dbname, charset='utf8',
            login_timeout=DB_LOGIN_TIMEOUT,
            timeout=DB_QUERY_TIMEOUT
        )
        cursor_biz = conn_biz.cursor()

        attributes = _read_row_api_attributes(cursor_biz, italia_epr_file_id, data_source, callback_url)
        is_api_task = attributes['data_source'] == 'api'
        resolved_callback_url = attributes['callback_url']
        resolved_biz_param = parse_biz_param(attributes['biz_param'])

        if is_api_task:
            # API 行优先使用库内说明书链接（与 PHP 生成阶段保持一致），避免回调用到上传前的旧快照；
            # source 行保持历史行为：只用调用方传入的 desc_file_oss_url
            desc_file_url_used = attributes.get('desc_file_url') or desc_file_oss_url
    except Exception as e:
        print('[WARN] 读取行内属性失败，按入参处理: {}'.format(e))
        _close_all()
        conn_biz = None
        cursor_biz = None

    print('[数据源] {}'.format('api（对接新系统）' if is_api_task else 'source（源库轮询）'))

    # ──────────────────────────────────────────────
    # 步骤 1：上传授权书到腾讯云 COS
    # ──────────────────────────────────────────────
    try:
        if not os.path.exists(local_file_path):
            raise FileNotFoundError('本地文件不存在: {}'.format(local_file_path))

        file_size = os.path.getsize(local_file_path)
        file_name = os.path.basename(local_file_path)
        file_name_without_ext = os.path.splitext(file_name)[0]
        auth_file_name = file_name_without_ext + '.pdf'

        print('[步骤1] 授权书文件大小: {} 字节，开始上传COS...'.format(file_size))

        cos_config = CosConfig(Region=region, SecretId=secret_id, SecretKey=secret_key)
        cos_client = CosS3Client(cos_config)

        if is_api_task:
            # API 行：上传到 API 专用桶 + 统一相对路径契约，按业务流水号隔离目录
            # 兜底标识与 PHP 侧一致（BusinessId/saas_id），保证说明书与授权书落进同一目录
            upload_bucket = resolve_api_flow_upload_bucket(api_bucket_name)
            sub_dir_source = (resolved_biz_param or {}).get('BusinessSerialNumber') or str(saas_id)
            cos_key = build_api_cos_key(file_name, sub_dir_source, str(saas_id))
        else:
            upload_bucket = bucket_name
            timestamp_folder = datetime.now().strftime('%Y%m%d%H%M%S')
            cos_key = os.path.join(cos_path, timestamp_folder, file_name).replace('\\', '/')

        cos_client.upload_file(
            Bucket=upload_bucket,
            LocalFilePath=local_file_path,
            Key=cos_key
        )

        if is_api_task:
            # API 行入库与回调统一用 OSS 相对路径（不带域名），与 PHP 侧 uploadApiFile 口径一致：
            # 域名由消费方（SaaS 侧）拼接，换桶/换域名不会脏数据
            auth_file_url = cos_key
        else:
            # source 行保持完整 URL（附件表 F_FilePath/FileUrl 用）
            auth_file_url = 'https://{}.cos.{}.myqcloud.com/{}'.format(upload_bucket, region, cos_key)

        print('[步骤1] 授权书上传成功，COS URL: {}'.format(auth_file_url))

    except Exception as e:
        _update_status_fail('步骤1 上传COS失败: {}'.format(e))
        return {'success': False, 'message': '步骤1 上传COS失败: {}'.format(e), 'error': str(e)}

    # ──────────────────────────────────────────────
    # 步骤 2 & 3 & 4：数据库操作（事务）
    # ──────────────────────────────────────────────
    try:
        if is_api_task:
            # ── API 行：不访问附件表、不更新 EPRRegInfo，只回写业务表 ──
            if conn_biz is None:
                conn_biz = pymssql.connect(
                    server=db_ip, port=db_port,
                    user=db_username, password=db_password,
                    database=dbname, charset='utf8',
                    login_timeout=DB_LOGIN_TIMEOUT,
                    timeout=DB_QUERY_TIMEOUT
                )
            if cursor_biz is None:
                cursor_biz = conn_biz.cursor()

            print('[步骤4] 更新 italia_epr_file（API 行），id={}...'.format(italia_epr_file_id))
            current_time = datetime.now().strftime('%Y-%m-%d %H:%M:%S')
            cursor_biz.execute(
                """
                UPDATE italia_epr_file
                SET auth_file_url     = %s,
                    auth_status       = 2,
                    auth_process_time = %s,
                    auth_error_msg    = NULL
                WHERE id = %s
                """,
                (auth_file_url, current_time, italia_epr_file_id)
            )
            conn_biz.commit()
            print('[步骤4] italia_epr_file 更新成功（API 行跳过附件表与 EPRRegInfo）')
        else:
            # ── source 行：历史行为完全一致 ──
            # 建立两个库的连接 source db = saas db = vat_db
            conn = pymssql.connect(
                server=db_ip, port=db_port,
                user=db_username, password=db_password,
                database=database, charset='utf8',
                login_timeout=DB_LOGIN_TIMEOUT,
                timeout=DB_QUERY_TIMEOUT
            )
            cursor = conn.cursor()

            # target db = rpa db = rpa
            if conn_biz is None:
                conn_biz = pymssql.connect(
                    server=db_ip, port=db_port,
                    user=db_username, password=db_password,
                    database=dbname, charset='utf8',
                    login_timeout=DB_LOGIN_TIMEOUT,
                    timeout=DB_QUERY_TIMEOUT
                )
            cursor_biz = conn_biz.cursor()

            # ── 步骤 2：用 saas_id 查旧附件 ID，存在则删除 ──
            print('[步骤2] 查询 saas_id={} 的旧附件...'.format(saas_id))
            cursor_biz.execute(
                """
                SELECT auth_file_id, desc_file_id
                FROM italia_epr_file
                WHERE saas_id = %s
                """,
                (saas_id,)
            )
            row = cursor_biz.fetchone()

            old_auth_file_id = row[0] if row and row[0] is not None else None
            old_desc_file_id = row[1] if row and row[1] is not None else None

            for old_fid, label in [(old_auth_file_id, '授权书'), (old_desc_file_id, '说明书')]:
                if old_fid:
                    cursor.execute(
                        "DELETE FROM Base_AnnexesFile WHERE F_Id = %s",
                        (old_fid,)
                    )
                    print('[步骤2] 已删除旧{}附件: {}'.format(label, old_fid))
                else:
                    print('[步骤2] 旧{}附件为空，跳过删除'.format(label))

            # ── 步骤 3：插入新附件记录 ────────────────────
            current_time = datetime.now().strftime('%Y-%m-%d %H:%M:%S')

            def insert_attachment(file_url, display_file_name, f_size, file_ext='pdf'):
                """向 Base_AnnexesFile 插入一条附件记录，返回生成的 F_Id。"""
                new_fid = str(uuid.uuid4()).lower()
                cursor.execute(
                    """
                    INSERT INTO Base_AnnexesFile (
                        F_Id, F_FolderId, F_FileName, F_FilePath, CreationName,
                        F_FileSize, F_FileExtensions, F_FileType, F_DownloadCount,
                        CreationDate, Creation_Id, InfoId, FileCover, FileUrl,
                        FileFrom, FileGroup, FileCategoryId
                    ) VALUES (
                        %s, %s, %s, %s, %s,
                        %s, %s, %s, %s,
                        %s, %s, %s, %s, %s,
                        %s, %s, %s
                    )
                    """,
                    (
                        new_fid,
                        '',
                        display_file_name,
                        file_url,
                        'System',
                        f_size,
                        '.{}'.format(file_ext),   # F_FileExtensions ← 动态
                        file_ext,                 # F_FileType       ← 动态
                        0,
                        current_time,
                        'System',
                        saas_id,
                        None,
                        file_url,
                        3,
                        2,
                        '1c199ef6-71ac-4f21-806a-f76407522e5c'
                    )
                )
                print('[步骤3] 附件插入成功，F_Id: {}，文件名: {}，类型: {}'.format(new_fid, display_file_name, file_ext))
                return new_fid

            # 插入授权书附件（pdf）
            auth_file_id = insert_attachment(auth_file_url, auth_file_name, file_size, file_ext='pdf')

            # 插入说明书附件（docx）
            desc_display_name = os.path.basename(desc_file_url_used.split('?')[0]) or 'desc.docx'
            if not desc_display_name.lower().endswith('.docx'):
                desc_display_name = os.path.splitext(desc_display_name)[0] + '.docx'
            desc_file_id = insert_attachment(desc_file_url_used, desc_display_name, 0, file_ext='docx')

            # ── 步骤 4：用主键 id 回写 italia_epr_file ───
            print('[步骤4] 更新 italia_epr_file，id={}...'.format(italia_epr_file_id))
            cursor_biz.execute(
                """
                UPDATE italia_epr_file
                SET auth_file_url     = %s,
                    auth_file_id      = %s,
                    desc_file_id      = %s,
                    auth_status       = 2,
                    auth_process_time = %s,
                    auth_error_msg    = NULL
                WHERE id = %s
                """,
                (auth_file_url, auth_file_id, desc_file_id, current_time, italia_epr_file_id)
            )
            print('[步骤4] italia_epr_file 更新成功')

            # ── 步骤 4.5：更新 EPRRegInfo 表 ──────
            # 三个关键字段都有值则为6，任意为空则为7
            push_status = 6 if (auth_file_id and desc_file_id and auth_file_url) else 7
            print('[步骤4.5] 更新 EPRRegInfo，PushTaxBureauStatus={}，saas_id={}...'.format(push_status, saas_id))
            cursor.execute(
                """
                UPDATE EPRRegInfo
                SET PushTaxBureauStatus = %s
                WHERE Id = %s
                """,
                (push_status, saas_id)
            )
            print('[步骤4.5] EPRRegInfo 更新成功，PushTaxBureauStatus={}'.format(push_status))

            # 提交两个事务
            conn.commit()
            conn_biz.commit()

    except Exception as e:
        if conn:
            try:
                conn.rollback()
            except Exception:
                pass
        if conn_biz:
            try:
                conn_biz.rollback()
            except Exception:
                pass
        _close_all()
        _update_status_fail('步骤2-4 数据库操作失败: {}'.format(e))
        return {'success': False, 'message': '步骤2-4 数据库操作失败: {}'.format(e), 'error': str(e)}

    finally:
        _close_all()

    # ──────────────────────────────────────────────
    # 步骤 4.9（仅 API 行）：统一异步结果回调，下发两份文件
    # ──────────────────────────────────────────────
    if is_api_task:
        # 重新读一次说明书链接，避免使用上传前的旧快照
        latest_desc_url = desc_file_url_used
        _conn = None
        _cur = None
        try:
            _conn = pymssql.connect(
                server=db_ip, port=db_port,
                user=db_username, password=db_password,
                database=dbname, charset='utf8',
                login_timeout=DB_LOGIN_TIMEOUT,
                timeout=DB_QUERY_TIMEOUT
            )
            _cur = _conn.cursor()
            _cur.execute("SELECT desc_file_url FROM italia_epr_file WHERE id = %s", (italia_epr_file_id,))
            _row = _cur.fetchone()
            if _row and _row[0]:
                latest_desc_url = _row[0]
        except Exception as e:
            print('[WARN] 重读说明书链接失败，使用已知值: {}'.format(e))
        finally:
            try:
                if _cur:
                    _cur.close()
            except Exception:
                pass
            try:
                if _conn:
                    _conn.close()
            except Exception:
                pass

        callback_files = _build_callback_files(latest_desc_url, auth_file_url, upload_bucket, region)

        if not callback_files:
            print('[WARN] 成功回调 files 为空（说明书与授权书链接均缺失），请检查 OSS 配置与库内链接')

        print('[步骤4.9] 发送成功回调，files={}'.format(json.dumps(callback_files, ensure_ascii=False)))

        notify_api_result(
            italia_epr_file_id,
            resolved_callback_url,
            True,
            files=callback_files,
            biz_param=resolved_biz_param,
            persist=_persist
        )

    # ──────────────────────────────────────────────
    # 步骤 5：发送企业微信通知
    # ──────────────────────────────────────────────
    try:
        notify_text = '{} {} 已成功生成文件'.format(serial_number, company_name)
        print('[步骤5] 发送企业微信通知: {}'.format(notify_text))

        wx_payload = {
            "msgtype": "text",
            "text": {
                "content": notify_text,
                "mentioned_mobile_list": [notify_phone]
            }
        }
        resp = requests.post(wx_webhook_url, json=wx_payload, timeout=10)
        resp_json = resp.json()

        if resp_json.get('errcode') == 0:
            print('[步骤5] 企业微信通知发送成功')
        else:
            print('[步骤5] 企业微信通知发送失败: {}'.format(resp_json))

    except Exception as e:
        print('[WARN] 企业微信通知异常: {}'.format(e))

    return {
        'success': True,
        'message': '授权书上传、附件入库、状态更新均成功',
        'auth_file_url': auth_file_url,
        'auth_file_id': auth_file_id,
        'desc_file_id': desc_file_id,
        'italia_epr_file_id': italia_epr_file_id,
    }
