"""HTTP orchestration boundary for data-research ingestion jobs.""" from __future__ import annotations import io import logging import uuid from flask import current_app, g, jsonify, request, send_file from app import db from app.api.data_development import bp from app.core.common.timezone_utils import now_china from app.core.data_research.errors import DataResearchError from app.models.result import failed, success logger = logging.getLogger(__name__) def get_ingestion_service(): from app.core.data_research.ingestion import IngestionService from app.core.data_research.repository import SqlAlchemyIngestionJobRepository return IngestionService( SqlAlchemyIngestionJobRepository(db.session), commit=db.session.commit, rollback=db.session.rollback, ) def get_database_source_registration_service(): from app.core.data_research.repository import ( SqlAlchemyIngestionSourceRepository, ) from app.core.data_research.sources import ( DatabaseSourceRegistrationService, ) from app.core.data_source.runtime import get_data_source_manager manager = get_data_source_manager() return DatabaseSourceRegistrationService( SqlAlchemyIngestionSourceRepository(db.session), definition_resolver=manager.definitions.get, commit=db.session.commit, rollback=db.session.rollback, ) def get_catalog_snapshot_repository(): from app.core.data_research.repository import ( SqlAlchemyCatalogSnapshotRepository, ) return SqlAlchemyCatalogSnapshotRepository(db.session) def get_catalog_ingestion_executor(): from app.core.data_research.catalog.execution import ( CatalogIngestionExecutor, ) from app.core.data_research.catalog.mysql import MySqlCatalogCollector from app.core.data_research.catalog.postgresql import ( PostgreSqlCatalogCollector, ) from app.core.data_research.catalog.service import CatalogCollectionService from app.core.data_research.errors import IngestionSourceInvalid from app.core.data_source.runtime import get_data_source_manager manager = get_data_source_manager() def collector_resolver(database_type): if database_type == "postgresql": return PostgreSqlCatalogCollector() if database_type == "mysql": return MySqlCatalogCollector() raise IngestionSourceInvalid( f"database type {database_type} is not supported" ) collector = CatalogCollectionService( manager, definition_resolver=manager.definitions.get, collector_resolver=collector_resolver, ) return CatalogIngestionExecutor( get_ingestion_service(), collector, get_catalog_snapshot_repository(), commit=db.session.commit, rollback=db.session.rollback, ) def get_data_element_service(): from app.core.data_research.data_elements import DataElementService from app.core.data_research.repository import SqlAlchemyDataElementRepository from app.core.events.outbox import enqueue_outbox return DataElementService( SqlAlchemyDataElementRepository(db.session), outbox_enqueue=lambda **event: enqueue_outbox(db.session, **event), ) def get_device_asset_service(): from app.core.data_research.device_asset_repository import ( SqlAlchemyDeviceAssetRepository, ) from app.core.data_research.device_assets import DeviceAssetService return DeviceAssetService( SqlAlchemyDeviceAssetRepository(db.session), commit=db.session.commit, rollback=db.session.rollback, ) def _device_ontology_authorizer(ontology_repository): from app.core.data_research.device_semantics import ( DeviceOntologyPublicationAuthorizer, ) from app.core.governance.responsibilities import ( ResponsibilityService, SqlAlchemyResponsibilityRepository, ) responsibilities = ResponsibilityService( SqlAlchemyResponsibilityRepository(db.session) ) return DeviceOntologyPublicationAuthorizer( ontology_repository, responsibility_lookup=responsibilities.get, ) def get_device_semantic_service(): from app.core.data_research.device_semantics import DeviceSemanticService from app.core.data_research.ontology.repository import ( SqlAlchemyOntologyRepository, ) return DeviceSemanticService( SqlAlchemyOntologyRepository(db.session), commit=db.session.commit, rollback=db.session.rollback, ) def get_device_semantic_code_service(): from app.core.data_research.device_semantic_repository import ( SqlAlchemyDeviceSemanticCodeRepository, ) from app.core.data_research.device_semantics import ( DeviceSemanticCodeService, ) from app.core.data_research.ontology.repository import ( SqlAlchemyOntologyRepository, ) ontology_repository = SqlAlchemyOntologyRepository(db.session) authorizer = _device_ontology_authorizer(ontology_repository) return DeviceSemanticCodeService( SqlAlchemyDeviceSemanticCodeRepository(db.session), review_authorizer=authorizer.assert_accountable, commit=db.session.commit, rollback=db.session.rollback, ) def get_device_entity_resolution_service(): from app.core.data_research.device_entity_repository import ( SqlAlchemyDeviceEntityResolutionRepository, ) from app.core.data_research.device_entity_resolution import ( DeviceEntityForbidden, DeviceEntityResolutionService, ) from app.core.governance.responsibilities import ( ResponsibilityService, SqlAlchemyResponsibilityRepository, ) responsibilities = ResponsibilityService( SqlAlchemyResponsibilityRepository(db.session) ) def assert_accountable(actor_uid): matrix = responsibilities.get( "device_mapping", "DEVICE_ENTITY_RESOLUTION", ) accountable = [ item for item in matrix.get("assignments", []) if item.get("responsibility_role") == "asset_manager" and item.get("raci_role") == "accountable" ] if len(accountable) != 1 or str(accountable[0].get("user_id")) != str( actor_uid ): raise DeviceEntityForbidden( "only the accountable device mapping asset manager may decide" ) return DeviceEntityResolutionService( SqlAlchemyDeviceEntityResolutionRepository(db.session), review_authorizer=assert_accountable, auto_merge_enabled=current_app.config.get( "DEVICE_ENTITY_AUTO_MERGE_ENABLED", False, ), commit=db.session.commit, rollback=db.session.rollback, ) def get_device_quality_service(): from app.core.data_research.device_quality import DeviceQualityService from app.core.data_research.device_quality_repository import ( SqlAlchemyDeviceQualityRepository, ) from app.core.data_research.errors import DeviceQualityForbidden from app.core.governance.responsibilities import ( ResponsibilityService, SqlAlchemyResponsibilityRepository, ) responsibilities = ResponsibilityService( SqlAlchemyResponsibilityRepository(db.session) ) def assert_accountable(actor_uid): matrix = responsibilities.get( "device_quality", "DEVICE_QUALITY", ) accountable = [ item for item in matrix.get("assignments", []) if item.get("responsibility_role") == "asset_manager" and item.get("raci_role") == "accountable" ] if len(accountable) != 1 or str(accountable[0].get("user_id")) != str( actor_uid ): raise DeviceQualityForbidden( "only the accountable device quality asset manager may publish" ) return DeviceQualityService( SqlAlchemyDeviceQualityRepository(db.session), publish_authorizer=assert_accountable, commit=db.session.commit, rollback=db.session.rollback, ) def get_quality_issue_service(): from app.core.data_research.errors import QualityIssueForbidden from app.core.data_research.quality_issue_repository import ( SqlAlchemyQualityIssueRepository, ) from app.core.data_research.quality_issues import QualityIssueService from app.core.governance.responsibilities import ( ResponsibilityService, SqlAlchemyResponsibilityRepository, ) repository = SqlAlchemyQualityIssueRepository(db.session) responsibilities = ResponsibilityService( SqlAlchemyResponsibilityRepository(db.session) ) def assert_accountable(actor_uid): matrix = responsibilities.get( "quality_issue", "DEVICE_QUALITY_ISSUES", ) accountable = [ item for item in matrix.get("assignments", []) if item.get("responsibility_role") == "asset_manager" and item.get("raci_role") == "accountable" ] if len(accountable) != 1 or str(accountable[0].get("user_id")) != str( actor_uid ): raise QualityIssueForbidden( "only the accountable quality issue asset manager may review" ) return QualityIssueService( repository, review_authorizer=assert_accountable, is_admin=lambda actor_uid: repository.user_has_role( actor_uid, "admin", ), commit=db.session.commit, rollback=db.session.rollback, ) def get_device_observability_service(): from app.core.data_research.device_observability import ( DeviceObservabilityService, ) from app.core.data_research.device_observability_repository import ( SqlAlchemyDeviceObservabilityRepository, ) return DeviceObservabilityService( SqlAlchemyDeviceObservabilityRepository(db.session), commit=db.session.commit, rollback=db.session.rollback, ) def get_candidate_decision_service(): from app.core.data_research.candidate_decisions import CandidateDecisionService from app.core.data_research.data_elements import DataElementService from app.core.data_research.repository import ( SqlAlchemyCandidateDecisionRepository, SqlAlchemyDataElementRepository, ) return CandidateDecisionService( SqlAlchemyCandidateDecisionRepository(db.session), data_elements=DataElementService(SqlAlchemyDataElementRepository(db.session)), ) def _artifact_storage(): from minio import Minio from app.core.data_research.artifacts import MinioArtifactStorage client = Minio( current_app.config["MINIO_HOST"], access_key=current_app.config["MINIO_USER"], secret_key=current_app.config["MINIO_PASSWORD"], secure=bool(current_app.config.get("MINIO_SECURE")), ) return MinioArtifactStorage(client, current_app.config["MINIO_BUCKET"]) def get_artifact_service(): from app.core.common.identifiers import new_governance_uid from app.core.data_research.artifacts import ( ArtifactService, SqlAlchemyArtifactRepository, ) from app.core.data_research.file_policy import FilePolicy return ArtifactService( SqlAlchemyArtifactRepository(db.session), _artifact_storage(), FilePolicy( max_bytes=int(current_app.config.get("DATA_RESEARCH_FILE_MAX_BYTES", 25 * 1024 * 1024)), max_pages=int(current_app.config.get("DATA_RESEARCH_FILE_MAX_PAGES", 500)), ), uid_factory=new_governance_uid, ) def get_evidence_service(): from app.core.data_research.artifacts import EvidenceService return EvidenceService(db.session, _artifact_storage()) def get_ontology_service(): from app.core.data_research.ontology.publication import ( OntologyApplicationService, OntologyPublicationService, ) from app.core.data_research.ontology.repository import SqlAlchemyOntologyRepository from app.core.events.outbox import enqueue_outbox repository = SqlAlchemyOntologyRepository(db.session) authorizer = _device_ontology_authorizer(repository) publication = OntologyPublicationService( repository, outbox_enqueue=lambda **event: enqueue_outbox(db.session, **event), publication_authorizer=authorizer, commit=db.session.commit, rollback=db.session.rollback, ) return OntologyApplicationService( repository, publication=publication, commit=db.session.commit, rollback=db.session.rollback, ) def get_ontology_dynamic_service(): from app.core.data_research.ontology.change_sets import ( SqlAlchemyDynamicOntologyService, ) return SqlAlchemyDynamicOntologyService(db.session) def get_ontology_exchange_service(): from app.core.data_research.ontology.exchange import OntologyExchangeService from app.core.data_research.ontology.repository import SqlAlchemyOntologyRepository return OntologyExchangeService( SqlAlchemyOntologyRepository(db.session), commit=db.session.commit, rollback=db.session.rollback, ) def get_semantic_query_service(): from app.core.data_research.ontology.query import ( Neo4jSemanticRepository, SemanticQueryService, ) from app.services.neo4j_driver import neo4j_driver return SemanticQueryService(Neo4jSemanticRepository(neo4j_driver)) def _identity(): return getattr(g, "current_user", {}) or {} def get_governance_metrics_service(): from app.core.data_research.governance_metric_repository import ( SqlAlchemyGovernanceMetricRepository, ) from app.core.data_research.governance_metrics import ( GovernanceMetricsService, ) return GovernanceMetricsService( SqlAlchemyGovernanceMetricRepository(db.session) ) def _governance_metric_access(): from app.core.data_research.governance_metrics import ( GovernanceMetricAccess, ) from app.core.knowledge.access import build_access_context context = build_access_context( db.session, identity=_identity(), requested_business_domains=None, correlation_id=str(uuid.uuid4()), ) return GovernanceMetricAccess( global_access=context.global_access, business_domain_uids=tuple(context.business_domain_uids), ) def _record(record): return { "uid": str(record.uid), "source_uid": str(record.source_uid), "artifact_uid": str(record.artifact_uid) if record.artifact_uid else None, "job_type": record.job_type, "parser_version": record.parser_version, "status": record.status, "parameters": dict(record.parameters or {}), "statistics": dict(record.statistics or {}), "last_error": record.last_error, "attempt_count": int(record.attempt_count or 0), "failure_stage": record.failure_stage, "actor_uid": record.actor_uid, "created_at": record.created_at.isoformat() if record.created_at else None, "updated_at": record.updated_at.isoformat() if record.updated_at else None, "started_at": record.started_at.isoformat() if record.started_at else None, "finished_at": record.finished_at.isoformat() if record.finished_at else None, } def _catalog_snapshot(record): return { "uid": str(record.uid), "job_uid": str(record.job_uid), "source_uid": str(record.source_uid), "attempt": int(record.attempt), "database_type": record.database_type, "content_hash": record.content_hash, "snapshot": dict(record.snapshot or {}), "evidence_count": int(record.evidence_count), "created_at": ( record.created_at.isoformat() if record.created_at else None ), } def _element(record): return { "uid": str(record.uid), "code": record.code, "status": record.status, "current_version": int(record.current_version), "snapshot": dict(record.snapshot or {}), "created_by": record.created_by, "updated_by": record.updated_by, } def _decision(record): return { "uid": str(record.uid), "candidate_uid": str(record.candidate_uid), "action": record.action, "data_element_uid": ( str(record.data_element_uid) if record.data_element_uid else None ), "evidence_uids": list(record.evidence_uids), "actor_uid": record.actor_uid, "reason": record.reason, } def _artifact(record): return { "uid": str(record.uid), "source_uid": str(record.source_uid), "filename": record.filename, "media_type": record.media_type, "size_bytes": int(record.size_bytes), "content_hash": record.content_hash, "parser_version": record.parser_version, } def _ontology(record): return { "uid": str(record.uid), "code": record.code, "name": record.name, "owner_uid": record.owner_uid, "status": record.status, "draft_revision": int(record.draft_revision), "active_version_uid": record.active_version_uid, "domain_links": [ {"domain_uid": link.domain_uid, "role": link.role} for link in record.domain_links ], } def _ontology_version(record): return { "uid": str(record.uid), "ontology_uid": str(record.ontology_uid), "version": int(record.version), "parent_version_uid": record.parent_version_uid, "status": record.status, "content_hash": record.content_hash, "graph_document": record.graph_document.to_dict(), "created_by": record.created_by, } def _iso(value): return value.isoformat() if value else None def _device_asset(record): return { "uid": str(record.uid), "asset_type": record.asset_type, "name": record.name, "status": record.status, "current_version": int(record.current_version), "location": record.location, "organization": record.organization, "responsible_person": record.responsible_person, "attributes": dict(record.attributes or {}), "created_by": record.created_by, "updated_by": record.updated_by, "created_at": _iso(record.created_at), "updated_at": _iso(record.updated_at), } def _device_asset_mapping(record): return { "uid": str(record.uid), "asset_uid": str(record.asset_uid), "source_uid": str(record.source_uid), "source_entity": record.source_entity, "asset_type": record.asset_type, "source_code": record.source_code, "source_updated_at": _iso(record.source_updated_at), "first_seen_at": _iso(record.first_seen_at), "last_seen_at": _iso(record.last_seen_at), } def _device_asset_detail(detail): data = _device_asset(detail.asset) data["source_mappings"] = [ _device_asset_mapping(mapping) for mapping in detail.mappings ] return data def _device_asset_version(record): return { "uid": str(record.uid), "asset_uid": str(record.asset_uid), "version": int(record.version), "snapshot": dict(record.snapshot or {}), "source_mapping_uid": str(record.source_mapping_uid), "actor_uid": record.actor_uid, "created_at": _iso(record.created_at), } def _device_asset_import_result(result): return { "records": [ { "action": item.action, "asset": _device_asset(item.asset), "source_mapping": _device_asset_mapping(item.mapping), } for item in result.items ], "created_count": int(result.created_count), "updated_count": int(result.updated_count), "unchanged_count": int(result.unchanged_count), } def _device_semantic_profile(result): if result is None: return { "bootstrapped": False, "ontology": None, "version": None, "profile": { "ready_to_publish": False, "class_count": 0, "relation_count": 0, "mapping_count": 0, "missing_classes": [], "missing_relations": [], "missing_mappings": [], }, } profile = result.profile return { "bootstrapped": True, "created": bool(result.created), "ontology": _ontology(result.ontology), "version": _ontology_version(result.version), "profile": { "ready_to_publish": bool(profile.ready_to_publish), "class_count": int(profile.class_count), "relation_count": int(profile.relation_count), "mapping_count": int(profile.mapping_count), "missing_classes": list(profile.missing_classes), "missing_relations": list(profile.missing_relations), "missing_mappings": list(profile.missing_mappings), }, } def _device_semantic_code(record): return { "uid": str(record.uid), "ontology_uid": str(record.ontology_uid), "code_type": record.code_type, "canonical_code": record.canonical_code, "canonical_name": record.canonical_name, "definition": record.definition, "status": record.status, "current_version": int(record.current_version), "source_mappings": [ dict(item) for item in record.source_mappings ], "evidence_uids": list(record.evidence_uids), "suggestion_source": record.suggestion_source, "confidence": record.confidence, "created_by": record.created_by, "updated_by": record.updated_by, "created_at": _iso(record.created_at), "updated_at": _iso(record.updated_at), } def _device_semantic_version(record): return { "uid": str(record.uid), "code_uid": str(record.code_uid), "version": int(record.version), "snapshot": dict(record.snapshot), "created_by": record.created_by, "created_at": _iso(record.created_at), } def _device_semantic_review(record): return { "uid": str(record.uid), "code_uid": str(record.code_uid), "version": int(record.version), "decision": record.decision, "reason": record.reason, "actor_uid": record.actor_uid, "created_at": _iso(record.created_at), } def _device_entity_candidate(record): return { "uid": str(record.uid), "left_asset_uid": str(record.left_asset_uid), "right_asset_uid": str(record.right_asset_uid), "canonical_asset_uid": ( str(record.canonical_asset_uid) if record.canonical_asset_uid else None ), "status": record.status, "suggestion_source": record.suggestion_source, "confidence": float(record.confidence), "explanation": [ dict(item) for item in record.explanation ], "evidence_uids": list(record.evidence_uids), "model_provider": record.model_provider, "model_name": record.model_name, "current_version": int(record.current_version), "created_by": record.created_by, "reviewed_by": record.reviewed_by, "created_at": _iso(record.created_at), "updated_at": _iso(record.updated_at), } def _device_entity_review(record): return { "uid": str(record.uid), "candidate_uid": str(record.candidate_uid), "version": int(record.version), "decision": record.decision, "reason": record.reason, "actor_uid": record.actor_uid, "created_at": _iso(record.created_at), } def _device_entity_merge(record): return { "uid": str(record.uid), "candidate_uid": str(record.candidate_uid), "canonical_asset_uid": str(record.canonical_asset_uid), "member_asset_uid": str(record.member_asset_uid), "review_uid": str(record.review_uid), "snapshot": dict(record.snapshot), "actor_uid": record.actor_uid, "created_at": _iso(record.created_at), } def _device_entity_rollback(record): return { "uid": str(record.uid), "merge_uid": str(record.merge_uid), "candidate_uid": str(record.candidate_uid), "reason": record.reason, "snapshot": dict(record.snapshot), "actor_uid": record.actor_uid, "created_at": _iso(record.created_at), } def _device_entity_generation(result): return { "records": [ _device_entity_candidate(record) for record in result.records ], "created_count": int(result.created_count), "existing_count": int(result.existing_count), "evaluated_pair_count": int(result.evaluated_pair_count), "auto_merged_count": int(result.auto_merged_count), "auto_merge_enabled": bool( current_app.config.get( "DEVICE_ENTITY_AUTO_MERGE_ENABLED", False, ) ), } def _device_quality_version(record): if record is None: return None return { "uid": str(record.uid), "profile_uid": str(record.profile_uid), "version": int(record.version), "status": record.status, "rules": [dict(item) for item in record.rules], "content_hash": record.content_hash, "created_by": record.created_by, "created_at": _iso(record.created_at), "published_by": record.published_by, "published_at": _iso(record.published_at), } def _device_quality_run(record): return { "uid": str(record.uid), "policy_version_uid": str(record.policy_version_uid), "policy_hash": record.policy_hash, "source_uid": ( str(record.source_uid) if record.source_uid else None ), "status": record.status, "total_assets": int(record.total_assets), "total_violations": int(record.total_violations), "score": float(record.score), "created_by": record.created_by, "created_at": _iso(record.created_at), } def _device_quality_rule_result(record): return { "uid": str(record.uid), "run_uid": str(record.run_uid), "rule_code": record.rule_code, "severity": record.severity, "weight": float(record.weight), "status": record.status, "evaluated_count": int(record.evaluated_count), "violation_count": int(record.violation_count), "sampled_count": int(record.sampled_count), "pass_rate": float(record.pass_rate), "weighted_score": float(record.weighted_score), "created_at": _iso(record.created_at), } def _device_quality_violation(record): return { "uid": str(record.uid), "run_uid": str(record.run_uid), "rule_code": record.rule_code, "severity": record.severity, "asset_uid": str(record.asset_uid), "field_name": record.field_name, "source_uid": ( str(record.source_uid) if record.source_uid else None ), "source_mapping_uid": ( str(record.source_mapping_uid) if record.source_mapping_uid else None ), "message": record.message, "evidence": dict(record.evidence or {}), "created_at": _iso(record.created_at), "expires_at": _iso(record.expires_at), } def _device_quality_asset_score(record): return { "uid": str(record.uid), "run_uid": str(record.run_uid), "asset_uid": str(record.asset_uid), "asset_type": record.asset_type, "status": record.status, "evaluated_rule_count": int(record.evaluated_rule_count), "violation_count": int(record.violation_count), "score": float(record.score), "created_at": _iso(record.created_at), } def _quality_issue(record): return { "uid": str(record.uid), "issue_code": record.issue_code, "source_violation_uid": str(record.source_violation_uid), "source_run_uid": str(record.source_run_uid), "rule_code": record.rule_code, "severity": record.severity, "priority": record.priority, "asset_uid": str(record.asset_uid), "field_name": record.field_name, "source_uid": ( str(record.source_uid) if record.source_uid else None ), "source_mapping_uid": ( str(record.source_mapping_uid) if record.source_mapping_uid else None ), "message": record.message, "evidence": dict(record.evidence or {}), "recurrence_key": record.recurrence_key, "occurrence_number": int(record.occurrence_number), "status": record.status, "assignee_uid": ( str(record.assignee_uid) if record.assignee_uid else None ), "due_at": _iso(record.due_at), "is_overdue": bool(record.is_overdue(now_china())), "current_version": int(record.current_version), "created_by": str(record.created_by), "updated_by": str(record.updated_by), "created_at": _iso(record.created_at), "updated_at": _iso(record.updated_at), "closed_at": _iso(record.closed_at), } def _quality_issue_remediation(record): if record is None: return None return { "uid": str(record.uid), "issue_uid": str(record.issue_uid), "round_number": int(record.round_number), "summary": record.summary, "evidence_refs": [dict(item) for item in record.evidence_refs], "submitted_by": str(record.submitted_by), "submitted_at": _iso(record.submitted_at), "review_status": record.review_status, "reviewed_by": ( str(record.reviewed_by) if record.reviewed_by else None ), "reviewed_at": _iso(record.reviewed_at), "review_note": record.review_note, "verification_run_uid": ( str(record.verification_run_uid) if record.verification_run_uid else None ), } def _quality_issue_timeline(record): return { "uid": str(record.uid), "issue_uid": str(record.issue_uid), "action": record.action, "from_status": record.from_status, "to_status": record.to_status, "actor_uid": str(record.actor_uid), "note": record.note, "payload": dict(record.payload or {}), "created_at": _iso(record.created_at), } def _device_operational_event(record): return { "uid": str(record.uid), "source_uid": str(record.source_uid), "source_entity": record.source_entity, "source_code": record.source_code, "event_type": record.event_type, "asset_uid": str(record.asset_uid), "component_uid": ( str(record.component_uid) if record.component_uid else None ), "title": record.title, "severity": record.severity, "status": record.status, "occurred_at": _iso(record.occurred_at), "ended_at": _iso(record.ended_at), "evidence_refs": [dict(item) for item in record.evidence_refs], "content_hash": record.content_hash, "created_by": str(record.created_by), "created_at": _iso(record.created_at), } def _device_evidence_relation(record): return { "uid": str(record.uid), "from_kind": record.from_kind, "from_uid": str(record.from_uid), "relation_type": record.relation_type, "to_kind": record.to_kind, "to_uid": str(record.to_uid), "evidence_refs": [dict(item) for item in record.evidence_refs], "source": record.source, "created_by": str(record.created_by), "created_at": _iso(record.created_at), } def _device_graph_node(record): return { "kind": record.kind, "uid": str(record.uid), "node_type": record.node_type, "label": record.label, "occurred_at": _iso(record.occurred_at), "evidence_refs": [dict(item) for item in record.evidence_refs], } def _device_evidence_graph(record): return { "anchor_kind": record.anchor_kind, "anchor_uid": str(record.anchor_uid), "nodes": [_device_graph_node(item) for item in record.nodes], "relations": [ _device_evidence_relation(item) for item in record.relations ], "truncated": bool(record.truncated), } def _root_cause_analysis(record): return { "analysis_status": record.analysis_status, "conclusion": record.conclusion, "anchor": _device_graph_node(record.anchor), "candidates": [ { "event_uid": str(item.event_uid), "event_type": item.event_type, "title": item.title, "occurred_at": _iso(item.occurred_at), "support_level": item.support_level, "path_node_uids": list(item.path_node_uids), "path_relation_uids": list(item.path_relation_uids), "path_relation_types": list(item.path_relation_types), "evidence_refs": [ dict(evidence) for evidence in item.evidence_refs ], } for item in record.candidates ], "graph": _device_evidence_graph(record.graph), "limitations": list(record.limitations), "generated_at": _iso(record.generated_at), } def _device_asset_page(name, *, default, maximum): from app.core.data_research.errors import DeviceAssetInvalid raw = request.args.get(name) try: value = default if raw in (None, "") else int(raw) except (TypeError, ValueError) as error: raise DeviceAssetInvalid(f"{name} must be an integer") from error if value < 1 or value > maximum: raise DeviceAssetInvalid( f"{name} must be between 1 and {maximum}" ) return value def _error(error): if isinstance(error, DataResearchError): return ( jsonify( failed( str(error), code=error.http_status, error={"code": error.code}, ) ), error.http_status, ) logger.exception("data-research ingestion request failed") return ( jsonify( failed( "数据采集任务处理失败", code=500, error={"code": "DATA_RESEARCH_ERROR"}, ) ), 500, ) @bp.route("/ingestion-jobs", methods=["POST"]) def create_ingestion_job(): payload = request.get_json(silent=True) or {} try: actor_uid = _identity().get("id") or _identity().get("sub") if payload.get("job_type") == "catalog_collect": get_database_source_registration_service().ensure( payload.get("source_uid"), actor_uid=actor_uid, ) record, created = get_ingestion_service().create_job( payload, actor_uid=actor_uid, ) return jsonify(success(_record(record))), 201 if created else 200 except Exception as error: return _error(error) @bp.route("/ingestion-jobs", methods=["GET"]) def list_ingestion_jobs(): filters = { name: request.args.get(name) for name in ("status", "source_uid") if request.args.get(name) } try: records = get_ingestion_service().list_jobs(filters) return jsonify( success({"records": [_record(item) for item in records], "total": len(records)}) ), 200 except Exception as error: return _error(error) @bp.route("/ingestion-jobs/", methods=["GET"]) def get_ingestion_job(job_uid): try: return jsonify(success(_record(get_ingestion_service().get_job(job_uid)))), 200 except Exception as error: return _error(error) @bp.route("/ingestion-jobs//execute", methods=["POST"]) def execute_ingestion_job(job_uid): try: record = get_catalog_ingestion_executor().execute(job_uid) return jsonify(success(_record(record))), 200 except Exception as error: return _error(error) @bp.route( "/ingestion-jobs//catalog-snapshots", methods=["GET"], ) def list_catalog_snapshots(job_uid): try: records = get_catalog_snapshot_repository().list(job_uid) return jsonify( success( { "records": [ _catalog_snapshot(record) for record in records ], "total": len(records), } ) ), 200 except Exception as error: return _error(error) @bp.route("/ingestion-jobs//evidence", methods=["GET"]) def list_ingestion_job_evidence(job_uid): try: records = get_evidence_service().list_for_job(job_uid) return jsonify( success({"records": records, "total": len(records)}) ), 200 except Exception as error: return _error(error) @bp.route("/ingestion-jobs//retry", methods=["POST"]) def retry_ingestion_job(job_uid): try: return jsonify(success(_record(get_ingestion_service().retry(job_uid)))), 200 except Exception as error: return _error(error) @bp.route("/ingestion-jobs//cancel", methods=["POST"]) def cancel_ingestion_job(job_uid): try: service = get_ingestion_service() record = service.get_job(job_uid) identity = _identity() permissions = set(identity.get("permissions") or []) actor_uid = identity.get("id") or identity.get("sub") if record.actor_uid != actor_uid and "ingestion:admin" not in permissions: return jsonify(failed("权限不足", code=403)), 403 return jsonify(success(_record(service.cancel(job_uid)))), 200 except Exception as error: return _error(error) @bp.route("/data-elements", methods=["GET"]) def list_data_elements(): try: records = get_data_element_service().list( status=request.args.get("status"), business_domain_uid=request.args.get("business_domain_uid"), ) return jsonify( success( { "records": [_element(item) for item in records], "total": len(records), } ) ), 200 except Exception as error: return _error(error) @bp.route("/device-assets", methods=["GET"]) def list_device_assets(): filters = { name: request.args.get(name) for name in ("keyword", "asset_type", "status", "source_uid") if request.args.get(name) } try: page = _device_asset_page( "page", default=1, maximum=1_000_000, ) page_size = _device_asset_page( "page_size", default=20, maximum=100, ) records, total = get_device_asset_service().search( filters, page=page, page_size=page_size, ) return jsonify( success( { "records": [ _device_asset_detail(record) for record in records ], "total": int(total), "page": page, "page_size": page_size, } ) ), 200 except Exception as error: return _error(error) @bp.route("/device-assets/import", methods=["POST"]) def import_device_assets(): try: result = get_device_asset_service().import_records( request.get_json(silent=True) or {}, actor_uid=_identity().get("id") or _identity().get("sub"), ) return jsonify(success(_device_asset_import_result(result))), 200 except Exception as error: return _error(error) @bp.route("/device-assets/", methods=["GET"]) def get_device_asset(asset_uid): try: return jsonify( success( _device_asset_detail( get_device_asset_service().get(asset_uid) ) ) ), 200 except Exception as error: return _error(error) @bp.route("/device-assets//versions", methods=["GET"]) def list_device_asset_versions(asset_uid): try: records = get_device_asset_service().versions(asset_uid) return jsonify( success( { "records": [ _device_asset_version(record) for record in records ], "total": len(records), } ) ), 200 except Exception as error: return _error(error) @bp.route("/device-semantics/bootstrap", methods=["POST"]) def bootstrap_device_semantics(): try: result = get_device_semantic_service().bootstrap( request.get_json(silent=True) or {}, actor_uid=_identity().get("id") or _identity().get("sub"), ) return ( jsonify(success(_device_semantic_profile(result))), 201 if result.created else 200, ) except Exception as error: return _error(error) @bp.route("/device-semantics/profile", methods=["GET"]) def get_device_semantic_profile(): try: result = get_device_semantic_service().profile() return jsonify(success(_device_semantic_profile(result))), 200 except Exception as error: return _error(error) @bp.route("/device-semantics/codes", methods=["GET"]) def list_device_semantic_codes(): filters = { name: request.args.get(name) for name in ("ontology_uid", "code_type", "status", "keyword") if request.args.get(name) } try: records, total = get_device_semantic_code_service().search( filters, page=request.args.get("page", 1), page_size=request.args.get("page_size", 20), ) return jsonify( success( { "records": [ _device_semantic_code(record) for record in records ], "total": int(total), "page": int(request.args.get("page", 1)), "page_size": int(request.args.get("page_size", 20)), } ) ), 200 except Exception as error: return _error(error) @bp.route("/device-semantics/codes", methods=["POST"]) def create_device_semantic_code(): try: record = get_device_semantic_code_service().create( request.get_json(silent=True) or {}, actor_uid=_identity().get("id") or _identity().get("sub"), ) return jsonify(success(_device_semantic_code(record))), 201 except Exception as error: return _error(error) @bp.route("/device-semantics/codes/", methods=["GET"]) def get_device_semantic_code(code_uid): try: record = get_device_semantic_code_service().get(code_uid) return jsonify(success(_device_semantic_code(record))), 200 except Exception as error: return _error(error) @bp.route( "/device-semantics/codes//revisions", methods=["POST"], ) def revise_device_semantic_code(code_uid): payload = request.get_json(silent=True) or {} try: record = get_device_semantic_code_service().revise( code_uid, payload, expected_version=payload.get("expected_version"), actor_uid=_identity().get("id") or _identity().get("sub"), ) return jsonify(success(_device_semantic_code(record))), 200 except Exception as error: return _error(error) @bp.route( "/device-semantics/codes//submit", methods=["POST"], ) def submit_device_semantic_code(code_uid): payload = request.get_json(silent=True) or {} try: record = get_device_semantic_code_service().submit( code_uid, expected_version=payload.get("expected_version"), actor_uid=_identity().get("id") or _identity().get("sub"), ) return jsonify(success(_device_semantic_code(record))), 200 except Exception as error: return _error(error) @bp.route( "/device-semantics/codes//review", methods=["POST"], ) def review_device_semantic_code(code_uid): try: record, review = get_device_semantic_code_service().review( code_uid, request.get_json(silent=True) or {}, actor_uid=_identity().get("id") or _identity().get("sub"), ) return jsonify( success( { "record": _device_semantic_code(record), "review": _device_semantic_review(review), } ) ), 200 except Exception as error: return _error(error) @bp.route( "/device-semantics/codes//versions", methods=["GET"], ) def list_device_semantic_code_versions(code_uid): try: records = get_device_semantic_code_service().versions(code_uid) return jsonify( success( { "records": [ _device_semantic_version(record) for record in records ], "total": len(records), } ) ), 200 except Exception as error: return _error(error) @bp.route( "/device-semantics/codes//reviews", methods=["GET"], ) def list_device_semantic_code_reviews(code_uid): try: records = get_device_semantic_code_service().reviews(code_uid) return jsonify( success( { "records": [ _device_semantic_review(record) for record in records ], "total": len(records), } ) ), 200 except Exception as error: return _error(error) @bp.route("/device-entities/candidates", methods=["GET"]) def list_device_entity_candidates(): filters = { name: request.args.get(name) for name in ("status", "suggestion_source") if request.args.get(name) } try: records, total = get_device_entity_resolution_service().search( filters, page=request.args.get("page", 1), page_size=request.args.get("page_size", 20), ) return jsonify( success( { "records": [ _device_entity_candidate(record) for record in records ], "total": int(total), "page": int(request.args.get("page", 1)), "page_size": int(request.args.get("page_size", 20)), "auto_merge_enabled": bool( current_app.config.get( "DEVICE_ENTITY_AUTO_MERGE_ENABLED", False, ) ), } ) ), 200 except Exception as error: return _error(error) @bp.route("/device-entities/candidates/generate", methods=["POST"]) def generate_device_entity_candidates(): try: result = get_device_entity_resolution_service().generate( request.get_json(silent=True) or {}, actor_uid=_identity().get("id") or _identity().get("sub"), ) status = 201 if result.created_count else 200 return jsonify(success(_device_entity_generation(result))), status except Exception as error: return _error(error) @bp.route("/device-entities/candidates", methods=["POST"]) def submit_device_entity_candidate(): try: record = ( get_device_entity_resolution_service().submit_ai_candidate( request.get_json(silent=True) or {}, actor_uid=_identity().get("id") or _identity().get("sub"), ) ) return jsonify(success(_device_entity_candidate(record))), 201 except Exception as error: return _error(error) @bp.route("/device-entities/candidates/", methods=["GET"]) def get_device_entity_candidate(candidate_uid): try: record = get_device_entity_resolution_service().get(candidate_uid) return jsonify(success(_device_entity_candidate(record))), 200 except Exception as error: return _error(error) @bp.route( "/device-entities/candidates//review", methods=["POST"], ) def review_device_entity_candidate(candidate_uid): try: candidate, review, merge = ( get_device_entity_resolution_service().review( candidate_uid, request.get_json(silent=True) or {}, actor_uid=_identity().get("id") or _identity().get("sub"), ) ) return jsonify( success( { "candidate": _device_entity_candidate(candidate), "review": _device_entity_review(review), "merge": ( _device_entity_merge(merge) if merge is not None else None ), } ) ), 200 except Exception as error: return _error(error) @bp.route( "/device-entities/candidates//reviews", methods=["GET"], ) def list_device_entity_reviews(candidate_uid): try: records = get_device_entity_resolution_service().reviews( candidate_uid ) return jsonify( success( { "records": [ _device_entity_review(record) for record in records ], "total": len(records), } ) ), 200 except Exception as error: return _error(error) @bp.route( "/device-entities/candidates//merges", methods=["GET"], ) def list_device_entity_merges(candidate_uid): try: records = get_device_entity_resolution_service().merges( candidate_uid ) return jsonify( success( { "records": [ _device_entity_merge(record) for record in records ], "total": len(records), } ) ), 200 except Exception as error: return _error(error) @bp.route( "/device-entities/merges//rollback", methods=["POST"], ) def rollback_device_entity_merge(merge_uid): try: candidate, rollback = ( get_device_entity_resolution_service().rollback( merge_uid, request.get_json(silent=True) or {}, actor_uid=_identity().get("id") or _identity().get("sub"), ) ) return jsonify( success( { "candidate": _device_entity_candidate(candidate), "rollback": _device_entity_rollback(rollback), } ) ), 200 except Exception as error: return _error(error) @bp.route( "/device-entities/merges//rollbacks", methods=["GET"], ) def list_device_entity_rollbacks(merge_uid): try: records = get_device_entity_resolution_service().rollbacks( merge_uid ) return jsonify( success( { "records": [ _device_entity_rollback(record) for record in records ], "total": len(records), } ) ), 200 except Exception as error: return _error(error) @bp.route("/device-quality/profile", methods=["GET"]) def get_device_quality_profile(): try: result = get_device_quality_service().profile() return jsonify( success( { "profile_uid": result["profile_uid"], "name": result["name"], "latest_version": _device_quality_version( result["latest_version"] ), "active_version": _device_quality_version( result["active_version"] ), } ) ), 200 except Exception as error: return _error(error) @bp.route("/device-quality/bootstrap", methods=["POST"]) def bootstrap_device_quality(): try: record = get_device_quality_service().bootstrap( actor_uid=_identity().get("id") or _identity().get("sub"), ) return jsonify(success(_device_quality_version(record))), 201 except Exception as error: return _error(error) @bp.route("/device-quality/profile/versions", methods=["POST"]) def revise_device_quality_profile(): try: body = request.get_json(silent=True) or {} record = get_device_quality_service().revise( rules=body.get("rules"), expected_version=body.get("expected_version"), actor_uid=_identity().get("id") or _identity().get("sub"), ) return jsonify(success(_device_quality_version(record))), 201 except Exception as error: return _error(error) @bp.route("/device-quality/profile/versions", methods=["GET"]) def list_device_quality_versions(): try: records = get_device_quality_service().versions() return jsonify( success( { "records": [ _device_quality_version(record) for record in records ], "total": len(records), } ) ), 200 except Exception as error: return _error(error) @bp.route( "/device-quality/profile/versions//publish", methods=["POST"], ) def publish_device_quality_version(version_uid): try: record = get_device_quality_service().publish( version_uid, actor_uid=_identity().get("id") or _identity().get("sub"), ) return jsonify(success(_device_quality_version(record))), 200 except Exception as error: return _error(error) @bp.route("/device-quality/runs", methods=["POST"]) def run_device_quality(): try: body = request.get_json(silent=True) or {} record = get_device_quality_service().run( actor_uid=_identity().get("id") or _identity().get("sub"), source_uid=body.get("source_uid"), ) return jsonify(success(_device_quality_run(record))), 201 except Exception as error: return _error(error) @bp.route("/device-quality/runs", methods=["GET"]) def list_device_quality_runs(): try: page = request.args.get("page", 1) page_size = request.args.get("page_size", 20) records, total = get_device_quality_service().runs( page=page, page_size=page_size, ) return jsonify( success( { "records": [ _device_quality_run(record) for record in records ], "total": int(total), "page": int(page), "page_size": int(page_size), } ) ), 200 except Exception as error: return _error(error) @bp.route("/device-quality/runs/", methods=["GET"]) def get_device_quality_run(run_uid): try: record, results = get_device_quality_service().get_run(run_uid) return jsonify( success( { **_device_quality_run(record), "rule_results": [ _device_quality_rule_result(item) for item in results ], } ) ), 200 except Exception as error: return _error(error) @bp.route( "/device-quality/runs//violations", methods=["GET"], ) def list_device_quality_violations(run_uid): try: page = request.args.get("page", 1) page_size = request.args.get("page_size", 20) records, total = get_device_quality_service().violations( run_uid, rule_code=request.args.get("rule_code"), page=page, page_size=page_size, ) return jsonify( success( { "records": [ _device_quality_violation(record) for record in records ], "total": int(total), "page": int(page), "page_size": int(page_size), } ) ), 200 except Exception as error: return _error(error) @bp.route( "/device-quality/runs//asset-scores", methods=["GET"], ) def list_device_quality_asset_scores(run_uid): try: page = request.args.get("page", 1) page_size = request.args.get("page_size", 20) records, total = get_device_quality_service().asset_scores( run_uid, page=page, page_size=page_size, ) return jsonify( success( { "records": [ _device_quality_asset_score(record) for record in records ], "total": int(total), "page": int(page), "page_size": int(page_size), } ) ), 200 except Exception as error: return _error(error) @bp.route("/device-quality/issues", methods=["GET"]) def list_quality_issues(): try: page = request.args.get("page", 1) page_size = request.args.get("page_size", 20) overdue_only = str( request.args.get("overdue_only", "") ).strip().lower() in {"1", "true", "yes"} records, total = get_quality_issue_service().issues( status=request.args.get("status"), assignee_uid=request.args.get("assignee_uid"), overdue_only=overdue_only, page=page, page_size=page_size, ) return jsonify( success( { "records": [_quality_issue(item) for item in records], "total": int(total), "page": int(page), "page_size": int(page_size), } ) ), 200 except Exception as error: return _error(error) @bp.route("/device-quality/issues/statistics", methods=["GET"]) def get_quality_issue_statistics(): try: return jsonify( success(get_quality_issue_service().statistics()) ), 200 except Exception as error: return _error(error) @bp.route("/device-quality/issues/assignees", methods=["GET"]) def list_quality_issue_assignees(): try: records = get_quality_issue_service().assignees() return jsonify( success({"records": list(records), "total": len(records)}) ), 200 except Exception as error: return _error(error) @bp.route("/device-quality/issues/import", methods=["POST"]) def import_quality_issues(): try: body = request.get_json(silent=True) or {} result = get_quality_issue_service().import_violations( violation_uids=body.get("violation_uids"), priority=body.get("priority"), due_at=body.get("due_at"), actor_uid=_identity().get("id") or _identity().get("sub"), ) payload = { "records": [_quality_issue(item) for item in result.records], "created_count": int(result.created_count), "existing_count": int(result.existing_count), } status = 201 if result.created_count else 200 return jsonify(success(payload)), status except Exception as error: return _error(error) @bp.route("/device-quality/issues/", methods=["GET"]) def get_quality_issue(issue_uid): try: issue, remediation = get_quality_issue_service().get(issue_uid) return jsonify( success( { **_quality_issue(issue), "latest_remediation": _quality_issue_remediation( remediation ), } ) ), 200 except Exception as error: return _error(error) @bp.route( "/device-quality/issues//timeline", methods=["GET"], ) def get_quality_issue_timeline(issue_uid): try: records = get_quality_issue_service().timeline(issue_uid) return jsonify( success( { "records": [ _quality_issue_timeline(item) for item in records ], "total": len(records), } ) ), 200 except Exception as error: return _error(error) @bp.route( "/device-quality/issues//assign", methods=["POST"], ) def assign_quality_issue(issue_uid): try: body = request.get_json(silent=True) or {} record = get_quality_issue_service().assign( issue_uid, assignee_uid=body.get("assignee_uid"), due_at=body.get("due_at"), expected_version=body.get("expected_version"), actor_uid=_identity().get("id") or _identity().get("sub"), ) return jsonify(success(_quality_issue(record))), 200 except Exception as error: return _error(error) @bp.route( "/device-quality/issues//start", methods=["POST"], ) def start_quality_issue(issue_uid): try: body = request.get_json(silent=True) or {} record = get_quality_issue_service().start( issue_uid, expected_version=body.get("expected_version"), actor_uid=_identity().get("id") or _identity().get("sub"), ) return jsonify(success(_quality_issue(record))), 200 except Exception as error: return _error(error) @bp.route( "/device-quality/issues//submit", methods=["POST"], ) def submit_quality_issue(issue_uid): try: body = request.get_json(silent=True) or {} record = get_quality_issue_service().submit( issue_uid, summary=body.get("summary"), evidence_refs=body.get("evidence_refs", []), expected_version=body.get("expected_version"), actor_uid=_identity().get("id") or _identity().get("sub"), ) return jsonify(success(_quality_issue(record))), 200 except Exception as error: return _error(error) @bp.route( "/device-quality/issues//review", methods=["POST"], ) def review_quality_issue(issue_uid): try: body = request.get_json(silent=True) or {} record = get_quality_issue_service().review( issue_uid, verification_result=body.get("verification_result"), note=body.get("note"), verification_run_uid=body.get("verification_run_uid"), expected_version=body.get("expected_version"), actor_uid=_identity().get("id") or _identity().get("sub"), ) return jsonify(success(_quality_issue(record))), 200 except Exception as error: return _error(error) @bp.route( "/device-quality/issues//reopen", methods=["POST"], ) def reopen_quality_issue(issue_uid): try: body = request.get_json(silent=True) or {} record = get_quality_issue_service().reopen( issue_uid, reason=body.get("reason"), due_at=body.get("due_at"), expected_version=body.get("expected_version"), actor_uid=_identity().get("id") or _identity().get("sub"), ) return jsonify(success(_quality_issue(record))), 200 except Exception as error: return _error(error) @bp.route("/device-observability/events", methods=["GET"]) def list_device_operational_events(): filters = { name: request.args.get(name) for name in ( "event_type", "severity", "status", "asset_uid", "source_uid", "keyword", ) if request.args.get(name) } try: page = request.args.get("page", 1) page_size = request.args.get("page_size", 20) records, total = get_device_observability_service().events( filters, page=page, page_size=page_size, ) return jsonify( success( { "records": [ _device_operational_event(item) for item in records ], "total": int(total), "page": int(page), "page_size": int(page_size), } ) ), 200 except Exception as error: return _error(error) @bp.route("/device-observability/import", methods=["POST"]) def import_device_operational_evidence(): try: result = get_device_observability_service().import_evidence( request.get_json(silent=True) or {}, actor_uid=_identity().get("id") or _identity().get("sub"), ) data = { "events": [ _device_operational_event(item) for item in result.events ], "relations": [ _device_evidence_relation(item) for item in result.relations ], "created_events": int(result.created_events), "existing_events": int(result.existing_events), "created_relations": int(result.created_relations), "existing_relations": int(result.existing_relations), } status = ( 201 if result.created_events or result.created_relations else 200 ) return jsonify(success(data)), status except Exception as error: return _error(error) @bp.route("/device-observability/graph", methods=["GET"]) def get_device_evidence_graph(): try: record = get_device_observability_service().graph( request.args.get("anchor_kind"), request.args.get("anchor_uid"), max_hops=request.args.get("max_hops", 2), ) return jsonify(success(_device_evidence_graph(record))), 200 except Exception as error: return _error(error) @bp.route("/device-observability/root-cause", methods=["GET"]) def get_device_root_cause(): try: record = get_device_observability_service().root_cause( request.args.get("anchor_kind"), request.args.get("anchor_uid"), max_hops=request.args.get("max_hops", 3), ) return jsonify(success(_root_cause_analysis(record))), 200 except Exception as error: return _error(error) @bp.route("/governance-metrics/summary", methods=["GET"]) def get_governance_metric_summary(): try: result = get_governance_metrics_service().summary( _governance_metric_access() ) return jsonify(success(result)), 200 except Exception as error: return _error(error) @bp.route("/governance-metrics/details", methods=["GET"]) def get_governance_metric_details(): try: result = get_governance_metrics_service().details( _governance_metric_access(), metric=request.args.get("metric"), state=request.args.get("state"), page=request.args.get("page", 1), page_size=request.args.get("page_size", 20), ) return jsonify(success(result)), 200 except Exception as error: return _error(error) @bp.route("/data-elements", methods=["POST"]) def create_data_element(): try: record = get_data_element_service().create_draft( request.get_json(silent=True) or {}, actor_uid=_identity().get("id") or _identity().get("sub"), ) db.session.commit() return jsonify(success(_element(record))), 201 except Exception as error: db.session.rollback() return _error(error) @bp.route("/data-elements//transitions", methods=["POST"]) def transition_data_element(element_uid): payload = request.get_json(silent=True) or {} identity = _identity() target_status = str(payload.get("target_status") or "") if ( target_status in {"published", "deprecated", "retired"} and "data-elements:publish" not in set(identity.get("permissions") or []) ): return jsonify(failed("权限不足", code=403)), 403 try: record = get_data_element_service().transition( element_uid, target_status, expected_version=int(payload.get("expected_version")), actor_uid=identity.get("id") or identity.get("sub"), ) db.session.commit() return jsonify(success(_element(record))), 200 except Exception as error: db.session.rollback() return _error(error) @bp.route("/candidate-decisions", methods=["POST"]) def decide_candidates(): payload = request.get_json(silent=True) or {} try: records = get_candidate_decision_service().decide( payload.get("decisions") or [], actor_uid=_identity().get("id") or _identity().get("sub"), ) return jsonify(success([_decision(record) for record in records])), 200 except Exception as error: return _error(error) @bp.route("/sources/files", methods=["POST"]) def upload_source_file(): uploaded = request.files.get("file") if uploaded is None: return jsonify(failed("缺少上传文件", code=400)), 400 try: record, created = get_artifact_service().store( request.form.get("source_uid"), uploaded.filename, uploaded.mimetype, uploaded.read(), request.form.get("parser_version") or "auto-v1", ) db.session.commit() return jsonify(success(_artifact(record))), 201 if created else 200 except Exception as error: db.session.rollback() return _error(error) @bp.route("/evidence/", methods=["GET"]) def get_evidence(evidence_uid): identity = _identity() try: service = get_evidence_service() if request.args.get("download") in {"1", "true", "yes"}: if "evidence:download" not in set(identity.get("permissions") or []): return jsonify(failed("权限不足", code=403)), 403 content, filename, media_type = service.download(evidence_uid) return send_file( io.BytesIO(content), mimetype=media_type, as_attachment=True, download_name=filename, ) preview = service.preview(evidence_uid) from app.core.data_research.artifacts import redact_excerpt preview["excerpt"] = redact_excerpt(preview.get("excerpt")) return jsonify(success(preview)), 200 except Exception as error: return _error(error) @bp.route("/ontologies", methods=["GET"]) def list_ontologies(): try: return jsonify(success([_ontology(item) for item in get_ontology_service().list()])), 200 except Exception as error: return _error(error) @bp.route("/ontologies", methods=["POST"]) def create_ontology(): try: record = get_ontology_service().create( request.get_json(silent=True) or {}, _identity().get("id") or _identity().get("sub"), ) return jsonify(success(_ontology(record))), 201 except Exception as error: return _error(error) @bp.route("/ontologies/", methods=["GET"]) def get_ontology(ontology_uid): try: return jsonify(success(_ontology(get_ontology_service().get(ontology_uid)))), 200 except Exception as error: return _error(error) @bp.route("/ontologies//versions", methods=["GET"]) def list_ontology_versions(ontology_uid): try: return jsonify( success( [ _ontology_version(item) for item in get_ontology_service().list_versions(ontology_uid) ] ) ), 200 except Exception as error: return _error(error) @bp.route("/ontologies//graph", methods=["GET"]) def get_ontology_graph(ontology_uid): try: service = get_ontology_service() ontology = service.get(ontology_uid) version = service.latest_version(ontology_uid) data = ( _ontology_version(version) if version is not None else { "uid": None, "ontology_uid": str(ontology_uid), "version": 0, "parent_version_uid": None, "status": "draft", "content_hash": None, "graph_document": { name: [] for name in ( "classes", "properties", "relations", "constraints", "domain_links", "element_mappings", ) }, "created_by": None, } ) response = jsonify(success(data)) response.headers["ETag"] = f'"{int(ontology.draft_revision)}"' return response, 200 except Exception as error: return _error(error) @bp.route("/ontologies//graph", methods=["PATCH"]) def save_ontology_graph(ontology_uid): etag = str(request.headers.get("If-Match") or "").strip().strip('"') if not etag.isdigit(): return jsonify(failed("缺少有效 If-Match 修订号", code=428)), 428 try: version = get_ontology_service().save_draft( ontology_uid, request.get_json(silent=True) or {}, int(etag), _identity().get("id") or _identity().get("sub"), ) response = jsonify(success(_ontology_version(version))) response.headers["ETag"] = f'"{int(etag) + 1}"' return response, 200 except Exception as error: return _error(error) @bp.route("/ontologies//validate", methods=["POST"]) def validate_ontology(ontology_uid): try: issues = get_ontology_service().validate(ontology_uid) data = [ {"code": item.code, "message": item.message, "path": item.path} for item in issues ] return jsonify(success(data)), 200 except Exception as error: return _error(error) @bp.route("/ontologies//publish", methods=["POST"]) def publish_ontology(ontology_uid): try: version = get_ontology_service().publish( ontology_uid, request.headers.get("Idempotency-Key") or f"{ontology_uid}:publish", _identity().get("id") or _identity().get("sub"), ) return jsonify(success(_ontology_version(version))), 200 except Exception as error: return _error(error) @bp.route("/ontologies//diff", methods=["GET"]) def diff_ontology(ontology_uid): try: return jsonify(success(get_ontology_service().diff( ontology_uid, request.args.get("left"), request.args.get("right") ))), 200 except Exception as error: return _error(error) @bp.route("/ontologies//rollback", methods=["POST"]) def rollback_ontology(ontology_uid): payload = request.get_json(silent=True) or {} try: version = get_ontology_service().rollback( ontology_uid, payload.get("target_version_uid"), int(payload.get("expected_revision")), _identity().get("id") or _identity().get("sub"), ) return jsonify(success(_ontology_version(version))), 201 except Exception as error: return _error(error) @bp.route("/ontologies//suggestions", methods=["POST"]) def generate_ontology_suggestions(ontology_uid): try: result = get_ontology_dynamic_service().generate( ontology_uid, request.get_json(silent=True) or {}, _identity().get("id") or _identity().get("sub"), ) records = ( result.get("suggestions") or [] if isinstance(result, dict) else result ) suggestions = [ item if isinstance(item, dict) else { "uid": item.uid, "kind": item.kind, "payload": item.payload, "evidence_uids": list(item.evidence_uids), "confidence": item.confidence, "source": item.source, "model_version": item.model_version, "prompt_version": item.prompt_version, } for item in records ] data = { "change_set_uid": ( result.get("change_set_uid") if isinstance(result, dict) else None ), "suggestions": suggestions, } return jsonify(success(data)), 201 except Exception as error: return _error(error) @bp.route( "/ontologies//change-sets//decisions", methods=["POST"], ) def decide_ontology_change_set(ontology_uid, change_set_uid): payload = request.get_json(silent=True) or {} try: result = get_ontology_dynamic_service().decide( ontology_uid, change_set_uid, payload.get("decisions") or [], _identity().get("id") or _identity().get("sub"), ) return jsonify(success(result)), 200 except Exception as error: return _error(error) @bp.route("/ontologies//export", methods=["GET"]) def export_ontology(ontology_uid): try: content, media_type, filename = get_ontology_exchange_service().export( ontology_uid, request.args.get("format") or "json" ) return send_file( io.BytesIO(content), mimetype=media_type, as_attachment=True, download_name=filename, ) except Exception as error: return _error(error) @bp.route("/ontologies/import", methods=["POST"]) def import_ontology(): try: uploaded = request.files.get("file") content = uploaded.read() if uploaded is not None else request.get_data(cache=False) result = get_ontology_exchange_service().import_document( content, request.args.get("format") or "json", _identity().get("id") or _identity().get("sub"), ) return jsonify(success(result)), 201 except Exception as error: return _error(error) @bp.route("/semantic/properties/", methods=["GET"]) def query_semantic_property(property_uid): domain = str(request.args.get("business_domain_uid") or "").strip() identity = _identity() scoped = set(identity.get("business_domains") or [domain]) if "*" not in scoped and domain not in scoped: return jsonify(failed("权限不足", code=403)), 403 try: result = get_semantic_query_service().trace_property( property_uid, allowed_domains={domain}, limit=int(request.args.get("limit") or 20), after_uid=request.args.get("after_uid"), ) return jsonify(success(result)), 200 except Exception as error: return _error(error)