routes.py 19 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597
  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 get_ontology_service():
  68. from app.core.data_research.ontology.publication import (
  69. OntologyApplicationService,
  70. OntologyPublicationService,
  71. )
  72. from app.core.data_research.ontology.repository import SqlAlchemyOntologyRepository
  73. from app.core.events.outbox import enqueue_outbox
  74. repository = SqlAlchemyOntologyRepository(db.session)
  75. publication = OntologyPublicationService(
  76. repository,
  77. outbox_enqueue=lambda **event: enqueue_outbox(db.session, **event),
  78. commit=db.session.commit,
  79. rollback=db.session.rollback,
  80. )
  81. return OntologyApplicationService(
  82. repository,
  83. publication=publication,
  84. commit=db.session.commit,
  85. rollback=db.session.rollback,
  86. )
  87. def get_ontology_dynamic_service():
  88. from app.core.data_research.ontology.change_sets import (
  89. SqlAlchemyDynamicOntologyService,
  90. )
  91. return SqlAlchemyDynamicOntologyService(db.session)
  92. def get_ontology_exchange_service():
  93. from app.core.data_research.ontology.exchange import OntologyExchangeService
  94. from app.core.data_research.ontology.repository import SqlAlchemyOntologyRepository
  95. return OntologyExchangeService(
  96. SqlAlchemyOntologyRepository(db.session),
  97. commit=db.session.commit,
  98. rollback=db.session.rollback,
  99. )
  100. def get_semantic_query_service():
  101. from app.core.data_research.ontology.query import (
  102. Neo4jSemanticRepository,
  103. SemanticQueryService,
  104. )
  105. from app.services.neo4j_driver import neo4j_driver
  106. return SemanticQueryService(Neo4jSemanticRepository(neo4j_driver))
  107. def _identity():
  108. return getattr(g, "current_user", {}) or {}
  109. def _record(record):
  110. return {
  111. "uid": str(record.uid),
  112. "source_uid": str(record.source_uid),
  113. "artifact_uid": str(record.artifact_uid) if record.artifact_uid else None,
  114. "job_type": record.job_type,
  115. "parser_version": record.parser_version,
  116. "status": record.status,
  117. "parameters": dict(record.parameters or {}),
  118. "statistics": dict(record.statistics or {}),
  119. "last_error": record.last_error,
  120. "actor_uid": record.actor_uid,
  121. "created_at": record.created_at.isoformat() if record.created_at else None,
  122. "updated_at": record.updated_at.isoformat() if record.updated_at else None,
  123. "started_at": record.started_at.isoformat() if record.started_at else None,
  124. "finished_at": record.finished_at.isoformat() if record.finished_at else None,
  125. }
  126. def _element(record):
  127. return {
  128. "uid": str(record.uid),
  129. "code": record.code,
  130. "status": record.status,
  131. "current_version": int(record.current_version),
  132. "snapshot": dict(record.snapshot or {}),
  133. "created_by": record.created_by,
  134. "updated_by": record.updated_by,
  135. }
  136. def _decision(record):
  137. return {
  138. "uid": str(record.uid),
  139. "candidate_uid": str(record.candidate_uid),
  140. "action": record.action,
  141. "data_element_uid": (
  142. str(record.data_element_uid) if record.data_element_uid else None
  143. ),
  144. "evidence_uids": list(record.evidence_uids),
  145. "actor_uid": record.actor_uid,
  146. "reason": record.reason,
  147. }
  148. def _artifact(record):
  149. return {
  150. "uid": str(record.uid),
  151. "source_uid": str(record.source_uid),
  152. "filename": record.filename,
  153. "media_type": record.media_type,
  154. "size_bytes": int(record.size_bytes),
  155. "content_hash": record.content_hash,
  156. "parser_version": record.parser_version,
  157. }
  158. def _ontology(record):
  159. return {
  160. "uid": str(record.uid),
  161. "code": record.code,
  162. "name": record.name,
  163. "owner_uid": record.owner_uid,
  164. "status": record.status,
  165. "draft_revision": int(record.draft_revision),
  166. "active_version_uid": record.active_version_uid,
  167. "domain_links": [
  168. {"domain_uid": link.domain_uid, "role": link.role}
  169. for link in record.domain_links
  170. ],
  171. }
  172. def _ontology_version(record):
  173. return {
  174. "uid": str(record.uid),
  175. "ontology_uid": str(record.ontology_uid),
  176. "version": int(record.version),
  177. "parent_version_uid": record.parent_version_uid,
  178. "status": record.status,
  179. "content_hash": record.content_hash,
  180. "graph_document": record.graph_document.to_dict(),
  181. "created_by": record.created_by,
  182. }
  183. def _error(error):
  184. if isinstance(error, DataResearchError):
  185. return (
  186. jsonify(
  187. failed(
  188. str(error),
  189. code=error.http_status,
  190. error={"code": error.code},
  191. )
  192. ),
  193. error.http_status,
  194. )
  195. logger.exception("data-research ingestion request failed")
  196. return (
  197. jsonify(
  198. failed(
  199. "数据采集任务处理失败",
  200. code=500,
  201. error={"code": "DATA_RESEARCH_ERROR"},
  202. )
  203. ),
  204. 500,
  205. )
  206. @bp.route("/ingestion-jobs", methods=["POST"])
  207. def create_ingestion_job():
  208. payload = request.get_json(silent=True) or {}
  209. try:
  210. record, created = get_ingestion_service().create_job(
  211. payload,
  212. actor_uid=_identity().get("id") or _identity().get("sub"),
  213. )
  214. return jsonify(success(_record(record))), 201 if created else 200
  215. except Exception as error:
  216. return _error(error)
  217. @bp.route("/ingestion-jobs", methods=["GET"])
  218. def list_ingestion_jobs():
  219. filters = {
  220. name: request.args.get(name)
  221. for name in ("status", "source_uid")
  222. if request.args.get(name)
  223. }
  224. try:
  225. records = get_ingestion_service().list_jobs(filters)
  226. return jsonify(
  227. success({"records": [_record(item) for item in records], "total": len(records)})
  228. ), 200
  229. except Exception as error:
  230. return _error(error)
  231. @bp.route("/ingestion-jobs/<job_uid>", methods=["GET"])
  232. def get_ingestion_job(job_uid):
  233. try:
  234. return jsonify(success(_record(get_ingestion_service().get_job(job_uid)))), 200
  235. except Exception as error:
  236. return _error(error)
  237. @bp.route("/ingestion-jobs/<job_uid>/retry", methods=["POST"])
  238. def retry_ingestion_job(job_uid):
  239. try:
  240. return jsonify(success(_record(get_ingestion_service().retry(job_uid)))), 200
  241. except Exception as error:
  242. return _error(error)
  243. @bp.route("/ingestion-jobs/<job_uid>/cancel", methods=["POST"])
  244. def cancel_ingestion_job(job_uid):
  245. try:
  246. service = get_ingestion_service()
  247. record = service.get_job(job_uid)
  248. identity = _identity()
  249. permissions = set(identity.get("permissions") or [])
  250. actor_uid = identity.get("id") or identity.get("sub")
  251. if record.actor_uid != actor_uid and "ingestion:admin" not in permissions:
  252. return jsonify(failed("权限不足", code=403)), 403
  253. return jsonify(success(_record(service.cancel(job_uid)))), 200
  254. except Exception as error:
  255. return _error(error)
  256. @bp.route("/data-elements", methods=["POST"])
  257. def create_data_element():
  258. try:
  259. record = get_data_element_service().create_draft(
  260. request.get_json(silent=True) or {},
  261. actor_uid=_identity().get("id") or _identity().get("sub"),
  262. )
  263. db.session.commit()
  264. return jsonify(success(_element(record))), 201
  265. except Exception as error:
  266. db.session.rollback()
  267. return _error(error)
  268. @bp.route("/data-elements/<element_uid>/transitions", methods=["POST"])
  269. def transition_data_element(element_uid):
  270. payload = request.get_json(silent=True) or {}
  271. identity = _identity()
  272. target_status = str(payload.get("target_status") or "")
  273. if (
  274. target_status in {"published", "deprecated", "retired"}
  275. and "data-elements:publish" not in set(identity.get("permissions") or [])
  276. ):
  277. return jsonify(failed("权限不足", code=403)), 403
  278. try:
  279. record = get_data_element_service().transition(
  280. element_uid,
  281. target_status,
  282. expected_version=int(payload.get("expected_version")),
  283. actor_uid=identity.get("id") or identity.get("sub"),
  284. )
  285. db.session.commit()
  286. return jsonify(success(_element(record))), 200
  287. except Exception as error:
  288. db.session.rollback()
  289. return _error(error)
  290. @bp.route("/candidate-decisions", methods=["POST"])
  291. def decide_candidates():
  292. payload = request.get_json(silent=True) or {}
  293. try:
  294. records = get_candidate_decision_service().decide(
  295. payload.get("decisions") or [],
  296. actor_uid=_identity().get("id") or _identity().get("sub"),
  297. )
  298. return jsonify(success([_decision(record) for record in records])), 200
  299. except Exception as error:
  300. return _error(error)
  301. @bp.route("/sources/files", methods=["POST"])
  302. def upload_source_file():
  303. uploaded = request.files.get("file")
  304. if uploaded is None:
  305. return jsonify(failed("缺少上传文件", code=400)), 400
  306. try:
  307. record, created = get_artifact_service().store(
  308. request.form.get("source_uid"),
  309. uploaded.filename,
  310. uploaded.mimetype,
  311. uploaded.read(),
  312. request.form.get("parser_version") or "auto-v1",
  313. )
  314. db.session.commit()
  315. return jsonify(success(_artifact(record))), 201 if created else 200
  316. except Exception as error:
  317. db.session.rollback()
  318. return _error(error)
  319. @bp.route("/evidence/<evidence_uid>", methods=["GET"])
  320. def get_evidence(evidence_uid):
  321. identity = _identity()
  322. try:
  323. service = get_evidence_service()
  324. if request.args.get("download") in {"1", "true", "yes"}:
  325. if "evidence:download" not in set(identity.get("permissions") or []):
  326. return jsonify(failed("权限不足", code=403)), 403
  327. content, filename, media_type = service.download(evidence_uid)
  328. return send_file(
  329. io.BytesIO(content),
  330. mimetype=media_type,
  331. as_attachment=True,
  332. download_name=filename,
  333. )
  334. preview = service.preview(evidence_uid)
  335. from app.core.data_research.artifacts import redact_excerpt
  336. preview["excerpt"] = redact_excerpt(preview.get("excerpt"))
  337. return jsonify(success(preview)), 200
  338. except Exception as error:
  339. return _error(error)
  340. @bp.route("/ontologies", methods=["GET"])
  341. def list_ontologies():
  342. try:
  343. return jsonify(success([_ontology(item) for item in get_ontology_service().list()])), 200
  344. except Exception as error:
  345. return _error(error)
  346. @bp.route("/ontologies", methods=["POST"])
  347. def create_ontology():
  348. try:
  349. record = get_ontology_service().create(
  350. request.get_json(silent=True) or {},
  351. _identity().get("id") or _identity().get("sub"),
  352. )
  353. return jsonify(success(_ontology(record))), 201
  354. except Exception as error:
  355. return _error(error)
  356. @bp.route("/ontologies/<ontology_uid>/graph", methods=["PATCH"])
  357. def save_ontology_graph(ontology_uid):
  358. etag = str(request.headers.get("If-Match") or "").strip().strip('"')
  359. if not etag.isdigit():
  360. return jsonify(failed("缺少有效 If-Match 修订号", code=428)), 428
  361. try:
  362. version = get_ontology_service().save_draft(
  363. ontology_uid,
  364. request.get_json(silent=True) or {},
  365. int(etag),
  366. _identity().get("id") or _identity().get("sub"),
  367. )
  368. response = jsonify(success(_ontology_version(version)))
  369. response.headers["ETag"] = f'"{int(etag) + 1}"'
  370. return response, 200
  371. except Exception as error:
  372. return _error(error)
  373. @bp.route("/ontologies/<ontology_uid>/validate", methods=["POST"])
  374. def validate_ontology(ontology_uid):
  375. try:
  376. issues = get_ontology_service().validate(ontology_uid)
  377. data = [
  378. {"code": item.code, "message": item.message, "path": item.path}
  379. for item in issues
  380. ]
  381. return jsonify(success(data)), 422 if data else 200
  382. except Exception as error:
  383. return _error(error)
  384. @bp.route("/ontologies/<ontology_uid>/publish", methods=["POST"])
  385. def publish_ontology(ontology_uid):
  386. try:
  387. version = get_ontology_service().publish(
  388. ontology_uid,
  389. request.headers.get("Idempotency-Key") or f"{ontology_uid}:publish",
  390. _identity().get("id") or _identity().get("sub"),
  391. )
  392. return jsonify(success(_ontology_version(version))), 200
  393. except Exception as error:
  394. return _error(error)
  395. @bp.route("/ontologies/<ontology_uid>/diff", methods=["GET"])
  396. def diff_ontology(ontology_uid):
  397. try:
  398. return jsonify(success(get_ontology_service().diff(
  399. ontology_uid, request.args.get("left"), request.args.get("right")
  400. ))), 200
  401. except Exception as error:
  402. return _error(error)
  403. @bp.route("/ontologies/<ontology_uid>/rollback", methods=["POST"])
  404. def rollback_ontology(ontology_uid):
  405. payload = request.get_json(silent=True) or {}
  406. try:
  407. version = get_ontology_service().rollback(
  408. ontology_uid,
  409. payload.get("target_version_uid"),
  410. int(payload.get("expected_revision")),
  411. _identity().get("id") or _identity().get("sub"),
  412. )
  413. return jsonify(success(_ontology_version(version))), 201
  414. except Exception as error:
  415. return _error(error)
  416. @bp.route("/ontologies/<ontology_uid>/suggestions", methods=["POST"])
  417. def generate_ontology_suggestions(ontology_uid):
  418. try:
  419. records = get_ontology_dynamic_service().generate(
  420. ontology_uid,
  421. request.get_json(silent=True) or {},
  422. _identity().get("id") or _identity().get("sub"),
  423. )
  424. data = [
  425. {
  426. "uid": item.uid,
  427. "kind": item.kind,
  428. "payload": item.payload,
  429. "evidence_uids": list(item.evidence_uids),
  430. "confidence": item.confidence,
  431. "source": item.source,
  432. "model_version": item.model_version,
  433. "prompt_version": item.prompt_version,
  434. }
  435. if not isinstance(item, dict)
  436. else item
  437. for item in records
  438. ]
  439. return jsonify(success(data)), 201
  440. except Exception as error:
  441. return _error(error)
  442. @bp.route(
  443. "/ontologies/<ontology_uid>/change-sets/<change_set_uid>/decisions",
  444. methods=["POST"],
  445. )
  446. def decide_ontology_change_set(ontology_uid, change_set_uid):
  447. payload = request.get_json(silent=True) or {}
  448. try:
  449. result = get_ontology_dynamic_service().decide(
  450. ontology_uid,
  451. change_set_uid,
  452. payload.get("decisions") or [],
  453. _identity().get("id") or _identity().get("sub"),
  454. )
  455. return jsonify(success(result)), 200
  456. except Exception as error:
  457. return _error(error)
  458. @bp.route("/ontologies/<ontology_uid>/export", methods=["GET"])
  459. def export_ontology(ontology_uid):
  460. try:
  461. content, media_type, filename = get_ontology_exchange_service().export(
  462. ontology_uid, request.args.get("format") or "json"
  463. )
  464. return send_file(
  465. io.BytesIO(content),
  466. mimetype=media_type,
  467. as_attachment=True,
  468. download_name=filename,
  469. )
  470. except Exception as error:
  471. return _error(error)
  472. @bp.route("/ontologies/import", methods=["POST"])
  473. def import_ontology():
  474. try:
  475. result = get_ontology_exchange_service().import_document(
  476. request.get_data(cache=False),
  477. request.args.get("format") or "json",
  478. _identity().get("id") or _identity().get("sub"),
  479. )
  480. return jsonify(success(result)), 201
  481. except Exception as error:
  482. return _error(error)
  483. @bp.route("/semantic/properties/<property_uid>", methods=["GET"])
  484. def query_semantic_property(property_uid):
  485. domain = str(request.args.get("business_domain_uid") or "").strip()
  486. identity = _identity()
  487. scoped = set(identity.get("business_domains") or [domain])
  488. if "*" not in scoped and domain not in scoped:
  489. return jsonify(failed("权限不足", code=403)), 403
  490. try:
  491. result = get_semantic_query_service().trace_property(
  492. property_uid,
  493. allowed_domains={domain},
  494. limit=int(request.args.get("limit") or 20),
  495. after_uid=request.args.get("after_uid"),
  496. )
  497. return jsonify(success(result)), 200
  498. except Exception as error:
  499. return _error(error)