mirror of
https://github.com/BerriAI/litellm.git
synced 2026-09-14 23:21:35 +00:00
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).
This commit is contained in:
parent
828a44ab26
commit
fcd285facd
6 changed files with 188 additions and 24 deletions
|
|
@ -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,
|
||||
)
|
||||
|
|
|
|||
|
|
@ -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)
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
||||
|
|
|
|||
|
|
@ -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,
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue