routes.py 21 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597
  1. """Governed control-plane endpoints for AI-authored data rules."""
  2. from __future__ import annotations
  3. import hashlib
  4. import json
  5. from datetime import UTC, datetime, timedelta
  6. from typing import Any
  7. from flask import current_app, g, jsonify, request
  8. from app import db
  9. from app.api.data_rules import bp
  10. from app.core.data_rules.authoring import (
  11. OpenAICompatibleRuleModel,
  12. RuleAuthoringAgent,
  13. )
  14. from app.core.data_rules.contracts import (
  15. dataflow_spec_hash,
  16. rule_spec_hash,
  17. standard_spec_hash,
  18. validate_dataflow_spec,
  19. validate_rule_spec,
  20. validate_standard_spec,
  21. )
  22. from app.core.data_rules.production_line import resolve_production_line
  23. from app.core.data_rules.publication import (
  24. GenerationReceiptSigner,
  25. LogicalRuleCompiler,
  26. PhysicalPlanPublicationService,
  27. RulePublicationService,
  28. RuleValidationRejected,
  29. ServerOwnedLogicalDryRunRunner,
  30. ServerOwnedPhysicalPreflightRunner,
  31. generation_receipt_claims,
  32. )
  33. from app.core.data_rules.release import ProductionLineReleaseService
  34. from app.core.data_rules.repository import DataRuleRepository
  35. from app.core.data_rules.schema_resolver import (
  36. Neo4jSchemaMetadataCatalog,
  37. SchemaResolver,
  38. )
  39. from app.models.result import failed, success
  40. _VALIDATORS = {
  41. "rule": (validate_rule_spec, rule_spec_hash),
  42. "standard": (validate_standard_spec, standard_spec_hash),
  43. "dataflow": (validate_dataflow_spec, dataflow_spec_hash),
  44. }
  45. def _body() -> dict[str, Any]:
  46. value = request.get_json(silent=True)
  47. if not isinstance(value, dict):
  48. raise ValueError("request body must be an object")
  49. return value
  50. def _bad_request(message: str = "规则请求无效"):
  51. return jsonify(failed(message, code=400)), 400
  52. def _closed_body(allowed: set[str]) -> dict[str, Any]:
  53. body = _body()
  54. if set(body) - allowed:
  55. raise ValueError("request contains unsupported fields")
  56. return body
  57. def _repository() -> DataRuleRepository:
  58. configured = current_app.extensions.get("data_rule_repository")
  59. if configured is not None:
  60. return configured
  61. return DataRuleRepository(db.session)
  62. def _release_service() -> ProductionLineReleaseService:
  63. configured = current_app.extensions.get("production_line_release_service")
  64. if configured is not None:
  65. return configured
  66. repository = _repository()
  67. resolver = current_app.extensions.get("data_rule_schema_resolver")
  68. if resolver is None:
  69. catalog = current_app.extensions.get("data_rule_metadata_catalog")
  70. resolver = SchemaResolver(catalog or Neo4jSchemaMetadataCatalog(), repository)
  71. return ProductionLineReleaseService(repository, schema_resolver=resolver)
  72. def _receipt_signer() -> GenerationReceiptSigner:
  73. configured = current_app.extensions.get("generation_receipt_signer")
  74. if configured is not None:
  75. return configured
  76. secret = current_app.config.get("RULE_GENERATION_RECEIPT_SECRET")
  77. if not isinstance(secret, str) or not secret.strip():
  78. raise RuntimeError(
  79. "dedicated generation receipt secret is not configured"
  80. )
  81. signer = GenerationReceiptSigner(secret)
  82. current_app.extensions["generation_receipt_signer"] = signer
  83. return signer
  84. def _rule_artifact_store():
  85. configured = current_app.extensions.get("rule_artifact_store")
  86. if configured is not None:
  87. return configured
  88. from minio import Minio
  89. from app.runner.artifacts import ArtifactStore
  90. store = ArtifactStore(
  91. Minio(
  92. current_app.config["MINIO_HOST"],
  93. access_key=current_app.config["MINIO_USER"],
  94. secret_key=current_app.config["MINIO_PASSWORD"],
  95. secure=bool(current_app.config["MINIO_SECURE"]),
  96. ),
  97. bucket=current_app.config["MINIO_BUCKET"],
  98. max_artifact_bytes=32 * 1024 * 1024,
  99. max_rows=100_000,
  100. memory_limit_bytes=256 * 1024 * 1024,
  101. max_ttl_seconds=3600,
  102. )
  103. current_app.extensions["rule_artifact_store"] = store
  104. return store
  105. def _publication_service() -> RulePublicationService:
  106. configured = current_app.extensions.get("rule_publication_service")
  107. if configured is not None:
  108. return configured
  109. test_runner = current_app.extensions.get("rule_validation_test_runner")
  110. if test_runner is None:
  111. try:
  112. test_runner = ServerOwnedLogicalDryRunRunner(
  113. _rule_artifact_store()
  114. )
  115. except Exception as exc:
  116. raise RuntimeError(
  117. "trusted rule validation test runner is not configured"
  118. ) from exc
  119. current_app.extensions["rule_validation_test_runner"] = test_runner
  120. service = RulePublicationService(
  121. _repository(),
  122. receipt_signer=_receipt_signer(),
  123. compiler=LogicalRuleCompiler(),
  124. test_runner=test_runner,
  125. )
  126. current_app.extensions["rule_publication_service"] = service
  127. return service
  128. def _physical_publication_service() -> PhysicalPlanPublicationService:
  129. configured = current_app.extensions.get(
  130. "physical_plan_publication_service"
  131. )
  132. if configured is not None:
  133. return configured
  134. test_runner = current_app.extensions.get("rule_physical_test_runner")
  135. if test_runner is None:
  136. try:
  137. from app.core.data_source.runtime import get_data_source_manager
  138. test_runner = ServerOwnedPhysicalPreflightRunner(
  139. _rule_artifact_store(),
  140. datasource_manager=get_data_source_manager(),
  141. )
  142. except Exception as exc:
  143. raise RuntimeError(
  144. "trusted physical plan test runner is not configured"
  145. ) from exc
  146. current_app.extensions["rule_physical_test_runner"] = test_runner
  147. service = PhysicalPlanPublicationService(
  148. _repository(), test_runner=test_runner
  149. )
  150. current_app.extensions["physical_plan_publication_service"] = service
  151. return service
  152. def _metadata_hash(value: Any) -> str:
  153. return hashlib.sha256(
  154. json.dumps(
  155. value,
  156. sort_keys=True,
  157. separators=(",", ":"),
  158. ensure_ascii=False,
  159. ).encode("utf-8")
  160. ).hexdigest()
  161. @bp.get("/capabilities")
  162. def capabilities():
  163. return jsonify(
  164. success(
  165. {
  166. "natural_language_authoring": True,
  167. "schema_constrained_candidates": True,
  168. "production_line_preview": True,
  169. "immutable_asset_versions": True,
  170. "server_side_publishing": True,
  171. "production_line_release": True,
  172. "data_factory_activation": False,
  173. }
  174. )
  175. )
  176. @bp.post("/validate")
  177. def validate_asset():
  178. try:
  179. body = _closed_body({"asset_type", "spec"})
  180. asset_type = body.get("asset_type")
  181. if asset_type not in _VALIDATORS:
  182. raise ValueError("unsupported asset type")
  183. validator, hasher = _VALIDATORS[asset_type]
  184. normalized = validator(body.get("spec"))
  185. return jsonify(
  186. success(
  187. {
  188. "asset_type": asset_type,
  189. "normalized": normalized,
  190. "spec_hash": hasher(normalized),
  191. }
  192. )
  193. )
  194. except (TypeError, ValueError):
  195. return _bad_request("规则定义无效")
  196. def _authoring_agent() -> RuleAuthoringAgent:
  197. configured = current_app.extensions.get("data_rule_authoring_agent")
  198. if configured is not None:
  199. return configured
  200. api_key = current_app.config.get("LLM_API_KEY") or current_app.config.get(
  201. "DEEPSEEK_API_KEY"
  202. )
  203. if not api_key:
  204. raise RuntimeError("rule authoring model is not configured")
  205. agent = RuleAuthoringAgent(model=OpenAICompatibleRuleModel())
  206. current_app.extensions["data_rule_authoring_agent"] = agent
  207. return agent
  208. @bp.post("/interpret")
  209. def interpret_rule():
  210. try:
  211. body = _closed_body(
  212. {"source_text", "authoring_surface", "context"}
  213. )
  214. receipt_signer = _receipt_signer()
  215. repository = _repository()
  216. validation_context = repository.resolve_validation_context(
  217. body.get("context", {})
  218. )
  219. result = _authoring_agent().interpret(
  220. source_text=body.get("source_text"),
  221. authoring_surface=body.get("authoring_surface"),
  222. context=validation_context,
  223. )
  224. audit = repository.record_generation_run(
  225. evidence=result,
  226. created_by=g.current_user["id"],
  227. validation_context=validation_context,
  228. )
  229. result = {
  230. **result,
  231. "generation_run_id": audit["id"],
  232. "correlation_id": audit["correlation_id"],
  233. }
  234. candidate = result.get("candidate")
  235. if (
  236. result.get("status") == "ready"
  237. and isinstance(candidate, dict)
  238. and candidate.get("candidate_type") == "rule"
  239. and isinstance(candidate.get("rule_spec"), dict)
  240. ):
  241. claims = generation_receipt_claims(
  242. generation_run_id=audit["id"],
  243. actor_uid=g.current_user["id"],
  244. source_text=result["source_text"],
  245. candidate_hash=result["candidate_hash"],
  246. rule_spec=candidate["rule_spec"],
  247. model_hash=result.get("model_hash")
  248. or _metadata_hash(
  249. {
  250. "provider": result.get("model_provider", "unknown"),
  251. "name": result.get("model_name", "unknown"),
  252. }
  253. ),
  254. prompt_hash=result.get("prompt_hash")
  255. or _metadata_hash(result.get("prompt_version", "unknown")),
  256. context_hash=result["context_hash"],
  257. expires_at=datetime.now(UTC)
  258. + timedelta(minutes=10),
  259. )
  260. result["generation_receipt"] = receipt_signer.issue(claims)
  261. db.session.commit()
  262. return jsonify(success(result))
  263. except (TypeError, ValueError):
  264. db.session.rollback()
  265. return _bad_request("自然语言规则描述无效")
  266. except RuntimeError:
  267. db.session.rollback()
  268. return jsonify(failed("AI 规则解析服务未配置", code=503)), 503
  269. except Exception:
  270. db.session.rollback()
  271. current_app.logger.exception("AI rule interpretation failed")
  272. return jsonify(failed("AI 规则解析暂时不可用", code=503)), 503
  273. @bp.post("/rule-versions")
  274. def create_rule_version():
  275. try:
  276. body = _closed_body(
  277. {
  278. "rule_spec",
  279. "source_text",
  280. "generation_receipt",
  281. "category",
  282. "source_language",
  283. "generated_kind",
  284. }
  285. )
  286. # Enforce V2 at the HTTP boundary even when a test/different
  287. # repository implementation is injected.
  288. rule_spec = validate_rule_spec(body.get("rule_spec"))
  289. result = _publication_service().create_draft(
  290. rule_spec=rule_spec,
  291. source_text=body.get("source_text"),
  292. category=body.get("category", "general"),
  293. source_language=body.get("source_language", "zh-CN"),
  294. generated_kind=body.get("generated_kind", "rulespec"),
  295. actor_uid=g.current_user["id"],
  296. generation_receipt=body.get("generation_receipt"),
  297. )
  298. db.session.commit()
  299. return jsonify(success(result, "规则版本创建成功")), 201
  300. except (TypeError, ValueError):
  301. db.session.rollback()
  302. return _bad_request("规则版本定义无效")
  303. except Exception:
  304. db.session.rollback()
  305. current_app.logger.exception("create rule version failed")
  306. return jsonify(failed("规则版本创建失败", code=500)), 500
  307. @bp.post("/rule-versions/<version_id>/publish")
  308. def publish_rule_version(version_id: str):
  309. try:
  310. if request.get_data(cache=True) and request.get_json(silent=True) not in (
  311. None,
  312. {},
  313. ):
  314. raise ValueError("publish request must not contain fields")
  315. result = _publication_service().publish(
  316. version_id, g.current_user["id"]
  317. )
  318. db.session.commit()
  319. return jsonify(success(result, "规则版本发布成功"))
  320. except (TypeError, ValueError):
  321. db.session.rollback()
  322. return jsonify(failed("规则版本无法发布", code=409)), 409
  323. except Exception:
  324. db.session.rollback()
  325. current_app.logger.exception("publish rule version failed")
  326. return jsonify(failed("规则版本发布失败", code=500)), 500
  327. @bp.post("/rule-versions/<version_id>/validate")
  328. def validate_rule_version(version_id: str):
  329. try:
  330. if request.get_data(cache=True) and request.get_json(silent=True) not in (
  331. None,
  332. {},
  333. ):
  334. raise ValueError("validate request must not contain fields")
  335. result = _publication_service().validate(
  336. version_id, g.current_user["id"]
  337. )
  338. db.session.commit()
  339. return jsonify(success(result, "规则版本编译验证成功"))
  340. except RuleValidationRejected:
  341. db.session.commit()
  342. return jsonify(failed("规则版本编译验证失败", code=409)), 409
  343. except (TypeError, ValueError):
  344. db.session.rollback()
  345. return jsonify(failed("规则版本编译验证失败", code=409)), 409
  346. except RuntimeError:
  347. db.session.rollback()
  348. return jsonify(failed("规则验证服务未配置", code=503)), 503
  349. @bp.post("/rule-versions/<version_id>/test")
  350. def test_rule_version(version_id: str):
  351. try:
  352. body = _closed_body({"plan_id"})
  353. result = _publication_service().test(
  354. version_id,
  355. g.current_user["id"],
  356. plan_id=body.get("plan_id"),
  357. )
  358. db.session.commit()
  359. return jsonify(success(result, "规则版本样本测试成功"))
  360. except (TypeError, ValueError):
  361. db.session.rollback()
  362. return jsonify(failed("规则版本样本测试失败", code=409)), 409
  363. except RuntimeError:
  364. db.session.rollback()
  365. return jsonify(failed("规则测试服务未配置", code=503)), 503
  366. @bp.get("/rule-versions/<version_id>/evidence")
  367. def rule_version_evidence(version_id: str):
  368. try:
  369. return jsonify(success(_publication_service().evidence(version_id)))
  370. except (TypeError, ValueError):
  371. return jsonify(failed("规则证据不存在", code=404)), 404
  372. except RuntimeError:
  373. return jsonify(failed("规则验证服务未配置", code=503)), 503
  374. @bp.get("/catalog")
  375. @bp.get("/catalog/rule-versions")
  376. def published_rule_catalog():
  377. try:
  378. if set(request.args) - {"query", "limit"}:
  379. raise ValueError("catalog query contains unsupported fields")
  380. query = request.args.get("query", "")
  381. limit = int(request.args.get("limit", "50"))
  382. return jsonify(
  383. success(
  384. {
  385. "items": _publication_service().catalog(
  386. query=query, limit=limit
  387. )
  388. }
  389. )
  390. )
  391. except (TypeError, ValueError):
  392. return _bad_request("规则目录查询无效")
  393. except RuntimeError:
  394. return jsonify(failed("规则目录服务未配置", code=503)), 503
  395. @bp.post("/execution-plans/<plan_id>/validate")
  396. def validate_physical_plan(plan_id: str):
  397. try:
  398. if request.get_data(cache=True) and request.get_json(silent=True) not in (
  399. None,
  400. {},
  401. ):
  402. raise ValueError("validate request must not contain fields")
  403. result = _physical_publication_service().validate(
  404. plan_id, g.current_user["id"]
  405. )
  406. db.session.commit()
  407. return jsonify(success(result, "物理执行计划编译证据已确认"))
  408. except (TypeError, ValueError):
  409. db.session.rollback()
  410. return jsonify(failed("物理执行计划验证失败", code=409)), 409
  411. except RuntimeError:
  412. db.session.rollback()
  413. return jsonify(failed("物理计划验证服务未配置", code=503)), 503
  414. @bp.post("/execution-plans/<plan_id>/test")
  415. def test_physical_plan(plan_id: str):
  416. try:
  417. if request.get_data(cache=True) and request.get_json(silent=True) not in (
  418. None,
  419. {},
  420. ):
  421. raise ValueError("test request must not contain fields")
  422. result = _physical_publication_service().test(
  423. plan_id, g.current_user["id"]
  424. )
  425. db.session.commit()
  426. return jsonify(success(result, "物理执行计划样本测试成功"))
  427. except (TypeError, ValueError):
  428. db.session.rollback()
  429. return jsonify(failed("物理执行计划样本测试失败", code=409)), 409
  430. except RuntimeError:
  431. db.session.rollback()
  432. return jsonify(failed("物理计划测试服务未配置", code=503)), 503
  433. @bp.post("/execution-plans/<plan_id>/publish")
  434. def publish_physical_plan(plan_id: str):
  435. try:
  436. if request.get_data(cache=True) and request.get_json(silent=True) not in (
  437. None,
  438. {},
  439. ):
  440. raise ValueError("publish request must not contain fields")
  441. result = _physical_publication_service().publish(
  442. plan_id, g.current_user["id"]
  443. )
  444. db.session.commit()
  445. return jsonify(success(result, "物理执行计划发布成功"))
  446. except (TypeError, ValueError):
  447. db.session.rollback()
  448. return jsonify(failed("物理执行计划无法发布", code=409)), 409
  449. except RuntimeError:
  450. db.session.rollback()
  451. return jsonify(failed("物理计划发布服务未配置", code=503)), 503
  452. @bp.post("/standard-versions")
  453. def create_standard_version():
  454. try:
  455. body = _closed_body({"standard_spec", "source_text"})
  456. result = _repository().create_standard_version(
  457. standard_spec=body.get("standard_spec"),
  458. source_text=body.get("source_text"),
  459. created_by=g.current_user["id"],
  460. )
  461. db.session.commit()
  462. return jsonify(success(result, "数据标准版本创建成功")), 201
  463. except (TypeError, ValueError):
  464. db.session.rollback()
  465. return _bad_request("数据标准版本定义无效")
  466. except Exception:
  467. db.session.rollback()
  468. current_app.logger.exception("create standard version failed")
  469. return jsonify(failed("数据标准版本创建失败", code=500)), 500
  470. @bp.post("/standard-versions/<version_id>/publish")
  471. def publish_standard_version(version_id: str):
  472. try:
  473. if request.get_data(cache=True) and request.get_json(silent=True) not in (
  474. None,
  475. {},
  476. ):
  477. raise ValueError("publish request must not contain fields")
  478. result = _repository().publish_standard_version(
  479. version_id=version_id,
  480. published_by=g.current_user["id"],
  481. )
  482. db.session.commit()
  483. return jsonify(success(result, "数据标准版本发布成功"))
  484. except (TypeError, ValueError):
  485. db.session.rollback()
  486. return jsonify(failed("数据标准版本无法发布", code=409)), 409
  487. except Exception:
  488. db.session.rollback()
  489. current_app.logger.exception("publish standard version failed")
  490. return jsonify(failed("数据标准版本发布失败", code=500)), 500
  491. @bp.post("/production-lines/resolve")
  492. def resolve_production_line_preview():
  493. try:
  494. body = _body()
  495. package = resolve_production_line(
  496. body.get("dataflow_spec"),
  497. body.get("standard_versions"),
  498. body.get("rule_versions"),
  499. component_binding_ids=body.get("component_binding_ids"),
  500. )
  501. return jsonify(
  502. success(
  503. {
  504. "preview": True,
  505. "release_ready": False,
  506. "package": package,
  507. }
  508. )
  509. )
  510. except (TypeError, ValueError):
  511. return _bad_request("数据生产线定义无效")
  512. @bp.post("/production-lines/<dataflow_uid>/release")
  513. def release_production_line(dataflow_uid: str):
  514. try:
  515. body = _closed_body(
  516. {
  517. "dataflow_spec",
  518. "source_text",
  519. }
  520. )
  521. result = _release_service().release(
  522. dataflow_uid=dataflow_uid,
  523. dataflow_spec=body.get("dataflow_spec"),
  524. source_text=body.get("source_text"),
  525. created_by=g.current_user["id"],
  526. )
  527. db.session.commit()
  528. return jsonify(success(result, "数据生产线发布成功")), 201
  529. except (TypeError, ValueError):
  530. db.session.rollback()
  531. return jsonify(failed("数据生产线无法发布", code=409)), 409
  532. except Exception:
  533. db.session.rollback()
  534. current_app.logger.exception("release production line failed")
  535. return jsonify(failed("数据生产线发布失败", code=500)), 500