import json
import re
import ssl
import time
import urllib.parse
import urllib.request
import pymssql
from datetime import datetime, timedelta
from concurrent.futures import ThreadPoolExecutor, TimeoutError as FutureTimeoutError
from qcloud_cos import CosConfig, CosS3Client

# ---- 超时配置（单位：秒），可根据实际网络情况调整 ----
DB_LOGIN_TIMEOUT = 10      # 建立数据库连接的超时时间
DB_QUERY_TIMEOUT = 30      # 单次查询/执行的超时时间
OVERALL_TIMEOUT = 60       # 整个函数执行的总超时时间（看门狗超时）

# ---- 异步结果回调配置（RPA 进程不读 .env，此处为部署常量）----
# ⚠️ 本文件与 rpa_save.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

# ---- API 流程 OSS（须与 PHP 侧 EPR_API_COS_BUCKET / EPR_API_OSS_PREFIX 一致）----
# API 桶为私有桶：下载对象必须带签名（见 sign_api_file_url）；source 流程用 file.usaeu.com 无需签名
API_FLOW_BUCKET_DEFAULT = 'usaeu-1259285998'
API_FLOW_REGION_DEFAULT = 'ap-guangzhou'
API_FLOW_OSS_BASE_URL = 'https://{}.cos.{}.myqcloud.com'.format(API_FLOW_BUCKET_DEFAULT, API_FLOW_REGION_DEFAULT)
# 签名有效期（秒）：须覆盖 RPA 取数后到实际下载完成的间隔
SIGN_EXPIRE_SECONDS = 600

# 与 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

# 授权书阶段认领次数上限（与认领 SQL 的 auth_count 门限一致），达到后判死并回调
AUTH_MAX_ATTEMPTS = 30

# 判死时的在途陈旧窗口（分钟）：auth_status=1 的行在此时长内视为「正在处理」，不回收，
# 避免把刚被认领、RPA 正在生成授权书的行判死（会与随后的成功回调冲突，产生两条终态回调）
AUTH_STALE_MINUTES = 30

# 单轮最多回收的耗尽行数（回收含回调投递，必须限制条数才能守住看门狗预算）
RELEASE_BATCH_LIMIT = 1

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')
        req = urllib.request.Request(
            url,
            data=data,
            headers={'Content-Type': 'application/json; charset=utf-8'},
            method='POST'
        )
        with urllib.request.urlopen(req, timeout=timeout, context=CALLBACK_SSL_CONTEXT) as resp:
            body = resp.read().decode('utf-8', errors='replace')
            return {
                'success': 200 <= resp.getcode() < 300,
                'http_code': resp.getcode(),
                'response': body[:500],
                'error': None,
            }
    except Exception as e:
        return {'success': False, 'http_code': None, 'response': None, 'error': str(e)}


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


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: 可选回调，每次投递尝试后调用 persist(result, attempts, payload)（用于写 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 resolve_signature_url(raw):
    """解析 API 文件对象（dict/list/JSON 字符串），取首个 fileUrl；普通字符串原样返回。"""
    if isinstance(raw, dict):
        return raw.get('fileUrl') or raw.get('FileUrl') or ''

    if isinstance(raw, list):
        return resolve_signature_url(raw[0]) if raw else ''

    if isinstance(raw, str) and raw.strip().startswith(('[', '{')):
        try:
            return resolve_signature_url(json.loads(raw))
        except Exception:
            return raw

    return raw or ''


def build_cos_client(secret_id, secret_key, region):
    """构建 COS 客户端（仅用于生成预签名 URL）；凭证/地域缺失或构建失败返回 None（调用方走无签名兜底）。"""
    if not (secret_id and secret_key and region):
        return None
    try:
        return CosS3Client(CosConfig(Region=region, SecretId=secret_id, SecretKey=secret_key))
    except Exception as e:
        print('构建COS客户端失败（走无签名兜底）: {}'.format(e))
        return None


