import json
import time
import urllib.parse

import pymssql
from datetime import datetime
from qcloud_cos import CosConfig, CosS3Client
import logging
import sys
import os
from urllib.parse import quote
from contextlib import closing
import requests

# 配置日志
logging.basicConfig(level=logging.INFO, stream=sys.stdout)

# ==================== 双系统对接常量 ====================
# 老系统(DataSource='source')文件解析接口（原值）
FILE_ANALYSIS_API = "https://sassservice.usaeu.com/Tool/ReadOtherCalculateFilePdf"
# 新系统(DataSource='api')文件解析接口：传参/返回与老接口完全一致；
# 新系统接口提供后替换此值即可（当前暂用老接口地址）。
NEW_FILE_ANALYSIS_API = FILE_ANALYSIS_API

# delivery 平台统一异步结果通知接口（§13 delivery/rpa/callback，与 ES 海牙/意大利 EPR 同模式）
# 仅 DataSource='api' 行保存成功后调用；地址来源：runner 传参 > 默认
DEFAULT_RESULT_CALLBACK_URL = "https://test-cloud.usaeu.com/prod-api/delivery/rpa/callback"

# 异步结果通知投递重试：首次 + 2 次重试，第 2、3 次尝试前分别等待 2s / 6s（覆盖瞬时网络/服务抖动）
DELIVERY_NOTIFY_MAX_ATTEMPTS = 3
DELIVERY_NOTIFY_BACKOFF_SECONDS = (2, 6)

# 通知 files[].type 文件类型枚举（CDS 专用，勿自定义）
CDS_FILE_TYPE_PVA = 'CDS_PVA_FILE'
CDS_FILE_TYPE_C79 = 'CDS_C79_FILE'

class TencentCOSUploader:
    """腾讯云COS上传器"""
    
    def __init__(self, secret_id, secret_key, region, bucket_name, token=None, scheme='https'):
        """
        初始化COS客户端
        
        Args:
            secret_id (str): 腾讯云SecretId
            secret_key (str): 腾讯云SecretKey  
            region (str): 地域,如'ap-guangzhou'
            bucket_name (str): 存储桶名称
            token (str, optional): 临时密钥token
            scheme (str): 协议,默认'https'
        """
        self.secret_id = secret_id
        self.secret_key = secret_key
        self.region = region
        self.bucket_name = bucket_name
        
        # 初始化COS客户端
        config = CosConfig(
            Region=region, 
            SecretId=secret_id, 
            SecretKey=secret_key, 
            Token=token, 
            Scheme=scheme
        )
        self.client = CosS3Client(config)
    
    def check_file_exists(self, cos_object_key):
        """
        检查文件在COS中是否已存在
        
        Args:
            cos_object_key (str): COS对象键
            
        Returns:
            bool: 文件是否存在
        """
        try:
            self.client.head_object(Bucket=self.bucket_name, Key=cos_object_key)
            return True
        except Exception:
            return False
    
    def upload_file(self, local_file_path, cos_object_key=None):
        """
        上传文件到腾讯云COS,如果文件已存在则跳过上传
        
        Args:
            local_file_path (str): 本地文件路径
            cos_object_key (str): COS对象键,如果不指定则使用文件名
            
        Returns:
            dict: 包含上传结果的字典
        """
        # 检查本地文件是否存在
        if not os.path.exists(local_file_path):
            return {
                'success': False,
                'error': f'本地文件不存在: {local_file_path}',
                'url': None,
                'skipped': False
            }
        
        # 如果没有指定COS对象键,使用文件名
        if cos_object_key is None:
            cos_object_key = f"uploads/{datetime.now().strftime('%Y%m%d')}/{os.path.basename(local_file_path)}"
        
        # 构建访问URL
        encoded_key = quote(cos_object_key, safe='/ ')
        file_url = f'https://{self.bucket_name}.cos.{self.region}.myqcloud.com/{encoded_key}'
        
        # 检查文件是否已存在
        if self.check_file_exists(cos_object_key):
            print(f"📁 文件已存在,跳过上传: {cos_object_key}")
            return {
                'success': True,
                'error': None,
                'url': file_url,
                'etag': None,
                'cos_key': cos_object_key,
                'skipped': True
            }
        
        try:
            print(f"🔄 开始上传文件: {local_file_path} -> {cos_object_key}")
            
            # 执行上传
            response = self.client.upload_file(
                Bucket=self.bucket_name,
                LocalFilePath=local_file_path,
                Key=cos_object_key,
                PartSize=1,
                MAXThread=5,
                EnableMD5=False
            )
            
            if response is None:
                return {
                    'success': False,
                    'error': '上传响应为空',
                    'url': None,
                    'skipped': False
                }
            
            # 获取ETag
            etag = response.get('ETag')
            
            print(f"✅ 文件上传成功: {local_file_path} -> {file_url}")
            
            return {
                'success': True,
                'error': None,
                'url': file_url,
                'etag': etag,
                'cos_key': cos_object_key,
                'skipped': False
            }
            
        except Exception as e:
            error_msg = f"上传失败: {str(e)}"
            print(f"❌ {error_msg}")
            return {
                'success': False,
                'error': error_msg,
                'url': None,
                'skipped': False
            }


