import pymssql
import os
import re
import json
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 程序固定配置，可通过参数覆盖）
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流程生成文件子目录（年份文件夹之后；海牙合并 PDF 与 030 文件同目录存放，文件命名不变）
API_FLOW_OSS_MODULE_DIR = 'es_vat_haiya'

# API流程上传桶（与 source 流程桶 vat-1259285998 不同）：COS 桶分流规则下 API 流程结果归 SaaS 共享桶，
# 默认 usaeu-1259285998（与 PHP .env.example API_FLOW_COS_BUCKET 同值）；RPA 进程环境与项目 .env 隔离，
# 不支持环境变量覆盖，覆盖通道为 api_bucket_name 参数
API_FLOW_BUCKET_DEFAULT = 'usaeu-1259285998'


def resolve_api_flow_upload_bucket(api_bucket_name=None):
    """
    解析 API 流程实际上传桶（COS 桶分流：API 流程结果归 SaaS 共享桶，与 source 流程桶不同）
    优先级：api_bucket_name 参数 > API_FLOW_BUCKET_DEFAULT

    参数:
        api_bucket_name: API 流程专用桶（可选）

    返回:
        str: API 流程实际上传桶
    """
    return api_bucket_name or API_FLOW_BUCKET_DEFAULT


# ---------------------------------------------------------------------------
# 2026-08-31 契约：API 流程结果入库/通知均为 OSS 相对路径（不带域名），不再生成 COS 签名 URL；
# COS 签名函数（build_cos_client/build_signed_cos_url/sign_cos_url）已随通知逻辑一并移除，
# rpa_030_get.py 内同名副本仍服务于取数流程（输入文件下载签名），本脚本不再需要。
# ---------------------------------------------------------------------------


def strip_cos_domain(url, bucket='', region=''):
    """
    将 COS 完整下载URL裁剪为对象相对路径（契约：入库/结果通知 files.url 为 OSS 相对路径、不含域名）
    - 相对路径（无 http(s):// 前缀）：原样返回
    - 完整 URL 命中 COS 桶域名（bucket 非空时校验，如 usaeu-1259285998.cos.ap-guangzhou.myqcloud.com）：
      取路径部分并剥离 ?sign= 等查询串（签名 URL / 无签名直连 URL 同规则）
    - 其他域名（如文件服务器）或无法识别：原样返回（不落入相对路径契约时透传兜底）

    参数:
        url: 完整URL或COS相对路径
        bucket: COS 桶名（为空时不做域名校验，直接裁剪）
        region: COS 地域,如 'ap-guangzhou'

    返回:
        str: OSS 相对路径
    """
    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('/'))


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, task_id):
    """
    查询任务记录信息，用于判断 API/source 流程分支

    参数:
        cursor: AsyncVATTasks 所在库的游标
        task_id: 任务ID

    返回:
        dict: {data_source, task_data, result_data}
        查询失败或任务不存在返回 None（按 source 流程处理，兼容未迁移 DataSource 列的旧环境）
    """
    try:
        cursor.execute(
            "SELECT DataSource, TaskData, ResultData FROM AsyncVATTasks WHERE Id = %s",
            (task_id,)
        )
        row = cursor.fetchone()
    except pymssql.ProgrammingError:
        # 旧环境未迁移 DataSource 列：无法判定，按 source 流程处理
        return None
    if not row:
        return None
    return {
        'data_source': (row[0] or '').strip().lower() if row[0] is not None else '',
        'task_data': row[1],
        'result_data': row[2],
    }


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/cos_key/cos_url，海牙文件 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 相对路径（不带域名、
      不生成签名，域名由接收方拼接下载）；有几个文件就接几个文件项（ES流程固定：海牙合并PDF + 030文件）
    - 失败: code=500, msg=错误原因, data=null（非 200 时 data 统一为 null，错误原因放 msg）

    重试: CALLBACK_MAX_RETRIES（总尝试=1+重试次数）+ CALLBACK_RETRY_DELAY 间隔，与 PHP 侧一致；
    2xx 视为成功，不捕获对方响应体中的业务失败。

    即时落库（2026-08-29 线上问题）：传入 persist 回调时，每次尝试结束后（含失败尝试）立即调用
    persist(result) 落库 AsyncVATTasks.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):
    """
    AsyncVATTasks 的 Notify* 列是否已执行迁移（add_notify_fields_to_async_tasks.sql）
    旧环境返回 False：跳过落库，与 PHP AsyncTaskManager::columnExists 逻辑一致
    """
    try:
        cursor.execute(
            "SELECT COUNT(*) FROM INFORMATION_SCHEMA.COLUMNS "
            "WHERE TABLE_NAME = 'AsyncVATTasks' 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):
    """
    保存结果通知的请求参数与结果到 AsyncVATTasks.Notify* 字段

    参数:
        cursor/conn: AsyncVATTasks 所在库的游标和连接
        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'AsyncVATTasks Notify* 列未迁移，跳过结果通知落库 (Id={task_id})')
            return False
        sql = """
            UPDATE AsyncVATTasks
            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'AsyncVATTasks 表中 Id={task_id} 的结果通知记录已保存（Notify* 字段）')
        return True
    except Exception as e:
        print(f"保存结果通知记录出错 (Id={task_id}): {e}")
        return False


