| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299 |
- """Governed control-plane endpoints for AI-authored data rules."""
- from __future__ import annotations
- from typing import Any
- from flask import current_app, g, jsonify, request
- from app import db
- from app.api.data_rules import bp
- from app.core.data_rules.authoring import (
- OpenAICompatibleRuleModel,
- RuleAuthoringAgent,
- )
- from app.core.data_rules.contracts import (
- dataflow_spec_hash,
- rule_spec_hash,
- standard_spec_hash,
- validate_dataflow_spec,
- validate_rule_spec,
- validate_standard_spec,
- )
- from app.core.data_rules.production_line import resolve_production_line
- from app.core.data_rules.release import ProductionLineReleaseService
- from app.core.data_rules.repository import DataRuleRepository
- from app.core.data_rules.schema_resolver import SchemaResolver
- from app.models.result import failed, success
- _VALIDATORS = {
- "rule": (validate_rule_spec, rule_spec_hash),
- "standard": (validate_standard_spec, standard_spec_hash),
- "dataflow": (validate_dataflow_spec, dataflow_spec_hash),
- }
- def _body() -> dict[str, Any]:
- value = request.get_json(silent=True)
- if not isinstance(value, dict):
- raise ValueError("request body must be an object")
- return value
- def _bad_request(message: str = "规则请求无效"):
- return jsonify(failed(message, code=400)), 400
- def _closed_body(allowed: set[str]) -> dict[str, Any]:
- body = _body()
- if set(body) - allowed:
- raise ValueError("request contains unsupported fields")
- return body
- def _repository() -> DataRuleRepository:
- configured = current_app.extensions.get("data_rule_repository")
- if configured is not None:
- return configured
- return DataRuleRepository(db.session)
- def _release_service() -> ProductionLineReleaseService:
- configured = current_app.extensions.get("production_line_release_service")
- if configured is not None:
- return configured
- repository = _repository()
- resolver = current_app.extensions.get("data_rule_schema_resolver")
- if resolver is None:
- catalog = current_app.extensions.get("data_rule_metadata_catalog")
- if catalog is not None:
- resolver = SchemaResolver(catalog, repository)
- return ProductionLineReleaseService(repository, schema_resolver=resolver)
- @bp.get("/capabilities")
- def capabilities():
- return jsonify(
- success(
- {
- "natural_language_authoring": True,
- "schema_constrained_candidates": True,
- "production_line_preview": True,
- "immutable_asset_versions": True,
- "server_side_publishing": True,
- "production_line_release": True,
- "data_factory_activation": False,
- }
- )
- )
- @bp.post("/validate")
- def validate_asset():
- try:
- body = _body()
- asset_type = body.get("asset_type")
- if asset_type not in _VALIDATORS:
- raise ValueError("unsupported asset type")
- validator, hasher = _VALIDATORS[asset_type]
- normalized = validator(body.get("spec"))
- return jsonify(
- success(
- {
- "asset_type": asset_type,
- "normalized": normalized,
- "spec_hash": hasher(normalized),
- }
- )
- )
- except (TypeError, ValueError):
- return _bad_request("规则定义无效")
- def _authoring_agent() -> RuleAuthoringAgent:
- configured = current_app.extensions.get("data_rule_authoring_agent")
- if configured is not None:
- return configured
- api_key = current_app.config.get("LLM_API_KEY") or current_app.config.get(
- "DEEPSEEK_API_KEY"
- )
- if not api_key:
- raise RuntimeError("rule authoring model is not configured")
- agent = RuleAuthoringAgent(model=OpenAICompatibleRuleModel())
- current_app.extensions["data_rule_authoring_agent"] = agent
- return agent
- @bp.post("/interpret")
- def interpret_rule():
- try:
- body = _body()
- result = _authoring_agent().interpret(
- source_text=body.get("source_text"),
- authoring_surface=body.get("authoring_surface"),
- context=body.get("context", {}),
- )
- audit = _repository().record_generation_run(evidence=result)
- db.session.commit()
- result = {
- **result,
- "generation_run_id": audit["id"],
- "correlation_id": audit["correlation_id"],
- }
- return jsonify(success(result))
- except (TypeError, ValueError):
- db.session.rollback()
- return _bad_request("自然语言规则描述无效")
- except RuntimeError:
- db.session.rollback()
- return jsonify(failed("AI 规则解析服务未配置", code=503)), 503
- except Exception:
- db.session.rollback()
- current_app.logger.exception("AI rule interpretation failed")
- return jsonify(failed("AI 规则解析暂时不可用", code=503)), 503
- @bp.post("/rule-versions")
- def create_rule_version():
- try:
- body = _closed_body(
- {
- "rule_spec",
- "source_text",
- "category",
- "source_language",
- "generated_kind",
- }
- )
- result = _repository().create_rule_version(
- rule_spec=body.get("rule_spec"),
- source_text=body.get("source_text"),
- category=body.get("category", "general"),
- source_language=body.get("source_language", "zh-CN"),
- generated_kind=body.get("generated_kind", "rulespec"),
- created_by=g.current_user["id"],
- )
- db.session.commit()
- return jsonify(success(result, "规则版本创建成功")), 201
- except (TypeError, ValueError):
- db.session.rollback()
- return _bad_request("规则版本定义无效")
- except Exception:
- db.session.rollback()
- current_app.logger.exception("create rule version failed")
- return jsonify(failed("规则版本创建失败", code=500)), 500
- @bp.post("/rule-versions/<version_id>/publish")
- def publish_rule_version(version_id: str):
- try:
- if request.get_data(cache=True) and request.get_json(silent=True) not in (
- None,
- {},
- ):
- raise ValueError("publish request must not contain fields")
- result = _repository().publish_rule_version(
- version_id=version_id,
- published_by=g.current_user["id"],
- )
- db.session.commit()
- return jsonify(success(result, "规则版本发布成功"))
- except (TypeError, ValueError):
- db.session.rollback()
- return jsonify(failed("规则版本无法发布", code=409)), 409
- except Exception:
- db.session.rollback()
- current_app.logger.exception("publish rule version failed")
- return jsonify(failed("规则版本发布失败", code=500)), 500
- @bp.post("/standard-versions")
- def create_standard_version():
- try:
- body = _closed_body({"standard_spec", "source_text"})
- result = _repository().create_standard_version(
- standard_spec=body.get("standard_spec"),
- source_text=body.get("source_text"),
- created_by=g.current_user["id"],
- )
- db.session.commit()
- return jsonify(success(result, "数据标准版本创建成功")), 201
- except (TypeError, ValueError):
- db.session.rollback()
- return _bad_request("数据标准版本定义无效")
- except Exception:
- db.session.rollback()
- current_app.logger.exception("create standard version failed")
- return jsonify(failed("数据标准版本创建失败", code=500)), 500
- @bp.post("/standard-versions/<version_id>/publish")
- def publish_standard_version(version_id: str):
- try:
- if request.get_data(cache=True) and request.get_json(silent=True) not in (
- None,
- {},
- ):
- raise ValueError("publish request must not contain fields")
- result = _repository().publish_standard_version(
- version_id=version_id,
- published_by=g.current_user["id"],
- )
- db.session.commit()
- return jsonify(success(result, "数据标准版本发布成功"))
- except (TypeError, ValueError):
- db.session.rollback()
- return jsonify(failed("数据标准版本无法发布", code=409)), 409
- except Exception:
- db.session.rollback()
- current_app.logger.exception("publish standard version failed")
- return jsonify(failed("数据标准版本发布失败", code=500)), 500
- @bp.post("/production-lines/resolve")
- def resolve_production_line_preview():
- try:
- body = _body()
- package = resolve_production_line(
- body.get("dataflow_spec"),
- body.get("standard_versions"),
- body.get("rule_versions"),
- component_binding_ids=body.get("component_binding_ids"),
- )
- return jsonify(
- success(
- {
- "preview": True,
- "release_ready": False,
- "package": package,
- }
- )
- )
- except (TypeError, ValueError):
- return _bad_request("数据生产线定义无效")
- @bp.post("/production-lines/<dataflow_uid>/release")
- def release_production_line(dataflow_uid: str):
- try:
- body = _closed_body(
- {
- "dataflow_spec",
- "source_text",
- }
- )
- result = _release_service().release(
- dataflow_uid=dataflow_uid,
- dataflow_spec=body.get("dataflow_spec"),
- source_text=body.get("source_text"),
- created_by=g.current_user["id"],
- )
- db.session.commit()
- return jsonify(success(result, "数据生产线发布成功")), 201
- except (TypeError, ValueError):
- db.session.rollback()
- return jsonify(failed("数据生产线无法发布", code=409)), 409
- except Exception:
- db.session.rollback()
- current_app.logger.exception("release production line failed")
- return jsonify(failed("数据生产线发布失败", code=500)), 500
|