def call_file_analysis_api(file_url, file_name, timeout=60, api_url=FILE_ANALYSIS_API):
    """
    调用文件解析API获取文件金额

    Args:
        file_url (str): 文件完整URL
        file_name (str): 文件名
        timeout (int): 请求超时时间(秒)
        api_url (str): 解析接口地址（按账号归属选择：source=FILE_ANALYSIS_API / api=NEW_FILE_ANALYSIS_API）

    Returns:
        dict: 包含API调用结果的字典
    """

    
    try:
        print(f"🔄 调用文件解析API: {file_name}")
        
        # 构建请求参数
        params = {
            'FileUrl': file_url,
            'FileName': file_name
        }
        
        # 构建完整URL
        from urllib.parse import urlencode
        full_url = f"{api_url}?{urlencode(params)}"
        print(f"📝 完整请求URL: {full_url}")
        
        # 发送GET请求
        response = requests.get(full_url, timeout=timeout)
        
        # 打印响应状态码和内容
        print(f"📝 响应状态码: {response.status_code}")
        print(f"📝 响应内容: {response.text[:500]}")  # 只打印前500字符
        
        # 检查HTTP状态码
        if response.status_code != 200:
            error_msg = f"HTTP错误: {response.status_code}, 内容: {response.text}"
            print(f"❌ {error_msg}")
            return {
                'success': False,
                'total_postponed': 0,
                'error': error_msg,
                'response': None
            }
        
        # 检查响应内容是否为空
        if not response.text or response.text.strip() == '':
            error_msg = "API返回空响应"
            print(f"❌ {error_msg}")
            return {
                'success': False,
                'total_postponed': 0,
                'error': error_msg,
                'response': None
            }
        
        # 尝试解析JSON
        try:
            result = response.json()
        except ValueError as json_error:
            error_msg = f"JSON解析失败: {str(json_error)}, 响应内容: {response.text[:200]}"
            print(f"❌ {error_msg}")
            return {
                'success': False,
                'total_postponed': 0,
                'error': error_msg,
                'response': None
            }
        
        # 检查API返回的业务状态
        if result.get('succeeded') == True:
            data = result.get('data', {})
            total_postponed = data.get('TotalPostponed', 0) if data else 0
            print(f"✅ API调用成功,获取到金额: {total_postponed}")
            return {
                'success': True,
                'total_postponed': total_postponed,
                'error': None,
                'response': result
            }
        else:
            error_message = result.get('message', '未知错误')
            print(f"❌ API调用失败: {error_message}")
            return {
                'success': False,
                'total_postponed': 0,
                'error': error_message,
                'response': result
            }
            
    except requests.exceptions.Timeout:
        error_msg = f"API请求超时(超过{timeout}秒)"
        print(f"❌ {error_msg}")
        return {
            'success': False,
            'total_postponed': 0,
            'error': error_msg,
            'response': None
        }
    except requests.exceptions.ConnectionError as e:
        error_msg = f"连接错误: {str(e)}"
        print(f"❌ {error_msg}")
        return {
            'success': False,
            'total_postponed': 0,
            'error': error_msg,
            'response': None
        }
    except requests.exceptions.RequestException as e:
        error_msg = f"API请求异常: {str(e)}"
        print(f"❌ {error_msg}")
        return {
            'success': False,
            'total_postponed': 0,
            'error': error_msg,
            'response': None
        }
    except Exception as e:
        error_msg = f"处理API响应时出错: {str(e)}"
        print(f"❌ {error_msg}")
        return {
            'success': False,
            'total_postponed': 0,
            'error': error_msg,
            'response': None
        }


