routes.py 22 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691
  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=["GET"])
  257. def list_data_elements():
  258. try:
  259. records = get_data_element_service().list(
  260. status=request.args.get("status"),
  261. business_domain_uid=request.args.get("business_domain_uid"),
  262. )
  263. return jsonify(
  264. success(
  265. {
  266. "records": [_element(item) for item in records],
  267. "total": len(records),
  268. }
  269. )
  270. ), 200
  271. except Exception as error:
  272. return _error(error)
  273. @bp.route("/data-elements", methods=["POST"])
  274. def create_data_element():
  275. try:
  276. record = get_data_element_service().create_draft(
  277. request.get_json(silent=True) or {},
  278. actor_uid=_identity().get("id") or _identity().get("sub"),
  279. )
  280. db.session.commit()
  281. return jsonify(success(_element(record))), 201
  282. except Exception as error:
  283. db.session.rollback()
  284. return _error(error)
  285. @bp.route("/data-elements/<element_uid>/transitions", methods=["POST"])
  286. def transition_data_element(element_uid):
  287. payload = request.get_json(silent=True) or {}
  288. identity = _identity()
  289. target_status = str(payload.get("target_status") or "")
  290. if (
  291. target_status in {"published", "deprecated", "retired"}
  292. and "data-elements:publish" not in set(identity.get("permissions") or [])
  293. ):
  294. return jsonify(failed("权限不足", code=403)), 403
  295. try:
  296. record = get_data_element_service().transition(
  297. element_uid,
  298. target_status,
  299. expected_version=int(payload.get("expected_version")),
  300. actor_uid=identity.get("id") or identity.get("sub"),
  301. )
  302. db.session.commit()
  303. return jsonify(success(_element(record))), 200
  304. except Exception as error:
  305. db.session.rollback()
  306. return _error(error)
  307. @bp.route("/candidate-decisions", methods=["POST"])
  308. def decide_candidates():
  309. payload = request.get_json(silent=True) or {}
  310. try:
  311. records = get_candidate_decision_service().decide(
  312. payload.get("decisions") or [],
  313. actor_uid=_identity().get("id") or _identity().get("sub"),
  314. )
  315. return jsonify(success([_decision(record) for record in records])), 200
  316. except Exception as error:
  317. return _error(error)
  318. @bp.route("/sources/files", methods=["POST"])
  319. def upload_source_file():
  320. uploaded = request.files.get("file")
  321. if uploaded is None:
  322. return jsonify(failed("缺少上传文件", code=400)), 400
  323. try:
  324. record, created = get_artifact_service().store(
  325. request.form.get("source_uid"),
  326. uploaded.filename,
  327. uploaded.mimetype,
  328. uploaded.read(),
  329. request.form.get("parser_version") or "auto-v1",
  330. )
  331. db.session.commit()
  332. return jsonify(success(_artifact(record))), 201 if created else 200
  333. except Exception as error:
  334. db.session.rollback()
  335. return _error(error)
  336. @bp.route("/evidence/<evidence_uid>", methods=["GET"])
  337. def get_evidence(evidence_uid):
  338. identity = _identity()
  339. try:
  340. service = get_evidence_service()
  341. if request.args.get("download") in {"1", "true", "yes"}:
  342. if "evidence:download" not in set(identity.get("permissions") or []):
  343. return jsonify(failed("权限不足", code=403)), 403
  344. content, filename, media_type = service.download(evidence_uid)
  345. return send_file(
  346. io.BytesIO(content),
  347. mimetype=media_type,
  348. as_attachment=True,
  349. download_name=filename,
  350. )
  351. preview = service.preview(evidence_uid)
  352. from app.core.data_research.artifacts import redact_excerpt
  353. preview["excerpt"] = redact_excerpt(preview.get("excerpt"))
  354. return jsonify(success(preview)), 200
  355. except Exception as error:
  356. return _error(error)
  357. @bp.route("/ontologies", methods=["GET"])
  358. def list_ontologies():
  359. try:
  360. return jsonify(success([_ontology(item) for item in get_ontology_service().list()])), 200
  361. except Exception as error:
  362. return _error(error)
  363. @bp.route("/ontologies", methods=["POST"])
  364. def create_ontology():
  365. try:
  366. record = get_ontology_service().create(
  367. request.get_json(silent=True) or {},
  368. _identity().get("id") or _identity().get("sub"),
  369. )
  370. return jsonify(success(_ontology(record))), 201
  371. except Exception as error:
  372. return _error(error)
  373. @bp.route("/ontologies/<ontology_uid>", methods=["GET"])
  374. def get_ontology(ontology_uid):
  375. try:
  376. return jsonify(success(_ontology(get_ontology_service().get(ontology_uid)))), 200
  377. except Exception as error:
  378. return _error(error)
  379. @bp.route("/ontologies/<ontology_uid>/versions", methods=["GET"])
  380. def list_ontology_versions(ontology_uid):
  381. try:
  382. return jsonify(
  383. success(
  384. [
  385. _ontology_version(item)
  386. for item in get_ontology_service().list_versions(ontology_uid)
  387. ]
  388. )
  389. ), 200
  390. except Exception as error:
  391. return _error(error)
  392. @bp.route("/ontologies/<ontology_uid>/graph", methods=["GET"])
  393. def get_ontology_graph(ontology_uid):
  394. try:
  395. service = get_ontology_service()
  396. ontology = service.get(ontology_uid)
  397. version = service.latest_version(ontology_uid)
  398. data = (
  399. _ontology_version(version)
  400. if version is not None
  401. else {
  402. "uid": None,
  403. "ontology_uid": str(ontology_uid),
  404. "version": 0,
  405. "parent_version_uid": None,
  406. "status": "draft",
  407. "content_hash": None,
  408. "graph_document": {
  409. name: []
  410. for name in (
  411. "classes",
  412. "properties",
  413. "relations",
  414. "constraints",
  415. "domain_links",
  416. "element_mappings",
  417. )
  418. },
  419. "created_by": None,
  420. }
  421. )
  422. response = jsonify(success(data))
  423. response.headers["ETag"] = f'"{int(ontology.draft_revision)}"'
  424. return response, 200
  425. except Exception as error:
  426. return _error(error)
  427. @bp.route("/ontologies/<ontology_uid>/graph", methods=["PATCH"])
  428. def save_ontology_graph(ontology_uid):
  429. etag = str(request.headers.get("If-Match") or "").strip().strip('"')
  430. if not etag.isdigit():
  431. return jsonify(failed("缺少有效 If-Match 修订号", code=428)), 428
  432. try:
  433. version = get_ontology_service().save_draft(
  434. ontology_uid,
  435. request.get_json(silent=True) or {},
  436. int(etag),
  437. _identity().get("id") or _identity().get("sub"),
  438. )
  439. response = jsonify(success(_ontology_version(version)))
  440. response.headers["ETag"] = f'"{int(etag) + 1}"'
  441. return response, 200
  442. except Exception as error:
  443. return _error(error)
  444. @bp.route("/ontologies/<ontology_uid>/validate", methods=["POST"])
  445. def validate_ontology(ontology_uid):
  446. try:
  447. issues = get_ontology_service().validate(ontology_uid)
  448. data = [
  449. {"code": item.code, "message": item.message, "path": item.path}
  450. for item in issues
  451. ]
  452. return jsonify(success(data)), 200
  453. except Exception as error:
  454. return _error(error)
  455. @bp.route("/ontologies/<ontology_uid>/publish", methods=["POST"])
  456. def publish_ontology(ontology_uid):
  457. try:
  458. version = get_ontology_service().publish(
  459. ontology_uid,
  460. request.headers.get("Idempotency-Key") or f"{ontology_uid}:publish",
  461. _identity().get("id") or _identity().get("sub"),
  462. )
  463. return jsonify(success(_ontology_version(version))), 200
  464. except Exception as error:
  465. return _error(error)
  466. @bp.route("/ontologies/<ontology_uid>/diff", methods=["GET"])
  467. def diff_ontology(ontology_uid):
  468. try:
  469. return jsonify(success(get_ontology_service().diff(
  470. ontology_uid, request.args.get("left"), request.args.get("right")
  471. ))), 200
  472. except Exception as error:
  473. return _error(error)
  474. @bp.route("/ontologies/<ontology_uid>/rollback", methods=["POST"])
  475. def rollback_ontology(ontology_uid):
  476. payload = request.get_json(silent=True) or {}
  477. try:
  478. version = get_ontology_service().rollback(
  479. ontology_uid,
  480. payload.get("target_version_uid"),
  481. int(payload.get("expected_revision")),
  482. _identity().get("id") or _identity().get("sub"),
  483. )
  484. return jsonify(success(_ontology_version(version))), 201
  485. except Exception as error:
  486. return _error(error)
  487. @bp.route("/ontologies/<ontology_uid>/suggestions", methods=["POST"])
  488. def generate_ontology_suggestions(ontology_uid):
  489. try:
  490. result = get_ontology_dynamic_service().generate(
  491. ontology_uid,
  492. request.get_json(silent=True) or {},
  493. _identity().get("id") or _identity().get("sub"),
  494. )
  495. records = (
  496. result.get("suggestions") or []
  497. if isinstance(result, dict)
  498. else result
  499. )
  500. suggestions = [
  501. item
  502. if isinstance(item, dict)
  503. else {
  504. "uid": item.uid,
  505. "kind": item.kind,
  506. "payload": item.payload,
  507. "evidence_uids": list(item.evidence_uids),
  508. "confidence": item.confidence,
  509. "source": item.source,
  510. "model_version": item.model_version,
  511. "prompt_version": item.prompt_version,
  512. }
  513. for item in records
  514. ]
  515. data = {
  516. "change_set_uid": (
  517. result.get("change_set_uid")
  518. if isinstance(result, dict)
  519. else None
  520. ),
  521. "suggestions": suggestions,
  522. }
  523. return jsonify(success(data)), 201
  524. except Exception as error:
  525. return _error(error)
  526. @bp.route(
  527. "/ontologies/<ontology_uid>/change-sets/<change_set_uid>/decisions",
  528. methods=["POST"],
  529. )
  530. def decide_ontology_change_set(ontology_uid, change_set_uid):
  531. payload = request.get_json(silent=True) or {}
  532. try:
  533. result = get_ontology_dynamic_service().decide(
  534. ontology_uid,
  535. change_set_uid,
  536. payload.get("decisions") or [],
  537. _identity().get("id") or _identity().get("sub"),
  538. )
  539. return jsonify(success(result)), 200
  540. except Exception as error:
  541. return _error(error)
  542. @bp.route("/ontologies/<ontology_uid>/export", methods=["GET"])
  543. def export_ontology(ontology_uid):
  544. try:
  545. content, media_type, filename = get_ontology_exchange_service().export(
  546. ontology_uid, request.args.get("format") or "json"
  547. )
  548. return send_file(
  549. io.BytesIO(content),
  550. mimetype=media_type,
  551. as_attachment=True,
  552. download_name=filename,
  553. )
  554. except Exception as error:
  555. return _error(error)
  556. @bp.route("/ontologies/import", methods=["POST"])
  557. def import_ontology():
  558. try:
  559. uploaded = request.files.get("file")
  560. content = uploaded.read() if uploaded is not None else request.get_data(cache=False)
  561. result = get_ontology_exchange_service().import_document(
  562. content,
  563. request.args.get("format") or "json",
  564. _identity().get("id") or _identity().get("sub"),
  565. )
  566. return jsonify(success(result)), 201
  567. except Exception as error:
  568. return _error(error)
  569. @bp.route("/semantic/properties/<property_uid>", methods=["GET"])
  570. def query_semantic_property(property_uid):
  571. domain = str(request.args.get("business_domain_uid") or "").strip()
  572. identity = _identity()
  573. scoped = set(identity.get("business_domains") or [domain])
  574. if "*" not in scoped and domain not in scoped:
  575. return jsonify(failed("权限不足", code=403)), 403
  576. try:
  577. result = get_semantic_query_service().trace_property(
  578. property_uid,
  579. allowed_domains={domain},
  580. limit=int(request.args.get("limit") or 20),
  581. after_uid=request.args.get("after_uid"),
  582. )
  583. return jsonify(success(result)), 200
  584. except Exception as error:
  585. return _error(error)