routes.py 20 KB

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