全部文档
当前文档

暂无内容

如果没有找到您期望的内容,请尝试其他搜索词

文档中心

最佳实践:Agent中集成MCP,自定义OAuth授权流程实例代码

最近更新时间:2026-09-08 19:37:48

自定义OAuth授权流程示例代码

# -*- coding: utf-8 -*-

# OAuth 极简版:纯 MCP SDK,单文件,复制即可运行

# 依赖:pip install "mcp>=1.9" # Python >= 3.11

import asyncio

import os

import webbrowser from http.server import BaseHTTPRequestHandler, HTTPServer from urllib.parse import parse_qs, urlparse

from mcp import ClientSession
from mcp.client.auth import OAuthClientProvider, TokenStorage
from mcp.client.streamable_http import streamablehttp_client
from mcp.shared.auth import (
OAuthClientInformationFull,
OAuthClientMetadata,
OAuthToken,
)

MCP_URL = os.getenv("MCP_URL", "") # 替换为你的 MCP 服务地址,或设置环境变量 MCP_URL
CALLBACK_PORT = 3000
CALLBACK_TIMEOUT = 120.0 # 等待授权回调的最长时间(秒)
SUCCESS_HTML = (
"<!DOCTYPE html><html><head><meta charset='utf-8'></head>"
"<body><h1>授权完成</h1><p>请回到终端继续。</p></body></html>"
)

1) token 存储

⚠️ 内存版仅用于演示,每次运行都要重新授权。

生产环境请换成持久化存储(如落盘到本地文件),避免反复授权、并正确保留 refresh_token。

class MemoryStorage(TokenStorage):
tokens: OAuthToken | None = None
client_info: OAuthClientInformationFull | None = None

async def get_tokens(self) -&gt; OAuthToken | None:
    return self.tokens

async def set_tokens(self, tokens: OAuthToken) -&gt; None:
    self.tokens = tokens

async def get_client_info(self) -&gt; OAuthClientInformationFull | None:
    return self.client_info

async def set_client_info(self, client_info: OAuthClientInformationFull) -&gt; None:
    self.client_info = client_info

2) 打开浏览器,完成金山云 IAM 登录授权

async def open_browser(auth_url: str):
print(f"请在浏览器中完成授权:\n{auth_url}\n")
webbrowser.open(auth_url)

3) 本地端口接收授权回调

async def wait_callback():
code, state = None, None

class Handler(BaseHTTPRequestHandler):
    def do_GET(self):
        nonlocal code, state
        q = parse_qs(urlparse(self.path).query)
        code, state = q.get("code", [None])[0], q.get("state", [None])[0]
        self.send_response(200)
        self.send_header("Content-Type", "text/html; charset=utf-8")
        self.end_headers()
        self.wfile.write(SUCCESS_HTML.encode("utf-8"))

    def log_message(self, *args):
        pass

try:
    server = HTTPServer(("127.0.0.1", CALLBACK_PORT), Handler)
except OSError as e:
    raise RuntimeError(
        f"无法 127.0.0.1:{CALLBACK_PORT},可能端口被占用:{e}"
    ) from e

server.timeout = 0.5
loop = asyncio.get_running_loop()
deadline = loop.time() + CALLBACK_TIMEOUT
try:
    while code is None:
        if loop.time() &gt; deadline:
            raise TimeoutError(
                f"等待授权回调超时({CALLBACK_TIMEOUT}s),"
                "请确认是否已在浏览器完成授权。"
            )
        await loop.run_in_executor(None, server.handle_request)
finally:
    server.server_close()
# state 由 MCP SDK 校验防 CSRF,本 demo 直接回传,不在此处额外校验
return code, state

oauth = OAuthClientProvider(
server_url=MCP_URL,
client_metadata=OAuthClientMetadata(
client_name="mcp-oauth-demo",
redirect_uris=[f"http://127.0.0.1:{CALLBACK_PORT}/callback"],
grant_types=["authorization_code", "refresh_token"],
response_types=["code"],
scope="ksc/mcpserver",
token_endpoint_auth_method="none", # 公共客户端 + PKCE,无需 client_secret
),
storage=MemoryStorage(),
redirect_handler=open_browser,
callback_handler=wait_callback,
)

