"""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_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 _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, "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 _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 _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: record, created = get_ingestion_service().create_job( payload, actor_uid=_identity().get("id") or _identity().get("sub"), ) 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//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=["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)