def ensure_total_postponed_column(cursor):
    """
    确保cds_download_log表中存在total_postponed字段
    如果需要手动添加,请执行以下SQL:
    
    ALTER TABLE cds_download_log ADD total_postponed DECIMAL(18,2) DEFAULT 0;
    
    EXEC sp_addextendedproperty 
        @name = N'MS_Description', 
        @value = N'文件金额', 
        @level0type = N'SCHEMA', @level0name = N'dbo', 
        @level1type = N'TABLE', @level1name = N'cds_download_log', 
        @level2type = N'COLUMN', @level2name = N'total_postponed';
    
    Args:
        cursor: 数据库游标
    """
    pass  # 不自动添加字段,由用户手动执行SQL
            

def update_cds_record(
                     account_id,
                     db_host='localhost',
                     db_database='vat',
                     db_user='sa',
                     db_password='password',
                     db_port=1433,
                     db_charset='utf8',
                     cos_secret_id=None,
                     cos_secret_key=None,
                     cos_region=None,
                     cos_bucket_name=None,
                     oss_key_path=None,
                     account_alias=None,
                     file_path=None,
                     file_date=None,
                     file_type=None,
                     delivery_callback_url=None):
    """
    更新CDS下载配置记录,支持自动上传文件到腾讯云COS并记录到日志表

    Args:
        account_id: 账号ID
        account_alias: 账号别名标识,对应对接系统中本条数据唯一标识
        file_path: 文件路径(本地路径或URL)
        file_date: 文件对应时间(6位字符,如202409)
        file_type: 文件类型(pva 或 c79)
        db_host: SQL Server服务器地址
        db_database: 数据库名称
        db_user: 数据库用户名
        db_password: 数据库密码
        db_port: 数据库端口,默认1433
        db_charset: 字符编码,默认'utf8'
        delivery_callback_url: delivery 统一异步结果通知地址(仅 DataSource='api' 行使用),
            runner 传参,缺省用 DEFAULT_RESULT_CALLBACK_URL
        其他参数: COS配置参数
    """
    connection = None
    cos_uploader = None
    upload_result = None
    api_result = None
    total_postponed = 0
    
    # 确保 account_id 是字符串(因为数据库是 VARCHAR)
    if account_id is not None:
        account_id = str(account_id)
    
    try:
        # 初始化COS上传器 (如果提供了COS配置)
        if all([cos_secret_id, cos_secret_key, cos_region, cos_bucket_name]):
            cos_uploader = TencentCOSUploader(
                secret_id=cos_secret_id,
                secret_key=cos_secret_key,
                region=cos_region,
                bucket_name=cos_bucket_name
            )
            print("✅ COS上传器初始化成功")
        else:
            print("⚠️  未提供完整的COS配置,将跳过文件上传")
        
        # 处理文件上传(如果提供了文件路径)
        final_file_path = file_path
        if file_path is not None and file_type is not None:
            # 检查本地文件是否存在
            if os.path.exists(file_path):
                # 如果配置了COS,则上传文件
                if cos_uploader:
                    print(f"🔄 处理{file_type.upper()}文件: {file_path}")
                    
                    # 生成COS对象键: oss_key_path/文件类型/账号ID/文件名
                    filename = os.path.basename(file_path)
                    cos_object_key = f"{oss_key_path}/{file_type.upper()}/{account_id}/{filename}"
                    
                    upload_result = cos_uploader.upload_file(file_path, cos_object_key)
                    
                    if upload_result['success']:
                        # 如果上传成功或已存在,用COS的URL替换本地路径
                        final_file_path = upload_result['url']
                        if upload_result.get('skipped', False):
                            print(f"📁 文件已存在,使用现有URL: {upload_result['url']}")
                        else:
                            print(f"✅ 文件上传成功: {file_path} -> {upload_result['url']}")
                    else:
                        # 如果上传失败,仍使用原路径
                        print(f"❌ 文件上传失败,使用原路径: {file_path}")
                else:
                    print(f"📁 文件路径: {file_path} (未上传到COS)")
            else:
                print(f"⚠️  警告: 文件不存在: {file_path}")
               
        # 连接SQL Server数据库
        connection = pymssql.connect(
            server=db_host,
            port=db_port,
            user=db_user,
            password=db_password,
            database=db_database,
            charset=db_charset,
            as_dict=True
        )
        
        cursor = connection.cursor()
        
        # 确保total_postponed字段存在
        ensure_total_postponed_column(cursor)

        # 0. 查询账号归属（DataSource/biz_param）：决定解析接口选择与是否发异步结果通知
        #    表缺列/查询失败/无行一律按 source(老系统) 处理,保持原行为
        row_source = query_account_source(cursor, account_id)
        data_source = str(get_field(row_source, 'DataSource') or 'source').lower()
        try:
            biz_param = json.loads(get_field(row_source, 'biz_param')) if get_field(row_source, 'biz_param') else None
        except (json.JSONDecodeError, TypeError):
            biz_param = None

        # 1. 插入文件日志记录到 cds_download_log(如果提供了文件信息)
        log_inserted = False
        if all([file_path, file_date, file_type]):
            # 检查是否已存在相同的记录（按三列唯一键 account_id+file_type+file_date 判重，
            # 与 UQ_cds_download_log 一致；account_alias 变更不影响判重，避免迁移后撞唯一索引）
            check_sql = """
            SELECT id FROM cds_download_log
            WHERE account_id = %s AND file_type = %s AND file_date = %s
            """
            cursor.execute(check_sql, (account_id, file_type, file_date))
            existing_record = cursor.fetchone()
            
            if existing_record:
                # 不需要处理已存在
                log_inserted = False
                print(f"✅ 不需要处理已存在cds_download_log记录ID: {existing_record['id']}")
            else:
                # 确认需要插入新记录,此时调用API获取文件金额（按账号归属选择解析接口：api行走新接口）
                if upload_result and upload_result['success']:
                    filename = os.path.basename(file_path)
                    api_url = NEW_FILE_ANALYSIS_API if data_source == 'api' else FILE_ANALYSIS_API
                    print(f"🔄 确认插入新记录,调用文件解析API获取金额... (接口: {api_url})")
                    api_result = call_file_analysis_api(final_file_path, filename, api_url=api_url)
                    if api_result['success']:
                        total_postponed = api_result['total_postponed']
                        print(f"✅ 成功获取文件金额: {total_postponed}")
                    else:
                        print(f"⚠️  API调用失败,金额将设置为0: {api_result['error']}")

                # 插入新记录(包含total_postponed字段)
                insert_log_sql = """
                INSERT INTO cds_download_log (account_id, file_type, file_path, file_date, account_alias, total_postponed)
                VALUES (%s, %s, %s, %s, %s, %s)
                """
                cursor.execute(insert_log_sql, (account_id, file_type, final_file_path, file_date, account_alias, total_postponed))
                log_inserted = True

                # 取新插入的日志行ID(通知 task_id) —— SCOPE_IDENTITY 与上一条 INSERT 同会话同作用域
                cursor.execute("SELECT SCOPE_IDENTITY() AS id")
                log_id_row = cursor.fetchone()
                log_id = log_id_row['id'] if log_id_row else None
                print(f"✅ 插入新的cds_download_log记录(id={log_id}),account_id: {account_id}, account_alias: {account_alias}, file_type: {file_type}, file_date: {file_date}, total_postponed: {total_postponed}")

        # 提交事务
        if log_inserted:
            connection.commit()

            # 2. api 行保存成功后发统一异步结果通知(old source 行不通知,老系统自助查询)
            notify_result = None
            if data_source == 'api' and log_id is not None:
                notify_url = resolve_notification_url(delivery_callback_url)
                filename = os.path.basename(file_path) if file_path else ''
                notify_payload = build_notification_payload(
                    log_id, account_id, account_alias, file_date, file_type,
                    total_postponed, final_file_path, filename, biz_param
                )
                print(f"🔄 发送异步结果通知(DataSource=api): {notify_url}")
                notify_result = send_async_result_notification(notify_url, notify_payload)
                save_cds_notify_log(cursor, log_id, account_id, account_alias, notify_url,
                                    notify_payload, notify_result)
                connection.commit()

            # 构建成功消息
            message_parts = []
            if log_inserted:
                message_parts.append(f"记录文件日志: {file_type}_{file_date}")
            if upload_result and upload_result['success']:
                if upload_result.get('skipped', False):
                    message_parts.append("文件已存在COS中")
                else:
                    message_parts.append("文件上传到COS成功")
            if api_result and api_result['success']:
                message_parts.append(f"文件金额: {total_postponed}")
            if notify_result is not None:
                message_parts.append(
                    "结果通知送达" if notify_result['success']
                    else f"结果通知失败: {notify_result['error']}"
                )

            success_message = f"操作完成 - {', '.join(message_parts)}"
            print(success_message)

            return {
                'success': True,
                'message': success_message,
                'log_inserted': log_inserted,
                'log_id': log_id,
                'upload_result': upload_result,
                'api_result': api_result,
                'notify_result': notify_result,
                'total_postponed': total_postponed,
                'final_file_path': final_file_path,
                'account_id': account_id
            }
        else:
            return {
                'success': False,
                'message': '没有需要插入的数据',
                'log_inserted': False,
                'account_id': account_id
            }
                
    except pymssql.Error as e:
        error_msg = f"数据库操作出错: {e}"
        print(f"❌ {error_msg}")
        if connection:
            connection.rollback()
        return {
            'success': False,
            'message': error_msg,
            'log_inserted': False,
            'upload_result': upload_result,
            'api_result': api_result,
            'account_id': account_id
        }
    except Exception as e:
        error_msg = f"处理过程中出错: {e}"
        print(f"❌ {error_msg}")
        if connection:
            connection.rollback()
        return {
            'success': False,
            'message': error_msg,
            'log_inserted': False,
            'upload_result': upload_result,
            'api_result': api_result,
            'account_id': account_id
        }
    finally:
        if connection:
            connection.close()