async def main():
if not MCP_URL:
raise RuntimeError(
"错误:未配置 MCP 服务地址,请设置环境变量 MCP_URL "
"(例如 export MCP_URL=https://your-mcp-host/mcp)"
)
# auth=oauth:401 后自动完成 服务发现 -> 客户端注册 -> 授权码 + PKCE -> 换取 token
async with streamablehttp_client(MCP_URL, auth=oauth) as (read, write, _):
async with ClientSession(read, write) as session:
await session.initialize()
tools = await session.list_tools()
print("可用工具:", [t.name for t in tools.tools])
result = await session.call_tool(
"eip-20160304-DescribeAddresses",
{"x_ksc_region": "cn-beijing-6", "MaxResults": 5},
)
print("调用结果:", result.content[0].text)

if name == "main":
asyncio.run(main())

在Agent中集成MCP,自定义OAuth授权流程实例代码

# -- coding: utf-8 -- # OAuth 鉴权 MCP 服务接入 ksadk Agent

# 依赖:pip install "ksadk[adk]" "mcp>=1.9" # Python >= 3.11

import asyncio

import os import webbrowser from http.server import BaseHTTPRequestHandler, HTTPServer from urllib.parse import parse_qs, urlparse

import google.genai.types as types
from google.adk.agents import Agent
from google.adk.models.lite_llm import LiteLlm
from google.adk.runners import Runner
from google.adk.sessions import InMemorySessionService
from google.adk.tools.mcp_tool.mcp_session_manager import StreamableHTTPConnectionParams
from google.adk.tools.mcp_tool.mcp_toolset import McpToolset
from mcp.client.auth import OAuthClientProvider, TokenStorage
from mcp.shared._httpx_utils import create_mcp_http_client
from mcp.shared.auth import (
OAuthClientInformationFull,
OAuthClientMetadata,
OAuthToken,
)

可变配置统一走环境变量,缺失时给出中文报错

MCP_URL = os.getenv("MCP_URL", "") # 替换为你的 MCP 服务地址,或设置环境变量 MCP_URL
CALLBACK_PORT = 3000
CALLBACK_TIMEOUT = 120.0 # 等待授权回调的最长时间(秒)
SUCCESS_HTML = (
"<!DOCTYPE html><html><head><meta charset='utf-8'></head>"
"<body><h1>授权完成</h1><p>请回到终端继续。</p></body></html>"
)

def _require_env(name: str) -> str:
value = os.getenv(name)
if not value:
raise RuntimeError(
f"错误:未配置 {name},请先设置环境变量 "
f"(例如 export {name}=xxxxx)"
)
return value

大模型网关配置(OpenAI 兼容接口)

OPENAI_BASE_URL = os.getenv("OPENAI_BASE_URL", "https://kspmas.ksyun.com/v1")
OPENAI_API_KEY = _require_env("OPENAI_API_KEY") # 勿硬编码,从环境变量读取
MODEL_ID = os.getenv("MODEL_ID", "openai/glm-5.2")
os.environ.setdefault("OPENAI_BASE_URL", OPENAI_BASE_URL)
os.environ.setdefault("OPENAI_API_KEY", OPENAI_API_KEY)

1) token 存储

⚠️ 内存版仅用于演示,每次运行都要重新授权。

生产环境请换成持久化存储,

避免反复授权、并正确保留 refresh_token。

class MemoryStorage(TokenStorage):
tokens: OAuthToken | None = None
client_info: OAuthClientInformationFull | None = None

async def get_tokens(self) -&gt; OAuthToken | None:
    return self.tokens

async def set_tokens(self, tokens: OAuthToken) -&gt; None:
    self.tokens = tokens

async def get_client_info(self) -&gt; OAuthClientInformationFull | None:
    return self.client_info

async def set_client_info(self, client_info: OAuthClientInformationFull) -&gt; None:
    self.client_info = client_info

2) 打开浏览器,完成金山云 IAM 登录授权

async def open_browser(auth_url: str):
print(f"请在浏览器中完成授权:\n{auth_url}\n")
webbrowser.open(auth_url)

3) 本地端口接收授权回调

async def wait_callback():
code, state = None, None

