routes.py 10 KB

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