import json
import pymssql
import os
import re
import time
import uuid
import urllib.parse
import urllib.request
import urllib.error
from datetime import datetime
from qcloud_cos import CosConfig
from qcloud_cos import CosS3Client

# API流程结果通知接口（通知 URL 由 RPA 程序固定配置，可通过参数覆盖，与 es_haiya 一致；
# 默认值与 PHP .env ES_HAGUE_RESULT_CALLBACK_URL / es_haiya 同为 test-cloud 回调地址）
DEFAULT_CALLBACK_URL = 'https://test-cloud.usaeu.com/prod-api/delivery/rpa/callback'
# HTTP 超时秒数，与 PHP AsyncResultNotifier 默认一致
CALLBACK_TIMEOUT = 10

# 发送失败重试次数（总尝试 = 1+重试次数，0 只发一次），与 PHP AsyncResultNotifier 默认值一致
CALLBACK_MAX_RETRIES = 3
# 每次重试间隔秒数
CALLBACK_RETRY_DELAY = 5

# API流程OSS相对路径前缀：生成文件 COS 键/入库/结果通知统一相对路径（年份文件夹 generatefile/2026/ 之前的部分），
# 按部署环境配置——测试环境 common-test/generatefile/、正式环境 common-prod/generatefile/；
# 与 PHP .env.example API_FLOW_OSS_PREFIX 同值；RPA 进程环境与项目 .env 隔离、不读环境变量，
# 部署到对应服务器时保持为对应环境的常量
API_FLOW_OSS_PREFIX = 'common-test/generatefile/'

# API流程生成文件子目录（年份文件夹之后；海牙整本/APODERAMIENTO/030 文件同目录存放，文件命名不变；
# 该子目录下按业务流水号再分一层唯一性子目录，见 upload_to_cos_and_update_db）
API_FLOW_OSS_MODULE_DIR = 'es_epr_haiya'

# API流程上传桶（与 source 流程桶 vat-1259285998 不同）：API 流程结果归 SaaS 共享桶，
# 默认 usaeu-1259285998（与 PHP .env.example API_FLOW_COS_BUCKET 同值）；RPA 进程环境与项目 .env 隔离，
# 不支持环境变量覆盖，覆盖通道为 api_bucket_name 参数（对齐 es_haiya）
API_FLOW_BUCKET_DEFAULT = 'usaeu-1259285998'


def resolve_api_flow_upload_bucket(api_bucket_name=None):
    """
    解析 API 流程实际上传桶（API 与 source 流程桶不同：API 结果归 SaaS 共享桶 usaeu-1259285998）
    优先级：api_bucket_name 参数 > API_FLOW_BUCKET_DEFAULT

    参数:
        api_bucket_name: API 流程专用桶（可选）

    返回:
        str: API 流程实际上传桶
    """
    return api_bucket_name or API_FLOW_BUCKET_DEFAULT


def strip_cos_domain(url, bucket='', region=''):
    """
    将 COS 完整下载URL裁剪为对象相对路径（2026-08-31 契约：入库/结果通知 files.url 为 OSS 相对路径、不含域名）
    - 相对路径（无 http(s):// 前缀）：原样返回
    - 完整 URL 命中 COS 桶域名（bucket 非空时校验，如 usaeu-1259285998.cos.ap-guangzhou.myqcloud.com）：
      取路径部分并剥离 ?sign= 等查询串（签名 URL / 无签名直连 URL 同规则）
    - 其他域名（如 vat 桶静态盖章要求文件、文件服务器）或无法识别：原样返回（透传兜底）

    参数:
        url: 完整URL或COS相对路径
        bucket: COS 桶名（为空时不做域名校验，直接裁剪）
        region: COS 地域,如 'ap-guangzhou'

    返回:
        str: OSS 相对路径或原URL
    """
    url = (url or '').strip()
    if not url:
        return ''
    if not re.match(r'^https?://', url, re.IGNORECASE):
        return url
    parsed = urllib.parse.urlparse(url)
    if bucket:
        domain = f'{bucket}.cos.{region}.myqcloud.com'
        if parsed.netloc.lower() != domain.lower():
            return url
    return urllib.parse.unquote(parsed.path.lstrip('/'))

# DB 写入死锁限时重试（vat_db 生产库堆表 40001 死锁高频，只能应用层重试缓解；
# 调用方按"一步一提交"执行，死锁只回滚本条语句，逐条重试保证不半截落库）
DB_WRITE_MAX_ATTEMPTS = 3              # 单条写入最大尝试次数（含首次）
DB_WRITE_BACKOFF_SECONDS = (2, 4)      # 重试退避间隔
DB_WRITE_HARD_DEADLINE_SECONDS = 15    # 单条写入重试硬超时：到点即放弃抛出，绝不卡死单


def _is_deadlock_error(err):
    """判断是否为可重试的死锁类错误（SQL Server 错误码 1205 / SQLSTATE 40001）"""
    text = str(err).lower()
    return '1205' in text or '40001' in text or 'deadlock' in text


