routes.py 72 KB

1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071727374757677787980818283848586878889909192939495969798991001011021031041051061071081091101111121131141151161171181191201211221231241251261271281291301311321331341351361371381391401411421431441451461471481491501511521531541551561571581591601611621631641651661671681691701711721731741751761771781791801811821831841851861871881891901911921931941951961971981992002012022032042052062072082092102112122132142152162172182192202212222232242252262272282292302312322332342352362372382392402412422432442452462472482492502512522532542552562572582592602612622632642652662672682692702712722732742752762772782792802812822832842852862872882892902912922932942952962972982993003013023033043053063073083093103113123133143153163173183193203213223233243253263273283293303313323333343353363373383393403413423433443453463473483493503513523533543553563573583593603613623633643653663673683693703713723733743753763773783793803813823833843853863873883893903913923933943953963973983994004014024034044054064074084094104114124134144154164174184194204214224234244254264274284294304314324334344354364374384394404414424434444454464474484494504514524534544554564574584594604614624634644654664674684694704714724734744754764774784794804814824834844854864874884894904914924934944954964974984995005015025035045055065075085095105115125135145155165175185195205215225235245255265275285295305315325335345355365375385395405415425435445455465475485495505515525535545555565575585595605615625635645655665675685695705715725735745755765775785795805815825835845855865875885895905915925935945955965975985996006016026036046056066076086096106116126136146156166176186196206216226236246256266276286296306316326336346356366376386396406416426436446456466476486496506516526536546556566576586596606616626636646656666676686696706716726736746756766776786796806816826836846856866876886896906916926936946956966976986997007017027037047057067077087097107117127137147157167177187197207217227237247257267277287297307317327337347357367377387397407417427437447457467477487497507517527537547557567577587597607617627637647657667677687697707717727737747757767777787797807817827837847857867877887897907917927937947957967977987998008018028038048058068078088098108118128138148158168178188198208218228238248258268278288298308318328338348358368378388398408418428438448458468478488498508518528538548558568578588598608618628638648658668678688698708718728738748758768778788798808818828838848858868878888898908918928938948958968978988999009019029039049059069079089099109119129139149159169179189199209219229239249259269279289299309319329339349359369379389399409419429439449459469479489499509519529539549559569579589599609619629639649659669679689699709719729739749759769779789799809819829839849859869879889899909919929939949959969979989991000100110021003100410051006100710081009101010111012101310141015101610171018101910201021102210231024102510261027102810291030103110321033103410351036103710381039104010411042104310441045104610471048104910501051105210531054105510561057105810591060106110621063106410651066106710681069107010711072107310741075107610771078107910801081108210831084108510861087108810891090109110921093109410951096109710981099110011011102110311041105110611071108110911101111111211131114111511161117111811191120112111221123112411251126112711281129113011311132113311341135113611371138113911401141114211431144114511461147114811491150115111521153115411551156115711581159116011611162116311641165116611671168116911701171117211731174117511761177117811791180118111821183118411851186118711881189119011911192119311941195119611971198119912001201120212031204120512061207120812091210121112121213121412151216121712181219122012211222122312241225122612271228122912301231123212331234123512361237123812391240124112421243124412451246124712481249125012511252125312541255125612571258125912601261126212631264126512661267126812691270127112721273127412751276127712781279128012811282128312841285128612871288128912901291129212931294129512961297129812991300130113021303130413051306130713081309131013111312131313141315131613171318131913201321132213231324132513261327132813291330133113321333133413351336133713381339134013411342134313441345134613471348134913501351135213531354135513561357135813591360136113621363136413651366136713681369137013711372137313741375137613771378137913801381138213831384138513861387138813891390139113921393139413951396139713981399140014011402140314041405140614071408140914101411141214131414141514161417141814191420142114221423142414251426142714281429143014311432143314341435143614371438143914401441144214431444144514461447144814491450145114521453145414551456145714581459146014611462146314641465146614671468146914701471147214731474147514761477147814791480148114821483148414851486148714881489149014911492149314941495149614971498149915001501150215031504150515061507150815091510151115121513151415151516151715181519152015211522152315241525152615271528152915301531153215331534153515361537153815391540154115421543154415451546154715481549155015511552155315541555155615571558155915601561156215631564156515661567156815691570157115721573157415751576157715781579158015811582158315841585158615871588158915901591159215931594159515961597159815991600160116021603160416051606160716081609161016111612161316141615161616171618161916201621162216231624162516261627162816291630163116321633163416351636163716381639164016411642164316441645164616471648164916501651165216531654165516561657165816591660166116621663166416651666166716681669167016711672167316741675167616771678167916801681168216831684168516861687168816891690169116921693169416951696169716981699170017011702170317041705170617071708170917101711171217131714171517161717171817191720172117221723172417251726172717281729173017311732173317341735173617371738173917401741174217431744174517461747174817491750175117521753175417551756175717581759176017611762176317641765176617671768176917701771177217731774177517761777177817791780178117821783178417851786178717881789179017911792179317941795179617971798179918001801180218031804180518061807180818091810181118121813181418151816181718181819182018211822182318241825182618271828182918301831183218331834183518361837183818391840184118421843184418451846184718481849185018511852185318541855185618571858185918601861186218631864186518661867186818691870187118721873187418751876187718781879188018811882188318841885188618871888188918901891189218931894189518961897189818991900190119021903190419051906190719081909191019111912191319141915191619171918191919201921192219231924192519261927192819291930193119321933193419351936193719381939194019411942194319441945194619471948194919501951195219531954195519561957195819591960196119621963196419651966196719681969197019711972197319741975197619771978197919801981198219831984198519861987198819891990199119921993199419951996199719981999200020012002200320042005200620072008200920102011201220132014201520162017201820192020202120222023202420252026202720282029203020312032203320342035203620372038203920402041204220432044204520462047204820492050205120522053205420552056205720582059206020612062206320642065206620672068206920702071207220732074207520762077207820792080208120822083208420852086208720882089209020912092209320942095209620972098209921002101210221032104210521062107210821092110211121122113211421152116211721182119212021212122212321242125212621272128212921302131213221332134213521362137213821392140214121422143214421452146214721482149215021512152215321542155215621572158215921602161216221632164216521662167216821692170217121722173217421752176217721782179218021812182218321842185218621872188218921902191219221932194219521962197219821992200220122022203220422052206220722082209221022112212221322142215221622172218221922202221222222232224222522262227222822292230223122322233223422352236223722382239224022412242224322442245224622472248224922502251225222532254225522562257225822592260226122622263226422652266226722682269227022712272227322742275227622772278227922802281228222832284228522862287228822892290229122922293229422952296229722982299230023012302
  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.common.timezone_utils import now_china
  9. from app.core.data_research.errors import DataResearchError
  10. from app.models.result import failed, success
  11. logger = logging.getLogger(__name__)
  12. def get_ingestion_service():
  13. from app.core.data_research.ingestion import IngestionService
  14. from app.core.data_research.repository import SqlAlchemyIngestionJobRepository
  15. return IngestionService(
  16. SqlAlchemyIngestionJobRepository(db.session),
  17. commit=db.session.commit,
  18. rollback=db.session.rollback,
  19. )
  20. def get_database_source_registration_service():
  21. from app.core.data_research.repository import (
  22. SqlAlchemyIngestionSourceRepository,
  23. )
  24. from app.core.data_research.sources import (
  25. DatabaseSourceRegistrationService,
  26. )
  27. from app.core.data_source.runtime import get_data_source_manager
  28. manager = get_data_source_manager()
  29. return DatabaseSourceRegistrationService(
  30. SqlAlchemyIngestionSourceRepository(db.session),
  31. definition_resolver=manager.definitions.get,
  32. commit=db.session.commit,
  33. rollback=db.session.rollback,
  34. )
  35. def get_catalog_snapshot_repository():
  36. from app.core.data_research.repository import (
  37. SqlAlchemyCatalogSnapshotRepository,
  38. )
  39. return SqlAlchemyCatalogSnapshotRepository(db.session)
  40. def get_catalog_ingestion_executor():
  41. from app.core.data_research.catalog.execution import (
  42. CatalogIngestionExecutor,
  43. )
  44. from app.core.data_research.catalog.mysql import MySqlCatalogCollector
  45. from app.core.data_research.catalog.postgresql import (
  46. PostgreSqlCatalogCollector,
  47. )
  48. from app.core.data_research.catalog.service import CatalogCollectionService
  49. from app.core.data_research.errors import IngestionSourceInvalid
  50. from app.core.data_source.runtime import get_data_source_manager
  51. manager = get_data_source_manager()
  52. def collector_resolver(database_type):
  53. if database_type == "postgresql":
  54. return PostgreSqlCatalogCollector()
  55. if database_type == "mysql":
  56. return MySqlCatalogCollector()
  57. raise IngestionSourceInvalid(
  58. f"database type {database_type} is not supported"
  59. )
  60. collector = CatalogCollectionService(
  61. manager,
  62. definition_resolver=manager.definitions.get,
  63. collector_resolver=collector_resolver,
  64. )
  65. return CatalogIngestionExecutor(
  66. get_ingestion_service(),
  67. collector,
  68. get_catalog_snapshot_repository(),
  69. commit=db.session.commit,
  70. rollback=db.session.rollback,
  71. )
  72. def get_data_element_service():
  73. from app.core.data_research.data_elements import DataElementService
  74. from app.core.data_research.repository import SqlAlchemyDataElementRepository
  75. from app.core.events.outbox import enqueue_outbox
  76. return DataElementService(
  77. SqlAlchemyDataElementRepository(db.session),
  78. outbox_enqueue=lambda **event: enqueue_outbox(db.session, **event),
  79. )
  80. def get_device_asset_service():
  81. from app.core.data_research.device_asset_repository import (
  82. SqlAlchemyDeviceAssetRepository,
  83. )
  84. from app.core.data_research.device_assets import DeviceAssetService
  85. return DeviceAssetService(
  86. SqlAlchemyDeviceAssetRepository(db.session),
  87. commit=db.session.commit,
  88. rollback=db.session.rollback,
  89. )
  90. def _device_ontology_authorizer(ontology_repository):
  91. from app.core.data_research.device_semantics import (
  92. DeviceOntologyPublicationAuthorizer,
  93. )
  94. from app.core.governance.responsibilities import (
  95. ResponsibilityService,
  96. SqlAlchemyResponsibilityRepository,
  97. )
  98. responsibilities = ResponsibilityService(
  99. SqlAlchemyResponsibilityRepository(db.session)
  100. )
  101. return DeviceOntologyPublicationAuthorizer(
  102. ontology_repository,
  103. responsibility_lookup=responsibilities.get,
  104. )
  105. def get_device_semantic_service():
  106. from app.core.data_research.device_semantics import DeviceSemanticService
  107. from app.core.data_research.ontology.repository import (
  108. SqlAlchemyOntologyRepository,
  109. )
  110. return DeviceSemanticService(
  111. SqlAlchemyOntologyRepository(db.session),
  112. commit=db.session.commit,
  113. rollback=db.session.rollback,
  114. )
  115. def get_device_semantic_code_service():
  116. from app.core.data_research.device_semantic_repository import (
  117. SqlAlchemyDeviceSemanticCodeRepository,
  118. )
  119. from app.core.data_research.device_semantics import (
  120. DeviceSemanticCodeService,
  121. )
  122. from app.core.data_research.ontology.repository import (
  123. SqlAlchemyOntologyRepository,
  124. )
  125. ontology_repository = SqlAlchemyOntologyRepository(db.session)
  126. authorizer = _device_ontology_authorizer(ontology_repository)
  127. return DeviceSemanticCodeService(
  128. SqlAlchemyDeviceSemanticCodeRepository(db.session),
  129. review_authorizer=authorizer.assert_accountable,
  130. commit=db.session.commit,
  131. rollback=db.session.rollback,
  132. )
  133. def get_device_entity_resolution_service():
  134. from app.core.data_research.device_entity_repository import (
  135. SqlAlchemyDeviceEntityResolutionRepository,
  136. )
  137. from app.core.data_research.device_entity_resolution import (
  138. DeviceEntityForbidden,
  139. DeviceEntityResolutionService,
  140. )
  141. from app.core.governance.responsibilities import (
  142. ResponsibilityService,
  143. SqlAlchemyResponsibilityRepository,
  144. )
  145. responsibilities = ResponsibilityService(
  146. SqlAlchemyResponsibilityRepository(db.session)
  147. )
  148. def assert_accountable(actor_uid):
  149. matrix = responsibilities.get(
  150. "device_mapping",
  151. "DEVICE_ENTITY_RESOLUTION",
  152. )
  153. accountable = [
  154. item
  155. for item in matrix.get("assignments", [])
  156. if item.get("responsibility_role") == "asset_manager"
  157. and item.get("raci_role") == "accountable"
  158. ]
  159. if len(accountable) != 1 or str(accountable[0].get("user_id")) != str(
  160. actor_uid
  161. ):
  162. raise DeviceEntityForbidden(
  163. "only the accountable device mapping asset manager may decide"
  164. )
  165. return DeviceEntityResolutionService(
  166. SqlAlchemyDeviceEntityResolutionRepository(db.session),
  167. review_authorizer=assert_accountable,
  168. auto_merge_enabled=current_app.config.get(
  169. "DEVICE_ENTITY_AUTO_MERGE_ENABLED",
  170. False,
  171. ),
  172. commit=db.session.commit,
  173. rollback=db.session.rollback,
  174. )
  175. def get_device_quality_service():
  176. from app.core.data_research.device_quality import DeviceQualityService
  177. from app.core.data_research.device_quality_repository import (
  178. SqlAlchemyDeviceQualityRepository,
  179. )
  180. from app.core.data_research.errors import DeviceQualityForbidden
  181. from app.core.governance.responsibilities import (
  182. ResponsibilityService,
  183. SqlAlchemyResponsibilityRepository,
  184. )
  185. responsibilities = ResponsibilityService(
  186. SqlAlchemyResponsibilityRepository(db.session)
  187. )
  188. def assert_accountable(actor_uid):
  189. matrix = responsibilities.get(
  190. "device_quality",
  191. "DEVICE_QUALITY",
  192. )
  193. accountable = [
  194. item
  195. for item in matrix.get("assignments", [])
  196. if item.get("responsibility_role") == "asset_manager"
  197. and item.get("raci_role") == "accountable"
  198. ]
  199. if len(accountable) != 1 or str(accountable[0].get("user_id")) != str(
  200. actor_uid
  201. ):
  202. raise DeviceQualityForbidden(
  203. "only the accountable device quality asset manager may publish"
  204. )
  205. return DeviceQualityService(
  206. SqlAlchemyDeviceQualityRepository(db.session),
  207. publish_authorizer=assert_accountable,
  208. commit=db.session.commit,
  209. rollback=db.session.rollback,
  210. )
  211. def get_quality_issue_service():
  212. from app.core.data_research.errors import QualityIssueForbidden
  213. from app.core.data_research.quality_issue_repository import (
  214. SqlAlchemyQualityIssueRepository,
  215. )
  216. from app.core.data_research.quality_issues import QualityIssueService
  217. from app.core.governance.responsibilities import (
  218. ResponsibilityService,
  219. SqlAlchemyResponsibilityRepository,
  220. )
  221. repository = SqlAlchemyQualityIssueRepository(db.session)
  222. responsibilities = ResponsibilityService(
  223. SqlAlchemyResponsibilityRepository(db.session)
  224. )
  225. def assert_accountable(actor_uid):
  226. matrix = responsibilities.get(
  227. "quality_issue",
  228. "DEVICE_QUALITY_ISSUES",
  229. )
  230. accountable = [
  231. item
  232. for item in matrix.get("assignments", [])
  233. if item.get("responsibility_role") == "asset_manager"
  234. and item.get("raci_role") == "accountable"
  235. ]
  236. if len(accountable) != 1 or str(accountable[0].get("user_id")) != str(
  237. actor_uid
  238. ):
  239. raise QualityIssueForbidden(
  240. "only the accountable quality issue asset manager may review"
  241. )
  242. return QualityIssueService(
  243. repository,
  244. review_authorizer=assert_accountable,
  245. is_admin=lambda actor_uid: repository.user_has_role(
  246. actor_uid,
  247. "admin",
  248. ),
  249. commit=db.session.commit,
  250. rollback=db.session.rollback,
  251. )
  252. def get_candidate_decision_service():
  253. from app.core.data_research.candidate_decisions import CandidateDecisionService
  254. from app.core.data_research.data_elements import DataElementService
  255. from app.core.data_research.repository import (
  256. SqlAlchemyCandidateDecisionRepository,
  257. SqlAlchemyDataElementRepository,
  258. )
  259. return CandidateDecisionService(
  260. SqlAlchemyCandidateDecisionRepository(db.session),
  261. data_elements=DataElementService(SqlAlchemyDataElementRepository(db.session)),
  262. )
  263. def _artifact_storage():
  264. from minio import Minio
  265. from app.core.data_research.artifacts import MinioArtifactStorage
  266. client = Minio(
  267. current_app.config["MINIO_HOST"],
  268. access_key=current_app.config["MINIO_USER"],
  269. secret_key=current_app.config["MINIO_PASSWORD"],
  270. secure=bool(current_app.config.get("MINIO_SECURE")),
  271. )
  272. return MinioArtifactStorage(client, current_app.config["MINIO_BUCKET"])
  273. def get_artifact_service():
  274. from app.core.common.identifiers import new_governance_uid
  275. from app.core.data_research.artifacts import (
  276. ArtifactService,
  277. SqlAlchemyArtifactRepository,
  278. )
  279. from app.core.data_research.file_policy import FilePolicy
  280. return ArtifactService(
  281. SqlAlchemyArtifactRepository(db.session),
  282. _artifact_storage(),
  283. FilePolicy(
  284. max_bytes=int(current_app.config.get("DATA_RESEARCH_FILE_MAX_BYTES", 25 * 1024 * 1024)),
  285. max_pages=int(current_app.config.get("DATA_RESEARCH_FILE_MAX_PAGES", 500)),
  286. ),
  287. uid_factory=new_governance_uid,
  288. )
  289. def get_evidence_service():
  290. from app.core.data_research.artifacts import EvidenceService
  291. return EvidenceService(db.session, _artifact_storage())
  292. def get_ontology_service():
  293. from app.core.data_research.ontology.publication import (
  294. OntologyApplicationService,
  295. OntologyPublicationService,
  296. )
  297. from app.core.data_research.ontology.repository import SqlAlchemyOntologyRepository
  298. from app.core.events.outbox import enqueue_outbox
  299. repository = SqlAlchemyOntologyRepository(db.session)
  300. authorizer = _device_ontology_authorizer(repository)
  301. publication = OntologyPublicationService(
  302. repository,
  303. outbox_enqueue=lambda **event: enqueue_outbox(db.session, **event),
  304. publication_authorizer=authorizer,
  305. commit=db.session.commit,
  306. rollback=db.session.rollback,
  307. )
  308. return OntologyApplicationService(
  309. repository,
  310. publication=publication,
  311. commit=db.session.commit,
  312. rollback=db.session.rollback,
  313. )
  314. def get_ontology_dynamic_service():
  315. from app.core.data_research.ontology.change_sets import (
  316. SqlAlchemyDynamicOntologyService,
  317. )
  318. return SqlAlchemyDynamicOntologyService(db.session)
  319. def get_ontology_exchange_service():
  320. from app.core.data_research.ontology.exchange import OntologyExchangeService
  321. from app.core.data_research.ontology.repository import SqlAlchemyOntologyRepository
  322. return OntologyExchangeService(
  323. SqlAlchemyOntologyRepository(db.session),
  324. commit=db.session.commit,
  325. rollback=db.session.rollback,
  326. )
  327. def get_semantic_query_service():
  328. from app.core.data_research.ontology.query import (
  329. Neo4jSemanticRepository,
  330. SemanticQueryService,
  331. )
  332. from app.services.neo4j_driver import neo4j_driver
  333. return SemanticQueryService(Neo4jSemanticRepository(neo4j_driver))
  334. def _identity():
  335. return getattr(g, "current_user", {}) or {}
  336. def _record(record):
  337. return {
  338. "uid": str(record.uid),
  339. "source_uid": str(record.source_uid),
  340. "artifact_uid": str(record.artifact_uid) if record.artifact_uid else None,
  341. "job_type": record.job_type,
  342. "parser_version": record.parser_version,
  343. "status": record.status,
  344. "parameters": dict(record.parameters or {}),
  345. "statistics": dict(record.statistics or {}),
  346. "last_error": record.last_error,
  347. "attempt_count": int(record.attempt_count or 0),
  348. "failure_stage": record.failure_stage,
  349. "actor_uid": record.actor_uid,
  350. "created_at": record.created_at.isoformat() if record.created_at else None,
  351. "updated_at": record.updated_at.isoformat() if record.updated_at else None,
  352. "started_at": record.started_at.isoformat() if record.started_at else None,
  353. "finished_at": record.finished_at.isoformat() if record.finished_at else None,
  354. }
  355. def _catalog_snapshot(record):
  356. return {
  357. "uid": str(record.uid),
  358. "job_uid": str(record.job_uid),
  359. "source_uid": str(record.source_uid),
  360. "attempt": int(record.attempt),
  361. "database_type": record.database_type,
  362. "content_hash": record.content_hash,
  363. "snapshot": dict(record.snapshot or {}),
  364. "evidence_count": int(record.evidence_count),
  365. "created_at": (
  366. record.created_at.isoformat()
  367. if record.created_at
  368. else None
  369. ),
  370. }
  371. def _element(record):
  372. return {
  373. "uid": str(record.uid),
  374. "code": record.code,
  375. "status": record.status,
  376. "current_version": int(record.current_version),
  377. "snapshot": dict(record.snapshot or {}),
  378. "created_by": record.created_by,
  379. "updated_by": record.updated_by,
  380. }
  381. def _decision(record):
  382. return {
  383. "uid": str(record.uid),
  384. "candidate_uid": str(record.candidate_uid),
  385. "action": record.action,
  386. "data_element_uid": (
  387. str(record.data_element_uid) if record.data_element_uid else None
  388. ),
  389. "evidence_uids": list(record.evidence_uids),
  390. "actor_uid": record.actor_uid,
  391. "reason": record.reason,
  392. }
  393. def _artifact(record):
  394. return {
  395. "uid": str(record.uid),
  396. "source_uid": str(record.source_uid),
  397. "filename": record.filename,
  398. "media_type": record.media_type,
  399. "size_bytes": int(record.size_bytes),
  400. "content_hash": record.content_hash,
  401. "parser_version": record.parser_version,
  402. }
  403. def _ontology(record):
  404. return {
  405. "uid": str(record.uid),
  406. "code": record.code,
  407. "name": record.name,
  408. "owner_uid": record.owner_uid,
  409. "status": record.status,
  410. "draft_revision": int(record.draft_revision),
  411. "active_version_uid": record.active_version_uid,
  412. "domain_links": [
  413. {"domain_uid": link.domain_uid, "role": link.role}
  414. for link in record.domain_links
  415. ],
  416. }
  417. def _ontology_version(record):
  418. return {
  419. "uid": str(record.uid),
  420. "ontology_uid": str(record.ontology_uid),
  421. "version": int(record.version),
  422. "parent_version_uid": record.parent_version_uid,
  423. "status": record.status,
  424. "content_hash": record.content_hash,
  425. "graph_document": record.graph_document.to_dict(),
  426. "created_by": record.created_by,
  427. }
  428. def _iso(value):
  429. return value.isoformat() if value else None
  430. def _device_asset(record):
  431. return {
  432. "uid": str(record.uid),
  433. "asset_type": record.asset_type,
  434. "name": record.name,
  435. "status": record.status,
  436. "current_version": int(record.current_version),
  437. "location": record.location,
  438. "organization": record.organization,
  439. "responsible_person": record.responsible_person,
  440. "attributes": dict(record.attributes or {}),
  441. "created_by": record.created_by,
  442. "updated_by": record.updated_by,
  443. "created_at": _iso(record.created_at),
  444. "updated_at": _iso(record.updated_at),
  445. }
  446. def _device_asset_mapping(record):
  447. return {
  448. "uid": str(record.uid),
  449. "asset_uid": str(record.asset_uid),
  450. "source_uid": str(record.source_uid),
  451. "source_entity": record.source_entity,
  452. "asset_type": record.asset_type,
  453. "source_code": record.source_code,
  454. "source_updated_at": _iso(record.source_updated_at),
  455. "first_seen_at": _iso(record.first_seen_at),
  456. "last_seen_at": _iso(record.last_seen_at),
  457. }
  458. def _device_asset_detail(detail):
  459. data = _device_asset(detail.asset)
  460. data["source_mappings"] = [
  461. _device_asset_mapping(mapping)
  462. for mapping in detail.mappings
  463. ]
  464. return data
  465. def _device_asset_version(record):
  466. return {
  467. "uid": str(record.uid),
  468. "asset_uid": str(record.asset_uid),
  469. "version": int(record.version),
  470. "snapshot": dict(record.snapshot or {}),
  471. "source_mapping_uid": str(record.source_mapping_uid),
  472. "actor_uid": record.actor_uid,
  473. "created_at": _iso(record.created_at),
  474. }
  475. def _device_asset_import_result(result):
  476. return {
  477. "records": [
  478. {
  479. "action": item.action,
  480. "asset": _device_asset(item.asset),
  481. "source_mapping": _device_asset_mapping(item.mapping),
  482. }
  483. for item in result.items
  484. ],
  485. "created_count": int(result.created_count),
  486. "updated_count": int(result.updated_count),
  487. "unchanged_count": int(result.unchanged_count),
  488. }
  489. def _device_semantic_profile(result):
  490. if result is None:
  491. return {
  492. "bootstrapped": False,
  493. "ontology": None,
  494. "version": None,
  495. "profile": {
  496. "ready_to_publish": False,
  497. "class_count": 0,
  498. "relation_count": 0,
  499. "mapping_count": 0,
  500. "missing_classes": [],
  501. "missing_relations": [],
  502. "missing_mappings": [],
  503. },
  504. }
  505. profile = result.profile
  506. return {
  507. "bootstrapped": True,
  508. "created": bool(result.created),
  509. "ontology": _ontology(result.ontology),
  510. "version": _ontology_version(result.version),
  511. "profile": {
  512. "ready_to_publish": bool(profile.ready_to_publish),
  513. "class_count": int(profile.class_count),
  514. "relation_count": int(profile.relation_count),
  515. "mapping_count": int(profile.mapping_count),
  516. "missing_classes": list(profile.missing_classes),
  517. "missing_relations": list(profile.missing_relations),
  518. "missing_mappings": list(profile.missing_mappings),
  519. },
  520. }
  521. def _device_semantic_code(record):
  522. return {
  523. "uid": str(record.uid),
  524. "ontology_uid": str(record.ontology_uid),
  525. "code_type": record.code_type,
  526. "canonical_code": record.canonical_code,
  527. "canonical_name": record.canonical_name,
  528. "definition": record.definition,
  529. "status": record.status,
  530. "current_version": int(record.current_version),
  531. "source_mappings": [
  532. dict(item) for item in record.source_mappings
  533. ],
  534. "evidence_uids": list(record.evidence_uids),
  535. "suggestion_source": record.suggestion_source,
  536. "confidence": record.confidence,
  537. "created_by": record.created_by,
  538. "updated_by": record.updated_by,
  539. "created_at": _iso(record.created_at),
  540. "updated_at": _iso(record.updated_at),
  541. }
  542. def _device_semantic_version(record):
  543. return {
  544. "uid": str(record.uid),
  545. "code_uid": str(record.code_uid),
  546. "version": int(record.version),
  547. "snapshot": dict(record.snapshot),
  548. "created_by": record.created_by,
  549. "created_at": _iso(record.created_at),
  550. }
  551. def _device_semantic_review(record):
  552. return {
  553. "uid": str(record.uid),
  554. "code_uid": str(record.code_uid),
  555. "version": int(record.version),
  556. "decision": record.decision,
  557. "reason": record.reason,
  558. "actor_uid": record.actor_uid,
  559. "created_at": _iso(record.created_at),
  560. }
  561. def _device_entity_candidate(record):
  562. return {
  563. "uid": str(record.uid),
  564. "left_asset_uid": str(record.left_asset_uid),
  565. "right_asset_uid": str(record.right_asset_uid),
  566. "canonical_asset_uid": (
  567. str(record.canonical_asset_uid)
  568. if record.canonical_asset_uid
  569. else None
  570. ),
  571. "status": record.status,
  572. "suggestion_source": record.suggestion_source,
  573. "confidence": float(record.confidence),
  574. "explanation": [
  575. dict(item) for item in record.explanation
  576. ],
  577. "evidence_uids": list(record.evidence_uids),
  578. "model_provider": record.model_provider,
  579. "model_name": record.model_name,
  580. "current_version": int(record.current_version),
  581. "created_by": record.created_by,
  582. "reviewed_by": record.reviewed_by,
  583. "created_at": _iso(record.created_at),
  584. "updated_at": _iso(record.updated_at),
  585. }
  586. def _device_entity_review(record):
  587. return {
  588. "uid": str(record.uid),
  589. "candidate_uid": str(record.candidate_uid),
  590. "version": int(record.version),
  591. "decision": record.decision,
  592. "reason": record.reason,
  593. "actor_uid": record.actor_uid,
  594. "created_at": _iso(record.created_at),
  595. }
  596. def _device_entity_merge(record):
  597. return {
  598. "uid": str(record.uid),
  599. "candidate_uid": str(record.candidate_uid),
  600. "canonical_asset_uid": str(record.canonical_asset_uid),
  601. "member_asset_uid": str(record.member_asset_uid),
  602. "review_uid": str(record.review_uid),
  603. "snapshot": dict(record.snapshot),
  604. "actor_uid": record.actor_uid,
  605. "created_at": _iso(record.created_at),
  606. }
  607. def _device_entity_rollback(record):
  608. return {
  609. "uid": str(record.uid),
  610. "merge_uid": str(record.merge_uid),
  611. "candidate_uid": str(record.candidate_uid),
  612. "reason": record.reason,
  613. "snapshot": dict(record.snapshot),
  614. "actor_uid": record.actor_uid,
  615. "created_at": _iso(record.created_at),
  616. }
  617. def _device_entity_generation(result):
  618. return {
  619. "records": [
  620. _device_entity_candidate(record)
  621. for record in result.records
  622. ],
  623. "created_count": int(result.created_count),
  624. "existing_count": int(result.existing_count),
  625. "evaluated_pair_count": int(result.evaluated_pair_count),
  626. "auto_merged_count": int(result.auto_merged_count),
  627. "auto_merge_enabled": bool(
  628. current_app.config.get(
  629. "DEVICE_ENTITY_AUTO_MERGE_ENABLED",
  630. False,
  631. )
  632. ),
  633. }
  634. def _device_quality_version(record):
  635. if record is None:
  636. return None
  637. return {
  638. "uid": str(record.uid),
  639. "profile_uid": str(record.profile_uid),
  640. "version": int(record.version),
  641. "status": record.status,
  642. "rules": [dict(item) for item in record.rules],
  643. "content_hash": record.content_hash,
  644. "created_by": record.created_by,
  645. "created_at": _iso(record.created_at),
  646. "published_by": record.published_by,
  647. "published_at": _iso(record.published_at),
  648. }
  649. def _device_quality_run(record):
  650. return {
  651. "uid": str(record.uid),
  652. "policy_version_uid": str(record.policy_version_uid),
  653. "policy_hash": record.policy_hash,
  654. "source_uid": (
  655. str(record.source_uid) if record.source_uid else None
  656. ),
  657. "status": record.status,
  658. "total_assets": int(record.total_assets),
  659. "total_violations": int(record.total_violations),
  660. "score": float(record.score),
  661. "created_by": record.created_by,
  662. "created_at": _iso(record.created_at),
  663. }
  664. def _device_quality_rule_result(record):
  665. return {
  666. "uid": str(record.uid),
  667. "run_uid": str(record.run_uid),
  668. "rule_code": record.rule_code,
  669. "severity": record.severity,
  670. "weight": float(record.weight),
  671. "status": record.status,
  672. "evaluated_count": int(record.evaluated_count),
  673. "violation_count": int(record.violation_count),
  674. "sampled_count": int(record.sampled_count),
  675. "pass_rate": float(record.pass_rate),
  676. "weighted_score": float(record.weighted_score),
  677. "created_at": _iso(record.created_at),
  678. }
  679. def _device_quality_violation(record):
  680. return {
  681. "uid": str(record.uid),
  682. "run_uid": str(record.run_uid),
  683. "rule_code": record.rule_code,
  684. "severity": record.severity,
  685. "asset_uid": str(record.asset_uid),
  686. "field_name": record.field_name,
  687. "source_uid": (
  688. str(record.source_uid) if record.source_uid else None
  689. ),
  690. "source_mapping_uid": (
  691. str(record.source_mapping_uid)
  692. if record.source_mapping_uid
  693. else None
  694. ),
  695. "message": record.message,
  696. "evidence": dict(record.evidence or {}),
  697. "created_at": _iso(record.created_at),
  698. "expires_at": _iso(record.expires_at),
  699. }
  700. def _device_quality_asset_score(record):
  701. return {
  702. "uid": str(record.uid),
  703. "run_uid": str(record.run_uid),
  704. "asset_uid": str(record.asset_uid),
  705. "asset_type": record.asset_type,
  706. "status": record.status,
  707. "evaluated_rule_count": int(record.evaluated_rule_count),
  708. "violation_count": int(record.violation_count),
  709. "score": float(record.score),
  710. "created_at": _iso(record.created_at),
  711. }
  712. def _quality_issue(record):
  713. return {
  714. "uid": str(record.uid),
  715. "issue_code": record.issue_code,
  716. "source_violation_uid": str(record.source_violation_uid),
  717. "source_run_uid": str(record.source_run_uid),
  718. "rule_code": record.rule_code,
  719. "severity": record.severity,
  720. "priority": record.priority,
  721. "asset_uid": str(record.asset_uid),
  722. "field_name": record.field_name,
  723. "source_uid": (
  724. str(record.source_uid) if record.source_uid else None
  725. ),
  726. "source_mapping_uid": (
  727. str(record.source_mapping_uid)
  728. if record.source_mapping_uid else None
  729. ),
  730. "message": record.message,
  731. "evidence": dict(record.evidence or {}),
  732. "recurrence_key": record.recurrence_key,
  733. "occurrence_number": int(record.occurrence_number),
  734. "status": record.status,
  735. "assignee_uid": (
  736. str(record.assignee_uid) if record.assignee_uid else None
  737. ),
  738. "due_at": _iso(record.due_at),
  739. "is_overdue": bool(record.is_overdue(now_china())),
  740. "current_version": int(record.current_version),
  741. "created_by": str(record.created_by),
  742. "updated_by": str(record.updated_by),
  743. "created_at": _iso(record.created_at),
  744. "updated_at": _iso(record.updated_at),
  745. "closed_at": _iso(record.closed_at),
  746. }
  747. def _quality_issue_remediation(record):
  748. if record is None:
  749. return None
  750. return {
  751. "uid": str(record.uid),
  752. "issue_uid": str(record.issue_uid),
  753. "round_number": int(record.round_number),
  754. "summary": record.summary,
  755. "evidence_refs": [dict(item) for item in record.evidence_refs],
  756. "submitted_by": str(record.submitted_by),
  757. "submitted_at": _iso(record.submitted_at),
  758. "review_status": record.review_status,
  759. "reviewed_by": (
  760. str(record.reviewed_by) if record.reviewed_by else None
  761. ),
  762. "reviewed_at": _iso(record.reviewed_at),
  763. "review_note": record.review_note,
  764. "verification_run_uid": (
  765. str(record.verification_run_uid)
  766. if record.verification_run_uid else None
  767. ),
  768. }
  769. def _quality_issue_timeline(record):
  770. return {
  771. "uid": str(record.uid),
  772. "issue_uid": str(record.issue_uid),
  773. "action": record.action,
  774. "from_status": record.from_status,
  775. "to_status": record.to_status,
  776. "actor_uid": str(record.actor_uid),
  777. "note": record.note,
  778. "payload": dict(record.payload or {}),
  779. "created_at": _iso(record.created_at),
  780. }
  781. def _device_asset_page(name, *, default, maximum):
  782. from app.core.data_research.errors import DeviceAssetInvalid
  783. raw = request.args.get(name)
  784. try:
  785. value = default if raw in (None, "") else int(raw)
  786. except (TypeError, ValueError) as error:
  787. raise DeviceAssetInvalid(f"{name} must be an integer") from error
  788. if value < 1 or value > maximum:
  789. raise DeviceAssetInvalid(
  790. f"{name} must be between 1 and {maximum}"
  791. )
  792. return value
  793. def _error(error):
  794. if isinstance(error, DataResearchError):
  795. return (
  796. jsonify(
  797. failed(
  798. str(error),
  799. code=error.http_status,
  800. error={"code": error.code},
  801. )
  802. ),
  803. error.http_status,
  804. )
  805. logger.exception("data-research ingestion request failed")
  806. return (
  807. jsonify(
  808. failed(
  809. "数据采集任务处理失败",
  810. code=500,
  811. error={"code": "DATA_RESEARCH_ERROR"},
  812. )
  813. ),
  814. 500,
  815. )
  816. @bp.route("/ingestion-jobs", methods=["POST"])
  817. def create_ingestion_job():
  818. payload = request.get_json(silent=True) or {}
  819. try:
  820. actor_uid = _identity().get("id") or _identity().get("sub")
  821. if payload.get("job_type") == "catalog_collect":
  822. get_database_source_registration_service().ensure(
  823. payload.get("source_uid"),
  824. actor_uid=actor_uid,
  825. )
  826. record, created = get_ingestion_service().create_job(
  827. payload,
  828. actor_uid=actor_uid,
  829. )
  830. return jsonify(success(_record(record))), 201 if created else 200
  831. except Exception as error:
  832. return _error(error)
  833. @bp.route("/ingestion-jobs", methods=["GET"])
  834. def list_ingestion_jobs():
  835. filters = {
  836. name: request.args.get(name)
  837. for name in ("status", "source_uid")
  838. if request.args.get(name)
  839. }
  840. try:
  841. records = get_ingestion_service().list_jobs(filters)
  842. return jsonify(
  843. success({"records": [_record(item) for item in records], "total": len(records)})
  844. ), 200
  845. except Exception as error:
  846. return _error(error)
  847. @bp.route("/ingestion-jobs/<job_uid>", methods=["GET"])
  848. def get_ingestion_job(job_uid):
  849. try:
  850. return jsonify(success(_record(get_ingestion_service().get_job(job_uid)))), 200
  851. except Exception as error:
  852. return _error(error)
  853. @bp.route("/ingestion-jobs/<job_uid>/execute", methods=["POST"])
  854. def execute_ingestion_job(job_uid):
  855. try:
  856. record = get_catalog_ingestion_executor().execute(job_uid)
  857. return jsonify(success(_record(record))), 200
  858. except Exception as error:
  859. return _error(error)
  860. @bp.route(
  861. "/ingestion-jobs/<job_uid>/catalog-snapshots",
  862. methods=["GET"],
  863. )
  864. def list_catalog_snapshots(job_uid):
  865. try:
  866. records = get_catalog_snapshot_repository().list(job_uid)
  867. return jsonify(
  868. success(
  869. {
  870. "records": [
  871. _catalog_snapshot(record)
  872. for record in records
  873. ],
  874. "total": len(records),
  875. }
  876. )
  877. ), 200
  878. except Exception as error:
  879. return _error(error)
  880. @bp.route("/ingestion-jobs/<job_uid>/evidence", methods=["GET"])
  881. def list_ingestion_job_evidence(job_uid):
  882. try:
  883. records = get_evidence_service().list_for_job(job_uid)
  884. return jsonify(
  885. success({"records": records, "total": len(records)})
  886. ), 200
  887. except Exception as error:
  888. return _error(error)
  889. @bp.route("/ingestion-jobs/<job_uid>/retry", methods=["POST"])
  890. def retry_ingestion_job(job_uid):
  891. try:
  892. return jsonify(success(_record(get_ingestion_service().retry(job_uid)))), 200
  893. except Exception as error:
  894. return _error(error)
  895. @bp.route("/ingestion-jobs/<job_uid>/cancel", methods=["POST"])
  896. def cancel_ingestion_job(job_uid):
  897. try:
  898. service = get_ingestion_service()
  899. record = service.get_job(job_uid)
  900. identity = _identity()
  901. permissions = set(identity.get("permissions") or [])
  902. actor_uid = identity.get("id") or identity.get("sub")
  903. if record.actor_uid != actor_uid and "ingestion:admin" not in permissions:
  904. return jsonify(failed("权限不足", code=403)), 403
  905. return jsonify(success(_record(service.cancel(job_uid)))), 200
  906. except Exception as error:
  907. return _error(error)
  908. @bp.route("/data-elements", methods=["GET"])
  909. def list_data_elements():
  910. try:
  911. records = get_data_element_service().list(
  912. status=request.args.get("status"),
  913. business_domain_uid=request.args.get("business_domain_uid"),
  914. )
  915. return jsonify(
  916. success(
  917. {
  918. "records": [_element(item) for item in records],
  919. "total": len(records),
  920. }
  921. )
  922. ), 200
  923. except Exception as error:
  924. return _error(error)
  925. @bp.route("/device-assets", methods=["GET"])
  926. def list_device_assets():
  927. filters = {
  928. name: request.args.get(name)
  929. for name in ("keyword", "asset_type", "status", "source_uid")
  930. if request.args.get(name)
  931. }
  932. try:
  933. page = _device_asset_page(
  934. "page",
  935. default=1,
  936. maximum=1_000_000,
  937. )
  938. page_size = _device_asset_page(
  939. "page_size",
  940. default=20,
  941. maximum=100,
  942. )
  943. records, total = get_device_asset_service().search(
  944. filters,
  945. page=page,
  946. page_size=page_size,
  947. )
  948. return jsonify(
  949. success(
  950. {
  951. "records": [
  952. _device_asset_detail(record)
  953. for record in records
  954. ],
  955. "total": int(total),
  956. "page": page,
  957. "page_size": page_size,
  958. }
  959. )
  960. ), 200
  961. except Exception as error:
  962. return _error(error)
  963. @bp.route("/device-assets/import", methods=["POST"])
  964. def import_device_assets():
  965. try:
  966. result = get_device_asset_service().import_records(
  967. request.get_json(silent=True) or {},
  968. actor_uid=_identity().get("id") or _identity().get("sub"),
  969. )
  970. return jsonify(success(_device_asset_import_result(result))), 200
  971. except Exception as error:
  972. return _error(error)
  973. @bp.route("/device-assets/<asset_uid>", methods=["GET"])
  974. def get_device_asset(asset_uid):
  975. try:
  976. return jsonify(
  977. success(
  978. _device_asset_detail(
  979. get_device_asset_service().get(asset_uid)
  980. )
  981. )
  982. ), 200
  983. except Exception as error:
  984. return _error(error)
  985. @bp.route("/device-assets/<asset_uid>/versions", methods=["GET"])
  986. def list_device_asset_versions(asset_uid):
  987. try:
  988. records = get_device_asset_service().versions(asset_uid)
  989. return jsonify(
  990. success(
  991. {
  992. "records": [
  993. _device_asset_version(record)
  994. for record in records
  995. ],
  996. "total": len(records),
  997. }
  998. )
  999. ), 200
  1000. except Exception as error:
  1001. return _error(error)
  1002. @bp.route("/device-semantics/bootstrap", methods=["POST"])
  1003. def bootstrap_device_semantics():
  1004. try:
  1005. result = get_device_semantic_service().bootstrap(
  1006. request.get_json(silent=True) or {},
  1007. actor_uid=_identity().get("id") or _identity().get("sub"),
  1008. )
  1009. return (
  1010. jsonify(success(_device_semantic_profile(result))),
  1011. 201 if result.created else 200,
  1012. )
  1013. except Exception as error:
  1014. return _error(error)
  1015. @bp.route("/device-semantics/profile", methods=["GET"])
  1016. def get_device_semantic_profile():
  1017. try:
  1018. result = get_device_semantic_service().profile()
  1019. return jsonify(success(_device_semantic_profile(result))), 200
  1020. except Exception as error:
  1021. return _error(error)
  1022. @bp.route("/device-semantics/codes", methods=["GET"])
  1023. def list_device_semantic_codes():
  1024. filters = {
  1025. name: request.args.get(name)
  1026. for name in ("ontology_uid", "code_type", "status", "keyword")
  1027. if request.args.get(name)
  1028. }
  1029. try:
  1030. records, total = get_device_semantic_code_service().search(
  1031. filters,
  1032. page=request.args.get("page", 1),
  1033. page_size=request.args.get("page_size", 20),
  1034. )
  1035. return jsonify(
  1036. success(
  1037. {
  1038. "records": [
  1039. _device_semantic_code(record)
  1040. for record in records
  1041. ],
  1042. "total": int(total),
  1043. "page": int(request.args.get("page", 1)),
  1044. "page_size": int(request.args.get("page_size", 20)),
  1045. }
  1046. )
  1047. ), 200
  1048. except Exception as error:
  1049. return _error(error)
  1050. @bp.route("/device-semantics/codes", methods=["POST"])
  1051. def create_device_semantic_code():
  1052. try:
  1053. record = get_device_semantic_code_service().create(
  1054. request.get_json(silent=True) or {},
  1055. actor_uid=_identity().get("id") or _identity().get("sub"),
  1056. )
  1057. return jsonify(success(_device_semantic_code(record))), 201
  1058. except Exception as error:
  1059. return _error(error)
  1060. @bp.route("/device-semantics/codes/<code_uid>", methods=["GET"])
  1061. def get_device_semantic_code(code_uid):
  1062. try:
  1063. record = get_device_semantic_code_service().get(code_uid)
  1064. return jsonify(success(_device_semantic_code(record))), 200
  1065. except Exception as error:
  1066. return _error(error)
  1067. @bp.route(
  1068. "/device-semantics/codes/<code_uid>/revisions",
  1069. methods=["POST"],
  1070. )
  1071. def revise_device_semantic_code(code_uid):
  1072. payload = request.get_json(silent=True) or {}
  1073. try:
  1074. record = get_device_semantic_code_service().revise(
  1075. code_uid,
  1076. payload,
  1077. expected_version=payload.get("expected_version"),
  1078. actor_uid=_identity().get("id") or _identity().get("sub"),
  1079. )
  1080. return jsonify(success(_device_semantic_code(record))), 200
  1081. except Exception as error:
  1082. return _error(error)
  1083. @bp.route(
  1084. "/device-semantics/codes/<code_uid>/submit",
  1085. methods=["POST"],
  1086. )
  1087. def submit_device_semantic_code(code_uid):
  1088. payload = request.get_json(silent=True) or {}
  1089. try:
  1090. record = get_device_semantic_code_service().submit(
  1091. code_uid,
  1092. expected_version=payload.get("expected_version"),
  1093. actor_uid=_identity().get("id") or _identity().get("sub"),
  1094. )
  1095. return jsonify(success(_device_semantic_code(record))), 200
  1096. except Exception as error:
  1097. return _error(error)
  1098. @bp.route(
  1099. "/device-semantics/codes/<code_uid>/review",
  1100. methods=["POST"],
  1101. )
  1102. def review_device_semantic_code(code_uid):
  1103. try:
  1104. record, review = get_device_semantic_code_service().review(
  1105. code_uid,
  1106. request.get_json(silent=True) or {},
  1107. actor_uid=_identity().get("id") or _identity().get("sub"),
  1108. )
  1109. return jsonify(
  1110. success(
  1111. {
  1112. "record": _device_semantic_code(record),
  1113. "review": _device_semantic_review(review),
  1114. }
  1115. )
  1116. ), 200
  1117. except Exception as error:
  1118. return _error(error)
  1119. @bp.route(
  1120. "/device-semantics/codes/<code_uid>/versions",
  1121. methods=["GET"],
  1122. )
  1123. def list_device_semantic_code_versions(code_uid):
  1124. try:
  1125. records = get_device_semantic_code_service().versions(code_uid)
  1126. return jsonify(
  1127. success(
  1128. {
  1129. "records": [
  1130. _device_semantic_version(record)
  1131. for record in records
  1132. ],
  1133. "total": len(records),
  1134. }
  1135. )
  1136. ), 200
  1137. except Exception as error:
  1138. return _error(error)
  1139. @bp.route(
  1140. "/device-semantics/codes/<code_uid>/reviews",
  1141. methods=["GET"],
  1142. )
  1143. def list_device_semantic_code_reviews(code_uid):
  1144. try:
  1145. records = get_device_semantic_code_service().reviews(code_uid)
  1146. return jsonify(
  1147. success(
  1148. {
  1149. "records": [
  1150. _device_semantic_review(record)
  1151. for record in records
  1152. ],
  1153. "total": len(records),
  1154. }
  1155. )
  1156. ), 200
  1157. except Exception as error:
  1158. return _error(error)
  1159. @bp.route("/device-entities/candidates", methods=["GET"])
  1160. def list_device_entity_candidates():
  1161. filters = {
  1162. name: request.args.get(name)
  1163. for name in ("status", "suggestion_source")
  1164. if request.args.get(name)
  1165. }
  1166. try:
  1167. records, total = get_device_entity_resolution_service().search(
  1168. filters,
  1169. page=request.args.get("page", 1),
  1170. page_size=request.args.get("page_size", 20),
  1171. )
  1172. return jsonify(
  1173. success(
  1174. {
  1175. "records": [
  1176. _device_entity_candidate(record)
  1177. for record in records
  1178. ],
  1179. "total": int(total),
  1180. "page": int(request.args.get("page", 1)),
  1181. "page_size": int(request.args.get("page_size", 20)),
  1182. "auto_merge_enabled": bool(
  1183. current_app.config.get(
  1184. "DEVICE_ENTITY_AUTO_MERGE_ENABLED",
  1185. False,
  1186. )
  1187. ),
  1188. }
  1189. )
  1190. ), 200
  1191. except Exception as error:
  1192. return _error(error)
  1193. @bp.route("/device-entities/candidates/generate", methods=["POST"])
  1194. def generate_device_entity_candidates():
  1195. try:
  1196. result = get_device_entity_resolution_service().generate(
  1197. request.get_json(silent=True) or {},
  1198. actor_uid=_identity().get("id") or _identity().get("sub"),
  1199. )
  1200. status = 201 if result.created_count else 200
  1201. return jsonify(success(_device_entity_generation(result))), status
  1202. except Exception as error:
  1203. return _error(error)
  1204. @bp.route("/device-entities/candidates", methods=["POST"])
  1205. def submit_device_entity_candidate():
  1206. try:
  1207. record = (
  1208. get_device_entity_resolution_service().submit_ai_candidate(
  1209. request.get_json(silent=True) or {},
  1210. actor_uid=_identity().get("id") or _identity().get("sub"),
  1211. )
  1212. )
  1213. return jsonify(success(_device_entity_candidate(record))), 201
  1214. except Exception as error:
  1215. return _error(error)
  1216. @bp.route("/device-entities/candidates/<candidate_uid>", methods=["GET"])
  1217. def get_device_entity_candidate(candidate_uid):
  1218. try:
  1219. record = get_device_entity_resolution_service().get(candidate_uid)
  1220. return jsonify(success(_device_entity_candidate(record))), 200
  1221. except Exception as error:
  1222. return _error(error)
  1223. @bp.route(
  1224. "/device-entities/candidates/<candidate_uid>/review",
  1225. methods=["POST"],
  1226. )
  1227. def review_device_entity_candidate(candidate_uid):
  1228. try:
  1229. candidate, review, merge = (
  1230. get_device_entity_resolution_service().review(
  1231. candidate_uid,
  1232. request.get_json(silent=True) or {},
  1233. actor_uid=_identity().get("id") or _identity().get("sub"),
  1234. )
  1235. )
  1236. return jsonify(
  1237. success(
  1238. {
  1239. "candidate": _device_entity_candidate(candidate),
  1240. "review": _device_entity_review(review),
  1241. "merge": (
  1242. _device_entity_merge(merge)
  1243. if merge is not None
  1244. else None
  1245. ),
  1246. }
  1247. )
  1248. ), 200
  1249. except Exception as error:
  1250. return _error(error)
  1251. @bp.route(
  1252. "/device-entities/candidates/<candidate_uid>/reviews",
  1253. methods=["GET"],
  1254. )
  1255. def list_device_entity_reviews(candidate_uid):
  1256. try:
  1257. records = get_device_entity_resolution_service().reviews(
  1258. candidate_uid
  1259. )
  1260. return jsonify(
  1261. success(
  1262. {
  1263. "records": [
  1264. _device_entity_review(record)
  1265. for record in records
  1266. ],
  1267. "total": len(records),
  1268. }
  1269. )
  1270. ), 200
  1271. except Exception as error:
  1272. return _error(error)
  1273. @bp.route(
  1274. "/device-entities/candidates/<candidate_uid>/merges",
  1275. methods=["GET"],
  1276. )
  1277. def list_device_entity_merges(candidate_uid):
  1278. try:
  1279. records = get_device_entity_resolution_service().merges(
  1280. candidate_uid
  1281. )
  1282. return jsonify(
  1283. success(
  1284. {
  1285. "records": [
  1286. _device_entity_merge(record)
  1287. for record in records
  1288. ],
  1289. "total": len(records),
  1290. }
  1291. )
  1292. ), 200
  1293. except Exception as error:
  1294. return _error(error)
  1295. @bp.route(
  1296. "/device-entities/merges/<merge_uid>/rollback",
  1297. methods=["POST"],
  1298. )
  1299. def rollback_device_entity_merge(merge_uid):
  1300. try:
  1301. candidate, rollback = (
  1302. get_device_entity_resolution_service().rollback(
  1303. merge_uid,
  1304. request.get_json(silent=True) or {},
  1305. actor_uid=_identity().get("id") or _identity().get("sub"),
  1306. )
  1307. )
  1308. return jsonify(
  1309. success(
  1310. {
  1311. "candidate": _device_entity_candidate(candidate),
  1312. "rollback": _device_entity_rollback(rollback),
  1313. }
  1314. )
  1315. ), 200
  1316. except Exception as error:
  1317. return _error(error)
  1318. @bp.route(
  1319. "/device-entities/merges/<merge_uid>/rollbacks",
  1320. methods=["GET"],
  1321. )
  1322. def list_device_entity_rollbacks(merge_uid):
  1323. try:
  1324. records = get_device_entity_resolution_service().rollbacks(
  1325. merge_uid
  1326. )
  1327. return jsonify(
  1328. success(
  1329. {
  1330. "records": [
  1331. _device_entity_rollback(record)
  1332. for record in records
  1333. ],
  1334. "total": len(records),
  1335. }
  1336. )
  1337. ), 200
  1338. except Exception as error:
  1339. return _error(error)
  1340. @bp.route("/device-quality/profile", methods=["GET"])
  1341. def get_device_quality_profile():
  1342. try:
  1343. result = get_device_quality_service().profile()
  1344. return jsonify(
  1345. success(
  1346. {
  1347. "profile_uid": result["profile_uid"],
  1348. "name": result["name"],
  1349. "latest_version": _device_quality_version(
  1350. result["latest_version"]
  1351. ),
  1352. "active_version": _device_quality_version(
  1353. result["active_version"]
  1354. ),
  1355. }
  1356. )
  1357. ), 200
  1358. except Exception as error:
  1359. return _error(error)
  1360. @bp.route("/device-quality/bootstrap", methods=["POST"])
  1361. def bootstrap_device_quality():
  1362. try:
  1363. record = get_device_quality_service().bootstrap(
  1364. actor_uid=_identity().get("id") or _identity().get("sub"),
  1365. )
  1366. return jsonify(success(_device_quality_version(record))), 201
  1367. except Exception as error:
  1368. return _error(error)
  1369. @bp.route("/device-quality/profile/versions", methods=["POST"])
  1370. def revise_device_quality_profile():
  1371. try:
  1372. body = request.get_json(silent=True) or {}
  1373. record = get_device_quality_service().revise(
  1374. rules=body.get("rules"),
  1375. expected_version=body.get("expected_version"),
  1376. actor_uid=_identity().get("id") or _identity().get("sub"),
  1377. )
  1378. return jsonify(success(_device_quality_version(record))), 201
  1379. except Exception as error:
  1380. return _error(error)
  1381. @bp.route("/device-quality/profile/versions", methods=["GET"])
  1382. def list_device_quality_versions():
  1383. try:
  1384. records = get_device_quality_service().versions()
  1385. return jsonify(
  1386. success(
  1387. {
  1388. "records": [
  1389. _device_quality_version(record)
  1390. for record in records
  1391. ],
  1392. "total": len(records),
  1393. }
  1394. )
  1395. ), 200
  1396. except Exception as error:
  1397. return _error(error)
  1398. @bp.route(
  1399. "/device-quality/profile/versions/<version_uid>/publish",
  1400. methods=["POST"],
  1401. )
  1402. def publish_device_quality_version(version_uid):
  1403. try:
  1404. record = get_device_quality_service().publish(
  1405. version_uid,
  1406. actor_uid=_identity().get("id") or _identity().get("sub"),
  1407. )
  1408. return jsonify(success(_device_quality_version(record))), 200
  1409. except Exception as error:
  1410. return _error(error)
  1411. @bp.route("/device-quality/runs", methods=["POST"])
  1412. def run_device_quality():
  1413. try:
  1414. body = request.get_json(silent=True) or {}
  1415. record = get_device_quality_service().run(
  1416. actor_uid=_identity().get("id") or _identity().get("sub"),
  1417. source_uid=body.get("source_uid"),
  1418. )
  1419. return jsonify(success(_device_quality_run(record))), 201
  1420. except Exception as error:
  1421. return _error(error)
  1422. @bp.route("/device-quality/runs", methods=["GET"])
  1423. def list_device_quality_runs():
  1424. try:
  1425. page = request.args.get("page", 1)
  1426. page_size = request.args.get("page_size", 20)
  1427. records, total = get_device_quality_service().runs(
  1428. page=page,
  1429. page_size=page_size,
  1430. )
  1431. return jsonify(
  1432. success(
  1433. {
  1434. "records": [
  1435. _device_quality_run(record)
  1436. for record in records
  1437. ],
  1438. "total": int(total),
  1439. "page": int(page),
  1440. "page_size": int(page_size),
  1441. }
  1442. )
  1443. ), 200
  1444. except Exception as error:
  1445. return _error(error)
  1446. @bp.route("/device-quality/runs/<run_uid>", methods=["GET"])
  1447. def get_device_quality_run(run_uid):
  1448. try:
  1449. record, results = get_device_quality_service().get_run(run_uid)
  1450. return jsonify(
  1451. success(
  1452. {
  1453. **_device_quality_run(record),
  1454. "rule_results": [
  1455. _device_quality_rule_result(item)
  1456. for item in results
  1457. ],
  1458. }
  1459. )
  1460. ), 200
  1461. except Exception as error:
  1462. return _error(error)
  1463. @bp.route(
  1464. "/device-quality/runs/<run_uid>/violations",
  1465. methods=["GET"],
  1466. )
  1467. def list_device_quality_violations(run_uid):
  1468. try:
  1469. page = request.args.get("page", 1)
  1470. page_size = request.args.get("page_size", 20)
  1471. records, total = get_device_quality_service().violations(
  1472. run_uid,
  1473. rule_code=request.args.get("rule_code"),
  1474. page=page,
  1475. page_size=page_size,
  1476. )
  1477. return jsonify(
  1478. success(
  1479. {
  1480. "records": [
  1481. _device_quality_violation(record)
  1482. for record in records
  1483. ],
  1484. "total": int(total),
  1485. "page": int(page),
  1486. "page_size": int(page_size),
  1487. }
  1488. )
  1489. ), 200
  1490. except Exception as error:
  1491. return _error(error)
  1492. @bp.route(
  1493. "/device-quality/runs/<run_uid>/asset-scores",
  1494. methods=["GET"],
  1495. )
  1496. def list_device_quality_asset_scores(run_uid):
  1497. try:
  1498. page = request.args.get("page", 1)
  1499. page_size = request.args.get("page_size", 20)
  1500. records, total = get_device_quality_service().asset_scores(
  1501. run_uid,
  1502. page=page,
  1503. page_size=page_size,
  1504. )
  1505. return jsonify(
  1506. success(
  1507. {
  1508. "records": [
  1509. _device_quality_asset_score(record)
  1510. for record in records
  1511. ],
  1512. "total": int(total),
  1513. "page": int(page),
  1514. "page_size": int(page_size),
  1515. }
  1516. )
  1517. ), 200
  1518. except Exception as error:
  1519. return _error(error)
  1520. @bp.route("/device-quality/issues", methods=["GET"])
  1521. def list_quality_issues():
  1522. try:
  1523. page = request.args.get("page", 1)
  1524. page_size = request.args.get("page_size", 20)
  1525. overdue_only = str(
  1526. request.args.get("overdue_only", "")
  1527. ).strip().lower() in {"1", "true", "yes"}
  1528. records, total = get_quality_issue_service().issues(
  1529. status=request.args.get("status"),
  1530. assignee_uid=request.args.get("assignee_uid"),
  1531. overdue_only=overdue_only,
  1532. page=page,
  1533. page_size=page_size,
  1534. )
  1535. return jsonify(
  1536. success(
  1537. {
  1538. "records": [_quality_issue(item) for item in records],
  1539. "total": int(total),
  1540. "page": int(page),
  1541. "page_size": int(page_size),
  1542. }
  1543. )
  1544. ), 200
  1545. except Exception as error:
  1546. return _error(error)
  1547. @bp.route("/device-quality/issues/statistics", methods=["GET"])
  1548. def get_quality_issue_statistics():
  1549. try:
  1550. return jsonify(
  1551. success(get_quality_issue_service().statistics())
  1552. ), 200
  1553. except Exception as error:
  1554. return _error(error)
  1555. @bp.route("/device-quality/issues/assignees", methods=["GET"])
  1556. def list_quality_issue_assignees():
  1557. try:
  1558. records = get_quality_issue_service().assignees()
  1559. return jsonify(
  1560. success({"records": list(records), "total": len(records)})
  1561. ), 200
  1562. except Exception as error:
  1563. return _error(error)
  1564. @bp.route("/device-quality/issues/import", methods=["POST"])
  1565. def import_quality_issues():
  1566. try:
  1567. body = request.get_json(silent=True) or {}
  1568. result = get_quality_issue_service().import_violations(
  1569. violation_uids=body.get("violation_uids"),
  1570. priority=body.get("priority"),
  1571. due_at=body.get("due_at"),
  1572. actor_uid=_identity().get("id") or _identity().get("sub"),
  1573. )
  1574. payload = {
  1575. "records": [_quality_issue(item) for item in result.records],
  1576. "created_count": int(result.created_count),
  1577. "existing_count": int(result.existing_count),
  1578. }
  1579. status = 201 if result.created_count else 200
  1580. return jsonify(success(payload)), status
  1581. except Exception as error:
  1582. return _error(error)
  1583. @bp.route("/device-quality/issues/<issue_uid>", methods=["GET"])
  1584. def get_quality_issue(issue_uid):
  1585. try:
  1586. issue, remediation = get_quality_issue_service().get(issue_uid)
  1587. return jsonify(
  1588. success(
  1589. {
  1590. **_quality_issue(issue),
  1591. "latest_remediation": _quality_issue_remediation(
  1592. remediation
  1593. ),
  1594. }
  1595. )
  1596. ), 200
  1597. except Exception as error:
  1598. return _error(error)
  1599. @bp.route(
  1600. "/device-quality/issues/<issue_uid>/timeline",
  1601. methods=["GET"],
  1602. )
  1603. def get_quality_issue_timeline(issue_uid):
  1604. try:
  1605. records = get_quality_issue_service().timeline(issue_uid)
  1606. return jsonify(
  1607. success(
  1608. {
  1609. "records": [
  1610. _quality_issue_timeline(item) for item in records
  1611. ],
  1612. "total": len(records),
  1613. }
  1614. )
  1615. ), 200
  1616. except Exception as error:
  1617. return _error(error)
  1618. @bp.route(
  1619. "/device-quality/issues/<issue_uid>/assign",
  1620. methods=["POST"],
  1621. )
  1622. def assign_quality_issue(issue_uid):
  1623. try:
  1624. body = request.get_json(silent=True) or {}
  1625. record = get_quality_issue_service().assign(
  1626. issue_uid,
  1627. assignee_uid=body.get("assignee_uid"),
  1628. due_at=body.get("due_at"),
  1629. expected_version=body.get("expected_version"),
  1630. actor_uid=_identity().get("id") or _identity().get("sub"),
  1631. )
  1632. return jsonify(success(_quality_issue(record))), 200
  1633. except Exception as error:
  1634. return _error(error)
  1635. @bp.route(
  1636. "/device-quality/issues/<issue_uid>/start",
  1637. methods=["POST"],
  1638. )
  1639. def start_quality_issue(issue_uid):
  1640. try:
  1641. body = request.get_json(silent=True) or {}
  1642. record = get_quality_issue_service().start(
  1643. issue_uid,
  1644. expected_version=body.get("expected_version"),
  1645. actor_uid=_identity().get("id") or _identity().get("sub"),
  1646. )
  1647. return jsonify(success(_quality_issue(record))), 200
  1648. except Exception as error:
  1649. return _error(error)
  1650. @bp.route(
  1651. "/device-quality/issues/<issue_uid>/submit",
  1652. methods=["POST"],
  1653. )
  1654. def submit_quality_issue(issue_uid):
  1655. try:
  1656. body = request.get_json(silent=True) or {}
  1657. record = get_quality_issue_service().submit(
  1658. issue_uid,
  1659. summary=body.get("summary"),
  1660. evidence_refs=body.get("evidence_refs", []),
  1661. expected_version=body.get("expected_version"),
  1662. actor_uid=_identity().get("id") or _identity().get("sub"),
  1663. )
  1664. return jsonify(success(_quality_issue(record))), 200
  1665. except Exception as error:
  1666. return _error(error)
  1667. @bp.route(
  1668. "/device-quality/issues/<issue_uid>/review",
  1669. methods=["POST"],
  1670. )
  1671. def review_quality_issue(issue_uid):
  1672. try:
  1673. body = request.get_json(silent=True) or {}
  1674. record = get_quality_issue_service().review(
  1675. issue_uid,
  1676. verification_result=body.get("verification_result"),
  1677. note=body.get("note"),
  1678. verification_run_uid=body.get("verification_run_uid"),
  1679. expected_version=body.get("expected_version"),
  1680. actor_uid=_identity().get("id") or _identity().get("sub"),
  1681. )
  1682. return jsonify(success(_quality_issue(record))), 200
  1683. except Exception as error:
  1684. return _error(error)
  1685. @bp.route(
  1686. "/device-quality/issues/<issue_uid>/reopen",
  1687. methods=["POST"],
  1688. )
  1689. def reopen_quality_issue(issue_uid):
  1690. try:
  1691. body = request.get_json(silent=True) or {}
  1692. record = get_quality_issue_service().reopen(
  1693. issue_uid,
  1694. reason=body.get("reason"),
  1695. due_at=body.get("due_at"),
  1696. expected_version=body.get("expected_version"),
  1697. actor_uid=_identity().get("id") or _identity().get("sub"),
  1698. )
  1699. return jsonify(success(_quality_issue(record))), 200
  1700. except Exception as error:
  1701. return _error(error)
  1702. @bp.route("/data-elements", methods=["POST"])
  1703. def create_data_element():
  1704. try:
  1705. record = get_data_element_service().create_draft(
  1706. request.get_json(silent=True) or {},
  1707. actor_uid=_identity().get("id") or _identity().get("sub"),
  1708. )
  1709. db.session.commit()
  1710. return jsonify(success(_element(record))), 201
  1711. except Exception as error:
  1712. db.session.rollback()
  1713. return _error(error)
  1714. @bp.route("/data-elements/<element_uid>/transitions", methods=["POST"])
  1715. def transition_data_element(element_uid):
  1716. payload = request.get_json(silent=True) or {}
  1717. identity = _identity()
  1718. target_status = str(payload.get("target_status") or "")
  1719. if (
  1720. target_status in {"published", "deprecated", "retired"}
  1721. and "data-elements:publish" not in set(identity.get("permissions") or [])
  1722. ):
  1723. return jsonify(failed("权限不足", code=403)), 403
  1724. try:
  1725. record = get_data_element_service().transition(
  1726. element_uid,
  1727. target_status,
  1728. expected_version=int(payload.get("expected_version")),
  1729. actor_uid=identity.get("id") or identity.get("sub"),
  1730. )
  1731. db.session.commit()
  1732. return jsonify(success(_element(record))), 200
  1733. except Exception as error:
  1734. db.session.rollback()
  1735. return _error(error)
  1736. @bp.route("/candidate-decisions", methods=["POST"])
  1737. def decide_candidates():
  1738. payload = request.get_json(silent=True) or {}
  1739. try:
  1740. records = get_candidate_decision_service().decide(
  1741. payload.get("decisions") or [],
  1742. actor_uid=_identity().get("id") or _identity().get("sub"),
  1743. )
  1744. return jsonify(success([_decision(record) for record in records])), 200
  1745. except Exception as error:
  1746. return _error(error)
  1747. @bp.route("/sources/files", methods=["POST"])
  1748. def upload_source_file():
  1749. uploaded = request.files.get("file")
  1750. if uploaded is None:
  1751. return jsonify(failed("缺少上传文件", code=400)), 400
  1752. try:
  1753. record, created = get_artifact_service().store(
  1754. request.form.get("source_uid"),
  1755. uploaded.filename,
  1756. uploaded.mimetype,
  1757. uploaded.read(),
  1758. request.form.get("parser_version") or "auto-v1",
  1759. )
  1760. db.session.commit()
  1761. return jsonify(success(_artifact(record))), 201 if created else 200
  1762. except Exception as error:
  1763. db.session.rollback()
  1764. return _error(error)
  1765. @bp.route("/evidence/<evidence_uid>", methods=["GET"])
  1766. def get_evidence(evidence_uid):
  1767. identity = _identity()
  1768. try:
  1769. service = get_evidence_service()
  1770. if request.args.get("download") in {"1", "true", "yes"}:
  1771. if "evidence:download" not in set(identity.get("permissions") or []):
  1772. return jsonify(failed("权限不足", code=403)), 403
  1773. content, filename, media_type = service.download(evidence_uid)
  1774. return send_file(
  1775. io.BytesIO(content),
  1776. mimetype=media_type,
  1777. as_attachment=True,
  1778. download_name=filename,
  1779. )
  1780. preview = service.preview(evidence_uid)
  1781. from app.core.data_research.artifacts import redact_excerpt
  1782. preview["excerpt"] = redact_excerpt(preview.get("excerpt"))
  1783. return jsonify(success(preview)), 200
  1784. except Exception as error:
  1785. return _error(error)
  1786. @bp.route("/ontologies", methods=["GET"])
  1787. def list_ontologies():
  1788. try:
  1789. return jsonify(success([_ontology(item) for item in get_ontology_service().list()])), 200
  1790. except Exception as error:
  1791. return _error(error)
  1792. @bp.route("/ontologies", methods=["POST"])
  1793. def create_ontology():
  1794. try:
  1795. record = get_ontology_service().create(
  1796. request.get_json(silent=True) or {},
  1797. _identity().get("id") or _identity().get("sub"),
  1798. )
  1799. return jsonify(success(_ontology(record))), 201
  1800. except Exception as error:
  1801. return _error(error)
  1802. @bp.route("/ontologies/<ontology_uid>", methods=["GET"])
  1803. def get_ontology(ontology_uid):
  1804. try:
  1805. return jsonify(success(_ontology(get_ontology_service().get(ontology_uid)))), 200
  1806. except Exception as error:
  1807. return _error(error)
  1808. @bp.route("/ontologies/<ontology_uid>/versions", methods=["GET"])
  1809. def list_ontology_versions(ontology_uid):
  1810. try:
  1811. return jsonify(
  1812. success(
  1813. [
  1814. _ontology_version(item)
  1815. for item in get_ontology_service().list_versions(ontology_uid)
  1816. ]
  1817. )
  1818. ), 200
  1819. except Exception as error:
  1820. return _error(error)
  1821. @bp.route("/ontologies/<ontology_uid>/graph", methods=["GET"])
  1822. def get_ontology_graph(ontology_uid):
  1823. try:
  1824. service = get_ontology_service()
  1825. ontology = service.get(ontology_uid)
  1826. version = service.latest_version(ontology_uid)
  1827. data = (
  1828. _ontology_version(version)
  1829. if version is not None
  1830. else {
  1831. "uid": None,
  1832. "ontology_uid": str(ontology_uid),
  1833. "version": 0,
  1834. "parent_version_uid": None,
  1835. "status": "draft",
  1836. "content_hash": None,
  1837. "graph_document": {
  1838. name: []
  1839. for name in (
  1840. "classes",
  1841. "properties",
  1842. "relations",
  1843. "constraints",
  1844. "domain_links",
  1845. "element_mappings",
  1846. )
  1847. },
  1848. "created_by": None,
  1849. }
  1850. )
  1851. response = jsonify(success(data))
  1852. response.headers["ETag"] = f'"{int(ontology.draft_revision)}"'
  1853. return response, 200
  1854. except Exception as error:
  1855. return _error(error)
  1856. @bp.route("/ontologies/<ontology_uid>/graph", methods=["PATCH"])
  1857. def save_ontology_graph(ontology_uid):
  1858. etag = str(request.headers.get("If-Match") or "").strip().strip('"')
  1859. if not etag.isdigit():
  1860. return jsonify(failed("缺少有效 If-Match 修订号", code=428)), 428
  1861. try:
  1862. version = get_ontology_service().save_draft(
  1863. ontology_uid,
  1864. request.get_json(silent=True) or {},
  1865. int(etag),
  1866. _identity().get("id") or _identity().get("sub"),
  1867. )
  1868. response = jsonify(success(_ontology_version(version)))
  1869. response.headers["ETag"] = f'"{int(etag) + 1}"'
  1870. return response, 200
  1871. except Exception as error:
  1872. return _error(error)
  1873. @bp.route("/ontologies/<ontology_uid>/validate", methods=["POST"])
  1874. def validate_ontology(ontology_uid):
  1875. try:
  1876. issues = get_ontology_service().validate(ontology_uid)
  1877. data = [
  1878. {"code": item.code, "message": item.message, "path": item.path}
  1879. for item in issues
  1880. ]
  1881. return jsonify(success(data)), 200
  1882. except Exception as error:
  1883. return _error(error)
  1884. @bp.route("/ontologies/<ontology_uid>/publish", methods=["POST"])
  1885. def publish_ontology(ontology_uid):
  1886. try:
  1887. version = get_ontology_service().publish(
  1888. ontology_uid,
  1889. request.headers.get("Idempotency-Key") or f"{ontology_uid}:publish",
  1890. _identity().get("id") or _identity().get("sub"),
  1891. )
  1892. return jsonify(success(_ontology_version(version))), 200
  1893. except Exception as error:
  1894. return _error(error)
  1895. @bp.route("/ontologies/<ontology_uid>/diff", methods=["GET"])
  1896. def diff_ontology(ontology_uid):
  1897. try:
  1898. return jsonify(success(get_ontology_service().diff(
  1899. ontology_uid, request.args.get("left"), request.args.get("right")
  1900. ))), 200
  1901. except Exception as error:
  1902. return _error(error)
  1903. @bp.route("/ontologies/<ontology_uid>/rollback", methods=["POST"])
  1904. def rollback_ontology(ontology_uid):
  1905. payload = request.get_json(silent=True) or {}
  1906. try:
  1907. version = get_ontology_service().rollback(
  1908. ontology_uid,
  1909. payload.get("target_version_uid"),
  1910. int(payload.get("expected_revision")),
  1911. _identity().get("id") or _identity().get("sub"),
  1912. )
  1913. return jsonify(success(_ontology_version(version))), 201
  1914. except Exception as error:
  1915. return _error(error)
  1916. @bp.route("/ontologies/<ontology_uid>/suggestions", methods=["POST"])
  1917. def generate_ontology_suggestions(ontology_uid):
  1918. try:
  1919. result = get_ontology_dynamic_service().generate(
  1920. ontology_uid,
  1921. request.get_json(silent=True) or {},
  1922. _identity().get("id") or _identity().get("sub"),
  1923. )
  1924. records = (
  1925. result.get("suggestions") or []
  1926. if isinstance(result, dict)
  1927. else result
  1928. )
  1929. suggestions = [
  1930. item
  1931. if isinstance(item, dict)
  1932. else {
  1933. "uid": item.uid,
  1934. "kind": item.kind,
  1935. "payload": item.payload,
  1936. "evidence_uids": list(item.evidence_uids),
  1937. "confidence": item.confidence,
  1938. "source": item.source,
  1939. "model_version": item.model_version,
  1940. "prompt_version": item.prompt_version,
  1941. }
  1942. for item in records
  1943. ]
  1944. data = {
  1945. "change_set_uid": (
  1946. result.get("change_set_uid")
  1947. if isinstance(result, dict)
  1948. else None
  1949. ),
  1950. "suggestions": suggestions,
  1951. }
  1952. return jsonify(success(data)), 201
  1953. except Exception as error:
  1954. return _error(error)
  1955. @bp.route(
  1956. "/ontologies/<ontology_uid>/change-sets/<change_set_uid>/decisions",
  1957. methods=["POST"],
  1958. )
  1959. def decide_ontology_change_set(ontology_uid, change_set_uid):
  1960. payload = request.get_json(silent=True) or {}
  1961. try:
  1962. result = get_ontology_dynamic_service().decide(
  1963. ontology_uid,
  1964. change_set_uid,
  1965. payload.get("decisions") or [],
  1966. _identity().get("id") or _identity().get("sub"),
  1967. )
  1968. return jsonify(success(result)), 200
  1969. except Exception as error:
  1970. return _error(error)
  1971. @bp.route("/ontologies/<ontology_uid>/export", methods=["GET"])
  1972. def export_ontology(ontology_uid):
  1973. try:
  1974. content, media_type, filename = get_ontology_exchange_service().export(
  1975. ontology_uid, request.args.get("format") or "json"
  1976. )
  1977. return send_file(
  1978. io.BytesIO(content),
  1979. mimetype=media_type,
  1980. as_attachment=True,
  1981. download_name=filename,
  1982. )
  1983. except Exception as error:
  1984. return _error(error)
  1985. @bp.route("/ontologies/import", methods=["POST"])
  1986. def import_ontology():
  1987. try:
  1988. uploaded = request.files.get("file")
  1989. content = uploaded.read() if uploaded is not None else request.get_data(cache=False)
  1990. result = get_ontology_exchange_service().import_document(
  1991. content,
  1992. request.args.get("format") or "json",
  1993. _identity().get("id") or _identity().get("sub"),
  1994. )
  1995. return jsonify(success(result)), 201
  1996. except Exception as error:
  1997. return _error(error)
  1998. @bp.route("/semantic/properties/<property_uid>", methods=["GET"])
  1999. def query_semantic_property(property_uid):
  2000. domain = str(request.args.get("business_domain_uid") or "").strip()
  2001. identity = _identity()
  2002. scoped = set(identity.get("business_domains") or [domain])
  2003. if "*" not in scoped and domain not in scoped:
  2004. return jsonify(failed("权限不足", code=403)), 403
  2005. try:
  2006. result = get_semantic_query_service().trace_property(
  2007. property_uid,
  2008. allowed_domains={domain},
  2009. limit=int(request.args.get("limit") or 20),
  2010. after_uid=request.args.get("after_uid"),
  2011. )
  2012. return jsonify(success(result)), 200
  2013. except Exception as error:
  2014. return _error(error)