# ==================== 双系统通知辅助函数（参照 uk_vat_reg rpa_save_step1.py） ====================

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 query_account_source(cursor, account_id):
    """
    查询账号行归属（DataSource / biz_param）。
    表缺少新列（旧库）或查询失败时返回空 dict，调用方按 source(老系统) 行处理。
    """
    if not account_id:
        return {}
    try:
        cursor.execute(
            "SELECT DataSource, biz_param FROM cds_download WHERE account_id = %s",
            (str(account_id),)
        )
        return cursor.fetchone() or {}
    except Exception as e:
        print(f"查询账号归属失败（按老系统处理）: {e}")
        return {}


def resolve_notification_url(delivery_callback_url, default_url=DEFAULT_RESULT_CALLBACK_URL):
    """通知地址来源：runner 传参 > 默认 delivery 地址"""
    if delivery_callback_url:
        return str(delivery_callback_url)
    return default_url


def file_basename_from_url(file_url):
    """从文件 URL / 相对路径提取文件名（去查询串、锚点与 URL 编码）"""
    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_notification_payload(log_id, account_id, account_alias, file_date, file_type,
                               total_postponed, file_url, file_name, biz_param):
    """
    构造 §13 统一异步结果通知（delivery/rpa/callback）请求体。
    data 含 CDS 业务字段(task_id=cds_download_log.id / account_id / account_alias / file_date /
    file_type / total_postponed)与 files([{url,name,type}], type=CDS_PVA_FILE/CDS_C79_FILE)。
    """
    file_type = str(file_type or '').upper()
    if file_type == 'PVA':
        file_type_enum = CDS_FILE_TYPE_PVA
    else:
        file_type_enum = CDS_FILE_TYPE_C79

    files = []
    if file_url:
        files.append({
            'url': str(file_url).strip(),
            'name': (file_name or '').strip() or file_basename_from_url(file_url),
            'type': file_type_enum,
        })

    return {
        'code': 200,
        'msg': 'success',
        'ProcessMode': 'async',
        'data': {
            'task_id': log_id,
            'account_id': account_id,
            'account_alias': account_alias or '',
            'file_date': file_date,
            'file_type': file_type,
            'total_postponed': total_postponed,
            'files': files,
        },
        'bizParam': biz_param if isinstance(biz_param, dict) else {},
    }


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_cds_notify_log(cursor, log_id, account_id, account_alias, notify_url,
                        payload, notify_result):
    """
    将异步结果通知的请求参数与投递结果写入 cds_download_notify_log（一行一次通知）。
    表不存在（旧库未迁移）/写入失败等错误仅告警，不影响主流程。调用方负责 commit。
    """
    try:
        cursor.execute(
            "INSERT INTO cds_download_notify_log "
            "(log_id, account_id, account_alias, notify_url, request_payload, response_text, "
            " http_code, attempts, status) "
            "VALUES (%s, %s, %s, %s, %s, %s, %s, %s, %s)",
            (
                log_id,
                account_id,
                account_alias,
                notify_url,
                json.dumps(payload, ensure_ascii=False),
                notify_result.get('response_text') or notify_result.get('error'),
                notify_result.get('http_code'),
                notify_result.get('attempts'),
                'success' if notify_result['success'] else 'failed',
            )
        )
        print(f"异步结果通知日志已写入 cds_download_notify_log(log_id={log_id})")
    except Exception as e:
        print(f"记录异步结果通知日志失败（不影响主流程）: {e}")





