routes.py 80 KB

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