Skip to content

Clients Using Storage Face Deadlock on Token Refresh for SSE #1326

Description

@Norcim133

Initial Checks

Description

tldr: in calls to an SSE server (like Atlassian) there is an infinite hang druing a request, caused by a deadlock; the deadlock occurs when token refresh yields a request in the middle of the client.stream method that is trying to monitor SSE events, not field requests

Resulting Behavior - After tokens expire or after a new instance is re-built from storage, the SSE server will not work until full token-clearing and re-auth

As long as a token refresh attempt is not made, the server will work.

There are two circumstances when a refresh attempt is made, causing the SSE call to hang:

  1. Tokens expire during a working SSE session
  2. When a Client uses a stale token from storage

Source of failure for each scenario

Failure Mode

  1. OAuthClientProvider lazy loads and is passed into aconnect_sse() without having checked refresh
  2. The OAuthClientProvider.async_auth_flow() generator enters the refresh logic
  3. It yields a refresh request and waits for a response (line ~544 in auth.py)
  4. aconnect_sse() uses client.stream() which is trying to open an SSE connection
  5. There's no response mechanism - SSE expects to start streaming events, not handle a token refresh request
  6. DEADLOCK: Auth generator waiting for refresh response that will never come

Temporary Success

  • First request in aconnect_sse() uses client.stream() yield 401 which sends everything back out of the SSE context
  • OAuthClientProvider is able to perform full re-auth
  • When next sent to aconnect_sse() fresh auth headers are availble and SSE stream succeeds
  • SSE is long-lived so works for as long as tokens don't need refresh

Root cause

In /mcp/client/sse.py, this call with aconnect_sse hangs infinitely.

async with anyio.create_task_group() as tg:
    try:
        logger.debug(f"Connecting to SSE endpoint: {remove_request_params(url)}")
        async with httpx_client_factory(
            headers=headers, auth=auth, timeout=httpx.Timeout(timeout, read=sse_read_timeout)
        ) as client:
            async with aconnect_sse(  #<======================== This hangs indefinitely with a 200 response
                client,
                "GET",
                url,
            ) as event_source:
                event_source.response.raise_for_status()

This is because OAuthClientProvider, which is attached to the client, on token refresh, tries to yield when in the middle of client.stream here:

            if not self.context.is_token_valid() and self.context.can_refresh_token():
                # Try to refresh token
                refresh_request = await self._refresh_token()
                refresh_response = yield refresh_request. #<======================== Causes client.stream() to hang

                if not await self._handle_refresh_response(refresh_response):
                    # Refresh failed, need full re-authentication
                    self._initialized = False

Other yields in the broader method occur after the stream got a 401 and so are not in the same context.

The Fix - Force token refresh check outside of the stream context

Any fix needs to remove the OAuthClientProvider from client before calling aconnect_sse. That way it won't yield the refresh request.

There is probably a more elgant/direct approach but the below worked. I had to create a separate client to avoid some lock conflict.

