import hashlib
import hmac
import pymssql
import requests
import json
import time
import urllib.parse
from datetime import datetime, date
from decimal import Decimal

# delivery 平台统一异步结果通知接口（与 ES 海牙 030 / 意大利 EPR 同模式）
# 仅 DataSource='api' 的记录在重试达上限时调用；runner 可传参覆盖（受理时 bizParam.callback_url 优先）
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)

# 腾讯云 COS 下载签名凭证：RPA 每个函数都是一次性独立运行的程序，不读环境变量，直接使用固定常量。
# API 流程落库的 upload_file_path* 是 COS 相对路径（私有桶需签名才能下载），
# 领取行时仅对 DataSource='api' 行转为带签名的临时下载 URL；source 流程行（file.usaeu.com）一律不处理。
SAAS_COS_SECRET_ID = 'AKIDjGZcOu9BwSikW1HHchW4PaVAYicGeldF'
SAAS_COS_SECRET_KEY = 'PXWGm24qUszrTpvlLh5E8mAI9m8F7k7E'
SAAS_COS_BUCKET = 'usaeu-1259285998'
SAAS_COS_REGION = 'ap-guangzhou'
# 签名下载 URL 有效期（秒）；COS 服务端仅校验"签名与 q-sign-time 一致 + 当前时间在窗口内"
COS_SIGN_EXPIRATION_SECONDS = 600


def convert_to_serializable(obj):
    """
    将不可序列化的对象转换为可序列化的格式
    """
    if isinstance(obj, (datetime, date)):
        return obj.strftime('%Y-%m-%d %H:%M:%S') if isinstance(obj, datetime) else obj.strftime('%Y-%m-%d')
    elif isinstance(obj, Decimal):
        return float(obj)
    elif isinstance(obj, bytes):
        return obj.decode('utf-8', errors='ignore')
    else:
        return obj


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 sign_cos_url(file_url, expiration_seconds=COS_SIGN_EXPIRATION_SECONDS):
    """
    为 COS 相对路径生成带签名的临时下载 URL。
    算法与 qcloud/cos-sdk-v5 Signature::createPresignedUrl 一致（q-sign-algorithm=sha1，
    q-header-list=host、q-url-param-list=空），已用 SDK 输出逐字符校验过。
    完整 http(s):// URL（source 流程 file.usaeu.com 等）无需签名，原样返回。
    """
    file_url = (file_url or '').strip()
    if not file_url:
        return file_url
    if file_url.lower().startswith(('http://', 'https://')):
        return file_url

    now = int(time.time())
    # q-sign-time 起始提前 60s（与 SDK 相同），结束 = 起始 + 有效期
    sign_time = f'{now - 60};{now - 60 + expiration_seconds}'
    host = f'{SAAS_COS_BUCKET}.cos.{SAAS_COS_REGION}.myqcloud.com'
    # 与 SDK encodeKey 相同：rawurlencode 后 %2F 还原为 /
    path = '/' + urllib.parse.quote(file_url.lstrip('/'), safe='/')

    http_string = 'get\n' + urllib.parse.unquote(path) + '\n\n' + 'host=' + host + '\n'
    string_to_sign = ('sha1\n' + sign_time + '\n'
                      + hashlib.sha1(http_string.encode('utf-8')).hexdigest() + '\n')
    sign_key = hmac.new(SAAS_COS_SECRET_KEY.encode('utf-8'), sign_time.encode('utf-8'),
                        hashlib.sha1).hexdigest()
    signature = hmac.new(sign_key.encode('utf-8'), string_to_sign.encode('utf-8'),
                         hashlib.sha1).hexdigest()

    authorization = (
        'q-sign-algorithm=sha1&q-ak=' + SAAS_COS_SECRET_ID
        + '&q-sign-time=' + sign_time + '&q-key-time=' + sign_time
        + '&q-header-list=host&q-url-param-list=&q-signature=' + signature
    )
    return 'https://' + host + path + '?sign=' + urllib.parse.quote(authorization, safe='')


