From 1453a133694e0206842cf8cab457f5804a544b53 Mon Sep 17 00:00:00 2001 From: joshua Date: Wed, 23 Sep 2026 05:47:43 +0000 Subject: [PATCH] fix(mcp): refresh waiters that observed a newer catalog revision Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> --- .../proxy/_experimental/mcp_server/catalog.py | 2 +- .../mcp_server/test_mcp_server_manager.py | 43 +++++++++++++++++++ 2 files changed, 44 insertions(+), 1 deletion(-) diff --git a/litellm/proxy/_experimental/mcp_server/catalog.py b/litellm/proxy/_experimental/mcp_server/catalog.py index 44cbb3c8203..a24935d6d54 100644 --- a/litellm/proxy/_experimental/mcp_server/catalog.py +++ b/litellm/proxy/_experimental/mcp_server/catalog.py @@ -125,7 +125,7 @@ class TargetCatalog: self._arrival_ticket += 1 arrival: Final = self._arrival_ticket async with self._refresh_lock: - if arrival > self._completed_ticket: + if arrival > self._completed_ticket or revision != self._applied_revision: try: await self._publish_refresh(revision, reuse_unchanged=True) except Exception as exc: diff --git a/tests/test_litellm/proxy/_experimental/mcp_server/test_mcp_server_manager.py b/tests/test_litellm/proxy/_experimental/mcp_server/test_mcp_server_manager.py index eb0e8cba9f8..aead498d408 100644 --- a/tests/test_litellm/proxy/_experimental/mcp_server/test_mcp_server_manager.py +++ b/tests/test_litellm/proxy/_experimental/mcp_server/test_mcp_server_manager.py @@ -15597,3 +15597,46 @@ async def test_catalog_reload_applies_revision_so_next_operation_skips_read(monk assert (await manager.catalog.list())["catalog-server"].name == "initial" read_rows.assert_awaited_once() assert read_revision.await_count == 2 + + +@pytest.mark.asyncio +async def test_catalog_waiter_that_observed_newer_revision_refreshes(monkeypatch): + entered = asyncio.Event() + release = asyncio.Event() + waiter_saw_new_revision = asyncio.Event() + current_revision = [5] + reads = [0] + + async def read_rows(**_kwargs): + reads[0] += 1 + if not entered.is_set(): + entered.set() + await release.wait() + return [_catalog_row()] + if reads[0] == 2: + return [_catalog_row()] + return [ + _catalog_row(), + _catalog_row("added").model_copy(update={"server_id": "added-server"}), + ] + + async def read_current_revision(**_kwargs): + if current_revision[0] == 6: + waiter_saw_new_revision.set() + return _revision_row(current_revision[0]) + + read_rows_mock = AsyncMock(side_effect=read_rows) + read_revision = AsyncMock(side_effect=read_current_revision) + _catalog_database(monkeypatch, read_rows_mock, read_revision) + manager = MCPServerManager() + first = asyncio.create_task(manager.catalog.list()) + await entered.wait() + middle = asyncio.create_task(manager.catalog.list()) + await asyncio.sleep(0) + current_revision[0] = 6 + last = asyncio.create_task(manager.catalog.list()) + await waiter_saw_new_revision.wait() + release.set() + _, _, servers = await asyncio.gather(first, middle, last) + assert "added-server" in servers + assert reads[0] == 3