def _execute_write_with_retry(cursor, sql, params, desc):
    """
    单条 DB 写入限时重试执行

    死锁类错误：按退避间隔重试本语句（调用方按"一步一提交"执行，死锁只回滚本条语句）；
    非死锁错误：立即抛出，不做无意义重试；
    次数上限 + 硬超时双重约束，绝不无限重试卡死机器人单线程。
    """
    deadline = time.time() + DB_WRITE_HARD_DEADLINE_SECONDS
    attempt = 0
    while True:
        attempt += 1
        try:
            cursor.execute(sql, params)
            return
        except Exception as err:
            if not _is_deadlock_error(err) or attempt >= DB_WRITE_MAX_ATTEMPTS or time.time() > deadline:
                raise
            backoff = DB_WRITE_BACKOFF_SECONDS[min(attempt - 1, len(DB_WRITE_BACKOFF_SECONDS) - 1)]
            print(f'{desc}死锁，第 {attempt} 次重试（退避 {backoff}s）: {err}')
            time.sleep(backoff)


def post_json(url, payload, timeout=CALLBACK_TIMEOUT):
    """
    POST JSON 请求

    参数:
        url: 请求地址
        payload: 请求体（dict，将被 JSON 编码）
        timeout: 超时秒数

    返回:
        dict: {success, http_code, response, error}
    """
    data = json.dumps(payload, ensure_ascii=False).encode('utf-8')
    req = urllib.request.Request(
        url,
        data=data,
        headers={'Content-Type': 'application/json; charset=utf-8'},
        method='POST'
    )
    try:
        with urllib.request.urlopen(req, timeout=timeout) as resp:
            body = resp.read().decode('utf-8', errors='replace')
            success = 200 <= resp.status < 300
            return {
                'success': success,
                'http_code': resp.status,
                'response': body,
                'error': None if success else f'HTTP {resp.status}',
            }
    except urllib.error.HTTPError as e:
        body = e.read().decode('utf-8', errors='replace')
        return {
            'success': False,
            'http_code': e.code,
            'response': body,
            'error': f'HTTP {e.code}',
        }
    except Exception as e:
        return {
            'success': False,
            'http_code': None,
            'response': None,
            'error': str(e),
        }


def query_task_info(cursor_old, task_id):
    """
    查询任务记录信息，用于通知接口取 bizParam 与 ResultData

    参数:
        cursor_old: EPRHaiyaProcessingTasks 所在库的游标
        task_id: 任务ID

    返回:
        dict: {task_data, result_data}；查询失败或任务不存在返回 None
    """
    try:
        cursor_old.execute(
            "SELECT TaskData, ResultData FROM EPRHaiyaProcessingTasks WHERE Id = %s",
            (task_id,)
        )
        row = cursor_old.fetchone()
    except Exception:
        return None
    if not row:
        return None
    return {
        'task_data': row[0],
        'result_data': row[1],
    }


def parse_biz_param(task_info):
    """
    从 TaskData JSON 中取 bizParam（受理时原样回传的业务参数，通知内容需含业务标识）

    参数:
        task_info: query_task_info 的返回

    返回:
        dict 或 None
    """
    try:
        task_data = json.loads(task_info['task_data']) if task_info and task_info.get('task_data') else {}
        if isinstance(task_data, dict):
            return task_data.get('bizParam')
    except (ValueError, TypeError):
        pass
    return None


def parse_result_data(task_info):
    """
    从 ResultData JSON 中取文件结果（file1~file4/cos_key/cos_url）

    参数:
        task_info: query_task_info 的返回

    返回:
        dict
    """
    try:
        result_data = json.loads(task_info['result_data']) if task_info and task_info.get('result_data') else {}
        return result_data if isinstance(result_data, dict) else {}
    except (ValueError, TypeError):
        return {}