def resolve_download_urls(row):
    """
    将领取行的 upload_file_path1/2/3 由 COS 相对路径替换为带签名的临时下载 URL。
    仅 DataSource='api' 行处理（API 流程落库的是 COS 相对路径，私有桶需签名才能下载）；
    source 流程行（file.usaeu.com 完整 URL）一律原样保留，不签名。空值或签名异常原样保留，不阻断领取。
    """
    # 显式按 DataSource 分支而非 URL 形态：两流程 bucket 不同，source 行不做任何签名
    if str(get_field(row, 'DataSource') or 'source').lower() != 'api':
        return row
    for field in ('upload_file_path1', 'upload_file_path2', 'upload_file_path3'):
        value = get_field(row, field)
        if not isinstance(value, str) or not value.strip():
            continue
        if value.strip().lower().startswith(('http://', 'https://')):
            continue
        try:
            row[field] = sign_cos_url(value)
        except Exception as e:
            print(f"生成 {field} 签名下载URL失败: {e}")
    return row


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 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"异步结果通知 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"异步结果通知 HTTP {response.status_code} 永久失败不重试: {url}")
                result['error'] = f"HTTP {response.status_code}"
                return result
            result['error'] = f"HTTP {response.status_code}"
            print(f"异步结果通知 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"异步结果通知失败（第 {attempt}/{DELIVERY_NOTIFY_MAX_ATTEMPTS} 次尝试）: {e}")
        if attempt < DELIVERY_NOTIFY_MAX_ATTEMPTS:
            wait_seconds = DELIVERY_NOTIFY_BACKOFF_SECONDS[attempt - 1]
            print(f"{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 send_wechat_work_notification(webhook_url, code, phone):
    """
    发送企业微信机器人消息

    参数:
        webhook_url: 企业微信机器人webhook地址
        code: 流水号
        phone: 负责人手机号
    """
    try:
        # 构造消息内容
        message = {
            "msgtype": "text",
            "text": {
                "content": f"流水号({code}) 请确认当前客户注册是否正确,已重试多次",
                "mentioned_mobile_list": [phone]  # @对应负责人
            }
        }

        # 发送POST请求
        headers = {'Content-Type': 'application/json'}
        response = requests.post(webhook_url, data=json.dumps(message), headers=headers)

        # 检查响应
        if response.status_code == 200:
            result = response.json()
            if result.get('errcode') == 0:
                print(f"企业微信消息发送成功: 流水号 {code}, 通知 {phone}")
                return True
            else:
                print(f"企业微信消息发送失败: {result.get('errmsg')}")
                return False
        else:
            print(f"企业微信消息发送失败: HTTP {response.status_code}")
            return False

    except Exception as e:
        print(f"发送企业微信消息时发生错误: {e}")
        return False


def query_and_update_uk_vat_register(server, database, username, password, db_saas_name, port=1433, delivery_callback_url=None):
    """
    连接 SQL Server 数据库查询 uk_vat_register 表的第一条记录,并将状态更新为 1
    如果失败次数达到10次,则:
      - source 行: 更新 SaaS 数据库中的 VATBusinessRecord 表并发送企业微信通知（维持现状）
      - api 行: 不更新 SaaS 源库（source_record_id 为 'API:{流水号}'），
        发送企业微信通知并通过 delivery 统一异步结果通知链接上报失败

    参数:
        server: 服务器地址 (例如: 'localhost' 或 '192.168.1.100')
        database: 数据库名称
        username: 用户名
        password: 密码
        db_saas_name: SaaS 数据库名称
        port: 端口号 (默认: 1433)
        delivery_callback_url (str, optional): delivery 统一异步结果通知地址，
            缺省用 DEFAULT_RESULT_CALLBACK_URL；仅 DataSource='api' 行调用

    返回:
        dict: 字典格式的数据,键为字段名,值为字段值;如果没有数据则返回 None
    """
    connection = None
    saas_connection = None

    # 企业微信机器人webhook地址
    WECHAT_WEBHOOK_URL = "https://qyapi.weixin.qq.com/cgi-bin/webhook/send?key=b48de199-044d-46b0-b789-95157197e555"

    try:
        # 建立连接到主数据库
        connection = pymssql.connect(
            server=server,
            user=username,
            password=password,
            database=database,
            port=port,
            charset='utf8',
            as_dict=True
        )

        cursor = connection.cursor()

        # 执行查询
        query = """SELECT TOP 1 *
                    FROM uk_vat_register
                    WHERE registration_status IN(0,1,3) AND registration_failure_count < %s
                    ORDER BY id ASC """
        cursor.execute(query, (10,))

        # 获取一条数据
        result = cursor.fetchone()

        if result:
            # 转换不可序列化的对象
            result = {key: convert_to_serializable(value) for key, value in result.items()}

            # 仅 DataSource='api' 行：upload_file_path* 是 COS 相对路径（私有桶需签名），
            # 换取带签名临时下载 URL；source 流程行（file.usaeu.com 完整 URL）一律原样保留
            resolve_download_urls(result)

            # 获取必要字段
            primary_key = get_field(result, 'id') or get_field(result, 'ID') or get_field(result, 'uuid')
            current_failure_count = get_field(result, 'registration_failure_count') or 0
            source_record_id = get_field(result, 'source_record_id')
            code = get_field(result, 'code')  # 流水号
            linked_sales_contact_phone = get_field(result, 'linked_sales_contact_phone')  # 负责人手机号
            data_source = str(get_field(result, 'DataSource') or 'source').lower()  # api / source（旧库无此列视为 source）

            if primary_key:
                # 计算新的失败次数
                new_failure_count = current_failure_count + 1

                # 如果失败次数达到10次,置终态并按来源分流通知
                if new_failure_count >= 10:
                    # 更新状态为 4 并增加失败计数
                    update_query = """UPDATE uk_vat_register
                                    SET registration_status = 4
                                    WHERE id = %s"""
                    cursor.execute(update_query, (primary_key,))
                    connection.commit()
                    print(f"成功更新记录 ID: {primary_key} 的状态为 4")

                    # ---------- source 行: 更新 SaaS 数据库（维持现状） ----------
                    if data_source != 'api' and source_record_id:
                        try:
                            # 建立连接到 SaaS 数据库
                            saas_connection = pymssql.connect(
                                server=server,
                                user=username,
                                password=password,
                                database=db_saas_name,
                                port=port,
                                charset='utf8'
                            )

                            saas_cursor = saas_connection.cursor()

                            # 更新 VATBusinessRecord 表
                            saas_update_query = """UPDATE VATBusinessRecord
                                                   SET PushTaxBureauStatus = -1
                                                   WHERE ID = %s"""
                                                    #    PushTaxBureauErrorMsg = %s

                            # error_msg = '重试次数达到限制,需要检查资料'
                            saas_cursor.execute(saas_update_query, (source_record_id))
                            saas_connection.commit()

                            print(f"失败次数达到10次,已更新 SaaS 数据库中 VATBusinessRecord 记录 ID: {source_record_id}")

                            saas_cursor.close()

                        except pymssql.Error as saas_error:
                            if saas_connection:
                                saas_connection.rollback()
                            print(f"更新 SaaS 数据库时发生错误: {saas_error}")
                        finally:
                            if saas_connection:
                                saas_connection.close()

                    # ---------- api 行: delivery 统一异步结果通知（失败上报，退避重试 + 日志落库） ----------
                    # 契约：非 200 失败通知 data 统一为 null
                    if data_source == 'api':
                        notify_url = resolve_notification_url(result, delivery_callback_url)
                        payload = {
                            'code': 500,
                            'msg': 'failed',
                            'ProcessMode': 'async',
                            'data': None,
                            'bizParam': parse_biz_param(result),
                        }
                        notify_result = send_async_result_notification(notify_url, payload)
                        save_delivery_notify_log(primary_key, notify_url, payload, notify_result,
                                                 server, port, username, password, database)

                    # 发送企业微信通知（两类行均发送）
                    if code and linked_sales_contact_phone:
                        send_wechat_work_notification(
                            webhook_url=WECHAT_WEBHOOK_URL,
                            code=code,
                            phone=linked_sales_contact_phone
                        )
                    else:
                        print(f"警告: 缺少流水号或负责人手机号,无法发送企业微信通知")
                        if not code:
                            print("  - 缺少流水号 (code)")
                        if not linked_sales_contact_phone:
                            print("  - 缺少负责人手机号 (linked_sales_contact_phone)")

                # 更新状态为 1 并增加失败计数
                update_query = """UPDATE uk_vat_register
                                  SET registration_status = 1,
                                      registration_failure_count = ISNULL(registration_failure_count, 0) + 1
                                  WHERE id = %s"""
                cursor.execute(update_query, (primary_key,))
                connection.commit()
                print(f"成功更新记录 ID: {primary_key} 的状态为 1, 失败次数: {new_failure_count}")
            else:
                print("警告: 无法确定主键字段,状态未更新")

        # 关闭连接
        cursor.close()
        connection.close()

        return result

    except pymssql.Error as e:
        if connection:
            connection.rollback()
        print(f"数据库错误: {e}")
        return None
    except Exception as e:
        if connection:
            connection.rollback()
        print(f"发生错误: {e}")
        return None
    finally:
        if connection:
            connection.close()
        if saas_connection:
            saas_connection.close()
