#!/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()