def notify_api_result(task_id, callback_url, success, error='', files=None, biz_param=None, persist=None):
    """
    调用异步结果通知接口（API流程专用），报告 RPA 030 生成成功/失败（发送失败自动重试）

    请求体信封与 PHP AsyncResultNotifier 一致：{code, msg, ProcessMode, data, bizParam}
    - 成功: code=200, msg='success', data={task_id, status:'success', files:[{url,name,type},...]}
      files：生成的文件数组（UNIFIED_API_DESIGN §5.1/§13.2 文件项规范），url=OSS 相对路径（不带域名、
      不生成签名，域名由接收方拼接；非 usaeu 桶对象如静态盖章要求文件保留完整 URL 透传）；文件槽顺序
      file1 盖章要求 → file2 海牙 → file3 APODERAMIENTO → file4 030（非空项）
    - 失败: code=500, msg=错误原因, data=null（非 200 时 data 统一为 null，错误原因放 msg）

    重试: CALLBACK_MAX_RETRIES（总尝试=1+重试次数）+ CALLBACK_RETRY_DELAY 间隔，与 PHP 侧一致；
    2xx 视为成功，不捕获对方响应体中的业务失败。

    即时落库（对齐 es_haiya 2026-08-29 线上问题修复）：传入 persist 回调时，每次尝试结束后（含失败尝试）
    立即调用 persist(result) 落库 EPRHaiyaProcessingTasks.Notify*。此前全部重试结束才落库一次，重试链阻塞
    可达 55 秒，期间流程被中断/超时杀进程时 Notify* 字段全空，失败信息无任何痕迹。NotifyStatus 受表 CHECK
    约束仅允许 skipped/success/failed，故无"发送中"中间态：首次失败尝试即写入 failed + 当次错误。

    参数:
        task_id: 任务ID
        callback_url: 结果通知接口地址（为空时用 DEFAULT_CALLBACK_URL）
        success: 是否成功
        error: 失败原因（success=False 时必填，写入 msg）
        files: 成功时的文件项数组 [{url,name,type}]（写入 data.files）
        biz_param: 受理时原样回传的业务参数
        persist: 可选落库回调 persist(result_dict)，每次尝试后调用一次；
                 回调自身异常仅打印，不影响通知重试主流程

    返回:
        dict: {success, status, attempts, request, http_code, response, error}
    """
    url = callback_url or DEFAULT_CALLBACK_URL
    if success:
        data = {
            'task_id': task_id,
            'status': 'success',
            'files': files or [],
        }
        msg = 'success'
    else:
        # 非 200：data 统一为 null，错误原因放 msg
        data = None
        msg = error or '未知错误'

    payload = {
        'code': 200 if success else 500,
        'msg': msg,
        'ProcessMode': 'async',
        'data': data,
        'bizParam': biz_param if biz_param else None,
    }
    result = {
        'success': False,
        'status': 'failed',
        'attempts': 0,
        'request': payload,
        'http_code': None,
        'response': None,
        'error': None,
    }
    total_attempts = CALLBACK_MAX_RETRIES + 1
    for attempt in range(1, total_attempts + 1):
        send_result = post_json(url, payload)
        result['attempts'] = attempt
        result['http_code'] = send_result['http_code']
        result['response'] = send_result['response']
        result['error'] = send_result['error']

        send_ok = send_result['success']
        if send_ok:
            result['success'] = True
            result['status'] = 'success'

        # 每次尝试后立即落库（含失败尝试）：重试链阻塞可达 55 秒，期间流程被中断/超时杀进程时，
        # 已发生的尝试信息仍保留在 Notify* 字段（此前全部重试结束才落一次，中断即全空）
        if persist:
            try:
                persist(dict(result))
            except Exception as persist_error:
                print(f"通知结果即时落库出错（不影响通知重试）: {persist_error}")

        if send_ok:
            print(f"结果通知接口调用成功: task_id={task_id}, status={'success' if success else 'failed'}, "
                  f"attempt={attempt}, http_code={send_result['http_code']}")
            return result

        print(f"结果通知接口调用失败: task_id={task_id}, status={'success' if success else 'failed'}, "
              f"attempt={attempt}/{total_attempts}, http_code={send_result['http_code']}, error={send_result['error']}")

        if attempt < total_attempts and CALLBACK_RETRY_DELAY > 0:
            time.sleep(CALLBACK_RETRY_DELAY)

    return result


def has_notify_columns(cursor):
    """
    EPRHaiyaProcessingTasks 的 Notify* 列是否已执行迁移（2026_09_09_add_notify_fields.sql）
    旧环境返回 False：跳过落库，与 PHP AsyncTaskManager::columnExists 逻辑一致
    """
    try:
        cursor.execute(
            "SELECT COUNT(*) FROM INFORMATION_SCHEMA.COLUMNS "
            "WHERE TABLE_NAME = 'EPRHaiyaProcessingTasks' AND COLUMN_NAME IN "
            "('NotifyRequest','NotifyResponse','NotifyStatus','NotifyAttempts','NotifyTime')"
        )
        row = cursor.fetchone()
        return row is not None and int(row[0]) == 5
    except Exception:
        return False


def save_notify_result(cursor, conn, task_id, notify_result):
    """
    保存结果通知的请求参数与结果到 EPRHaiyaProcessingTasks.Notify* 字段

    参数:
        cursor/conn: EPRHaiyaProcessingTasks 所在库的游标和连接
        task_id: 任务ID
        notify_result: notify_api_result 的返回

    返回:
        bool: 是否保存成功（列未迁移返回 False；失败不抛异常，不影响通知主流程）
    """
    if not notify_result or not notify_result.get('request'):
        return False
    try:
        if not has_notify_columns(cursor):
            print(f'EPRHaiyaProcessingTasks Notify* 列未迁移，跳过结果通知落库 (Id={task_id})')
            return False
        sql = """
            UPDATE EPRHaiyaProcessingTasks
            SET NotifyRequest = %s,
                NotifyResponse = %s,
                NotifyStatus = %s,
                NotifyAttempts = %s,
                NotifyTime = %s
            WHERE Id = %s
        """
        cursor.execute(sql, (
            json.dumps(notify_result['request'], ensure_ascii=False),
            json.dumps({
                'http_code': notify_result.get('http_code'),
                'response': notify_result.get('response'),
                'error': notify_result.get('error'),
            }, ensure_ascii=False),
            notify_result.get('status', 'failed'),
            notify_result.get('attempts', 0),
            datetime.now().strftime('%Y-%m-%d %H:%M:%S'),
            task_id,
        ))
        conn.commit()
        print(f'EPRHaiyaProcessingTasks 表中 Id={task_id} 的结果通知记录已保存（Notify* 字段）')
        return True
    except Exception as e:
        print(f"保存结果通知记录出错 (Id={task_id}): {e}")
        return False


