"""Ai-DOP 第三方对接模拟器 · 本地控制服务(任务书 P1-B / P1-G)。 - 只绑定 127.0.0.1,页面与 API 仅本机可用; - Secret/密码/Token 只从环境变量读取,任何 API 不回显; - 所有对外目标 host 走白名单硬阻断(默认拒绝); - 运行记录(RunStore)写入前统一脱敏; - 内嵌管理 Mock API(providers.mock_api),默认端口 8018,可启停/切鉴权/注入错误; - 三层报告:代码层 SUPPORTED/UNSUPPORTED、配置层 READY/NOT_*、数据层 NOT_RUN/ACCEPTED/TRANSFORM_PENDING/RUN_FAILED。 启动:python tools/integration-simulator/app.py (详见 README.md) """ from __future__ import annotations import asyncio import os import socket import sys import time import uuid from pathlib import Path SIM_ROOT = Path(__file__).resolve().parent if str(SIM_ROOT) not in sys.path: sys.path.insert(0, str(SIM_ROOT)) import uvicorn from fastapi import FastAPI, HTTPException from fastapi.responses import FileResponse from fastapi.staticfiles import StaticFiles from pydantic import BaseModel from catalog import sample_adapters as sa from clients import aidop_probe from clients.inbound_client import InboundClient from clients.redaction import mask_access_key, redact_obj from providers import mock_api, mock_db DEFAULT_CONTROL_PORT = 8017 MAX_RUNS = 200 # 状态枚举(任务书 §4) CONFIG_STATES = ("READY", "CONFIG_NOT_REGISTERED", "CONFIG_NOT_ENABLED", "GRANT_MISSING", "SOURCE_UNREACHABLE", "RUN_FAILED") def _env(name: str) -> str: return os.environ.get(name, "") or "" def _secrets() -> list[str]: return [v for v in (_env("AIDOP_SIM_INBOUND_SECRET"), _env("AIDOP_SIM_MYSQL_PASSWORD"), _env("AIDOP_SIM_ADMIN_TOKEN")) if v] def load_config() -> dict: cfg = { "aidopBaseUrl": "http://127.0.0.1:5005", "mockApiPort": 8018, "mysqlDatabase": mock_db.SIM_DB_DEFAULT, "allowedTargetHosts": ["127.0.0.1", "localhost"], } cfg_file = SIM_ROOT / "config.json" if cfg_file.exists(): import json with cfg_file.open("r", encoding="utf-8") as f: cfg.update(json.load(f)) return cfg CONFIG = load_config() MOCK_STATE = mock_api.MockApiState() MOCK_APP = mock_api.create_app(MOCK_STATE) class RunStore: def __init__(self, limit: int = MAX_RUNS): self.limit = limit self.runs: list[dict] = [] def add(self, run_id: str, object_code: str, channel: str, request_summary: dict, response: object, layers: dict, followup_sql: list[str] | None = None, gaps: list[str] | None = None, scenario_id: str | None = None, nodes: list[dict] | None = None) -> dict: entry = { "runId": run_id, "time": time.strftime("%Y-%m-%dT%H:%M:%S%z"), "objectCode": object_code, "channel": channel, "scenarioId": scenario_id, "nodes": nodes, "request": redact_obj(request_summary, _secrets()), "response": redact_obj(response, _secrets()), "layers": layers, "followupSql": followup_sql or [], "knownGaps": gaps or [], } self.runs.append(entry) if len(self.runs) > self.limit: del self.runs[: len(self.runs) - self.limit] return entry def recent(self) -> list[dict]: return list(reversed(self.runs[-50:])) RUNS = RunStore() app = FastAPI(title="Ai-DOP Integration Simulator", version="1.0.0") # --------------------------------------------------------------------------- # Mock API 内嵌启停 # --------------------------------------------------------------------------- class _MockServerHandle: def __init__(self): self.server: uvicorn.Server | None = None self.task: asyncio.Task | None = None self.last_error: str = "" MOCK_HANDLE = _MockServerHandle() def _port_in_use(port: int) -> bool: """真实探测本机端口是否已被占用。 不依赖 uvicorn 内部属性:0.30+ 的 Server 已无 startup_error, 读它会抛 AttributeError 并把已启动的服务误报成 500。 """ with socket.socket(socket.AF_INET, socket.SOCK_STREAM) as probe: probe.settimeout(0.3) return probe.connect_ex(("127.0.0.1", int(port))) == 0 async def _start_mock_api() -> dict: port = int(CONFIG["mockApiPort"]) if MOCK_HANDLE.task and not MOCK_HANDLE.task.done(): return {"running": True, "port": port, "alreadyRunning": True} # 端口预检必须在建 task 之前:否则失败路径会留下一个已赋值的句柄, # 使 _mock_running() 与实际监听状态不一致。 if _port_in_use(port): MOCK_HANDLE.last_error = f"端口 {port} 已被占用(可能是另一个模拟器实例)" return {"running": False, "error": MOCK_HANDLE.last_error} config = uvicorn.Config(MOCK_APP, host="127.0.0.1", port=port, log_level="warning") server = uvicorn.Server(config) MOCK_HANDLE.server = server MOCK_HANDLE.task = asyncio.get_running_loop().create_task(server.serve()) # 轮询库的公开状态位,上限 ~2s:比固定 sleep 既不会过早返回未就绪,也不会白等。 for _ in range(40): await asyncio.sleep(0.05) if server.started: MOCK_HANDLE.last_error = "" return {"running": True, "port": port} if server.should_exit or MOCK_HANDLE.task.done(): break MOCK_HANDLE.last_error = f"Mock API 在端口 {port} 上启动失败" failed_task = MOCK_HANDLE.task MOCK_HANDLE.server = None MOCK_HANDLE.task = None if failed_task and not failed_task.done(): failed_task.cancel() return {"running": False, "error": MOCK_HANDLE.last_error} async def _stop_mock_api() -> dict: if MOCK_HANDLE.server and MOCK_HANDLE.task and not MOCK_HANDLE.task.done(): MOCK_HANDLE.server.should_exit = True try: await asyncio.wait_for(MOCK_HANDLE.task, timeout=5) except asyncio.TimeoutError: MOCK_HANDLE.task.cancel() return {"running": False} return {"running": False, "alreadyStopped": True} def _mock_running() -> bool: return bool(MOCK_HANDLE.task and not MOCK_HANDLE.task.done() and MOCK_HANDLE.server and MOCK_HANDLE.server.started) # --------------------------------------------------------------------------- # 基础模型 # --------------------------------------------------------------------------- class RowsPayload(BaseModel): objectCode: str runId: str | None = None rows: list[dict] | None = None class InboundExecPayload(BaseModel): objectCode: str | None = None entityCode: str | None = None operation: str = "push" # schema|push|receipt|bulk|snapshot_open|snapshot_commit|digest runId: str | None = None rows: list[dict] | None = None contractVersion: str | None = None # 用例与操作参数 case: str | None = None # happy|missing-idem|tamper-body|expired-ts|nonce-replay|same-key-replay|same-key-diff-body|bad-signature syncBatchId: str | None = None snapshotId: str | None = None force: bool = False bizDate: str | None = None class ScenarioPayload(BaseModel): scenarioId: str runId: str | None = None channelOverride: dict[str, str] | None = None # --------------------------------------------------------------------------- # 公共辅助 # --------------------------------------------------------------------------- def _require_obj(object_code: str) -> dict: obj = sa.find_object(object_code) if not obj: raise HTTPException(status_code=404, detail=f"unknown objectCode: {object_code}") return obj def _build_inbound_client() -> InboundClient: base = CONFIG["aidopBaseUrl"] try: aidop_probe.assert_host_allowed(base, CONFIG["allowedTargetHosts"]) except aidop_probe.HostNotAllowedError as ex: raise HTTPException(status_code=400, detail=str(ex)) access_key = _env("AIDOP_SIM_INBOUND_ACCESS_KEY") secret = _env("AIDOP_SIM_INBOUND_SECRET") missing = [n for n, v in (("AIDOP_SIM_INBOUND_ACCESS_KEY", access_key), ("AIDOP_SIM_INBOUND_SECRET", secret)) if not v] if missing: raise HTTPException(status_code=400, detail=f"缺少环境变量: {', '.join(missing)}(Secret 只从环境变量读取,不落盘)") return InboundClient(base, access_key, secret) def _resolve_entity(payload: InboundExecPayload) -> tuple[dict | None, str]: """解析本次执行的对象与 entityCode。仅给 entityCode 时允许无目录对象直发(rows 必填)。""" if payload.objectCode: obj = _require_obj(payload.objectCode) else: obj = _require_obj_by_entity(payload.entityCode or "") if payload.entityCode: entity = payload.entityCode.upper() else: if not obj: raise HTTPException(status_code=400, detail="need objectCode or entityCode") inbound = obj.get("apiInbound") or {} if not inbound.get("supported"): raise HTTPException(status_code=400, detail=f"{obj['objectCode']} 不支持 API_INBOUND: {inbound.get('reason', '无契约')}") entity = inbound["entityCode"] return obj, entity def _require_obj_by_entity(entity_code: str) -> dict | None: for obj in sa.load_catalog()["objects"]: if (obj.get("apiInbound") or {}).get("entityCode") == entity_code.upper(): return obj return None def _layers(code: str, config: str, data: str) -> dict: return {"code": code, "config": config, "data": data} def _followup_sql(obj: dict, channel: str, run_id: str) -> list[str]: prefix = sa.biz_prefix(run_id) db = (obj.get("dbSync") or {}).get("table") sqls = [ "SELECT source_code, status, COUNT(*) cnt, MAX(end_time) latest FROM mdp_sync_log " "WHERE source_code LIKE 'SIM_MODE3%' GROUP BY source_code, status;", "SELECT COUNT(*) FROM mdp_inbound_request WHERE create_time >= CURDATE();", ] if channel == "DB_SYNC" and db: sqls.insert(0, f"-- 模拟库核对:\nSELECT COUNT(*), MAX(sourceUpdatedAt) FROM `{db}` WHERE bizKey LIKE '{prefix}%';") sqls.append(f"SELECT * FROM `{db.replace('sim_', 'mdp_stg_sim_')}` WHERE raw_data LIKE '%{prefix}%' LIMIT 20;") if channel == "API_INBOUND": entity = (obj.get("apiInbound") or {}).get("entityCode") sqls.append(f"SELECT status, accepted_rows, rejected_rows, sync_batch_id FROM mdp_inbound_request " f"WHERE entity_code='{entity}' ORDER BY id DESC LIMIT 10;") sqls.append("-- 标准层/KPI 下游核验属第二阶段自动对账范围;第一阶段报告 TRANSFORM_PENDING,不猜测成功。") return sqls # --------------------------------------------------------------------------- # 状态与目录 # --------------------------------------------------------------------------- @app.get("/api/status") def status(): inbound_key = _env("AIDOP_SIM_INBOUND_ACCESS_KEY") return { "controlPort": DEFAULT_CONTROL_PORT, "aidopBaseUrl": CONFIG["aidopBaseUrl"], "allowedTargetHosts": CONFIG["allowedTargetHosts"], "aidopReachable": None, # 由 /api/precheck 实测,避免每次状态查询打真实请求 "mockApi": {"running": _mock_running(), "port": CONFIG["mockApiPort"], "authMode": MOCK_STATE.settings.auth_mode, "forceStatus": MOCK_STATE.settings.force_status, "lastError": MOCK_HANDLE.last_error}, "mysql": { "database": _env("AIDOP_SIM_MYSQL_DATABASE") or mock_db.SIM_DB_DEFAULT, "passwordLoaded": bool(_env("AIDOP_SIM_MYSQL_PASSWORD")), }, "inbound": { "accessKeyLoaded": bool(inbound_key), "accessKeyMasked": mask_access_key(inbound_key), "secretLoaded": bool(_env("AIDOP_SIM_INBOUND_SECRET")), }, "adminTokenLoaded": bool(_env("AIDOP_SIM_ADMIN_TOKEN")), "objectCounts": _channel_counts(), } def _channel_counts() -> dict: counts = {"DB_SYNC": 0, "API_PULL": 0, "API_INBOUND": 0, "total": 0} for obj in sa.load_catalog()["objects"]: counts["total"] += 1 if (obj.get("dbSync") or {}).get("supported"): counts["DB_SYNC"] += 1 if (obj.get("apiPull") or {}).get("supported"): counts["API_PULL"] += 1 if (obj.get("apiInbound") or {}).get("supported"): counts["API_INBOUND"] += 1 return counts @app.get("/api/catalog") def catalog(): return sa.load_catalog() @app.get("/api/scenarios") def scenarios(): return sa.load_scenarios() @app.get("/api/sample/{object_code}") def get_sample(object_code: str, runId: str | None = None): obj = _require_obj(object_code) run_id = runId or sa.new_run_id() inbound = obj.get("apiInbound") or {} result: dict = {"objectCode": obj["objectCode"], "runId": run_id, "canonical": sa.load_canonical_rows(obj), "dbSync": None, "apiPull": None, "apiInbound": None} if (obj.get("dbSync") or {}).get("supported"): result["dbSync"] = sa.adapt_db_rows(obj, run_id) if (obj.get("apiPull") or {}).get("supported"): result["apiPull"] = sa.adapt_api_rows(obj, run_id) if inbound.get("supported"): adapted = sa.adapt_inbound_rows(obj, run_id) result["apiInbound"] = { "entityCode": inbound["entityCode"], "contractVersion": inbound.get("contractVersion", "v1"), "rows": adapted["rows"], "missingRequired": adapted["missing_required"], "envelope": sa.inbound_envelope(adapted["rows"]), } return result # --------------------------------------------------------------------------- # Mock API 控制 # --------------------------------------------------------------------------- @app.post("/api/mock-api/start") async def mock_api_start(): return await _start_mock_api() @app.post("/api/mock-api/stop") async def mock_api_stop(): return await _stop_mock_api() @app.get("/api/mock-api/settings") def mock_api_settings(): s = MOCK_STATE.settings return {"running": _mock_running(), "authMode": s.auth_mode, "delaySeconds": s.delay_seconds, "forceStatus": s.force_status, "responsePath": s.response_path, "apiKeyHeader": s.api_key_header, "overrides": sorted(MOCK_STATE.overrides.keys())} @app.post("/api/mock-api/control") def mock_api_control(payload: dict): """统一控制入口:auth_mode/token/basic/apikey/delay/force_status/response_path/override。""" s = MOCK_STATE.settings if "auth_mode" in payload: mode = str(payload["auth_mode"]).upper() if mode not in ("NONE", "TOKEN", "BASIC", "APIKEY", "OAUTH2"): raise HTTPException(status_code=400, detail=f"unknown auth_mode: {mode}") s.auth_mode = mode for key, attr in (("token", "token"), ("basic_user", "basic_user"), ("basic_password", "basic_password"), ("api_key_header", "api_key_header"), ("api_key_value", "api_key_value"), ("response_path", "response_path")): if payload.get(key): setattr(s, attr, str(payload[key])) if "delay_seconds" in payload: s.delay_seconds = max(0.0, float(payload["delay_seconds"])) if "force_status" in payload: s.force_status = int(payload["force_status"]) if payload["force_status"] else None if payload.get("override_path") and isinstance(payload.get("override_rows"), list): MOCK_STATE.overrides[payload["override_path"]] = payload["override_rows"] if payload.get("clear_overrides"): MOCK_STATE.overrides.clear() return {"ok": True, "settings": mock_api_settings()} @app.get("/api/mock-api/audit") def mock_api_audit(): return {"entries": redact_obj(MOCK_STATE.audit[-50:], _secrets())} # --------------------------------------------------------------------------- # 模拟 MySQL 控制 # --------------------------------------------------------------------------- @app.post("/api/db/ensure") def db_ensure(payload: RowsPayload): """初始化模拟 schema 并确保对象源表存在(DDL 由样例推断)。""" obj = _require_obj(payload.objectCode) db_spec = obj.get("dbSync") or {} if not db_spec.get("supported"): raise HTTPException(status_code=400, detail=f"{obj['objectCode']} 不支持 DB_SYNC") if not _env("AIDOP_SIM_MYSQL_PASSWORD"): raise HTTPException(status_code=400, detail="缺少环境变量 AIDOP_SIM_MYSQL_PASSWORD") run_id = payload.runId or sa.new_run_id() rows = payload.rows if payload.rows is not None else sa.adapt_db_rows(obj, run_id) try: result = mock_db.init_schema({db_spec["table"]: { "rows": rows, "businessKey": db_spec.get("businessKey"), "incrementColumn": db_spec.get("incrementColumn", "sourceUpdatedAt")}}) return {"ok": True, "runId": run_id, **result} except mock_db.BlockedDatabaseError as ex: raise HTTPException(status_code=400, detail=str(ex)) except Exception as ex: raise HTTPException(status_code=502, detail=f"MySQL 不可达或执行失败: {type(ex).__name__}: {ex}") @app.post("/api/db/load") def db_load(payload: RowsPayload): obj = _require_obj(payload.objectCode) db_spec = obj.get("dbSync") or {} if not db_spec.get("supported"): raise HTTPException(status_code=400, detail=f"{obj['objectCode']} 不支持 DB_SYNC") run_id = payload.runId or sa.new_run_id() rows = payload.rows if payload.rows is not None else sa.adapt_db_rows(obj, run_id) try: count = mock_db.upsert_rows(db_spec["table"], rows) stats = mock_db.table_stats(db_spec["table"], db_spec.get("incrementColumn", "sourceUpdatedAt")) except Exception as ex: raise HTTPException(status_code=502, detail=f"装载失败(先执行 /api/db/ensure?): {type(ex).__name__}: {ex}") layers = _layers("SUPPORTED", "READY", "ACCEPTED") RUNS.add(run_id, obj["objectCode"], "DB_SYNC", {"action": "load", "table": db_spec["table"], "rowCount": count}, stats, layers, _followup_sql(obj, "DB_SYNC", run_id), gaps=["Ai-DOP 侧同步须由 MdpDbPullExecutor 真实执行;本操作只装载第三方模拟源库"]) return {"ok": True, "runId": run_id, "loaded": count, "stats": stats, "layers": layers} @app.post("/api/db/reset") def db_reset(payload: RowsPayload): obj = _require_obj(payload.objectCode) table = (obj.get("dbSync") or {}).get("table") if not table: raise HTTPException(status_code=400, detail="object has no dbSync table") try: mock_db.reset_table(table) return {"ok": True, "table": table} except Exception as ex: raise HTTPException(status_code=502, detail=f"{type(ex).__name__}: {ex}") @app.post("/api/db/update") def db_update(payload: dict): obj = _require_obj(payload.get("objectCode", "")) db_spec = obj.get("dbSync") or {} table = db_spec.get("table") updates = payload.get("updates") or {} where = payload.get("where") or {} if not table or not updates or not where: raise HTTPException(status_code=400, detail="need objectCode/updates/where") try: affected = mock_db.update_rows(table, updates, where, db_spec.get("incrementColumn", "sourceUpdatedAt"), payload.get("newIncrement")) return {"ok": True, "table": table, "affected": affected} except Exception as ex: raise HTTPException(status_code=502, detail=f"{type(ex).__name__}: {ex}") @app.get("/api/db/stats/{object_code}") def db_stats(object_code: str): obj = _require_obj(object_code) db_spec = obj.get("dbSync") or {} if not db_spec.get("table"): raise HTTPException(status_code=400, detail="object has no dbSync table") try: return mock_db.table_stats(db_spec["table"], db_spec.get("incrementColumn", "sourceUpdatedAt")) except Exception as ex: raise HTTPException(status_code=502, detail=f"{type(ex).__name__}: {ex}") @app.get("/api/db/aidop-config/{object_code}") def db_aidop_config(object_code: str): obj = _require_obj(object_code) base = f"http://127.0.0.1:{CONFIG['mockApiPort']}" return mock_db.generate_aidop_config(obj, base) # --------------------------------------------------------------------------- # 预检与三层报告(P1-G) # --------------------------------------------------------------------------- @app.get("/api/precheck") def precheck(): base = CONFIG["aidopBaseUrl"] reach = aidop_probe.check_reachable(base, CONFIG["allowedTargetHosts"]) inbound_ready = bool(_env("AIDOP_SIM_INBOUND_ACCESS_KEY") and _env("AIDOP_SIM_INBOUND_SECRET")) env_missing = [n for n in ("AIDOP_SIM_INBOUND_ACCESS_KEY", "AIDOP_SIM_INBOUND_SECRET", "AIDOP_SIM_MYSQL_PASSWORD") if not _env(n)] return { "aidopBaseUrl": base, "aidop": reach, "mockApi": {"running": _mock_running(), "port": CONFIG["mockApiPort"]}, "mysql": {"database": _env("AIDOP_SIM_MYSQL_DATABASE") or mock_db.SIM_DB_DEFAULT, "passwordLoaded": bool(_env("AIDOP_SIM_MYSQL_PASSWORD"))}, "inboundCredentialsLoaded": inbound_ready, "envMissing": env_missing, "testSources": mock_db.SIM_SOURCES, "note": "DB_SYNC/API_PULL 的来源、实体、租户登记须人工在 Ai-DOP 配置页核对(生成清单见 /api/db/aidop-config/{objectCode})", } @app.get("/api/precheck/inbound") def precheck_inbound(): """26 个代码契约逐一 schema 预检(真实签名,只读)。""" client = _build_inbound_client() codes = sorted({(o.get("apiInbound") or {}).get("entityCode") for o in sa.load_catalog()["objects"] if (o.get("apiInbound") or {}).get("supported")}) results = aidop_probe.precheck_all_inbound(client, codes, CONFIG["allowedTargetHosts"]) return {"accessKey": mask_access_key(client.key), "count": len(results), "results": results} @app.get("/api/precheck/inbound/{entity_code}") def precheck_inbound_one(entity_code: str): client = _build_inbound_client() return aidop_probe.precheck_inbound_entity(client, entity_code.upper(), CONFIG["allowedTargetHosts"]) # --------------------------------------------------------------------------- # API_INBOUND 执行(P1-E) # --------------------------------------------------------------------------- @app.post("/api/inbound/exec") def inbound_exec(payload: InboundExecPayload): obj, entity = _resolve_entity(payload) client = _build_inbound_client() run_id = payload.runId or sa.new_run_id() contract_version = payload.contractVersion if obj is None: # 仅按 entityCode 直发(无目录样例),rows 必填 rows = payload.rows if rows is None and payload.operation in ("push", "bulk"): raise HTTPException(status_code=400, detail="entityCode-only exec requires rows") envelope = sa.inbound_envelope(rows or []) else: inbound = obj.get("apiInbound") or {} contract_version = contract_version or inbound.get("contractVersion") rows = payload.rows if payload.rows is not None else sa.adapt_inbound_rows(obj, run_id)["rows"] envelope = sa.inbound_envelope(rows) object_code = obj["objectCode"] if obj else entity op = payload.operation result: dict = {"runId": run_id, "entityCode": entity, "operation": op, "case": payload.case} def record(status: int, body, headers, layers: dict, extra: dict | None = None): result.update({"httpStatus": status, "response": body, "requestHeaders": headers, **(extra or {})}) RUNS.add(run_id, object_code, "API_INBOUND", {"operation": op, "case": payload.case, "entityCode": entity, "rowCount": len(rows or []), "contractVersion": contract_version}, {"httpStatus": status, "body": body, "headers": headers}, layers, _followup_sql(obj or {"objectCode": entity, "apiInbound": {"entityCode": entity}}, "API_INBOUND", run_id)) return result # 预检(配置层) pre = aidop_probe.precheck_inbound_entity(client, entity, CONFIG["allowedTargetHosts"], contract_version) config_state = pre["configState"] result["precheck"] = {"configState": config_state, "httpStatus": pre["httpStatus"]} if op == "schema": status, body, headers = client.op_schema(entity, contract_version) layers = _layers("SUPPORTED", aidop_probe.classify_inbound_precheck(status, body), "NOT_RUN") return record(status, body, headers, layers) if config_state != "READY" and op in ("push", "bulk", "snapshot_open"): layers = _layers("SUPPORTED", config_state, "NOT_RUN") raise HTTPException(status_code=409, detail={ "message": f"配置层 {config_state}:普通推数只对 READY 实体执行(不伪造成功)", "precheck": pre, "layers": layers}) if op == "push": case = payload.case or "happy" if case == "happy": status, body, headers = client.op_push(entity, envelope, contract_version=contract_version) elif case == "missing-idem": status, body, headers = client.case_missing_idem(entity, envelope) elif case == "tamper-body": status, body, headers = client.case_tamper_body(entity, envelope) elif case == "expired-ts": status, body, headers = client.case_expired_ts(entity, envelope) elif case == "bad-signature": status, body, headers = client.case_bad_signature(entity, envelope) elif case == "nonce-replay": (s1, d1), (s2, d2) = client.case_nonce_replay(entity, envelope) status, body, headers = s2, {"first": [s1, d1], "replay": [s2, d2]}, {} elif case == "same-key-replay": (s1, d1), (s2, d2) = client.case_same_key_replay(entity, envelope) status, body, headers = s2, {"first": [s1, d1], "replay": [s2, d2]}, {} elif case == "same-key-diff-body": env2 = sa.inbound_envelope([{**r, "sourceUpdatedAt": sa.now_iso()} for r in (rows or [])]) (s1, d1), (s2, d2) = client.case_same_key_diff_body(entity, envelope, env2) status, body, headers = s2, {"first": [s1, d1], "conflict": [s2, d2]}, {} else: raise HTTPException(status_code=400, detail=f"unknown case: {case}") data_state = "ACCEPTED" if status == 202 else ("RUN_FAILED" if status >= 400 else "NOT_RUN") layers = _layers("SUPPORTED", config_state, data_state) return record(status, body, headers, layers, extra={"gaps": ["HTTP 202 仅代表摄取受理;贴源/标准层核验见 followupSql(TRANSFORM_PENDING)"]}) if op == "receipt": if not payload.syncBatchId: raise HTTPException(status_code=400, detail="receipt requires syncBatchId") status, body, headers = client.op_receipt(payload.syncBatchId) layers = _layers("SUPPORTED", config_state, "ACCEPTED" if status == 200 else "RUN_FAILED") return record(status, body, headers, layers) if op == "bulk": status, body, headers = client.op_bulk(entity, rows or [], contract_version=contract_version) layers = _layers("SUPPORTED", config_state, "ACCEPTED" if status == 202 else "RUN_FAILED") return record(status, body, headers, layers) if op == "snapshot_open": status, body, headers = client.op_snapshot_open(entity) layers = _layers("SUPPORTED", config_state, "ACCEPTED" if status in (200, 201) else "RUN_FAILED") return record(status, body, headers, layers) if op == "snapshot_commit": if not payload.snapshotId: raise HTTPException(status_code=400, detail="snapshot_commit requires snapshotId") status, body, headers = client.op_snapshot_commit(entity, payload.snapshotId, payload.force) layers = _layers("SUPPORTED", config_state, "ACCEPTED" if status in (200, 202) else "RUN_FAILED") return record(status, body, headers, layers) if op == "digest": biz_date = payload.bizDate or time.strftime("%Y-%m-%d") status, body, headers = client.op_digest(entity, biz_date) layers = _layers("SUPPORTED", config_state, "ACCEPTED" if status == 200 else "RUN_FAILED") return record(status, body, headers, layers) raise HTTPException(status_code=400, detail=f"unknown operation: {op}") # --------------------------------------------------------------------------- # 场景执行(P1-F) # --------------------------------------------------------------------------- @app.post("/api/scenarios/run") def scenario_run(payload: ScenarioPayload): scenario = next((s for s in sa.load_scenarios()["scenarios"] if s["scenarioId"] == payload.scenarioId), None) if not scenario: raise HTTPException(status_code=404, detail=f"unknown scenario: {payload.scenarioId}") run_id = payload.runId or sa.new_run_id() overrides = payload.channelOverride or {} nodes_out = [] for node in scenario["nodes"]: obj = sa.find_object(node["objectCode"]) channel = overrides.get(node["objectCode"], node["channel"]) node_result = {"objectCode": node["objectCode"], "channel": channel, "relationKey": node.get("relationKey")} if not obj: node_result.update({"state": "CODE_UNSUPPORTED", "message": "目录中无此对象"}) nodes_out.append(node_result) continue spec = {"DB_SYNC": obj.get("dbSync"), "API_PULL": obj.get("apiPull"), "API_INBOUND": obj.get("apiInbound")}.get(channel) if not spec or not spec.get("supported"): reason = (spec or {}).get("reason") or node.get("inboundNote") or "该通道不支持" node_result.update({"state": "CODE_UNSUPPORTED", "message": reason}) nodes_out.append(node_result) continue try: if channel == "DB_SYNC": rows = sa.adapt_db_rows(obj, run_id) if not _env("AIDOP_SIM_MYSQL_PASSWORD"): raise RuntimeError("缺少 AIDOP_SIM_MYSQL_PASSWORD") mock_db.init_schema({spec["table"]: {"rows": rows, "businessKey": spec.get("businessKey"), "incrementColumn": spec.get("incrementColumn")}}) count = mock_db.upsert_rows(spec["table"], rows) node_result.update({"state": "ACCEPTED", "loaded": count, "table": spec["table"]}) elif channel == "API_PULL": if not _mock_running(): raise RuntimeError("Mock API 未启动(先 POST /api/mock-api/start)") rows = sa.adapt_api_rows(obj, run_id) MOCK_STATE.overrides[spec["path"]] = rows node_result.update({"state": "ACCEPTED", "overridePath": spec["path"], "rowCount": len(rows), "note": "已在 Mock API 挂载适配样例;由 Ai-DOP MdpApiPullExecutor 真实拉取"}) else: client = _build_inbound_client() entity = spec["entityCode"] version = node.get("contractVersion") or spec.get("contractVersion") pre = aidop_probe.precheck_inbound_entity(client, entity, CONFIG["allowedTargetHosts"], version) if pre["configState"] != "READY": node_result.update({"state": pre["configState"], "message": f"配置层 {pre['configState']},未推数(不伪造成功)"}) else: rows = sa.adapt_inbound_rows(obj, run_id)["rows"] status, body, _ = client.op_push(entity, sa.inbound_envelope(rows), contract_version=version) node_result.update({"state": "ACCEPTED" if status == 202 else "RUN_FAILED", "httpStatus": status, "response": body}) except HTTPException: raise except aidop_probe.HostNotAllowedError as ex: node_result.update({"state": "SOURCE_UNREACHABLE", "message": str(ex)}) except Exception as ex: node_result.update({"state": "RUN_FAILED", "message": f"{type(ex).__name__}: {ex}"}) nodes_out.append(node_result) unsupported = [n["objectCode"] for n in nodes_out if n["state"] == "CODE_UNSUPPORTED"] layers = _layers("SUPPORTED", "READY" if all(n["state"] not in ("SOURCE_UNREACHABLE",) for n in nodes_out) else "SOURCE_UNREACHABLE", "ACCEPTED" if all(n["state"] in ("ACCEPTED", "CODE_UNSUPPORTED") for n in nodes_out) else "RUN_FAILED") entry = RUNS.add(run_id, ",".join(n["objectCode"] for n in nodes_out), f"SCENARIO:{payload.scenarioId}", {"scenario": scenario["name"]}, {"nodes": nodes_out}, layers, _followup_sql({"objectCode": "SCENARIO", "dbSync": {}}, "SCENARIO", run_id), gaps=([f"以下节点该通道不支持:{', '.join(unsupported)}(链路不完整,不得宣称全链路跑通)"] if unsupported else []), scenario_id=payload.scenarioId, nodes=nodes_out) return entry # --------------------------------------------------------------------------- # 运行记录 # --------------------------------------------------------------------------- @app.get("/api/runs") def runs(): return {"runs": RUNS.recent()} # --------------------------------------------------------------------------- # 静态页面 # --------------------------------------------------------------------------- @app.get("/") def index(): return FileResponse(SIM_ROOT / "static" / "index.html") app.mount("/static", StaticFiles(directory=SIM_ROOT / "static"), name="static") def main() -> None: print(f"[simulator] control panel: http://127.0.0.1:{DEFAULT_CONTROL_PORT}/") print(f"[simulator] aidop target : {CONFIG['aidopBaseUrl']} (allowed hosts: {CONFIG['allowedTargetHosts']})") uvicorn.run(app, host="127.0.0.1", port=DEFAULT_CONTROL_PORT, log_level="info") if __name__ == "__main__": main()