|
|
@@ -0,0 +1,751 @@
|
|
|
+"""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 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()
|
|
|
+
|
|
|
+
|
|
|
+async def _start_mock_api() -> dict:
|
|
|
+ if MOCK_HANDLE.task and not MOCK_HANDLE.task.done():
|
|
|
+ return {"running": True, "port": CONFIG["mockApiPort"], "alreadyRunning": True}
|
|
|
+ config = uvicorn.Config(MOCK_APP, host="127.0.0.1", port=int(CONFIG["mockApiPort"]), log_level="warning")
|
|
|
+ server = uvicorn.Server(config)
|
|
|
+ MOCK_HANDLE.server = server
|
|
|
+ MOCK_HANDLE.task = asyncio.get_running_loop().create_task(server.serve())
|
|
|
+ await asyncio.sleep(0.4)
|
|
|
+ if server.startup_error:
|
|
|
+ MOCK_HANDLE.last_error = str(server.startup_error)
|
|
|
+ return {"running": False, "error": MOCK_HANDLE.last_error}
|
|
|
+ return {"running": True, "port": CONFIG["mockApiPort"]}
|
|
|
+
|
|
|
+
|
|
|
+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()
|