def fail_api_task(cursor_old, conn_old, task_id, error, callback_url, task_info):
    """
    API任务失败处理：pdf_result=3 + 调用结果通知接口告知失败（对齐 es_haiya fail_api_task）
    置 3 终结任务——rpa_030_get 捞取条件为 pdf_result IN (0,1)，置 3 后不再被反复重捞；
    通知失败不影响任务状态更新

    参数:
        cursor_old/conn_old: EPRHaiyaProcessingTasks 所在库的游标和连接
        task_id: 任务ID
        error: 失败原因
        callback_url: 结果通知接口地址
        task_info: query_task_info 的返回
    """
    try:
        sql = """
            UPDATE EPRHaiyaProcessingTasks
            SET pdf_result = 3
            WHERE Id = %s
        """
        cursor_old.execute(sql, (task_id,))
        conn_old.commit()
        print(f'EPRHaiyaProcessingTasks 表中 Id={task_id} 的记录已标记 030 失败(pdf_result=3)')
    except Exception as e:
        print(f"更新任务失败状态出错: {e}")

    try:
        # 通知即时落库：每次尝试后（含失败）写 Notify* 字段，重试中途流程被中断也不丢已发生的尝试信息
        notify_api_result(
            task_id, callback_url, False, error, biz_param=parse_biz_param(task_info),
            persist=lambda notify_state: save_notify_result(cursor_old, conn_old, task_id, notify_state)
        )
    except Exception as e:
        print(f"调用结果通知接口出错: {e}")


