routes.py 26 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822
  1. """HTTP orchestration boundary for data-research ingestion jobs."""
  2. from __future__ import annotations
  3. import io
  4. import logging
  5. from flask import current_app, g, jsonify, request, send_file
  6. from app import db
  7. from app.api.data_development import bp
  8. from app.core.data_research.errors import DataResearchError
  9. from app.models.result import failed, success
  10. logger = logging.getLogger(__name__)
  11. def get_ingestion_service():
  12. from app.core.data_research.ingestion import IngestionService
  13. from app.core.data_research.repository import SqlAlchemyIngestionJobRepository
  14. return IngestionService(
  15. SqlAlchemyIngestionJobRepository(db.session),
  16. commit=db.session.commit,
  17. rollback=db.session.rollback,
  18. )
  19. def get_database_source_registration_service():
  20. from app.core.data_research.repository import (
  21. SqlAlchemyIngestionSourceRepository,
  22. )
  23. from app.core.data_research.sources import (
  24. DatabaseSourceRegistrationService,
  25. )
  26. from app.core.data_source.runtime import get_data_source_manager
  27. manager = get_data_source_manager()
  28. return DatabaseSourceRegistrationService(
  29. SqlAlchemyIngestionSourceRepository(db.session),
  30. definition_resolver=manager.definitions.get,
  31. commit=db.session.commit,
  32. rollback=db.session.rollback,
  33. )
  34. def get_catalog_snapshot_repository():
  35. from app.core.data_research.repository import (
  36. SqlAlchemyCatalogSnapshotRepository,
  37. )
  38. return SqlAlchemyCatalogSnapshotRepository(db.session)
  39. def get_catalog_ingestion_executor():
  40. from app.core.data_research.catalog.execution import (
  41. CatalogIngestionExecutor,
  42. )
  43. from app.core.data_research.catalog.mysql import MySqlCatalogCollector
  44. from app.core.data_research.catalog.postgresql import (
  45. PostgreSqlCatalogCollector,
  46. )
  47. from app.core.data_research.catalog.service import CatalogCollectionService
  48. from app.core.data_research.errors import IngestionSourceInvalid
  49. from app.core.data_source.runtime import get_data_source_manager
  50. manager = get_data_source_manager()
  51. def collector_resolver(database_type):
  52. if database_type == "postgresql":
  53. return PostgreSqlCatalogCollector()
  54. if database_type == "mysql":
  55. return MySqlCatalogCollector()
  56. raise IngestionSourceInvalid(
  57. f"database type {database_type} is not supported"
  58. )
  59. collector = CatalogCollectionService(
  60. manager,
  61. definition_resolver=manager.definitions.get,
  62. collector_resolver=collector_resolver,
  63. )
  64. return CatalogIngestionExecutor(
  65. get_ingestion_service(),
  66. collector,
  67. get_catalog_snapshot_repository(),
  68. commit=db.session.commit,
  69. rollback=db.session.rollback,
  70. )
  71. def get_data_element_service():
  72. from app.core.data_research.data_elements import DataElementService
  73. from app.core.data_research.repository import SqlAlchemyDataElementRepository
  74. from app.core.events.outbox import enqueue_outbox
  75. return DataElementService(
  76. SqlAlchemyDataElementRepository(db.session),
  77. outbox_enqueue=lambda **event: enqueue_outbox(db.session, **event),
  78. )
  79. def get_candidate_decision_service():
  80. from app.core.data_research.candidate_decisions import CandidateDecisionService
  81. from app.core.data_research.data_elements import DataElementService
  82. from app.core.data_research.repository import (
  83. SqlAlchemyCandidateDecisionRepository,
  84. SqlAlchemyDataElementRepository,
  85. )
  86. return CandidateDecisionService(
  87. SqlAlchemyCandidateDecisionRepository(db.session),
  88. data_elements=DataElementService(SqlAlchemyDataElementRepository(db.session)),
  89. )
  90. def _artifact_storage():
  91. from minio import Minio
  92. from app.core.data_research.artifacts import MinioArtifactStorage
  93. client = 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.get("MINIO_SECURE")),
  98. )
  99. return MinioArtifactStorage(client, current_app.config["MINIO_BUCKET"])
  100. def get_artifact_service():
  101. from app.core.common.identifiers import new_governance_uid
  102. from app.core.data_research.artifacts import (
  103. ArtifactService,
  104. SqlAlchemyArtifactRepository,
  105. )
  106. from app.core.data_research.file_policy import FilePolicy
  107. return ArtifactService(
  108. SqlAlchemyArtifactRepository(db.session),
  109. _artifact_storage(),
  110. FilePolicy(
  111. max_bytes=int(current_app.config.get("DATA_RESEARCH_FILE_MAX_BYTES", 25 * 1024 * 1024)),
  112. max_pages=int(current_app.config.get("DATA_RESEARCH_FILE_MAX_PAGES", 500)),
  113. ),
  114. uid_factory=new_governance_uid,
  115. )
  116. def get_evidence_service():
  117. from app.core.data_research.artifacts import EvidenceService
  118. return EvidenceService(db.session, _artifact_storage())
  119. def get_ontology_service():
  120. from app.core.data_research.ontology.publication import (
  121. OntologyApplicationService,
  122. OntologyPublicationService,
  123. )
  124. from app.core.data_research.ontology.repository import SqlAlchemyOntologyRepository
  125. from app.core.events.outbox import enqueue_outbox
  126. repository = SqlAlchemyOntologyRepository(db.session)
  127. publication = OntologyPublicationService(
  128. repository,
  129. outbox_enqueue=lambda **event: enqueue_outbox(db.session, **event),
  130. commit=db.session.commit,
  131. rollback=db.session.rollback,
  132. )
  133. return OntologyApplicationService(
  134. repository,
  135. publication=publication,
  136. commit=db.session.commit,
  137. rollback=db.session.rollback,
  138. )
  139. def get_ontology_dynamic_service():
  140. from app.core.data_research.ontology.change_sets import (
  141. SqlAlchemyDynamicOntologyService,
  142. )
  143. return SqlAlchemyDynamicOntologyService(db.session)
  144. def get_ontology_exchange_service():
  145. from app.core.data_research.ontology.exchange import OntologyExchangeService
  146. from app.core.data_research.ontology.repository import SqlAlchemyOntologyRepository
  147. return OntologyExchangeService(
  148. SqlAlchemyOntologyRepository(db.session),
  149. commit=db.session.commit,
  150. rollback=db.session.rollback,
  151. )
  152. def get_semantic_query_service():
  153. from app.core.data_research.ontology.query import (
  154. Neo4jSemanticRepository,
  155. SemanticQueryService,
  156. )
  157. from app.services.neo4j_driver import neo4j_driver
  158. return SemanticQueryService(Neo4jSemanticRepository(neo4j_driver))
  159. def _identity():
  160. return getattr(g, "current_user", {}) or {}
  161. def _record(record):
  162. return {
  163. "uid": str(record.uid),
  164. "source_uid": str(record.source_uid),
  165. "artifact_uid": str(record.artifact_uid) if record.artifact_uid else None,
  166. "job_type": record.job_type,
  167. "parser_version": record.parser_version,
  168. "status": record.status,
  169. "parameters": dict(record.parameters or {}),
  170. "statistics": dict(record.statistics or {}),
  171. "last_error": record.last_error,
  172. "attempt_count": int(record.attempt_count or 0),
  173. "failure_stage": record.failure_stage,
  174. "actor_uid": record.actor_uid,
  175. "created_at": record.created_at.isoformat() if record.created_at else None,
  176. "updated_at": record.updated_at.isoformat() if record.updated_at else None,
  177. "started_at": record.started_at.isoformat() if record.started_at else None,
  178. "finished_at": record.finished_at.isoformat() if record.finished_at else None,
  179. }
  180. def _catalog_snapshot(record):
  181. return {
  182. "uid": str(record.uid),
  183. "job_uid": str(record.job_uid),
  184. "source_uid": str(record.source_uid),
  185. "attempt": int(record.attempt),
  186. "database_type": record.database_type,
  187. "content_hash": record.content_hash,
  188. "snapshot": dict(record.snapshot or {}),
  189. "evidence_count": int(record.evidence_count),
  190. "created_at": (
  191. record.created_at.isoformat()
  192. if record.created_at
  193. else None
  194. ),
  195. }
  196. def _element(record):
  197. return {
  198. "uid": str(record.uid),
  199. "code": record.code,
  200. "status": record.status,
  201. "current_version": int(record.current_version),
  202. "snapshot": dict(record.snapshot or {}),
  203. "created_by": record.created_by,
  204. "updated_by": record.updated_by,
  205. }
  206. def _decision(record):
  207. return {
  208. "uid": str(record.uid),
  209. "candidate_uid": str(record.candidate_uid),
  210. "action": record.action,
  211. "data_element_uid": (
  212. str(record.data_element_uid) if record.data_element_uid else None
  213. ),
  214. "evidence_uids": list(record.evidence_uids),
  215. "actor_uid": record.actor_uid,
  216. "reason": record.reason,
  217. }
  218. def _artifact(record):
  219. return {
  220. "uid": str(record.uid),
  221. "source_uid": str(record.source_uid),
  222. "filename": record.filename,
  223. "media_type": record.media_type,
  224. "size_bytes": int(record.size_bytes),
  225. "content_hash": record.content_hash,
  226. "parser_version": record.parser_version,
  227. }
  228. def _ontology(record):
  229. return {
  230. "uid": str(record.uid),
  231. "code": record.code,
  232. "name": record.name,
  233. "owner_uid": record.owner_uid,
  234. "status": record.status,
  235. "draft_revision": int(record.draft_revision),
  236. "active_version_uid": record.active_version_uid,
  237. "domain_links": [
  238. {"domain_uid": link.domain_uid, "role": link.role}
  239. for link in record.domain_links
  240. ],
  241. }
  242. def _ontology_version(record):
  243. return {
  244. "uid": str(record.uid),
  245. "ontology_uid": str(record.ontology_uid),
  246. "version": int(record.version),
  247. "parent_version_uid": record.parent_version_uid,
  248. "status": record.status,
  249. "content_hash": record.content_hash,
  250. "graph_document": record.graph_document.to_dict(),
  251. "created_by": record.created_by,
  252. }
  253. def _error(error):
  254. if isinstance(error, DataResearchError):
  255. return (
  256. jsonify(
  257. failed(
  258. str(error),
  259. code=error.http_status,
  260. error={"code": error.code},
  261. )
  262. ),
  263. error.http_status,
  264. )
  265. logger.exception("data-research ingestion request failed")
  266. return (
  267. jsonify(
  268. failed(
  269. "数据采集任务处理失败",
  270. code=500,
  271. error={"code": "DATA_RESEARCH_ERROR"},
  272. )
  273. ),
  274. 500,
  275. )
  276. @bp.route("/ingestion-jobs", methods=["POST"])
  277. def create_ingestion_job():
  278. payload = request.get_json(silent=True) or {}
  279. try:
  280. actor_uid = _identity().get("id") or _identity().get("sub")
  281. if payload.get("job_type") == "catalog_collect":
  282. get_database_source_registration_service().ensure(
  283. payload.get("source_uid"),
  284. actor_uid=actor_uid,
  285. )
  286. record, created = get_ingestion_service().create_job(
  287. payload,
  288. actor_uid=actor_uid,
  289. )
  290. return jsonify(success(_record(record))), 201 if created else 200
  291. except Exception as error:
  292. return _error(error)
  293. @bp.route("/ingestion-jobs", methods=["GET"])
  294. def list_ingestion_jobs():
  295. filters = {
  296. name: request.args.get(name)
  297. for name in ("status", "source_uid")
  298. if request.args.get(name)
  299. }
  300. try:
  301. records = get_ingestion_service().list_jobs(filters)
  302. return jsonify(
  303. success({"records": [_record(item) for item in records], "total": len(records)})
  304. ), 200
  305. except Exception as error:
  306. return _error(error)
  307. @bp.route("/ingestion-jobs/<job_uid>", methods=["GET"])
  308. def get_ingestion_job(job_uid):
  309. try:
  310. return jsonify(success(_record(get_ingestion_service().get_job(job_uid)))), 200
  311. except Exception as error:
  312. return _error(error)
  313. @bp.route("/ingestion-jobs/<job_uid>/execute", methods=["POST"])
  314. def execute_ingestion_job(job_uid):
  315. try:
  316. record = get_catalog_ingestion_executor().execute(job_uid)
  317. return jsonify(success(_record(record))), 200
  318. except Exception as error:
  319. return _error(error)
  320. @bp.route(
  321. "/ingestion-jobs/<job_uid>/catalog-snapshots",
  322. methods=["GET"],
  323. )
  324. def list_catalog_snapshots(job_uid):
  325. try:
  326. records = get_catalog_snapshot_repository().list(job_uid)
  327. return jsonify(
  328. success(
  329. {
  330. "records": [
  331. _catalog_snapshot(record)
  332. for record in records
  333. ],
  334. "total": len(records),
  335. }
  336. )
  337. ), 200
  338. except Exception as error:
  339. return _error(error)
  340. @bp.route("/ingestion-jobs/<job_uid>/evidence", methods=["GET"])
  341. def list_ingestion_job_evidence(job_uid):
  342. try:
  343. records = get_evidence_service().list_for_job(job_uid)
  344. return jsonify(
  345. success({"records": records, "total": len(records)})
  346. ), 200
  347. except Exception as error:
  348. return _error(error)
  349. @bp.route("/ingestion-jobs/<job_uid>/retry", methods=["POST"])
  350. def retry_ingestion_job(job_uid):
  351. try:
  352. return jsonify(success(_record(get_ingestion_service().retry(job_uid)))), 200
  353. except Exception as error:
  354. return _error(error)
  355. @bp.route("/ingestion-jobs/<job_uid>/cancel", methods=["POST"])
  356. def cancel_ingestion_job(job_uid):
  357. try:
  358. service = get_ingestion_service()
  359. record = service.get_job(job_uid)
  360. identity = _identity()
  361. permissions = set(identity.get("permissions") or [])
  362. actor_uid = identity.get("id") or identity.get("sub")
  363. if record.actor_uid != actor_uid and "ingestion:admin" not in permissions:
  364. return jsonify(failed("权限不足", code=403)), 403
  365. return jsonify(success(_record(service.cancel(job_uid)))), 200
  366. except Exception as error:
  367. return _error(error)
  368. @bp.route("/data-elements", methods=["GET"])
  369. def list_data_elements():
  370. try:
  371. records = get_data_element_service().list(
  372. status=request.args.get("status"),
  373. business_domain_uid=request.args.get("business_domain_uid"),
  374. )
  375. return jsonify(
  376. success(
  377. {
  378. "records": [_element(item) for item in records],
  379. "total": len(records),
  380. }
  381. )
  382. ), 200
  383. except Exception as error:
  384. return _error(error)
  385. @bp.route("/data-elements", methods=["POST"])
  386. def create_data_element():
  387. try:
  388. record = get_data_element_service().create_draft(
  389. request.get_json(silent=True) or {},
  390. actor_uid=_identity().get("id") or _identity().get("sub"),
  391. )
  392. db.session.commit()
  393. return jsonify(success(_element(record))), 201
  394. except Exception as error:
  395. db.session.rollback()
  396. return _error(error)
  397. @bp.route("/data-elements/<element_uid>/transitions", methods=["POST"])
  398. def transition_data_element(element_uid):
  399. payload = request.get_json(silent=True) or {}
  400. identity = _identity()
  401. target_status = str(payload.get("target_status") or "")
  402. if (
  403. target_status in {"published", "deprecated", "retired"}
  404. and "data-elements:publish" not in set(identity.get("permissions") or [])
  405. ):
  406. return jsonify(failed("权限不足", code=403)), 403
  407. try:
  408. record = get_data_element_service().transition(
  409. element_uid,
  410. target_status,
  411. expected_version=int(payload.get("expected_version")),
  412. actor_uid=identity.get("id") or identity.get("sub"),
  413. )
  414. db.session.commit()
  415. return jsonify(success(_element(record))), 200
  416. except Exception as error:
  417. db.session.rollback()
  418. return _error(error)
  419. @bp.route("/candidate-decisions", methods=["POST"])
  420. def decide_candidates():
  421. payload = request.get_json(silent=True) or {}
  422. try:
  423. records = get_candidate_decision_service().decide(
  424. payload.get("decisions") or [],
  425. actor_uid=_identity().get("id") or _identity().get("sub"),
  426. )
  427. return jsonify(success([_decision(record) for record in records])), 200
  428. except Exception as error:
  429. return _error(error)
  430. @bp.route("/sources/files", methods=["POST"])
  431. def upload_source_file():
  432. uploaded = request.files.get("file")
  433. if uploaded is None:
  434. return jsonify(failed("缺少上传文件", code=400)), 400
  435. try:
  436. record, created = get_artifact_service().store(
  437. request.form.get("source_uid"),
  438. uploaded.filename,
  439. uploaded.mimetype,
  440. uploaded.read(),
  441. request.form.get("parser_version") or "auto-v1",
  442. )
  443. db.session.commit()
  444. return jsonify(success(_artifact(record))), 201 if created else 200
  445. except Exception as error:
  446. db.session.rollback()
  447. return _error(error)
  448. @bp.route("/evidence/<evidence_uid>", methods=["GET"])
  449. def get_evidence(evidence_uid):
  450. identity = _identity()
  451. try:
  452. service = get_evidence_service()
  453. if request.args.get("download") in {"1", "true", "yes"}:
  454. if "evidence:download" not in set(identity.get("permissions") or []):
  455. return jsonify(failed("权限不足", code=403)), 403
  456. content, filename, media_type = service.download(evidence_uid)
  457. return send_file(
  458. io.BytesIO(content),
  459. mimetype=media_type,
  460. as_attachment=True,
  461. download_name=filename,
  462. )
  463. preview = service.preview(evidence_uid)
  464. from app.core.data_research.artifacts import redact_excerpt
  465. preview["excerpt"] = redact_excerpt(preview.get("excerpt"))
  466. return jsonify(success(preview)), 200
  467. except Exception as error:
  468. return _error(error)
  469. @bp.route("/ontologies", methods=["GET"])
  470. def list_ontologies():
  471. try:
  472. return jsonify(success([_ontology(item) for item in get_ontology_service().list()])), 200
  473. except Exception as error:
  474. return _error(error)
  475. @bp.route("/ontologies", methods=["POST"])
  476. def create_ontology():
  477. try:
  478. record = get_ontology_service().create(
  479. request.get_json(silent=True) or {},
  480. _identity().get("id") or _identity().get("sub"),
  481. )
  482. return jsonify(success(_ontology(record))), 201
  483. except Exception as error:
  484. return _error(error)
  485. @bp.route("/ontologies/<ontology_uid>", methods=["GET"])
  486. def get_ontology(ontology_uid):
  487. try:
  488. return jsonify(success(_ontology(get_ontology_service().get(ontology_uid)))), 200
  489. except Exception as error:
  490. return _error(error)
  491. @bp.route("/ontologies/<ontology_uid>/versions", methods=["GET"])
  492. def list_ontology_versions(ontology_uid):
  493. try:
  494. return jsonify(
  495. success(
  496. [
  497. _ontology_version(item)
  498. for item in get_ontology_service().list_versions(ontology_uid)
  499. ]
  500. )
  501. ), 200
  502. except Exception as error:
  503. return _error(error)
  504. @bp.route("/ontologies/<ontology_uid>/graph", methods=["GET"])
  505. def get_ontology_graph(ontology_uid):
  506. try:
  507. service = get_ontology_service()
  508. ontology = service.get(ontology_uid)
  509. version = service.latest_version(ontology_uid)
  510. data = (
  511. _ontology_version(version)
  512. if version is not None
  513. else {
  514. "uid": None,
  515. "ontology_uid": str(ontology_uid),
  516. "version": 0,
  517. "parent_version_uid": None,
  518. "status": "draft",
  519. "content_hash": None,
  520. "graph_document": {
  521. name: []
  522. for name in (
  523. "classes",
  524. "properties",
  525. "relations",
  526. "constraints",
  527. "domain_links",
  528. "element_mappings",
  529. )
  530. },
  531. "created_by": None,
  532. }
  533. )
  534. response = jsonify(success(data))
  535. response.headers["ETag"] = f'"{int(ontology.draft_revision)}"'
  536. return response, 200
  537. except Exception as error:
  538. return _error(error)
  539. @bp.route("/ontologies/<ontology_uid>/graph", methods=["PATCH"])
  540. def save_ontology_graph(ontology_uid):
  541. etag = str(request.headers.get("If-Match") or "").strip().strip('"')
  542. if not etag.isdigit():
  543. return jsonify(failed("缺少有效 If-Match 修订号", code=428)), 428
  544. try:
  545. version = get_ontology_service().save_draft(
  546. ontology_uid,
  547. request.get_json(silent=True) or {},
  548. int(etag),
  549. _identity().get("id") or _identity().get("sub"),
  550. )
  551. response = jsonify(success(_ontology_version(version)))
  552. response.headers["ETag"] = f'"{int(etag) + 1}"'
  553. return response, 200
  554. except Exception as error:
  555. return _error(error)
  556. @bp.route("/ontologies/<ontology_uid>/validate", methods=["POST"])
  557. def validate_ontology(ontology_uid):
  558. try:
  559. issues = get_ontology_service().validate(ontology_uid)
  560. data = [
  561. {"code": item.code, "message": item.message, "path": item.path}
  562. for item in issues
  563. ]
  564. return jsonify(success(data)), 200
  565. except Exception as error:
  566. return _error(error)
  567. @bp.route("/ontologies/<ontology_uid>/publish", methods=["POST"])
  568. def publish_ontology(ontology_uid):
  569. try:
  570. version = get_ontology_service().publish(
  571. ontology_uid,
  572. request.headers.get("Idempotency-Key") or f"{ontology_uid}:publish",
  573. _identity().get("id") or _identity().get("sub"),
  574. )
  575. return jsonify(success(_ontology_version(version))), 200
  576. except Exception as error:
  577. return _error(error)
  578. @bp.route("/ontologies/<ontology_uid>/diff", methods=["GET"])
  579. def diff_ontology(ontology_uid):
  580. try:
  581. return jsonify(success(get_ontology_service().diff(
  582. ontology_uid, request.args.get("left"), request.args.get("right")
  583. ))), 200
  584. except Exception as error:
  585. return _error(error)
  586. @bp.route("/ontologies/<ontology_uid>/rollback", methods=["POST"])
  587. def rollback_ontology(ontology_uid):
  588. payload = request.get_json(silent=True) or {}
  589. try:
  590. version = get_ontology_service().rollback(
  591. ontology_uid,
  592. payload.get("target_version_uid"),
  593. int(payload.get("expected_revision")),
  594. _identity().get("id") or _identity().get("sub"),
  595. )
  596. return jsonify(success(_ontology_version(version))), 201
  597. except Exception as error:
  598. return _error(error)
  599. @bp.route("/ontologies/<ontology_uid>/suggestions", methods=["POST"])
  600. def generate_ontology_suggestions(ontology_uid):
  601. try:
  602. result = get_ontology_dynamic_service().generate(
  603. ontology_uid,
  604. request.get_json(silent=True) or {},
  605. _identity().get("id") or _identity().get("sub"),
  606. )
  607. records = (
  608. result.get("suggestions") or []
  609. if isinstance(result, dict)
  610. else result
  611. )
  612. suggestions = [
  613. item
  614. if isinstance(item, dict)
  615. else {
  616. "uid": item.uid,
  617. "kind": item.kind,
  618. "payload": item.payload,
  619. "evidence_uids": list(item.evidence_uids),
  620. "confidence": item.confidence,
  621. "source": item.source,
  622. "model_version": item.model_version,
  623. "prompt_version": item.prompt_version,
  624. }
  625. for item in records
  626. ]
  627. data = {
  628. "change_set_uid": (
  629. result.get("change_set_uid")
  630. if isinstance(result, dict)
  631. else None
  632. ),
  633. "suggestions": suggestions,
  634. }
  635. return jsonify(success(data)), 201
  636. except Exception as error:
  637. return _error(error)
  638. @bp.route(
  639. "/ontologies/<ontology_uid>/change-sets/<change_set_uid>/decisions",
  640. methods=["POST"],
  641. )
  642. def decide_ontology_change_set(ontology_uid, change_set_uid):
  643. payload = request.get_json(silent=True) or {}
  644. try:
  645. result = get_ontology_dynamic_service().decide(
  646. ontology_uid,
  647. change_set_uid,
  648. payload.get("decisions") or [],
  649. _identity().get("id") or _identity().get("sub"),
  650. )
  651. return jsonify(success(result)), 200
  652. except Exception as error:
  653. return _error(error)
  654. @bp.route("/ontologies/<ontology_uid>/export", methods=["GET"])
  655. def export_ontology(ontology_uid):
  656. try:
  657. content, media_type, filename = get_ontology_exchange_service().export(
  658. ontology_uid, request.args.get("format") or "json"
  659. )
  660. return send_file(
  661. io.BytesIO(content),
  662. mimetype=media_type,
  663. as_attachment=True,
  664. download_name=filename,
  665. )
  666. except Exception as error:
  667. return _error(error)
  668. @bp.route("/ontologies/import", methods=["POST"])
  669. def import_ontology():
  670. try:
  671. uploaded = request.files.get("file")
  672. content = uploaded.read() if uploaded is not None else request.get_data(cache=False)
  673. result = get_ontology_exchange_service().import_document(
  674. content,
  675. request.args.get("format") or "json",
  676. _identity().get("id") or _identity().get("sub"),
  677. )
  678. return jsonify(success(result)), 201
  679. except Exception as error:
  680. return _error(error)
  681. @bp.route("/semantic/properties/<property_uid>", methods=["GET"])
  682. def query_semantic_property(property_uid):
  683. domain = str(request.args.get("business_domain_uid") or "").strip()
  684. identity = _identity()
  685. scoped = set(identity.get("business_domains") or [domain])
  686. if "*" not in scoped and domain not in scoped:
  687. return jsonify(failed("权限不足", code=403)), 403
  688. try:
  689. result = get_semantic_query_service().trace_property(
  690. property_uid,
  691. allowed_domains={domain},
  692. limit=int(request.args.get("limit") or 20),
  693. after_uid=request.args.get("after_uid"),
  694. )
  695. return jsonify(success(result)), 200
  696. except Exception as error:
  697. return _error(error)