| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751 |
- """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()
|