"""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("//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("//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, )