routes.py 10 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304
  1. """Governed control-plane endpoints for AI-authored data rules."""
  2. from __future__ import annotations
  3. from typing import Any
  4. from flask import current_app, g, jsonify, request
  5. from app import db
  6. from app.api.data_rules import bp
  7. from app.core.data_rules.authoring import (
  8. OpenAICompatibleRuleModel,
  9. RuleAuthoringAgent,
  10. )
  11. from app.core.data_rules.contracts import (
  12. dataflow_spec_hash,
  13. rule_spec_hash,
  14. standard_spec_hash,
  15. validate_dataflow_spec,
  16. validate_rule_spec,
  17. validate_standard_spec,
  18. )
  19. from app.core.data_rules.production_line import resolve_production_line
  20. from app.core.data_rules.release import ProductionLineReleaseService
  21. from app.core.data_rules.repository import DataRuleRepository
  22. from app.core.data_rules.schema_resolver import (
  23. Neo4jSchemaMetadataCatalog,
  24. SchemaResolver,
  25. )
  26. from app.models.result import failed, success
  27. _VALIDATORS = {
  28. "rule": (validate_rule_spec, rule_spec_hash),
  29. "standard": (validate_standard_spec, standard_spec_hash),
  30. "dataflow": (validate_dataflow_spec, dataflow_spec_hash),
  31. }
  32. def _body() -> dict[str, Any]:
  33. value = request.get_json(silent=True)
  34. if not isinstance(value, dict):
  35. raise ValueError("request body must be an object")
  36. return value
  37. def _bad_request(message: str = "规则请求无效"):
  38. return jsonify(failed(message, code=400)), 400
  39. def _closed_body(allowed: set[str]) -> dict[str, Any]:
  40. body = _body()
  41. if set(body) - allowed:
  42. raise ValueError("request contains unsupported fields")
  43. return body
  44. def _repository() -> DataRuleRepository:
  45. configured = current_app.extensions.get("data_rule_repository")
  46. if configured is not None:
  47. return configured
  48. return DataRuleRepository(db.session)
  49. def _release_service() -> ProductionLineReleaseService:
  50. configured = current_app.extensions.get("production_line_release_service")
  51. if configured is not None:
  52. return configured
  53. repository = _repository()
  54. resolver = current_app.extensions.get("data_rule_schema_resolver")
  55. if resolver is None:
  56. catalog = current_app.extensions.get("data_rule_metadata_catalog")
  57. resolver = SchemaResolver(catalog or Neo4jSchemaMetadataCatalog(), repository)
  58. return ProductionLineReleaseService(repository, schema_resolver=resolver)
  59. @bp.get("/capabilities")
  60. def capabilities():
  61. return jsonify(
  62. success(
  63. {
  64. "natural_language_authoring": True,
  65. "schema_constrained_candidates": True,
  66. "production_line_preview": True,
  67. "immutable_asset_versions": True,
  68. "server_side_publishing": True,
  69. "production_line_release": True,
  70. "data_factory_activation": False,
  71. }
  72. )
  73. )
  74. @bp.post("/validate")
  75. def validate_asset():
  76. try:
  77. body = _body()
  78. asset_type = body.get("asset_type")
  79. if asset_type not in _VALIDATORS:
  80. raise ValueError("unsupported asset type")
  81. validator, hasher = _VALIDATORS[asset_type]
  82. normalized = validator(body.get("spec"))
  83. return jsonify(
  84. success(
  85. {
  86. "asset_type": asset_type,
  87. "normalized": normalized,
  88. "spec_hash": hasher(normalized),
  89. }
  90. )
  91. )
  92. except (TypeError, ValueError):
  93. return _bad_request("规则定义无效")
  94. def _authoring_agent() -> RuleAuthoringAgent:
  95. configured = current_app.extensions.get("data_rule_authoring_agent")
  96. if configured is not None:
  97. return configured
  98. api_key = current_app.config.get("LLM_API_KEY") or current_app.config.get(
  99. "DEEPSEEK_API_KEY"
  100. )
  101. if not api_key:
  102. raise RuntimeError("rule authoring model is not configured")
  103. agent = RuleAuthoringAgent(model=OpenAICompatibleRuleModel())
  104. current_app.extensions["data_rule_authoring_agent"] = agent
  105. return agent
  106. @bp.post("/interpret")
  107. def interpret_rule():
  108. try:
  109. body = _body()
  110. result = _authoring_agent().interpret(
  111. source_text=body.get("source_text"),
  112. authoring_surface=body.get("authoring_surface"),
  113. context=body.get("context", {}),
  114. )
  115. audit = _repository().record_generation_run(evidence=result)
  116. db.session.commit()
  117. result = {
  118. **result,
  119. "generation_run_id": audit["id"],
  120. "correlation_id": audit["correlation_id"],
  121. }
  122. return jsonify(success(result))
  123. except (TypeError, ValueError):
  124. db.session.rollback()
  125. return _bad_request("自然语言规则描述无效")
  126. except RuntimeError:
  127. db.session.rollback()
  128. return jsonify(failed("AI 规则解析服务未配置", code=503)), 503
  129. except Exception:
  130. db.session.rollback()
  131. current_app.logger.exception("AI rule interpretation failed")
  132. return jsonify(failed("AI 规则解析暂时不可用", code=503)), 503
  133. @bp.post("/rule-versions")
  134. def create_rule_version():
  135. try:
  136. body = _closed_body(
  137. {
  138. "rule_spec",
  139. "source_text",
  140. "category",
  141. "source_language",
  142. "generated_kind",
  143. }
  144. )
  145. # Enforce V2 at the HTTP boundary even when a test/different
  146. # repository implementation is injected.
  147. rule_spec = validate_rule_spec(body.get("rule_spec"))
  148. result = _repository().create_rule_version(
  149. rule_spec=rule_spec,
  150. source_text=body.get("source_text"),
  151. category=body.get("category", "general"),
  152. source_language=body.get("source_language", "zh-CN"),
  153. generated_kind=body.get("generated_kind", "rulespec"),
  154. created_by=g.current_user["id"],
  155. )
  156. db.session.commit()
  157. return jsonify(success(result, "规则版本创建成功")), 201
  158. except (TypeError, ValueError):
  159. db.session.rollback()
  160. return _bad_request("规则版本定义无效")
  161. except Exception:
  162. db.session.rollback()
  163. current_app.logger.exception("create rule version failed")
  164. return jsonify(failed("规则版本创建失败", code=500)), 500
  165. @bp.post("/rule-versions/<version_id>/publish")
  166. def publish_rule_version(version_id: str):
  167. try:
  168. if request.get_data(cache=True) and request.get_json(silent=True) not in (
  169. None,
  170. {},
  171. ):
  172. raise ValueError("publish request must not contain fields")
  173. result = _repository().publish_rule_version(
  174. version_id=version_id,
  175. published_by=g.current_user["id"],
  176. )
  177. db.session.commit()
  178. return jsonify(success(result, "规则版本发布成功"))
  179. except (TypeError, ValueError):
  180. db.session.rollback()
  181. return jsonify(failed("规则版本无法发布", code=409)), 409
  182. except Exception:
  183. db.session.rollback()
  184. current_app.logger.exception("publish rule version failed")
  185. return jsonify(failed("规则版本发布失败", code=500)), 500
  186. @bp.post("/standard-versions")
  187. def create_standard_version():
  188. try:
  189. body = _closed_body({"standard_spec", "source_text"})
  190. result = _repository().create_standard_version(
  191. standard_spec=body.get("standard_spec"),
  192. source_text=body.get("source_text"),
  193. created_by=g.current_user["id"],
  194. )
  195. db.session.commit()
  196. return jsonify(success(result, "数据标准版本创建成功")), 201
  197. except (TypeError, ValueError):
  198. db.session.rollback()
  199. return _bad_request("数据标准版本定义无效")
  200. except Exception:
  201. db.session.rollback()
  202. current_app.logger.exception("create standard version failed")
  203. return jsonify(failed("数据标准版本创建失败", code=500)), 500
  204. @bp.post("/standard-versions/<version_id>/publish")
  205. def publish_standard_version(version_id: str):
  206. try:
  207. if request.get_data(cache=True) and request.get_json(silent=True) not in (
  208. None,
  209. {},
  210. ):
  211. raise ValueError("publish request must not contain fields")
  212. result = _repository().publish_standard_version(
  213. version_id=version_id,
  214. published_by=g.current_user["id"],
  215. )
  216. db.session.commit()
  217. return jsonify(success(result, "数据标准版本发布成功"))
  218. except (TypeError, ValueError):
  219. db.session.rollback()
  220. return jsonify(failed("数据标准版本无法发布", code=409)), 409
  221. except Exception:
  222. db.session.rollback()
  223. current_app.logger.exception("publish standard version failed")
  224. return jsonify(failed("数据标准版本发布失败", code=500)), 500
  225. @bp.post("/production-lines/resolve")
  226. def resolve_production_line_preview():
  227. try:
  228. body = _body()
  229. package = resolve_production_line(
  230. body.get("dataflow_spec"),
  231. body.get("standard_versions"),
  232. body.get("rule_versions"),
  233. component_binding_ids=body.get("component_binding_ids"),
  234. )
  235. return jsonify(
  236. success(
  237. {
  238. "preview": True,
  239. "release_ready": False,
  240. "package": package,
  241. }
  242. )
  243. )
  244. except (TypeError, ValueError):
  245. return _bad_request("数据生产线定义无效")
  246. @bp.post("/production-lines/<dataflow_uid>/release")
  247. def release_production_line(dataflow_uid: str):
  248. try:
  249. body = _closed_body(
  250. {
  251. "dataflow_spec",
  252. "source_text",
  253. }
  254. )
  255. result = _release_service().release(
  256. dataflow_uid=dataflow_uid,
  257. dataflow_spec=body.get("dataflow_spec"),
  258. source_text=body.get("source_text"),
  259. created_by=g.current_user["id"],
  260. )
  261. db.session.commit()
  262. return jsonify(success(result, "数据生产线发布成功")), 201
  263. except (TypeError, ValueError):
  264. db.session.rollback()
  265. return jsonify(failed("数据生产线无法发布", code=409)), 409
  266. except Exception:
  267. db.session.rollback()
  268. current_app.logger.exception("release production line failed")
  269. return jsonify(failed("数据生产线发布失败", code=500)), 500