app.py 33 KB

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