routes.py 26 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742
  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.production_line import resolve_production_line
  23. from app.core.data_rules.publication import (
  24. GenerationReceiptSigner,
  25. LogicalRuleCompiler,
  26. PhysicalPlanPublicationService,
  27. RulePublicationService,
  28. RuleValidationRejected,
  29. ServerOwnedLogicalDryRunRunner,
  30. ServerOwnedPhysicalPreflightRunner,
  31. generation_receipt_claims,
  32. )
  33. from app.core.data_rules.release import ProductionLineReleaseService
  34. from app.core.data_rules.repository import DataRuleRepository
  35. from app.core.data_rules.schema_resolver import (
  36. Neo4jSchemaMetadataCatalog,
  37. SchemaResolver,
  38. )
  39. from app.models.result import failed, success
  40. _VALIDATORS = {
  41. "rule": (validate_rule_spec, rule_spec_hash),
  42. "standard": (validate_standard_spec, standard_spec_hash),
  43. "dataflow": (validate_dataflow_spec, dataflow_spec_hash),
  44. }
  45. def _body() -> dict[str, Any]:
  46. value = request.get_json(silent=True)
  47. if not isinstance(value, dict):
  48. raise ValueError("request body must be an object")
  49. return value
  50. def _bad_request(message: str = "规则请求无效"):
  51. return jsonify(failed(message, code=400)), 400
  52. def _closed_body(allowed: set[str]) -> dict[str, Any]:
  53. body = _body()
  54. if set(body) - allowed:
  55. raise ValueError("request contains unsupported fields")
  56. return body
  57. def _repository() -> DataRuleRepository:
  58. configured = current_app.extensions.get("data_rule_repository")
  59. if configured is not None:
  60. return configured
  61. return DataRuleRepository(db.session)
  62. def _release_service() -> ProductionLineReleaseService:
  63. configured = current_app.extensions.get("production_line_release_service")
  64. if configured is not None:
  65. return configured
  66. repository = _repository()
  67. resolver = current_app.extensions.get("data_rule_schema_resolver")
  68. if resolver is None:
  69. catalog = current_app.extensions.get("data_rule_metadata_catalog")
  70. resolver = SchemaResolver(
  71. catalog or Neo4jSchemaMetadataCatalog(), repository
  72. )
  73. return ProductionLineReleaseService(repository, schema_resolver=resolver)
  74. def _receipt_signer() -> GenerationReceiptSigner:
  75. configured = current_app.extensions.get("generation_receipt_signer")
  76. if configured is not None:
  77. return configured
  78. secret = current_app.config.get("RULE_GENERATION_RECEIPT_SECRET")
  79. if not isinstance(secret, str) or not secret.strip():
  80. raise RuntimeError(
  81. "dedicated generation receipt secret is not configured"
  82. )
  83. signer = GenerationReceiptSigner(secret)
  84. current_app.extensions["generation_receipt_signer"] = signer
  85. return signer
  86. def _rule_artifact_store():
  87. configured = current_app.extensions.get("rule_artifact_store")
  88. if configured is not None:
  89. return configured
  90. from minio import Minio
  91. from app.runner.artifacts import ArtifactStore
  92. store = ArtifactStore(
  93. Minio(
  94. current_app.config["MINIO_HOST"],
  95. access_key=current_app.config["MINIO_USER"],
  96. secret_key=current_app.config["MINIO_PASSWORD"],
  97. secure=bool(current_app.config["MINIO_SECURE"]),
  98. ),
  99. bucket=current_app.config["MINIO_BUCKET"],
  100. max_artifact_bytes=32 * 1024 * 1024,
  101. max_rows=100_000,
  102. memory_limit_bytes=256 * 1024 * 1024,
  103. max_ttl_seconds=3600,
  104. )
  105. current_app.extensions["rule_artifact_store"] = store
  106. return store
  107. def _publication_service() -> RulePublicationService:
  108. configured = current_app.extensions.get("rule_publication_service")
  109. if configured is not None:
  110. return configured
  111. test_runner = current_app.extensions.get("rule_validation_test_runner")
  112. if test_runner is None:
  113. try:
  114. test_runner = ServerOwnedLogicalDryRunRunner(
  115. _rule_artifact_store()
  116. )
  117. except Exception as exc:
  118. raise RuntimeError(
  119. "trusted rule validation test runner is not configured"
  120. ) from exc
  121. current_app.extensions["rule_validation_test_runner"] = test_runner
  122. service = RulePublicationService(
  123. _repository(),
  124. receipt_signer=_receipt_signer(),
  125. compiler=LogicalRuleCompiler(),
  126. test_runner=test_runner,
  127. )
  128. current_app.extensions["rule_publication_service"] = service
  129. return service
  130. def _physical_publication_service() -> PhysicalPlanPublicationService:
  131. configured = current_app.extensions.get(
  132. "physical_plan_publication_service"
  133. )
  134. if configured is not None:
  135. return configured
  136. test_runner = current_app.extensions.get("rule_physical_test_runner")
  137. if test_runner is None:
  138. try:
  139. from app.core.data_source.runtime import get_data_source_manager
  140. test_runner = ServerOwnedPhysicalPreflightRunner(
  141. _rule_artifact_store(),
  142. datasource_manager=get_data_source_manager(),
  143. )
  144. except Exception as exc:
  145. raise RuntimeError(
  146. "trusted physical plan test runner is not configured"
  147. ) from exc
  148. current_app.extensions["rule_physical_test_runner"] = test_runner
  149. service = PhysicalPlanPublicationService(
  150. _repository(), test_runner=test_runner
  151. )
  152. current_app.extensions["physical_plan_publication_service"] = service
  153. return service
  154. def _metadata_hash(value: Any) -> str:
  155. return hashlib.sha256(
  156. json.dumps(
  157. value,
  158. sort_keys=True,
  159. separators=(",", ":"),
  160. ensure_ascii=False,
  161. ).encode("utf-8")
  162. ).hexdigest()
  163. @bp.get("/capabilities")
  164. def capabilities():
  165. return jsonify(
  166. success(
  167. {
  168. "natural_language_authoring": True,
  169. "schema_constrained_candidates": True,
  170. "production_line_preview": True,
  171. "immutable_asset_versions": True,
  172. "server_side_publishing": True,
  173. "production_line_release": True,
  174. "data_factory_activation": False,
  175. }
  176. )
  177. )
  178. @bp.post("/production-lines/draft-identity")
  179. def create_production_line_draft_identity():
  180. """Persist a short-lived, single-use, actor-bound DataFlow reservation."""
  181. try:
  182. _closed_body(set())
  183. result = _repository().reserve_dataflow_draft(
  184. actor_uid=g.current_user["id"]
  185. )
  186. db.session.commit()
  187. return (
  188. jsonify(
  189. success(
  190. result,
  191. "生产线草稿身份已创建",
  192. code=201,
  193. )
  194. ),
  195. 201,
  196. )
  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 DataFlow draft reservation failed")
  203. return jsonify(failed("生产线草稿身份暂时不可用", code=503)), 503
  204. @bp.post("/validate")
  205. def validate_asset():
  206. try:
  207. body = _closed_body({"asset_type", "spec"})
  208. asset_type = body.get("asset_type")
  209. if asset_type not in _VALIDATORS:
  210. raise ValueError("unsupported asset type")
  211. validator, hasher = _VALIDATORS[asset_type]
  212. normalized = validator(body.get("spec"))
  213. return jsonify(
  214. success(
  215. {
  216. "asset_type": asset_type,
  217. "normalized": normalized,
  218. "spec_hash": hasher(normalized),
  219. }
  220. )
  221. )
  222. except (TypeError, ValueError):
  223. return _bad_request("规则定义无效")
  224. def _authoring_agent() -> RuleAuthoringAgent:
  225. configured = current_app.extensions.get("data_rule_authoring_agent")
  226. if configured is not None:
  227. return configured
  228. api_key = current_app.config.get("LLM_API_KEY") or current_app.config.get(
  229. "DEEPSEEK_API_KEY"
  230. )
  231. if not api_key:
  232. raise RuntimeError("rule authoring model is not configured")
  233. agent = RuleAuthoringAgent(model=OpenAICompatibleRuleModel())
  234. current_app.extensions["data_rule_authoring_agent"] = agent
  235. return agent
  236. @bp.post("/interpret")
  237. def interpret_rule():
  238. try:
  239. body = _closed_body({"source_text", "authoring_surface", "context"})
  240. receipt_signer = _receipt_signer()
  241. repository = _repository()
  242. validation_context = repository.resolve_validation_context(
  243. body.get("context", {})
  244. )
  245. result = _authoring_agent().interpret(
  246. source_text=body.get("source_text"),
  247. authoring_surface=body.get("authoring_surface"),
  248. context=validation_context,
  249. )
  250. candidate = result.get("candidate")
  251. if (
  252. body.get("authoring_surface") == "data_standard"
  253. and (
  254. not isinstance(candidate, dict)
  255. or candidate.get("candidate_type") != "rule"
  256. or not isinstance(candidate.get("rule_spec"), dict)
  257. )
  258. ):
  259. raise ValueError(
  260. "data_standard authoring must produce a rule candidate"
  261. )
  262. audit = repository.record_generation_run(
  263. evidence=result,
  264. created_by=g.current_user["id"],
  265. validation_context=validation_context,
  266. )
  267. result = {
  268. **result,
  269. "generation_run_id": audit["id"],
  270. "correlation_id": audit["correlation_id"],
  271. }
  272. if (
  273. result.get("status") == "ready"
  274. and isinstance(candidate, dict)
  275. and candidate.get("candidate_type") == "rule"
  276. and isinstance(candidate.get("rule_spec"), dict)
  277. ):
  278. claims = generation_receipt_claims(
  279. generation_run_id=audit["id"],
  280. actor_uid=g.current_user["id"],
  281. source_text=result["source_text"],
  282. candidate_hash=result["candidate_hash"],
  283. rule_spec=candidate["rule_spec"],
  284. model_hash=result.get("model_hash")
  285. or _metadata_hash(
  286. {
  287. "provider": result.get("model_provider", "unknown"),
  288. "name": result.get("model_name", "unknown"),
  289. }
  290. ),
  291. prompt_hash=result.get("prompt_hash")
  292. or _metadata_hash(result.get("prompt_version", "unknown")),
  293. context_hash=result["context_hash"],
  294. expires_at=datetime.now(UTC) + timedelta(minutes=10),
  295. )
  296. result["generation_receipt"] = receipt_signer.issue(claims)
  297. db.session.commit()
  298. return jsonify(success(result))
  299. except (TypeError, ValueError):
  300. db.session.rollback()
  301. return _bad_request("自然语言规则描述无效")
  302. except RuntimeError:
  303. db.session.rollback()
  304. return jsonify(failed("AI 规则解析服务未配置", code=503)), 503
  305. except Exception:
  306. db.session.rollback()
  307. current_app.logger.exception("AI rule interpretation failed")
  308. return jsonify(failed("AI 规则解析暂时不可用", code=503)), 503
  309. @bp.post("/rule-versions")
  310. def create_rule_version():
  311. try:
  312. body = _closed_body(
  313. {
  314. "rule_spec",
  315. "source_text",
  316. "generation_receipt",
  317. "category",
  318. "source_language",
  319. "generated_kind",
  320. }
  321. )
  322. # Enforce V2 at the HTTP boundary even when a test/different
  323. # repository implementation is injected.
  324. rule_spec = validate_rule_spec(body.get("rule_spec"))
  325. result = _publication_service().create_draft(
  326. rule_spec=rule_spec,
  327. source_text=body.get("source_text"),
  328. category=body.get("category", "general"),
  329. source_language=body.get("source_language", "zh-CN"),
  330. generated_kind=body.get("generated_kind", "rulespec"),
  331. actor_uid=g.current_user["id"],
  332. generation_receipt=body.get("generation_receipt"),
  333. )
  334. db.session.commit()
  335. return jsonify(success(result, "规则版本创建成功")), 201
  336. except (TypeError, ValueError):
  337. db.session.rollback()
  338. return _bad_request("规则版本定义无效")
  339. except Exception:
  340. db.session.rollback()
  341. current_app.logger.exception("create rule version failed")
  342. return jsonify(failed("规则版本创建失败", code=500)), 500
  343. @bp.post("/rule-versions/<version_id>/publish")
  344. def publish_rule_version(version_id: str):
  345. try:
  346. if request.get_data(cache=True) and request.get_json(
  347. silent=True
  348. ) not in (
  349. None,
  350. {},
  351. ):
  352. raise ValueError("publish request must not contain fields")
  353. result = _publication_service().publish(
  354. version_id, g.current_user["id"]
  355. )
  356. db.session.commit()
  357. return jsonify(success(result, "规则版本发布成功"))
  358. except (TypeError, ValueError):
  359. db.session.rollback()
  360. return jsonify(failed("规则版本无法发布", code=409)), 409
  361. except Exception:
  362. db.session.rollback()
  363. current_app.logger.exception("publish rule version failed")
  364. return jsonify(failed("规则版本发布失败", code=500)), 500
  365. @bp.post("/rule-versions/<version_id>/validate")
  366. def validate_rule_version(version_id: str):
  367. try:
  368. if request.get_data(cache=True) and request.get_json(
  369. silent=True
  370. ) not in (
  371. None,
  372. {},
  373. ):
  374. raise ValueError("validate request must not contain fields")
  375. result = _publication_service().validate(
  376. version_id, g.current_user["id"]
  377. )
  378. db.session.commit()
  379. return jsonify(success(result, "规则版本编译验证成功"))
  380. except RuleValidationRejected:
  381. db.session.commit()
  382. return jsonify(failed("规则版本编译验证失败", code=409)), 409
  383. except (TypeError, ValueError):
  384. db.session.rollback()
  385. return jsonify(failed("规则版本编译验证失败", code=409)), 409
  386. except RuntimeError:
  387. db.session.rollback()
  388. return jsonify(failed("规则验证服务未配置", code=503)), 503
  389. @bp.post("/rule-versions/<version_id>/test")
  390. def test_rule_version(version_id: str):
  391. try:
  392. body = _closed_body({"plan_id"})
  393. result = _publication_service().test(
  394. version_id,
  395. g.current_user["id"],
  396. plan_id=body.get("plan_id"),
  397. )
  398. db.session.commit()
  399. return jsonify(success(result, "规则版本样本测试成功"))
  400. except (TypeError, ValueError):
  401. db.session.rollback()
  402. return jsonify(failed("规则版本样本测试失败", code=409)), 409
  403. except RuntimeError:
  404. db.session.rollback()
  405. return jsonify(failed("规则测试服务未配置", code=503)), 503
  406. @bp.get("/rule-versions/<version_id>/evidence")
  407. def rule_version_evidence(version_id: str):
  408. try:
  409. return jsonify(
  410. success(
  411. _repository().get_asset_evidence(
  412. asset_type="rule", version_id=version_id
  413. )
  414. )
  415. )
  416. except (TypeError, ValueError):
  417. return jsonify(failed("规则证据不存在", code=404)), 404
  418. except Exception:
  419. current_app.logger.exception("load rule evidence failed")
  420. return jsonify(failed("规则证据暂时不可用", code=503)), 503
  421. @bp.get("/catalog/assets/<asset_type>/<version_id>/evidence")
  422. def catalog_asset_evidence(asset_type: str, version_id: str):
  423. try:
  424. return jsonify(
  425. success(
  426. _repository().get_asset_evidence(
  427. asset_type=asset_type,
  428. version_id=version_id,
  429. )
  430. )
  431. )
  432. except (TypeError, ValueError):
  433. return jsonify(failed("资产证据不存在", code=404)), 404
  434. except Exception:
  435. current_app.logger.exception("load catalog asset evidence failed")
  436. return jsonify(failed("资产证据暂时不可用", code=503)), 503
  437. def _catalog_schema_context():
  438. raw_inputs = request.args.get("input_schema_refs")
  439. output = request.args.get("output_schema_ref")
  440. if raw_inputs is None and output is None:
  441. return None, None
  442. if raw_inputs is None or output is None:
  443. raise ValueError("catalog schema context is incomplete")
  444. inputs = json.loads(raw_inputs)
  445. if not isinstance(inputs, list):
  446. raise ValueError("catalog input_schema_refs must be an array")
  447. return inputs, output
  448. @bp.get("/catalog/assets/<asset_type>/<version_id>")
  449. def catalog_asset(asset_type: str, version_id: str):
  450. try:
  451. if set(request.args) - {
  452. "input_schema_refs",
  453. "output_schema_ref",
  454. }:
  455. raise ValueError("catalog asset query contains unsupported fields")
  456. inputs, output = _catalog_schema_context()
  457. return jsonify(
  458. success(
  459. _repository().get_published_asset(
  460. asset_type=asset_type,
  461. version_id=version_id,
  462. input_schema_refs=inputs,
  463. output_schema_ref=output,
  464. )
  465. )
  466. )
  467. except (TypeError, ValueError, json.JSONDecodeError):
  468. return jsonify(failed("已发布资产不存在或上下文无效", code=404)), 404
  469. except Exception:
  470. current_app.logger.exception("load exact catalog asset failed")
  471. return jsonify(failed("规则目录暂时不可用", code=503)), 503
  472. @bp.get("/catalog")
  473. @bp.get("/catalog/rule-versions")
  474. def published_rule_catalog():
  475. try:
  476. legacy_rule_alias = request.path.endswith("/rule-versions")
  477. allowed = {
  478. "query",
  479. "limit",
  480. "offset",
  481. "input_schema_refs",
  482. "output_schema_ref",
  483. }
  484. if not legacy_rule_alias:
  485. allowed.add("asset_type")
  486. if set(request.args) - allowed:
  487. raise ValueError("catalog query contains unsupported fields")
  488. query = request.args.get("query", "")
  489. limit = int(request.args.get("limit", "50"))
  490. offset = int(request.args.get("offset", "0"))
  491. asset_type = (
  492. "rule"
  493. if legacy_rule_alias
  494. else request.args.get("asset_type") or None
  495. )
  496. if asset_type not in {None, "rule", "standard"}:
  497. raise ValueError("catalog asset_type is invalid")
  498. if len(query) > 200 or limit < 1 or limit > 100:
  499. raise ValueError("catalog bounds are invalid")
  500. if offset < 0 or offset > 1_000_000:
  501. raise ValueError("catalog offset is invalid")
  502. inputs, output = _catalog_schema_context()
  503. catalog_args = {
  504. "query": query,
  505. "asset_type": asset_type,
  506. "limit": limit,
  507. "offset": offset,
  508. }
  509. if inputs is not None:
  510. catalog_args.update(
  511. {
  512. "input_schema_refs": inputs,
  513. "output_schema_ref": output,
  514. }
  515. )
  516. return jsonify(
  517. success(_repository().search_published_assets(**catalog_args))
  518. )
  519. except (TypeError, ValueError):
  520. return _bad_request("规则目录查询无效")
  521. except Exception:
  522. current_app.logger.exception("load rule catalog failed")
  523. return jsonify(failed("规则目录暂时不可用", code=503)), 503
  524. @bp.post("/execution-plans/<plan_id>/validate")
  525. def validate_physical_plan(plan_id: str):
  526. try:
  527. if request.get_data(cache=True) and request.get_json(
  528. silent=True
  529. ) not in (
  530. None,
  531. {},
  532. ):
  533. raise ValueError("validate request must not contain fields")
  534. result = _physical_publication_service().validate(
  535. plan_id, g.current_user["id"]
  536. )
  537. db.session.commit()
  538. return jsonify(success(result, "物理执行计划编译证据已确认"))
  539. except (TypeError, ValueError):
  540. db.session.rollback()
  541. return jsonify(failed("物理执行计划验证失败", code=409)), 409
  542. except RuntimeError:
  543. db.session.rollback()
  544. return jsonify(failed("物理计划验证服务未配置", code=503)), 503
  545. @bp.post("/execution-plans/<plan_id>/test")
  546. def test_physical_plan(plan_id: str):
  547. try:
  548. if request.get_data(cache=True) and request.get_json(
  549. silent=True
  550. ) not in (
  551. None,
  552. {},
  553. ):
  554. raise ValueError("test request must not contain fields")
  555. result = _physical_publication_service().test(
  556. plan_id, g.current_user["id"]
  557. )
  558. db.session.commit()
  559. return jsonify(success(result, "物理执行计划样本测试成功"))
  560. except (TypeError, ValueError):
  561. db.session.rollback()
  562. return jsonify(failed("物理执行计划样本测试失败", code=409)), 409
  563. except RuntimeError:
  564. db.session.rollback()
  565. return jsonify(failed("物理计划测试服务未配置", code=503)), 503
  566. @bp.post("/execution-plans/<plan_id>/publish")
  567. def publish_physical_plan(plan_id: str):
  568. try:
  569. if request.get_data(cache=True) and request.get_json(
  570. silent=True
  571. ) not in (
  572. None,
  573. {},
  574. ):
  575. raise ValueError("publish request must not contain fields")
  576. result = _physical_publication_service().publish(
  577. plan_id, g.current_user["id"]
  578. )
  579. db.session.commit()
  580. return jsonify(success(result, "物理执行计划发布成功"))
  581. except (TypeError, ValueError):
  582. db.session.rollback()
  583. return jsonify(failed("物理执行计划无法发布", code=409)), 409
  584. except RuntimeError:
  585. db.session.rollback()
  586. return jsonify(failed("物理计划发布服务未配置", code=503)), 503
  587. @bp.post("/standard-versions")
  588. def create_standard_version():
  589. try:
  590. body = _closed_body({"standard_spec", "source_text"})
  591. result = _repository().create_standard_version(
  592. standard_spec=body.get("standard_spec"),
  593. source_text=body.get("source_text"),
  594. created_by=g.current_user["id"],
  595. )
  596. db.session.commit()
  597. return jsonify(success(result, "数据标准版本创建成功")), 201
  598. except (TypeError, ValueError):
  599. db.session.rollback()
  600. return _bad_request("数据标准版本定义无效")
  601. except Exception:
  602. db.session.rollback()
  603. current_app.logger.exception("create standard version failed")
  604. return jsonify(failed("数据标准版本创建失败", code=500)), 500
  605. @bp.post("/standard-versions/<version_id>/publish")
  606. def publish_standard_version(version_id: str):
  607. try:
  608. if request.get_data(cache=True) and request.get_json(
  609. silent=True
  610. ) not in (
  611. None,
  612. {},
  613. ):
  614. raise ValueError("publish request must not contain fields")
  615. result = _repository().publish_standard_version(
  616. version_id=version_id,
  617. published_by=g.current_user["id"],
  618. )
  619. db.session.commit()
  620. return jsonify(success(result, "数据标准版本发布成功"))
  621. except (TypeError, ValueError):
  622. db.session.rollback()
  623. return jsonify(failed("数据标准版本无法发布", code=409)), 409
  624. except Exception:
  625. db.session.rollback()
  626. current_app.logger.exception("publish standard version failed")
  627. return jsonify(failed("数据标准版本发布失败", code=500)), 500
  628. @bp.post("/production-lines/resolve")
  629. def resolve_production_line_preview():
  630. try:
  631. body = _body()
  632. package = resolve_production_line(
  633. body.get("dataflow_spec"),
  634. body.get("standard_versions"),
  635. body.get("rule_versions"),
  636. component_binding_ids=body.get("component_binding_ids"),
  637. )
  638. return jsonify(
  639. success(
  640. {
  641. "preview": True,
  642. "release_ready": False,
  643. "package": package,
  644. }
  645. )
  646. )
  647. except (TypeError, ValueError):
  648. return _bad_request("数据生产线定义无效")
  649. @bp.post("/production-lines/<dataflow_uid>/release")
  650. def release_production_line(dataflow_uid: str):
  651. try:
  652. body = _closed_body(
  653. {
  654. "dataflow_spec",
  655. "source_text",
  656. }
  657. )
  658. result = _release_service().release(
  659. dataflow_uid=dataflow_uid,
  660. dataflow_spec=body.get("dataflow_spec"),
  661. source_text=body.get("source_text"),
  662. created_by=g.current_user["id"],
  663. )
  664. db.session.commit()
  665. return jsonify(success(result, "数据生产线发布成功")), 201
  666. except (TypeError, ValueError):
  667. db.session.rollback()
  668. return jsonify(failed("数据生产线无法发布", code=409)), 409
  669. except Exception:
  670. db.session.rollback()
  671. current_app.logger.exception("release production line failed")
  672. return jsonify(failed("数据生产线发布失败", code=500)), 500