#!/usr/bin/env python3
# -*- coding: utf-8 -*-
"""
顺丰物流轨迹数据管道
用法: python3 shipment_pipeline.py          # 从MCP拉数据+顺丰API
      python3 shipment_pipeline.py --mock   # 用mock数据测试

输出: 数据/shipment_status.json
"""
import json, hashlib, base64, time, uuid, urllib.request, urllib.parse, os, sys

ROOT = os.path.dirname(os.path.dirname(os.path.abspath(__file__)))
DATA_DIR = os.path.join(ROOT, '数据')
OUTPUT = os.path.join(DATA_DIR, 'shipment_status.json')

# 顺丰凭证
SF_PARTNER_ID = "BRSWK5CQ08HV"
SF_SECRET = "V7JDT0azGX93pnIXRyZe9zvlzY6fk1rC"
SF_TOKEN_URL = "https://sfapi.sf-express.com/oauth2/accessToken"
SF_API_URL = "https://bspgw.sf-express.com/std/service"

def sf_get_token():
    """获取顺丰 accessToken"""
    body = urllib.parse.urlencode({
        "partnerID": SF_PARTNER_ID,
        "secret": SF_SECRET,
        "grantType": "password"
    }).encode("utf-8")
    req = urllib.request.Request(SF_TOKEN_URL, data=body,
        headers={"Content-Type": "application/x-www-form-urlencoded"})
    with urllib.request.urlopen(req, timeout=30) as resp:
        result = json.loads(resp.read().decode("utf-8"))
    if result.get("apiResultCode") != "A1000":
        raise Exception(f"SF token failed: {result}")
    return result["accessToken"]

def sf_query_routes(tracking_numbers, access_token):
    """批量查询物流轨迹（最多10个单号）"""
    msg_data = json.dumps({
        "language": "zh-CN",
        "trackingType": 1,
        "trackingNumber": tracking_numbers,
        "methodType": 1
    }, ensure_ascii=False)

    timestamp = str(int(time.time()))
    request_id = str(uuid.uuid4()).replace("-", "")

    body = urllib.parse.urlencode({
        "partnerID": SF_PARTNER_ID,
        "requestID": request_id,
        "serviceCode": "EXP_RECE_SEARCH_ROUTES",
        "timestamp": timestamp,
        "accessToken": access_token,
        "msgData": msg_data
    }).encode("utf-8")

    req = urllib.request.Request(SF_API_URL, data=body,
        headers={"Content-Type": "application/x-www-form-urlencoded;charset=utf-8"})
    with urllib.request.urlopen(req, timeout=30) as resp:
        result = json.loads(resp.read().decode("utf-8"))

    if result.get("apiResultCode") != "A1000":
        return {"error": result.get("apiErrorMsg", "unknown"), "code": result.get("apiResultCode")}

    inner = json.loads(result.get("apiResultData", "{}"))
    return inner

def parse_routes(routes_data):
    """解析路由信息，提取关键状态"""
    if not routes_data.get("success"):
        return {"routes": [], "status": "error", "error": routes_data.get("errorMsg", "")}

    route_resps = routes_data.get("msgData", {}).get("routeResps", [])
    results = {}
    for resp in route_resps:
        mail_no = resp.get("mailNo", "")
        routes = resp.get("routes", [])
        reason = resp.get("reasonRemark", "")

        # 判断状态
        status = "unknown"
        last_route = None
        if routes:
            last_route = routes[-1]
            remark = last_route.get("remark", "")
            if "签收" in remark:
                status = "delivered"
            elif "派送" in remark or last_route.get("opCode") == "30":
                status = "delivering"
            elif "揽收" in remark or last_route.get("opCode") == "54":
                status = "picked_up"
            else:
                status = "in_transit"
        elif reason:
            status = "no_data"  # 可能是旧单号或权限问题

        results[mail_no] = {
            "mailNo": mail_no,
            "status": status,
            "statusLabel": {
                "delivered": "已签收",
                "delivering": "派送中",
                "picked_up": "已揽收",
                "in_transit": "运输中",
                "no_data": "暂无数据",
                "unknown": "未知",
                "error": "查询异常"
            }.get(status, "未知"),
            "lastRoute": last_route,
            "routes": routes,
            "reason": reason,
            "routeCount": len(routes)
        }
    return results

# 离职人员订单归属映射：离职销售的历史订单划转到接替人
RESIGN_REASSIGN = {
    "张立娅": ("刘子研", "liuziyan@biori.com"),
    "黄明月": ("刘子研", "liuziyan@biori.com"),
}

