routes.py 32 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855856857858859860861862863864865866867868869870871872873874875876877878879880881882883884885886887888889890891892893894895896897898899900901902903904905906907908909910911912913914915916917918919920921922923924925926927928929930931932933934935936937938939940941942943944945946947948949950951952953954955956957958959960961962963964965966967968969970971972973974975976977978979980981982983984985986987988989990991992993994995996997998999100010011002100310041005100610071008
  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_device_asset_service():
  80. from app.core.data_research.device_asset_repository import (
  81. SqlAlchemyDeviceAssetRepository,
  82. )
  83. from app.core.data_research.device_assets import DeviceAssetService
  84. return DeviceAssetService(
  85. SqlAlchemyDeviceAssetRepository(db.session),
  86. commit=db.session.commit,
  87. rollback=db.session.rollback,
  88. )
  89. def get_candidate_decision_service():
  90. from app.core.data_research.candidate_decisions import CandidateDecisionService
  91. from app.core.data_research.data_elements import DataElementService
  92. from app.core.data_research.repository import (
  93. SqlAlchemyCandidateDecisionRepository,
  94. SqlAlchemyDataElementRepository,
  95. )
  96. return CandidateDecisionService(
  97. SqlAlchemyCandidateDecisionRepository(db.session),
  98. data_elements=DataElementService(SqlAlchemyDataElementRepository(db.session)),
  99. )
  100. def _artifact_storage():
  101. from minio import Minio
  102. from app.core.data_research.artifacts import MinioArtifactStorage
  103. client = Minio(
  104. current_app.config["MINIO_HOST"],
  105. access_key=current_app.config["MINIO_USER"],
  106. secret_key=current_app.config["MINIO_PASSWORD"],
  107. secure=bool(current_app.config.get("MINIO_SECURE")),
  108. )
  109. return MinioArtifactStorage(client, current_app.config["MINIO_BUCKET"])
  110. def get_artifact_service():
  111. from app.core.common.identifiers import new_governance_uid
  112. from app.core.data_research.artifacts import (
  113. ArtifactService,
  114. SqlAlchemyArtifactRepository,
  115. )
  116. from app.core.data_research.file_policy import FilePolicy
  117. return ArtifactService(
  118. SqlAlchemyArtifactRepository(db.session),
  119. _artifact_storage(),
  120. FilePolicy(
  121. max_bytes=int(current_app.config.get("DATA_RESEARCH_FILE_MAX_BYTES", 25 * 1024 * 1024)),
  122. max_pages=int(current_app.config.get("DATA_RESEARCH_FILE_MAX_PAGES", 500)),
  123. ),
  124. uid_factory=new_governance_uid,
  125. )
  126. def get_evidence_service():
  127. from app.core.data_research.artifacts import EvidenceService
  128. return EvidenceService(db.session, _artifact_storage())
  129. def get_ontology_service():
  130. from app.core.data_research.ontology.publication import (
  131. OntologyApplicationService,
  132. OntologyPublicationService,
  133. )
  134. from app.core.data_research.ontology.repository import SqlAlchemyOntologyRepository
  135. from app.core.events.outbox import enqueue_outbox
  136. repository = SqlAlchemyOntologyRepository(db.session)
  137. publication = OntologyPublicationService(
  138. repository,
  139. outbox_enqueue=lambda **event: enqueue_outbox(db.session, **event),
  140. commit=db.session.commit,
  141. rollback=db.session.rollback,
  142. )
  143. return OntologyApplicationService(
  144. repository,
  145. publication=publication,
  146. commit=db.session.commit,
  147. rollback=db.session.rollback,
  148. )
  149. def get_ontology_dynamic_service():
  150. from app.core.data_research.ontology.change_sets import (
  151. SqlAlchemyDynamicOntologyService,
  152. )
  153. return SqlAlchemyDynamicOntologyService(db.session)
  154. def get_ontology_exchange_service():
  155. from app.core.data_research.ontology.exchange import OntologyExchangeService
  156. from app.core.data_research.ontology.repository import SqlAlchemyOntologyRepository
  157. return OntologyExchangeService(
  158. SqlAlchemyOntologyRepository(db.session),
  159. commit=db.session.commit,
  160. rollback=db.session.rollback,
  161. )
  162. def get_semantic_query_service():
  163. from app.core.data_research.ontology.query import (
  164. Neo4jSemanticRepository,
  165. SemanticQueryService,
  166. )
  167. from app.services.neo4j_driver import neo4j_driver
  168. return SemanticQueryService(Neo4jSemanticRepository(neo4j_driver))
  169. def _identity():
  170. return getattr(g, "current_user", {}) or {}
  171. def _record(record):
  172. return {
  173. "uid": str(record.uid),
  174. "source_uid": str(record.source_uid),
  175. "artifact_uid": str(record.artifact_uid) if record.artifact_uid else None,
  176. "job_type": record.job_type,
  177. "parser_version": record.parser_version,
  178. "status": record.status,
  179. "parameters": dict(record.parameters or {}),
  180. "statistics": dict(record.statistics or {}),
  181. "last_error": record.last_error,
  182. "attempt_count": int(record.attempt_count or 0),
  183. "failure_stage": record.failure_stage,
  184. "actor_uid": record.actor_uid,
  185. "created_at": record.created_at.isoformat() if record.created_at else None,
  186. "updated_at": record.updated_at.isoformat() if record.updated_at else None,
  187. "started_at": record.started_at.isoformat() if record.started_at else None,
  188. "finished_at": record.finished_at.isoformat() if record.finished_at else None,
  189. }
  190. def _catalog_snapshot(record):
  191. return {
  192. "uid": str(record.uid),
  193. "job_uid": str(record.job_uid),
  194. "source_uid": str(record.source_uid),
  195. "attempt": int(record.attempt),
  196. "database_type": record.database_type,
  197. "content_hash": record.content_hash,
  198. "snapshot": dict(record.snapshot or {}),
  199. "evidence_count": int(record.evidence_count),
  200. "created_at": (
  201. record.created_at.isoformat()
  202. if record.created_at
  203. else None
  204. ),
  205. }
  206. def _element(record):
  207. return {
  208. "uid": str(record.uid),
  209. "code": record.code,
  210. "status": record.status,
  211. "current_version": int(record.current_version),
  212. "snapshot": dict(record.snapshot or {}),
  213. "created_by": record.created_by,
  214. "updated_by": record.updated_by,
  215. }
  216. def _decision(record):
  217. return {
  218. "uid": str(record.uid),
  219. "candidate_uid": str(record.candidate_uid),
  220. "action": record.action,
  221. "data_element_uid": (
  222. str(record.data_element_uid) if record.data_element_uid else None
  223. ),
  224. "evidence_uids": list(record.evidence_uids),
  225. "actor_uid": record.actor_uid,
  226. "reason": record.reason,
  227. }
  228. def _artifact(record):
  229. return {
  230. "uid": str(record.uid),
  231. "source_uid": str(record.source_uid),
  232. "filename": record.filename,
  233. "media_type": record.media_type,
  234. "size_bytes": int(record.size_bytes),
  235. "content_hash": record.content_hash,
  236. "parser_version": record.parser_version,
  237. }
  238. def _ontology(record):
  239. return {
  240. "uid": str(record.uid),
  241. "code": record.code,
  242. "name": record.name,
  243. "owner_uid": record.owner_uid,
  244. "status": record.status,
  245. "draft_revision": int(record.draft_revision),
  246. "active_version_uid": record.active_version_uid,
  247. "domain_links": [
  248. {"domain_uid": link.domain_uid, "role": link.role}
  249. for link in record.domain_links
  250. ],
  251. }
  252. def _ontology_version(record):
  253. return {
  254. "uid": str(record.uid),
  255. "ontology_uid": str(record.ontology_uid),
  256. "version": int(record.version),
  257. "parent_version_uid": record.parent_version_uid,
  258. "status": record.status,
  259. "content_hash": record.content_hash,
  260. "graph_document": record.graph_document.to_dict(),
  261. "created_by": record.created_by,
  262. }
  263. def _iso(value):
  264. return value.isoformat() if value else None
  265. def _device_asset(record):
  266. return {
  267. "uid": str(record.uid),
  268. "asset_type": record.asset_type,
  269. "name": record.name,
  270. "status": record.status,
  271. "current_version": int(record.current_version),
  272. "location": record.location,
  273. "organization": record.organization,
  274. "responsible_person": record.responsible_person,
  275. "attributes": dict(record.attributes or {}),
  276. "created_by": record.created_by,
  277. "updated_by": record.updated_by,
  278. "created_at": _iso(record.created_at),
  279. "updated_at": _iso(record.updated_at),
  280. }
  281. def _device_asset_mapping(record):
  282. return {
  283. "uid": str(record.uid),
  284. "asset_uid": str(record.asset_uid),
  285. "source_uid": str(record.source_uid),
  286. "source_entity": record.source_entity,
  287. "asset_type": record.asset_type,
  288. "source_code": record.source_code,
  289. "source_updated_at": _iso(record.source_updated_at),
  290. "first_seen_at": _iso(record.first_seen_at),
  291. "last_seen_at": _iso(record.last_seen_at),
  292. }
  293. def _device_asset_detail(detail):
  294. data = _device_asset(detail.asset)
  295. data["source_mappings"] = [
  296. _device_asset_mapping(mapping)
  297. for mapping in detail.mappings
  298. ]
  299. return data
  300. def _device_asset_version(record):
  301. return {
  302. "uid": str(record.uid),
  303. "asset_uid": str(record.asset_uid),
  304. "version": int(record.version),
  305. "snapshot": dict(record.snapshot or {}),
  306. "source_mapping_uid": str(record.source_mapping_uid),
  307. "actor_uid": record.actor_uid,
  308. "created_at": _iso(record.created_at),
  309. }
  310. def _device_asset_import_result(result):
  311. return {
  312. "records": [
  313. {
  314. "action": item.action,
  315. "asset": _device_asset(item.asset),
  316. "source_mapping": _device_asset_mapping(item.mapping),
  317. }
  318. for item in result.items
  319. ],
  320. "created_count": int(result.created_count),
  321. "updated_count": int(result.updated_count),
  322. "unchanged_count": int(result.unchanged_count),
  323. }
  324. def _device_asset_page(name, *, default, maximum):
  325. from app.core.data_research.errors import DeviceAssetInvalid
  326. raw = request.args.get(name)
  327. try:
  328. value = default if raw in (None, "") else int(raw)
  329. except (TypeError, ValueError) as error:
  330. raise DeviceAssetInvalid(f"{name} must be an integer") from error
  331. if value < 1 or value > maximum:
  332. raise DeviceAssetInvalid(
  333. f"{name} must be between 1 and {maximum}"
  334. )
  335. return value
  336. def _error(error):
  337. if isinstance(error, DataResearchError):
  338. return (
  339. jsonify(
  340. failed(
  341. str(error),
  342. code=error.http_status,
  343. error={"code": error.code},
  344. )
  345. ),
  346. error.http_status,
  347. )
  348. logger.exception("data-research ingestion request failed")
  349. return (
  350. jsonify(
  351. failed(
  352. "数据采集任务处理失败",
  353. code=500,
  354. error={"code": "DATA_RESEARCH_ERROR"},
  355. )
  356. ),
  357. 500,
  358. )
  359. @bp.route("/ingestion-jobs", methods=["POST"])
  360. def create_ingestion_job():
  361. payload = request.get_json(silent=True) or {}
  362. try:
  363. actor_uid = _identity().get("id") or _identity().get("sub")
  364. if payload.get("job_type") == "catalog_collect":
  365. get_database_source_registration_service().ensure(
  366. payload.get("source_uid"),
  367. actor_uid=actor_uid,
  368. )
  369. record, created = get_ingestion_service().create_job(
  370. payload,
  371. actor_uid=actor_uid,
  372. )
  373. return jsonify(success(_record(record))), 201 if created else 200
  374. except Exception as error:
  375. return _error(error)
  376. @bp.route("/ingestion-jobs", methods=["GET"])
  377. def list_ingestion_jobs():
  378. filters = {
  379. name: request.args.get(name)
  380. for name in ("status", "source_uid")
  381. if request.args.get(name)
  382. }
  383. try:
  384. records = get_ingestion_service().list_jobs(filters)
  385. return jsonify(
  386. success({"records": [_record(item) for item in records], "total": len(records)})
  387. ), 200
  388. except Exception as error:
  389. return _error(error)
  390. @bp.route("/ingestion-jobs/<job_uid>", methods=["GET"])
  391. def get_ingestion_job(job_uid):
  392. try:
  393. return jsonify(success(_record(get_ingestion_service().get_job(job_uid)))), 200
  394. except Exception as error:
  395. return _error(error)
  396. @bp.route("/ingestion-jobs/<job_uid>/execute", methods=["POST"])
  397. def execute_ingestion_job(job_uid):
  398. try:
  399. record = get_catalog_ingestion_executor().execute(job_uid)
  400. return jsonify(success(_record(record))), 200
  401. except Exception as error:
  402. return _error(error)
  403. @bp.route(
  404. "/ingestion-jobs/<job_uid>/catalog-snapshots",
  405. methods=["GET"],
  406. )
  407. def list_catalog_snapshots(job_uid):
  408. try:
  409. records = get_catalog_snapshot_repository().list(job_uid)
  410. return jsonify(
  411. success(
  412. {
  413. "records": [
  414. _catalog_snapshot(record)
  415. for record in records
  416. ],
  417. "total": len(records),
  418. }
  419. )
  420. ), 200
  421. except Exception as error:
  422. return _error(error)
  423. @bp.route("/ingestion-jobs/<job_uid>/evidence", methods=["GET"])
  424. def list_ingestion_job_evidence(job_uid):
  425. try:
  426. records = get_evidence_service().list_for_job(job_uid)
  427. return jsonify(
  428. success({"records": records, "total": len(records)})
  429. ), 200
  430. except Exception as error:
  431. return _error(error)
  432. @bp.route("/ingestion-jobs/<job_uid>/retry", methods=["POST"])
  433. def retry_ingestion_job(job_uid):
  434. try:
  435. return jsonify(success(_record(get_ingestion_service().retry(job_uid)))), 200
  436. except Exception as error:
  437. return _error(error)
  438. @bp.route("/ingestion-jobs/<job_uid>/cancel", methods=["POST"])
  439. def cancel_ingestion_job(job_uid):
  440. try:
  441. service = get_ingestion_service()
  442. record = service.get_job(job_uid)
  443. identity = _identity()
  444. permissions = set(identity.get("permissions") or [])
  445. actor_uid = identity.get("id") or identity.get("sub")
  446. if record.actor_uid != actor_uid and "ingestion:admin" not in permissions:
  447. return jsonify(failed("权限不足", code=403)), 403
  448. return jsonify(success(_record(service.cancel(job_uid)))), 200
  449. except Exception as error:
  450. return _error(error)
  451. @bp.route("/data-elements", methods=["GET"])
  452. def list_data_elements():
  453. try:
  454. records = get_data_element_service().list(
  455. status=request.args.get("status"),
  456. business_domain_uid=request.args.get("business_domain_uid"),
  457. )
  458. return jsonify(
  459. success(
  460. {
  461. "records": [_element(item) for item in records],
  462. "total": len(records),
  463. }
  464. )
  465. ), 200
  466. except Exception as error:
  467. return _error(error)
  468. @bp.route("/device-assets", methods=["GET"])
  469. def list_device_assets():
  470. filters = {
  471. name: request.args.get(name)
  472. for name in ("keyword", "asset_type", "status", "source_uid")
  473. if request.args.get(name)
  474. }
  475. try:
  476. page = _device_asset_page(
  477. "page",
  478. default=1,
  479. maximum=1_000_000,
  480. )
  481. page_size = _device_asset_page(
  482. "page_size",
  483. default=20,
  484. maximum=100,
  485. )
  486. records, total = get_device_asset_service().search(
  487. filters,
  488. page=page,
  489. page_size=page_size,
  490. )
  491. return jsonify(
  492. success(
  493. {
  494. "records": [
  495. _device_asset_detail(record)
  496. for record in records
  497. ],
  498. "total": int(total),
  499. "page": page,
  500. "page_size": page_size,
  501. }
  502. )
  503. ), 200
  504. except Exception as error:
  505. return _error(error)
  506. @bp.route("/device-assets/import", methods=["POST"])
  507. def import_device_assets():
  508. try:
  509. result = get_device_asset_service().import_records(
  510. request.get_json(silent=True) or {},
  511. actor_uid=_identity().get("id") or _identity().get("sub"),
  512. )
  513. return jsonify(success(_device_asset_import_result(result))), 200
  514. except Exception as error:
  515. return _error(error)
  516. @bp.route("/device-assets/<asset_uid>", methods=["GET"])
  517. def get_device_asset(asset_uid):
  518. try:
  519. return jsonify(
  520. success(
  521. _device_asset_detail(
  522. get_device_asset_service().get(asset_uid)
  523. )
  524. )
  525. ), 200
  526. except Exception as error:
  527. return _error(error)
  528. @bp.route("/device-assets/<asset_uid>/versions", methods=["GET"])
  529. def list_device_asset_versions(asset_uid):
  530. try:
  531. records = get_device_asset_service().versions(asset_uid)
  532. return jsonify(
  533. success(
  534. {
  535. "records": [
  536. _device_asset_version(record)
  537. for record in records
  538. ],
  539. "total": len(records),
  540. }
  541. )
  542. ), 200
  543. except Exception as error:
  544. return _error(error)
  545. @bp.route("/data-elements", methods=["POST"])
  546. def create_data_element():
  547. try:
  548. record = get_data_element_service().create_draft(
  549. request.get_json(silent=True) or {},
  550. actor_uid=_identity().get("id") or _identity().get("sub"),
  551. )
  552. db.session.commit()
  553. return jsonify(success(_element(record))), 201
  554. except Exception as error:
  555. db.session.rollback()
  556. return _error(error)
  557. @bp.route("/data-elements/<element_uid>/transitions", methods=["POST"])
  558. def transition_data_element(element_uid):
  559. payload = request.get_json(silent=True) or {}
  560. identity = _identity()
  561. target_status = str(payload.get("target_status") or "")
  562. if (
  563. target_status in {"published", "deprecated", "retired"}
  564. and "data-elements:publish" not in set(identity.get("permissions") or [])
  565. ):
  566. return jsonify(failed("权限不足", code=403)), 403
  567. try:
  568. record = get_data_element_service().transition(
  569. element_uid,
  570. target_status,
  571. expected_version=int(payload.get("expected_version")),
  572. actor_uid=identity.get("id") or identity.get("sub"),
  573. )
  574. db.session.commit()
  575. return jsonify(success(_element(record))), 200
  576. except Exception as error:
  577. db.session.rollback()
  578. return _error(error)
  579. @bp.route("/candidate-decisions", methods=["POST"])
  580. def decide_candidates():
  581. payload = request.get_json(silent=True) or {}
  582. try:
  583. records = get_candidate_decision_service().decide(
  584. payload.get("decisions") or [],
  585. actor_uid=_identity().get("id") or _identity().get("sub"),
  586. )
  587. return jsonify(success([_decision(record) for record in records])), 200
  588. except Exception as error:
  589. return _error(error)
  590. @bp.route("/sources/files", methods=["POST"])
  591. def upload_source_file():
  592. uploaded = request.files.get("file")
  593. if uploaded is None:
  594. return jsonify(failed("缺少上传文件", code=400)), 400
  595. try:
  596. record, created = get_artifact_service().store(
  597. request.form.get("source_uid"),
  598. uploaded.filename,
  599. uploaded.mimetype,
  600. uploaded.read(),
  601. request.form.get("parser_version") or "auto-v1",
  602. )
  603. db.session.commit()
  604. return jsonify(success(_artifact(record))), 201 if created else 200
  605. except Exception as error:
  606. db.session.rollback()
  607. return _error(error)
  608. @bp.route("/evidence/<evidence_uid>", methods=["GET"])
  609. def get_evidence(evidence_uid):
  610. identity = _identity()
  611. try:
  612. service = get_evidence_service()
  613. if request.args.get("download") in {"1", "true", "yes"}:
  614. if "evidence:download" not in set(identity.get("permissions") or []):
  615. return jsonify(failed("权限不足", code=403)), 403
  616. content, filename, media_type = service.download(evidence_uid)
  617. return send_file(
  618. io.BytesIO(content),
  619. mimetype=media_type,
  620. as_attachment=True,
  621. download_name=filename,
  622. )
  623. preview = service.preview(evidence_uid)
  624. from app.core.data_research.artifacts import redact_excerpt
  625. preview["excerpt"] = redact_excerpt(preview.get("excerpt"))
  626. return jsonify(success(preview)), 200
  627. except Exception as error:
  628. return _error(error)
  629. @bp.route("/ontologies", methods=["GET"])
  630. def list_ontologies():
  631. try:
  632. return jsonify(success([_ontology(item) for item in get_ontology_service().list()])), 200
  633. except Exception as error:
  634. return _error(error)
  635. @bp.route("/ontologies", methods=["POST"])
  636. def create_ontology():
  637. try:
  638. record = get_ontology_service().create(
  639. request.get_json(silent=True) or {},
  640. _identity().get("id") or _identity().get("sub"),
  641. )
  642. return jsonify(success(_ontology(record))), 201
  643. except Exception as error:
  644. return _error(error)
  645. @bp.route("/ontologies/<ontology_uid>", methods=["GET"])
  646. def get_ontology(ontology_uid):
  647. try:
  648. return jsonify(success(_ontology(get_ontology_service().get(ontology_uid)))), 200
  649. except Exception as error:
  650. return _error(error)
  651. @bp.route("/ontologies/<ontology_uid>/versions", methods=["GET"])
  652. def list_ontology_versions(ontology_uid):
  653. try:
  654. return jsonify(
  655. success(
  656. [
  657. _ontology_version(item)
  658. for item in get_ontology_service().list_versions(ontology_uid)
  659. ]
  660. )
  661. ), 200
  662. except Exception as error:
  663. return _error(error)
  664. @bp.route("/ontologies/<ontology_uid>/graph", methods=["GET"])
  665. def get_ontology_graph(ontology_uid):
  666. try:
  667. service = get_ontology_service()
  668. ontology = service.get(ontology_uid)
  669. version = service.latest_version(ontology_uid)
  670. data = (
  671. _ontology_version(version)
  672. if version is not None
  673. else {
  674. "uid": None,
  675. "ontology_uid": str(ontology_uid),
  676. "version": 0,
  677. "parent_version_uid": None,
  678. "status": "draft",
  679. "content_hash": None,
  680. "graph_document": {
  681. name: []
  682. for name in (
  683. "classes",
  684. "properties",
  685. "relations",
  686. "constraints",
  687. "domain_links",
  688. "element_mappings",
  689. )
  690. },
  691. "created_by": None,
  692. }
  693. )
  694. response = jsonify(success(data))
  695. response.headers["ETag"] = f'"{int(ontology.draft_revision)}"'
  696. return response, 200
  697. except Exception as error:
  698. return _error(error)
  699. @bp.route("/ontologies/<ontology_uid>/graph", methods=["PATCH"])
  700. def save_ontology_graph(ontology_uid):
  701. etag = str(request.headers.get("If-Match") or "").strip().strip('"')
  702. if not etag.isdigit():
  703. return jsonify(failed("缺少有效 If-Match 修订号", code=428)), 428
  704. try:
  705. version = get_ontology_service().save_draft(
  706. ontology_uid,
  707. request.get_json(silent=True) or {},
  708. int(etag),
  709. _identity().get("id") or _identity().get("sub"),
  710. )
  711. response = jsonify(success(_ontology_version(version)))
  712. response.headers["ETag"] = f'"{int(etag) + 1}"'
  713. return response, 200
  714. except Exception as error:
  715. return _error(error)
  716. @bp.route("/ontologies/<ontology_uid>/validate", methods=["POST"])
  717. def validate_ontology(ontology_uid):
  718. try:
  719. issues = get_ontology_service().validate(ontology_uid)
  720. data = [
  721. {"code": item.code, "message": item.message, "path": item.path}
  722. for item in issues
  723. ]
  724. return jsonify(success(data)), 200
  725. except Exception as error:
  726. return _error(error)
  727. @bp.route("/ontologies/<ontology_uid>/publish", methods=["POST"])
  728. def publish_ontology(ontology_uid):
  729. try:
  730. version = get_ontology_service().publish(
  731. ontology_uid,
  732. request.headers.get("Idempotency-Key") or f"{ontology_uid}:publish",
  733. _identity().get("id") or _identity().get("sub"),
  734. )
  735. return jsonify(success(_ontology_version(version))), 200
  736. except Exception as error:
  737. return _error(error)
  738. @bp.route("/ontologies/<ontology_uid>/diff", methods=["GET"])
  739. def diff_ontology(ontology_uid):
  740. try:
  741. return jsonify(success(get_ontology_service().diff(
  742. ontology_uid, request.args.get("left"), request.args.get("right")
  743. ))), 200
  744. except Exception as error:
  745. return _error(error)
  746. @bp.route("/ontologies/<ontology_uid>/rollback", methods=["POST"])
  747. def rollback_ontology(ontology_uid):
  748. payload = request.get_json(silent=True) or {}
  749. try:
  750. version = get_ontology_service().rollback(
  751. ontology_uid,
  752. payload.get("target_version_uid"),
  753. int(payload.get("expected_revision")),
  754. _identity().get("id") or _identity().get("sub"),
  755. )
  756. return jsonify(success(_ontology_version(version))), 201
  757. except Exception as error:
  758. return _error(error)
  759. @bp.route("/ontologies/<ontology_uid>/suggestions", methods=["POST"])
  760. def generate_ontology_suggestions(ontology_uid):
  761. try:
  762. result = get_ontology_dynamic_service().generate(
  763. ontology_uid,
  764. request.get_json(silent=True) or {},
  765. _identity().get("id") or _identity().get("sub"),
  766. )
  767. records = (
  768. result.get("suggestions") or []
  769. if isinstance(result, dict)
  770. else result
  771. )
  772. suggestions = [
  773. item
  774. if isinstance(item, dict)
  775. else {
  776. "uid": item.uid,
  777. "kind": item.kind,
  778. "payload": item.payload,
  779. "evidence_uids": list(item.evidence_uids),
  780. "confidence": item.confidence,
  781. "source": item.source,
  782. "model_version": item.model_version,
  783. "prompt_version": item.prompt_version,
  784. }
  785. for item in records
  786. ]
  787. data = {
  788. "change_set_uid": (
  789. result.get("change_set_uid")
  790. if isinstance(result, dict)
  791. else None
  792. ),
  793. "suggestions": suggestions,
  794. }
  795. return jsonify(success(data)), 201
  796. except Exception as error:
  797. return _error(error)
  798. @bp.route(
  799. "/ontologies/<ontology_uid>/change-sets/<change_set_uid>/decisions",
  800. methods=["POST"],
  801. )
  802. def decide_ontology_change_set(ontology_uid, change_set_uid):
  803. payload = request.get_json(silent=True) or {}
  804. try:
  805. result = get_ontology_dynamic_service().decide(
  806. ontology_uid,
  807. change_set_uid,
  808. payload.get("decisions") or [],
  809. _identity().get("id") or _identity().get("sub"),
  810. )
  811. return jsonify(success(result)), 200
  812. except Exception as error:
  813. return _error(error)
  814. @bp.route("/ontologies/<ontology_uid>/export", methods=["GET"])
  815. def export_ontology(ontology_uid):
  816. try:
  817. content, media_type, filename = get_ontology_exchange_service().export(
  818. ontology_uid, request.args.get("format") or "json"
  819. )
  820. return send_file(
  821. io.BytesIO(content),
  822. mimetype=media_type,
  823. as_attachment=True,
  824. download_name=filename,
  825. )
  826. except Exception as error:
  827. return _error(error)
  828. @bp.route("/ontologies/import", methods=["POST"])
  829. def import_ontology():
  830. try:
  831. uploaded = request.files.get("file")
  832. content = uploaded.read() if uploaded is not None else request.get_data(cache=False)
  833. result = get_ontology_exchange_service().import_document(
  834. content,
  835. request.args.get("format") or "json",
  836. _identity().get("id") or _identity().get("sub"),
  837. )
  838. return jsonify(success(result)), 201
  839. except Exception as error:
  840. return _error(error)
  841. @bp.route("/semantic/properties/<property_uid>", methods=["GET"])
  842. def query_semantic_property(property_uid):
  843. domain = str(request.args.get("business_domain_uid") or "").strip()
  844. identity = _identity()
  845. scoped = set(identity.get("business_domains") or [domain])
  846. if "*" not in scoped and domain not in scoped:
  847. return jsonify(failed("权限不足", code=403)), 403
  848. try:
  849. result = get_semantic_query_service().trace_property(
  850. property_uid,
  851. allowed_domains={domain},
  852. limit=int(request.args.get("limit") or 20),
  853. after_uid=request.args.get("after_uid"),
  854. )
  855. return jsonify(success(result)), 200
  856. except Exception as error:
  857. return _error(error)