"""HTTP orchestration boundary for data-research ingestion jobs.""" from __future__ import annotations import io import logging from flask import current_app, g, jsonify, request, send_file from app import db from app.api.data_development import bp 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_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) publication = OntologyPublicationService( repository, outbox_enqueue=lambda **event: enqueue_outbox(db.session, **event), 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 _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 _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("/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)