| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227 |
- """HTTP boundary for secret-free external data-source management."""
- import logging
- from flask import g, jsonify, request
- from app.api.data_source import bp
- from app.core.data_source.errors import DataSourceError
- from app.core.data_source.redaction import (
- redact_mapping,
- sanitize_exception,
- )
- from app.models.result import failed, success
- logger = logging.getLogger(__name__)
- def get_data_source_service():
- from app.core.data_source.runtime import get_data_source_manager
- from app.core.data_source.service import build_data_source_service
- return build_data_source_service(get_data_source_manager())
- def _actor_uid():
- identity = getattr(g, "current_user", {}) or {}
- return identity.get("id") or identity.get("sub")
- def _error_response(error):
- if isinstance(error, DataSourceError):
- logger.warning(
- "数据源操作失败: code=%s message=%s",
- error.code,
- sanitize_exception(error),
- )
- return (
- jsonify(
- failed(
- str(error),
- code=error.http_status,
- error={"code": error.code},
- )
- ),
- error.http_status,
- )
- logger.error(
- "数据源操作异常: %s",
- sanitize_exception(error),
- )
- return (
- jsonify(
- failed(
- "数据源操作失败",
- code=500,
- error={"code": "DATASOURCE_ERROR"},
- )
- ),
- 500,
- )
- @bp.route("/save", methods=["POST"])
- def data_source_save():
- payload = request.get_json(silent=True) or {}
- logger.debug("保存数据源请求: %s", redact_mapping(payload))
- try:
- service = get_data_source_service()
- definition, created = service.save(
- payload,
- actor_uid=_actor_uid(),
- )
- status = 201 if created else 200
- return jsonify(success(service.serialize(definition))), status
- except Exception as error:
- return _error_response(error)
- @bp.route("/list", methods=["POST"])
- def data_source_list():
- payload = request.get_json(silent=True) or {}
- try:
- service = get_data_source_service()
- definitions = service.list(payload)
- items = [service.serialize(item) for item in definitions]
- return jsonify(
- success({"data_source": items, "total": len(items)})
- ), 200
- except Exception as error:
- return _error_response(error)
- @bp.route("/delete", methods=["POST"])
- def data_source_delete():
- payload = request.get_json(silent=True) or {}
- logger.debug("删除数据源请求: %s", redact_mapping(payload))
- try:
- result = get_data_source_service().delete(
- payload.get("uid"),
- actor_uid=_actor_uid(),
- )
- return jsonify(success(result)), 200
- except Exception as error:
- return _error_response(error)
- @bp.route("/conntest", methods=["POST"])
- def data_source_conn_test():
- payload = request.get_json(silent=True) or {}
- logger.debug("测试数据源连接请求: %s", redact_mapping(payload))
- try:
- result = get_data_source_service().test_connection(payload)
- return jsonify(success(result)), 200
- except Exception as error:
- return _error_response(error)
- @bp.route("/valid", methods=["POST"])
- def data_source_connstr_valid():
- payload = request.get_json(silent=True) or {}
- logger.debug("验证数据源连接请求: %s", redact_mapping(payload))
- try:
- result = get_data_source_service().test_connection(payload)
- return jsonify(success({"exists": False, **result})), 200
- except Exception as error:
- return _error_response(error)
- @bp.route("/parse", methods=["POST"])
- def data_source_connstr_parse():
- return (
- jsonify(
- failed(
- "连接字符串快捷解析已停用",
- code=410,
- error={"code": "DATASOURCE_PARSE_RETIRED"},
- )
- ),
- 410,
- )
- @bp.route("/pools", methods=["GET"])
- def data_source_pool_list():
- try:
- service = get_data_source_service()
- items = [
- service.serialize_pool_status(status)
- for status in service.pool_statuses()
- ]
- return jsonify(success({"pools": items, "total": len(items)})), 200
- except Exception as error:
- return _error_response(error)
- @bp.route("/<data_source_uid>/pool", methods=["GET"])
- def data_source_pool_status(data_source_uid):
- try:
- service = get_data_source_service()
- status = service.pool_status(data_source_uid)
- serializer = getattr(
- service,
- "serialize_pool_status",
- None,
- )
- data = (
- serializer(status)
- if serializer is not None
- else {
- "data_source_uid": status.data_source_uid,
- "credential_version": status.credential_version,
- "pool_state": status.pool_state,
- "pool_size": status.pool_size,
- "checked_out": status.checked_out,
- "checked_in": status.checked_in,
- "overflow": status.overflow,
- "leases": status.leases,
- }
- )
- return jsonify(success(data)), 200
- except Exception as error:
- return _error_response(error)
- @bp.route("/<data_source_uid>/pool/invalidate", methods=["POST"])
- def data_source_pool_invalidate(data_source_uid):
- payload = request.get_json(silent=True) or {}
- reason = payload.get("reason")
- if reason not in {
- "admin_reset",
- "configuration_changed",
- "credential_rotated",
- }:
- return (
- jsonify(
- failed(
- "连接池失效原因无效",
- code=400,
- error={"code": "DATASOURCE_CONFIGURATION_INVALID"},
- )
- ),
- 400,
- )
- try:
- result = get_data_source_service().invalidate_pool(
- data_source_uid,
- reason=reason,
- actor_uid=_actor_uid(),
- )
- return jsonify(success(result)), 200
- except Exception as error:
- return _error_response(error)
- @bp.route("/graph", methods=["POST"])
- def data_source_graph_relationship():
- return (
- jsonify(
- failed(
- "该功能尚未实现",
- code=501,
- error={"code": "DATASOURCE_GRAPH_NOT_IMPLEMENTED"},
- )
- ),
- 501,
- )
|