def upload_to_cos_and_update_db(db_ip, db_port, db_username, db_password, database, dbname,
                                  secret_id, secret_key, region, bucket_name, cos_path,
                                  local_file_path, EPRRegInfoId,
                                  async_vat_tasks_id, callback_url=None, api_bucket_name=None):
    """
    上传本地文件到腾讯云COS,并更新数据库中的相关记录

    API任务（任务 DataSource='api'，公共接口受理的外部项目数据）：
        - 上传桶为 SaaS 共享桶（api_bucket_name 参数 > 默认 usaeu-1259285998，与 source 桶不同——COS 桶分流）
        - COS 对象键 = {API_FLOW_OSS_PREFIX}{当前年}/es_epr_haiya/{业务流水号}/{文件名}（相对路径，
          2026-08-31 契约 + 2026-09-02 唯一性子目录：不含时间戳子文件夹，业务流水号取
          bizParam.BusinessSerialNumber、兜底任务ID；与 PHP 生成的海牙/APODERAMIENTO 同子目录，文件命名不变）
        - 只回填 EPRHaiyaProcessingTasks 任务表（pdf030_fid / pdf_result / pdf_result_url、
          ResultData.file4=030 的相对路径与跨任务重定向）
        - 不写 vat_db 附件表 Base_AnnexesFile、不更新 EPRRegInfo（外部项目无对应业务记录）
        - 030 处理完成/失败时调用异步结果通知接口（与 es_haiya 一致，data.files 为文件数组、
          url=OSS 相对路径不带域名；URL 默认 DEFAULT_CALLBACK_URL，可传 callback_url 覆盖；
          发送失败自动重试 + 每次尝试即时落库 Notify* 字段，通知失败不影响主流程）
        - 失败置死（对齐 es_haiya fail_api_task）：本地文件不存在或上传/回填失败时置 pdf_result=3
          终结任务（rpa_030_get 捞取条件 pdf_result IN (0,1)，置 3 后不再被反复重捞）+ 通知调用方；
          例外：任务回填已提交（pdf_result=2）后的辅助步骤失败不回改 3，只通知异常
    source任务（原逻辑不变）：
    source任务（原逻辑不变）：
        - 写 vat_db 附件表 + EPRRegInfo(PushTaxBureauStatus=6) + 回填任务表（不改 ResultData）

    参数:
        db_ip: 数据库服务器IP
        db_port: 数据库端口
        db_username: 数据库用户名
        db_password: 数据库密码
        database: 数据库名(用于附件表Base_AnnexesFile)
        dbname: 数据库名(用于EPRHaiyaProcessingTasks表)
        secret_id: 腾讯云SecretId
        secret_key: 腾讯云SecretKey
        region: 腾讯云地域,如 'ap-guangzhou'
        bucket_name: COS存储桶名称
        cos_path: COS目录路径,如 'vat/spain/'
        local_file_path: 本地文件路径
        EPRRegInfoId: EPRRegInfo表的ID
        async_vat_tasks_id: 当前EPRHaiyaProcessingTasks任务ID
        callback_url: 结果通知接口地址（API任务成功后/失败时调用；为空时用 DEFAULT_CALLBACK_URL）
        api_bucket_name: API 流程上传桶（可选；为空时用默认值 usaeu-1259285998）

    返回:
        dict: 包含操作结果的字典
    """
    conn = None
    cursor = None
    conn_old = None
    cursor_old = None
    is_api_task = False  # 预初始化，异常通知分支可安全引用
    task_info = None  # 预初始化：API 任务 key 唯一性子目录与通知均需任务数据
    current_task_committed = False  # 当前任务回填(pdf_result=2)是否已提交；异常分支据此避免把成功态覆盖成 3

    try:
        # 先连接任务表数据库(EPRHaiyaProcessingTasks)并判定任务来源（上传路径/桶按来源分流，
        # 判定前置到上传/文件检查之前——API 任务失败需置死+通知调用方，必须先知道任务来源；对齐 es_haiya）
        conn_old = pymssql.connect(
            server=db_ip,
            port=db_port,
            user=db_username,
            password=db_password,
            database=dbname,
            charset='utf8'
        )
        cursor_old = conn_old.cursor()

        # 查询任务 EPRBusinessRecordId
        cursor_old.execute(
            "SELECT EPRBusinessRecordId FROM EPRHaiyaProcessingTasks WHERE Id = %s",
            (async_vat_tasks_id,)
        )
        bid_row = cursor_old.fetchone()
        if not bid_row or not bid_row[0]:
            return {
                'success': False,
                'message': f'未找到任务(Id={async_vat_tasks_id})的 EPRBusinessRecordId'
            }
        epr_business_record_id = bid_row[0]

        # 判定任务来源：API 任务（DataSource='api'，公共接口受理，外部项目数据）
        # 只回填任务表字段，不写 vat_db 附件表（Base_AnnexesFile）与 EPRRegInfo；
        # source 任务保持原逻辑不变。DataSource 列缺失（旧环境）时按 source 处理。
        is_api_task = False
        try:
            cursor_old.execute(
                "SELECT DataSource FROM EPRHaiyaProcessingTasks WHERE Id = %s",
                (async_vat_tasks_id,)
            )
            ds_row = cursor_old.fetchone()
            data_source = (ds_row[0] or '').strip().lower() if ds_row and ds_row[0] else ''
            is_api_task = (data_source == 'api')
            print(f'任务来源: {"API(仅回填任务表)" if is_api_task else "SOURCE(原逻辑)"}')
        except Exception as ds_err:
            # DataSource 列不存在等异常时保守按 source 流程处理
            print(f'读取DataSource失败({ds_err}),按source流程处理')

        # 检查本地文件是否存在（判定来源之后，对齐 es_haiya：API 任务置死 pdf_result=3 + 通知调用方，
        # 任务终结不再被 rpa_030_get 以 pdf_result IN (0,1) 反复重捞；source 任务维持原返回）
        if not os.path.exists(local_file_path):
            message = f'本地文件不存在: {local_file_path}'
            if is_api_task:
                task_info = task_info or query_task_info(cursor_old, async_vat_tasks_id)
                fail_api_task(cursor_old, conn_old, async_vat_tasks_id, message, callback_url, task_info)
            return {
                'success': False,
                'message': message
            }

        # 计算文件大小(字节)
        file_size = os.path.getsize(local_file_path)
        print(f'文件大小: {file_size} 字节')

        # 配置腾讯云COS
        config = CosConfig(
            Region=region,
            SecretId=secret_id,
            SecretKey=secret_key
        )
        client = CosS3Client(config)

        # 获取文件名
        file_name = os.path.basename(local_file_path)

        # 构建COS对象键与上传桶
        # API 任务：统一相对路径 {API_FLOW_OSS_PREFIX}{当前年}/es_epr_haiya/{业务流水号}/{文件名}
        # （2026-08-31 契约，2026-09-02 追加唯一性子目录：文件名不含唯一因子（海牙为"海牙待认证_{公司名}.pdf"），
        # 子目录 = bizParam.BusinessSerialNumber（受理校验必填），兜底 EPRRegInfoId（API 任务正常为 null 不触发），
        # 与 PHP uploadToTencentCOS 一致；海牙/APODERAMIENTO/030 同子目录存放，文件命名不变），上传 SaaS 共享桶；
        # source 任务：cos_path + 时间戳文件夹 + 文件名、source 桶（原逻辑不变）
        if is_api_task:
            upload_bucket = resolve_api_flow_upload_bucket(api_bucket_name)
            if upload_bucket != bucket_name:
                print(f'API流程上传桶分流: {bucket_name} -> {upload_bucket}')
            # API 任务提前取任务信息（key 唯一性子目录需要 bizParam.BusinessSerialNumber；通知分支复用）
            task_info = task_info or query_task_info(cursor_old, async_vat_tasks_id)
            biz_param = parse_biz_param(task_info) or {}
            biz_serial = (biz_param.get('BusinessSerialNumber') or '').strip()
            # 兜底用任务ID（EPRRegInfoId 在 API 任务中为 null 不可靠；BusinessSerialNumber 受理校验必填，兜底不触发）
            sub_dir = biz_serial if biz_serial else str(async_vat_tasks_id or '').strip()
            # 目录名清洗：控制字符/Windows非法字符->下划线（对齐 PHP sanitize_business_code），
            # 空值兜底 unknown（BusinessSerialNumber 受理校验必填，正常不会为空）
            sub_dir = re.sub(r'[\x00-\x1f\x7f/\\:*?"<>|]', '_', sub_dir) if sub_dir else 'unknown'
            year = datetime.now().strftime('%Y')
            cos_key = os.path.join(
                API_FLOW_OSS_PREFIX.rstrip('/'),
                year,
                API_FLOW_OSS_MODULE_DIR,
                sub_dir,
                file_name
            ).replace('\\', '/')
        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
        print(f'开始上传文件: {local_file_path} -> {cos_key}')
        response = client.upload_file(
            Bucket=upload_bucket,
            LocalFilePath=local_file_path,
            Key=cos_key
        )

        # 构建结果地址（2026-08-31 契约：API 任务入库/通知均用 OSS 相对路径；source 任务拼完整 URL，原逻辑不变）
        if is_api_task:
            cos_url = cos_key
        else:
            cos_url = f'https://{upload_bucket}.cos.{region}.myqcloud.com/{cos_key}'
        print(f'文件上传成功,COS地址: {cos_url}')

        # 连接数据库(用于附件表Base_AnnexesFile)
        conn = pymssql.connect(
            server=db_ip,
            port=db_port,
            user=db_username,
            password=db_password,
            database=database,
            charset='utf8'
        )
        cursor = conn.cursor()

        # 查询并删除同一业务记录下所有旧030附件(覆盖301海牙/352两套流程；仅 source 任务)
        # 说明:301与352分属不同 EPRRegInfo 行(eri.ID 不同),但共用同一 EPRBusinessRecordId;
        #       旧逻辑按 EPRRegInfoId + TOP 1 只能删本流程上一条,跨流程旧030文件会残留;
        #       业务要求两套流程只保留一个最新030文件,故按 EPRBusinessRecordId 汇总全部旧 pdf030_fid 删除。

        if not is_api_task:
            print(f'查询同一业务记录(EPRBusinessRecordId={epr_business_record_id})下所有旧 pdf030_fid(含301/352)...')
            cursor_old.execute(
                """
                SELECT [pdf030_fid]
                FROM EPRHaiyaProcessingTasks
                WHERE EPRBusinessRecordId = %s
                AND [pdf030_fid] IS NOT NULL
                """,
                (epr_business_record_id,)
            )
            old_fids = list(dict.fromkeys(r[0] for r in cursor_old.fetchall() if r[0]))

            if old_fids:
                placeholders = ','.join(['%s'] * len(old_fids))
                _execute_write_with_retry(
                    cursor,
                    f"DELETE FROM Base_AnnexesFile WHERE F_Id IN ({placeholders})",
                    tuple(old_fids),
                    '删除旧030附件'
                )
                conn.commit()
                print(f'成功删除 {len(old_fids)} 条旧030附件记录: {old_fids}')
            else:
                print('未找到任何旧任务的 pdf030_fid,跳过删除步骤')

        # 生成新的文件ID(GUID)
        file_id = str(uuid.uuid4()).lower()

        if not is_api_task:
            # 第二步:插入新的附件记录（仅 source 流程）
            # 从URL中提取文件名(去掉.pdf后缀)
            file_name_without_ext = os.path.splitext(file_name)[0]
            file_name_with_pdf = file_name_without_ext + '.pdf'

            print(f'准备插入新附件记录, File ID: {file_id}')

            sql_insert_attachment = """
                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
                )
            """

            current_time = datetime.now().strftime('%Y-%m-%d %H:%M:%S')

            _execute_write_with_retry(cursor, sql_insert_attachment, (
                file_id,                    # F_Id
                '',                         # F_FolderId
                file_name_with_pdf,         # F_FileName
                cos_url,                    # F_FilePath
                'System',                   # CreationName
                file_size,                  # F_FileSize
                '.pdf',                     # F_FileExtensions
                'pdf',                      # F_FileType
                0,                          # F_DownloadCount
                current_time,               # CreationDate
                'System',                   # Creation_Id
                EPRRegInfoId,              # InfoId
                None,                       # FileCover
                cos_url,                    # FileUrl
                3,                          # FileFrom
                2,                          # FileGroup
                '1c199ef6-71ac-4f21-806a-f76407522e5c'  # FileCategoryId
            ), '附件表插入')
            conn.commit()
            print(f'附件表插入成功, File ID: {file_id}')

            # 更新EPRRegInfo表,设置PushTaxBureauStatus=6
            sql_update_business = """
                UPDATE EPRRegInfo
                SET PushTaxBureauStatus = %s
                WHERE ID = %s
            """
            _execute_write_with_retry(cursor, sql_update_business, (6, EPRRegInfoId), 'EPRRegInfo状态更新')
            conn.commit()
            print(f'EPRRegInfo表更新成功')
        else:
            # API 任务：不写 vat_db 附件表与 EPRRegInfo（外部项目数据，无对应业务记录）
            print('API任务: 跳过附件表插入与EPRRegInfo状态更新')

        # 更新EPRHaiyaProcessingTasks表
        # 1. 更新当前任务（精确匹配 Id）
        sql_update_current_task = """
            UPDATE EPRHaiyaProcessingTasks
            SET pdf_result_url = %s,
                pdf_result = %s,
                [pdf030_fid] = %s
            WHERE Id = %s
        """
        _execute_write_with_retry(
            cursor_old,
            sql_update_current_task,
            (cos_url, 2, file_id, async_vat_tasks_id),
            '任务表回填'
        )
        conn_old.commit()
        current_task_committed = True
        print(f'任务表 EPRHaiyaProcessingTasks 中 Id={async_vat_tasks_id} 的记录更新成功')

        # [API 补充] 把 030 的 OSS 相对路径 merge 进 ResultData.file4（统一契约 file4=030文件，
        # 2026-08-31 契约：相对路径不带域名），保留已有字段（file1=盖章要求/file2=海牙/file3=APODERAMIENTO
        # 由 PHP 写入）。source 任务 ResultData 为其他结构（附件表+列字段），不在此改动；merge 失败不阻断主流程。
        if is_api_task:
            try:
                merge_file4_into_result_data(cursor_old, async_vat_tasks_id, cos_url)
            except Exception as merge_err:
                print(f'当前任务 file4 merge 失败(不阻断主流程): {merge_err}')

        # 2. 将同一 VATBusinessRecordId 下其他未完成的任务标记为 2（作废/覆盖）
        #    API 任务无 EPRRegInfoId（NULL），跳过此步骤（WHERE EPRRegInfoId=NULL 永远不匹配）
        if not is_api_task:
            sql_invalidate_others = """
                UPDATE EPRHaiyaProcessingTasks
                SET pdf_result = 2,
                    [pdf030_fid] = %s
                WHERE EPRRegInfoId = %s
                AND Id != %s
                AND pdf_result IN (0, 1)
            """
            _execute_write_with_retry(
                cursor_old,
                sql_invalidate_others,
                (file_id, EPRRegInfoId, async_vat_tasks_id),
                '旧任务作废'
            )
            conn_old.commit()
            updated_rows = cursor_old.rowcount
            print(f'已将 {updated_rows} 条同业务记录下的旧任务状态设为 2（作废）')

        # 3. 将其他任务(含跨流程301/352)的旧 pdf030_fid 指向最新文件,避免悬空指针（原逻辑不变，单条 UPDATE）
        #    只改 pdf030_fid、不动 pdf_result;不触碰尚未生成030文件的任务(pdf030_fid IS NULL),
        #    也不取消另一流程进行中的任务(上面 invalidate 仅处理同流程待处理任务)
        _execute_write_with_retry(
            cursor_old,
            """
            UPDATE EPRHaiyaProcessingTasks
            SET [pdf030_fid] = %s
            WHERE EPRBusinessRecordId = %s
            AND Id != %s
            AND [pdf030_fid] IS NOT NULL
            """,
            (file_id, epr_business_record_id, async_vat_tasks_id),
            '跨任务pdf030_fid重定向'
        )
        conn_old.commit()
        redirected_rows = cursor_old.rowcount
        print(f'已将 {redirected_rows} 条其他任务的 pdf030_fid 指向最新文件(避免悬空)')

        # 3b. [API 补充] 被重定向的 API 任务同步把 file4 merge 进 ResultData(与当前任务同一份最新030文件)
        #     独立于原流程之外，DataSource 列缺失或 merge 失败均不阻断主流程
        try:
            cursor_old.execute(
                """
                SELECT Id
                FROM EPRHaiyaProcessingTasks
                WHERE EPRBusinessRecordId = %s
                AND Id != %s
                AND [pdf030_fid] IS NOT NULL
                AND DataSource = 'api'
                """,
                (epr_business_record_id, async_vat_tasks_id)
            )
            api_redirect_ids = [r[0] for r in cursor_old.fetchall()]
            for api_task_id in api_redirect_ids:
                merge_file4_into_result_data(cursor_old, api_task_id, cos_url)
            if api_redirect_ids:
                print(f'已将 {len(api_redirect_ids)} 条被重定向 API 任务的 ResultData.file4 指向最新030文件')
        except Exception as merge_err:
            print(f'API任务file4 merge跳过(不阻断主流程): {merge_err}')

        # 提交事务（API 任务 conn 无写入，仅提交任务表事务）
        if is_api_task:
            conn_old.commit()
        else:
            conn.commit()
            conn_old.commit()

        # [API 补充] 030 处理完成，调用异步结果通知接口（与 es_haiya 一致；
        # 通知 data.files = file1~file4 非空文件项 [{url,name,type}]，url 为 OSS 相对路径（2026-08-31 契约，
        # 不带域名；非 usaeu 桶对象如静态盖章要求文件完整 URL 透传）；
        # 发送失败自动重试 + 每次尝试即时落库 Notify*（重试中途被中断也不丢痕迹）；通知失败不影响主流程）
        callback_result = None
        if is_api_task:
            try:
                task_info = task_info or query_task_info(cursor_old, async_vat_tasks_id)
                result_data = parse_result_data(task_info)
                files = []
                slot_files = [
                    ('file1', '其他推送文件'),      # 盖章要求静态 PDF（API 流程为同桶相对路径原样透传；回退时 vat 完整 URL 透传）
                    ('file2', '海牙文件-已认证'),   # 海牙整本
                    ('file3', '授权书'),            # APODERAMIENTO
                    ('file4', '其他推送文件'),      # 030 文件
                ]
                for slot, file_type in slot_files:
                    file_value = (result_data.get(slot) or '').strip()
                    if not file_value:
                        continue
                    file_value = strip_cos_domain(file_value, upload_bucket, region)
                    files.append({
                        'url': file_value,
                        'name': os.path.basename(file_value) or file_value,
                        'type': file_type,
                    })
                callback_result = notify_api_result(
                    async_vat_tasks_id, callback_url, True,
                    files=files, biz_param=parse_biz_param(task_info),
                    persist=lambda notify_state: save_notify_result(
                        cursor_old, conn_old, async_vat_tasks_id, notify_state)
                )
                # 兜底再落一次最终态（persist 每次已写，此处防 persist 回调异常被吞的边角场景）
                if callback_result:
                    save_notify_result(cursor_old, conn_old, async_vat_tasks_id, callback_result)
            except Exception as notify_err:
                print(f'结果通知调用失败(不阻断主流程): {notify_err}')

        return {
            'success': True,
            'message': '文件上传和数据库更新成功',
            'cos_url': cos_url,
            'file_id': file_id,
            'file_size': file_size,
            'EPRRegInfoId': EPRRegInfoId,
            'callback_success': callback_result.get('success', False) if callback_result else None,
            'callback_status': callback_result.get('status') if callback_result else None,
            'callback_attempts': callback_result.get('attempts') if callback_result else None,
            'callback_response': callback_result.get('response') if callback_result else None,
            'callback_error': callback_result.get('error') if callback_result else None,
        }

    except Exception as e:
        # 如果发生错误,回滚事务
        if conn:
            conn.rollback()
        if conn_old:
            conn_old.rollback()

        # [对齐 es_haiya fail_api_task] API 任务失败：置 pdf_result=3 终结任务（不再被 rpa_030_get 以
        # pdf_result IN (0,1) 反复重捞）+ 通知调用方（重试 + Notify* 即时落库）；
        # 例外：当前任务回填已提交（pdf_result=2）后的辅助步骤失败（如跨任务 pdf030_fid 重定向死锁耗尽）
        # 不回改 3——030 文件已上传且任务行已完整，置 3 会把成功态改死，只通知异常供人工核查
        if is_api_task:
            try:
                task_info = task_info or query_task_info(cursor_old, async_vat_tasks_id)
                if current_task_committed:
                    notify_api_result(
                        async_vat_tasks_id, callback_url, False,
                        f'030文件已上传且回填成功，但后续辅助步骤失败: {e}',
                        biz_param=parse_biz_param(task_info),
                        persist=lambda notify_state: save_notify_result(
                            cursor_old, conn_old, async_vat_tasks_id, notify_state)
                    )
                else:
                    fail_api_task(cursor_old, conn_old, async_vat_tasks_id, str(e), callback_url, task_info)
            except Exception as notify_err:
                print(f'结果通知调用失败(不阻断主流程): {notify_err}')

        return {
            'success': False,
            'message': f'操作失败: {str(e)}',
            'error': str(e)
        }
        
    finally:
        # 关闭游标和连接
        if cursor:
            cursor.close()
        if conn:
            conn.close()
        if cursor_old:
            cursor_old.close()
        if conn_old:
            conn_old.close()