async with anyio.create_task_group() as tg:
    try:
        logger.debug(f"Connecting to SSE endpoint: {remove_request_params(url)}")
        async with httpx_client_factory(
            headers=headers, auth=auth, timeout=httpx.Timeout(timeout, read=sse_read_timeout)
        ) as client:
            
            #----- Start of fix -------
            auth_headers = {}
            if auth:
                # Initialize auth to load stored tokens
                if hasattr(auth, '_initialize') and not getattr(auth, '_initialized', False):
                    await auth._initialize()

                # Check if tokens need refresh
                if hasattr(auth, 'context') and auth.context:
                    # Check if token is expired or needs refresh
                    if not auth.context.is_token_valid() and auth.context.can_refresh_token():
                        logger.debug("Token needs refresh before SSE connection")

                        # Use a separate client for refresh to avoid lock issues
                        async with httpx.AsyncClient() as refresh_client:
                            # Create a dummy request to trigger the refresh flow
                            dummy_request = httpx.Request("GET", url)

                            # Run the auth flow to refresh tokens
                            auth_gen = auth.async_auth_flow(dummy_request)
                            try:
                                # Start the generator
                                refresh_request = await auth_gen.asend(None)

                                # Execute the refresh request
                                refresh_response = await refresh_client.request(
                                    refresh_request.method,
                                    refresh_request.url,
                                    data=refresh_request.content,
                                    headers=refresh_request.headers,
                                )

                                # Send response back to generator
                                await auth_gen.asend(refresh_response)
                                logger.debug("Token refreshed successfully before SSE connection")
                            except StopAsyncIteration:
                                # Normal completion of auth flow
                                pass
                            except Exception as e:
                                logger.error(f"Failed to refresh token before SSE: {e}")

                    # Extract the token after potential refresh
                    if auth.context.current_tokens and auth.context.current_tokens.access_token:
                        token = auth.context.current_tokens.access_token
                        auth_headers['Authorization'] = f"Bearer {token}"

            # CRITICAL: Remove auth from client to prevent generator interference
            client.auth = None
            #----- Resume existing code -------
            async with aconnect_sse(
                client,
                "GET",
                url,
                headers=auth_headers,    #<============================== Must add this too
            ) as event_source:
                event_source.response.raise_for_status()

