import json
import time
import urllib.parse

import pymssql
import redis
import requests
from datetime import date, datetime

# delivery 平台统一异步结果通知接口（与 ES 海牙 030 / 意大利 EPR 同模式）
# 第二次通知（ISSUED_INFO 下号信息）：申请 EORI 完成后调用，仅 DataSource='api' 行；
# 且仅在 need_eori=1（需要申请 EORI）时发送——不需要申请 EORI 时由 step2 发送
DEFAULT_RESULT_CALLBACK_URL = "https://test-cloud.usaeu.com/prod-api/delivery/rpa/autoRegisterCallback"

# 异步结果通知投递重试：首次 + 2 次重试，第 2、3 次尝试前分别等待 2s / 6s（覆盖瞬时网络/服务抖动）
DELIVERY_NOTIFY_MAX_ATTEMPTS = 3
DELIVERY_NOTIFY_BACKOFF_SECONDS = (2, 6)

# Redis 队列：source 行通知（流程结束），消费方按 source_record_id 处理
REDIS_HOST = 'gz-crs-h0jgdbez.sql.tencentcdb.com'
REDIS_PORT = 22599
REDIS_PASSWORD = 'v2yYCpAdXqeXNTgK'
REDIS_DB = 12
REDIS_QUEUE_VAT_REG_OVER = 'factory:meiou:queue:VatRegOver'

# 通知类型（回执类型）：第二次通知=下号信息
RECEIPT_TYPE_ISSUED_INFO = 'ISSUED_INFO'


def get_field(row, name):
    """按列名大小写不敏感取值（兼容 pymssql as_dict 键名大小写差异）"""
    if not row:
        return None
    for key, value in row.items():
        if key.lower() == name.lower():
            return value
    return None


def date_to_str(value):
    """
    日期列值转 'yyyy-MM-dd' 字符串（通知契约格式）。
    pymssql 读回 SQL date/datetime 列为 datetime.date/datetime 对象，直接进 payload 会
    json.dumps 抛 TypeError（Object of type date is not JSON serializable），必须先转字符串。
    """
    if value is None or value == '':
        return ''
    if isinstance(value, (datetime, date)):
        return value.strftime('%Y-%m-%d')
    return str(value)


def parse_biz_param(row):
    """解析行内 biz_param JSON；失败时回退为行内字段构造的最小 bizParam"""
    raw = get_field(row, 'biz_param')
    if isinstance(raw, str) and raw:
        try:
            parsed = json.loads(raw)
            if isinstance(parsed, dict):
                return parsed
        except json.JSONDecodeError:
            pass
    return {
        'BusinessSerialNumber': get_field(row, 'code'),
        'source_record_id': get_field(row, 'source_record_id'),
    }


def resolve_notification_url(row, delivery_callback_url):
    """通知地址优先级：受理 bizParam.callback_url > runner 传参 > 默认 delivery 地址"""
    biz = parse_biz_param(row)
    if isinstance(biz, dict) and biz.get('callback_url'):
        return biz['callback_url']
    if delivery_callback_url:
        return delivery_callback_url
    return DEFAULT_RESULT_CALLBACK_URL


def file_basename_from_url(file_url):
    """从文件 URL / 相对路径提取文件名（去查询串、锚点与 URL 编码），供异步结果通知 name 字段使用"""
    file_url = (file_url or '').strip()
    if not file_url:
        return ''
    path = urllib.parse.urlparse(file_url).path
    return urllib.parse.unquote(path.rsplit('/', 1)[-1])


def build_files(entries):
    """
    构造通知 files 数组：[{url, name, type}]。
    参数 entries: [(url, name, type), ...]；url 空串/None 过滤；name 缺省取 basename。无文件返回 None。
    """
    files = []
    for url, name, file_type in entries:
        if not isinstance(url, str) or not url.strip():
            continue
        files.append({
            'url': url.strip(),
            'name': (name or '').strip() or file_basename_from_url(url),
            'type': file_type,
        })
    return files if files else None