def build_signed_cos_url(client, bucket, region, url_or_key, expire_seconds=SIGN_EXPIRE_SECONDS):
    """
    生成 COS 带签名下载 URL（API 桶为私有桶，对象下载需带验证信息；与 es_haiya 的
    build_signed_cos_url / PHP CosDownloader::signBucketUrl 同口径）。

    - 相对路径：按 {bucket}.cos.{region} 域名拼 key 后签名
    - 完整 URL 命中本桶域名且未带签名：取对象键重新签名（获取执行时刻的新鲜签名）
    - 已是签名 URL / 其他域名 / 本地路径：原样返回（source 流程的 file.usaeu.com 不走这里）
    - 签名失败：原样返回（公读对象仍可直连，由调用方兜底）
    """
    url_or_key = (url_or_key or '').strip()
    if not url_or_key:
        return ''

    domain = '{}.cos.{}.myqcloud.com'.format(bucket, region)

    if re.match(r'^https?://', url_or_key, re.IGNORECASE):
        parsed = urllib.parse.urlparse(url_or_key)
        if parsed.netloc.lower() != domain.lower():
            return url_or_key
        if 'sign=' in parsed.query or 'q-sign-' in parsed.query:
            return url_or_key
        key = urllib.parse.unquote(parsed.path.lstrip('/'))
    else:
        key = url_or_key.lstrip('/')

    try:
        signed = client.get_presigned_url(Bucket=bucket, Key=key, Expired=expire_seconds, Method='get')
        if signed:
            return signed
    except Exception as e:
        print('生成签名URL失败（兜底无签名直连）: {}'.format(e))

    return url_or_key


def sign_api_file_url(url_or_key, secret_id, secret_key, bucket, region):
    """按 RPA 程序传入的凭证对 API 桶文件签名；凭证缺失/不完整时原样返回（绝不抛出）。"""
    client = build_cos_client(secret_id, secret_key, region)
    if client is None:
        return url_or_key

    return build_signed_cos_url(client, bucket, region, url_or_key)


def resolve_api_signed_file_url(raw):
    """
    API 行的法人签名文件：直接取接口传值（不查附件库）。
    相对路径按 API 桶域名补全为完整 URL；绝对 URL 原样返回。
    调用方随后按需对 API 桶对象签名（私有桶下载必须带签名）。
    """
    url = (resolve_signature_url(raw) or '').strip()

    if not url:
        return ''

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

    return '{}/{}'.format(API_FLOW_OSS_BASE_URL.rstrip('/'), url.lstrip('/'))


def process_datetime_fields(result):
    """
    处理结果中的datetime字段，将datetime对象转换为字符串

    参数:
        result: 查询结果字典

    返回:
        None: 直接修改原字典
    """
    for key, value in result.items():
        if isinstance(value, datetime):
            result[key] = value.strftime('%Y-%m-%d %H:%M:%S')


def _has_notify_columns(cursor):
    """探测 5 个 notify_* 列是否全部存在（迁移未执行时跳过落库）。"""
    global _notify_columns_probe_warned
    try:
        placeholders = ', '.join(['%s'] * len(NOTIFY_COLUMNS))
        cursor.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.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(cursor, conn, italia_epr_file_id, result, attempts, payload):
    """复用当前连接写入回调投递明细（notify_* 列），失败只告警不抛出。"""
    try:
        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


def _mark_auth_failed(cursor, row_id, error_msg):
    """把行的授权书阶段置为终态失败（auth_status=3），并记录处理时间与重试计数。"""
    cursor.execute(
        """
        UPDATE italia_epr_file
        SET auth_status       = 3,
            auth_error_msg    = %s,
            auth_process_time = GETDATE(),
            auth_count        = ISNULL(auth_count, 0) + 1
        WHERE id = %s
        """,
        (error_msg[:500], row_id)
    )

    return cursor.rowcount


