routes.py 39 KB

1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071727374757677787980818283848586878889909192939495969798991001011021031041051061071081091101111121131141151161171181191201211221231241251261271281291301311321331341351361371381391401411421431441451461471481491501511521531541551561571581591601611621631641651661671681691701711721731741751761771781791801811821831841851861871881891901911921931941951961971981992002012022032042052062072082092102112122132142152162172182192202212222232242252262272282292302312322332342352362372382392402412422432442452462472482492502512522532542552562572582592602612622632642652662672682692702712722732742752762772782792802812822832842852862872882892902912922932942952962972982993003013023033043053063073083093103113123133143153163173183193203213223233243253263273283293303313323333343353363373383393403413423433443453463473483493503513523533543553563573583593603613623633643653663673683693703713723733743753763773783793803813823833843853863873883893903913923933943953963973983994004014024034044054064074084094104114124134144154164174184194204214224234244254264274284294304314324334344354364374384394404414424434444454464474484494504514524534544554564574584594604614624634644654664674684694704714724734744754764774784794804814824834844854864874884894904914924934944954964974984995005015025035045055065075085095105115125135145155165175185195205215225235245255265275285295305315325335345355365375385395405415425435445455465475485495505515525535545555565575585595605615625635645655665675685695705715725735745755765775785795805815825835845855865875885895905915925935945955965975985996006016026036046056066076086096106116126136146156166176186196206216226236246256266276286296306316326336346356366376386396406416426436446456466476486496506516526536546556566576586596606616626636646656666676686696706716726736746756766776786796806816826836846856866876886896906916926936946956966976986997007017027037047057067077087097107117127137147157167177187197207217227237247257267277287297307317327337347357367377387397407417427437447457467477487497507517527537547557567577587597607617627637647657667677687697707717727737747757767777787797807817827837847857867877887897907917927937947957967977987998008018028038048058068078088098108118128138148158168178188198208218228238248258268278288298308318328338348358368378388398408418428438448458468478488498508518528538548558568578588598608618628638648658668678688698708718728738748758768778788798808818828838848858868878888898908918928938948958968978988999009019029039049059069079089099109119129139149159169179189199209219229239249259269279289299309319329339349359369379389399409419429439449459469479489499509519529539549559569579589599609619629639649659669679689699709719729739749759769779789799809819829839849859869879889899909919929939949959969979989991000100110021003100410051006100710081009101010111012101310141015101610171018101910201021102210231024102510261027102810291030103110321033103410351036103710381039104010411042104310441045104610471048104910501051105210531054105510561057105810591060106110621063106410651066106710681069107010711072107310741075
  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.data_rules.authoring import (
  11. OpenAICompatibleRuleModel,
  12. RuleAuthoringAgent,
  13. )
  14. from app.core.data_rules.contracts import (
  15. dataflow_spec_hash,
  16. rule_spec_hash,
  17. standard_spec_hash,
  18. validate_dataflow_spec,
  19. validate_rule_spec,
  20. validate_standard_spec,
  21. )
  22. from app.core.data_rules.deployment import (
  23. DataFlowDeploymentService,
  24. DeploymentOperationError,
  25. )
  26. from app.core.data_rules.production_line import resolve_production_line
  27. from app.core.data_rules.publication import (
  28. GenerationReceiptSigner,
  29. LogicalRuleCompiler,
  30. PhysicalPlanPublicationService,
  31. RulePublicationService,
  32. RuleValidationRejected,
  33. ServerOwnedLogicalDryRunRunner,
  34. ServerOwnedPhysicalPreflightRunner,
  35. generation_receipt_claims,
  36. )
  37. from app.core.data_rules.release import ProductionLineReleaseService
  38. from app.core.data_rules.repository import DataRuleRepository
  39. from app.core.data_rules.schema_resolver import (
  40. Neo4jSchemaMetadataCatalog,
  41. SchemaResolver,
  42. )
  43. from app.core.orchestration.engines.kestra import KestraAdapter
  44. from app.models.result import failed, success
  45. from app.runner.auth import TaskTokenIssuer
  46. _VALIDATORS = {
  47. "rule": (validate_rule_spec, rule_spec_hash),
  48. "standard": (validate_standard_spec, standard_spec_hash),
  49. "dataflow": (validate_dataflow_spec, dataflow_spec_hash),
  50. }
  51. def _body() -> dict[str, Any]:
  52. value = request.get_json(silent=True)
  53. if not isinstance(value, dict):
  54. raise ValueError("request body must be an object")
  55. return value
  56. def _bad_request(message: str = "规则请求无效"):
  57. return jsonify(failed(message, code=400)), 400
  58. def _closed_body(allowed: set[str]) -> dict[str, Any]:
  59. body = _body()
  60. if set(body) - allowed:
  61. raise ValueError("request contains unsupported fields")
  62. return body
  63. def _repository() -> DataRuleRepository:
  64. configured = current_app.extensions.get("data_rule_repository")
  65. if configured is not None:
  66. return configured
  67. return DataRuleRepository(db.session)
  68. def _release_service() -> ProductionLineReleaseService:
  69. configured = current_app.extensions.get("production_line_release_service")
  70. if configured is not None:
  71. return configured
  72. repository = _repository()
  73. resolver = current_app.extensions.get("data_rule_schema_resolver")
  74. if resolver is None:
  75. catalog = current_app.extensions.get("data_rule_metadata_catalog")
  76. resolver = SchemaResolver(catalog or Neo4jSchemaMetadataCatalog(), repository)
  77. return ProductionLineReleaseService(repository, schema_resolver=resolver)
  78. def _schema_resolver() -> SchemaResolver:
  79. configured = current_app.extensions.get("data_rule_schema_resolver")
  80. if configured is not None:
  81. return configured
  82. repository = _repository()
  83. catalog = current_app.extensions.get("data_rule_metadata_catalog")
  84. return SchemaResolver(catalog or Neo4jSchemaMetadataCatalog(), repository)
  85. def _deployment_readiness() -> dict[str, Any]:
  86. reasons = []
  87. if not current_app.config.get("KESTRA_BASE_URL"):
  88. reasons.append("kestra_base_url_missing")
  89. if not current_app.config.get("KESTRA_USERNAME"):
  90. reasons.append("kestra_username_missing")
  91. if not current_app.config.get("KESTRA_PASSWORD"):
  92. reasons.append("kestra_password_missing")
  93. secret = current_app.config.get("RUNNER_TASK_TOKEN_SECRET")
  94. if not isinstance(secret, str) or len(secret.encode("utf-8")) < 32:
  95. reasons.append("runner_task_token_secret_invalid")
  96. gate_open = bool(current_app.config.get("DATA_FACTORY_ACTIVATION_ENABLED", False))
  97. if not gate_open:
  98. reasons.append("post_acceptance_activation_gate_closed")
  99. return {
  100. "ready": not reasons,
  101. "gate_open": gate_open,
  102. "reasons": reasons,
  103. "server_owned": True,
  104. }
  105. def _deployment_service() -> DataFlowDeploymentService:
  106. configured = current_app.extensions.get("dataflow_deployment_service")
  107. if configured is not None:
  108. return configured
  109. readiness = _deployment_readiness()
  110. infrastructure_reasons = [
  111. reason
  112. for reason in readiness["reasons"]
  113. if reason != "post_acceptance_activation_gate_closed"
  114. ]
  115. if infrastructure_reasons:
  116. raise RuntimeError("data factory deployment infrastructure is not ready")
  117. engine = KestraAdapter(
  118. current_app.config["KESTRA_BASE_URL"],
  119. current_app.config["KESTRA_USERNAME"],
  120. current_app.config["KESTRA_PASSWORD"],
  121. tenant_id=current_app.config.get("KESTRA_TENANT_ID", "main"),
  122. timeout=int(current_app.config.get("KESTRA_HTTP_TIMEOUT_SECONDS", 30)),
  123. )
  124. service = DataFlowDeploymentService(
  125. _repository(),
  126. engine=engine,
  127. token_issuer=TaskTokenIssuer(current_app.config["RUNNER_TASK_TOKEN_SECRET"]),
  128. )
  129. current_app.extensions["dataflow_deployment_service"] = service
  130. return service
  131. def _activation_gate_ready() -> bool:
  132. if current_app.extensions.get("dataflow_deployment_service") is not None:
  133. return bool(current_app.config.get("DATA_FACTORY_ACTIVATION_ENABLED", False))
  134. return bool(_deployment_readiness()["ready"])
  135. def _receipt_signer() -> GenerationReceiptSigner:
  136. configured = current_app.extensions.get("generation_receipt_signer")
  137. if configured is not None:
  138. return configured
  139. secret = current_app.config.get("RULE_GENERATION_RECEIPT_SECRET")
  140. if not isinstance(secret, str) or not secret.strip():
  141. raise RuntimeError("dedicated generation receipt secret is not configured")
  142. signer = GenerationReceiptSigner(secret)
  143. current_app.extensions["generation_receipt_signer"] = signer
  144. return signer
  145. def _rule_artifact_store():
  146. configured = current_app.extensions.get("rule_artifact_store")
  147. if configured is not None:
  148. return configured
  149. from minio import Minio
  150. from app.runner.artifacts import ArtifactStore
  151. store = ArtifactStore(
  152. Minio(
  153. current_app.config["MINIO_HOST"],
  154. access_key=current_app.config["MINIO_USER"],
  155. secret_key=current_app.config["MINIO_PASSWORD"],
  156. secure=bool(current_app.config["MINIO_SECURE"]),
  157. ),
  158. bucket=current_app.config["MINIO_BUCKET"],
  159. max_artifact_bytes=32 * 1024 * 1024,
  160. max_rows=100_000,
  161. memory_limit_bytes=256 * 1024 * 1024,
  162. max_ttl_seconds=3600,
  163. )
  164. current_app.extensions["rule_artifact_store"] = store
  165. return store
  166. def _publication_service() -> RulePublicationService:
  167. configured = current_app.extensions.get("rule_publication_service")
  168. if configured is not None:
  169. return configured
  170. test_runner = current_app.extensions.get("rule_validation_test_runner")
  171. if test_runner is None:
  172. try:
  173. test_runner = ServerOwnedLogicalDryRunRunner(_rule_artifact_store())
  174. except Exception as exc:
  175. raise RuntimeError(
  176. "trusted rule validation test runner is not configured"
  177. ) from exc
  178. current_app.extensions["rule_validation_test_runner"] = test_runner
  179. service = RulePublicationService(
  180. _repository(),
  181. receipt_signer=_receipt_signer(),
  182. compiler=LogicalRuleCompiler(),
  183. test_runner=test_runner,
  184. )
  185. current_app.extensions["rule_publication_service"] = service
  186. return service
  187. def _physical_publication_service() -> PhysicalPlanPublicationService:
  188. configured = current_app.extensions.get("physical_plan_publication_service")
  189. if configured is not None:
  190. return configured
  191. test_runner = current_app.extensions.get("rule_physical_test_runner")
  192. if test_runner is None:
  193. try:
  194. from app.core.data_source.runtime import get_data_source_manager
  195. test_runner = ServerOwnedPhysicalPreflightRunner(
  196. _rule_artifact_store(),
  197. datasource_manager=get_data_source_manager(),
  198. )
  199. except Exception as exc:
  200. raise RuntimeError(
  201. "trusted physical plan test runner is not configured"
  202. ) from exc
  203. current_app.extensions["rule_physical_test_runner"] = test_runner
  204. service = PhysicalPlanPublicationService(_repository(), test_runner=test_runner)
  205. current_app.extensions["physical_plan_publication_service"] = service
  206. return service
  207. def _metadata_hash(value: Any) -> str:
  208. return hashlib.sha256(
  209. json.dumps(
  210. value,
  211. sort_keys=True,
  212. separators=(",", ":"),
  213. ensure_ascii=False,
  214. ).encode("utf-8")
  215. ).hexdigest()
  216. @bp.get("/capabilities")
  217. def capabilities():
  218. readiness = _deployment_readiness()
  219. if current_app.extensions.get("dataflow_deployment_service") is not None:
  220. readiness = {
  221. **readiness,
  222. "ready": bool(readiness["gate_open"]),
  223. "reasons": [
  224. reason
  225. for reason in readiness["reasons"]
  226. if reason == "post_acceptance_activation_gate_closed"
  227. ],
  228. }
  229. return jsonify(
  230. success(
  231. {
  232. "natural_language_authoring": True,
  233. "schema_constrained_candidates": True,
  234. "production_line_preview": True,
  235. "immutable_asset_versions": True,
  236. "server_side_publishing": True,
  237. "production_line_release": True,
  238. "data_factory_activation": readiness["ready"],
  239. "data_factory_readiness": readiness,
  240. }
  241. )
  242. )
  243. @bp.post("/production-lines/draft-identity")
  244. def create_production_line_draft_identity():
  245. """Persist a short-lived, single-use, actor-bound DataFlow reservation."""
  246. try:
  247. _closed_body(set())
  248. result = _repository().reserve_dataflow_draft(actor_uid=g.current_user["id"])
  249. db.session.commit()
  250. return (
  251. jsonify(
  252. success(
  253. result,
  254. "生产线草稿身份已创建",
  255. code=201,
  256. )
  257. ),
  258. 201,
  259. )
  260. except (TypeError, ValueError):
  261. db.session.rollback()
  262. return _bad_request("生产线草稿身份请求无效")
  263. except Exception:
  264. db.session.rollback()
  265. current_app.logger.exception("create DataFlow draft reservation failed")
  266. return jsonify(failed("生产线草稿身份暂时不可用", code=503)), 503
  267. @bp.post("/validate")
  268. def validate_asset():
  269. try:
  270. body = _closed_body({"asset_type", "spec"})
  271. asset_type = body.get("asset_type")
  272. if asset_type not in _VALIDATORS:
  273. raise ValueError("unsupported asset type")
  274. validator, hasher = _VALIDATORS[asset_type]
  275. normalized = validator(body.get("spec"))
  276. return jsonify(
  277. success(
  278. {
  279. "asset_type": asset_type,
  280. "normalized": normalized,
  281. "spec_hash": hasher(normalized),
  282. }
  283. )
  284. )
  285. except (TypeError, ValueError):
  286. return _bad_request("规则定义无效")
  287. def _authoring_agent() -> RuleAuthoringAgent:
  288. configured = current_app.extensions.get("data_rule_authoring_agent")
  289. if configured is not None:
  290. return configured
  291. api_key = current_app.config.get("LLM_API_KEY") or current_app.config.get(
  292. "DEEPSEEK_API_KEY"
  293. )
  294. if not api_key:
  295. raise RuntimeError("rule authoring model is not configured")
  296. agent = RuleAuthoringAgent(model=OpenAICompatibleRuleModel())
  297. current_app.extensions["data_rule_authoring_agent"] = agent
  298. return agent
  299. @bp.post("/interpret")
  300. def interpret_rule():
  301. try:
  302. body = _closed_body({"source_text", "authoring_surface", "context"})
  303. receipt_signer = _receipt_signer()
  304. repository = _repository()
  305. validation_context = repository.resolve_validation_context(
  306. body.get("context", {})
  307. )
  308. result = _authoring_agent().interpret(
  309. source_text=body.get("source_text"),
  310. authoring_surface=body.get("authoring_surface"),
  311. context=validation_context,
  312. )
  313. candidate = result.get("candidate")
  314. if body.get("authoring_surface") == "data_standard" and (
  315. not isinstance(candidate, dict)
  316. or candidate.get("candidate_type") != "rule"
  317. or not isinstance(candidate.get("rule_spec"), dict)
  318. ):
  319. raise ValueError("data_standard authoring must produce a rule candidate")
  320. audit = repository.record_generation_run(
  321. evidence=result,
  322. created_by=g.current_user["id"],
  323. validation_context=validation_context,
  324. )
  325. result = {
  326. **result,
  327. "generation_run_id": audit["id"],
  328. "correlation_id": audit["correlation_id"],
  329. }
  330. if (
  331. result.get("status") == "ready"
  332. and isinstance(candidate, dict)
  333. and candidate.get("candidate_type") == "rule"
  334. and isinstance(candidate.get("rule_spec"), dict)
  335. ):
  336. claims = generation_receipt_claims(
  337. generation_run_id=audit["id"],
  338. actor_uid=g.current_user["id"],
  339. source_text=result["source_text"],
  340. candidate_hash=result["candidate_hash"],
  341. rule_spec=candidate["rule_spec"],
  342. model_hash=result.get("model_hash")
  343. or _metadata_hash(
  344. {
  345. "provider": result.get("model_provider", "unknown"),
  346. "name": result.get("model_name", "unknown"),
  347. }
  348. ),
  349. prompt_hash=result.get("prompt_hash")
  350. or _metadata_hash(result.get("prompt_version", "unknown")),
  351. context_hash=result["context_hash"],
  352. expires_at=datetime.now(UTC) + timedelta(minutes=10),
  353. )
  354. result["generation_receipt"] = receipt_signer.issue(claims)
  355. db.session.commit()
  356. return jsonify(success(result))
  357. except (TypeError, ValueError):
  358. db.session.rollback()
  359. return _bad_request("自然语言规则描述无效")
  360. except RuntimeError:
  361. db.session.rollback()
  362. return jsonify(failed("AI 规则解析服务未配置", code=503)), 503
  363. except Exception:
  364. db.session.rollback()
  365. current_app.logger.exception("AI rule interpretation failed")
  366. return jsonify(failed("AI 规则解析暂时不可用", code=503)), 503
  367. @bp.post("/rule-versions")
  368. def create_rule_version():
  369. try:
  370. body = _closed_body(
  371. {
  372. "rule_spec",
  373. "source_text",
  374. "generation_receipt",
  375. "category",
  376. "source_language",
  377. "generated_kind",
  378. }
  379. )
  380. # Enforce V2 at the HTTP boundary even when a test/different
  381. # repository implementation is injected.
  382. rule_spec = validate_rule_spec(body.get("rule_spec"))
  383. result = _publication_service().create_draft(
  384. rule_spec=rule_spec,
  385. source_text=body.get("source_text"),
  386. category=body.get("category", "general"),
  387. source_language=body.get("source_language", "zh-CN"),
  388. generated_kind=body.get("generated_kind", "rulespec"),
  389. actor_uid=g.current_user["id"],
  390. generation_receipt=body.get("generation_receipt"),
  391. )
  392. db.session.commit()
  393. return jsonify(success(result, "规则版本创建成功")), 201
  394. except (TypeError, ValueError):
  395. db.session.rollback()
  396. return _bad_request("规则版本定义无效")
  397. except Exception:
  398. db.session.rollback()
  399. current_app.logger.exception("create rule version failed")
  400. return jsonify(failed("规则版本创建失败", code=500)), 500
  401. @bp.post("/rule-versions/<version_id>/publish")
  402. def publish_rule_version(version_id: str):
  403. try:
  404. if request.get_data(cache=True) and request.get_json(silent=True) not in (
  405. None,
  406. {},
  407. ):
  408. raise ValueError("publish request must not contain fields")
  409. result = _publication_service().publish(version_id, g.current_user["id"])
  410. db.session.commit()
  411. return jsonify(success(result, "规则版本发布成功"))
  412. except (TypeError, ValueError):
  413. db.session.rollback()
  414. return jsonify(failed("规则版本无法发布", code=409)), 409
  415. except Exception:
  416. db.session.rollback()
  417. current_app.logger.exception("publish rule version failed")
  418. return jsonify(failed("规则版本发布失败", code=500)), 500
  419. @bp.post("/rule-versions/<version_id>/validate")
  420. def validate_rule_version(version_id: str):
  421. try:
  422. if request.get_data(cache=True) and request.get_json(silent=True) not in (
  423. None,
  424. {},
  425. ):
  426. raise ValueError("validate request must not contain fields")
  427. result = _publication_service().validate(version_id, g.current_user["id"])
  428. db.session.commit()
  429. return jsonify(success(result, "规则版本编译验证成功"))
  430. except RuleValidationRejected:
  431. db.session.commit()
  432. return jsonify(failed("规则版本编译验证失败", code=409)), 409
  433. except (TypeError, ValueError):
  434. db.session.rollback()
  435. return jsonify(failed("规则版本编译验证失败", code=409)), 409
  436. except RuntimeError:
  437. db.session.rollback()
  438. return jsonify(failed("规则验证服务未配置", code=503)), 503
  439. @bp.post("/rule-versions/<version_id>/test")
  440. def test_rule_version(version_id: str):
  441. try:
  442. body = _closed_body({"plan_id"})
  443. result = _publication_service().test(
  444. version_id,
  445. g.current_user["id"],
  446. plan_id=body.get("plan_id"),
  447. )
  448. db.session.commit()
  449. return jsonify(success(result, "规则版本样本测试成功"))
  450. except (TypeError, ValueError):
  451. db.session.rollback()
  452. return jsonify(failed("规则版本样本测试失败", code=409)), 409
  453. except RuntimeError:
  454. db.session.rollback()
  455. return jsonify(failed("规则测试服务未配置", code=503)), 503
  456. @bp.get("/rule-versions/<version_id>/evidence")
  457. def rule_version_evidence(version_id: str):
  458. try:
  459. return jsonify(
  460. success(
  461. _repository().get_asset_evidence(
  462. asset_type="rule", version_id=version_id
  463. )
  464. )
  465. )
  466. except (TypeError, ValueError):
  467. return jsonify(failed("规则证据不存在", code=404)), 404
  468. except Exception:
  469. current_app.logger.exception("load rule evidence failed")
  470. return jsonify(failed("规则证据暂时不可用", code=503)), 503
  471. @bp.get("/catalog/assets/<asset_type>/<version_id>/evidence")
  472. def catalog_asset_evidence(asset_type: str, version_id: str):
  473. try:
  474. return jsonify(
  475. success(
  476. _repository().get_asset_evidence(
  477. asset_type=asset_type,
  478. version_id=version_id,
  479. )
  480. )
  481. )
  482. except (TypeError, ValueError):
  483. return jsonify(failed("资产证据不存在", code=404)), 404
  484. except Exception:
  485. current_app.logger.exception("load catalog asset evidence failed")
  486. return jsonify(failed("资产证据暂时不可用", code=503)), 503
  487. def _catalog_schema_context():
  488. raw_inputs = request.args.get("input_schema_refs")
  489. output = request.args.get("output_schema_ref")
  490. if raw_inputs is None and output is None:
  491. return None, None
  492. if raw_inputs is None or output is None:
  493. raise ValueError("catalog schema context is incomplete")
  494. inputs = json.loads(raw_inputs)
  495. if not isinstance(inputs, list):
  496. raise ValueError("catalog input_schema_refs must be an array")
  497. return inputs, output
  498. @bp.get("/catalog/assets/<asset_type>/<version_id>")
  499. def catalog_asset(asset_type: str, version_id: str):
  500. try:
  501. if set(request.args) - {
  502. "input_schema_refs",
  503. "output_schema_ref",
  504. }:
  505. raise ValueError("catalog asset query contains unsupported fields")
  506. inputs, output = _catalog_schema_context()
  507. result = _repository().get_published_asset(
  508. asset_type=asset_type,
  509. version_id=version_id,
  510. input_schema_refs=inputs,
  511. output_schema_ref=output,
  512. schema_resolver=_schema_resolver() if inputs is not None else None,
  513. )
  514. if inputs is not None:
  515. # SchemaResolver is a read-through snapshot boundary. Persist the
  516. # IDs returned to this response before making them observable.
  517. db.session.commit()
  518. return jsonify(success(result))
  519. except (TypeError, ValueError, json.JSONDecodeError):
  520. db.session.rollback()
  521. return jsonify(failed("已发布资产不存在或上下文无效", code=404)), 404
  522. except Exception:
  523. db.session.rollback()
  524. current_app.logger.exception("load exact catalog asset failed")
  525. return jsonify(failed("规则目录暂时不可用", code=503)), 503
  526. @bp.get("/catalog")
  527. @bp.get("/catalog/rule-versions")
  528. def published_rule_catalog():
  529. try:
  530. legacy_rule_alias = request.path.endswith("/rule-versions")
  531. allowed = {
  532. "query",
  533. "limit",
  534. "offset",
  535. "input_schema_refs",
  536. "output_schema_ref",
  537. }
  538. if not legacy_rule_alias:
  539. allowed.add("asset_type")
  540. if set(request.args) - allowed:
  541. raise ValueError("catalog query contains unsupported fields")
  542. query = request.args.get("query", "")
  543. limit = int(request.args.get("limit", "50"))
  544. offset = int(request.args.get("offset", "0"))
  545. asset_type = (
  546. "rule" if legacy_rule_alias else request.args.get("asset_type") or None
  547. )
  548. if asset_type not in {None, "rule", "standard"}:
  549. raise ValueError("catalog asset_type is invalid")
  550. if len(query) > 200 or limit < 1 or limit > 100:
  551. raise ValueError("catalog bounds are invalid")
  552. if offset < 0 or offset > 1_000_000:
  553. raise ValueError("catalog offset is invalid")
  554. inputs, output = _catalog_schema_context()
  555. catalog_args = {
  556. "query": query,
  557. "asset_type": asset_type,
  558. "limit": limit,
  559. "offset": offset,
  560. }
  561. if inputs is not None:
  562. catalog_args.update(
  563. {
  564. "input_schema_refs": inputs,
  565. "output_schema_ref": output,
  566. "schema_resolver": _schema_resolver(),
  567. }
  568. )
  569. result = _repository().search_published_assets(**catalog_args)
  570. if inputs is not None:
  571. db.session.commit()
  572. return jsonify(success(result))
  573. except (TypeError, ValueError):
  574. db.session.rollback()
  575. return _bad_request("规则目录查询无效")
  576. except Exception:
  577. db.session.rollback()
  578. current_app.logger.exception("load rule catalog failed")
  579. return jsonify(failed("规则目录暂时不可用", code=503)), 503
  580. @bp.post("/execution-plans/<plan_id>/validate")
  581. def validate_physical_plan(plan_id: str):
  582. try:
  583. if request.get_data(cache=True) and request.get_json(silent=True) not in (
  584. None,
  585. {},
  586. ):
  587. raise ValueError("validate request must not contain fields")
  588. result = _physical_publication_service().validate(plan_id, g.current_user["id"])
  589. db.session.commit()
  590. return jsonify(success(result, "物理执行计划编译证据已确认"))
  591. except (TypeError, ValueError):
  592. db.session.rollback()
  593. return jsonify(failed("物理执行计划验证失败", code=409)), 409
  594. except RuntimeError:
  595. db.session.rollback()
  596. return jsonify(failed("物理计划验证服务未配置", code=503)), 503
  597. @bp.post("/execution-plans/<plan_id>/test")
  598. def test_physical_plan(plan_id: str):
  599. try:
  600. if request.get_data(cache=True) and request.get_json(silent=True) not in (
  601. None,
  602. {},
  603. ):
  604. raise ValueError("test request must not contain fields")
  605. result = _physical_publication_service().test(plan_id, g.current_user["id"])
  606. db.session.commit()
  607. return jsonify(success(result, "物理执行计划样本测试成功"))
  608. except (TypeError, ValueError):
  609. db.session.rollback()
  610. return jsonify(failed("物理执行计划样本测试失败", code=409)), 409
  611. except RuntimeError:
  612. db.session.rollback()
  613. return jsonify(failed("物理计划测试服务未配置", code=503)), 503
  614. @bp.post("/execution-plans/<plan_id>/publish")
  615. def publish_physical_plan(plan_id: str):
  616. try:
  617. if request.get_data(cache=True) and request.get_json(silent=True) not in (
  618. None,
  619. {},
  620. ):
  621. raise ValueError("publish request must not contain fields")
  622. result = _physical_publication_service().publish(plan_id, g.current_user["id"])
  623. db.session.commit()
  624. return jsonify(success(result, "物理执行计划发布成功"))
  625. except (TypeError, ValueError):
  626. db.session.rollback()
  627. return jsonify(failed("物理执行计划无法发布", code=409)), 409
  628. except RuntimeError:
  629. db.session.rollback()
  630. return jsonify(failed("物理计划发布服务未配置", code=503)), 503
  631. @bp.post("/standard-versions")
  632. def create_standard_version():
  633. try:
  634. body = _closed_body({"standard_spec", "source_text"})
  635. result = _repository().create_standard_version(
  636. standard_spec=body.get("standard_spec"),
  637. source_text=body.get("source_text"),
  638. created_by=g.current_user["id"],
  639. )
  640. db.session.commit()
  641. return jsonify(success(result, "数据标准版本创建成功")), 201
  642. except (TypeError, ValueError):
  643. db.session.rollback()
  644. return _bad_request("数据标准版本定义无效")
  645. except Exception:
  646. db.session.rollback()
  647. current_app.logger.exception("create standard version failed")
  648. return jsonify(failed("数据标准版本创建失败", code=500)), 500
  649. @bp.post("/standard-versions/<version_id>/publish")
  650. def publish_standard_version(version_id: str):
  651. try:
  652. if request.get_data(cache=True) and request.get_json(silent=True) not in (
  653. None,
  654. {},
  655. ):
  656. raise ValueError("publish request must not contain fields")
  657. result = _repository().publish_standard_version(
  658. version_id=version_id,
  659. published_by=g.current_user["id"],
  660. )
  661. db.session.commit()
  662. return jsonify(success(result, "数据标准版本发布成功"))
  663. except (TypeError, ValueError):
  664. db.session.rollback()
  665. return jsonify(failed("数据标准版本无法发布", code=409)), 409
  666. except Exception:
  667. db.session.rollback()
  668. current_app.logger.exception("publish standard version failed")
  669. return jsonify(failed("数据标准版本发布失败", code=500)), 500
  670. @bp.post("/production-lines/resolve")
  671. def resolve_production_line_preview():
  672. try:
  673. body = _body()
  674. package = resolve_production_line(
  675. body.get("dataflow_spec"),
  676. body.get("standard_versions"),
  677. body.get("rule_versions"),
  678. component_binding_ids=body.get("component_binding_ids"),
  679. )
  680. return jsonify(
  681. success(
  682. {
  683. "preview": True,
  684. "release_ready": False,
  685. "package": package,
  686. }
  687. )
  688. )
  689. except (TypeError, ValueError):
  690. return _bad_request("数据生产线定义无效")
  691. @bp.post("/production-lines/<dataflow_uid>/release")
  692. def release_production_line(dataflow_uid: str):
  693. try:
  694. body = _closed_body(
  695. {
  696. "dataflow_spec",
  697. "source_text",
  698. }
  699. )
  700. result = _release_service().release(
  701. dataflow_uid=dataflow_uid,
  702. dataflow_spec=body.get("dataflow_spec"),
  703. source_text=body.get("source_text"),
  704. created_by=g.current_user["id"],
  705. )
  706. db.session.commit()
  707. return jsonify(success(result, "数据生产线发布成功")), 201
  708. except (TypeError, ValueError):
  709. db.session.rollback()
  710. return jsonify(failed("数据生产线无法发布", code=409)), 409
  711. except Exception:
  712. db.session.rollback()
  713. current_app.logger.exception("release production line failed")
  714. return jsonify(failed("数据生产线发布失败", code=500)), 500
  715. @bp.get("/deployments")
  716. def list_dataflow_deployments():
  717. try:
  718. if set(request.args) - {"environment"}:
  719. raise ValueError("deployment query contains unsupported fields")
  720. return jsonify(
  721. success(
  722. {
  723. "items": _deployment_service().list(
  724. environment=request.args.get("environment") or None
  725. ),
  726. "readiness": _deployment_readiness(),
  727. }
  728. )
  729. )
  730. except (TypeError, ValueError):
  731. return _bad_request("数据工厂投产查询无效")
  732. except RuntimeError:
  733. return jsonify(failed("数据工厂投产服务未配置", code=503)), 503
  734. @bp.get("/deployments/candidates")
  735. def list_dataflow_deployment_candidates():
  736. try:
  737. if set(request.args) != {"environment"}:
  738. raise ValueError("candidate environment is required")
  739. return jsonify(
  740. success(
  741. {
  742. "items": _deployment_service().candidates(
  743. environment=request.args["environment"]
  744. )
  745. }
  746. )
  747. )
  748. except (TypeError, ValueError):
  749. return _bad_request("可入厂生产线查询无效")
  750. except RuntimeError:
  751. return jsonify(failed("数据工厂投产服务未配置", code=503)), 503
  752. @bp.post("/deployments")
  753. def create_dataflow_deployment():
  754. try:
  755. body = _closed_body(
  756. {
  757. "dataflow_version_id",
  758. "environment",
  759. "binding_snapshot",
  760. "schedule_plan",
  761. "reason",
  762. "idempotency_key",
  763. "correlation_id",
  764. }
  765. )
  766. result = _deployment_service().create(
  767. body.get("dataflow_version_id"),
  768. binding_snapshot=body.get("binding_snapshot"),
  769. environment=body.get("environment"),
  770. schedule_plan=body.get("schedule_plan"),
  771. actor_uid=g.current_user["id"],
  772. reason=body.get("reason"),
  773. idempotency_key=body.get("idempotency_key"),
  774. correlation_id=body.get("correlation_id"),
  775. )
  776. db.session.commit()
  777. return jsonify(success(result, "数据生产线已进入数据工厂")), 201
  778. except (TypeError, ValueError) as exc:
  779. db.session.rollback()
  780. return _deployment_conflict(exc, "数据工厂投产定义无效")
  781. except RuntimeError:
  782. db.session.rollback()
  783. return jsonify(failed("数据工厂投产服务未配置", code=503)), 503
  784. except Exception:
  785. db.session.rollback()
  786. current_app.logger.exception("create dataflow deployment failed")
  787. return jsonify(failed("数据工厂投产创建失败", code=500)), 500
  788. def _deployment_action_body(extra: set[str] | None = None):
  789. return _closed_body(
  790. {
  791. "reason",
  792. "idempotency_key",
  793. "correlation_id",
  794. *(extra or set()),
  795. }
  796. )
  797. def _deployment_conflict(exc: Exception, message: str):
  798. code = getattr(exc, "code", "terminal_conflict")
  799. disposition = getattr(exc, "idempotency_key_disposition", "rotate")
  800. error = {
  801. "code": code,
  802. "idempotency_key_disposition": disposition,
  803. "reconcile_required": code == "operation_unknown",
  804. }
  805. blocking_operation = getattr(exc, "blocking_operation", None)
  806. if blocking_operation is not None:
  807. error["blocking_operation"] = blocking_operation
  808. return (
  809. jsonify(
  810. failed(
  811. message,
  812. code=409,
  813. error=error,
  814. )
  815. ),
  816. 409,
  817. )
  818. @bp.post("/deployments/<deployment_id>/deploy-disabled")
  819. def deploy_dataflow_disabled(deployment_id: str):
  820. try:
  821. body = _deployment_action_body()
  822. result = _deployment_service().deploy_disabled(
  823. deployment_id,
  824. g.current_user["id"],
  825. reason=body.get("reason", "deploy disabled"),
  826. idempotency_key=body.get("idempotency_key"),
  827. correlation_id=body.get("correlation_id"),
  828. )
  829. db.session.commit()
  830. return jsonify(success(result, "生产线已禁用部署"))
  831. except DeploymentOperationError as exc:
  832. db.session.rollback()
  833. return _deployment_conflict(exc, "生产线禁用部署状态冲突")
  834. except (TypeError, ValueError) as exc:
  835. db.session.rollback()
  836. return _deployment_conflict(exc, "生产线无法禁用部署")
  837. except RuntimeError:
  838. db.session.rollback()
  839. return jsonify(failed("调度引擎部署结果异常", code=503)), 503
  840. @bp.post("/deployments/<deployment_id>/canary")
  841. def run_dataflow_canary(deployment_id: str):
  842. try:
  843. body = _deployment_action_body({"inputs"})
  844. result = _deployment_service().run_canary(
  845. deployment_id,
  846. body.get("inputs", {}),
  847. g.current_user["id"],
  848. reason=body.get("reason", "trial production"),
  849. idempotency_key=body.get("idempotency_key"),
  850. correlation_id=body.get("correlation_id"),
  851. )
  852. db.session.commit()
  853. return jsonify(success(result, "试生产已完成"))
  854. except DeploymentOperationError as exc:
  855. db.session.rollback()
  856. return _deployment_conflict(exc, "试生产操作状态冲突")
  857. except (TypeError, ValueError) as exc:
  858. db.session.rollback()
  859. return _deployment_conflict(exc, "试生产请求无效")
  860. except RuntimeError:
  861. db.session.rollback()
  862. return jsonify(failed("试生产执行结果异常", code=503)), 503
  863. @bp.post("/deployments/<deployment_id>/activate")
  864. def activate_dataflow_deployment(deployment_id: str):
  865. if not _activation_gate_ready():
  866. return jsonify(failed("数据工厂激活门禁尚未开放", code=503)), 503
  867. try:
  868. body = _deployment_action_body({"evidence_id"})
  869. result = _deployment_service().activate(
  870. deployment_id,
  871. body.get("evidence_id"),
  872. g.current_user["id"],
  873. reason=body.get("reason", "mass production activation"),
  874. idempotency_key=body.get("idempotency_key"),
  875. correlation_id=body.get("correlation_id"),
  876. )
  877. db.session.commit()
  878. return jsonify(success(result, "生产线已转入批量生产"))
  879. except DeploymentOperationError as exc:
  880. db.session.rollback()
  881. return _deployment_conflict(exc, "生产线激活状态冲突")
  882. except (TypeError, ValueError) as exc:
  883. db.session.rollback()
  884. return _deployment_conflict(exc, "生产线无法激活")
  885. except RuntimeError:
  886. db.session.rollback()
  887. return jsonify(failed("生产线激活结果异常", code=503)), 503
  888. @bp.post("/deployments/<deployment_id>/execute")
  889. def execute_active_dataflow_deployment(deployment_id: str):
  890. try:
  891. body = _deployment_action_body({"inputs"})
  892. result = _deployment_service().execute_active(
  893. deployment_id,
  894. body.get("inputs", {}),
  895. g.current_user["id"],
  896. reason=body.get("reason", "manual production execution"),
  897. idempotency_key=body.get("idempotency_key"),
  898. correlation_id=body.get("correlation_id"),
  899. )
  900. db.session.commit()
  901. return jsonify(success(result, "生产数据处理已完成"))
  902. except DeploymentOperationError as exc:
  903. db.session.rollback()
  904. return _deployment_conflict(exc, "生产执行操作状态冲突")
  905. except (TypeError, ValueError) as exc:
  906. db.session.rollback()
  907. return _deployment_conflict(exc, "生产执行请求无效")
  908. except RuntimeError:
  909. db.session.rollback()
  910. return jsonify(failed("生产数据处理执行异常", code=503)), 503
  911. @bp.post("/deployments/<deployment_id>/rollback")
  912. def rollback_dataflow_deployment(deployment_id: str):
  913. try:
  914. body = _deployment_action_body()
  915. result = _deployment_service().rollback(
  916. deployment_id,
  917. g.current_user["id"],
  918. reason=body.get("reason", "restore previous production line"),
  919. idempotency_key=body.get("idempotency_key"),
  920. correlation_id=body.get("correlation_id"),
  921. )
  922. db.session.commit()
  923. return jsonify(success(result, "已恢复上一条稳定生产线"))
  924. except DeploymentOperationError as exc:
  925. db.session.rollback()
  926. return _deployment_conflict(exc, "生产线回滚状态冲突")
  927. except (TypeError, ValueError) as exc:
  928. db.session.rollback()
  929. return _deployment_conflict(exc, "生产线无法回滚")
  930. except RuntimeError:
  931. db.session.rollback()
  932. return jsonify(failed("生产线回滚结果异常", code=503)), 503
  933. def _reconcile_dataflow_deployment(deployment_id: str, action: str):
  934. try:
  935. body = _closed_body({"idempotency_key"})
  936. result = _deployment_service().reconcile(
  937. deployment_id,
  938. action,
  939. g.current_user["id"],
  940. idempotency_key=body.get("idempotency_key"),
  941. )
  942. db.session.commit()
  943. return jsonify(success(result, "投产未知结果已完成对账"))
  944. except PermissionError:
  945. db.session.rollback()
  946. return jsonify(failed("权限不足", code=403)), 403
  947. except DeploymentOperationError as exc:
  948. db.session.rollback()
  949. return _deployment_conflict(exc, "投产对账状态冲突")
  950. except (TypeError, ValueError) as exc:
  951. db.session.rollback()
  952. return _deployment_conflict(exc, "投产未知结果无法对账")
  953. except RuntimeError:
  954. db.session.rollback()
  955. return jsonify(failed("投产对账服务异常", code=503)), 503
  956. @bp.post("/deployments/<deployment_id>/reconcile-deploy")
  957. def reconcile_dataflow_deploy(deployment_id: str):
  958. return _reconcile_dataflow_deployment(deployment_id, "deploy_disabled")
  959. @bp.post("/deployments/<deployment_id>/reconcile-activate")
  960. def reconcile_dataflow_activate(deployment_id: str):
  961. return _reconcile_dataflow_deployment(deployment_id, "activate")
  962. @bp.post("/deployments/<deployment_id>/reconcile-rollback")
  963. def reconcile_dataflow_rollback(deployment_id: str):
  964. return _reconcile_dataflow_deployment(deployment_id, "rollback")
  965. @bp.post("/deployments/<deployment_id>/reconcile-execute")
  966. def reconcile_dataflow_execute(deployment_id: str):
  967. return _reconcile_dataflow_deployment(deployment_id, "execute_active")