Since you removed auth from the client, you have to insert the headers later for the POST call as well:

                    async def post_writer(endpoint_url: str):
                        try:
                            async with write_stream_reader:
                                async for session_message in write_stream_reader:
                                    logger.debug(f"Sending client message: {session_message}")
                                    response = await client.post(
                                        endpoint_url,
                                        json=session_message.message.model_dump(
                                            by_alias=True,
                                            mode="json",
                                            exclude_none=True,
                                        ),
                                        headers=auth_headers,  #<======================== Cause client.stream() to hang

Example Code

Python & MCP Python SDK

1.12.4 (but same code as current version)

Activity

  1. someposer commented on Sep 3, 2025

    @someposer

    I am experiencing the same issue. When I try to reuse a token, the client hangs inside of aconnect_sse.

    Below is some sample code that uses the Simple Auth Client Example code as a base to demonstrate the issue.

    #!/usr/bin/env python3
    """
    Simple MCP client example with OAuth authentication support.
    
    This client connects to an MCP server using streamable HTTP transport with OAuth.
    
    """
    
    import asyncio
    import threading
    import time
    import webbrowser
    from http.server import BaseHTTPRequestHandler, HTTPServer
    from urllib.parse import parse_qs, urlparse
    
    from mcp.client.auth import OAuthClientProvider, TokenStorage
    from mcp.client.session import ClientSession
    from mcp.client.sse import sse_client
    from mcp.shared.auth import OAuthClientInformationFull, OAuthClientMetadata, OAuthToken
    
    
    class InMemoryTokenStorage(TokenStorage):
        """Simple in-memory token storage implementation."""
    
        def __init__(self):
            self._tokens: OAuthToken | None = None
            self._client_info: OAuthClientInformationFull | None = None
    
        async def get_tokens(self) -> OAuthToken | None:
            return self._tokens
    
        async def set_tokens(self, tokens: OAuthToken) -> None:
            self._tokens = tokens
    
        async def get_client_info(self) -> OAuthClientInformationFull | None:
            return self._client_info
    
        async def set_client_info(self, client_info: OAuthClientInformationFull) -> None:
            self._client_info = client_info
    
    
    token_storage = InMemoryTokenStorage()
    
    
    class CallbackHandler(BaseHTTPRequestHandler):
        """Simple HTTP handler to capture OAuth callback."""
    
        def __init__(self, request, client_address, server, callback_data):
            """Initialize with callback data storage."""
            self.callback_data = callback_data
            super().__init__(request, client_address, server)
    
        def do_GET(self):
            """Handle GET request from OAuth redirect."""
            parsed = urlparse(self.path)
            query_params = parse_qs(parsed.query)
    
            if "code" in query_params:
                self.callback_data["authorization_code"] = query_params["code"][0]
                self.callback_data["state"] = query_params.get("state", [None])[0]
                self.send_response(200)
                self.send_header("Content-type", "text/html")
                self.end_headers()
                self.wfile.write(b"""
                <html>
                <body>
                    <h1>Authorization Successful!</h1>
                    <p>You can close this window and return to the terminal.</p>
                    <script>setTimeout(() => window.close(), 2000);</script>
                </body>
                </html>
                """)
            elif "error" in query_params:
                self.callback_data["error"] = query_params["error"][0]
                self.send_response(400)
                self.send_header("Content-type", "text/html")
                self.end_headers()
                self.wfile.write(
                    f"""
                <html>
                <body>
                    <h1>Authorization Failed</h1>
                    <p>Error: {query_params["error"][0]}</p>
                    <p>You can close this window and return to the terminal.</p>
                </body>
                </html>
                """.encode()
                )
            else:
                self.send_response(404)
                self.end_headers()
    
        def log_message(self, format, *args):
            """Suppress default logging."""
            pass
    
    
    class CallbackServer:
        """Simple server to handle OAuth callbacks."""
    
        def __init__(self, port=3000):
            self.port = port
            self.server = None
            self.thread = None
            self.callback_data = {"authorization_code": None, "state": None, "error": None}
    
        def _create_handler_with_data(self):
            """Create a handler class with access to callback data."""
            callback_data = self.callback_data
    
            class DataCallbackHandler(CallbackHandler):
                def __init__(self, request, client_address, server):
                    super().__init__(request, client_address, server, callback_data)
    
            return DataCallbackHandler
    
        def start(self):
            """Start the callback server in a background thread."""
            handler_class = self._create_handler_with_data()
            self.server = HTTPServer(("localhost", self.port), handler_class)
            self.thread = threading.Thread(target=self.server.serve_forever, daemon=True)
            self.thread.start()
            print(f"🖥️  Started callback server on http://localhost:{self.port}")
    
        def stop(self):
            """Stop the callback server."""
            if self.server:
                self.server.shutdown()
                self.server.server_close()
            if self.thread:
                self.thread.join(timeout=1)
    
        def wait_for_callback(self, timeout=300):
            """Wait for OAuth callback with timeout."""
            start_time = time.time()
            while time.time() - start_time < timeout:
                if self.callback_data["authorization_code"]:
                    return self.callback_data["authorization_code"]
                elif self.callback_data["error"]:
                    raise Exception(f"OAuth error: {self.callback_data['error']}")
                time.sleep(0.1)
            raise Exception("Timeout waiting for OAuth callback")
    
        def get_state(self):
            """Get the received state parameter."""
            return self.callback_data["state"]
    
    
    class SimpleAuthClient:
        """Simple MCP client with auth support (SSE only)."""
    
        def __init__(self, server_url: str):
            self.server_url = server_url
            self.session: ClientSession | None = None
    
        async def connect(self):
            """Connect to the MCP server using SSE transport only."""
            print(f"🔗 Attempting to connect to {self.server_url}...")
    
            try:
                callback_server = CallbackServer(port=3030)
                callback_server.start()
    
                async def callback_handler() -> tuple[str, str | None]:
                    """Wait for OAuth callback and return auth code and state."""
                    print("⏳ Waiting for authorization callback...")
                    try:
                        auth_code = callback_server.wait_for_callback(timeout=300)
                        return auth_code, callback_server.get_state()
                    finally:
                        callback_server.stop()
    
                client_metadata_dict = {
                    "client_name": "Simple Auth Client",
                    "redirect_uris": ["http://localhost:3030/callback"],
                    "grant_types": ["authorization_code", "refresh_token"],
                    "response_types": ["code"],
                    "token_endpoint_auth_method": "client_secret_post",
                }
    
                async def _default_redirect_handler(authorization_url: str) -> None:
                    """Default redirect handler that opens the URL in a browser."""
                    print(f"Opening browser for authorization: {authorization_url}")
                    webbrowser.open(authorization_url)
    
                # Create OAuth authentication handler using the new interface
                server_base_url = self.server_url.replace("/sse", "")
                oauth_auth = OAuthClientProvider(
                    server_url=server_base_url,
                    client_metadata=OAuthClientMetadata.model_validate(client_metadata_dict),
                    storage=token_storage,
                    redirect_handler=_default_redirect_handler,
                    callback_handler=callback_handler,
                )
    
                # Only SSE transport is supported in this version
                print("📡 Opening SSE transport connection with auth...")
                async with sse_client(
                    url=self.server_url,
                    auth=oauth_auth,
                    timeout=60,
                ) as (read_stream, write_stream):
                    await self._run_session(read_stream, write_stream, None)
    
            except Exception as e:
                print(f"❌ Failed to connect: {e}")
                import traceback
    
                traceback.print_exc()
    
        async def _run_session(self, read_stream, write_stream, get_session_id):
            """Run the MCP session with the given streams."""
            print("🤝 Initializing MCP session...")
            async with ClientSession(read_stream, write_stream) as session:
                self.session = session
                print("⚡ Starting session initialization...")
                await session.initialize()
                print("✨ Session initialization complete!")
    
                print(f"\n✅ Connected to MCP server at {self.server_url}")
                if get_session_id:
                    session_id = get_session_id()
                    if session_id:
                        print(f"Session ID: {session_id}")
    
                # Count the tools
                print("🔄 Counting available tools...")
                result = await self.session.list_tools()
                if hasattr(result, "tools") and result.tools:
                    tool_count = len(result.tools)
                    print(f"🛠️  Available tools: {tool_count}")
                else:
                    print("No tools available")
    
    
    async def main():
        """Main entry point."""
        server_url = "https://mcp.example.com/sse"
    
        print("🚀 Simple MCP Auth Client")
        print(f"Connecting to: {server_url}")
    
        # Start connection flow - OAuth will be handled automatically
        client = SimpleAuthClient(server_url)
        await client.connect()
    
        print(f"\n{'-' * 30}")
        print("Create a new client and connect again, reusing the token")
        client = SimpleAuthClient(server_url)
        await client.connect()  # !!! Notice that this will hang when trying to connect with auth...
        print("Token reuse was successful!")
    
    
    def cli():
        """CLI entry point for uv script."""
        asyncio.run(main())
    
    
    if __name__ == "__main__":
        cli()
  2. delatt commented on Sep 22, 2025

    @delatt

    Just adding that I also run into this issue and need to perform the re-auth every time

  3. added
    bugSomething isn't working
    ready for workEnough information for someone to start working on
    P2Moderate issues affecting some users, edge cases, potentially valuable feature
    and removed
    needs confirmationNeeds confirmation that the PR is actually required or needed.
    on Oct 3, 2025
  4. peisuke commented on May 22, 2026

    @peisuke

    Hi maintainers — happy to take this if it's available.

    The fix I'd like to propose narrows the lock scope in
    OAuthClientProvider.async_auth_flow so the response = yield request
    runs outside any lock (a GET SSE long-poll then no longer holds
    context.lock for its entire lifetime). Token refresh stays
    single-flight via a new dedicated refresh_lock.

    I have the change + tests ready locally on a fork. Should I open a PR
    against main, or would a v1.x target be preferred?

  5. alibeg-begow commented on Aug 7, 2026

    @alibeg-begow

    Not a direct fix for the SSE deadlock, but related: if anyone's auth_flow refresh logic also needs to survive concurrent callers (not just the SSE stream conflict), I built a small double-checked-locking coordinator for httpx/requests auth refresh: https://github.com/alibeg-begow/singleflight_auth — might be useful groundwork for whatever the eventual fix here looks like.

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

No one assigned

    Labels

    P2Moderate issues affecting some users, edge cases, potentially valuable featurebugSomething isn't workingready for workEnough information for someone to start working on

    Type

    No type

    Projects

    No projects

      Relationships

      None yet

      Development

      No branches or pull requests

      Issue actions