def send_async_result_notification(url, payload):
    """
    发送 delivery 平台统一异步结果通知
    2xx 视为送达；网络错误/超时/5xx 退避重试，最多 DELIVERY_NOTIFY_MAX_ATTEMPTS 次；
    4xx（408/429 瞬时状态除外）视为永久失败（契约错误）不重试；失败不影响主流程。

    返回:
        dict: {success, http_code, response_text, attempts, error}
    """
    headers = {'Content-Type': 'application/json; charset=utf-8'}
    request_json = json.dumps(payload, ensure_ascii=False)
    result = {
        'success': False,
        'http_code': None,
        'response_text': None,
        'attempts': 0,
        'error': None,
    }
    for attempt in range(1, DELIVERY_NOTIFY_MAX_ATTEMPTS + 1):
        result['attempts'] = attempt
        try:
            response = requests.post(
                url,
                data=request_json,
                headers=headers,
                timeout=10
            )
            result['http_code'] = response.status_code
            result['response_text'] = response.text
            if 200 <= response.status_code < 300:
                print(f"[{datetime.now().strftime('%Y-%m-%d %H:%M:%S')}] 异步结果通知 HTTP {response.status_code} 送达（第 {attempt} 次尝试）: {url}")
                result['success'] = True
                result['error'] = None
                return result
            if 400 <= response.status_code < 500 and response.status_code not in (408, 429):
                print(f"[{datetime.now().strftime('%Y-%m-%d %H:%M:%S')}] 异步结果通知 HTTP {response.status_code} 永久失败不重试: {url}")
                result['error'] = f"HTTP {response.status_code}"
                return result
            result['error'] = f"HTTP {response.status_code}"
            print(f"[{datetime.now().strftime('%Y-%m-%d %H:%M:%S')}] 异步结果通知 HTTP {response.status_code}（第 {attempt}/{DELIVERY_NOTIFY_MAX_ATTEMPTS} 次尝试）: {url}")
        except Exception as e:
            # 网络异常：清除上一次尝试的响应，保证日志只反映最后一次尝试的真实状态
            result['error'] = str(e)
            result['http_code'] = None
            result['response_text'] = None
            print(f"[{datetime.now().strftime('%Y-%m-%d %H:%M:%S')}] 异步结果通知失败（第 {attempt}/{DELIVERY_NOTIFY_MAX_ATTEMPTS} 次尝试）: {e}")
        if attempt < DELIVERY_NOTIFY_MAX_ATTEMPTS:
            wait_seconds = DELIVERY_NOTIFY_BACKOFF_SECONDS[attempt - 1]
            print(f"[{datetime.now().strftime('%Y-%m-%d %H:%M:%S')}] {wait_seconds} 秒后重试...")
            time.sleep(wait_seconds)
    return result


def save_delivery_notify_log(record_id, notify_url, payload, notify_result,
                             db_ip, db_port, db_username, db_password, database):
    """
    将 delivery 异步结果通知的请求参数与投递结果写入 uk_vat_register.delivery_notify_log（JSON）。
    列不存在（旧库未迁移）等错误仅告警，不影响主流程。
    """
    log = {
        'url': notify_url,
        'status': 'success' if notify_result['success'] else 'failed',
        'http_code': notify_result.get('http_code'),
        'attempts': notify_result.get('attempts'),
        'request': payload,
        'response': notify_result.get('response_text') or notify_result.get('error'),
    }
    conn = None
    cursor = None
    try:
        conn = pymssql.connect(
            server=db_ip,
            port=db_port,
            user=db_username,
            password=db_password,
            database=database,
            charset='utf8'
        )
        cursor = conn.cursor()
        cursor.execute(
            "UPDATE uk_vat_register SET delivery_notify_log = %s WHERE id = %s",
            (json.dumps(log, ensure_ascii=False), record_id)
        )
        conn.commit()
        print(f"异步结果通知日志已写入 uk_vat_register(id={record_id})")
    except Exception as e:
        print(f"记录异步结果通知日志失败（不影响主流程）: {e}")
    finally:
        if cursor:
            cursor.close()
        if conn:
            conn.close()