def merge_file4_into_result_data(cursor_old, task_id, file4_value):
    """
    把 030 文件的 OSS 相对路径 merge 进任务 ResultData.file4（保留已有字段）

    仅用于 API 任务（ResultData 为 PHP 写入的 {file1,file2,file3,file4,...} JSON 对象）；
    source 任务的 ResultData 是其他结构，调用方需先判定 DataSource='api' 再调用。

    参数:
        cursor_old: EPRHaiyaProcessingTasks 所在库的游标
        task_id: 任务ID
        file4_value: 030 文件的 OSS 相对路径（API 任务 cos_url=cos_key，2026-08-31 契约不带域名）
    """
    cursor_old.execute(
        "SELECT ResultData FROM EPRHaiyaProcessingTasks WHERE Id = %s",
        (task_id,)
    )
    row = cursor_old.fetchone()
    result_data = {}
    if row and row[0]:
        try:
            parsed = json.loads(row[0])
            if isinstance(parsed, dict):
                result_data = parsed
        except (ValueError, TypeError):
            result_data = {}

    result_data['file4'] = file4_value

    cursor_old.execute(
        "UPDATE EPRHaiyaProcessingTasks SET ResultData = %s WHERE Id = %s",
        (json.dumps(result_data, ensure_ascii=False), task_id)
    )
    print(f'已把 030 文件 merge 进任务 Id={task_id} 的 ResultData.file4: {file4_value}')