def resolve_result_file_url(result_data, bucket='', region=''):
    """
    从 ResultData 解析海牙合并PDF的OSS相对路径（契约：入库/结果通知 files.url 为相对路径、不含域名）
    优先 cos_url，缺失时用 file1（OSS相对路径）；两者为完整 URL 时统一裁剪为相对路径（strip_cos_domain）
    """
    url = (result_data.get('cos_url') or '').strip()
    if url:
        return strip_cos_domain(url, bucket, region)
    file1 = (result_data.get('file1') or '').strip()
    return strip_cos_domain(file1, bucket, region)


def save_api_task_result(cursor, conn, task_id, cos_url, task_info, callback_url, file_size, secret_id=None, secret_key=None, bucket=None, region=None):
    """
    API流程 030 保存：不写 vat_db 附件表
    1. 更新 AsyncVATTasks：pdf030_result=2、pdf030_result_url=cos_url，
       ResultData merge 030 结果（与 PHP completePdf030 行为一致）
    2. 调用结果通知接口告知成功（data.files：海牙文件相对路径 + 030文件相对路径 + bizParam；
       files.url 为 OSS 相对路径（不带域名、不签名），域名由接收方拼接，2026-08-31 契约）

    参数:
        cursor/conn: AsyncVATTasks 所在库的游标和连接
        task_id: 任务ID
        cos_url: 030文件OSS相对路径（入库与通知 files.url 共用）
        task_info: query_task_info 的返回
        callback_url: 结果通知接口地址
        file_size: 文件大小
        secret_id: 腾讯云 SecretId（保留兼容：结果通知不再签名，不再使用）
        secret_key: 腾讯云 SecretKey（保留兼容：不再使用）
        bucket: COS 桶名（相对路径裁剪校验用）
        region: COS 地域,如 'ap-guangzhou'

    返回:
        dict: 操作结果（含回调结果）
    """
    current_time = datetime.now().strftime('%Y-%m-%d %H:%M:%S')
    result_data = parse_result_data(task_info)
    result_data.update({
        'pdf030_url': cos_url,
        'pdf030_result': 2,
        'pdf030_completed_at': current_time,
    })
    new_result_json = json.dumps(result_data, ensure_ascii=False)

    # 更新任务记录（rpa库）
    sql = """
        UPDATE AsyncVATTasks
        SET pdf030_result = 2,
            pdf030_result_url = %s,
            ResultData = %s
        WHERE Id = %s
    """
    cursor.execute(sql, (cos_url, new_result_json, task_id))
    conn.commit()
    print(f'AsyncVATTasks 表中 Id={task_id} 的记录更新成功（API流程）')

    # 调用结果通知接口：组装 files 文件数组（按 §5.4 槽位顺序：海牙合并PDF → 030文件）
    # files.url 为 OSS 相对路径（不带域名、不生成签名），域名由接收方按自身规则拼接（2026-08-31 契约）
    biz_param = parse_biz_param(task_info)
    files = []
    hague_url = resolve_result_file_url(result_data, bucket, region)
    if hague_url:
        file_name = os.path.basename(hague_url) or hague_url
        files.append({
            'url': hague_url,
            'name': file_name,
            'type': '海牙文件-已认证',
        })
    files.append({
        'url': cos_url,
        'name': os.path.basename(cos_url) or cos_url,
        'type': '其他推送文件',
    })
    # 通知即时落库：每次尝试后（含失败）写 Notify* 字段，重试中途流程被中断也不丢已发生的尝试信息
    notify_result = notify_api_result(
        task_id, callback_url, True, files=files, biz_param=biz_param,
        persist=lambda notify_state: save_notify_result(cursor, conn, task_id, notify_state)
    )
    # 兜底再落一次最终态（persist 每次已写，此处防 persist 回调异常被吞的边角场景）
    save_notify_result(cursor, conn, task_id, notify_result)

    return {
        'success': True,
        'message': '文件上传和数据库更新成功（API流程，未写附件表）',
        'cos_url': cos_url,
        'file_size': file_size,
        'async_vat_tasks_id': task_id,
        'callback_success': notify_result.get('success', False),
        'callback_status': notify_result.get('status'),
        'callback_attempts': notify_result.get('attempts'),
        'callback_response': notify_result.get('response'),
        'callback_error': notify_result.get('error'),
    }