def save_notify2_columns(record_id, payload, notify_result,
                         db_ip, db_port, db_username, db_password, database):
    """
    将轮2（ISSUED_INFO 下号信息）通知的请求/响应/状态/次数/时间写入 uk_vat_register 的
    Notify2Request / Notify2Response / Notify2Status / Notify2Attempts / Notify2Time。
    列不存在（旧库未迁移）等错误仅告警，不影响主流程。
    """
    conn = None
    cursor = None
    try:
        conn = pymssql.connect(
            server=db_ip,
            port=db_port,
            user=db_username,
            password=db_password,
            database=database,
            charset='utf8'
        )
        cursor = conn.cursor()
        cursor.execute(
            "UPDATE uk_vat_register SET Notify2Request = %s, Notify2Response = %s,"
            " Notify2Status = %s, Notify2Attempts = %s, Notify2Time = GETDATE() WHERE id = %s",
            (
                json.dumps(payload, ensure_ascii=False),
                notify_result.get('response_text') or notify_result.get('error'),
                'success' if notify_result['success'] else 'failed',
                notify_result.get('attempts'),
                record_id,
            )
        )
        conn.commit()
        print(f"轮2通知结果已写入 uk_vat_register Notify2* 列(id={record_id})")
    except Exception as e:
        print(f"写入 Notify2* 列失败（不影响主流程）: {e}")
    finally:
        if cursor:
            cursor.close()
        if conn:
            conn.close()


def push_to_redis_queue(source_record_id):
    """
    向 Redis 队列推送注册流程结束消息（source 行专用）
    返回: bool 推送成功返回 True，失败返回 False
    """
    redis_client = None
    try:
        redis_client = redis.Redis(
            host=REDIS_HOST,
            port=REDIS_PORT,
            password=REDIS_PASSWORD,
            db=REDIS_DB,
            decode_responses=True
        )
        redis_client.ping()
        data = {"Id": source_record_id}
        json_data = json.dumps(data)
        result = redis_client.lpush(REDIS_QUEUE_VAT_REG_OVER, json_data)
        print(f"成功推送消息到队列: {json_data}")
        print(f"当前队列长度: {result}")
        return True
    except redis.ConnectionError as e:
        print(f"Redis 连接失败: {e}")
        return False
    except Exception as e:
        print(f"推送失败: {e}")
        return False
    finally:
        if redis_client:
            try:
                redis_client.close()
            except Exception:
                pass


def query_record(record_id, db_ip, db_port, db_username, db_password, database):
    """
    查询记录（DataSource / need_eori / source_record_id / biz_param 及下号相关列）
    表结构缺少新列（旧库）时对应字段返回空，调用方按 source 行 / need_eori=0 处理
    """
    conn = None
    cursor = None
    try:
        conn = pymssql.connect(
            server=db_ip,
            port=db_port,
            user=db_username,
            password=db_password,
            database=database,
            charset='utf8',
            as_dict=True
        )
        cursor = conn.cursor()
        cursor.execute(
            "SELECT DataSource, biz_param, code, source_record_id, need_eori,"
            " vat_number, declaration_deadline, declaration_start_date,"
            " declaration_end_date, vatnumber_registration_date "
            "FROM uk_vat_register WHERE id = %s",
            (record_id,)
        )
        row = cursor.fetchone()
        return row if row else {}
    except pymssql.Error as e:
        print(f"查询记录失败（按 source 行处理）: {e}")
        return {}
    finally:
        if cursor:
            cursor.close()
        if conn:
            conn.close()


