From 30a24aab687bf460f4aa7d2a4dba056f5a081adb Mon Sep 17 00:00:00 2001 From: liangzai <2440983361@qq.com> Date: Sat, 27 Jun 2026 23:56:18 +0800 Subject: [PATCH] f --- .../api/api_application_service.go | 3 + internal/application/api/errors.go | 2 + .../domains/api/services/processors/errors.go | 9 +- .../processors/ivyz/ivyz18hy_processor.go | 7 +- .../processors/ivyz/ivyz28hy_processor.go | 7 +- .../processors/ivyz/ivyz38sr_processor.go | 7 +- .../processors/ivyz/ivyz48sr_processor.go | 7 +- scripts/batch_flxg7e8f_analyze.py | 573 ++++++++++++++++++ scripts/decrypt_api_calls.py | 199 ++++++ scripts/export_combmy01_api_calls.sql | 27 + 10 files changed, 813 insertions(+), 28 deletions(-) create mode 100644 scripts/batch_flxg7e8f_analyze.py create mode 100644 scripts/decrypt_api_calls.py create mode 100644 scripts/export_combmy01_api_calls.sql diff --git a/internal/application/api/api_application_service.go b/internal/application/api/api_application_service.go index 63746b1..976e207 100644 --- a/internal/application/api/api_application_service.go +++ b/internal/application/api/api_application_service.go @@ -425,6 +425,9 @@ func (s *ApiApplicationServiceImpl) callExternalApi(ctx context.Context, cmd *co zap.String("error_type", mappedErrorType), zap.Error(err)) + if errors.Is(err, processors.ErrContactBusiness) { + return "", ErrContactBusiness + } if mappedErrorType == entities.ApiCallErrorInvalidParam { return "", ErrInvalidParam } diff --git a/internal/application/api/errors.go b/internal/application/api/errors.go index 6a6e653..c60330f 100644 --- a/internal/application/api/errors.go +++ b/internal/application/api/errors.go @@ -25,6 +25,7 @@ var ( ErrSubscriptionExpired = errors.New("订阅已过期") ErrSubscriptionSuspended = errors.New("订阅已暂停") ErrBusiness = errors.New("业务失败") + ErrContactBusiness = errors.New("请联系商务咨询") ErrSubordinateLinkNotFound = errors.New("非子账号,无法使用master_accessid") ErrSubordinateParentMismatch = errors.New("master_accessid与主账号不匹配") ErrMissingMgmtKey = errors.New("缺少管理密钥") @@ -55,6 +56,7 @@ var ErrorCodeMap = map[error]int{ ErrSubscriptionExpired: 1008, ErrSubscriptionSuspended: 1008, ErrBusiness: 2001, + ErrContactBusiness: 2001, ErrSubordinateLinkNotFound: 1301, ErrSubordinateParentMismatch: 1302, ErrMissingMgmtKey: 1010, diff --git a/internal/domains/api/services/processors/errors.go b/internal/domains/api/services/processors/errors.go index 87556e2..7830f69 100644 --- a/internal/domains/api/services/processors/errors.go +++ b/internal/domains/api/services/processors/errors.go @@ -4,8 +4,9 @@ import "errors" // 处理器相关错误类型 var ( - ErrDatasource = errors.New("数据源异常") - ErrSystem = errors.New("系统异常") - ErrInvalidParam = errors.New("参数校验不正确") - ErrNotFound = errors.New("查询为空") + ErrDatasource = errors.New("数据源异常") + ErrSystem = errors.New("系统异常") + ErrInvalidParam = errors.New("参数校验不正确") + ErrNotFound = errors.New("查询为空") + ErrContactBusiness = errors.New("请联系商务咨询") ) diff --git a/internal/domains/api/services/processors/ivyz/ivyz18hy_processor.go b/internal/domains/api/services/processors/ivyz/ivyz18hy_processor.go index 8c872ec..317dc89 100644 --- a/internal/domains/api/services/processors/ivyz/ivyz18hy_processor.go +++ b/internal/domains/api/services/processors/ivyz/ivyz18hy_processor.go @@ -21,12 +21,7 @@ func ProcessIVYZ18HYRequest(ctx context.Context, params []byte, deps *processors return nil, errors.Join(processors.ErrInvalidParam, err) } - fixedData := map[string]interface{}{"msg": "请联系商务咨询"} - fixedRespBytes, err := json.Marshal(fixedData) - if err != nil { - return nil, errors.Join(processors.ErrSystem, err) - } - return fixedRespBytes, nil + return nil, processors.ErrContactBusiness authDate := "" if len(paramsDto.AuthDate) >= 8 { diff --git a/internal/domains/api/services/processors/ivyz/ivyz28hy_processor.go b/internal/domains/api/services/processors/ivyz/ivyz28hy_processor.go index b8cdc57..619f812 100644 --- a/internal/domains/api/services/processors/ivyz/ivyz28hy_processor.go +++ b/internal/domains/api/services/processors/ivyz/ivyz28hy_processor.go @@ -21,12 +21,7 @@ func ProcessIVYZ28HYRequest(ctx context.Context, params []byte, deps *processors return nil, errors.Join(processors.ErrInvalidParam, err) } - fixedData := map[string]interface{}{"msg": "请联系商务咨询"} - fixedRespBytes, err := json.Marshal(fixedData) - if err != nil { - return nil, errors.Join(processors.ErrSystem, err) - } - return fixedRespBytes, nil + return nil, processors.ErrContactBusiness reqParams := map[string]interface{}{ "key": "", diff --git a/internal/domains/api/services/processors/ivyz/ivyz38sr_processor.go b/internal/domains/api/services/processors/ivyz/ivyz38sr_processor.go index 8067526..a5c4ead 100644 --- a/internal/domains/api/services/processors/ivyz/ivyz38sr_processor.go +++ b/internal/domains/api/services/processors/ivyz/ivyz38sr_processor.go @@ -21,12 +21,7 @@ func ProcessIVYZ38SRRequest(ctx context.Context, params []byte, deps *processors return nil, errors.Join(processors.ErrInvalidParam, err) } - fixedData := map[string]interface{}{"msg": "请联系商务咨询"} - fixedRespBytes, err := json.Marshal(fixedData) - if err != nil { - return nil, errors.Join(processors.ErrSystem, err) - } - return fixedRespBytes, nil + return nil, processors.ErrContactBusiness reqParams := map[string]interface{}{ "key": "", diff --git a/internal/domains/api/services/processors/ivyz/ivyz48sr_processor.go b/internal/domains/api/services/processors/ivyz/ivyz48sr_processor.go index 8480ec2..b1f8a18 100644 --- a/internal/domains/api/services/processors/ivyz/ivyz48sr_processor.go +++ b/internal/domains/api/services/processors/ivyz/ivyz48sr_processor.go @@ -21,12 +21,7 @@ func ProcessIVYZ48SRRequest(ctx context.Context, params []byte, deps *processors return nil, errors.Join(processors.ErrInvalidParam, err) } - fixedData := map[string]interface{}{"msg": "请联系商务咨询"} - fixedRespBytes, err := json.Marshal(fixedData) - if err != nil { - return nil, errors.Join(processors.ErrSystem, err) - } - return fixedRespBytes, nil + return nil, processors.ErrContactBusiness reqParams := map[string]interface{}{ "key": "", diff --git a/scripts/batch_flxg7e8f_analyze.py b/scripts/batch_flxg7e8f_analyze.py new file mode 100644 index 0000000..d1b4d37 --- /dev/null +++ b/scripts/batch_flxg7e8f_analyze.py @@ -0,0 +1,573 @@ +#!/usr/bin/env python3 +""" +批量调用远程 FLXG7E8F,分析执行/失信/限高三类案件中的金额字段。 + +金额字段含义(见 个人司法涉诉查询_返回字段说明): + 执行案件 implement.cases[]: + n_sqzxbdje 申请执行标的金额 + n_sjdwje 实际到位金额(已还款/已执行到位,repaidAmount 应对应此字段) + n_wzxje 未执行金额(未还款/未执行到位,正确的「未还款金额」应对应此字段) + n_jabdje 结案标的金额(民事结案金额,不能当作还款金额) + + 失信 breachCaseList[]: + estimatedJudgementAmount 判决金额估计(非未还款字段) + 无 n_wzxje;可尝试按案号关联 implement 案件取 n_wzxje + + 限高 consumptionRestrictionList[]: + 无金额字段;可尝试按案号关联 implement 案件取 n_wzxje + +依赖: pip install pycryptodome requests + +示例: + python batch_flxg7e8f_analyze.py \\ + --input ..\\query_1-2026-06-27_32042_decrypted.json \\ + --access-id <你的AccessId> \\ + --secret-key f507c58537aac8227b5f4a99cbd317ff \\ + --output ..\\flxg7e8f_analysis.json \\ + --limit 3 + +获取 access_id: + SELECT access_id, secret_key FROM api_users + WHERE user_id = '2acfef92-0700-4638-b386-bd03a8ad09c3'; +""" + +from __future__ import annotations + +import argparse +import base64 +import json +import os +import re +import sys +import time +from pathlib import Path +from typing import Any + +try: + import requests + from Crypto.Cipher import AES + from Crypto.Random import get_random_bytes +except ImportError: + print("请先安装: pip install pycryptodome requests", file=sys.stderr) + sys.exit(1) + + +DEFAULT_BASE_URL = "https://api.tianyuanapi.com" +API_NAME = "FLXG7E8F" + + +def aes_encrypt(plain_text: bytes, key_hex: str) -> str: + key = bytes.fromhex(key_hex.strip()) + block = 16 + pad = block - len(plain_text) % block + plain_text = plain_text + bytes([pad]) * pad + iv = get_random_bytes(16) + ct = AES.new(key, AES.MODE_CBC, iv).encrypt(plain_text) + return base64.b64encode(iv + ct).decode() + + +def aes_decrypt(cipher_b64: str, key_hex: str) -> Any: + key = bytes.fromhex(key_hex.strip()) + raw = base64.b64decode(cipher_b64.strip()) + iv, ciphertext = raw[:16], raw[16:] + plain = AES.new(key, AES.MODE_CBC, iv).decrypt(ciphertext) + pad = plain[-1] + plain = plain[:-pad] + return json.loads(plain.decode("utf-8")) + + +def to_float(val: Any) -> float | None: + if val is None or val == "": + return None + if isinstance(val, (int, float)): + return float(val) + if isinstance(val, str): + s = val.strip().replace(",", "") + if not s: + return None + try: + return float(s) + except ValueError: + return None + return None + + +def normalize_case_no(case_no: str) -> str: + if not case_no: + return "" + s = case_no.strip().upper() + s = re.sub(r"\s+", "", s) + return s + + +def pick_amount_fields(case_map: dict) -> dict[str, Any]: + sqzx = to_float(case_map.get("n_sqzxbdje")) + sjdw = to_float(case_map.get("n_sjdwje")) + wzx = to_float(case_map.get("n_wzxje")) + jabd = to_float(case_map.get("n_jabdje")) + + computed_unpaid = None + if sqzx is not None and sjdw is not None: + computed_unpaid = max(sqzx - sjdw, 0.0) + + has_correct_unpaid_field = wzx is not None + has_correct_repaid_field = sjdw is not None + + wrong_if_use_jabdje_as_repaid = ( + not has_correct_repaid_field and jabd is not None and jabd > 0 + ) + + return { + "n_sqzxbdje": sqzx, + "n_sjdwje": sjdw, + "n_wzxje": wzx, + "n_jabdje": jabd, + "computed_unpaid_sqzx_minus_sjdw": computed_unpaid, + "has_correct_unpaid_field": has_correct_unpaid_field, + "has_correct_repaid_field": has_correct_repaid_field, + "wrong_if_use_jabdje_as_repaid": wrong_if_use_jabdje_as_repaid, + "unpaid_amount": wzx if wzx is not None else computed_unpaid, + } + + +def build_implement_index(implement_cases: list[dict]) -> dict[str, dict]: + index: dict[str, dict] = {} + for case in implement_cases: + case_no = normalize_case_no(str(case.get("c_ah") or "")) + if not case_no: + continue + amounts = pick_amount_fields(case) + index[case_no] = { + "case_number": case.get("c_ah"), + "case_status": case.get("n_ajjzjd"), + "filing_date": case.get("d_larq"), + "close_date": case.get("d_jarq"), + **amounts, + } + return index + + +def is_case_fully_paid(amounts: dict) -> bool: + """执行案件是否已全部到位。""" + wzx = amounts.get("n_wzxje") + if wzx is not None: + return wzx <= 0 + sqzx = amounts.get("n_sqzxbdje") + sjdw = amounts.get("n_sjdwje") + if sqzx is not None and sjdw is not None: + return sjdw >= sqzx + unpaid = amounts.get("unpaid_amount") + if unpaid is not None: + return unpaid <= 0 + return True + + +def has_outstanding_unpaid(amounts: dict) -> bool: + wzx = amounts.get("n_wzxje") + if wzx is not None: + return wzx > 0 + unpaid = amounts.get("unpaid_amount") + if unpaid is not None: + return unpaid > 0 + sqzx = amounts.get("n_sqzxbdje") + sjdw = amounts.get("n_sjdwje") + if sqzx is not None and sjdw is not None and sqzx > sjdw: + return True + return False + + +def classify_analysis(analysis: dict) -> tuple[str, str]: + """ + 判定是否有「金额映射问题」需要关注: + NO_ISSUE - 无执行/失信/限高,或金额已全部到位 + HAS_OUTSTANDING_UNPAID - 存在未结/未执行金额,需用 n_wzxje / n_sjdwje 正确映射 + """ + if not analysis.get("hit"): + return "NO_ISSUE", "未命中司法数据" + + implement_cases = analysis.get("implement_cases") or [] + breach_cases = analysis.get("breach_cases") or [] + limit_cases = analysis.get("limit_cases") or [] + + if not implement_cases and not breach_cases and not limit_cases: + return "NO_ISSUE", "无执行/失信/限高案件" + + outstanding_details: list[str] = [] + + for case in implement_cases: + if has_outstanding_unpaid(case): + outstanding_details.append(f"执行 {case.get('case_number')} 未执行金额={case.get('n_wzxje') or case.get('unpaid_amount')}") + + for case in breach_cases: + fulfill = (case.get("fulfill_status") or "").strip() + linked = case.get("linked_implement") or {} + if fulfill == "全部未履行": + if linked and has_outstanding_unpaid(linked): + outstanding_details.append(f"失信 {case.get('case_number')} 关联未执行={linked.get('n_wzxje')}") + elif not linked: + outstanding_details.append(f"失信 {case.get('case_number')} 全部未履行(无关联执行明细)") + elif linked and has_outstanding_unpaid(linked): + outstanding_details.append(f"失信 {case.get('case_number')} 关联未执行={linked.get('n_wzxje')}") + + for case in limit_cases: + linked = case.get("linked_implement") or {} + if linked and has_outstanding_unpaid(linked): + outstanding_details.append(f"限高 {case.get('case_number')} 关联未执行={linked.get('n_wzxje')}") + + if outstanding_details: + return "HAS_OUTSTANDING_UNPAID", "; ".join(outstanding_details) + + return "NO_ISSUE", "有执行/失信/限高案件,但金额已全部到位或无未结金额" + + +def analyze_judicial_data(judicial_data: dict) -> dict[str, Any]: + if not judicial_data or judicial_data == -1: + return {"hit": False, "message": "未命中司法数据"} + + lawsuit_stat = judicial_data.get("lawsuitStat") or {} + implement_section = lawsuit_stat.get("implement") or {} + implement_raw = implement_section.get("cases") or [] + + implement_cases = [] + for item in implement_raw: + if isinstance(item, dict): + case_no = item.get("c_ah") + amounts = pick_amount_fields(item) + implement_cases.append( + { + "case_number": case_no, + "case_status": item.get("n_ajjzjd"), + "filing_date": item.get("d_larq"), + "close_date": item.get("d_jarq"), + **amounts, + } + ) + + implement_index = build_implement_index( + [c for c in implement_raw if isinstance(c, dict)] + ) + + breach_cases = [] + for item in judicial_data.get("breachCaseList") or []: + if not isinstance(item, dict): + continue + case_no = normalize_case_no(str(item.get("caseNumber") or "")) + linked = implement_index.get(case_no) + est = to_float(item.get("estimatedJudgementAmount")) + breach_cases.append( + { + "case_number": item.get("caseNumber"), + "fulfill_status": item.get("fulfillStatus"), + "estimated_judgement_amount": est, + "issue_date": item.get("issueDate"), + "linked_implement": linked, + "has_correct_unpaid_via_link": bool( + linked and linked.get("has_correct_unpaid_field") + ), + "unpaid_amount": linked.get("unpaid_amount") if linked else None, + } + ) + + limit_cases = [] + for item in judicial_data.get("consumptionRestrictionList") or []: + if not isinstance(item, dict): + continue + case_no = normalize_case_no(str(item.get("caseNumber") or "")) + linked = implement_index.get(case_no) + limit_cases.append( + { + "case_number": item.get("caseNumber"), + "issue_date": item.get("issueDate"), + "file_date": item.get("fileDate"), + "court": item.get("executiveCourt"), + "linked_implement": linked, + "has_correct_unpaid_via_link": bool( + linked and linked.get("has_correct_unpaid_field") + ), + "unpaid_amount": linked.get("unpaid_amount") if linked else None, + } + ) + + implement_with_outstanding = [c for c in implement_cases if has_outstanding_unpaid(c)] + breach_with_outstanding = [ + c + for c in breach_cases + if (c.get("fulfill_status") == "全部未履行") + or (c.get("linked_implement") and has_outstanding_unpaid(c["linked_implement"])) + ] + limit_with_outstanding = [ + c + for c in limit_cases + if c.get("linked_implement") and has_outstanding_unpaid(c["linked_implement"]) + ] + + result = { + "hit": True, + "summary": { + "implement_count": len(implement_cases), + "breach_count": len(breach_cases), + "limit_count": len(limit_cases), + "implement_with_outstanding_unpaid": len(implement_with_outstanding), + "breach_with_outstanding": len(breach_with_outstanding), + "limit_with_outstanding": len(limit_with_outstanding), + }, + "implement_cases": implement_cases, + "breach_cases": breach_cases, + "limit_cases": limit_cases, + } + flag, reason = classify_analysis(result) + result["flag"] = flag + result["flag_reason"] = reason + return result + + +def call_flxg7e8f( + base_url: str, + access_id: str, + secret_key: str, + name: str, + id_card: str, + timeout: int = 60, +) -> tuple[bool, dict[str, Any]]: + url = f"{base_url.rstrip('/')}/api/v1/{API_NAME}" + params = {"name": name, "id_card": id_card} + encrypted = aes_encrypt(json.dumps(params, ensure_ascii=False).encode("utf-8"), secret_key) + payload = {"data": encrypted, "options": {"json": True}} + + resp = requests.post( + url, + headers={"Access-Id": access_id, "Content-Type": "application/json"}, + json=payload, + timeout=timeout, + ) + resp.raise_for_status() + body = resp.json() + + if body.get("code") != 0: + return False, { + "error_code": body.get("code"), + "error_message": body.get("message"), + "transaction_id": body.get("transaction_id"), + } + + decrypted = aes_decrypt(body["data"], secret_key) + return True, { + "transaction_id": body.get("transaction_id"), + "response": decrypted, + } + + +def load_input_records(path: Path) -> list[dict]: + data = json.loads(path.read_text(encoding="utf-8-sig")) + if not isinstance(data, list): + raise ValueError("输入 JSON 应为数组") + return data + + +def dedupe_persons(records: list[dict]) -> list[dict]: + seen: set[str] = set() + persons = [] + for rec in records: + rp = rec.get("request_params") or {} + if isinstance(rp, str): + continue + id_card = (rp.get("id_card") or "").strip().upper() + name = (rp.get("name") or "").strip() + if not id_card or not name: + continue + if id_card in seen: + continue + seen.add(id_card) + persons.append( + { + "name": name, + "id_card": id_card, + "source_ids": [rec.get("id")], + "source_transaction_ids": [rec.get("transaction_id")], + } + ) + # attach all source ids for duplicates + id_to_person = {p["id_card"]: p for p in persons} + for rec in records: + rp = rec.get("request_params") or {} + if isinstance(rp, str): + continue + id_card = (rp.get("id_card") or "").strip().upper() + if id_card in id_to_person: + p = id_to_person[id_card] + rid, rtx = rec.get("id"), rec.get("transaction_id") + if rid and rid not in p["source_ids"]: + p["source_ids"].append(rid) + if rtx and rtx not in p["source_transaction_ids"]: + p["source_transaction_ids"].append(rtx) + return persons + + +def main() -> None: + parser = argparse.ArgumentParser(description="批量 FLXG7E8F + 未还款金额分析") + parser.add_argument("--input", required=True, help="解密后的 COMBMY01 JSON") + parser.add_argument("--output", default="flxg7e8f_analysis.json") + parser.add_argument("--raw-output", default="flxg7e8f_raw.json", help="FLXG7E8F 解密后源数据汇总") + parser.add_argument("--access-id", default=os.getenv("TYAPI_ACCESS_ID")) + parser.add_argument("--secret-key", default=os.getenv("TYAPI_SECRET_KEY")) + parser.add_argument("--base-url", default=os.getenv("TYAPI_BASE_URL", DEFAULT_BASE_URL)) + parser.add_argument("--delay", type=float, default=0.5, help="每次请求间隔秒数") + parser.add_argument("--limit", type=int, default=0, help="仅处理前 N 个唯一身份证(0=全部)") + parser.add_argument("--resume", help="已有结果 JSON,跳过已成功的人员") + parser.add_argument("--dry-run", action="store_true", help="只列出待调用人员,不请求 API") + args = parser.parse_args() + + if not args.secret_key: + print("请提供 --secret-key 或环境变量 TYAPI_SECRET_KEY", file=sys.stderr) + sys.exit(1) + + input_path = Path(args.input) + records = load_input_records(input_path) + persons = dedupe_persons(records) + if args.limit > 0: + persons = persons[: args.limit] + + print(f"输入记录 {len(records)} 条, 去重后 {len(persons)} 人") + + existing: dict[str, dict] = {} + if args.resume and Path(args.resume).is_file(): + prev = json.loads(Path(args.resume).read_text(encoding="utf-8")) + for item in prev.get("persons", []): + if item.get("api_success"): + existing[item["id_card"]] = item + + if args.dry_run: + for p in persons: + print(f"{p['name']}\t{p['id_card']}\tsources={len(p['source_ids'])}") + return + + if not args.access_id: + print("请提供 --access-id 或环境变量 TYAPI_ACCESS_ID", file=sys.stderr) + sys.exit(1) + + results = [] + raw_records = [] + ok_count = fail_count = no_issue_count = outstanding_count = 0 + + for idx, person in enumerate(persons, start=1): + id_card = person["id_card"] + if id_card in existing: + cached = existing[id_card] + results.append(cached) + if cached.get("raw_response"): + raw_records.append( + { + "name": cached.get("name"), + "id_card": id_card, + "transaction_id": cached.get("transaction_id"), + "raw_response": cached.get("raw_response"), + } + ) + print(f"[{idx}/{len(persons)}] 跳过(已存在) {person['name']} {id_card}") + continue + + print(f"[{idx}/{len(persons)}] 调用 {person['name']} {id_card} ...", flush=True) + item = { + "name": person["name"], + "id_card": id_card, + "source_ids": person["source_ids"], + "source_transaction_ids": person["source_transaction_ids"], + "api_success": False, + } + + try: + success, api_result = call_flxg7e8f( + args.base_url, args.access_id, args.secret_key, person["name"], id_card + ) + item["api_success"] = success + if not success: + item.update(api_result) + fail_count += 1 + else: + item["transaction_id"] = api_result.get("transaction_id") + item["raw_response"] = api_result.get("response") + raw_records.append( + { + "name": person["name"], + "id_card": id_card, + "transaction_id": item["transaction_id"], + "source_ids": person["source_ids"], + "raw_response": item["raw_response"], + } + ) + judicial = (api_result.get("response") or {}).get("judicial_data") + analysis = analyze_judicial_data(judicial) + item["analysis"] = analysis + item["flag"] = analysis.get("flag", "NO_ISSUE") + item["flag_reason"] = analysis.get("flag_reason", "") + ok_count += 1 + if item["flag"] == "HAS_OUTSTANDING_UNPAID": + outstanding_count += 1 + else: + no_issue_count += 1 + except Exception as exc: + item["api_success"] = False + item["error_message"] = str(exc) + fail_count += 1 + + results.append(item) + + # 增量保存,防止中断丢数据 + out_path = Path(args.output) + raw_path = Path(args.raw_output) + payload = { + "meta": { + "input": str(input_path), + "total_records": len(records), + "unique_persons": len(persons), + "processed": len(results), + }, + "persons": results, + } + out_path.write_text(json.dumps(payload, ensure_ascii=False, indent=2), encoding="utf-8") + raw_path.write_text( + json.dumps( + { + "meta": payload["meta"], + "records": raw_records, + }, + ensure_ascii=False, + indent=2, + ), + encoding="utf-8", + ) + + if idx < len(persons) and args.delay > 0: + time.sleep(args.delay) + + summary = { + "api_ok": ok_count, + "api_fail": fail_count, + "no_issue": no_issue_count, + "has_outstanding_unpaid": outstanding_count, + } + + final = { + "meta": { + "input": str(input_path), + "total_records": len(records), + "unique_persons": len(persons), + "summary": summary, + "raw_output": str(Path(args.raw_output).resolve()), + }, + "persons": results, + } + Path(args.output).write_text(json.dumps(final, ensure_ascii=False, indent=2), encoding="utf-8") + Path(args.raw_output).write_text( + json.dumps({"meta": final["meta"], "records": raw_records}, ensure_ascii=False, indent=2), + encoding="utf-8", + ) + + print("\n=== 完成 ===") + print(json.dumps(summary, ensure_ascii=False, indent=2)) + print(f"分析结果: {Path(args.output).resolve()}") + print(f"源数据: {Path(args.raw_output).resolve()}") + + +if __name__ == "__main__": + main() diff --git a/scripts/decrypt_api_calls.py b/scripts/decrypt_api_calls.py new file mode 100644 index 0000000..fb1b03b --- /dev/null +++ b/scripts/decrypt_api_calls.py @@ -0,0 +1,199 @@ +#!/usr/bin/env python3 +""" +将 SQL 导出的 API 调用记录(CSV 或 JSON)解密 request_params,输出 JSON 文件。 + +依赖: pip install pycryptodome + +用法: + python decrypt_api_calls.py -i combmy01_calls.csv -o combmy01_decrypted.json + python decrypt_api_calls.py -i query_1-2026-06-27_32042.json -o combmy01_decrypted.json + +输入必须包含字段: request_params, secret_key(或通过 --secret-key 传入) +可选字段会原样写入输出: id, transaction_id, user_id, product_code, status, cost, start_at, end_at, created_at +输出默认不包含 secret_key。 +""" + +from __future__ import annotations + +import argparse +import base64 +import csv +import json +import sys +from pathlib import Path +from typing import Any + +try: + from Crypto.Cipher import AES +except ImportError: + print("请先安装依赖: pip install pycryptodome", file=sys.stderr) + sys.exit(1) + + +OPTIONAL_FIELDS = ( + "id", + "transaction_id", + "user_id", + "product_code", + "status", + "cost", + "start_at", + "end_at", + "created_at", +) + + +def aes_decrypt(cipher_b64: str, key_hex: str) -> Any: + """与 tyapi-server internal/shared/crypto/crypto.go AesDecrypt 逻辑一致。""" + key = bytes.fromhex(key_hex.strip()) + raw = base64.b64decode(cipher_b64.strip()) + if len(raw) < 16: + raise ValueError("ciphertext too short") + + iv, ciphertext = raw[:16], raw[16:] + plain = AES.new(key, AES.MODE_CBC, iv).decrypt(ciphertext) + + pad = plain[-1] + if pad < 1 or pad > 16 or pad > len(plain): + raise ValueError("invalid padding") + for i in range(pad): + if plain[-1 - i] != pad: + raise ValueError("invalid padding") + plain = plain[:-pad] + + text = plain.decode("utf-8") + return json.loads(text) + + +def normalize_row(row: dict) -> dict[str, str]: + """兼容带 BOM 或大小写不一致的表头;值统一转为字符串。""" + out: dict[str, str] = {} + for k, v in row.items(): + key = (k or "").strip().lower().lstrip("\ufeff") + if v is None: + out[key] = "" + elif isinstance(v, (dict, list)): + out[key] = json.dumps(v, ensure_ascii=False) + else: + out[key] = str(v).strip() + return out + + +def process_row(row: dict, default_secret_key: str | None) -> dict[str, Any]: + row = normalize_row(row) + + encrypted = row.get("request_params", "") + secret_key = row.get("secret_key") or default_secret_key or "" + + if not encrypted: + raise ValueError("request_params 为空") + if not secret_key: + raise ValueError("缺少 secret_key(输入字段或 --secret-key 参数)") + + item: dict[str, Any] = {} + for field in OPTIONAL_FIELDS: + if field in row and row[field] != "": + item[field] = row[field] + + try: + item["request_params"] = aes_decrypt(encrypted, secret_key) + except Exception as exc: + item["request_params"] = None + item["decrypt_error"] = str(exc) + + return item + + +def read_text(path: Path) -> str: + for encoding in ("utf-8-sig", "utf-8", "gbk"): + try: + return path.read_text(encoding=encoding) + except UnicodeDecodeError: + continue + raise ValueError(f"无法读取文件编码: {path}") + + +def read_rows(path: Path) -> list[dict]: + suffix = path.suffix.lower() + if suffix == ".json": + data = json.loads(read_text(path)) + if isinstance(data, list): + return data + if isinstance(data, dict): + for key in ("items", "data", "rows", "records"): + if key in data and isinstance(data[key], list): + return data[key] + raise ValueError("JSON 根节点应为数组,或包含 items/data/rows/records 数组字段") + raise ValueError("不支持的 JSON 结构") + + if suffix == ".csv": + for encoding in ("utf-8-sig", "utf-8", "gbk"): + try: + with path.open("r", encoding=encoding, newline="") as f: + return list(csv.DictReader(f)) + except UnicodeDecodeError: + continue + raise ValueError(f"无法读取 CSV 编码: {path}") + + # 无扩展名时自动探测 + text = read_text(path).lstrip() + if text.startswith("[") or text.startswith("{"): + data = json.loads(text) + if isinstance(data, list): + return data + if isinstance(data, dict): + for key in ("items", "data", "rows", "records"): + if key in data and isinstance(data[key], list): + return data[key] + raise ValueError("不支持的 JSON 结构") + + for encoding in ("utf-8-sig", "utf-8", "gbk"): + try: + with path.open("r", encoding=encoding, newline="") as f: + return list(csv.DictReader(f)) + except UnicodeDecodeError: + continue + raise ValueError(f"无法识别输入格式: {path}") + + +def main() -> None: + parser = argparse.ArgumentParser(description="解密 API 调用记录中的 request_params(支持 CSV / JSON)") + parser.add_argument("-i", "--input", required=True, help="SQL 导出的 CSV 或 JSON 文件路径") + parser.add_argument("-o", "--output", default="api_calls_decrypted.json", help="输出 JSON 文件路径") + parser.add_argument( + "--secret-key", + help="全局 secret_key(若 CSV 无 secret_key 列,或每行相同可省略 CSV 中该列)", + ) + args = parser.parse_args() + + input_path = Path(args.input) + if not input_path.is_file(): + print(f"输入文件不存在: {input_path}", file=sys.stderr) + sys.exit(1) + + rows = read_rows(input_path) + if not rows: + print("输入文件无数据行", file=sys.stderr) + sys.exit(1) + + result = [] + failed = 0 + for idx, row in enumerate(rows, start=1): + try: + result.append(process_row(row, args.secret_key)) + if result[-1].get("decrypt_error"): + failed += 1 + except Exception as exc: + failed += 1 + result.append({"row": idx, "decrypt_error": str(exc), "request_params": None}) + + output_path = Path(args.output) + output_path.write_text(json.dumps(result, ensure_ascii=False, indent=2), encoding="utf-8") + + ok = len(result) - failed + print(f"完成: 共 {len(result)} 条, 成功解密 {ok} 条, 失败 {failed} 条") + print(f"输出: {output_path.resolve()}") + + +if __name__ == "__main__": + main() diff --git a/scripts/export_combmy01_api_calls.sql b/scripts/export_combmy01_api_calls.sql new file mode 100644 index 0000000..caa6e67 --- /dev/null +++ b/scripts/export_combmy01_api_calls.sql @@ -0,0 +1,27 @@ +-- 导出 COMBMY01 成功调用记录(含解密密钥) +-- 用法(MySQL 命令行导出 CSV): +-- mysql -h -u -p --default-character-set=utf8mb4 < export_combmy01_api_calls.sql > combmy01_calls.csv +-- +-- 或在 Navicat / DBeaver 中执行本 SQL,再「导出结果为 CSV」。 +-- 若 JOIN 报错表不存在,把 product 改成 products 试一下。 + +SELECT + ac.id, + ac.transaction_id, + ac.access_id, + ac.user_id, + p.code AS product_code, + ac.request_params, + au.secret_key, + ac.status, + ac.cost, + ac.start_at, + ac.end_at, + ac.created_at +FROM api_calls ac +INNER JOIN product p ON ac.product_id = p.id +INNER JOIN api_users au ON ac.user_id = au.user_id +WHERE ac.user_id = '2acfef92-0700-4638-b386-bd03a8ad09c3' + AND p.code = 'COMBMY01' + AND ac.status = 'success' +ORDER BY ac.created_at DESC;