From fcd285facdb6cf7452f56b9af888b16b32fb4f38 Mon Sep 17 00:00:00 2001 From: Krrish Dholakia Date: Fri, 10 Jul 2026 21:35:58 -0700 Subject: [PATCH] fix(claude-code): retry transient fetch failures, don't abort sync on one A sync fans out one request per skill folder; on a repo with dozens of skills, a single request occasionally hitting a transient connection blip was expected, not exceptional. Previously any one failure aborted the entire asyncio.gather with an unhandled exception, so admin retries on a big repo were coin flips even after the concurrency fix. Two changes: - _http_get retries transient httpx errors (not real HTTP status codes) with a short backoff before giving up. - A single skill that still fails after retries is skipped rather than fatal, and counted via a new persisted skipped_count column so an admin can actually see whether every skill loaded, instead of an opaque all-or-nothing success/error. Verified live: 10/10 registrations succeeded for both a manifest-based repo (alirezarezvani/claude-skills, 83 skills) and the dir-scan fallback path that was actually flaky before (garrytan/gstack, 53 skills). --- .../claude_code_marketplace_sources.py | 1 + .../claude_code_marketplace_sync.py | 99 ++++++++++++---- litellm/proxy/schema.prisma | 1 + litellm/types/proxy/claude_code_endpoints.py | 1 + ...aude_code_marketplace_sources_endpoints.py | 1 + .../test_claude_code_marketplace_sync.py | 109 ++++++++++++++++++ 6 files changed, 188 insertions(+), 24 deletions(-) diff --git a/litellm/proxy/anthropic_endpoints/claude_code_endpoints/claude_code_marketplace_sources.py b/litellm/proxy/anthropic_endpoints/claude_code_endpoints/claude_code_marketplace_sources.py index d954e648273..54a16c7e028 100644 --- a/litellm/proxy/anthropic_endpoints/claude_code_endpoints/claude_code_marketplace_sources.py +++ b/litellm/proxy/anthropic_endpoints/claude_code_endpoints/claude_code_marketplace_sources.py @@ -113,6 +113,7 @@ def _to_marketplace_source_response(marketplace, plugin_count: Optional[int]) -> sync_error=marketplace.sync_error, last_synced_at=marketplace.last_synced_at.isoformat() if marketplace.last_synced_at else None, plugin_count=plugin_count, + skipped_count=marketplace.skipped_count, created_at=marketplace.created_at.isoformat() if marketplace.created_at else None, updated_at=marketplace.updated_at.isoformat() if marketplace.updated_at else None, ) diff --git a/litellm/proxy/anthropic_endpoints/claude_code_endpoints/claude_code_marketplace_sync.py b/litellm/proxy/anthropic_endpoints/claude_code_endpoints/claude_code_marketplace_sync.py index b788738626b..37e82ba230c 100644 --- a/litellm/proxy/anthropic_endpoints/claude_code_endpoints/claude_code_marketplace_sync.py +++ b/litellm/proxy/anthropic_endpoints/claude_code_endpoints/claude_code_marketplace_sync.py @@ -56,6 +56,14 @@ DEFAULT_SYNC_TIMEOUT_SECONDS = 10.0 _MAX_CONCURRENT_GITHUB_FETCHES = 6 _github_fetch_semaphore = asyncio.Semaphore(_MAX_CONCURRENT_GITHUB_FETCHES) +# A sync fans out one request per skill folder; on a repo with dozens of +# skills, at least one request occasionally hitting a transient connection +# blip is expected, not exceptional. Retrying here means one flaky request +# doesn't cost the whole sync - see _fetch_skill_entry for what happens if +# a request still fails after retries (skipped, not fatal). +_MAX_HTTP_RETRIES = 2 +_RETRY_BACKOFF_SECONDS = 0.3 + SourceHost = Literal["github", "gitlab", "bitbucket", "url"] SyncErrorReason = Literal["unreachable", "http_error", "invalid_json", "invalid_schema"] @@ -102,6 +110,7 @@ class SyncResult: status: Literal["success", "error"] error: str | None plugin_count: int + skipped_count: int = 0 @dataclass(frozen=True, slots=True) @@ -177,16 +186,23 @@ def _git_clone_url(resolved: ResolvedSource) -> str: async def _http_get(client: AsyncHTTPHandler, url: str, *, timeout: float) -> httpx.Response: - try: - # async_safe_get is SSRF-guarded (resolves + validates every redirect - # hop) - required here because the target URL is admin-supplied. - # Bounded by _github_fetch_semaphore: see its module-level comment. - async with _github_fetch_semaphore: - return await async_safe_get(client, url, headers={}, timeout=timeout) - except SSRFError as exc: - raise MarketplaceSyncError(reason="unreachable", detail=str(exc)) from exc - except httpx.HTTPError as exc: - raise MarketplaceSyncError(reason="unreachable", detail=str(exc)) from exc + last_exc: httpx.HTTPError | None = None + for attempt in range(_MAX_HTTP_RETRIES + 1): + try: + # async_safe_get is SSRF-guarded (resolves + validates every + # redirect hop) - required here because the target URL is + # admin-supplied. Bounded by _github_fetch_semaphore: see its + # module-level comment. + async with _github_fetch_semaphore: + return await async_safe_get(client, url, headers={}, timeout=timeout) + except SSRFError as exc: + # A policy rejection, not a transient failure - retrying won't help. + raise MarketplaceSyncError(reason="unreachable", detail=str(exc)) from exc + except httpx.HTTPError as exc: + last_exc = exc + if attempt < _MAX_HTTP_RETRIES: + await asyncio.sleep(_RETRY_BACKOFF_SECONDS * (attempt + 1)) + raise MarketplaceSyncError(reason="unreachable", detail=str(last_exc)) from last_exc def _parse_json_body(response: httpx.Response) -> object: @@ -385,7 +401,18 @@ async def _list_github_contents( async def _skill_md_exists(client: AsyncHTTPHandler, repo: str, branch: str, dir_path: str, *, timeout: float) -> bool: url = f"https://raw.githubusercontent.com/{repo}/{branch}/{dir_path}/SKILL.md" - response = await _http_get(client, url, timeout=timeout) + try: + response = await _http_get(client, url, timeout=timeout) + except MarketplaceSyncError as exc: + # This is a discovery-phase existence check, not a confirmed skill's + # content fetch - after _http_get's own retries are exhausted, treat + # it as "not found" rather than aborting the entire sync over one + # candidate folder. _fetch_skill_entry is where a CONFIRMED skill's + # failure gets tracked and surfaced via skipped_count. + verbose_proxy_logger.warning( + "skill-marketplace-sync: could not check %s after retries (%s), treating as not found", url, exc + ) + return False return response.status_code == 200 @@ -498,7 +525,17 @@ async def _fetch_skill_entry( timeout: float, ) -> ResolvedPluginEntry | None: url = f"https://raw.githubusercontent.com/{repo}/{branch}/{doc.skill_md_path}" - response = await _http_get(client, url, timeout=timeout) + try: + response = await _http_get(client, url, timeout=timeout) + except MarketplaceSyncError as exc: + # A CONFIRMED skill (its SKILL.md's existence already passed a + # separate check) still failed to fetch after _http_get's retries. + # Skip it rather than aborting every other skill's import - the + # caller counts this against skipped_count so it's visible, not silent. + verbose_proxy_logger.warning( + "skill-marketplace-sync: could not fetch %s after retries (%s), skipping this skill", url, exc + ) + return None if response.status_code != 200: return None @@ -524,19 +561,21 @@ async def _fetch_entries_for_docs( docs: tuple[_DiscoveredSkillDoc, ...], *, timeout: float, -) -> tuple[ResolvedPluginEntry, ...]: +) -> tuple[tuple[ResolvedPluginEntry, ...], int]: + """Returns (successfully-fetched entries, count skipped after retries).""" fetched = await asyncio.gather( *( _fetch_skill_entry(client, repo, branch, marketplace_name, doc, timeout=timeout) for doc in docs ) ) - return tuple(entry for entry in fetched if entry is not None) + entries = tuple(entry for entry in fetched if entry is not None) + return entries, len(docs) - len(entries) async def _fetch_github_fallback_entries( client: AsyncHTTPHandler, resolved: ResolvedSource, branch: str, marketplace_name: str -) -> tuple[tuple[ResolvedPluginEntry, ...], MarketplaceSourceType]: +) -> tuple[tuple[ResolvedPluginEntry, ...], MarketplaceSourceType, int]: """Called once the repo has no ``.claude-plugin/marketplace.json``. Tries, in order: an explicit single-plugin ``.claude-plugin/plugin.json`` skill list, the ``skills/*/SKILL.md`` (or one level of category nesting) @@ -551,22 +590,24 @@ async def _fetch_github_fallback_entries( plugin_manifest = _parse_plugin_manifest(plugin_json_response) plugin_docs = _plugin_manifest_skill_docs(plugin_manifest) if plugin_docs: - entries = await _fetch_entries_for_docs( + entries, skipped = await _fetch_entries_for_docs( client, repo, branch, marketplace_name, plugin_docs, timeout=timeout ) if entries: - return entries, "claude_plugin_json" + return entries, "claude_plugin_json", skipped skills_dir_docs = await _discover_github_skill_docs(client, repo, branch, timeout=timeout) if not skills_dir_docs: skills_dir_docs = await _discover_root_skill_docs(client, repo, branch, timeout=timeout) - entries = await _fetch_entries_for_docs(client, repo, branch, marketplace_name, skills_dir_docs, timeout=timeout) - return entries, "skills_dir" + entries, skipped = await _fetch_entries_for_docs( + client, repo, branch, marketplace_name, skills_dir_docs, timeout=timeout + ) + return entries, "skills_dir", skipped async def _fetch_marketplace_entries( marketplace_row: MarketplaceRow, -) -> tuple[tuple[ResolvedPluginEntry, ...], MarketplaceSourceType]: +) -> tuple[tuple[ResolvedPluginEntry, ...], MarketplaceSourceType, int]: if not marketplace_row.source_ref: raise MarketplaceSyncError(reason="invalid_schema", detail="marketplace has no source_ref to sync from") @@ -583,7 +624,7 @@ async def _fetch_marketplace_entries( if manifest_response.status_code == 200: manifest = _parse_marketplace_manifest(manifest_response) entries = _build_manifest_plugin_entries(marketplace_row.name, resolved, manifest) - return entries, "claude_marketplace_json" + return entries, "claude_marketplace_json", len(manifest.plugins) - len(entries) if manifest_response.status_code == 404 and resolved.host == "github": return await _fetch_github_fallback_entries(client, resolved, branch, marketplace_row.name) @@ -672,6 +713,7 @@ async def _record_sync_success( prisma_client: Any, # noqa: ANN401 # prisma_client has no importable type stubs, see _upsert_plugin_entries marketplace_row: MarketplaceRow, source_type: MarketplaceSourceType, + skipped_count: int, ) -> None: await SkillMarketplaceRepository(prisma_client).table.update( where={"id": marketplace_row.id}, @@ -679,6 +721,7 @@ async def _record_sync_success( "sync_status": "success", "sync_error": None, "source_type": source_type, + "skipped_count": skipped_count, "last_synced_at": datetime.now(timezone.utc), }, ) @@ -709,13 +752,21 @@ async def resolve_and_sync( (``sync_status="error"``) and reflected in the returned ``SyncResult``. """ try: - entries, source_type = await _fetch_marketplace_entries(marketplace_row) + entries, source_type, skipped_count = await _fetch_marketplace_entries(marketplace_row) except MarketplaceSyncError as exc: verbose_proxy_logger.warning("skill-marketplace-sync: failed to sync %r: %s", marketplace_row.name, exc) await _record_sync_failure(prisma_client, marketplace_row, str(exc)) return SyncResult(status="error", error=str(exc), plugin_count=0) + if skipped_count: + verbose_proxy_logger.warning( + "skill-marketplace-sync: %r imported %d skill(s), skipped %d after retries", + marketplace_row.name, + len(entries), + skipped_count, + ) + await _upsert_plugin_entries(prisma_client, marketplace_row.id, entries) await _soft_disable_stale_plugins(prisma_client, marketplace_row.id, entries) - await _record_sync_success(prisma_client, marketplace_row, source_type) - return SyncResult(status="success", error=None, plugin_count=len(entries)) + await _record_sync_success(prisma_client, marketplace_row, source_type, skipped_count) + return SyncResult(status="success", error=None, plugin_count=len(entries), skipped_count=skipped_count) diff --git a/litellm/proxy/schema.prisma b/litellm/proxy/schema.prisma index 1439861ee09..220632ad5a2 100644 --- a/litellm/proxy/schema.prisma +++ b/litellm/proxy/schema.prisma @@ -1284,6 +1284,7 @@ model LiteLLM_SkillMarketplaceTable { enabled Boolean @default(true) sync_status String @default("pending") // "pending" | "success" | "error" sync_error String? + skipped_count Int @default(0) // skills that failed to fetch after retries during the last sync last_synced_at DateTime? created_at DateTime? @default(now()) updated_at DateTime? @default(now()) @updatedAt diff --git a/litellm/types/proxy/claude_code_endpoints.py b/litellm/types/proxy/claude_code_endpoints.py index d79fe63c647..38722da9e44 100644 --- a/litellm/types/proxy/claude_code_endpoints.py +++ b/litellm/types/proxy/claude_code_endpoints.py @@ -180,6 +180,7 @@ class MarketplaceSourceResponse(BaseModel): sync_error: Optional[str] = None last_synced_at: Optional[str] = None plugin_count: Optional[int] = None + skipped_count: int = 0 created_at: Optional[str] = None updated_at: Optional[str] = None diff --git a/tests/test_litellm/proxy/anthropic_endpoints/test_claude_code_marketplace_sources_endpoints.py b/tests/test_litellm/proxy/anthropic_endpoints/test_claude_code_marketplace_sources_endpoints.py index 5a880837150..184918177ec 100644 --- a/tests/test_litellm/proxy/anthropic_endpoints/test_claude_code_marketplace_sources_endpoints.py +++ b/tests/test_litellm/proxy/anthropic_endpoints/test_claude_code_marketplace_sources_endpoints.py @@ -63,6 +63,7 @@ class _FakeTable: "branch": "main", "enabled": True, "sync_error": None, + "skipped_count": 0, "last_synced_at": None, "created_at": None, "updated_at": None, diff --git a/tests/test_litellm/proxy/anthropic_endpoints/test_claude_code_marketplace_sync.py b/tests/test_litellm/proxy/anthropic_endpoints/test_claude_code_marketplace_sync.py index f217d5e2431..fd73317e50b 100644 --- a/tests/test_litellm/proxy/anthropic_endpoints/test_claude_code_marketplace_sync.py +++ b/tests/test_litellm/proxy/anthropic_endpoints/test_claude_code_marketplace_sync.py @@ -7,6 +7,7 @@ GitHub-only skills/ directory-scan fallback), error classification, and the idempotent/soft-disable upsert semantics. """ +import asyncio import json import uuid from datetime import datetime, timedelta @@ -263,6 +264,114 @@ async def test_resolve_and_sync_falls_back_to_skills_directory_scan(monkeypatch) } +@pytest.mark.asyncio +async def test_resolve_and_sync_retries_transient_failures_then_succeeds(monkeypatch): + """Regression test: a request that fails with a connection error on its + first attempt but succeeds on retry must not be treated as a permanent + failure - this is the exact shape of the flakiness a single unlucky + request among many concurrent fetches used to cause. Uses the contents + listing (a single call site) so the retry count is unambiguous, unlike + a SKILL.md URL which is hit once for discovery and again for content.""" + client = _make_fake_prisma_client() + marketplace = await _create_marketplace(client, name="vercel-skills", source_ref="vercel-labs/skills") + + manifest_url = "https://raw.githubusercontent.com/vercel-labs/skills/main/.claude-plugin/marketplace.json" + plugin_json_url = "https://raw.githubusercontent.com/vercel-labs/skills/main/.claude-plugin/plugin.json" + contents_url = "https://api.github.com/repos/vercel-labs/skills/contents/skills?ref=main" + skill_md_url = "https://raw.githubusercontent.com/vercel-labs/skills/main/skills/find-skills/SKILL.md" + + contents_attempts = {"count": 0} + + async def _get(http_client, url, **kwargs): + if url == manifest_url: + return httpx.Response(404) + if url == plugin_json_url: + return httpx.Response(404) + if url == contents_url: + contents_attempts["count"] += 1 + if contents_attempts["count"] == 1: + raise httpx.ConnectError("connection reset", request=httpx.Request("GET", url)) + return httpx.Response( + 200, + json=[{"name": "find-skills", "path": "skills/find-skills", "type": "dir"}], + ) + if url == skill_md_url: + return httpx.Response(200, text=_FIND_SKILLS_SKILL_MD) + raise AssertionError(f"unexpected url requested: {url}") + + monkeypatch.setattr(sync_module, "async_safe_get", _get) + _real_sleep = asyncio.sleep + monkeypatch.setattr(sync_module.asyncio, "sleep", lambda *_a, **_kw: _real_sleep(0)) + + result = await resolve_and_sync(client, marketplace) + + assert result.status == "success" + assert result.plugin_count == 1 + assert result.skipped_count == 0 + assert contents_attempts["count"] == 2, "expected exactly one retry after the first transient failure" + + +@pytest.mark.asyncio +async def test_resolve_and_sync_skips_skill_that_fails_after_retries_exhausted(monkeypatch): + """Regression test: a skill whose SKILL.md is confirmed to exist (discovery + succeeds) but whose content fetch keeps failing even after retries must be + skipped (and counted via skipped_count), not abort the whole sync - the + other, healthy skill must still import successfully.""" + client = _make_fake_prisma_client() + marketplace = await _create_marketplace(client, name="vercel-skills", source_ref="vercel-labs/skills") + + manifest_url = "https://raw.githubusercontent.com/vercel-labs/skills/main/.claude-plugin/marketplace.json" + plugin_json_url = "https://raw.githubusercontent.com/vercel-labs/skills/main/.claude-plugin/plugin.json" + contents_url = "https://api.github.com/repos/vercel-labs/skills/contents/skills?ref=main" + healthy_skill_md_url = "https://raw.githubusercontent.com/vercel-labs/skills/main/skills/find-skills/SKILL.md" + flaky_skill_md_url = "https://raw.githubusercontent.com/vercel-labs/skills/main/skills/flaky-skill/SKILL.md" + + # First call to flaky_skill_md_url is the discovery-phase existence check + # (must succeed so this candidate is confirmed as a real skill); every + # call after that is the content-fetch phase, which keeps failing. + flaky_attempts = {"count": 0} + + async def _get(http_client, url, **kwargs): + if url == manifest_url: + return httpx.Response(404) + if url == plugin_json_url: + return httpx.Response(404) + if url == contents_url: + return httpx.Response( + 200, + json=[ + {"name": "find-skills", "path": "skills/find-skills", "type": "dir"}, + {"name": "flaky-skill", "path": "skills/flaky-skill", "type": "dir"}, + ], + ) + if url == healthy_skill_md_url: + return httpx.Response(200, text=_FIND_SKILLS_SKILL_MD) + if url == flaky_skill_md_url: + flaky_attempts["count"] += 1 + if flaky_attempts["count"] == 1: + return httpx.Response(200, text=_FIND_SKILLS_SKILL_MD) + raise httpx.ConnectError("connection reset", request=httpx.Request("GET", url)) + raise AssertionError(f"unexpected url requested: {url}") + + monkeypatch.setattr(sync_module, "async_safe_get", _get) + _real_sleep = asyncio.sleep + monkeypatch.setattr(sync_module.asyncio, "sleep", lambda *_a, **_kw: _real_sleep(0)) + + result = await resolve_and_sync(client, marketplace) + + assert result.status == "success" + assert result.plugin_count == 1 + assert result.skipped_count == 1 + + plugins = await client.db.litellm_claudecodeplugintable.find_many(where={"marketplace_id": marketplace.id}) + assert len(plugins) == 1 + assert plugins[0].name == "vercel-skills--find-skills" + + refreshed_marketplace = await client.db.litellm_skillmarketplacetable.find_unique(where={"id": marketplace.id}) + assert refreshed_marketplace.sync_status == "success" + assert refreshed_marketplace.skipped_count == 1 + + @pytest.mark.asyncio async def test_resolve_and_sync_reads_plugin_json_explicit_skill_list(monkeypatch): """Regression test: a repo with no marketplace.json but a single-plugin