class Handler(BaseHTTPRequestHandler):
    def do_GET(self):
        nonlocal code, state
        q = parse_qs(urlparse(self.path).query)
        code, state = q.get("code", [None])[0], q.get("state", [None])[0]
        self.send_response(200)
        self.send_header("Content-Type", "text/html; charset=utf-8")
        self.end_headers()
        self.wfile.write(SUCCESS_HTML.encode("utf-8"))

    def log_message(self, *args):
        pass

try:
    server = HTTPServer(("127.0.0.1", CALLBACK_PORT), Handler)
except OSError as e:
    raise RuntimeError(
        f"无法调用 127.0.0.1:{CALLBACK_PORT},可能端口被占用:{e}"
    ) from e

server.timeout = 0.5
loop = asyncio.get_running_loop()
deadline = loop.time() + CALLBACK_TIMEOUT
try:
    while code is None:
        if loop.time() &gt; deadline:
            raise TimeoutError(
                f"等待授权回调超时({CALLBACK_TIMEOUT}s),"
                "请确认是否已在浏览器完成授权。"
            )
        await loop.run_in_executor(None, server.handle_request)
finally:
    server.server_close()
# state 由 MCP SDK 校验防 CSRF,本 demo 直接回传,不在此处额外校验
return code, state

4) OAuth 提供者:401 后自动完成 服务发现 -> 客户端注册 -> 授权码+PKCE -> 换取 token

oauth = OAuthClientProvider(
server_url=MCP_URL,
client_metadata=OAuthClientMetadata(
client_name="ksadk-oauth-demo",
redirect_uris=[f"http://127.0.0.1:{CALLBACK_PORT}/callback"],
grant_types=["authorization_code", "refresh_token"],
response_types=["code"],
scope="ksc/mcpserver",
token_endpoint_auth_method="none", # 公共客户端 + PKCE,无需 client_secret
),
storage=MemoryStorage(),
redirect_handler=open_browser,
callback_handler=wait_callback,
)

5) 通过 httpx_client_factory 把 OAuth 挂进 MCP 连接

注意:底层调用工厂时会显式传 auth=None,需忽略它、强制使用 OAuth

def oauth_httpx_factory(headers=None, timeout=None, auth=None):
return create_mcp_http_client(headers=headers, timeout=timeout, auth=oauth)

mcp_toolset = McpToolset(
connection_params=StreamableHTTPConnectionParams(
url=MCP_URL,
timeout=300.0,
httpx_client_factory=oauth_httpx_factory,
)
)

agent = Agent(
name="ksc_agent",
model=LiteLlm(model=MODEL_ID),
description="管理金山云资源 Agent",
instruction="你是一个金山云资源管理助手。查询和操作资源必须调用 MCP 工具,不要凭记忆回答。",
tools=[mcp_toolset],
)

async def main():
if not MCP_URL:
raise RuntimeError(
"错误:未配置 MCP 服务地址,请设置环境变量 MCP_URL "
"(例如 export MCP_URL=https://your-mcp-host/mcp)"
)
session_service = InMemorySessionService()
runner = Runner(agent=agent, app_name="oauth_demo", session_service=session_service)
session = await session_service.create_session(app_name="oauth_demo", user_id="demo")

message = types.Content(
    role="user", parts=[types.Part(text="查询 cn-beijing-6 的 EIP 个数")]
)
try:
    async for event in runner.run_async(
        user_id="demo", session_id=session.id, new_message=message
    ):
        if not event.content or not event.content.parts:
            continue
        for part in event.content.parts:
            if part.function_call:  # 模型发起的工具调用
                fc = part.function_call
                print(f"\n[工具调用] {fc.name}  参数: {dict(fc.args or {})}")
            elif part.function_response:  # MCP 工具返回的原始结果
                fr = part.function_response
                text = str(fr.response)
                print(
                    f"[工具返回] {fr.name}  结果: "
                    f"{text[:500]}{' …' if len(text) &gt; 500 else ''}"
                )
            elif part.text and event.is_final_response():  # 最终回答
                print(f"\n[最终回答] {part.text}")
finally:
    await mcp_toolset.close()

if name == "main":
asyncio.run(main())

文档导读
纯净模式常规模式

纯净模式

点击可全屏预览文档内容
文档反馈