def _release_exhausted_rows(cursor, conn):
    """
    回收已被反复认领但始终未完成的 API 行：置 auth_status=3 并发送失败回调。

    判定条件（缺一不可，避免误杀在途行）：
      - data_source='api' 且 desc_status=2（说明书已生成，确实卡在授权书阶段）
      - auth_status=0（从未认领）或 auth_status=1 且 auth_process_time 已超过 AUTH_STALE_MINUTES
      - ISNULL(auth_count,0) >= AUTH_MAX_ATTEMPTS
    单轮最多回收 RELEASE_BATCH_LIMIT 行；每行先 commit 再回调（避免回滚后回调已送达的不一致）。
    """
    cursor.execute(
        """
        SELECT TOP {} id, code, biz_param, callback_url
        FROM italia_epr_file
        WHERE data_source = 'api'
          AND desc_status = 2
          AND ISNULL(auth_count, 0) >= %s
          AND (
                auth_status = 0
                OR (auth_status = 1 AND (auth_process_time IS NULL
                                         OR auth_process_time < DATEADD(MINUTE, -%s, GETDATE())))
              )
        ORDER BY id ASC
        """.format(RELEASE_BATCH_LIMIT),
        (AUTH_MAX_ATTEMPTS, AUTH_STALE_MINUTES)
    )
    exhausted = cursor.fetchall() or []

    for row in exhausted:
        error_msg = '授权书阶段重试次数已达上限（{} 次），任务判死'.format(AUTH_MAX_ATTEMPTS)
        affected = _mark_auth_failed(cursor, row.get('id'), error_msg)
        conn.commit()

        if affected is None or affected == 0:
            # 已被其它参与者改掉，不再重复通知
            print('[回收] API 行 id={} 状态已被其它参与者更新，跳过回调'.format(row.get('id')))
            continue

        # 先落库再通知：行已终态，回调只投递一次（max_retries=0），失败仅影响本次通知，不会重扫重发
        notify_api_result(
            row.get('id'),
            (row.get('callback_url') or '').strip() or DEFAULT_CALLBACK_URL,
            False,
            error_msg,
            biz_param=parse_biz_param(row),
            persist=lambda result, attempts, payload: _persist_notify_result(
                cursor, conn, row.get('id'), result, attempts, payload
            ),
            max_retries=0
        )
        print('[回收] API 行 id={} 重试耗尽，已置 auth_status=3 并回调'.format(row.get('id')))


