routes.py 11 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332
  1. """HTTP orchestration boundary for data-research ingestion jobs."""
  2. from __future__ import annotations
  3. import io
  4. import logging
  5. from flask import current_app, g, jsonify, request, send_file
  6. from app import db
  7. from app.api.data_development import bp
  8. from app.core.data_research.errors import DataResearchError
  9. from app.models.result import failed, success
  10. logger = logging.getLogger(__name__)
  11. def get_ingestion_service():
  12. from app.core.data_research.ingestion import IngestionService
  13. from app.core.data_research.repository import SqlAlchemyIngestionJobRepository
  14. return IngestionService(
  15. SqlAlchemyIngestionJobRepository(db.session),
  16. commit=db.session.commit,
  17. rollback=db.session.rollback,
  18. )
  19. def get_data_element_service():
  20. from app.core.data_research.data_elements import DataElementService
  21. from app.core.data_research.repository import SqlAlchemyDataElementRepository
  22. from app.core.events.outbox import enqueue_outbox
  23. return DataElementService(
  24. SqlAlchemyDataElementRepository(db.session),
  25. outbox_enqueue=lambda **event: enqueue_outbox(db.session, **event),
  26. )
  27. def get_candidate_decision_service():
  28. from app.core.data_research.candidate_decisions import CandidateDecisionService
  29. from app.core.data_research.data_elements import DataElementService
  30. from app.core.data_research.repository import (
  31. SqlAlchemyCandidateDecisionRepository,
  32. SqlAlchemyDataElementRepository,
  33. )
  34. return CandidateDecisionService(
  35. SqlAlchemyCandidateDecisionRepository(db.session),
  36. data_elements=DataElementService(SqlAlchemyDataElementRepository(db.session)),
  37. )
  38. def _artifact_storage():
  39. from minio import Minio
  40. from app.core.data_research.artifacts import MinioArtifactStorage
  41. client = Minio(
  42. current_app.config["MINIO_HOST"],
  43. access_key=current_app.config["MINIO_USER"],
  44. secret_key=current_app.config["MINIO_PASSWORD"],
  45. secure=bool(current_app.config.get("MINIO_SECURE")),
  46. )
  47. return MinioArtifactStorage(client, current_app.config["MINIO_BUCKET"])
  48. def get_artifact_service():
  49. from app.core.common.identifiers import new_governance_uid
  50. from app.core.data_research.artifacts import (
  51. ArtifactService,
  52. SqlAlchemyArtifactRepository,
  53. )
  54. from app.core.data_research.file_policy import FilePolicy
  55. return ArtifactService(
  56. SqlAlchemyArtifactRepository(db.session),
  57. _artifact_storage(),
  58. FilePolicy(
  59. max_bytes=int(current_app.config.get("DATA_RESEARCH_FILE_MAX_BYTES", 25 * 1024 * 1024)),
  60. max_pages=int(current_app.config.get("DATA_RESEARCH_FILE_MAX_PAGES", 500)),
  61. ),
  62. uid_factory=new_governance_uid,
  63. )
  64. def get_evidence_service():
  65. from app.core.data_research.artifacts import EvidenceService
  66. return EvidenceService(db.session, _artifact_storage())
  67. def _identity():
  68. return getattr(g, "current_user", {}) or {}
  69. def _record(record):
  70. return {
  71. "uid": str(record.uid),
  72. "source_uid": str(record.source_uid),
  73. "artifact_uid": str(record.artifact_uid) if record.artifact_uid else None,
  74. "job_type": record.job_type,
  75. "parser_version": record.parser_version,
  76. "status": record.status,
  77. "parameters": dict(record.parameters or {}),
  78. "statistics": dict(record.statistics or {}),
  79. "last_error": record.last_error,
  80. "actor_uid": record.actor_uid,
  81. "created_at": record.created_at.isoformat() if record.created_at else None,
  82. "updated_at": record.updated_at.isoformat() if record.updated_at else None,
  83. "started_at": record.started_at.isoformat() if record.started_at else None,
  84. "finished_at": record.finished_at.isoformat() if record.finished_at else None,
  85. }
  86. def _element(record):
  87. return {
  88. "uid": str(record.uid),
  89. "code": record.code,
  90. "status": record.status,
  91. "current_version": int(record.current_version),
  92. "snapshot": dict(record.snapshot or {}),
  93. "created_by": record.created_by,
  94. "updated_by": record.updated_by,
  95. }
  96. def _decision(record):
  97. return {
  98. "uid": str(record.uid),
  99. "candidate_uid": str(record.candidate_uid),
  100. "action": record.action,
  101. "data_element_uid": (
  102. str(record.data_element_uid) if record.data_element_uid else None
  103. ),
  104. "evidence_uids": list(record.evidence_uids),
  105. "actor_uid": record.actor_uid,
  106. "reason": record.reason,
  107. }
  108. def _artifact(record):
  109. return {
  110. "uid": str(record.uid),
  111. "source_uid": str(record.source_uid),
  112. "filename": record.filename,
  113. "media_type": record.media_type,
  114. "size_bytes": int(record.size_bytes),
  115. "content_hash": record.content_hash,
  116. "parser_version": record.parser_version,
  117. }
  118. def _error(error):
  119. if isinstance(error, DataResearchError):
  120. return (
  121. jsonify(
  122. failed(
  123. str(error),
  124. code=error.http_status,
  125. error={"code": error.code},
  126. )
  127. ),
  128. error.http_status,
  129. )
  130. logger.exception("data-research ingestion request failed")
  131. return (
  132. jsonify(
  133. failed(
  134. "数据采集任务处理失败",
  135. code=500,
  136. error={"code": "DATA_RESEARCH_ERROR"},
  137. )
  138. ),
  139. 500,
  140. )
  141. @bp.route("/ingestion-jobs", methods=["POST"])
  142. def create_ingestion_job():
  143. payload = request.get_json(silent=True) or {}
  144. try:
  145. record, created = get_ingestion_service().create_job(
  146. payload,
  147. actor_uid=_identity().get("id") or _identity().get("sub"),
  148. )
  149. return jsonify(success(_record(record))), 201 if created else 200
  150. except Exception as error:
  151. return _error(error)
  152. @bp.route("/ingestion-jobs", methods=["GET"])
  153. def list_ingestion_jobs():
  154. filters = {
  155. name: request.args.get(name)
  156. for name in ("status", "source_uid")
  157. if request.args.get(name)
  158. }
  159. try:
  160. records = get_ingestion_service().list_jobs(filters)
  161. return jsonify(
  162. success({"records": [_record(item) for item in records], "total": len(records)})
  163. ), 200
  164. except Exception as error:
  165. return _error(error)
  166. @bp.route("/ingestion-jobs/<job_uid>", methods=["GET"])
  167. def get_ingestion_job(job_uid):
  168. try:
  169. return jsonify(success(_record(get_ingestion_service().get_job(job_uid)))), 200
  170. except Exception as error:
  171. return _error(error)
  172. @bp.route("/ingestion-jobs/<job_uid>/retry", methods=["POST"])
  173. def retry_ingestion_job(job_uid):
  174. try:
  175. return jsonify(success(_record(get_ingestion_service().retry(job_uid)))), 200
  176. except Exception as error:
  177. return _error(error)
  178. @bp.route("/ingestion-jobs/<job_uid>/cancel", methods=["POST"])
  179. def cancel_ingestion_job(job_uid):
  180. try:
  181. service = get_ingestion_service()
  182. record = service.get_job(job_uid)
  183. identity = _identity()
  184. permissions = set(identity.get("permissions") or [])
  185. actor_uid = identity.get("id") or identity.get("sub")
  186. if record.actor_uid != actor_uid and "ingestion:admin" not in permissions:
  187. return jsonify(failed("权限不足", code=403)), 403
  188. return jsonify(success(_record(service.cancel(job_uid)))), 200
  189. except Exception as error:
  190. return _error(error)
  191. @bp.route("/data-elements", methods=["POST"])
  192. def create_data_element():
  193. try:
  194. record = get_data_element_service().create_draft(
  195. request.get_json(silent=True) or {},
  196. actor_uid=_identity().get("id") or _identity().get("sub"),
  197. )
  198. db.session.commit()
  199. return jsonify(success(_element(record))), 201
  200. except Exception as error:
  201. db.session.rollback()
  202. return _error(error)
  203. @bp.route("/data-elements/<element_uid>/transitions", methods=["POST"])
  204. def transition_data_element(element_uid):
  205. payload = request.get_json(silent=True) or {}
  206. identity = _identity()
  207. target_status = str(payload.get("target_status") or "")
  208. if (
  209. target_status in {"published", "deprecated", "retired"}
  210. and "data-elements:publish" not in set(identity.get("permissions") or [])
  211. ):
  212. return jsonify(failed("权限不足", code=403)), 403
  213. try:
  214. record = get_data_element_service().transition(
  215. element_uid,
  216. target_status,
  217. expected_version=int(payload.get("expected_version")),
  218. actor_uid=identity.get("id") or identity.get("sub"),
  219. )
  220. db.session.commit()
  221. return jsonify(success(_element(record))), 200
  222. except Exception as error:
  223. db.session.rollback()
  224. return _error(error)
  225. @bp.route("/candidate-decisions", methods=["POST"])
  226. def decide_candidates():
  227. payload = request.get_json(silent=True) or {}
  228. try:
  229. records = get_candidate_decision_service().decide(
  230. payload.get("decisions") or [],
  231. actor_uid=_identity().get("id") or _identity().get("sub"),
  232. )
  233. return jsonify(success([_decision(record) for record in records])), 200
  234. except Exception as error:
  235. return _error(error)
  236. @bp.route("/sources/files", methods=["POST"])
  237. def upload_source_file():
  238. uploaded = request.files.get("file")
  239. if uploaded is None:
  240. return jsonify(failed("缺少上传文件", code=400)), 400
  241. try:
  242. record, created = get_artifact_service().store(
  243. request.form.get("source_uid"),
  244. uploaded.filename,
  245. uploaded.mimetype,
  246. uploaded.read(),
  247. request.form.get("parser_version") or "auto-v1",
  248. )
  249. db.session.commit()
  250. return jsonify(success(_artifact(record))), 201 if created else 200
  251. except Exception as error:
  252. db.session.rollback()
  253. return _error(error)
  254. @bp.route("/evidence/<evidence_uid>", methods=["GET"])
  255. def get_evidence(evidence_uid):
  256. identity = _identity()
  257. try:
  258. service = get_evidence_service()
  259. if request.args.get("download") in {"1", "true", "yes"}:
  260. if "evidence:download" not in set(identity.get("permissions") or []):
  261. return jsonify(failed("权限不足", code=403)), 403
  262. content, filename, media_type = service.download(evidence_uid)
  263. return send_file(
  264. io.BytesIO(content),
  265. mimetype=media_type,
  266. as_attachment=True,
  267. download_name=filename,
  268. )
  269. preview = service.preview(evidence_uid)
  270. from app.core.data_research.artifacts import redact_excerpt
  271. preview["excerpt"] = redact_excerpt(preview.get("excerpt"))
  272. return jsonify(success(preview)), 200
  273. except Exception as error:
  274. return _error(error)