app.py 34 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783
  1. """Ai-DOP 第三方对接模拟器 · 本地控制服务(任务书 P1-B / P1-G)。
  2. - 只绑定 127.0.0.1,页面与 API 仅本机可用;
  3. - Secret/密码/Token 只从环境变量读取,任何 API 不回显;
  4. - 所有对外目标 host 走白名单硬阻断(默认拒绝);
  5. - 运行记录(RunStore)写入前统一脱敏;
  6. - 内嵌管理 Mock API(providers.mock_api),默认端口 8018,可启停/切鉴权/注入错误;
  7. - 三层报告:代码层 SUPPORTED/UNSUPPORTED、配置层 READY/NOT_*、数据层 NOT_RUN/ACCEPTED/TRANSFORM_PENDING/RUN_FAILED。
  8. 启动:python tools/integration-simulator/app.py (详见 README.md)
  9. """
  10. from __future__ import annotations
  11. import asyncio
  12. import os
  13. import socket
  14. import sys
  15. import time
  16. import uuid
  17. from pathlib import Path
  18. SIM_ROOT = Path(__file__).resolve().parent
  19. if str(SIM_ROOT) not in sys.path:
  20. sys.path.insert(0, str(SIM_ROOT))
  21. import uvicorn
  22. from fastapi import FastAPI, HTTPException
  23. from fastapi.responses import FileResponse
  24. from fastapi.staticfiles import StaticFiles
  25. from pydantic import BaseModel
  26. from catalog import sample_adapters as sa
  27. from clients import aidop_probe
  28. from clients.inbound_client import InboundClient
  29. from clients.redaction import mask_access_key, redact_obj
  30. from providers import mock_api, mock_db
  31. DEFAULT_CONTROL_PORT = 8017
  32. MAX_RUNS = 200
  33. # 状态枚举(任务书 §4)
  34. CONFIG_STATES = ("READY", "CONFIG_NOT_REGISTERED", "CONFIG_NOT_ENABLED",
  35. "GRANT_MISSING", "SOURCE_UNREACHABLE", "RUN_FAILED")
  36. def _env(name: str) -> str:
  37. return os.environ.get(name, "") or ""
  38. def _secrets() -> list[str]:
  39. return [v for v in (_env("AIDOP_SIM_INBOUND_SECRET"), _env("AIDOP_SIM_MYSQL_PASSWORD"),
  40. _env("AIDOP_SIM_ADMIN_TOKEN")) if v]
  41. def load_config() -> dict:
  42. cfg = {
  43. "aidopBaseUrl": "http://127.0.0.1:5005",
  44. "mockApiPort": 8018,
  45. "mysqlDatabase": mock_db.SIM_DB_DEFAULT,
  46. "allowedTargetHosts": ["127.0.0.1", "localhost"],
  47. }
  48. cfg_file = SIM_ROOT / "config.json"
  49. if cfg_file.exists():
  50. import json
  51. with cfg_file.open("r", encoding="utf-8") as f:
  52. cfg.update(json.load(f))
  53. return cfg
  54. CONFIG = load_config()
  55. MOCK_STATE = mock_api.MockApiState()
  56. MOCK_APP = mock_api.create_app(MOCK_STATE)
  57. class RunStore:
  58. def __init__(self, limit: int = MAX_RUNS):
  59. self.limit = limit
  60. self.runs: list[dict] = []
  61. def add(self, run_id: str, object_code: str, channel: str, request_summary: dict,
  62. response: object, layers: dict, followup_sql: list[str] | None = None,
  63. gaps: list[str] | None = None, scenario_id: str | None = None,
  64. nodes: list[dict] | None = None) -> dict:
  65. entry = {
  66. "runId": run_id,
  67. "time": time.strftime("%Y-%m-%dT%H:%M:%S%z"),
  68. "objectCode": object_code,
  69. "channel": channel,
  70. "scenarioId": scenario_id,
  71. "nodes": nodes,
  72. "request": redact_obj(request_summary, _secrets()),
  73. "response": redact_obj(response, _secrets()),
  74. "layers": layers,
  75. "followupSql": followup_sql or [],
  76. "knownGaps": gaps or [],
  77. }
  78. self.runs.append(entry)
  79. if len(self.runs) > self.limit:
  80. del self.runs[: len(self.runs) - self.limit]
  81. return entry
  82. def recent(self) -> list[dict]:
  83. return list(reversed(self.runs[-50:]))
  84. RUNS = RunStore()
  85. app = FastAPI(title="Ai-DOP Integration Simulator", version="1.0.0")
  86. # ---------------------------------------------------------------------------
  87. # Mock API 内嵌启停
  88. # ---------------------------------------------------------------------------
  89. class _MockServerHandle:
  90. def __init__(self):
  91. self.server: uvicorn.Server | None = None
  92. self.task: asyncio.Task | None = None
  93. self.last_error: str = ""
  94. MOCK_HANDLE = _MockServerHandle()
  95. def _port_in_use(port: int) -> bool:
  96. """真实探测本机端口是否已被占用。
  97. 不依赖 uvicorn 内部属性:0.30+ 的 Server 已无 startup_error,
  98. 读它会抛 AttributeError 并把已启动的服务误报成 500。
  99. """
  100. with socket.socket(socket.AF_INET, socket.SOCK_STREAM) as probe:
  101. probe.settimeout(0.3)
  102. return probe.connect_ex(("127.0.0.1", int(port))) == 0
  103. async def _start_mock_api() -> dict:
  104. port = int(CONFIG["mockApiPort"])
  105. if MOCK_HANDLE.task and not MOCK_HANDLE.task.done():
  106. return {"running": True, "port": port, "alreadyRunning": True}
  107. # 端口预检必须在建 task 之前:否则失败路径会留下一个已赋值的句柄,
  108. # 使 _mock_running() 与实际监听状态不一致。
  109. if _port_in_use(port):
  110. MOCK_HANDLE.last_error = f"端口 {port} 已被占用(可能是另一个模拟器实例)"
  111. return {"running": False, "error": MOCK_HANDLE.last_error}
  112. config = uvicorn.Config(MOCK_APP, host="127.0.0.1", port=port, log_level="warning")
  113. server = uvicorn.Server(config)
  114. MOCK_HANDLE.server = server
  115. MOCK_HANDLE.task = asyncio.get_running_loop().create_task(server.serve())
  116. # 轮询库的公开状态位,上限 ~2s:比固定 sleep 既不会过早返回未就绪,也不会白等。
  117. for _ in range(40):
  118. await asyncio.sleep(0.05)
  119. if server.started:
  120. MOCK_HANDLE.last_error = ""
  121. return {"running": True, "port": port}
  122. if server.should_exit or MOCK_HANDLE.task.done():
  123. break
  124. MOCK_HANDLE.last_error = f"Mock API 在端口 {port} 上启动失败"
  125. failed_task = MOCK_HANDLE.task
  126. MOCK_HANDLE.server = None
  127. MOCK_HANDLE.task = None
  128. if failed_task and not failed_task.done():
  129. failed_task.cancel()
  130. return {"running": False, "error": MOCK_HANDLE.last_error}
  131. async def _stop_mock_api() -> dict:
  132. if MOCK_HANDLE.server and MOCK_HANDLE.task and not MOCK_HANDLE.task.done():
  133. MOCK_HANDLE.server.should_exit = True
  134. try:
  135. await asyncio.wait_for(MOCK_HANDLE.task, timeout=5)
  136. except asyncio.TimeoutError:
  137. MOCK_HANDLE.task.cancel()
  138. return {"running": False}
  139. return {"running": False, "alreadyStopped": True}
  140. def _mock_running() -> bool:
  141. return bool(MOCK_HANDLE.task and not MOCK_HANDLE.task.done() and
  142. MOCK_HANDLE.server and MOCK_HANDLE.server.started)
  143. # ---------------------------------------------------------------------------
  144. # 基础模型
  145. # ---------------------------------------------------------------------------
  146. class RowsPayload(BaseModel):
  147. objectCode: str
  148. runId: str | None = None
  149. rows: list[dict] | None = None
  150. class InboundExecPayload(BaseModel):
  151. objectCode: str | None = None
  152. entityCode: str | None = None
  153. operation: str = "push" # schema|push|receipt|bulk|snapshot_open|snapshot_commit|digest
  154. runId: str | None = None
  155. rows: list[dict] | None = None
  156. contractVersion: str | None = None
  157. # 用例与操作参数
  158. case: str | None = None # happy|missing-idem|tamper-body|expired-ts|nonce-replay|same-key-replay|same-key-diff-body|bad-signature
  159. syncBatchId: str | None = None
  160. snapshotId: str | None = None
  161. force: bool = False
  162. bizDate: str | None = None
  163. class ScenarioPayload(BaseModel):
  164. scenarioId: str
  165. runId: str | None = None
  166. channelOverride: dict[str, str] | None = None
  167. # ---------------------------------------------------------------------------
  168. # 公共辅助
  169. # ---------------------------------------------------------------------------
  170. def _require_obj(object_code: str) -> dict:
  171. obj = sa.find_object(object_code)
  172. if not obj:
  173. raise HTTPException(status_code=404, detail=f"unknown objectCode: {object_code}")
  174. return obj
  175. def _build_inbound_client() -> InboundClient:
  176. base = CONFIG["aidopBaseUrl"]
  177. try:
  178. aidop_probe.assert_host_allowed(base, CONFIG["allowedTargetHosts"])
  179. except aidop_probe.HostNotAllowedError as ex:
  180. raise HTTPException(status_code=400, detail=str(ex))
  181. access_key = _env("AIDOP_SIM_INBOUND_ACCESS_KEY")
  182. secret = _env("AIDOP_SIM_INBOUND_SECRET")
  183. missing = [n for n, v in (("AIDOP_SIM_INBOUND_ACCESS_KEY", access_key),
  184. ("AIDOP_SIM_INBOUND_SECRET", secret)) if not v]
  185. if missing:
  186. raise HTTPException(status_code=400,
  187. detail=f"缺少环境变量: {', '.join(missing)}(Secret 只从环境变量读取,不落盘)")
  188. return InboundClient(base, access_key, secret)
  189. def _resolve_entity(payload: InboundExecPayload) -> tuple[dict | None, str]:
  190. """解析本次执行的对象与 entityCode。仅给 entityCode 时允许无目录对象直发(rows 必填)。"""
  191. if payload.objectCode:
  192. obj = _require_obj(payload.objectCode)
  193. else:
  194. obj = _require_obj_by_entity(payload.entityCode or "")
  195. if payload.entityCode:
  196. entity = payload.entityCode.upper()
  197. else:
  198. if not obj:
  199. raise HTTPException(status_code=400, detail="need objectCode or entityCode")
  200. inbound = obj.get("apiInbound") or {}
  201. if not inbound.get("supported"):
  202. raise HTTPException(status_code=400,
  203. detail=f"{obj['objectCode']} 不支持 API_INBOUND: {inbound.get('reason', '无契约')}")
  204. entity = inbound["entityCode"]
  205. return obj, entity
  206. def _require_obj_by_entity(entity_code: str) -> dict | None:
  207. for obj in sa.load_catalog()["objects"]:
  208. if (obj.get("apiInbound") or {}).get("entityCode") == entity_code.upper():
  209. return obj
  210. return None
  211. def _layers(code: str, config: str, data: str) -> dict:
  212. return {"code": code, "config": config, "data": data}
  213. def _followup_sql(obj: dict, channel: str, run_id: str) -> list[str]:
  214. prefix = sa.biz_prefix(run_id)
  215. db = (obj.get("dbSync") or {}).get("table")
  216. sqls = [
  217. "SELECT source_code, status, COUNT(*) cnt, MAX(end_time) latest FROM mdp_sync_log "
  218. "WHERE source_code LIKE 'SIM_MODE3%' GROUP BY source_code, status;",
  219. "SELECT COUNT(*) FROM mdp_inbound_request WHERE create_time >= CURDATE();",
  220. ]
  221. if channel == "DB_SYNC" and db:
  222. sqls.insert(0, f"-- 模拟库核对:\nSELECT COUNT(*), MAX(sourceUpdatedAt) FROM `{db}` WHERE bizKey LIKE '{prefix}%';")
  223. sqls.append(f"SELECT * FROM `{db.replace('sim_', 'mdp_stg_sim_')}` WHERE raw_data LIKE '%{prefix}%' LIMIT 20;")
  224. if channel == "API_INBOUND":
  225. entity = (obj.get("apiInbound") or {}).get("entityCode")
  226. sqls.append(f"SELECT status, accepted_rows, rejected_rows, sync_batch_id FROM mdp_inbound_request "
  227. f"WHERE entity_code='{entity}' ORDER BY id DESC LIMIT 10;")
  228. sqls.append("-- 标准层/KPI 下游核验属第二阶段自动对账范围;第一阶段报告 TRANSFORM_PENDING,不猜测成功。")
  229. return sqls
  230. # ---------------------------------------------------------------------------
  231. # 状态与目录
  232. # ---------------------------------------------------------------------------
  233. @app.get("/api/status")
  234. def status():
  235. inbound_key = _env("AIDOP_SIM_INBOUND_ACCESS_KEY")
  236. return {
  237. "controlPort": DEFAULT_CONTROL_PORT,
  238. "aidopBaseUrl": CONFIG["aidopBaseUrl"],
  239. "allowedTargetHosts": CONFIG["allowedTargetHosts"],
  240. "aidopReachable": None, # 由 /api/precheck 实测,避免每次状态查询打真实请求
  241. "mockApi": {"running": _mock_running(), "port": CONFIG["mockApiPort"],
  242. "authMode": MOCK_STATE.settings.auth_mode,
  243. "forceStatus": MOCK_STATE.settings.force_status,
  244. "lastError": MOCK_HANDLE.last_error},
  245. "mysql": {
  246. "database": _env("AIDOP_SIM_MYSQL_DATABASE") or mock_db.SIM_DB_DEFAULT,
  247. "passwordLoaded": bool(_env("AIDOP_SIM_MYSQL_PASSWORD")),
  248. },
  249. "inbound": {
  250. "accessKeyLoaded": bool(inbound_key),
  251. "accessKeyMasked": mask_access_key(inbound_key),
  252. "secretLoaded": bool(_env("AIDOP_SIM_INBOUND_SECRET")),
  253. },
  254. "adminTokenLoaded": bool(_env("AIDOP_SIM_ADMIN_TOKEN")),
  255. "objectCounts": _channel_counts(),
  256. }
  257. def _channel_counts() -> dict:
  258. counts = {"DB_SYNC": 0, "API_PULL": 0, "API_INBOUND": 0, "total": 0}
  259. for obj in sa.load_catalog()["objects"]:
  260. counts["total"] += 1
  261. if (obj.get("dbSync") or {}).get("supported"):
  262. counts["DB_SYNC"] += 1
  263. if (obj.get("apiPull") or {}).get("supported"):
  264. counts["API_PULL"] += 1
  265. if (obj.get("apiInbound") or {}).get("supported"):
  266. counts["API_INBOUND"] += 1
  267. return counts
  268. @app.get("/api/catalog")
  269. def catalog():
  270. return sa.load_catalog()
  271. @app.get("/api/scenarios")
  272. def scenarios():
  273. return sa.load_scenarios()
  274. @app.get("/api/sample/{object_code}")
  275. def get_sample(object_code: str, runId: str | None = None):
  276. obj = _require_obj(object_code)
  277. run_id = runId or sa.new_run_id()
  278. inbound = obj.get("apiInbound") or {}
  279. result: dict = {"objectCode": obj["objectCode"], "runId": run_id,
  280. "canonical": sa.load_canonical_rows(obj),
  281. "dbSync": None, "apiPull": None, "apiInbound": None}
  282. if (obj.get("dbSync") or {}).get("supported"):
  283. result["dbSync"] = sa.adapt_db_rows(obj, run_id)
  284. if (obj.get("apiPull") or {}).get("supported"):
  285. result["apiPull"] = sa.adapt_api_rows(obj, run_id)
  286. if inbound.get("supported"):
  287. adapted = sa.adapt_inbound_rows(obj, run_id)
  288. result["apiInbound"] = {
  289. "entityCode": inbound["entityCode"],
  290. "contractVersion": inbound.get("contractVersion", "v1"),
  291. "rows": adapted["rows"],
  292. "missingRequired": adapted["missing_required"],
  293. "envelope": sa.inbound_envelope(adapted["rows"]),
  294. }
  295. return result
  296. # ---------------------------------------------------------------------------
  297. # Mock API 控制
  298. # ---------------------------------------------------------------------------
  299. @app.post("/api/mock-api/start")
  300. async def mock_api_start():
  301. return await _start_mock_api()
  302. @app.post("/api/mock-api/stop")
  303. async def mock_api_stop():
  304. return await _stop_mock_api()
  305. @app.get("/api/mock-api/settings")
  306. def mock_api_settings():
  307. s = MOCK_STATE.settings
  308. return {"running": _mock_running(), "authMode": s.auth_mode, "delaySeconds": s.delay_seconds,
  309. "forceStatus": s.force_status, "responsePath": s.response_path,
  310. "apiKeyHeader": s.api_key_header, "overrides": sorted(MOCK_STATE.overrides.keys())}
  311. @app.post("/api/mock-api/control")
  312. def mock_api_control(payload: dict):
  313. """统一控制入口:auth_mode/token/basic/apikey/delay/force_status/response_path/override。"""
  314. s = MOCK_STATE.settings
  315. if "auth_mode" in payload:
  316. mode = str(payload["auth_mode"]).upper()
  317. if mode not in ("NONE", "TOKEN", "BASIC", "APIKEY", "OAUTH2"):
  318. raise HTTPException(status_code=400, detail=f"unknown auth_mode: {mode}")
  319. s.auth_mode = mode
  320. for key, attr in (("token", "token"), ("basic_user", "basic_user"),
  321. ("basic_password", "basic_password"), ("api_key_header", "api_key_header"),
  322. ("api_key_value", "api_key_value"), ("response_path", "response_path")):
  323. if payload.get(key):
  324. setattr(s, attr, str(payload[key]))
  325. if "delay_seconds" in payload:
  326. s.delay_seconds = max(0.0, float(payload["delay_seconds"]))
  327. if "force_status" in payload:
  328. s.force_status = int(payload["force_status"]) if payload["force_status"] else None
  329. if payload.get("override_path") and isinstance(payload.get("override_rows"), list):
  330. MOCK_STATE.overrides[payload["override_path"]] = payload["override_rows"]
  331. if payload.get("clear_overrides"):
  332. MOCK_STATE.overrides.clear()
  333. return {"ok": True, "settings": mock_api_settings()}
  334. @app.get("/api/mock-api/audit")
  335. def mock_api_audit():
  336. return {"entries": redact_obj(MOCK_STATE.audit[-50:], _secrets())}
  337. # ---------------------------------------------------------------------------
  338. # 模拟 MySQL 控制
  339. # ---------------------------------------------------------------------------
  340. @app.post("/api/db/ensure")
  341. def db_ensure(payload: RowsPayload):
  342. """初始化模拟 schema 并确保对象源表存在(DDL 由样例推断)。"""
  343. obj = _require_obj(payload.objectCode)
  344. db_spec = obj.get("dbSync") or {}
  345. if not db_spec.get("supported"):
  346. raise HTTPException(status_code=400, detail=f"{obj['objectCode']} 不支持 DB_SYNC")
  347. if not _env("AIDOP_SIM_MYSQL_PASSWORD"):
  348. raise HTTPException(status_code=400, detail="缺少环境变量 AIDOP_SIM_MYSQL_PASSWORD")
  349. run_id = payload.runId or sa.new_run_id()
  350. rows = payload.rows if payload.rows is not None else sa.adapt_db_rows(obj, run_id)
  351. try:
  352. result = mock_db.init_schema({db_spec["table"]: {
  353. "rows": rows, "businessKey": db_spec.get("businessKey"),
  354. "incrementColumn": db_spec.get("incrementColumn", "sourceUpdatedAt")}})
  355. return {"ok": True, "runId": run_id, **result}
  356. except mock_db.BlockedDatabaseError as ex:
  357. raise HTTPException(status_code=400, detail=str(ex))
  358. except Exception as ex:
  359. raise HTTPException(status_code=502, detail=f"MySQL 不可达或执行失败: {type(ex).__name__}: {ex}")
  360. @app.post("/api/db/load")
  361. def db_load(payload: RowsPayload):
  362. obj = _require_obj(payload.objectCode)
  363. db_spec = obj.get("dbSync") or {}
  364. if not db_spec.get("supported"):
  365. raise HTTPException(status_code=400, detail=f"{obj['objectCode']} 不支持 DB_SYNC")
  366. run_id = payload.runId or sa.new_run_id()
  367. rows = payload.rows if payload.rows is not None else sa.adapt_db_rows(obj, run_id)
  368. try:
  369. count = mock_db.upsert_rows(db_spec["table"], rows)
  370. stats = mock_db.table_stats(db_spec["table"], db_spec.get("incrementColumn", "sourceUpdatedAt"))
  371. except Exception as ex:
  372. raise HTTPException(status_code=502, detail=f"装载失败(先执行 /api/db/ensure?): {type(ex).__name__}: {ex}")
  373. layers = _layers("SUPPORTED", "READY", "ACCEPTED")
  374. RUNS.add(run_id, obj["objectCode"], "DB_SYNC",
  375. {"action": "load", "table": db_spec["table"], "rowCount": count},
  376. stats, layers, _followup_sql(obj, "DB_SYNC", run_id),
  377. gaps=["Ai-DOP 侧同步须由 MdpDbPullExecutor 真实执行;本操作只装载第三方模拟源库"])
  378. return {"ok": True, "runId": run_id, "loaded": count, "stats": stats, "layers": layers}
  379. @app.post("/api/db/reset")
  380. def db_reset(payload: RowsPayload):
  381. obj = _require_obj(payload.objectCode)
  382. table = (obj.get("dbSync") or {}).get("table")
  383. if not table:
  384. raise HTTPException(status_code=400, detail="object has no dbSync table")
  385. try:
  386. mock_db.reset_table(table)
  387. return {"ok": True, "table": table}
  388. except Exception as ex:
  389. raise HTTPException(status_code=502, detail=f"{type(ex).__name__}: {ex}")
  390. @app.post("/api/db/update")
  391. def db_update(payload: dict):
  392. obj = _require_obj(payload.get("objectCode", ""))
  393. db_spec = obj.get("dbSync") or {}
  394. table = db_spec.get("table")
  395. updates = payload.get("updates") or {}
  396. where = payload.get("where") or {}
  397. if not table or not updates or not where:
  398. raise HTTPException(status_code=400, detail="need objectCode/updates/where")
  399. try:
  400. affected = mock_db.update_rows(table, updates, where,
  401. db_spec.get("incrementColumn", "sourceUpdatedAt"),
  402. payload.get("newIncrement"))
  403. return {"ok": True, "table": table, "affected": affected}
  404. except Exception as ex:
  405. raise HTTPException(status_code=502, detail=f"{type(ex).__name__}: {ex}")
  406. @app.get("/api/db/stats/{object_code}")
  407. def db_stats(object_code: str):
  408. obj = _require_obj(object_code)
  409. db_spec = obj.get("dbSync") or {}
  410. if not db_spec.get("table"):
  411. raise HTTPException(status_code=400, detail="object has no dbSync table")
  412. try:
  413. return mock_db.table_stats(db_spec["table"], db_spec.get("incrementColumn", "sourceUpdatedAt"))
  414. except Exception as ex:
  415. raise HTTPException(status_code=502, detail=f"{type(ex).__name__}: {ex}")
  416. @app.get("/api/db/aidop-config/{object_code}")
  417. def db_aidop_config(object_code: str):
  418. obj = _require_obj(object_code)
  419. base = f"http://127.0.0.1:{CONFIG['mockApiPort']}"
  420. return mock_db.generate_aidop_config(obj, base)
  421. # ---------------------------------------------------------------------------
  422. # 预检与三层报告(P1-G)
  423. # ---------------------------------------------------------------------------
  424. @app.get("/api/precheck")
  425. def precheck():
  426. base = CONFIG["aidopBaseUrl"]
  427. reach = aidop_probe.check_reachable(base, CONFIG["allowedTargetHosts"])
  428. inbound_ready = bool(_env("AIDOP_SIM_INBOUND_ACCESS_KEY") and _env("AIDOP_SIM_INBOUND_SECRET"))
  429. env_missing = [n for n in ("AIDOP_SIM_INBOUND_ACCESS_KEY", "AIDOP_SIM_INBOUND_SECRET",
  430. "AIDOP_SIM_MYSQL_PASSWORD") if not _env(n)]
  431. return {
  432. "aidopBaseUrl": base,
  433. "aidop": reach,
  434. "mockApi": {"running": _mock_running(), "port": CONFIG["mockApiPort"]},
  435. "mysql": {"database": _env("AIDOP_SIM_MYSQL_DATABASE") or mock_db.SIM_DB_DEFAULT,
  436. "passwordLoaded": bool(_env("AIDOP_SIM_MYSQL_PASSWORD"))},
  437. "inboundCredentialsLoaded": inbound_ready,
  438. "envMissing": env_missing,
  439. "testSources": mock_db.SIM_SOURCES,
  440. "note": "DB_SYNC/API_PULL 的来源、实体、租户登记须人工在 Ai-DOP 配置页核对(生成清单见 /api/db/aidop-config/{objectCode})",
  441. }
  442. @app.get("/api/precheck/inbound")
  443. def precheck_inbound():
  444. """26 个代码契约逐一 schema 预检(真实签名,只读)。"""
  445. client = _build_inbound_client()
  446. codes = sorted({(o.get("apiInbound") or {}).get("entityCode")
  447. for o in sa.load_catalog()["objects"]
  448. if (o.get("apiInbound") or {}).get("supported")})
  449. results = aidop_probe.precheck_all_inbound(client, codes, CONFIG["allowedTargetHosts"])
  450. return {"accessKey": mask_access_key(client.key), "count": len(results), "results": results}
  451. @app.get("/api/precheck/inbound/{entity_code}")
  452. def precheck_inbound_one(entity_code: str):
  453. client = _build_inbound_client()
  454. return aidop_probe.precheck_inbound_entity(client, entity_code.upper(), CONFIG["allowedTargetHosts"])
  455. # ---------------------------------------------------------------------------
  456. # API_INBOUND 执行(P1-E)
  457. # ---------------------------------------------------------------------------
  458. @app.post("/api/inbound/exec")
  459. def inbound_exec(payload: InboundExecPayload):
  460. obj, entity = _resolve_entity(payload)
  461. client = _build_inbound_client()
  462. run_id = payload.runId or sa.new_run_id()
  463. contract_version = payload.contractVersion
  464. if obj is None:
  465. # 仅按 entityCode 直发(无目录样例),rows 必填
  466. rows = payload.rows
  467. if rows is None and payload.operation in ("push", "bulk"):
  468. raise HTTPException(status_code=400, detail="entityCode-only exec requires rows")
  469. envelope = sa.inbound_envelope(rows or [])
  470. else:
  471. inbound = obj.get("apiInbound") or {}
  472. contract_version = contract_version or inbound.get("contractVersion")
  473. rows = payload.rows if payload.rows is not None else sa.adapt_inbound_rows(obj, run_id)["rows"]
  474. envelope = sa.inbound_envelope(rows)
  475. object_code = obj["objectCode"] if obj else entity
  476. op = payload.operation
  477. result: dict = {"runId": run_id, "entityCode": entity, "operation": op, "case": payload.case}
  478. def record(status: int, body, headers, layers: dict, extra: dict | None = None):
  479. result.update({"httpStatus": status, "response": body, "requestHeaders": headers, **(extra or {})})
  480. RUNS.add(run_id, object_code, "API_INBOUND",
  481. {"operation": op, "case": payload.case, "entityCode": entity,
  482. "rowCount": len(rows or []), "contractVersion": contract_version},
  483. {"httpStatus": status, "body": body, "headers": headers},
  484. layers, _followup_sql(obj or {"objectCode": entity, "apiInbound": {"entityCode": entity}},
  485. "API_INBOUND", run_id))
  486. return result
  487. # 预检(配置层)
  488. pre = aidop_probe.precheck_inbound_entity(client, entity, CONFIG["allowedTargetHosts"], contract_version)
  489. config_state = pre["configState"]
  490. result["precheck"] = {"configState": config_state, "httpStatus": pre["httpStatus"]}
  491. if op == "schema":
  492. status, body, headers = client.op_schema(entity, contract_version)
  493. layers = _layers("SUPPORTED", aidop_probe.classify_inbound_precheck(status, body), "NOT_RUN")
  494. return record(status, body, headers, layers)
  495. if config_state != "READY" and op in ("push", "bulk", "snapshot_open"):
  496. layers = _layers("SUPPORTED", config_state, "NOT_RUN")
  497. raise HTTPException(status_code=409, detail={
  498. "message": f"配置层 {config_state}:普通推数只对 READY 实体执行(不伪造成功)",
  499. "precheck": pre, "layers": layers})
  500. if op == "push":
  501. case = payload.case or "happy"
  502. if case == "happy":
  503. status, body, headers = client.op_push(entity, envelope, contract_version=contract_version)
  504. elif case == "missing-idem":
  505. status, body, headers = client.case_missing_idem(entity, envelope)
  506. elif case == "tamper-body":
  507. status, body, headers = client.case_tamper_body(entity, envelope)
  508. elif case == "expired-ts":
  509. status, body, headers = client.case_expired_ts(entity, envelope)
  510. elif case == "bad-signature":
  511. status, body, headers = client.case_bad_signature(entity, envelope)
  512. elif case == "nonce-replay":
  513. (s1, d1), (s2, d2) = client.case_nonce_replay(entity, envelope)
  514. status, body, headers = s2, {"first": [s1, d1], "replay": [s2, d2]}, {}
  515. elif case == "same-key-replay":
  516. (s1, d1), (s2, d2) = client.case_same_key_replay(entity, envelope)
  517. status, body, headers = s2, {"first": [s1, d1], "replay": [s2, d2]}, {}
  518. elif case == "same-key-diff-body":
  519. env2 = sa.inbound_envelope([{**r, "sourceUpdatedAt": sa.now_iso()} for r in (rows or [])])
  520. (s1, d1), (s2, d2) = client.case_same_key_diff_body(entity, envelope, env2)
  521. status, body, headers = s2, {"first": [s1, d1], "conflict": [s2, d2]}, {}
  522. else:
  523. raise HTTPException(status_code=400, detail=f"unknown case: {case}")
  524. data_state = "ACCEPTED" if status == 202 else ("RUN_FAILED" if status >= 400 else "NOT_RUN")
  525. layers = _layers("SUPPORTED", config_state, data_state)
  526. return record(status, body, headers, layers,
  527. extra={"gaps": ["HTTP 202 仅代表摄取受理;贴源/标准层核验见 followupSql(TRANSFORM_PENDING)"]})
  528. if op == "receipt":
  529. if not payload.syncBatchId:
  530. raise HTTPException(status_code=400, detail="receipt requires syncBatchId")
  531. status, body, headers = client.op_receipt(payload.syncBatchId)
  532. layers = _layers("SUPPORTED", config_state, "ACCEPTED" if status == 200 else "RUN_FAILED")
  533. return record(status, body, headers, layers)
  534. if op == "bulk":
  535. status, body, headers = client.op_bulk(entity, rows or [], contract_version=contract_version)
  536. layers = _layers("SUPPORTED", config_state, "ACCEPTED" if status == 202 else "RUN_FAILED")
  537. return record(status, body, headers, layers)
  538. if op == "snapshot_open":
  539. status, body, headers = client.op_snapshot_open(entity)
  540. layers = _layers("SUPPORTED", config_state, "ACCEPTED" if status in (200, 201) else "RUN_FAILED")
  541. return record(status, body, headers, layers)
  542. if op == "snapshot_commit":
  543. if not payload.snapshotId:
  544. raise HTTPException(status_code=400, detail="snapshot_commit requires snapshotId")
  545. status, body, headers = client.op_snapshot_commit(entity, payload.snapshotId, payload.force)
  546. layers = _layers("SUPPORTED", config_state, "ACCEPTED" if status in (200, 202) else "RUN_FAILED")
  547. return record(status, body, headers, layers)
  548. if op == "digest":
  549. biz_date = payload.bizDate or time.strftime("%Y-%m-%d")
  550. status, body, headers = client.op_digest(entity, biz_date)
  551. layers = _layers("SUPPORTED", config_state, "ACCEPTED" if status == 200 else "RUN_FAILED")
  552. return record(status, body, headers, layers)
  553. raise HTTPException(status_code=400, detail=f"unknown operation: {op}")
  554. # ---------------------------------------------------------------------------
  555. # 场景执行(P1-F)
  556. # ---------------------------------------------------------------------------
  557. @app.post("/api/scenarios/run")
  558. def scenario_run(payload: ScenarioPayload):
  559. scenario = next((s for s in sa.load_scenarios()["scenarios"]
  560. if s["scenarioId"] == payload.scenarioId), None)
  561. if not scenario:
  562. raise HTTPException(status_code=404, detail=f"unknown scenario: {payload.scenarioId}")
  563. run_id = payload.runId or sa.new_run_id()
  564. overrides = payload.channelOverride or {}
  565. nodes_out = []
  566. for node in scenario["nodes"]:
  567. obj = sa.find_object(node["objectCode"])
  568. channel = overrides.get(node["objectCode"], node["channel"])
  569. node_result = {"objectCode": node["objectCode"], "channel": channel,
  570. "relationKey": node.get("relationKey")}
  571. if not obj:
  572. node_result.update({"state": "CODE_UNSUPPORTED", "message": "目录中无此对象"})
  573. nodes_out.append(node_result)
  574. continue
  575. spec = {"DB_SYNC": obj.get("dbSync"), "API_PULL": obj.get("apiPull"),
  576. "API_INBOUND": obj.get("apiInbound")}.get(channel)
  577. if not spec or not spec.get("supported"):
  578. reason = (spec or {}).get("reason") or node.get("inboundNote") or "该通道不支持"
  579. node_result.update({"state": "CODE_UNSUPPORTED", "message": reason})
  580. nodes_out.append(node_result)
  581. continue
  582. try:
  583. if channel == "DB_SYNC":
  584. rows = sa.adapt_db_rows(obj, run_id)
  585. if not _env("AIDOP_SIM_MYSQL_PASSWORD"):
  586. raise RuntimeError("缺少 AIDOP_SIM_MYSQL_PASSWORD")
  587. mock_db.init_schema({spec["table"]: {"rows": rows, "businessKey": spec.get("businessKey"),
  588. "incrementColumn": spec.get("incrementColumn")}})
  589. count = mock_db.upsert_rows(spec["table"], rows)
  590. node_result.update({"state": "ACCEPTED", "loaded": count, "table": spec["table"]})
  591. elif channel == "API_PULL":
  592. if not _mock_running():
  593. raise RuntimeError("Mock API 未启动(先 POST /api/mock-api/start)")
  594. rows = sa.adapt_api_rows(obj, run_id)
  595. MOCK_STATE.overrides[spec["path"]] = rows
  596. node_result.update({"state": "ACCEPTED", "overridePath": spec["path"], "rowCount": len(rows),
  597. "note": "已在 Mock API 挂载适配样例;由 Ai-DOP MdpApiPullExecutor 真实拉取"})
  598. else:
  599. client = _build_inbound_client()
  600. entity = spec["entityCode"]
  601. version = node.get("contractVersion") or spec.get("contractVersion")
  602. pre = aidop_probe.precheck_inbound_entity(client, entity, CONFIG["allowedTargetHosts"], version)
  603. if pre["configState"] != "READY":
  604. node_result.update({"state": pre["configState"],
  605. "message": f"配置层 {pre['configState']},未推数(不伪造成功)"})
  606. else:
  607. rows = sa.adapt_inbound_rows(obj, run_id)["rows"]
  608. status, body, _ = client.op_push(entity, sa.inbound_envelope(rows),
  609. contract_version=version)
  610. node_result.update({"state": "ACCEPTED" if status == 202 else "RUN_FAILED",
  611. "httpStatus": status, "response": body})
  612. except HTTPException:
  613. raise
  614. except aidop_probe.HostNotAllowedError as ex:
  615. node_result.update({"state": "SOURCE_UNREACHABLE", "message": str(ex)})
  616. except Exception as ex:
  617. node_result.update({"state": "RUN_FAILED", "message": f"{type(ex).__name__}: {ex}"})
  618. nodes_out.append(node_result)
  619. unsupported = [n["objectCode"] for n in nodes_out if n["state"] == "CODE_UNSUPPORTED"]
  620. layers = _layers("SUPPORTED",
  621. "READY" if all(n["state"] not in ("SOURCE_UNREACHABLE",) for n in nodes_out) else "SOURCE_UNREACHABLE",
  622. "ACCEPTED" if all(n["state"] in ("ACCEPTED", "CODE_UNSUPPORTED") for n in nodes_out) else "RUN_FAILED")
  623. entry = RUNS.add(run_id, ",".join(n["objectCode"] for n in nodes_out), f"SCENARIO:{payload.scenarioId}",
  624. {"scenario": scenario["name"]}, {"nodes": nodes_out}, layers,
  625. _followup_sql({"objectCode": "SCENARIO", "dbSync": {}}, "SCENARIO", run_id),
  626. gaps=([f"以下节点该通道不支持:{', '.join(unsupported)}(链路不完整,不得宣称全链路跑通)"]
  627. if unsupported else []),
  628. scenario_id=payload.scenarioId, nodes=nodes_out)
  629. return entry
  630. # ---------------------------------------------------------------------------
  631. # 运行记录
  632. # ---------------------------------------------------------------------------
  633. @app.get("/api/runs")
  634. def runs():
  635. return {"runs": RUNS.recent()}
  636. # ---------------------------------------------------------------------------
  637. # 静态页面
  638. # ---------------------------------------------------------------------------
  639. @app.get("/")
  640. def index():
  641. return FileResponse(SIM_ROOT / "static" / "index.html")
  642. app.mount("/static", StaticFiles(directory=SIM_ROOT / "static"), name="static")
  643. def main() -> None:
  644. print(f"[simulator] control panel: http://127.0.0.1:{DEFAULT_CONTROL_PORT}/")
  645. print(f"[simulator] aidop target : {CONFIG['aidopBaseUrl']} (allowed hosts: {CONFIG['allowedTargetHosts']})")
  646. uvicorn.run(app, host="127.0.0.1", port=DEFAULT_CONTROL_PORT, log_level="info")
  647. if __name__ == "__main__":
  648. main()