import pymssql
import requests
import json
from datetime import datetime, date
from decimal import Decimal

# delivery 平台统一异步结果通知接口（与 ES 海牙 030 / 意大利 EPR 同模式）
# 仅 DataSource='api' 的记录在重试达上限时调用；runner 可传参覆盖（受理时 bizParam.callback_url 优先）
DEFAULT_RESULT_CALLBACK_URL = "http://192.168.1.211:8080/delivery/rpa/callback"


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 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 平台统一异步结果通知
    best-effort：2xx 视为送达，单次不重试，失败不影响主流程
    """
    try:
        headers = {'Content-Type': 'application/json; charset=utf-8'}
        response = requests.post(
            url,
            data=json.dumps(payload, ensure_ascii=False),
            headers=headers,
            timeout=10
        )
        print(f"异步结果通知 HTTP {response.status_code}: {url}")
        return 200 <= response.status_code < 300
    except Exception as e:
        print(f"异步结果通知失败（不影响主流程）: {e}")
        return False


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

            # 获取必要字段
            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 统一异步结果通知（失败上报） ----------
                    if data_source == 'api':
                        notify_url = resolve_notification_url(result, delivery_callback_url)
                        payload = {
                            'code': 500,
                            'msg': 'failed',
                            'ProcessMode': 'async',
                            'data': {
                                'record_id': primary_key,
                                'status': 'failed',
                                'error': '重试次数达到限制,需要检查资料',
                            },
                            'bizParam': parse_biz_param(result),
                        }
                        send_async_result_notification(notify_url, payload)

                    # 发送企业微信通知（两类行均发送）
                    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()
