routes.py 51 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855856857858859860861862863864865866867868869870871872873874875876877878879880881882883884885886887888889890891892893894895896897898899900901902903904905906907908909910911912913914915916917918919920921922923924925926927928929930931932933934935936937938939940941942943944945946947948949950951952953954955956957958959960961962963964965966967968969970971972973974975976977978979980981982983984985986987988989990991992993994995996997998999100010011002100310041005100610071008100910101011101210131014101510161017101810191020102110221023102410251026102710281029103010311032103310341035103610371038103910401041104210431044104510461047104810491050105110521053105410551056105710581059106010611062106310641065106610671068106910701071107210731074107510761077107810791080108110821083108410851086108710881089109010911092109310941095109610971098109911001101110211031104110511061107110811091110111111121113111411151116111711181119112011211122112311241125112611271128112911301131113211331134113511361137113811391140114111421143114411451146114711481149115011511152115311541155115611571158115911601161116211631164116511661167116811691170117111721173117411751176117711781179118011811182118311841185118611871188118911901191119211931194119511961197119811991200120112021203120412051206120712081209121012111212121312141215121612171218121912201221122212231224122512261227122812291230123112321233123412351236123712381239124012411242124312441245124612471248124912501251125212531254125512561257125812591260126112621263126412651266126712681269127012711272127312741275127612771278127912801281128212831284128512861287128812891290129112921293129412951296129712981299130013011302130313041305130613071308130913101311131213131314131513161317131813191320132113221323132413251326132713281329133013311332133313341335133613371338133913401341134213431344134513461347134813491350135113521353135413551356135713581359136013611362136313641365136613671368136913701371137213731374137513761377137813791380138113821383138413851386138713881389139013911392139313941395139613971398139914001401140214031404140514061407140814091410141114121413141414151416141714181419142014211422142314241425142614271428142914301431143214331434143514361437143814391440144114421443144414451446144714481449145014511452145314541455145614571458145914601461146214631464146514661467146814691470147114721473147414751476147714781479148014811482148314841485148614871488148914901491149214931494149514961497149814991500150115021503150415051506150715081509151015111512151315141515151615171518151915201521152215231524152515261527152815291530153115321533153415351536153715381539154015411542154315441545154615471548154915501551155215531554155515561557155815591560156115621563156415651566156715681569157015711572157315741575157615771578157915801581158215831584158515861587158815891590159115921593159415951596159715981599160016011602160316041605160616071608160916101611161216131614161516161617161816191620162116221623162416251626162716281629163016311632163316341635163616371638163916401641164216431644164516461647
  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 _device_ontology_authorizer(ontology_repository):
  90. from app.core.data_research.device_semantics import (
  91. DeviceOntologyPublicationAuthorizer,
  92. )
  93. from app.core.governance.responsibilities import (
  94. ResponsibilityService,
  95. SqlAlchemyResponsibilityRepository,
  96. )
  97. responsibilities = ResponsibilityService(
  98. SqlAlchemyResponsibilityRepository(db.session)
  99. )
  100. return DeviceOntologyPublicationAuthorizer(
  101. ontology_repository,
  102. responsibility_lookup=responsibilities.get,
  103. )
  104. def get_device_semantic_service():
  105. from app.core.data_research.device_semantics import DeviceSemanticService
  106. from app.core.data_research.ontology.repository import (
  107. SqlAlchemyOntologyRepository,
  108. )
  109. return DeviceSemanticService(
  110. SqlAlchemyOntologyRepository(db.session),
  111. commit=db.session.commit,
  112. rollback=db.session.rollback,
  113. )
  114. def get_device_semantic_code_service():
  115. from app.core.data_research.device_semantic_repository import (
  116. SqlAlchemyDeviceSemanticCodeRepository,
  117. )
  118. from app.core.data_research.device_semantics import (
  119. DeviceSemanticCodeService,
  120. )
  121. from app.core.data_research.ontology.repository import (
  122. SqlAlchemyOntologyRepository,
  123. )
  124. ontology_repository = SqlAlchemyOntologyRepository(db.session)
  125. authorizer = _device_ontology_authorizer(ontology_repository)
  126. return DeviceSemanticCodeService(
  127. SqlAlchemyDeviceSemanticCodeRepository(db.session),
  128. review_authorizer=authorizer.assert_accountable,
  129. commit=db.session.commit,
  130. rollback=db.session.rollback,
  131. )
  132. def get_device_entity_resolution_service():
  133. from app.core.data_research.device_entity_repository import (
  134. SqlAlchemyDeviceEntityResolutionRepository,
  135. )
  136. from app.core.data_research.device_entity_resolution import (
  137. DeviceEntityForbidden,
  138. DeviceEntityResolutionService,
  139. )
  140. from app.core.governance.responsibilities import (
  141. ResponsibilityService,
  142. SqlAlchemyResponsibilityRepository,
  143. )
  144. responsibilities = ResponsibilityService(
  145. SqlAlchemyResponsibilityRepository(db.session)
  146. )
  147. def assert_accountable(actor_uid):
  148. matrix = responsibilities.get(
  149. "device_mapping",
  150. "DEVICE_ENTITY_RESOLUTION",
  151. )
  152. accountable = [
  153. item
  154. for item in matrix.get("assignments", [])
  155. if item.get("responsibility_role") == "asset_manager"
  156. and item.get("raci_role") == "accountable"
  157. ]
  158. if len(accountable) != 1 or str(accountable[0].get("user_id")) != str(
  159. actor_uid
  160. ):
  161. raise DeviceEntityForbidden(
  162. "only the accountable device mapping asset manager may decide"
  163. )
  164. return DeviceEntityResolutionService(
  165. SqlAlchemyDeviceEntityResolutionRepository(db.session),
  166. review_authorizer=assert_accountable,
  167. auto_merge_enabled=current_app.config.get(
  168. "DEVICE_ENTITY_AUTO_MERGE_ENABLED",
  169. False,
  170. ),
  171. commit=db.session.commit,
  172. rollback=db.session.rollback,
  173. )
  174. def get_candidate_decision_service():
  175. from app.core.data_research.candidate_decisions import CandidateDecisionService
  176. from app.core.data_research.data_elements import DataElementService
  177. from app.core.data_research.repository import (
  178. SqlAlchemyCandidateDecisionRepository,
  179. SqlAlchemyDataElementRepository,
  180. )
  181. return CandidateDecisionService(
  182. SqlAlchemyCandidateDecisionRepository(db.session),
  183. data_elements=DataElementService(SqlAlchemyDataElementRepository(db.session)),
  184. )
  185. def _artifact_storage():
  186. from minio import Minio
  187. from app.core.data_research.artifacts import MinioArtifactStorage
  188. client = Minio(
  189. current_app.config["MINIO_HOST"],
  190. access_key=current_app.config["MINIO_USER"],
  191. secret_key=current_app.config["MINIO_PASSWORD"],
  192. secure=bool(current_app.config.get("MINIO_SECURE")),
  193. )
  194. return MinioArtifactStorage(client, current_app.config["MINIO_BUCKET"])
  195. def get_artifact_service():
  196. from app.core.common.identifiers import new_governance_uid
  197. from app.core.data_research.artifacts import (
  198. ArtifactService,
  199. SqlAlchemyArtifactRepository,
  200. )
  201. from app.core.data_research.file_policy import FilePolicy
  202. return ArtifactService(
  203. SqlAlchemyArtifactRepository(db.session),
  204. _artifact_storage(),
  205. FilePolicy(
  206. max_bytes=int(current_app.config.get("DATA_RESEARCH_FILE_MAX_BYTES", 25 * 1024 * 1024)),
  207. max_pages=int(current_app.config.get("DATA_RESEARCH_FILE_MAX_PAGES", 500)),
  208. ),
  209. uid_factory=new_governance_uid,
  210. )
  211. def get_evidence_service():
  212. from app.core.data_research.artifacts import EvidenceService
  213. return EvidenceService(db.session, _artifact_storage())
  214. def get_ontology_service():
  215. from app.core.data_research.ontology.publication import (
  216. OntologyApplicationService,
  217. OntologyPublicationService,
  218. )
  219. from app.core.data_research.ontology.repository import SqlAlchemyOntologyRepository
  220. from app.core.events.outbox import enqueue_outbox
  221. repository = SqlAlchemyOntologyRepository(db.session)
  222. authorizer = _device_ontology_authorizer(repository)
  223. publication = OntologyPublicationService(
  224. repository,
  225. outbox_enqueue=lambda **event: enqueue_outbox(db.session, **event),
  226. publication_authorizer=authorizer,
  227. commit=db.session.commit,
  228. rollback=db.session.rollback,
  229. )
  230. return OntologyApplicationService(
  231. repository,
  232. publication=publication,
  233. commit=db.session.commit,
  234. rollback=db.session.rollback,
  235. )
  236. def get_ontology_dynamic_service():
  237. from app.core.data_research.ontology.change_sets import (
  238. SqlAlchemyDynamicOntologyService,
  239. )
  240. return SqlAlchemyDynamicOntologyService(db.session)
  241. def get_ontology_exchange_service():
  242. from app.core.data_research.ontology.exchange import OntologyExchangeService
  243. from app.core.data_research.ontology.repository import SqlAlchemyOntologyRepository
  244. return OntologyExchangeService(
  245. SqlAlchemyOntologyRepository(db.session),
  246. commit=db.session.commit,
  247. rollback=db.session.rollback,
  248. )
  249. def get_semantic_query_service():
  250. from app.core.data_research.ontology.query import (
  251. Neo4jSemanticRepository,
  252. SemanticQueryService,
  253. )
  254. from app.services.neo4j_driver import neo4j_driver
  255. return SemanticQueryService(Neo4jSemanticRepository(neo4j_driver))
  256. def _identity():
  257. return getattr(g, "current_user", {}) or {}
  258. def _record(record):
  259. return {
  260. "uid": str(record.uid),
  261. "source_uid": str(record.source_uid),
  262. "artifact_uid": str(record.artifact_uid) if record.artifact_uid else None,
  263. "job_type": record.job_type,
  264. "parser_version": record.parser_version,
  265. "status": record.status,
  266. "parameters": dict(record.parameters or {}),
  267. "statistics": dict(record.statistics or {}),
  268. "last_error": record.last_error,
  269. "attempt_count": int(record.attempt_count or 0),
  270. "failure_stage": record.failure_stage,
  271. "actor_uid": record.actor_uid,
  272. "created_at": record.created_at.isoformat() if record.created_at else None,
  273. "updated_at": record.updated_at.isoformat() if record.updated_at else None,
  274. "started_at": record.started_at.isoformat() if record.started_at else None,
  275. "finished_at": record.finished_at.isoformat() if record.finished_at else None,
  276. }
  277. def _catalog_snapshot(record):
  278. return {
  279. "uid": str(record.uid),
  280. "job_uid": str(record.job_uid),
  281. "source_uid": str(record.source_uid),
  282. "attempt": int(record.attempt),
  283. "database_type": record.database_type,
  284. "content_hash": record.content_hash,
  285. "snapshot": dict(record.snapshot or {}),
  286. "evidence_count": int(record.evidence_count),
  287. "created_at": (
  288. record.created_at.isoformat()
  289. if record.created_at
  290. else None
  291. ),
  292. }
  293. def _element(record):
  294. return {
  295. "uid": str(record.uid),
  296. "code": record.code,
  297. "status": record.status,
  298. "current_version": int(record.current_version),
  299. "snapshot": dict(record.snapshot or {}),
  300. "created_by": record.created_by,
  301. "updated_by": record.updated_by,
  302. }
  303. def _decision(record):
  304. return {
  305. "uid": str(record.uid),
  306. "candidate_uid": str(record.candidate_uid),
  307. "action": record.action,
  308. "data_element_uid": (
  309. str(record.data_element_uid) if record.data_element_uid else None
  310. ),
  311. "evidence_uids": list(record.evidence_uids),
  312. "actor_uid": record.actor_uid,
  313. "reason": record.reason,
  314. }
  315. def _artifact(record):
  316. return {
  317. "uid": str(record.uid),
  318. "source_uid": str(record.source_uid),
  319. "filename": record.filename,
  320. "media_type": record.media_type,
  321. "size_bytes": int(record.size_bytes),
  322. "content_hash": record.content_hash,
  323. "parser_version": record.parser_version,
  324. }
  325. def _ontology(record):
  326. return {
  327. "uid": str(record.uid),
  328. "code": record.code,
  329. "name": record.name,
  330. "owner_uid": record.owner_uid,
  331. "status": record.status,
  332. "draft_revision": int(record.draft_revision),
  333. "active_version_uid": record.active_version_uid,
  334. "domain_links": [
  335. {"domain_uid": link.domain_uid, "role": link.role}
  336. for link in record.domain_links
  337. ],
  338. }
  339. def _ontology_version(record):
  340. return {
  341. "uid": str(record.uid),
  342. "ontology_uid": str(record.ontology_uid),
  343. "version": int(record.version),
  344. "parent_version_uid": record.parent_version_uid,
  345. "status": record.status,
  346. "content_hash": record.content_hash,
  347. "graph_document": record.graph_document.to_dict(),
  348. "created_by": record.created_by,
  349. }
  350. def _iso(value):
  351. return value.isoformat() if value else None
  352. def _device_asset(record):
  353. return {
  354. "uid": str(record.uid),
  355. "asset_type": record.asset_type,
  356. "name": record.name,
  357. "status": record.status,
  358. "current_version": int(record.current_version),
  359. "location": record.location,
  360. "organization": record.organization,
  361. "responsible_person": record.responsible_person,
  362. "attributes": dict(record.attributes or {}),
  363. "created_by": record.created_by,
  364. "updated_by": record.updated_by,
  365. "created_at": _iso(record.created_at),
  366. "updated_at": _iso(record.updated_at),
  367. }
  368. def _device_asset_mapping(record):
  369. return {
  370. "uid": str(record.uid),
  371. "asset_uid": str(record.asset_uid),
  372. "source_uid": str(record.source_uid),
  373. "source_entity": record.source_entity,
  374. "asset_type": record.asset_type,
  375. "source_code": record.source_code,
  376. "source_updated_at": _iso(record.source_updated_at),
  377. "first_seen_at": _iso(record.first_seen_at),
  378. "last_seen_at": _iso(record.last_seen_at),
  379. }
  380. def _device_asset_detail(detail):
  381. data = _device_asset(detail.asset)
  382. data["source_mappings"] = [
  383. _device_asset_mapping(mapping)
  384. for mapping in detail.mappings
  385. ]
  386. return data
  387. def _device_asset_version(record):
  388. return {
  389. "uid": str(record.uid),
  390. "asset_uid": str(record.asset_uid),
  391. "version": int(record.version),
  392. "snapshot": dict(record.snapshot or {}),
  393. "source_mapping_uid": str(record.source_mapping_uid),
  394. "actor_uid": record.actor_uid,
  395. "created_at": _iso(record.created_at),
  396. }
  397. def _device_asset_import_result(result):
  398. return {
  399. "records": [
  400. {
  401. "action": item.action,
  402. "asset": _device_asset(item.asset),
  403. "source_mapping": _device_asset_mapping(item.mapping),
  404. }
  405. for item in result.items
  406. ],
  407. "created_count": int(result.created_count),
  408. "updated_count": int(result.updated_count),
  409. "unchanged_count": int(result.unchanged_count),
  410. }
  411. def _device_semantic_profile(result):
  412. if result is None:
  413. return {
  414. "bootstrapped": False,
  415. "ontology": None,
  416. "version": None,
  417. "profile": {
  418. "ready_to_publish": False,
  419. "class_count": 0,
  420. "relation_count": 0,
  421. "mapping_count": 0,
  422. "missing_classes": [],
  423. "missing_relations": [],
  424. "missing_mappings": [],
  425. },
  426. }
  427. profile = result.profile
  428. return {
  429. "bootstrapped": True,
  430. "created": bool(result.created),
  431. "ontology": _ontology(result.ontology),
  432. "version": _ontology_version(result.version),
  433. "profile": {
  434. "ready_to_publish": bool(profile.ready_to_publish),
  435. "class_count": int(profile.class_count),
  436. "relation_count": int(profile.relation_count),
  437. "mapping_count": int(profile.mapping_count),
  438. "missing_classes": list(profile.missing_classes),
  439. "missing_relations": list(profile.missing_relations),
  440. "missing_mappings": list(profile.missing_mappings),
  441. },
  442. }
  443. def _device_semantic_code(record):
  444. return {
  445. "uid": str(record.uid),
  446. "ontology_uid": str(record.ontology_uid),
  447. "code_type": record.code_type,
  448. "canonical_code": record.canonical_code,
  449. "canonical_name": record.canonical_name,
  450. "definition": record.definition,
  451. "status": record.status,
  452. "current_version": int(record.current_version),
  453. "source_mappings": [
  454. dict(item) for item in record.source_mappings
  455. ],
  456. "evidence_uids": list(record.evidence_uids),
  457. "suggestion_source": record.suggestion_source,
  458. "confidence": record.confidence,
  459. "created_by": record.created_by,
  460. "updated_by": record.updated_by,
  461. "created_at": _iso(record.created_at),
  462. "updated_at": _iso(record.updated_at),
  463. }
  464. def _device_semantic_version(record):
  465. return {
  466. "uid": str(record.uid),
  467. "code_uid": str(record.code_uid),
  468. "version": int(record.version),
  469. "snapshot": dict(record.snapshot),
  470. "created_by": record.created_by,
  471. "created_at": _iso(record.created_at),
  472. }
  473. def _device_semantic_review(record):
  474. return {
  475. "uid": str(record.uid),
  476. "code_uid": str(record.code_uid),
  477. "version": int(record.version),
  478. "decision": record.decision,
  479. "reason": record.reason,
  480. "actor_uid": record.actor_uid,
  481. "created_at": _iso(record.created_at),
  482. }
  483. def _device_entity_candidate(record):
  484. return {
  485. "uid": str(record.uid),
  486. "left_asset_uid": str(record.left_asset_uid),
  487. "right_asset_uid": str(record.right_asset_uid),
  488. "canonical_asset_uid": (
  489. str(record.canonical_asset_uid)
  490. if record.canonical_asset_uid
  491. else None
  492. ),
  493. "status": record.status,
  494. "suggestion_source": record.suggestion_source,
  495. "confidence": float(record.confidence),
  496. "explanation": [
  497. dict(item) for item in record.explanation
  498. ],
  499. "evidence_uids": list(record.evidence_uids),
  500. "model_provider": record.model_provider,
  501. "model_name": record.model_name,
  502. "current_version": int(record.current_version),
  503. "created_by": record.created_by,
  504. "reviewed_by": record.reviewed_by,
  505. "created_at": _iso(record.created_at),
  506. "updated_at": _iso(record.updated_at),
  507. }
  508. def _device_entity_review(record):
  509. return {
  510. "uid": str(record.uid),
  511. "candidate_uid": str(record.candidate_uid),
  512. "version": int(record.version),
  513. "decision": record.decision,
  514. "reason": record.reason,
  515. "actor_uid": record.actor_uid,
  516. "created_at": _iso(record.created_at),
  517. }
  518. def _device_entity_merge(record):
  519. return {
  520. "uid": str(record.uid),
  521. "candidate_uid": str(record.candidate_uid),
  522. "canonical_asset_uid": str(record.canonical_asset_uid),
  523. "member_asset_uid": str(record.member_asset_uid),
  524. "review_uid": str(record.review_uid),
  525. "snapshot": dict(record.snapshot),
  526. "actor_uid": record.actor_uid,
  527. "created_at": _iso(record.created_at),
  528. }
  529. def _device_entity_rollback(record):
  530. return {
  531. "uid": str(record.uid),
  532. "merge_uid": str(record.merge_uid),
  533. "candidate_uid": str(record.candidate_uid),
  534. "reason": record.reason,
  535. "snapshot": dict(record.snapshot),
  536. "actor_uid": record.actor_uid,
  537. "created_at": _iso(record.created_at),
  538. }
  539. def _device_entity_generation(result):
  540. return {
  541. "records": [
  542. _device_entity_candidate(record)
  543. for record in result.records
  544. ],
  545. "created_count": int(result.created_count),
  546. "existing_count": int(result.existing_count),
  547. "evaluated_pair_count": int(result.evaluated_pair_count),
  548. "auto_merged_count": int(result.auto_merged_count),
  549. "auto_merge_enabled": bool(
  550. current_app.config.get(
  551. "DEVICE_ENTITY_AUTO_MERGE_ENABLED",
  552. False,
  553. )
  554. ),
  555. }
  556. def _device_asset_page(name, *, default, maximum):
  557. from app.core.data_research.errors import DeviceAssetInvalid
  558. raw = request.args.get(name)
  559. try:
  560. value = default if raw in (None, "") else int(raw)
  561. except (TypeError, ValueError) as error:
  562. raise DeviceAssetInvalid(f"{name} must be an integer") from error
  563. if value < 1 or value > maximum:
  564. raise DeviceAssetInvalid(
  565. f"{name} must be between 1 and {maximum}"
  566. )
  567. return value
  568. def _error(error):
  569. if isinstance(error, DataResearchError):
  570. return (
  571. jsonify(
  572. failed(
  573. str(error),
  574. code=error.http_status,
  575. error={"code": error.code},
  576. )
  577. ),
  578. error.http_status,
  579. )
  580. logger.exception("data-research ingestion request failed")
  581. return (
  582. jsonify(
  583. failed(
  584. "数据采集任务处理失败",
  585. code=500,
  586. error={"code": "DATA_RESEARCH_ERROR"},
  587. )
  588. ),
  589. 500,
  590. )
  591. @bp.route("/ingestion-jobs", methods=["POST"])
  592. def create_ingestion_job():
  593. payload = request.get_json(silent=True) or {}
  594. try:
  595. actor_uid = _identity().get("id") or _identity().get("sub")
  596. if payload.get("job_type") == "catalog_collect":
  597. get_database_source_registration_service().ensure(
  598. payload.get("source_uid"),
  599. actor_uid=actor_uid,
  600. )
  601. record, created = get_ingestion_service().create_job(
  602. payload,
  603. actor_uid=actor_uid,
  604. )
  605. return jsonify(success(_record(record))), 201 if created else 200
  606. except Exception as error:
  607. return _error(error)
  608. @bp.route("/ingestion-jobs", methods=["GET"])
  609. def list_ingestion_jobs():
  610. filters = {
  611. name: request.args.get(name)
  612. for name in ("status", "source_uid")
  613. if request.args.get(name)
  614. }
  615. try:
  616. records = get_ingestion_service().list_jobs(filters)
  617. return jsonify(
  618. success({"records": [_record(item) for item in records], "total": len(records)})
  619. ), 200
  620. except Exception as error:
  621. return _error(error)
  622. @bp.route("/ingestion-jobs/<job_uid>", methods=["GET"])
  623. def get_ingestion_job(job_uid):
  624. try:
  625. return jsonify(success(_record(get_ingestion_service().get_job(job_uid)))), 200
  626. except Exception as error:
  627. return _error(error)
  628. @bp.route("/ingestion-jobs/<job_uid>/execute", methods=["POST"])
  629. def execute_ingestion_job(job_uid):
  630. try:
  631. record = get_catalog_ingestion_executor().execute(job_uid)
  632. return jsonify(success(_record(record))), 200
  633. except Exception as error:
  634. return _error(error)
  635. @bp.route(
  636. "/ingestion-jobs/<job_uid>/catalog-snapshots",
  637. methods=["GET"],
  638. )
  639. def list_catalog_snapshots(job_uid):
  640. try:
  641. records = get_catalog_snapshot_repository().list(job_uid)
  642. return jsonify(
  643. success(
  644. {
  645. "records": [
  646. _catalog_snapshot(record)
  647. for record in records
  648. ],
  649. "total": len(records),
  650. }
  651. )
  652. ), 200
  653. except Exception as error:
  654. return _error(error)
  655. @bp.route("/ingestion-jobs/<job_uid>/evidence", methods=["GET"])
  656. def list_ingestion_job_evidence(job_uid):
  657. try:
  658. records = get_evidence_service().list_for_job(job_uid)
  659. return jsonify(
  660. success({"records": records, "total": len(records)})
  661. ), 200
  662. except Exception as error:
  663. return _error(error)
  664. @bp.route("/ingestion-jobs/<job_uid>/retry", methods=["POST"])
  665. def retry_ingestion_job(job_uid):
  666. try:
  667. return jsonify(success(_record(get_ingestion_service().retry(job_uid)))), 200
  668. except Exception as error:
  669. return _error(error)
  670. @bp.route("/ingestion-jobs/<job_uid>/cancel", methods=["POST"])
  671. def cancel_ingestion_job(job_uid):
  672. try:
  673. service = get_ingestion_service()
  674. record = service.get_job(job_uid)
  675. identity = _identity()
  676. permissions = set(identity.get("permissions") or [])
  677. actor_uid = identity.get("id") or identity.get("sub")
  678. if record.actor_uid != actor_uid and "ingestion:admin" not in permissions:
  679. return jsonify(failed("权限不足", code=403)), 403
  680. return jsonify(success(_record(service.cancel(job_uid)))), 200
  681. except Exception as error:
  682. return _error(error)
  683. @bp.route("/data-elements", methods=["GET"])
  684. def list_data_elements():
  685. try:
  686. records = get_data_element_service().list(
  687. status=request.args.get("status"),
  688. business_domain_uid=request.args.get("business_domain_uid"),
  689. )
  690. return jsonify(
  691. success(
  692. {
  693. "records": [_element(item) for item in records],
  694. "total": len(records),
  695. }
  696. )
  697. ), 200
  698. except Exception as error:
  699. return _error(error)
  700. @bp.route("/device-assets", methods=["GET"])
  701. def list_device_assets():
  702. filters = {
  703. name: request.args.get(name)
  704. for name in ("keyword", "asset_type", "status", "source_uid")
  705. if request.args.get(name)
  706. }
  707. try:
  708. page = _device_asset_page(
  709. "page",
  710. default=1,
  711. maximum=1_000_000,
  712. )
  713. page_size = _device_asset_page(
  714. "page_size",
  715. default=20,
  716. maximum=100,
  717. )
  718. records, total = get_device_asset_service().search(
  719. filters,
  720. page=page,
  721. page_size=page_size,
  722. )
  723. return jsonify(
  724. success(
  725. {
  726. "records": [
  727. _device_asset_detail(record)
  728. for record in records
  729. ],
  730. "total": int(total),
  731. "page": page,
  732. "page_size": page_size,
  733. }
  734. )
  735. ), 200
  736. except Exception as error:
  737. return _error(error)
  738. @bp.route("/device-assets/import", methods=["POST"])
  739. def import_device_assets():
  740. try:
  741. result = get_device_asset_service().import_records(
  742. request.get_json(silent=True) or {},
  743. actor_uid=_identity().get("id") or _identity().get("sub"),
  744. )
  745. return jsonify(success(_device_asset_import_result(result))), 200
  746. except Exception as error:
  747. return _error(error)
  748. @bp.route("/device-assets/<asset_uid>", methods=["GET"])
  749. def get_device_asset(asset_uid):
  750. try:
  751. return jsonify(
  752. success(
  753. _device_asset_detail(
  754. get_device_asset_service().get(asset_uid)
  755. )
  756. )
  757. ), 200
  758. except Exception as error:
  759. return _error(error)
  760. @bp.route("/device-assets/<asset_uid>/versions", methods=["GET"])
  761. def list_device_asset_versions(asset_uid):
  762. try:
  763. records = get_device_asset_service().versions(asset_uid)
  764. return jsonify(
  765. success(
  766. {
  767. "records": [
  768. _device_asset_version(record)
  769. for record in records
  770. ],
  771. "total": len(records),
  772. }
  773. )
  774. ), 200
  775. except Exception as error:
  776. return _error(error)
  777. @bp.route("/device-semantics/bootstrap", methods=["POST"])
  778. def bootstrap_device_semantics():
  779. try:
  780. result = get_device_semantic_service().bootstrap(
  781. request.get_json(silent=True) or {},
  782. actor_uid=_identity().get("id") or _identity().get("sub"),
  783. )
  784. return (
  785. jsonify(success(_device_semantic_profile(result))),
  786. 201 if result.created else 200,
  787. )
  788. except Exception as error:
  789. return _error(error)
  790. @bp.route("/device-semantics/profile", methods=["GET"])
  791. def get_device_semantic_profile():
  792. try:
  793. result = get_device_semantic_service().profile()
  794. return jsonify(success(_device_semantic_profile(result))), 200
  795. except Exception as error:
  796. return _error(error)
  797. @bp.route("/device-semantics/codes", methods=["GET"])
  798. def list_device_semantic_codes():
  799. filters = {
  800. name: request.args.get(name)
  801. for name in ("ontology_uid", "code_type", "status", "keyword")
  802. if request.args.get(name)
  803. }
  804. try:
  805. records, total = get_device_semantic_code_service().search(
  806. filters,
  807. page=request.args.get("page", 1),
  808. page_size=request.args.get("page_size", 20),
  809. )
  810. return jsonify(
  811. success(
  812. {
  813. "records": [
  814. _device_semantic_code(record)
  815. for record in records
  816. ],
  817. "total": int(total),
  818. "page": int(request.args.get("page", 1)),
  819. "page_size": int(request.args.get("page_size", 20)),
  820. }
  821. )
  822. ), 200
  823. except Exception as error:
  824. return _error(error)
  825. @bp.route("/device-semantics/codes", methods=["POST"])
  826. def create_device_semantic_code():
  827. try:
  828. record = get_device_semantic_code_service().create(
  829. request.get_json(silent=True) or {},
  830. actor_uid=_identity().get("id") or _identity().get("sub"),
  831. )
  832. return jsonify(success(_device_semantic_code(record))), 201
  833. except Exception as error:
  834. return _error(error)
  835. @bp.route("/device-semantics/codes/<code_uid>", methods=["GET"])
  836. def get_device_semantic_code(code_uid):
  837. try:
  838. record = get_device_semantic_code_service().get(code_uid)
  839. return jsonify(success(_device_semantic_code(record))), 200
  840. except Exception as error:
  841. return _error(error)
  842. @bp.route(
  843. "/device-semantics/codes/<code_uid>/revisions",
  844. methods=["POST"],
  845. )
  846. def revise_device_semantic_code(code_uid):
  847. payload = request.get_json(silent=True) or {}
  848. try:
  849. record = get_device_semantic_code_service().revise(
  850. code_uid,
  851. payload,
  852. expected_version=payload.get("expected_version"),
  853. actor_uid=_identity().get("id") or _identity().get("sub"),
  854. )
  855. return jsonify(success(_device_semantic_code(record))), 200
  856. except Exception as error:
  857. return _error(error)
  858. @bp.route(
  859. "/device-semantics/codes/<code_uid>/submit",
  860. methods=["POST"],
  861. )
  862. def submit_device_semantic_code(code_uid):
  863. payload = request.get_json(silent=True) or {}
  864. try:
  865. record = get_device_semantic_code_service().submit(
  866. code_uid,
  867. expected_version=payload.get("expected_version"),
  868. actor_uid=_identity().get("id") or _identity().get("sub"),
  869. )
  870. return jsonify(success(_device_semantic_code(record))), 200
  871. except Exception as error:
  872. return _error(error)
  873. @bp.route(
  874. "/device-semantics/codes/<code_uid>/review",
  875. methods=["POST"],
  876. )
  877. def review_device_semantic_code(code_uid):
  878. try:
  879. record, review = get_device_semantic_code_service().review(
  880. code_uid,
  881. request.get_json(silent=True) or {},
  882. actor_uid=_identity().get("id") or _identity().get("sub"),
  883. )
  884. return jsonify(
  885. success(
  886. {
  887. "record": _device_semantic_code(record),
  888. "review": _device_semantic_review(review),
  889. }
  890. )
  891. ), 200
  892. except Exception as error:
  893. return _error(error)
  894. @bp.route(
  895. "/device-semantics/codes/<code_uid>/versions",
  896. methods=["GET"],
  897. )
  898. def list_device_semantic_code_versions(code_uid):
  899. try:
  900. records = get_device_semantic_code_service().versions(code_uid)
  901. return jsonify(
  902. success(
  903. {
  904. "records": [
  905. _device_semantic_version(record)
  906. for record in records
  907. ],
  908. "total": len(records),
  909. }
  910. )
  911. ), 200
  912. except Exception as error:
  913. return _error(error)
  914. @bp.route(
  915. "/device-semantics/codes/<code_uid>/reviews",
  916. methods=["GET"],
  917. )
  918. def list_device_semantic_code_reviews(code_uid):
  919. try:
  920. records = get_device_semantic_code_service().reviews(code_uid)
  921. return jsonify(
  922. success(
  923. {
  924. "records": [
  925. _device_semantic_review(record)
  926. for record in records
  927. ],
  928. "total": len(records),
  929. }
  930. )
  931. ), 200
  932. except Exception as error:
  933. return _error(error)
  934. @bp.route("/device-entities/candidates", methods=["GET"])
  935. def list_device_entity_candidates():
  936. filters = {
  937. name: request.args.get(name)
  938. for name in ("status", "suggestion_source")
  939. if request.args.get(name)
  940. }
  941. try:
  942. records, total = get_device_entity_resolution_service().search(
  943. filters,
  944. page=request.args.get("page", 1),
  945. page_size=request.args.get("page_size", 20),
  946. )
  947. return jsonify(
  948. success(
  949. {
  950. "records": [
  951. _device_entity_candidate(record)
  952. for record in records
  953. ],
  954. "total": int(total),
  955. "page": int(request.args.get("page", 1)),
  956. "page_size": int(request.args.get("page_size", 20)),
  957. "auto_merge_enabled": bool(
  958. current_app.config.get(
  959. "DEVICE_ENTITY_AUTO_MERGE_ENABLED",
  960. False,
  961. )
  962. ),
  963. }
  964. )
  965. ), 200
  966. except Exception as error:
  967. return _error(error)
  968. @bp.route("/device-entities/candidates/generate", methods=["POST"])
  969. def generate_device_entity_candidates():
  970. try:
  971. result = get_device_entity_resolution_service().generate(
  972. request.get_json(silent=True) or {},
  973. actor_uid=_identity().get("id") or _identity().get("sub"),
  974. )
  975. status = 201 if result.created_count else 200
  976. return jsonify(success(_device_entity_generation(result))), status
  977. except Exception as error:
  978. return _error(error)
  979. @bp.route("/device-entities/candidates", methods=["POST"])
  980. def submit_device_entity_candidate():
  981. try:
  982. record = (
  983. get_device_entity_resolution_service().submit_ai_candidate(
  984. request.get_json(silent=True) or {},
  985. actor_uid=_identity().get("id") or _identity().get("sub"),
  986. )
  987. )
  988. return jsonify(success(_device_entity_candidate(record))), 201
  989. except Exception as error:
  990. return _error(error)
  991. @bp.route("/device-entities/candidates/<candidate_uid>", methods=["GET"])
  992. def get_device_entity_candidate(candidate_uid):
  993. try:
  994. record = get_device_entity_resolution_service().get(candidate_uid)
  995. return jsonify(success(_device_entity_candidate(record))), 200
  996. except Exception as error:
  997. return _error(error)
  998. @bp.route(
  999. "/device-entities/candidates/<candidate_uid>/review",
  1000. methods=["POST"],
  1001. )
  1002. def review_device_entity_candidate(candidate_uid):
  1003. try:
  1004. candidate, review, merge = (
  1005. get_device_entity_resolution_service().review(
  1006. candidate_uid,
  1007. request.get_json(silent=True) or {},
  1008. actor_uid=_identity().get("id") or _identity().get("sub"),
  1009. )
  1010. )
  1011. return jsonify(
  1012. success(
  1013. {
  1014. "candidate": _device_entity_candidate(candidate),
  1015. "review": _device_entity_review(review),
  1016. "merge": (
  1017. _device_entity_merge(merge)
  1018. if merge is not None
  1019. else None
  1020. ),
  1021. }
  1022. )
  1023. ), 200
  1024. except Exception as error:
  1025. return _error(error)
  1026. @bp.route(
  1027. "/device-entities/candidates/<candidate_uid>/reviews",
  1028. methods=["GET"],
  1029. )
  1030. def list_device_entity_reviews(candidate_uid):
  1031. try:
  1032. records = get_device_entity_resolution_service().reviews(
  1033. candidate_uid
  1034. )
  1035. return jsonify(
  1036. success(
  1037. {
  1038. "records": [
  1039. _device_entity_review(record)
  1040. for record in records
  1041. ],
  1042. "total": len(records),
  1043. }
  1044. )
  1045. ), 200
  1046. except Exception as error:
  1047. return _error(error)
  1048. @bp.route(
  1049. "/device-entities/candidates/<candidate_uid>/merges",
  1050. methods=["GET"],
  1051. )
  1052. def list_device_entity_merges(candidate_uid):
  1053. try:
  1054. records = get_device_entity_resolution_service().merges(
  1055. candidate_uid
  1056. )
  1057. return jsonify(
  1058. success(
  1059. {
  1060. "records": [
  1061. _device_entity_merge(record)
  1062. for record in records
  1063. ],
  1064. "total": len(records),
  1065. }
  1066. )
  1067. ), 200
  1068. except Exception as error:
  1069. return _error(error)
  1070. @bp.route(
  1071. "/device-entities/merges/<merge_uid>/rollback",
  1072. methods=["POST"],
  1073. )
  1074. def rollback_device_entity_merge(merge_uid):
  1075. try:
  1076. candidate, rollback = (
  1077. get_device_entity_resolution_service().rollback(
  1078. merge_uid,
  1079. request.get_json(silent=True) or {},
  1080. actor_uid=_identity().get("id") or _identity().get("sub"),
  1081. )
  1082. )
  1083. return jsonify(
  1084. success(
  1085. {
  1086. "candidate": _device_entity_candidate(candidate),
  1087. "rollback": _device_entity_rollback(rollback),
  1088. }
  1089. )
  1090. ), 200
  1091. except Exception as error:
  1092. return _error(error)
  1093. @bp.route(
  1094. "/device-entities/merges/<merge_uid>/rollbacks",
  1095. methods=["GET"],
  1096. )
  1097. def list_device_entity_rollbacks(merge_uid):
  1098. try:
  1099. records = get_device_entity_resolution_service().rollbacks(
  1100. merge_uid
  1101. )
  1102. return jsonify(
  1103. success(
  1104. {
  1105. "records": [
  1106. _device_entity_rollback(record)
  1107. for record in records
  1108. ],
  1109. "total": len(records),
  1110. }
  1111. )
  1112. ), 200
  1113. except Exception as error:
  1114. return _error(error)
  1115. @bp.route("/data-elements", methods=["POST"])
  1116. def create_data_element():
  1117. try:
  1118. record = get_data_element_service().create_draft(
  1119. request.get_json(silent=True) or {},
  1120. actor_uid=_identity().get("id") or _identity().get("sub"),
  1121. )
  1122. db.session.commit()
  1123. return jsonify(success(_element(record))), 201
  1124. except Exception as error:
  1125. db.session.rollback()
  1126. return _error(error)
  1127. @bp.route("/data-elements/<element_uid>/transitions", methods=["POST"])
  1128. def transition_data_element(element_uid):
  1129. payload = request.get_json(silent=True) or {}
  1130. identity = _identity()
  1131. target_status = str(payload.get("target_status") or "")
  1132. if (
  1133. target_status in {"published", "deprecated", "retired"}
  1134. and "data-elements:publish" not in set(identity.get("permissions") or [])
  1135. ):
  1136. return jsonify(failed("权限不足", code=403)), 403
  1137. try:
  1138. record = get_data_element_service().transition(
  1139. element_uid,
  1140. target_status,
  1141. expected_version=int(payload.get("expected_version")),
  1142. actor_uid=identity.get("id") or identity.get("sub"),
  1143. )
  1144. db.session.commit()
  1145. return jsonify(success(_element(record))), 200
  1146. except Exception as error:
  1147. db.session.rollback()
  1148. return _error(error)
  1149. @bp.route("/candidate-decisions", methods=["POST"])
  1150. def decide_candidates():
  1151. payload = request.get_json(silent=True) or {}
  1152. try:
  1153. records = get_candidate_decision_service().decide(
  1154. payload.get("decisions") or [],
  1155. actor_uid=_identity().get("id") or _identity().get("sub"),
  1156. )
  1157. return jsonify(success([_decision(record) for record in records])), 200
  1158. except Exception as error:
  1159. return _error(error)
  1160. @bp.route("/sources/files", methods=["POST"])
  1161. def upload_source_file():
  1162. uploaded = request.files.get("file")
  1163. if uploaded is None:
  1164. return jsonify(failed("缺少上传文件", code=400)), 400
  1165. try:
  1166. record, created = get_artifact_service().store(
  1167. request.form.get("source_uid"),
  1168. uploaded.filename,
  1169. uploaded.mimetype,
  1170. uploaded.read(),
  1171. request.form.get("parser_version") or "auto-v1",
  1172. )
  1173. db.session.commit()
  1174. return jsonify(success(_artifact(record))), 201 if created else 200
  1175. except Exception as error:
  1176. db.session.rollback()
  1177. return _error(error)
  1178. @bp.route("/evidence/<evidence_uid>", methods=["GET"])
  1179. def get_evidence(evidence_uid):
  1180. identity = _identity()
  1181. try:
  1182. service = get_evidence_service()
  1183. if request.args.get("download") in {"1", "true", "yes"}:
  1184. if "evidence:download" not in set(identity.get("permissions") or []):
  1185. return jsonify(failed("权限不足", code=403)), 403
  1186. content, filename, media_type = service.download(evidence_uid)
  1187. return send_file(
  1188. io.BytesIO(content),
  1189. mimetype=media_type,
  1190. as_attachment=True,
  1191. download_name=filename,
  1192. )
  1193. preview = service.preview(evidence_uid)
  1194. from app.core.data_research.artifacts import redact_excerpt
  1195. preview["excerpt"] = redact_excerpt(preview.get("excerpt"))
  1196. return jsonify(success(preview)), 200
  1197. except Exception as error:
  1198. return _error(error)
  1199. @bp.route("/ontologies", methods=["GET"])
  1200. def list_ontologies():
  1201. try:
  1202. return jsonify(success([_ontology(item) for item in get_ontology_service().list()])), 200
  1203. except Exception as error:
  1204. return _error(error)
  1205. @bp.route("/ontologies", methods=["POST"])
  1206. def create_ontology():
  1207. try:
  1208. record = get_ontology_service().create(
  1209. request.get_json(silent=True) or {},
  1210. _identity().get("id") or _identity().get("sub"),
  1211. )
  1212. return jsonify(success(_ontology(record))), 201
  1213. except Exception as error:
  1214. return _error(error)
  1215. @bp.route("/ontologies/<ontology_uid>", methods=["GET"])
  1216. def get_ontology(ontology_uid):
  1217. try:
  1218. return jsonify(success(_ontology(get_ontology_service().get(ontology_uid)))), 200
  1219. except Exception as error:
  1220. return _error(error)
  1221. @bp.route("/ontologies/<ontology_uid>/versions", methods=["GET"])
  1222. def list_ontology_versions(ontology_uid):
  1223. try:
  1224. return jsonify(
  1225. success(
  1226. [
  1227. _ontology_version(item)
  1228. for item in get_ontology_service().list_versions(ontology_uid)
  1229. ]
  1230. )
  1231. ), 200
  1232. except Exception as error:
  1233. return _error(error)
  1234. @bp.route("/ontologies/<ontology_uid>/graph", methods=["GET"])
  1235. def get_ontology_graph(ontology_uid):
  1236. try:
  1237. service = get_ontology_service()
  1238. ontology = service.get(ontology_uid)
  1239. version = service.latest_version(ontology_uid)
  1240. data = (
  1241. _ontology_version(version)
  1242. if version is not None
  1243. else {
  1244. "uid": None,
  1245. "ontology_uid": str(ontology_uid),
  1246. "version": 0,
  1247. "parent_version_uid": None,
  1248. "status": "draft",
  1249. "content_hash": None,
  1250. "graph_document": {
  1251. name: []
  1252. for name in (
  1253. "classes",
  1254. "properties",
  1255. "relations",
  1256. "constraints",
  1257. "domain_links",
  1258. "element_mappings",
  1259. )
  1260. },
  1261. "created_by": None,
  1262. }
  1263. )
  1264. response = jsonify(success(data))
  1265. response.headers["ETag"] = f'"{int(ontology.draft_revision)}"'
  1266. return response, 200
  1267. except Exception as error:
  1268. return _error(error)
  1269. @bp.route("/ontologies/<ontology_uid>/graph", methods=["PATCH"])
  1270. def save_ontology_graph(ontology_uid):
  1271. etag = str(request.headers.get("If-Match") or "").strip().strip('"')
  1272. if not etag.isdigit():
  1273. return jsonify(failed("缺少有效 If-Match 修订号", code=428)), 428
  1274. try:
  1275. version = get_ontology_service().save_draft(
  1276. ontology_uid,
  1277. request.get_json(silent=True) or {},
  1278. int(etag),
  1279. _identity().get("id") or _identity().get("sub"),
  1280. )
  1281. response = jsonify(success(_ontology_version(version)))
  1282. response.headers["ETag"] = f'"{int(etag) + 1}"'
  1283. return response, 200
  1284. except Exception as error:
  1285. return _error(error)
  1286. @bp.route("/ontologies/<ontology_uid>/validate", methods=["POST"])
  1287. def validate_ontology(ontology_uid):
  1288. try:
  1289. issues = get_ontology_service().validate(ontology_uid)
  1290. data = [
  1291. {"code": item.code, "message": item.message, "path": item.path}
  1292. for item in issues
  1293. ]
  1294. return jsonify(success(data)), 200
  1295. except Exception as error:
  1296. return _error(error)
  1297. @bp.route("/ontologies/<ontology_uid>/publish", methods=["POST"])
  1298. def publish_ontology(ontology_uid):
  1299. try:
  1300. version = get_ontology_service().publish(
  1301. ontology_uid,
  1302. request.headers.get("Idempotency-Key") or f"{ontology_uid}:publish",
  1303. _identity().get("id") or _identity().get("sub"),
  1304. )
  1305. return jsonify(success(_ontology_version(version))), 200
  1306. except Exception as error:
  1307. return _error(error)
  1308. @bp.route("/ontologies/<ontology_uid>/diff", methods=["GET"])
  1309. def diff_ontology(ontology_uid):
  1310. try:
  1311. return jsonify(success(get_ontology_service().diff(
  1312. ontology_uid, request.args.get("left"), request.args.get("right")
  1313. ))), 200
  1314. except Exception as error:
  1315. return _error(error)
  1316. @bp.route("/ontologies/<ontology_uid>/rollback", methods=["POST"])
  1317. def rollback_ontology(ontology_uid):
  1318. payload = request.get_json(silent=True) or {}
  1319. try:
  1320. version = get_ontology_service().rollback(
  1321. ontology_uid,
  1322. payload.get("target_version_uid"),
  1323. int(payload.get("expected_revision")),
  1324. _identity().get("id") or _identity().get("sub"),
  1325. )
  1326. return jsonify(success(_ontology_version(version))), 201
  1327. except Exception as error:
  1328. return _error(error)
  1329. @bp.route("/ontologies/<ontology_uid>/suggestions", methods=["POST"])
  1330. def generate_ontology_suggestions(ontology_uid):
  1331. try:
  1332. result = get_ontology_dynamic_service().generate(
  1333. ontology_uid,
  1334. request.get_json(silent=True) or {},
  1335. _identity().get("id") or _identity().get("sub"),
  1336. )
  1337. records = (
  1338. result.get("suggestions") or []
  1339. if isinstance(result, dict)
  1340. else result
  1341. )
  1342. suggestions = [
  1343. item
  1344. if isinstance(item, dict)
  1345. else {
  1346. "uid": item.uid,
  1347. "kind": item.kind,
  1348. "payload": item.payload,
  1349. "evidence_uids": list(item.evidence_uids),
  1350. "confidence": item.confidence,
  1351. "source": item.source,
  1352. "model_version": item.model_version,
  1353. "prompt_version": item.prompt_version,
  1354. }
  1355. for item in records
  1356. ]
  1357. data = {
  1358. "change_set_uid": (
  1359. result.get("change_set_uid")
  1360. if isinstance(result, dict)
  1361. else None
  1362. ),
  1363. "suggestions": suggestions,
  1364. }
  1365. return jsonify(success(data)), 201
  1366. except Exception as error:
  1367. return _error(error)
  1368. @bp.route(
  1369. "/ontologies/<ontology_uid>/change-sets/<change_set_uid>/decisions",
  1370. methods=["POST"],
  1371. )
  1372. def decide_ontology_change_set(ontology_uid, change_set_uid):
  1373. payload = request.get_json(silent=True) or {}
  1374. try:
  1375. result = get_ontology_dynamic_service().decide(
  1376. ontology_uid,
  1377. change_set_uid,
  1378. payload.get("decisions") or [],
  1379. _identity().get("id") or _identity().get("sub"),
  1380. )
  1381. return jsonify(success(result)), 200
  1382. except Exception as error:
  1383. return _error(error)
  1384. @bp.route("/ontologies/<ontology_uid>/export", methods=["GET"])
  1385. def export_ontology(ontology_uid):
  1386. try:
  1387. content, media_type, filename = get_ontology_exchange_service().export(
  1388. ontology_uid, request.args.get("format") or "json"
  1389. )
  1390. return send_file(
  1391. io.BytesIO(content),
  1392. mimetype=media_type,
  1393. as_attachment=True,
  1394. download_name=filename,
  1395. )
  1396. except Exception as error:
  1397. return _error(error)
  1398. @bp.route("/ontologies/import", methods=["POST"])
  1399. def import_ontology():
  1400. try:
  1401. uploaded = request.files.get("file")
  1402. content = uploaded.read() if uploaded is not None else request.get_data(cache=False)
  1403. result = get_ontology_exchange_service().import_document(
  1404. content,
  1405. request.args.get("format") or "json",
  1406. _identity().get("id") or _identity().get("sub"),
  1407. )
  1408. return jsonify(success(result)), 201
  1409. except Exception as error:
  1410. return _error(error)
  1411. @bp.route("/semantic/properties/<property_uid>", methods=["GET"])
  1412. def query_semantic_property(property_uid):
  1413. domain = str(request.args.get("business_domain_uid") or "").strip()
  1414. identity = _identity()
  1415. scoped = set(identity.get("business_domains") or [domain])
  1416. if "*" not in scoped and domain not in scoped:
  1417. return jsonify(failed("权限不足", code=403)), 403
  1418. try:
  1419. result = get_semantic_query_service().trace_property(
  1420. property_uid,
  1421. allowed_domains={domain},
  1422. limit=int(request.args.get("limit") or 20),
  1423. after_uid=request.args.get("after_uid"),
  1424. )
  1425. return jsonify(success(result)), 200
  1426. except Exception as error:
  1427. return _error(error)