"""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.connectors.errors import ( ConnectorAuthenticationError, ConnectorConfigurationError, ConnectorError, ) 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, ConnectorError): logger.warning("连接器操作失败: category=%s", error.category) return jsonify( failed( "连接器操作失败", code=error.http_status, error={ "code": "CONNECTOR_ERROR", "category": error.category, "retryable": error.retryable, }, ) ), error.http_status 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(): from app import db from app.core.connectors.repository import ConnectorRepository payload = request.get_json(silent=True) or {} try: graph = ConnectorRepository(db.session).graph( source_uid=payload.get("source_uid"), business_domain_uid=payload.get("business_domain_uid"), process_key=payload.get("process_key"), run_uid=payload.get("run_uid"), limit=payload.get("limit", 500), ) return jsonify(success(graph)), 200 except Exception as error: return _error_response(error) def _connector_registry(): from flask import current_app from app.core.connectors.builtin import register_builtin_connectors from app.core.connectors.registry import ConnectorRegistry from app.core.connectors.secrets import EnvironmentSecretResolver from app.core.data_source.runtime import get_data_source_manager registry = current_app.extensions.get("connector_registry") if registry is not None: return registry registry = ConnectorRegistry() register_builtin_connectors( registry, connection_provider=get_data_source_manager().connect, secret_resolver=current_app.config.get("CONNECTOR_SECRET_RESOLVER") or EnvironmentSecretResolver(), ) current_app.extensions["connector_registry"] = registry return registry @bp.route("/connectors/manifests", methods=["GET"]) def connector_manifests(): try: return jsonify(success({"manifests": _connector_registry().manifests()})), 200 except Exception as error: return _error_response(error) @bp.route("/connectors/config/validate", methods=["POST"]) def connector_config_validate(): payload = request.get_json(silent=True) or {} try: config = _connector_registry().validate( payload.get("connector_id"), payload.get("version"), payload.get("config") ) return jsonify(success({"valid": True, "config_keys": sorted(config)})), 200 except Exception as error: return _error_response(error) @bp.route("/connectors///health", methods=["POST"]) def connector_health(connector_id, version): payload = request.get_json(silent=True) or {} try: connector = _connector_registry().resolve(connector_id, version) return jsonify( success(vars(connector.health(payload.get("config") or {}))) ), 200 except Exception as error: return _error_response(error) @bp.route("/connectors///compatibility", methods=["GET"]) def connector_compatibility(connector_id, version): try: connector = _connector_registry().resolve(connector_id, version) return jsonify(success(vars(connector.compatibility()))), 200 except Exception as error: return _error_response(error) @bp.route("/connectors/runs", methods=["GET"]) def connector_runs_list(): from app import db from app.core.connectors.repository import ConnectorRepository try: repository = ConnectorRepository(db.session, _actor_uid()) return jsonify( success({"runs": repository.list(request.args.get("limit", 100))}) ), 200 except Exception as error: db.session.rollback() return _error_response(error) def _require_human_dry_run(payload): if payload.get("dry_run") is not True: raise ConnectorConfigurationError("human connector runs must use dry_run=true") def _require_collection_operation(payload): if payload.get("operation") not in { "discover", "snapshot", "incremental", "lineage", "profile", "evidence", }: raise ConnectorConfigurationError("connector collection operation is invalid") @bp.route("/connectors/runs", methods=["POST"]) def connector_run_create(): from app import db from app.core.connectors.repository import ConnectorRepository try: repository = ConnectorRepository( db.session, _actor_uid(), run_access="human" ) from dataclasses import asdict from app.core.connectors.runtime import ConnectorRuntime from app.core.connectors.sdk import OperationRequest payload = request.get_json(silent=True) or {} _require_human_dry_run(payload) _require_collection_operation(payload) registry = _connector_registry() for manifest in ( registry.resolve( payload.get("connector_id"), payload.get("version") ).manifest, ): repository.register_manifest(manifest, _actor_uid()) db.session.commit() operation_request = OperationRequest( source_uid=str(payload.get("source_uid") or ""), operation=str(payload.get("operation") or ""), config=payload.get("config") or {}, scope=payload.get("scope") or {}, cursor=payload.get("cursor") or {}, checkpoint=payload.get("checkpoint") or {}, idempotency_key=payload.get("idempotency_key"), dry_run=bool(payload.get("dry_run", False)), process_key=str(payload.get("process_key") or "dry-run"), ) result = ConnectorRuntime(registry, store=repository).execute( payload.get("connector_id"), payload.get("version"), operation_request ) return jsonify(success(asdict(result), code=201)), 201 except Exception as error: db.session.rollback() return _error_response(error) @bp.route("/connectors/machine/runs", methods=["POST"]) def connector_machine_run(): """Machine-only run boundary; a human bearer token is never accepted here.""" from app import db from app.core.connectors.identity import ConnectorIdentityRepository from app.core.connectors.repository import ConnectorRepository from app.core.connectors.runtime import ConnectorRuntime from app.core.connectors.sdk import OperationRequest payload = request.get_json(silent=True) or {} token = request.headers.get("X-Connector-Credential", "").strip() if not token: return _error_response( ConnectorAuthenticationError("machine credential is required") ) try: for field in ( "connector_id", "version", "source_uid", "business_domain_uid", "environment", "process_key", "operation", ): if not str(payload.get(field) or "").strip(): raise ConnectorConfigurationError( "machine connector binding is incomplete" ) _require_collection_operation(payload) if "config" in payload: raise ConnectorConfigurationError( "machine connector config is supplied by the approved source binding" ) identity = ConnectorIdentityRepository(db.session).authenticate( token, connector_id=str(payload.get("connector_id") or ""), connector_version=str(payload.get("version") or ""), source_uid=str(payload.get("source_uid") or ""), business_domain_uid=str(payload.get("business_domain_uid") or ""), environment=str(payload.get("environment") or ""), operation=str(payload.get("operation") or ""), scope=payload.get("scope") or {}, ) registry = _connector_registry() manifest = registry.resolve( payload.get("connector_id"), payload.get("version") ).manifest if not identity.get("source_binding_uid"): raise ConnectorConfigurationError( "machine connector principal has no approved source binding" ) repository = ConnectorRepository( db.session, identity["principal_uid"], run_access="machine", principal_uid=identity["principal_uid"], ) repository.register_manifest(manifest, identity["principal_uid"]) db.session.commit() operation_request = OperationRequest( source_uid=str(payload.get("source_uid")), operation=str(payload.get("operation")), config=identity["approved_config"], scope=payload.get("scope") or {}, cursor=payload.get("cursor") or {}, checkpoint=payload.get("checkpoint") or {}, idempotency_key=payload.get("idempotency_key"), dry_run=bool(payload.get("dry_run", False)), principal_uid=identity["principal_uid"], business_domain_uid=str(payload.get("business_domain_uid")), environment=str(payload.get("environment")), process_key=str(payload.get("process_key") or ""), source_binding_uid=identity["source_binding_uid"], source_binding_version=identity["source_binding_version"], ) result = ConnectorRuntime(registry, store=repository).execute( payload.get("connector_id"), payload.get("version"), operation_request ) from dataclasses import asdict return jsonify(success(asdict(result), code=201)), 201 except Exception as error: db.session.rollback() return _error_response(error) @bp.route("/connectors/runs//cancel", methods=["POST"]) def connector_run_cancel(idempotency_key): from app import db from app.core.connectors.repository import ConnectorRepository from app.core.connectors.runtime import ConnectorRuntime try: result = ConnectorRuntime( _connector_registry(), store=ConnectorRepository( db.session, _actor_uid(), run_access="human" ), ).cancel(idempotency_key) return jsonify(success(result)), 200 except Exception as error: db.session.rollback() return _error_response(error) @bp.route("/connectors/runs//resume", methods=["POST"]) def connector_run_resume(idempotency_key): from app import db from app.core.connectors.repository import ConnectorRepository from app.core.connectors.runtime import ConnectorRuntime try: result = ConnectorRuntime( _connector_registry(), store=ConnectorRepository( db.session, _actor_uid(), run_access="human" ), ).resume(idempotency_key) return jsonify(success(result)), 200 except Exception as error: db.session.rollback() return _error_response(error) def _machine_run_action(idempotency_key, action): from app import db from app.core.connectors.identity import ConnectorIdentityRepository from app.core.connectors.repository import ConnectorRepository from app.core.connectors.runtime import ConnectorRuntime token = request.headers.get("X-Connector-Credential", "").strip() if not token: return _error_response( ConnectorAuthenticationError("machine credential is required") ) try: record = ConnectorRepository(db.session).get(idempotency_key) if not record or not record.get("principal_uid") or record.get("dry_run"): raise ConnectorConfigurationError("machine connector run was not found") identity = ConnectorIdentityRepository(db.session).authenticate( token, connector_id=record["connector_id"], connector_version=record["connector_version"], source_uid=record["source_uid"], business_domain_uid=record["business_domain_uid"], environment=record["environment"], operation=action, scope=record.get("scope") or {}, ) if ( identity["principal_uid"] != record["principal_uid"] or identity.get("source_binding_uid") != record.get("source_binding_uid") or identity.get("source_binding_version") != record.get("source_binding_version") ): raise ConnectorAuthenticationError( "machine credential run binding was rejected" ) repository = ConnectorRepository( db.session, identity["principal_uid"], run_access="machine", principal_uid=identity["principal_uid"], ) runtime = ConnectorRuntime(_connector_registry(), store=repository) result = getattr(runtime, action)(idempotency_key) if action == "resume": from dataclasses import asdict result = asdict(result) return jsonify(success(result)), 200 except Exception as error: db.session.rollback() return _error_response(error) @bp.route( "/connectors/machine/runs//cancel", methods=["POST"] ) def connector_machine_run_cancel(idempotency_key): """Cancel one bound machine run with a fresh one-time credential.""" return _machine_run_action(idempotency_key, "cancel") @bp.route( "/connectors/machine/runs//resume", methods=["POST"] ) def connector_machine_run_resume(idempotency_key): """Resume one bound machine run with a fresh one-time credential.""" return _machine_run_action(idempotency_key, "resume") @bp.route("/connectors/source-bindings", methods=["GET"]) def connector_source_bindings_list(): from app import db from app.core.connectors.bindings import ConnectorSourceBindingRepository try: items = ConnectorSourceBindingRepository(db.session).list_public( request.args.get("limit", 100) ) return jsonify(success({"source_bindings": items})), 200 except Exception as error: db.session.rollback() return _error_response(error) @bp.route("/connectors/source-bindings", methods=["POST"]) def connector_source_binding_approve(): from app import db from app.core.connectors.bindings import ConnectorSourceBindingRepository from app.core.connectors.repository import ConnectorRepository payload = request.get_json(silent=True) or {} try: registry = _connector_registry() manifest = registry.resolve( payload.get("connector_id"), payload.get("version") ).manifest approved_config = registry.validate( payload.get("connector_id"), payload.get("version"), payload.get("approved_config") or {}, ) ConnectorRepository(db.session, _actor_uid()).register_manifest( manifest, _actor_uid() ) result = ConnectorSourceBindingRepository(db.session).approve( connector_id=payload.get("connector_id"), connector_version=payload.get("version"), source_uid=payload.get("source_uid"), business_domain_uid=payload.get("business_domain_uid"), environment=payload.get("environment"), approved_config=approved_config, approved_by=_actor_uid(), binding_uid=payload.get("binding_uid"), ) return jsonify(success(result, code=201)), 201 except Exception as error: db.session.rollback() return _error_response(error) @bp.route("/connectors/source-bindings//revoke", methods=["POST"]) def connector_source_binding_revoke(binding_uid): from app import db from app.core.connectors.bindings import ConnectorSourceBindingRepository try: revoked = ConnectorSourceBindingRepository(db.session).revoke( binding_uid, _actor_uid() ) return jsonify(success(revoked)), 200 except Exception as error: db.session.rollback() return _error_response(error) @bp.route("/connectors/principals", methods=["POST"]) def connector_principal_create(): from app import db from app.core.connectors.identity import ConnectorIdentityRepository payload = request.get_json(silent=True) or {} try: _connector_registry().resolve( payload.get("connector_id"), payload.get("version", "1.0.0") ) uid = ConnectorIdentityRepository(db.session).create_principal( connector_id=payload.get("connector_id"), connector_version=payload.get("version", "1.0.0"), source_uid=payload.get("source_uid"), business_domain_uid=payload.get("business_domain_uid"), environment=payload.get("environment"), operations=payload.get("operations") or (), scopes=payload.get("scopes") or {}, actor_uid=_actor_uid(), source_binding_uid=payload.get("source_binding_uid"), source_binding_version=payload.get("source_binding_version"), ) return jsonify(success({"principal_uid": uid}, code=201)), 201 except Exception as error: db.session.rollback() return _error_response(error) @bp.route("/connectors/principals//credentials", methods=["POST"]) def connector_credential_issue(principal_uid): from app import db from app.core.connectors.identity import ConnectorIdentityRepository payload = request.get_json(silent=True) or {} try: result = ConnectorIdentityRepository(db.session).issue( principal_uid, ttl_seconds=payload.get("ttl_seconds", 900), actor_uid=_actor_uid(), ) response = jsonify(success(result, code=201)) response.headers["Cache-Control"] = "no-store" return response, 201 except Exception as error: db.session.rollback() return _error_response(error) @bp.route("/connectors/credentials//", methods=["POST"]) def connector_credential_action(credential_uid, action): from app import db from app.core.connectors.identity import ConnectorIdentityRepository payload = request.get_json(silent=True) or {} try: repository = ConnectorIdentityRepository(db.session) if action == "revoke": result = {"revoked": repository.revoke(credential_uid, _actor_uid())} elif action == "rotate": result = repository.rotate( credential_uid, ttl_seconds=payload.get("ttl_seconds", 900), actor_uid=_actor_uid(), ) else: raise ConnectorConfigurationError("credential action is invalid") response = jsonify(success(result)) response.headers["Cache-Control"] = "no-store" return response, 200 except Exception as error: db.session.rollback() return _error_response(error)