def update_database(db_ip, db_port, db_username, db_password, database,
                    record_id, submit_data, mtd_account=None, mtd_password=None,
                    mtd_key=None, vat_number=None, declaration_deadline=None,
                    declaration_start_date=None, declaration_end_date=None,
                    vatnumber_registration_date=None):
    """
    保存申请 EORI 结果：mtd_*/下号相关列与 submit_backup_data（本步提交数据快照）落库
    """
    conn = None
    cursor = None

    try:
        conn = pymssql.connect(
            server=db_ip,
            port=db_port,
            user=db_username,
            password=db_password,
            database=database,
            charset='utf8'
        )
        cursor = conn.cursor()
        sql = """
            UPDATE uk_vat_register
            SET submit_backup_data = %s,
                mtd_account = %s,
                mtd_password = %s,
                mtd_key = %s,
                vat_number = %s,
                declaration_deadline = %s,
                declaration_start_date = %s,
                declaration_end_date = %s,
                vatnumber_registration_date = %s,
                updated_date = GETDATE()
            WHERE id = %s
        """
        cursor.execute(sql, (
            submit_data,
            mtd_account,
            mtd_password,
            mtd_key,
            vat_number,
            declaration_deadline,
            declaration_start_date,
            declaration_end_date,
            vatnumber_registration_date,
            record_id
        ))
        conn.commit()

    finally:
        if cursor:
            cursor.close()
        if conn:
            conn.close()