def _query_and_update_italia_epr_file_impl(ip, port, username, password, database, vat_database,
                                           secret_id=None, secret_key=None, api_bucket_name=None, api_region=None):
    """
    实际执行逻辑（原函数体），内部会被主函数用超时机制包裹调用。
    """
    conn = None
    cursor = None

    try:
        # 连接数据库（加上登录超时和查询超时）
        conn = pymssql.connect(
            server=ip,
            port=port,
            user=username,
            password=password,
            database=database,
            login_timeout=DB_LOGIN_TIMEOUT,
            timeout=DB_QUERY_TIMEOUT
        )
        cursor = conn.cursor(as_dict=True)

        # 先回收重试耗尽的 API 行（判死 + 失败回调），避免它们一直挡住队列
        try:
            _release_exhausted_rows(cursor, conn)
        except Exception as e:
            print('[WARN] 回收重试耗尽的 API 行时出错（已忽略）: {}'.format(e))
            try:
                conn.rollback()
            except Exception:
                pass

        # 查询一条符合条件的记录
        query_sql = """
        SELECT TOP 1 *
        FROM italia_epr_file
        WHERE auth_status IN (0, 1)
        AND desc_status = 2 AND (auth_count IS NULL OR auth_count < %s)
        ORDER BY id ASC
        """
        cursor.execute(query_sql, (AUTH_MAX_ATTEMPTS,))
        result = cursor.fetchone()

        if result:
            # 数据源判定：API 行（data_source='api'）不访问附件库；字段缺失/判死时发送失败回调
            is_api_task = str(result.get('data_source') or '').strip().lower() == 'api'

            # 定义必须非空的字段列表
            required_fields = [
                'company_address_en',
                'company_name_en',
                'company_city_en',
                'legal_rep_birthplace',
                'legal_rep_birth_day',
                'legal_rep_birth_month',
                'legal_rep_birth_year',
                'doc_day',
                'doc_month',
                'doc_year',
            ]

            # 检查所有必填字段是否为空（纯空白字符串视为空）
            empty_fields = []
            for field in required_fields:
                value = result.get(field)
                if value is None or str(value).strip() == '':
                    empty_fields.append(field)

            # 如果有必填字段为空，设置auth_status为3并写入错误信息
            if empty_fields:
                error_msg = '必填字段为空: {}'.format(', '.join(empty_fields))
                _mark_auth_failed(cursor, result.get('id'), error_msg)
                conn.commit()

                # API 行：同步发送失败回调（source 行无回调地址，跳过）
                if is_api_task:
                    notify_api_result(
                        result.get('id'),
                        (result.get('callback_url') or '').strip() or DEFAULT_CALLBACK_URL,
                        False,
                        error_msg,
                        biz_param=parse_biz_param(result),
                        persist=lambda r, attempts, payload: _persist_notify_result(
                            cursor, conn, result.get('id'), r, attempts, payload
                        )
                    )

                return None

            # ---- 字段处理 ----

            # datetime字段转字符串
            process_datetime_fields(result)

            # company_address_en、company_city_en 转大写
            for field in ('company_address_en', 'company_city_en'):
                if result.get(field):
                    result[field] = result[field].strip().upper()

            # legal_rep_birthplace 处理：每个单词首字母大写其余小写，去除空格拼接
            if result.get('legal_rep_birthplace'):
                parts = result['legal_rep_birthplace'].strip().split()
                result['legal_rep_birthplace'] = ''.join(
                    part.capitalize() for part in parts
                )

            # 拼接 legal_rep_birth：day/month/year
            result['legal_rep_birth'] = '{}/{}/{}'.format(
                str(result.get('legal_rep_birth_day', '')).zfill(2),
                str(result.get('legal_rep_birth_month', '')).zfill(2),
                result.get('legal_rep_birth_year', '')
            )

            # 拼接 doc_time：day/month/year
            result['doc_time'] = '{}/{}/{}'.format(
                str(result.get('doc_day', '')).zfill(2),
                str(result.get('doc_month', '')).zfill(2),
                result.get('doc_year', '')
            )

            # 处理 legal_rep_signature_file_id
            if is_api_task:
                # API 行：受理时传的就是文件 URL（或 [{fileUrl}]），不查附件库
                signature_url = resolve_api_signed_file_url(result.get('legal_rep_signature_file_id'))

                # API 桶为私有桶：交给 RPA 下载前必须先签名（凭证由 RPA 程序传入；未传则原样兜底，
                # 公读对象仍可直连）。签名 URL 含凭证信息，勿写日志/入库。
                if signature_url:
                    signature_url = sign_api_file_url(
                        signature_url,
                        secret_id,
                        secret_key,
                        (api_bucket_name or '').strip() or API_FLOW_BUCKET_DEFAULT,
                        (api_region or '').strip() or API_FLOW_REGION_DEFAULT
                    )

                result['legal_rep_signature_file_id'] = signature_url
            elif result.get('legal_rep_signature_file_id'):
                file_conn = None
                file_cursor = None
                try:
                    file_conn = pymssql.connect(
                        server=ip,
                        port=port,
                        user=username,
                        password=password,
                        database=vat_database,
                        login_timeout=DB_LOGIN_TIMEOUT,
                        timeout=DB_QUERY_TIMEOUT
                    )
                    file_cursor = file_conn.cursor()
                    file_cursor.execute(
                        "SELECT F_FilePath FROM Base_AnnexesFile WHERE F_Id = %s",
                        (result['legal_rep_signature_file_id'],)
                    )
                    file_result = file_cursor.fetchone()
                    if file_result and file_result[0]:
                        result['legal_rep_signature_file_id'] = file_result[0].replace(
                            'G:/fileAnnexes', 'https://file.usaeu.com'
                        )
                    else:
                        result['legal_rep_signature_file_id'] = ''
                except Exception as e:
                    print("处理legal_rep_signature_file_id时发生错误: {}".format(e))
                    result['legal_rep_signature_file_id'] = ''
                finally:
                    if file_cursor:
                        file_cursor.close()
                    if file_conn:
                        file_conn.close()
            else:
                result['legal_rep_signature_file_id'] = ''

            # 更新 auth_status 为 1（同时刷新在途时间戳，供回收判定判陈旧）
            update_sql = """
                UPDATE italia_epr_file
                SET auth_status = 1,
                    auth_process_time = GETDATE(),
                    auth_count = ISNULL(auth_count, 0) + 1
                WHERE id = %s
                """
            cursor.execute(update_sql, (result.get('id'),))
            conn.commit()

            return result
        else:
            return None

    except Exception as e:
        print("数据库操作错误: {}".format(e))
        raise

    finally:
        # close() 本身也可能因为连接处于"半死"状态而阻塞，
        # 这里不做额外超时保护（Python没有原生手段强制中断阻塞的C扩展调用），
        # 但外层的看门狗超时会保证整个任务不会无限期卡住调用方。
        try:
            if cursor:
                cursor.close()
        except Exception as e:
            print("关闭cursor时出错（已忽略）: {}".format(e))
        try:
            if conn:
                conn.close()
        except Exception as e:
            print("关闭连接时出错（已忽略）: {}".format(e))