def reassign_creator(o):
    """把离职人员的创建者/邮箱映射到接替人"""
    creator = o.get("创建者", "") or ""
    email = o.get("创建者邮箱", "") or ""
    if creator in RESIGN_REASSIGN:
        o["创建者"], o["创建者邮箱"] = RESIGN_REASSIGN[creator]
    elif email.lower() in ("zhangliya@biori.com", "huangmingyue@biori.com"):
        o["创建者"], o["创建者邮箱"] = RESIGN_REASSIGN["张立娅"] if email.lower() == "zhangliya@biori.com" else RESIGN_REASSIGN["黄明月"]
    return o


def process_orders(orders_data, token):
    """处理所有订单（含无物流单号的），批量查轨迹"""
    # 收集全部订单+有物流单号的去重
    all_orders = []
    seen_tn = set()
    tracking_orders = []  # 只用于SF批量查询
    for o in orders_data:
        o = reassign_creator(o)  # 离职人员划转
        tn = (o.get("物流单号") or "").strip()
        all_orders.append(o)
        if tn and tn not in seen_tn:
            seen_tn.add(tn)
            tracking_orders.append(o)

    print(f"全量: {len(all_orders)}条, 有物流单号: {len(tracking_orders)}条")

    # 批量查询SF（仅对有物流单号的）
    all_results = {}
    batch_size = 10
    for i in range(0, len(tracking_orders), batch_size):
        batch = tracking_orders[i:i+batch_size]
        tracking_nums = [o["物流单号"].strip() for o in batch]
        try:
            result = sf_query_routes(tracking_nums, token)
            if isinstance(result, dict) and "error" not in result:
                routes = parse_routes(result)
                all_results.update(routes)
            else:
                for tn in tracking_nums:
                    all_results[tn] = {"mailNo": tn, "status": "error", "error": str(result)}
        except Exception as e:
            for tn in tracking_nums:
                all_results[tn] = {"mailNo": tn, "status": "error", "error": str(e)}
        if i + batch_size < len(tracking_orders):
            time.sleep(0.5)

    # 输出全部订单
    output = []
    for o in all_orders:
        tn = (o.get("物流单号") or "").strip()
        route_info = all_results.get(tn, {}) if tn else {}
        output.append({
            "orderId": o.get("单据编号", ""),
            "customer": o.get("购货单位#", ""),
            "mcpCustomer": o.get("销售订单关联客户", ""),
            "dept": o.get("销售部门", ""),
            "createDate": o.get("创建时间", ""),
            "recipient": o.get("收货人姓名", ""),
            "phone": o.get("收货人联系电话#", ""),
            "address": o.get("收货地址", ""),
            "trackingNo": tn,
            "amount": o.get("金额", ""),
            "taxTotal": o.get("金额", ""),
            "creator": o.get("创建者", ""),
            "creatorEmail": o.get("创建者邮箱", ""),
            "contractNo": o.get("合同编号#", ""),
            "shipment": route_info
        })
    return output

def main():
    if "--mock" in sys.argv:
        # Mock数据测试
        mock = [{
            "单据编号": "ZHBR260724027",
            "购货单位#": "圣湘生物",
            "销售订单关联客户": "圣湘生物科技股份有限公司",
            "物料名称": "2×FastAmpli Premix-UNG Ⅳ",
            "货号#": "M2181-4-NA-095203",
            "出货批号#": "20260417",
            "销售部门": "诊断原料大客户销售部",
            "单据类型": "标准销售订单",
            "状态": "未开始",
            "创建时间": "2026-07-24",
            "收货人姓名": "蔡琦",
            "收货人联系电话#": "17610556866",
            "收货地址": "湖南省长沙市岳麓区",
            "物流单号": "SF5112860018518"
        }]
        token = sf_get_token()
        output = process_orders(mock, token)
        print(f"Mock结果: {json.dumps(output, ensure_ascii=False, indent=2)[:500]}")
        return

    # 实际运行: 从MCP输出文件读取
    input_file = os.path.join(DATA_DIR, 'mcp_orders_input.json')
    if not os.path.exists(input_file):
        print(f"请先运行MCP拉取数据: 将订单JSON放入 {input_file}")
        print("或者使用 --mock 测试")
        sys.exit(1)

    with open(input_file, encoding='utf-8') as f:
        orders_data = json.load(f)

    token = sf_get_token()
    output = process_orders(orders_data, token)

    os.makedirs(DATA_DIR, exist_ok=True)
    with open(OUTPUT, 'w', encoding='utf-8') as f:
        json.dump({
            "updated_at": time.strftime("%Y-%m-%d %H:%M:%S"),
            "total": len(output),
            "orders": output
        }, f, ensure_ascii=False, indent=2)

    print(f"✅ 写入 {OUTPUT} ({len(output)} 条)")

if __name__ == '__main__':
    main()