def send_callback_request(
    record_id,
    db_ip,
    db_port,
    db_username,
    db_password,
    database,
    mtd_account=None,
    mtd_password=None,
    mtd_key=None,
    vat_number=None,
    declaration_deadline=None,
    declaration_start_date=None,
    declaration_end_date=None,
    vatnumber_registration_date=None,
    eori_number=None,
    vat_certificate_file=None,
    eori_screenshot_file=None,
    delivery_callback_url=None
):
    """
    第三步：申请 EORI 完成后保存结果（写库 + 按来源分流）。
    写库（mtd_*、vat_number、declaration_*、vatnumber_registration_date 与 submit_backup_data 快照）后分流：
      - source 行：本步为流程终点，推送 Redis 队列（VatRegOver，Id=source_record_id）；
      - api 行且 need_eori=1：发第二次异步结果通知 ISSUED_INFO（下号信息，字段从表取值）；
      - api 行且 need_eori=0：不发（第二步已发）。
    source 与 api 是两套完全不同的处理流程，不存在回退或兼容。

    参数:
        record_id (int): 数据库记录ID
        db_ip (str): 数据库IP地址
        db_port (int): 数据库端口
        db_username (str): 数据库用户名
        db_password (str): 数据库密码
        database (str): 数据库名称
        mtd_account (str, optional): MTD账号
        mtd_password (str, optional): MTD密码
        mtd_key (str, optional): MTD秘钥
        vat_number (str, optional): VAT税号
        declaration_deadline (str, optional): 申报截止时间（如 2026-12-07）
        declaration_start_date (str, optional): 首次申报开始时间（如 2026-08-20）
        declaration_end_date (str, optional): 首次申报结束时间（如 2026-10-31）
        vatnumber_registration_date (str, optional): 税号有效开始期（如 2026-08-20）
        eori_number (str, optional): EORI号，缺省固定拼'GB'+vatNumber+'000'
        vat_certificate_file (str or list, optional): VAT税号证书文件（source 行完整 URL / api 行 OSS 相对路径，字符串或列表）
        eori_screenshot_file (str or list, optional): 申请EORI结果截图文件（source 行完整 URL / api 行 OSS 相对路径，字符串或列表）
        delivery_callback_url (str, optional): delivery 统一异步结果通知地址（仅 api 行使用），
            缺省用 DEFAULT_RESULT_CALLBACK_URL

    返回:
        dict: {'notified': bool, 'error': str}
    """
    if isinstance(vat_certificate_file, str):
        vat_certificate_array = [vat_certificate_file] if vat_certificate_file else []
    elif isinstance(vat_certificate_file, list):
        vat_certificate_array = vat_certificate_file
    else:
        vat_certificate_array = []

    if isinstance(eori_screenshot_file, str):
        eori_screenshot_array = [eori_screenshot_file] if eori_screenshot_file else []
    elif isinstance(eori_screenshot_file, list):
        eori_screenshot_array = eori_screenshot_file
    else:
        eori_screenshot_array = []

    # 记录来源：source 行推 Redis 队列，api 行调用异步结果通知，两套处理完全不同
    row_source = query_record(record_id, db_ip, db_port, db_username, db_password, database)
    data_source = str(get_field(row_source, 'DataSource') or 'source').lower()
    source_record_id = get_field(row_source, 'source_record_id')
    need_eori = int(get_field(row_source, 'need_eori') or 0)

    result = {'notified': False, 'error': None}

    try:
        update_database(
            db_ip,
            db_port,
            db_username,
            db_password,
            database,
            record_id,
            json.dumps({
                'mtd_account': mtd_account,
                'mtd_password': mtd_password,
                'mtd_key': mtd_key,
                'vat_number': vat_number,
                'declaration_deadline': declaration_deadline,
                'declaration_start_date': declaration_start_date,
                'declaration_end_date': declaration_end_date,
                'vatnumber_registration_date': vatnumber_registration_date,
                'eori_number': eori_number,
                'vat_certificate_file': vat_certificate_array,
                'eori_screenshot_file': eori_screenshot_array,
            }, ensure_ascii=False),
            mtd_account,
            mtd_password,
            mtd_key,
            vat_number,
            declaration_deadline,
            declaration_start_date,
            declaration_end_date,
            vatnumber_registration_date
        )
        print(f"[{datetime.now().strftime('%Y-%m-%d %H:%M:%S')}] 数据库更新成功, 记录ID: {record_id}")
    except Exception as e:
        result['error'] = f"数据库更新失败: {str(e)}"
        print(f"[{datetime.now().strftime('%Y-%m-%d %H:%M:%S')}] 错误: {result['error']}")
        return result

    # ---------- source 行：本步即流程终点，推 Redis 结束队列 ----------
    if data_source != 'api':
        if source_record_id:
            push_to_redis_queue(source_record_id)
        else:
            print(f"警告: source 行缺少 source_record_id, 无法推送 Redis 队列 (id={record_id})")
        return result

    # ---------- api 行：need_eori=1 时发第二次通知，need_eori=0 时不发（第二步已发） ----------
    if need_eori != 1:
        print(f"该记录不需要申请 EORI（need_eori=0），第三步不发通知（id={record_id}）")
        return result

    # 通知字段从表取值（刚写库，保证通知书 = 库值）
    row_after = query_record(record_id, db_ip, db_port, db_username, db_password, database)
    notify_url = resolve_notification_url(row_after, delivery_callback_url)
    biz_param = parse_biz_param(row_after)
    vat_number_db = get_field(row_after, 'vat_number')
    payload = {
        'code': 200,
        'msg': 'success',
        'ProcessMode': 'async',
        'data': {
            'receiptType': RECEIPT_TYPE_ISSUED_INFO,
            'vatNumber': vat_number_db or '',
            'declarationDeadline': date_to_str(get_field(row_after, 'declaration_deadline')),
            'firstDeclarationPeriodStart': date_to_str(get_field(row_after, 'declaration_start_date')),
            'firstDeclarationPeriodEnd': date_to_str(get_field(row_after, 'declaration_end_date')),
            'vatEffectiveDate': date_to_str(get_field(row_after, 'vatnumber_registration_date')),
            'eoriNumber': eori_number or ('GB' + (vat_number_db or '') + '000'),
            'files': build_files([
                (url, '', 'VAT_CERTIFICATE_FILE') for url in vat_certificate_array
            ] + [
                (url, '', 'EORI_APPLICATION_RESULT_SCREENSHOT') for url in eori_screenshot_array
            ]),
        },
        'bizParam': biz_param,
    }
    notify_result = send_async_result_notification(notify_url, payload)
    save_notify2_columns(record_id, payload, notify_result,
                         db_ip, db_port, db_username, db_password, database)
    save_delivery_notify_log(record_id, notify_url, payload, notify_result,
                             db_ip, db_port, db_username, db_password, database)
    result['notified'] = True
    return result
