routes.py 78 KB

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