Skip to content
Open
25 changes: 19 additions & 6 deletions lib/crewai-files/src/crewai_files/cache/cleanup.py
Original file line number Diff line number Diff line change
Expand Up @@ -17,6 +17,19 @@
logger = logging.getLogger(__name__)


def _uploader_or_none(provider: ProviderType) -> FileUploader | None:
"""Return the provider's uploader, or None when it is unavailable.

get_uploader raises ValueError for an unknown or unconfigured provider.
Cleanup skips such providers rather than aborting the whole pass, so that
error is treated as "no uploader available" here.
"""
try:
return get_uploader(provider)
except ValueError:
return None


def _safe_delete(
uploader: FileUploader,
file_id: str,
Expand Down Expand Up @@ -70,7 +83,7 @@ def cleanup_uploaded_files(

if delete_from_provider:
for provider, uploads in provider_uploads.items():
uploader = get_uploader(provider)
uploader = _uploader_or_none(provider)
if uploader is None:
logger.warning(
f"No uploader available for {provider}, skipping cleanup"
Expand Down Expand Up @@ -116,7 +129,7 @@ def cleanup_expired_files(

if delete_from_provider:
for upload in expired_entries:
uploader = get_uploader(upload.provider)
uploader = _uploader_or_none(upload.provider)
if uploader is not None:
try:
uploader.delete(upload.file_id)
Expand Down Expand Up @@ -144,7 +157,7 @@ def cleanup_provider_files(
Number of files deleted.
"""
deleted = 0
uploader = get_uploader(provider)
uploader = _uploader_or_none(provider)

if uploader is None:
logger.warning(f"No uploader available for {provider}")
Expand Down Expand Up @@ -247,7 +260,7 @@ async def delete_one(file_uploader: FileUploader, cached: CachedUpload) -> bool:

tasks: list[asyncio.Task[bool]] = []
for provider, uploads in provider_uploads.items():
uploader = get_uploader(provider)
uploader = _uploader_or_none(provider)
if uploader is None:
logger.warning(
f"No uploader available for {provider}, skipping cleanup"
Expand Down Expand Up @@ -298,7 +311,7 @@ async def acleanup_expired_files(
async def delete_expired(cached: CachedUpload) -> None:
"""Delete an expired file with semaphore limiting."""
async with semaphore:
file_uploader = get_uploader(cached.provider)
file_uploader = _uploader_or_none(cached.provider)
if file_uploader is not None:
try:
await file_uploader.adelete(cached.file_id)
Expand Down Expand Up @@ -334,7 +347,7 @@ async def acleanup_provider_files(
Number of files deleted.
"""
deleted = 0
uploader = get_uploader(provider)
uploader = _uploader_or_none(provider)

if uploader is None:
logger.warning(f"No uploader available for {provider}")
Expand Down
30 changes: 14 additions & 16 deletions lib/crewai-files/src/crewai_files/resolution/resolver.py
Original file line number Diff line number Diff line change
Expand Up @@ -307,10 +307,6 @@ def _resolve_via_upload(
)

uploader = self._get_uploader(provider)
if uploader is None:
logger.debug(f"No uploader available for {provider}")
return None

result = self._upload_with_retry(uploader, file, provider, context.size)
if result is None:
return None
Expand Down Expand Up @@ -483,6 +479,12 @@ async def resolve_single(

output: dict[str, ResolvedFile] = {}
for item in gather_results:
# A lookup failure (unknown provider, unconfigured Bedrock, or a
# missing provider SDK) applies to every file in the batch, since
# they share one provider. Surface it instead of silently dropping
# files, matching the sync resolve_files path.
if isinstance(item, (ValueError, ImportError)):
raise item
if isinstance(item, BaseException):
logger.error(f"Resolution failed: {item}")
continue
Expand Down Expand Up @@ -524,10 +526,6 @@ async def _aresolve_via_upload(
)

uploader = self._get_uploader(provider)
if uploader is None:
logger.debug(f"No uploader available for {provider}")
return None

result = await self._aupload_with_retry(uploader, file, provider, context.size)
if result is None:
return None
Expand Down Expand Up @@ -612,23 +610,23 @@ async def _aupload_with_retry(
)
return None

def _get_uploader(self, provider: ProviderType) -> FileUploader | None:
def _get_uploader(self, provider: ProviderType) -> FileUploader:
"""Get or create an uploader for a provider.

Args:
provider: Provider name.

Returns:
FileUploader instance or None if not available.
FileUploader instance for the provider.

Raises:
ValueError: If the provider is unknown or not configured.
ImportError: If the provider's SDK is not installed.
"""
if provider not in self._uploaders:
uploader = get_uploader(provider)
if uploader is not None:
self._uploaders[provider] = uploader
else:
return None
self._uploaders[provider] = get_uploader(provider)

return self._uploaders.get(provider)
return self._uploaders[provider]

def get_cached_uploads(self, provider: ProviderType) -> list[CachedUpload]:
"""Get all cached uploads for a provider.
Expand Down
25 changes: 15 additions & 10 deletions lib/crewai-files/src/crewai_files/uploaders/factory.py
Original file line number Diff line number Diff line change
Expand Up @@ -134,7 +134,12 @@ def get_uploader(
**kwargs: Additional arguments passed to the uploader constructor.

Returns:
FileUploader instance for the provider, or None if not supported.
FileUploader instance for the provider.

Raises:
ValueError: If the provider is unknown, or Bedrock is selected without a
configured S3 bucket (CREWAI_BEDROCK_S3_BUCKET or bucket_name).
ImportError: If the selected provider's SDK is not installed.
"""
provider_lower = provider.lower()

Expand Down Expand Up @@ -188,15 +193,13 @@ def get_uploader(
if "bedrock" in provider_lower or "aws" in provider_lower:
import os

if (
not os.environ.get("CREWAI_BEDROCK_S3_BUCKET")
and "bucket_name" not in kwargs
if not os.environ.get("CREWAI_BEDROCK_S3_BUCKET") and not kwargs.get(
"bucket_name"
):
logger.debug(
"Bedrock S3 uploader not configured. "
"Set CREWAI_BEDROCK_S3_BUCKET environment variable to enable."
raise ValueError(
"Bedrock file uploads are not configured. Set the "
"CREWAI_BEDROCK_S3_BUCKET environment variable or pass bucket_name."
)
raise
try:
from crewai_files.uploaders.bedrock import BedrockFileUploader

Expand All @@ -212,5 +215,7 @@ def get_uploader(
logger.warning("boto3 not installed. Install with: pip install boto3")
raise

logger.debug(f"No file uploader available for provider: {provider}")
raise
raise ValueError(
f"No file uploader available for provider: {provider!r}. Supported "
"providers: gemini/google, anthropic/claude, openai/gpt/azure, bedrock/aws."
)
39 changes: 39 additions & 0 deletions lib/crewai-files/tests/test_factory.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,39 @@
"""Tests for get_uploader."""

from crewai_files.uploaders import get_uploader
from crewai_files.uploaders.openai import OpenAIFileUploader
import pytest


def test_get_uploader_returns_uploader_for_configured_provider():
# Happy path: a configured provider returns its uploader instance rather
# than raising or returning None
uploader = get_uploader("openai", api_key="test-key")
assert isinstance(uploader, OpenAIFileUploader)


def test_get_uploader_raises_for_unknown_provider():
# Regression for #7282: an unsupported provider must raise a clear
# ValueError, not the opaque "RuntimeError: No active exception to reraise"
# a bare `raise` produced, and not a silent None that hides the
# misconfiguration behind an inline fallback
with pytest.raises(ValueError, match="No file uploader available"):
get_uploader("does-not-exist")


def test_get_uploader_raises_for_unconfigured_bedrock(monkeypatch):
# Bedrock without a configured S3 bucket must raise a ValueError that names
# the missing configuration, not RuntimeError and not a silent None
monkeypatch.delenv("CREWAI_BEDROCK_S3_BUCKET", raising=False)
with pytest.raises(ValueError, match="CREWAI_BEDROCK_S3_BUCKET"):
get_uploader("bedrock")


def test_get_uploader_raises_for_bedrock_with_falsy_bucket_name(monkeypatch):
# An explicit falsy bucket_name (None or "") is unconfigured just like an
# absent one, so the guard keys on the value, not key presence
monkeypatch.delenv("CREWAI_BEDROCK_S3_BUCKET", raising=False)
with pytest.raises(ValueError, match="CREWAI_BEDROCK_S3_BUCKET"):
get_uploader("bedrock", bucket_name=None)
with pytest.raises(ValueError, match="CREWAI_BEDROCK_S3_BUCKET"):
get_uploader("bedrock", bucket_name="")