From c7f8b4c7021249260cb24289f369328911cdb545 Mon Sep 17 00:00:00 2001 From: break Date: Wed, 22 Jul 2026 19:09:28 +0900 Subject: [PATCH] =?UTF-8?q?=E6=89=B9=E9=87=8F=E6=8D=A2=E8=B4=A7=E8=84=9A?= =?UTF-8?q?=E6=9C=AC?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- docs/批量换货脚本/功能总结.md | 49 ++ scripts/batch_exchange/README.md | 114 ++++ scripts/batch_exchange/batch_exchange.py | 613 ++++++++++++++++++ scripts/batch_exchange/exchanges.example.csv | 38 ++ scripts/batch_exchange/test_batch_exchange.py | 172 +++++ 5 files changed, 986 insertions(+) create mode 100644 docs/批量换货脚本/功能总结.md create mode 100644 scripts/batch_exchange/README.md create mode 100755 scripts/batch_exchange/batch_exchange.py create mode 100644 scripts/batch_exchange/exchanges.example.csv create mode 100644 scripts/batch_exchange/test_batch_exchange.py diff --git a/docs/批量换货脚本/功能总结.md b/docs/批量换货脚本/功能总结.md new file mode 100644 index 0000000..3c46d83 --- /dev/null +++ b/docs/批量换货脚本/功能总结.md @@ -0,0 +1,49 @@ +# 批量换货脚本功能总结 + +## 功能说明 + +新增 Python 运维脚本 `scripts/batch_exchange/batch_exchange.py`,用于读取旧资产标识、新资产标识两列 CSV,并逐行调用后台 `POST /api/admin/exchanges` 接口。 + +脚本复用现有换货业务事务,不直接操作数据库,不新增换货接口,也不改变现有 Handler、Service 或数据模型。 + +## 固定业务约束 + +- 换货流程固定为 `flow_type=direct`,不支持物流换货。 +- 数据迁移固定为 `migrate_data=true`,不提供关闭选项。 +- 每组资产调用一次创建换货接口;现有接口会在同一事务内创建并立即完成直接换货。 +- 资产类型通过 `--asset-type` 按批次指定为 `iot_card` 或 `device`,新旧资产类型仍由接口校验一致性。 +- 换货原因默认使用 `批量直接换货`,支持按批次覆盖原因和备注。 + +## 数据迁移范围 + +数据迁移完全复用现有换货服务,包含: + +- 资产钱包余额。 +- 套餐使用记录。 +- 累计充值字段。 +- 资产标签。 + +此外,直接换货原有流程仍会处理客户绑定切换以及新旧资产状态更新。 + +## 安全控制 + +- 默认仅预演,显式增加 `--execute` 后才会真实换货。 +- 请求前校验 CSV 必须正好两列且新旧标识非空。 +- 拦截同一行新旧资产相同、重复旧资产、重复新资产。 +- 拦截同一资产在批次内同时作为旧资产和新资产,避免顺序执行改变后续行的资产状态。 +- 接口返回成功时仍要求响应中的 `migration_completed=true`,否则结果记为失败并要求人工核对。 +- 每行结果立即写入结果 CSV,中断后保留已处理记录。 +- POST 请求不自动重试,避免接口已成功但客户端未收到响应时产生误操作。 +- Token 和密码支持环境变量传入,不写入结果文件。 + +## 输出 + +结果 CSV 包含: + +- CSV 原始行号、新旧资产标识。 +- 成功或失败状态。 +- HTTP 状态码、业务错误码和接口消息。 +- 换货单 ID、换货单号。 +- 迁移完成状态和迁移余额。 + +全部成功时退出码为 `0`;存在失败或认证失效导致中途停止时退出码为 `2`;参数、CSV 或登录错误时退出码为 `1`。 diff --git a/scripts/batch_exchange/README.md b/scripts/batch_exchange/README.md new file mode 100644 index 0000000..014e28e --- /dev/null +++ b/scripts/batch_exchange/README.md @@ -0,0 +1,114 @@ +# 批量换货脚本 + +该脚本读取两列 CSV,逐行调用 `POST /api/admin/exchanges` 执行换货。脚本固定使用: + +- `flow_type=direct`:只支持直接换货,接口创建换货单后会立即完成换货。 +- `migrate_data=true`:必须执行全量数据迁移,不能通过参数关闭。 + +全量迁移由现有换货接口在同一事务中执行,包括钱包余额、套餐使用记录、累计充值字段和资产标签。脚本仅使用 Python 标准库,不需要安装依赖。 + +默认只预演,必须增加 `--execute` 才会真实换货。 + +## CSV 格式 + +首行表头可选,第一列填写旧资产标识,第二列填写新资产标识: + +```csv +old_identifier,new_identifier +89860000000000000001,89860000000000000101 +89860000000000000002,89860000000000000102 +``` + +支持表头: + +- 英文:`old_identifier,new_identifier` 或 `old_asset_identifier,new_asset_identifier` +- 中文:`旧资产标识,新资产标识` 或 `旧资产,新资产` + +旧、新资产标识均使用换货接口已有的识别规则:物联网卡支持 ICCID、接入号、虚拟号;设备支持虚拟号、IMEI、SN。 + +脚本会在请求前拦截: + +- 列数不是两列,或任一列为空。 +- 同一行新旧资产相同。 +- 旧资产重复,或新资产重复。 +- 同一资产在本批次中既作为旧资产又作为新资产。该情况会受到执行顺序影响,因此整批终止。 + +## 预演 + +物联网卡示例: + +```bash +python3 scripts/batch_exchange/batch_exchange.py \ + --base-url https://cmp-api.example.com \ + --csv scripts/batch_exchange/exchanges.example.csv \ + --asset-type iot_card +``` + +设备批次将 `--asset-type` 改为 `device`。同一个 CSV 批次只能使用一种资产类型;新资产必须与旧资产类型一致,否则接口会拒绝该行。 + +预演只校验 CSV 并展示最多五条请求示例,不需要 Token,也不会调用接口。输出中会明确显示 `direct` 和 `migrate_data=true`。 + +## 使用已有 Token 执行 + +推荐通过环境变量传递 Token,避免进入命令历史: + +```bash +JUNHONG_ADMIN_TOKEN='<后台Access Token>' \ +python3 scripts/batch_exchange/batch_exchange.py \ + --base-url https://cmp-api.example.com \ + --csv /path/to/exchanges.csv \ + --asset-type iot_card \ + --execute +``` + +## 使用账号自动登录后执行 + +未提供 Token 时,可以使用后台账号调用 `/api/admin/login` 自动获取 Token: + +```bash +JUNHONG_ADMIN_USERNAME='<后台账号>' \ +JUNHONG_ADMIN_PASSWORD='<后台密码>' \ +python3 scripts/batch_exchange/batch_exchange.py \ + --base-url https://cmp-api.example.com \ + --csv /path/to/exchanges.csv \ + --asset-type iot_card \ + --execute +``` + +也可以使用 `--token`、`--username`、`--password` 参数。密码优先通过环境变量或交互输入,避免保存在 Shell 历史中。 + +## 换货原因和备注 + +脚本默认使用 `批量直接换货` 作为换货原因。可按整批覆盖原因和备注: + +```bash +python3 scripts/batch_exchange/batch_exchange.py \ + --base-url https://cmp-api.example.com \ + --csv /path/to/exchanges.csv \ + --asset-type iot_card \ + --exchange-reason '故障卡批量换货' \ + --remark '2026年7月批次' \ + --execute +``` + +## 结果文件 + +默认在输入 CSV 同目录生成: + +```text +原文件名_换货结果_YYYYMMDD_HHMMSS.csv +``` + +结果包含新旧资产标识、成功或失败状态、HTTP 状态码、业务错误码、错误消息、换货单 ID、换货单号、迁移完成状态和迁移余额。每处理一条都会立即刷新文件,中途中断时已完成的结果不会丢失。 + +为避免覆盖历史结果,`--output` 指定的文件已经存在时脚本会直接报错。 + +## 其他参数和执行约束 + +- `--timeout`:单次请求超时秒数,默认 `30`。 +- `--interval`:每次换货后的等待秒数,默认 `0.2`。 +- `--execute`:显式开启真实换货。 + +每组资产单独调用一次接口,脚本不会自动重试 POST 请求。接口可能已经成功提交事务但客户端没有收到响应,自动重试可能造成误判;失败项应先查询换货单或资产状态,再决定是否单独重跑。 + +脚本不会回滚前面已经成功的行。执行前应先预演并确认完整映射;执行后根据结果 CSV 逐条核对失败项。 diff --git a/scripts/batch_exchange/batch_exchange.py b/scripts/batch_exchange/batch_exchange.py new file mode 100755 index 0000000..789386d --- /dev/null +++ b/scripts/batch_exchange/batch_exchange.py @@ -0,0 +1,613 @@ +#!/usr/bin/env python3 +"""批量直接换货脚本:读取双列 CSV,逐行创建必须迁移数据的直接换货单。 + +默认只执行预演。只有显式传入 --execute 时,才会调用 +POST /api/admin/exchanges。脚本仅使用 Python 标准库,不需要安装第三方依赖。 +""" +from __future__ import annotations + +import argparse +import csv +import json +import os +import sys +import time +from dataclasses import dataclass +from datetime import datetime +from getpass import getpass +from pathlib import Path +from typing import Any +from urllib.error import HTTPError, URLError +from urllib.request import Request, urlopen + + +EXCHANGE_PATH = "/api/admin/exchanges" +LOGIN_PATH = "/api/admin/login" +AUTH_ERROR_CODES = {1002, 1003, 1004} +ASSET_TYPES = {"iot_card", "device"} +DEFAULT_EXCHANGE_REASON = "批量直接换货" +OLD_HEADER_NAMES = { + "old_identifier", + "old_asset_identifier", + "旧资产标识", + "旧资产", +} +NEW_HEADER_NAMES = { + "new_identifier", + "new_asset_identifier", + "新资产标识", + "新资产", +} + + +@dataclass(frozen=True) +class ExchangeInput: + """保存 CSV 中的一组新旧资产标识及原始行号。""" + + line_no: int + old_identifier: str + new_identifier: str + + +@dataclass(frozen=True) +class HTTPResult: + """保存一次 HTTP 请求的响应信息。""" + + status: int + body: dict[str, Any] | None + raw_body: str + + +@dataclass(frozen=True) +class ExchangeOutcome: + """保存单组资产的换货结果及结果文件行。""" + + row: dict[str, object] + success: bool + http_status: int | None + code: int | str | None + message: str + exchange_no: str + + +class RequestFailedError(Exception): + """表示请求尚未获得可解析的 HTTP 响应。""" + + +class AdminAPIClient: + """调用后台认证和换货接口的轻量客户端。""" + + def __init__(self, base_url: str, timeout: float) -> None: + self.base_url = base_url.rstrip("/") + self.timeout = timeout + + def login(self, username: str, password: str) -> str: + """使用后台账号登录并返回 Access Token。""" + result = self._post_json( + LOGIN_PATH, + {"username": username, "password": password, "device": "web"}, + token=None, + ) + code = response_code(result.body) + if not is_success(result.status, code): + raise RequestFailedError( + f"登录失败:HTTP {result.status},code={display_value(code)}," + f"msg={response_message(result.body, result.raw_body)}" + ) + + data = result.body.get("data") if result.body else None + token = data.get("access_token") if isinstance(data, dict) else None + if not isinstance(token, str) or not token.strip(): + raise RequestFailedError("登录响应中缺少 data.access_token") + return token.strip() + + def create_exchange( + self, + token: str, + exchange: ExchangeInput, + asset_type: str, + exchange_reason: str, + remark: str | None, + ) -> HTTPResult: + """创建并立即完成一组必须迁移数据的直接换货。""" + payload = build_exchange_payload(exchange, asset_type, exchange_reason, remark) + return self._post_json(EXCHANGE_PATH, payload, token=token) + + def _post_json(self, path: str, payload: dict[str, Any], token: str | None) -> HTTPResult: + body = json.dumps(payload, ensure_ascii=False, separators=(",", ":")).encode("utf-8") + headers = { + "Accept": "application/json", + "Content-Type": "application/json", + "User-Agent": "junhong-batch-exchange/1.0", + } + if token: + headers["Authorization"] = f"Bearer {token}" + + request = Request( + url=self.base_url + path, + data=body, + headers=headers, + method="POST", + ) + try: + with urlopen(request, timeout=self.timeout) as response: + raw_body = response.read().decode("utf-8", errors="replace") + return HTTPResult( + status=response.status, + body=parse_json_object(raw_body), + raw_body=raw_body, + ) + except HTTPError as exc: + raw_body = exc.read().decode("utf-8", errors="replace") + return HTTPResult( + status=exc.code, + body=parse_json_object(raw_body), + raw_body=raw_body, + ) + except (URLError, TimeoutError, OSError) as exc: + raise RequestFailedError(f"请求失败:{exc}") from exc + + +def parse_args() -> argparse.Namespace: + """解析命令行参数。""" + parser = argparse.ArgumentParser( + description="读取双列 CSV,逐行调用 /api/admin/exchanges 执行直接换货和数据迁移", + ) + parser.add_argument( + "--base-url", + default=os.getenv("JUNHONG_ADMIN_BASE_URL", ""), + help="接口 Base URL,例如 https://cmp-api.example.com;也可使用 JUNHONG_ADMIN_BASE_URL", + ) + parser.add_argument( + "--csv", + required=True, + help="双列资产 CSV 文件路径,依次为旧资产标识、新资产标识,首行可有表头", + ) + parser.add_argument( + "--asset-type", + required=True, + choices=sorted(ASSET_TYPES), + help="本批资产类型:iot_card(物联网卡)或 device(设备)", + ) + parser.add_argument( + "--exchange-reason", + default=DEFAULT_EXCHANGE_REASON, + help=f"本批换货原因(默认:{DEFAULT_EXCHANGE_REASON})", + ) + parser.add_argument("--remark", default="", help="本批换货备注;默认不传") + parser.add_argument( + "--token", + default=os.getenv("JUNHONG_ADMIN_TOKEN", ""), + help="后台 Access Token;也可使用 JUNHONG_ADMIN_TOKEN", + ) + parser.add_argument( + "--username", + default=os.getenv("JUNHONG_ADMIN_USERNAME", ""), + help="未提供 Token 时用于自动登录;也可使用 JUNHONG_ADMIN_USERNAME", + ) + parser.add_argument( + "--password", + default=os.getenv("JUNHONG_ADMIN_PASSWORD", ""), + help="后台登录密码;建议使用 JUNHONG_ADMIN_PASSWORD,避免进入命令历史", + ) + parser.add_argument("--output", default="", help="结果 CSV 路径;默认输出到输入文件同目录") + parser.add_argument("--timeout", type=float, default=30.0, help="单次请求超时秒数(默认 30)") + parser.add_argument("--interval", type=float, default=0.2, help="每次换货后的间隔秒数(默认 0.2)") + parser.add_argument( + "--execute", + action="store_true", + help="真实调用接口换货;不传时只校验 CSV 并预览请求", + ) + return parser.parse_args() + + +def load_exchanges(csv_path: Path) -> list[ExchangeInput]: + """读取双列 CSV,并在执行前拦截可能破坏批次映射的数据。""" + if not csv_path.exists(): + raise ValueError(f"找不到 CSV 文件:{csv_path}") + if not csv_path.is_file(): + raise ValueError(f"CSV 路径不是文件:{csv_path}") + + exchanges: list[ExchangeInput] = [] + old_line_by_identifier: dict[str, int] = {} + new_line_by_identifier: dict[str, int] = {} + errors: list[str] = [] + first_nonempty_seen = False + + with csv_path.open("r", encoding="utf-8-sig", newline="") as file: + reader = csv.reader(file) + for line_no, row in enumerate(reader, start=1): + if not any(value.strip() for value in row): + continue + if len(row) != 2: + errors.append(f"第 {line_no} 行必须正好有两列,实际读取到 {len(row)} 列") + continue + + old_identifier, new_identifier = (value.strip() for value in row) + if not first_nonempty_seen: + first_nonempty_seen = True + if is_header_row(old_identifier, new_identifier): + continue + + if not old_identifier or not new_identifier: + errors.append(f"第 {line_no} 行的旧资产标识和新资产标识均不能为空") + continue + if len(old_identifier) > 100 or len(new_identifier) > 100: + errors.append(f"第 {line_no} 行的资产标识不能超过 100 个字符") + continue + if old_identifier == new_identifier: + errors.append(f"第 {line_no} 行的新旧资产标识相同:{old_identifier}") + continue + if old_identifier in old_line_by_identifier: + errors.append( + f"第 {line_no} 行旧资产与第 {old_line_by_identifier[old_identifier]} 行重复:" + f"{old_identifier}" + ) + continue + if new_identifier in new_line_by_identifier: + errors.append( + f"第 {line_no} 行新资产与第 {new_line_by_identifier[new_identifier]} 行重复:" + f"{new_identifier}" + ) + continue + + old_line_by_identifier[old_identifier] = line_no + new_line_by_identifier[new_identifier] = line_no + exchanges.append( + ExchangeInput( + line_no=line_no, + old_identifier=old_identifier, + new_identifier=new_identifier, + ) + ) + + for identifier in old_line_by_identifier.keys() & new_line_by_identifier.keys(): + errors.append( + f"资产 {identifier} 同时作为第 {old_line_by_identifier[identifier]} 行旧资产和" + f"第 {new_line_by_identifier[identifier]} 行新资产,批次执行顺序会改变其状态" + ) + + if errors: + preview = "\n".join(f" - {error}" for error in errors[:20]) + if len(errors) > 20: + preview += f"\n - 其余 {len(errors) - 20} 个错误已省略" + raise ValueError(f"CSV 校验失败,请修正后重试:\n{preview}") + if not exchanges: + raise ValueError("CSV 中没有有效的换货资产映射") + return exchanges + + +def is_header_row(old_identifier: str, new_identifier: str) -> bool: + """判断首个非空行是否为支持的双列表头。""" + return old_identifier.lower() in OLD_HEADER_NAMES and new_identifier.lower() in NEW_HEADER_NAMES + + +def parse_json_object(raw_body: str) -> dict[str, Any] | None: + """尝试把响应正文解析为 JSON 对象。""" + if not raw_body.strip(): + return None + try: + value = json.loads(raw_body) + except json.JSONDecodeError: + return None + return value if isinstance(value, dict) else None + + +def response_code(body: dict[str, Any] | None) -> int | str | None: + """读取统一响应中的业务错误码。""" + if not body: + return None + code = body.get("code") + if isinstance(code, bool): + return int(code) + if isinstance(code, int): + return code + if isinstance(code, str): + stripped = code.strip() + return int(stripped) if stripped.isdigit() else stripped + return None + + +def response_message(body: dict[str, Any] | None, raw_body: str) -> str: + """读取统一响应消息,非 JSON 响应则保留截断后的正文。""" + if body: + message = body.get("msg", body.get("message", "")) + if message is not None and str(message).strip(): + return str(message).strip() + text = raw_body.strip().replace("\r", " ").replace("\n", " ") + return text[:500] if text else "接口未返回错误信息" + + +def is_success(http_status: int, code: int | str | None) -> bool: + """同时校验 HTTP 状态码和业务响应码。""" + return 200 <= http_status < 300 and str(code) == "0" + + +def display_value(value: object) -> str: + """把可能为空的字段转为适合日志和 CSV 的文本。""" + return "" if value is None else str(value) + + +def resolve_output_path(input_path: Path, output_arg: str) -> Path: + """生成本批次的结果文件路径。""" + if output_arg.strip(): + return Path(output_arg).expanduser().resolve() + timestamp = datetime.now().strftime("%Y%m%d_%H%M%S") + return input_path.with_name(f"{input_path.stem}_换货结果_{timestamp}.csv") + + +def resolve_token(args: argparse.Namespace, client: AdminAPIClient) -> str: + """优先使用现有 Token,否则使用后台账号自动登录。""" + token = args.token.strip() + if token: + return token + + username = args.username.strip() + if not username: + raise ValueError( + "真实执行需要 --token,或同时提供 --username/--password;" + "也可通过 JUNHONG_ADMIN_TOKEN 等环境变量配置" + ) + + password = args.password + if not password and sys.stdin.isatty(): + password = getpass("请输入后台登录密码:") + if not password: + raise ValueError("使用账号登录时必须提供密码") + + print(f"正在使用后台账号 {username!r} 获取 Access Token...") + return client.login(username, password) + + +def build_exchange_payload( + exchange: ExchangeInput, + asset_type: str, + exchange_reason: str, + remark: str | None, +) -> dict[str, Any]: + """构造预演和真实执行共同使用的固定直接换货请求。""" + payload: dict[str, Any] = { + "old_asset_type": asset_type, + "old_identifier": exchange.old_identifier, + "flow_type": "direct", + "new_identifier": exchange.new_identifier, + "migrate_data": True, + "exchange_reason": exchange_reason, + } + if remark is not None: + payload["remark"] = remark + return payload + + +def preview( + exchanges: list[ExchangeInput], + base_url: str, + asset_type: str, + exchange_reason: str, + remark: str | None, +) -> int: + """输出预演信息,不发送任何 HTTP 请求。""" + print("预演完成:未发送任何 HTTP 请求。") + print(f"接口地址:{base_url.rstrip('/')}{EXCHANGE_PATH}") + print(f"换货数量:{len(exchanges)}") + print(f"资产类型:{asset_type}") + print("换货流程:direct(直接换货)") + print("数据迁移:true(必须迁移)") + print(f"换货原因:{exchange_reason}") + print("请求示例:") + for exchange in exchanges[:5]: + payload = build_exchange_payload(exchange, asset_type, exchange_reason, remark) + print(f" 第 {exchange.line_no} 行:{json.dumps(payload, ensure_ascii=False)}") + if len(exchanges) > 5: + print(f" 其余 {len(exchanges) - 5} 条已省略") + print("确认资产映射无误后,增加 --execute 才会真实换货。") + return 0 + + +def exchange_asset( + client: AdminAPIClient, + token: str, + exchange: ExchangeInput, + asset_type: str, + exchange_reason: str, + remark: str | None, +) -> ExchangeOutcome: + """调用一次换货接口,并转换为统一的结果记录。""" + try: + result = client.create_exchange(token, exchange, asset_type, exchange_reason, remark) + code = response_code(result.body) + message = response_message(result.body, result.raw_body) + data = result.body.get("data") if result.body else None + exchange_data = data if isinstance(data, dict) else {} + request_succeeded = is_success(result.status, code) + migration_completed = exchange_data.get("migration_completed") is True + success = request_succeeded and migration_completed + if request_succeeded and not migration_completed: + message = "接口返回成功,但响应未确认数据迁移完成,请人工核对该换货单" + exchange_no = display_value(exchange_data.get("exchange_no")) + return ExchangeOutcome( + row={ + "line_no": exchange.line_no, + "old_identifier": exchange.old_identifier, + "new_identifier": exchange.new_identifier, + "status": "成功" if success else "失败", + "http_status": result.status, + "code": display_value(code), + "msg": message, + "exchange_id": display_value(exchange_data.get("id")), + "exchange_no": exchange_no, + "migration_completed": display_value(exchange_data.get("migration_completed")), + "migration_balance": display_value(exchange_data.get("migration_balance")), + }, + success=success, + http_status=result.status, + code=code, + message=message, + exchange_no=exchange_no, + ) + except RequestFailedError as exc: + message = str(exc) + return ExchangeOutcome( + row={ + "line_no": exchange.line_no, + "old_identifier": exchange.old_identifier, + "new_identifier": exchange.new_identifier, + "status": "失败", + "http_status": "", + "code": "", + "msg": message, + "exchange_id": "", + "exchange_no": "", + "migration_completed": "", + "migration_balance": "", + }, + success=False, + http_status=None, + code=None, + message=message, + exchange_no="", + ) + + +def execute( + exchanges: list[ExchangeInput], + client: AdminAPIClient, + token: str, + asset_type: str, + exchange_reason: str, + remark: str | None, + output_path: Path, + interval: float, +) -> int: + """顺序执行换货,并把每条结果立即写入 CSV。""" + output_path.parent.mkdir(parents=True, exist_ok=True) + success_count = 0 + failed_count = 0 + total = len(exchanges) + + with output_path.open("w", encoding="utf-8-sig", newline="") as file: + writer = csv.DictWriter( + file, + fieldnames=[ + "line_no", + "old_identifier", + "new_identifier", + "status", + "http_status", + "code", + "msg", + "exchange_id", + "exchange_no", + "migration_completed", + "migration_balance", + ], + ) + writer.writeheader() + file.flush() + + for index, exchange in enumerate(exchanges, start=1): + outcome = exchange_asset( + client, + token, + exchange, + asset_type, + exchange_reason, + remark, + ) + writer.writerow(outcome.row) + if outcome.success: + success_count += 1 + print( + f"[{index}/{total}] 成功:{exchange.old_identifier} -> " + f"{exchange.new_identifier},换货单号={outcome.exchange_no}" + ) + else: + failed_count += 1 + print( + f"[{index}/{total}] 失败:{exchange.old_identifier} -> " + f"{exchange.new_identifier},HTTP={display_value(outcome.http_status)}," + f"code={display_value(outcome.code)},msg={outcome.message}", + file=sys.stderr, + ) + + file.flush() + + if outcome.http_status == 401 or outcome.code in AUTH_ERROR_CODES: + print("认证已失效,停止后续换货;已处理结果已保存。", file=sys.stderr) + break + if interval > 0 and index < total: + time.sleep(interval) + + print() + print(f"执行结束:成功 {success_count} 条,失败 {failed_count} 条。") + print(f"结果文件:{output_path}") + return 0 if failed_count == 0 and success_count == total else 2 + + +def main() -> int: + """校验参数,执行预演或真实批量换货。""" + args = parse_args() + try: + base_url = args.base_url.strip() + if not base_url: + raise ValueError("必须通过 --base-url 或 JUNHONG_ADMIN_BASE_URL 配置接口地址") + if not base_url.startswith(("http://", "https://")): + raise ValueError("base-url 必须以 http:// 或 https:// 开头") + exchange_reason = args.exchange_reason.strip() + if not exchange_reason: + raise ValueError("exchange-reason 不能为空") + if len(exchange_reason) > 100: + raise ValueError("exchange-reason 不能超过 100 个字符") + remark = args.remark.strip() or None + if remark is not None and len(remark) > 500: + raise ValueError("remark 不能超过 500 个字符") + if args.timeout <= 0: + raise ValueError("timeout 必须大于 0") + if args.interval < 0: + raise ValueError("interval 不能小于 0") + + input_path = Path(args.csv).expanduser().resolve() + exchanges = load_exchanges(input_path) + if not args.execute: + return preview( + exchanges, + base_url, + args.asset_type, + exchange_reason, + remark, + ) + + output_path = resolve_output_path(input_path, args.output) + if output_path == input_path: + raise ValueError("结果文件不能与输入 CSV 使用同一路径") + if output_path.exists(): + raise ValueError(f"结果文件已存在,请更换 --output 路径:{output_path}") + client = AdminAPIClient(base_url, args.timeout) + token = resolve_token(args, client) + + print( + f"即将真实换货:{len(exchanges)} 组,资产类型={args.asset_type}," + "流程=direct,迁移数据=true" + ) + print(f"接口地址:{base_url.rstrip('/')}{EXCHANGE_PATH}") + print(f"结果文件:{output_path}") + return execute( + exchanges=exchanges, + client=client, + token=token, + asset_type=args.asset_type, + exchange_reason=exchange_reason, + remark=remark, + output_path=output_path, + interval=args.interval, + ) + except (ValueError, RequestFailedError) as exc: + print(f"错误:{exc}", file=sys.stderr) + return 1 + except KeyboardInterrupt: + print("\n用户中断执行;已写入的结果会保留。", file=sys.stderr) + return 130 + + +if __name__ == "__main__": + sys.exit(main()) diff --git a/scripts/batch_exchange/exchanges.example.csv b/scripts/batch_exchange/exchanges.example.csv new file mode 100644 index 0000000..5e45e68 --- /dev/null +++ b/scripts/batch_exchange/exchanges.example.csv @@ -0,0 +1,38 @@ +old_identifier,new_identifier +99840868838,89861590172420400342 +99840868841,89861590172420400343 +99840868842,89861590172420400344 +99840868843,89861590172420400345 +99840868844,89861590172420400346 +99840868845,89861590172420400347 +99840868847,89861590172420400348 +99840868848,89861590172420400349 +99840868850,89861590172420400351 +99840868851,89861590172420400352 +99840868852,89861590172420400353 +99840868854,89861590172420400354 +99840868856,89861590172420400356 +99840868863,89861590172420400363 +99840868865,89861590172420400364 +99820416468,89861590172420400365 +99840868917,89861590172420400319 +99840868918,89861590172420400320 +99840868919,89861590172420400321 +99840868920,89861590172420400322 +99840868922,89861590172420400323 +99840868923,89861590172420400324 +99840868924,89861590172420400325 +99840868925,89861590172420400326 +99840868938,89861590172420400327 +99840868939,89861590172420400328 +99840868940,89861590172420400329 +99840868941,89861590172420400330 +99840868942,89861590172420400331 +99840868943,89861590172420400332 +99840868944,89861590172420400333 +99840868945,89861590172420400334 +99840868946,89861590172420400335 +99840868948,89861590172420400336 +99840868949,89861590172420400337 +99840868951,89861590172420400338 +99840868952,89861590172420400339 diff --git a/scripts/batch_exchange/test_batch_exchange.py b/scripts/batch_exchange/test_batch_exchange.py new file mode 100644 index 0000000..2027aa2 --- /dev/null +++ b/scripts/batch_exchange/test_batch_exchange.py @@ -0,0 +1,172 @@ +"""批量换货脚本的输入校验测试。""" +from __future__ import annotations + +import importlib.util +import sys +import tempfile +import unittest +from pathlib import Path + + +SCRIPT_PATH = Path(__file__).with_name("batch_exchange.py") +SPEC = importlib.util.spec_from_file_location("batch_exchange", SCRIPT_PATH) +assert SPEC is not None and SPEC.loader is not None +batch_exchange = importlib.util.module_from_spec(SPEC) +sys.modules[SPEC.name] = batch_exchange +SPEC.loader.exec_module(batch_exchange) + + +class LoadExchangesTest(unittest.TestCase): + """验证 CSV 映射校验不会把危险批次交给接口执行。""" + + def write_csv(self, content: str) -> Path: + """在临时目录中创建待解析的 CSV。""" + directory = tempfile.TemporaryDirectory() + self.addCleanup(directory.cleanup) + path = Path(directory.name) / "exchanges.csv" + path.write_text(content, encoding="utf-8") + return path + + def test_loads_supported_header_and_rows(self) -> None: + """支持标准表头并保留原始行号。""" + path = self.write_csv("old_identifier,new_identifier\nold-1,new-1\nold-2,new-2\n") + + exchanges = batch_exchange.load_exchanges(path) + + self.assertEqual( + exchanges, + [ + batch_exchange.ExchangeInput(2, "old-1", "new-1"), + batch_exchange.ExchangeInput(3, "old-2", "new-2"), + ], + ) + + def test_rejects_duplicate_old_identifier(self) -> None: + """同一旧资产不能在一批中换出两次。""" + path = self.write_csv("old-1,new-1\nold-1,new-2\n") + + with self.assertRaisesRegex(ValueError, "旧资产.*重复"): + batch_exchange.load_exchanges(path) + + def test_rejects_duplicate_new_identifier(self) -> None: + """同一新资产不能接收两次换货。""" + path = self.write_csv("old-1,new-1\nold-2,new-1\n") + + with self.assertRaisesRegex(ValueError, "新资产.*重复"): + batch_exchange.load_exchanges(path) + + def test_rejects_same_identifier_in_one_row(self) -> None: + """禁止资产换给自身。""" + path = self.write_csv("old-1,old-1\n") + + with self.assertRaisesRegex(ValueError, "新旧资产标识相同"): + batch_exchange.load_exchanges(path) + + def test_rejects_cross_row_asset_reuse(self) -> None: + """禁止资产在同批中同时作为旧资产和新资产。""" + path = self.write_csv("old-1,new-1\nnew-1,new-2\n") + + with self.assertRaisesRegex(ValueError, "同时作为.*旧资产.*新资产"): + batch_exchange.load_exchanges(path) + + def test_rejects_missing_column_value(self) -> None: + """双列中的空值必须在请求前被拦截。""" + path = self.write_csv("old_identifier,new_identifier\nold-1,\n") + + with self.assertRaisesRegex(ValueError, "均不能为空"): + batch_exchange.load_exchanges(path) + + +class BuildExchangePayloadTest(unittest.TestCase): + """验证脚本不能关闭直接换货或数据迁移。""" + + def test_forces_direct_flow_and_migration(self) -> None: + """请求固定为 direct 且 migrate_data 为 true。""" + exchange = batch_exchange.ExchangeInput(2, "old-1", "new-1") + + payload = batch_exchange.build_exchange_payload( + exchange, + "iot_card", + "批量直接换货", + None, + ) + + self.assertEqual(payload["flow_type"], "direct") + self.assertIs(payload["migrate_data"], True) + self.assertEqual(payload["old_asset_type"], "iot_card") + + +class ExchangeAssetTest(unittest.TestCase): + """验证结果必须同时满足接口成功和迁移完成。""" + + class FakeClient: + """返回指定响应的换货客户端替身。""" + + def __init__(self, result: object) -> None: + self.result = result + + def create_exchange(self, *args: object) -> object: + """返回测试预设的接口响应。""" + return self.result + + def test_marks_completed_migration_as_success(self) -> None: + """业务成功且迁移完成时结果为成功。""" + result = batch_exchange.HTTPResult( + status=200, + body={ + "code": 0, + "msg": "success", + "data": { + "id": 123, + "exchange_no": "EX123", + "migration_completed": True, + "migration_balance": 100, + }, + }, + raw_body="", + ) + + outcome = batch_exchange.exchange_asset( + self.FakeClient(result), + "token", + batch_exchange.ExchangeInput(2, "old-1", "new-1"), + "iot_card", + "批量直接换货", + None, + ) + + self.assertTrue(outcome.success) + self.assertEqual(outcome.exchange_no, "EX123") + + def test_rejects_success_response_without_completed_migration(self) -> None: + """业务成功但未确认迁移完成时结果仍为失败。""" + result = batch_exchange.HTTPResult( + status=200, + body={ + "code": 0, + "msg": "success", + "data": { + "id": 123, + "exchange_no": "EX123", + "migration_completed": False, + "migration_balance": 0, + }, + }, + raw_body="", + ) + + outcome = batch_exchange.exchange_asset( + self.FakeClient(result), + "token", + batch_exchange.ExchangeInput(2, "old-1", "new-1"), + "iot_card", + "批量直接换货", + None, + ) + + self.assertFalse(outcome.success) + self.assertIn("未确认数据迁移完成", outcome.message) + + +if __name__ == "__main__": + unittest.main()