"""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 ( Neo4jSchemaMetadataCatalog, 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") resolver = SchemaResolver(catalog or Neo4jSchemaMetadataCatalog(), 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", } ) # Enforce V2 at the HTTP boundary even when a test/different # repository implementation is injected. rule_spec = validate_rule_spec(body.get("rule_spec")) result = _repository().create_rule_version( rule_spec=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//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//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//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