routes.py 23 KB

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