def query_and_update_italia_epr_file(ip, port, username, password, database,
                                      overall_timeout=OVERALL_TIMEOUT,
                                      vat_database='vat_db',
                                      secret_id=None, secret_key=None,
                                      api_bucket_name=None, api_region=None):
    """
    查询italia_epr_file表中符合条件的一条记录，并更新其状态。
    带总超时保护：如果整体执行时间超过 overall_timeout 秒仍未完成
    （例如连接中途卡死、网络异常导致close()阻塞等），会主动放弃等待并抛出超时异常，
    调用方可以据此记录日志、告警或跳过本次任务，而不会导致主程序/调度器被永久卡住。

    数据源分流（data_source）：
      - api：签名文件直接取接口传值（不查附件库），并按传入凭证生成**带签名的下载 URL**
        （API 桶为私有桶，缺签名下载会 403）；字段缺失 / 认领耗尽时发送失败回调
      - source（默认）：行为与历史版本完全一致（file.usaeu.com，无需签名）

    参数:
        ip: 数据库IP地址
        port: 数据库端口
        username: 数据库用户名
        password: 数据库密码
        database: 数据库名称
        overall_timeout: 整体执行超时时间（秒），默认使用 OVERALL_TIMEOUT
        vat_database: 附件表所在数据库名（仅 source 数据源使用，默认 'vat_db'）
        secret_id/secret_key: 腾讯云凭证（仅 API 行签名用；缺省则不签名、原样直连兜底）
        api_bucket_name: API 桶名（缺省用部署常量 API_FLOW_BUCKET_DEFAULT）
        api_region: COS 地域（缺省用部署常量 API_FLOW_REGION_DEFAULT）

    返回:
        dict: 查询到的数据字典，如果没有查到则返回None

    异常:
        TimeoutError: 整体执行超过 overall_timeout 仍未完成时抛出
        (注意：由于底层是C扩展的阻塞调用，超时后线程可能仍在后台运行，
         无法被强制杀死，仅代表当前调用方不再等待)
    """
    # 兼容历史签名：早期版本第 6 个位置参数是 vat_database，此处按类型自动纠正，
    # 避免老调用方把库名传成 overall_timeout（会得到 "str 与 float 比较" 的 TypeError）
    if isinstance(overall_timeout, str):
        vat_database, overall_timeout = overall_timeout, OVERALL_TIMEOUT

    # 不用 with：with 退出时会 shutdown(wait=True) 等 worker 结束，
    # 使超时形同虚设；超时后直接放弃等待，与 docstring 承诺一致（可能残留 1 个线程/连接）
    executor = ThreadPoolExecutor(max_workers=1)
    future = executor.submit(
        _query_and_update_italia_epr_file_impl,
        ip, port, username, password, database, vat_database,
        secret_id, secret_key, api_bucket_name, api_region
    )
    try:
        return future.result(timeout=overall_timeout)
    except FutureTimeoutError:
        print("数据库操作整体超时（超过 {} 秒），主动放弃等待".format(overall_timeout))
        raise TimeoutError(
            "query_and_update_italia_epr_file 执行超过 {} 秒，可能是数据库连接中途异常".format(overall_timeout)
        )
    finally:
        executor.shutdown(wait=False)