def fail_api_task(cursor, conn, task_id, error, callback_url, task_info):
    """
    API流程失败处理：pdf030_result=3 + 调用结果通知接口告知失败
    通知失败不影响任务状态更新（与 PHP AsyncResultNotifier 一致）

    参数:
        cursor/conn: AsyncVATTasks 所在库的游标和连接
        task_id: 任务ID
        error: 失败原因
        callback_url: 结果通知接口地址
        task_info: query_task_info 的返回
    """
    try:
        sql = """
            UPDATE AsyncVATTasks
            SET pdf030_result = 3
            WHERE Id = %s
        """
        cursor.execute(sql, (task_id,))
        conn.commit()
        print(f'AsyncVATTasks 表中 Id={task_id} 的记录已标记 030 失败')
    except Exception as e:
        print(f"更新任务失败状态出错: {e}")

    try:
        # 通知即时落库：每次尝试后（含失败）写 Notify* 字段，重试中途流程被中断也不丢已发生的尝试信息
        notify_result = notify_api_result(
            task_id, callback_url, False, error, biz_param=parse_biz_param(task_info),
            persist=lambda notify_state: save_notify_result(cursor, conn, task_id, notify_state)
        )
        # 兜底再落一次最终态（persist 每次已写，此处防 persist 回调异常被吞的边角场景）
        save_notify_result(cursor, conn, task_id, notify_result)
    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, vat_business_record_id,
                                  async_vat_tasks_id, callback_url=None, api_bucket_name=None):
    """
    上传本地文件到腾讯云COS,并更新数据库中的相关记录

    API流程（任务 DataSource='api'）：
        - 上传桶按 COS 桶分流（API 与 source 流程桶不同）：api_bucket_name 参数（可选）>
          默认 usaeu-1259285998（SaaS 共享桶）；
        - COS 对象键 = {API_FLOW_OSS_PREFIX}{当前年}/es_vat_haiya/{文件名}（相对路径，2026-08-31 契约，
          不再含时间戳子文件夹；海牙合并PDF与030文件同目录，文件命名不变）
        - 入库（pdf030_result_url / ResultData.cos_url）与结果通知 files.url 均为 OSS 相对路径
          （不带域名、不生成 COS 签名）
        - 不写 vat_db 附件表（Base_AnnexesFile）与 VATRegInfo
        - 仅更新 AsyncVATTasks 任务记录（pdf030_result / pdf030_result_url / ResultData）
        - 调用异步结果通知接口，将成功或失败结果发给调用方（URL 默认 DEFAULT_CALLBACK_URL，可传 callback_url 覆盖）
    source流程（原逻辑不变）：写 vat_db 附件表 + VATRegInfo + 更新 AsyncVATTasks

    参数:
        db_ip: 数据库服务器IP
        db_port: 数据库端口
        db_username: 数据库用户名
        db_password: 数据库密码
        database: 数据库名(用于附件表Base_AnnexesFile)
        dbname: 数据库名(用于AsyncVATTasks表)
        secret_id: 腾讯云SecretId
        secret_key: 腾讯云SecretKey
        region: 腾讯云地域,如 'ap-guangzhou'
        bucket_name: COS存储桶名称（source流程桶）
        cos_path: COS目录路径,如 'vat/spain/'
        local_file_path: 本地文件路径
        vat_business_record_id: VATBusinessRecord表的ID
        attachment_id: 附件InfoId
        async_vat_tasks_id: 当前AsyncVATTasks任务ID
        callback_url: 结果通知接口地址（API流程使用；为空时用 DEFAULT_CALLBACK_URL）
        api_bucket_name: API流程上传桶（可选，API流程使用；为空时用默认值 usaeu-1259285998）

    返回:
        dict: 包含操作结果的字典
    """
    conn = None
    cursor = None
    conn_old = None
    cursor_old = None
    task_info = None
    is_api_flow = False

    try:
        # 连接rpa数据库(用于AsyncVATTasks表)，先查任务信息判断 API/source 流程分支
        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()
        task_info = query_task_info(cursor_old, async_vat_tasks_id)
        is_api_flow = bool(task_info and task_info.get('data_source') == 'api')

        # COS 桶分流：API 与 source 流程桶不同——API 流程结果归 SaaS 共享桶（usaeu-1259285998），source 流程沿用 bucket_name
        upload_bucket = bucket_name
        if is_api_flow:
            upload_bucket = resolve_api_flow_upload_bucket(api_bucket_name)
            if upload_bucket != bucket_name:
                print(f'API流程上传桶分流: {bucket_name} -> {upload_bucket}')

        # 检查本地文件是否存在
        if not os.path.exists(local_file_path):
            message = f'本地文件不存在: {local_file_path}'
            if is_api_flow:
                # API流程：标记失败并通知结果接口
                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_vat_haiya/{文件名}（2026-08-31 契约，
        # 不含时间戳子文件夹，海牙合并PDF与030文件同目录）；source 流程：cos_path + 时间戳文件夹 + 文件名（原逻辑不变）
        if is_api_flow:
            year = datetime.now().strftime('%Y')
            cos_key = os.path.join(
                API_FLOW_OSS_PREFIX.rstrip('/'),
                year,
                API_FLOW_OSS_MODULE_DIR,
                file_name
            ).replace('\\', '/')
        else:
            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_flow:
            cos_url = cos_key
        else:
            cos_url = f'https://{upload_bucket}.cos.{region}.myqcloud.com/{cos_key}'
        print(f'文件上传成功,COS地址: {cos_url}')

        # API流程：不写 vat_db 附件表，只更新任务记录 + 调用结果通知接口
        # （bucket/region 用于通知 files.url 域名裁剪校验；结果不再生成 COS 签名）
        if is_api_flow:
            return save_api_task_result(
                cursor_old, conn_old, async_vat_tasks_id, cos_url, task_info, callback_url, file_size,
                bucket=upload_bucket, region=region
            )

        # 连接数据库(用于附件表Base_AnnexesFile)
        conn = pymssql.connect(
            server=db_ip,
            port=db_port,
            user=db_username,
            password=db_password,
            database=database,
            charset='utf8'
        )
        cursor = conn.cursor()

        # 第一步:查询并删除旧的附件记录
        print(f'查询旧任务的030AttachmentFid...')
        sql_query_old = """
            SELECT TOP 1 Id, [030AttachmentFid]
            FROM AsyncVATTasks
            WHERE VATBusinessRecordId = %s
            AND Id < %s
            ORDER BY Id DESC
        """
        cursor_old.execute(sql_query_old, (vat_business_record_id, async_vat_tasks_id))
        old_task = cursor_old.fetchone()

        if old_task and old_task[1] is not None:  # 明确排除 NULL
            old_attachment_fid = old_task[1]
            print(f'找到旧任务附件ID: {old_attachment_fid}, 准备删除...')

            sql_delete_attachment = """
                DELETE FROM Base_AnnexesFile
                WHERE F_Id = %s
            """
            cursor.execute(sql_delete_attachment, (old_attachment_fid,))
            print(f'成功删除旧附件记录: {old_attachment_fid}')
        else:
            print('未找到旧任务的 030AttachmentFid 或其为 NULL，跳过删除步骤')

        # 第二步:插入新的附件记录
        # 生成新的文件ID(GUID)
        file_id = str(uuid.uuid4()).lower()

        # 从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')

        cursor.execute(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
            vat_business_record_id,              # InfoId
            None,                       # FileCover
            cos_url,                    # FileUrl
            3,                          # FileFrom
            2,                          # FileGroup
            '1c199ef6-71ac-4f21-806a-f76407522e5c'  # FileCategoryId
        ))
        print(f'附件表插入成功, File ID: {file_id}')

        # 更新VATBusinessRecord表,设置PushTaxBureauStatus=7
        sql_update_business = """
            UPDATE VATRegInfo
            SET PushTaxBureauStatus = %s
            WHERE ID = %s
        """
        cursor.execute(sql_update_business, (6, vat_business_record_id))
        print(f'VATBusinessRecord表更新成功')

        # 更新AsyncVATTasks表
        # 1. 更新当前任务（精确匹配 Id）
        sql_update_current_task = """
            UPDATE AsyncVATTasks
            SET pdf030_result_url = %s,
                pdf030_result = %s,
                [030AttachmentFid] = %s
            WHERE Id = %s
        """
        cursor_old.execute(
            sql_update_current_task,
            (cos_url, 2, file_id, async_vat_tasks_id)
        )
        print(f'AsyncVATTasks 表中 Id={async_vat_tasks_id} 的记录更新成功')

        # 2. 将同一 VATBusinessRecordId 下其他未完成的任务标记为 2（作废/覆盖）
        sql_invalidate_others = """
            UPDATE AsyncVATTasks
            SET pdf030_result = 2,
                [030AttachmentFid] = %s
            WHERE VATBusinessRecordId = %s
            AND Id != %s
            AND pdf030_result IN (0, 1)
        """
        cursor_old.execute(
            sql_invalidate_others,
            (file_id, vat_business_record_id, async_vat_tasks_id)
        )
        updated_rows = cursor_old.rowcount
        print(f'已将 {updated_rows} 条同业务记录下的旧任务状态设为 2（作废）')

        # 提交事务
        conn.commit()
        conn_old.commit()

        return {
            'success': True,
            'message': '文件上传和数据库更新成功',
            'cos_url': cos_url,
            'file_id': file_id,
            'file_size': file_size,
            'vat_business_record_id': vat_business_record_id
        }

    except Exception as e:
        # 如果发生错误,回滚事务
        if conn:
            conn.rollback()
        if conn_old:
            conn_old.rollback()

        # API流程失败：通知结果接口
        if is_api_flow and task_info:
            try:
                fail_api_task(cursor_old, conn_old, async_vat_tasks_id, str(e), callback_url, task_info)
            except Exception as notify_error:
                print(f"API流程失败通知出错: {notify_error}")

        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()
