From 9b2159d488ad13fe728435ccaf5efc65c6feb2e7 Mon Sep 17 00:00:00 2001 From: mrinal Date: Tue, 29 Sep 2026 18:19:37 +0000 Subject: [PATCH] feat(s3_v2): add s3_partition_granularity option for hourly S3 folders Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> --- litellm/constants.py | 1 + litellm/integrations/callback_configs.json | 7 + litellm/integrations/s3.py | 17 +- litellm/integrations/s3_v2.py | 50 +++- litellm/litellm_core_utils/litellm_logging.py | 5 + litellm/proxy/_types.py | 1 + litellm/types/integrations/s3_v2.py | 4 + tests/unit/integrations/test_s3.py | 7 + tests/unit/integrations/test_s3_v2.py | 224 ++++++++++++++++++ .../test_litellm_logging.py | 3 + .../src/components/settings.test.tsx | 47 ++++ .../src/components/settings.tsx | 4 +- 12 files changed, 362 insertions(+), 8 deletions(-) diff --git a/litellm/constants.py b/litellm/constants.py index 39c10d71709..f3f0398a368 100644 --- a/litellm/constants.py +++ b/litellm/constants.py @@ -57,6 +57,7 @@ S3_PREFIX_DIGEST_CHARS: Final = 16 # s3 allows 2048 bytes of combined metadata headers, which Content-Disposition counts against MAX_S3_OBJECT_DOWNLOAD_FILENAME_BYTES: Final = 1024 S3_LOG_PROMPTS_ONLY_ENV_VAR: Final = "S3_LOG_PROMPTS_ONLY" +S3_PARTITION_GRANULARITY_ENV_VAR: Final = "S3_PARTITION_GRANULARITY" MAX_FILE_LIST_LIMIT: Final = 10000 DEFAULT_SQS_FLUSH_INTERVAL_SECONDS: Final = int(os.getenv("DEFAULT_SQS_FLUSH_INTERVAL_SECONDS", 10)) DEFAULT_NUM_WORKERS_LITELLM_PROXY: Final = int(os.getenv("DEFAULT_NUM_WORKERS_LITELLM_PROXY", 1)) diff --git a/litellm/integrations/callback_configs.json b/litellm/integrations/callback_configs.json index 190c283d087..38928f67f42 100644 --- a/litellm/integrations/callback_configs.json +++ b/litellm/integrations/callback_configs.json @@ -498,6 +498,13 @@ "ui_name": "Log Prompts Only", "description": "Log request messages to S3 but drop the model response from each logged object", "required": false + }, + "s3_partition_granularity": { + "type": "select", + "ui_name": "Folder Partitioning", + "description": "day writes one folder per date, hour adds an hour folder below each date (s3_v2 only)", + "options": ["day", "hour"], + "required": false } }, "description": "S3 Bucket (AWS) Logging Integration" diff --git a/litellm/integrations/s3.py b/litellm/integrations/s3.py index f330ca8e0ac..129fceb40bf 100644 --- a/litellm/integrations/s3.py +++ b/litellm/integrations/s3.py @@ -16,8 +16,10 @@ from litellm.constants import ( MAX_S3_OBJECT_KEY_BYTES, S3_BOUNDED_OBJECT_KEY_HEAD_BYTES, S3_LOG_PROMPTS_ONLY_ENV_VAR, + S3_PARTITION_GRANULARITY_ENV_VAR, S3_PREFIX_DIGEST_CHARS, ) +from litellm.types.integrations.s3_v2 import S3PartitionGranularity from litellm.types.utils import StandardLoggingPayload _S3_BOOL: Final = TypeAdapter(bool) @@ -36,6 +38,18 @@ def resolve_s3_log_prompts_only(configured: object, environ: Mapping[str, str] | return True +def resolve_s3_partition_granularity( + configured: object, environ: Mapping[str, str] | None = None +) -> S3PartitionGranularity: + env: Final = os.environ if environ is None else environ + raw: Final = env.get(S3_PARTITION_GRANULARITY_ENV_VAR) if configured is None else configured + if raw == "hour": + return "hour" + if raw is not None and raw not in ("", "day"): + verbose_logger.warning("s3 logging: s3_partition_granularity=%r is not one of day, hour, using day", raw) + return "day" + + def _resolve_positive_int(setting: str, configured: object, fallback: int, *, reject_bool: bool) -> int: if configured is None or configured == "": return fallback @@ -371,10 +385,11 @@ def get_s3_object_key( prefix: str, start_time: datetime, s3_file_name: str, + partition_granularity: S3PartitionGranularity = "day", ) -> str: sanitized_s3_file_name: Final = s3_file_name.replace("/", "_").replace(":", "_") configured_prefix: Final = (s3_path.rstrip("/") + "/" if s3_path else "") + prefix - date_segment: Final = start_time.strftime("%Y-%m-%d") + "/" + date_segment: Final = start_time.strftime("%Y-%m-%d/%H/" if partition_granularity == "hour" else "%Y-%m-%d/") # we need the s3 key to include the time, so we log cache hits too s3_object_key: Final = configured_prefix + date_segment + sanitized_s3_file_name + ".json" if len(s3_object_key.encode("utf-8")) <= MAX_S3_OBJECT_KEY_BYTES: diff --git a/litellm/integrations/s3_v2.py b/litellm/integrations/s3_v2.py index 88d7906cc4b..6119cbdc87b 100644 --- a/litellm/integrations/s3_v2.py +++ b/litellm/integrations/s3_v2.py @@ -9,6 +9,7 @@ NOTE 1: S3 does not provide a BATCH PUT API endpoint; by default each element is import asyncio import contextvars import logging +import os import re import time from collections.abc import Awaitable, Callable, Mapping @@ -28,6 +29,7 @@ from litellm.constants import ( DEFAULT_S3_FLUSH_INTERVAL_SECONDS, DEFAULT_S3_MAX_ADAPTIVE_CONCURRENCY, DEFAULT_S3_MAX_CONCURRENT_UPLOADS, + S3_PARTITION_GRANULARITY_ENV_VAR, ) from litellm.integrations.adaptive_concurrency import AdaptiveConcurrencyLimiter, PutSample from litellm.integrations.s3 import ( @@ -42,6 +44,7 @@ from litellm.integrations.s3 import ( resolve_s3_max_concurrent_uploads, resolve_s3_max_queue_size, resolve_s3_max_retry_age_seconds, + resolve_s3_partition_granularity, resolve_sse_params, ) from litellm.litellm_core_utils.aws_partition import get_aws_dns_suffix @@ -53,7 +56,7 @@ from litellm.llms.custom_httpx.http_handler import ( get_async_httpx_client, httpxSpecialProvider, ) -from litellm.types.integrations.s3_v2 import s3BatchLoggingElement +from litellm.types.integrations.s3_v2 import S3PartitionGranularity, s3BatchLoggingElement from litellm.types.utils import StandardAuditLogPayload, StandardLoggingPayload from .custom_batch_logger import CustomBatchLogger @@ -119,6 +122,8 @@ class S3Logger(CustomBatchLogger, BaseAWSLLM): _upload_limiter: asyncio.Semaphore | AdaptiveConcurrencyLimiter | None = None s3_drop_on_terminal_error: bool = True s3_max_retry_age_seconds: int | None = 3600 + s3_partition_granularity: object = None + _partition_granularity_cache: tuple[object, S3PartitionGranularity] | None = None def __init__( self, @@ -147,6 +152,7 @@ class S3Logger(CustomBatchLogger, BaseAWSLLM): s3_server_side_encryption: str | None = None, s3_sse_kms_key_id: str | None = None, s3_log_prompts_only: bool | None = None, + s3_partition_granularity: str | None = None, s3_max_concurrent_uploads: int = DEFAULT_S3_MAX_CONCURRENT_UPLOADS, s3_max_queue_size: int | None = None, s3_max_retry_age_seconds: int | None = 3600, @@ -195,6 +201,7 @@ class S3Logger(CustomBatchLogger, BaseAWSLLM): s3_server_side_encryption=s3_server_side_encryption, s3_sse_kms_key_id=s3_sse_kms_key_id, s3_log_prompts_only=s3_log_prompts_only, + s3_partition_granularity=s3_partition_granularity, s3_max_concurrent_uploads=s3_max_concurrent_uploads, s3_max_queue_size=s3_max_queue_size, s3_max_retry_age_seconds=s3_max_retry_age_seconds, @@ -271,6 +278,7 @@ class S3Logger(CustomBatchLogger, BaseAWSLLM): s3_server_side_encryption: str | None = None, s3_sse_kms_key_id: str | None = None, s3_log_prompts_only: bool | None = None, + s3_partition_granularity: str | None = None, s3_max_concurrent_uploads: int = DEFAULT_S3_MAX_CONCURRENT_UPLOADS, s3_max_queue_size: int | None = None, s3_max_retry_age_seconds: int | None = 3600, @@ -331,6 +339,11 @@ class S3Logger(CustomBatchLogger, BaseAWSLLM): params.get("s3_log_prompts_only") if s3_log_prompts_only is None else s3_log_prompts_only ) + self.s3_partition_granularity = ( + params.get("s3_partition_granularity") if s3_partition_granularity is None else s3_partition_granularity + ) + self._partition_granularity_cache = None + self.s3_server_side_encryption, self.s3_sse_kms_key_id = resolve_sse_params( params.get("s3_server_side_encryption") or s3_server_side_encryption, params.get("s3_sse_kms_key_id") or s3_sse_kms_key_id, @@ -482,6 +495,7 @@ class S3Logger(CustomBatchLogger, BaseAWSLLM): "audit_logs/", now, f"{now.strftime('%H-%M-%S')}_{audit_log_id}", + partition_granularity=self.resolve_partition_granularity(), ) element: Final = s3BatchLoggingElement( @@ -758,6 +772,19 @@ class S3Logger(CustomBatchLogger, BaseAWSLLM): ), ) + def resolve_partition_granularity(self) -> S3PartitionGranularity: + raw: Final = ( + os.environ.get(S3_PARTITION_GRANULARITY_ENV_VAR) + if self.s3_partition_granularity is None + else self.s3_partition_granularity + ) + cached: Final = self._partition_granularity_cache + if cached is not None and cached[0] == raw: + return cached[1] + resolved: Final = resolve_s3_partition_granularity(raw) + self._partition_granularity_cache = (raw, resolved) + return resolved + def create_s3_batch_logging_element( self, start_time: datetime, @@ -803,11 +830,22 @@ class S3Logger(CustomBatchLogger, BaseAWSLLM): prefix_path, s3_file_name, ) - s3_object_key: Final = get_s3_object_key( - s3_path=cast(str | None, self.s3_path) or "", - prefix=prefix_path, - start_time=start_time, - s3_file_name=s3_file_name, + + def object_key(partition_granularity: S3PartitionGranularity) -> str: + return get_s3_object_key( + s3_path=cast(str | None, self.s3_path) or "", + prefix=prefix_path, + start_time=start_time, + s3_file_name=s3_file_name, + partition_granularity=partition_granularity, + ) + + cold_storage_object_key: Final = standard_logging_payload.get("metadata", {}).get("cold_storage_object_key") + s3_object_key: Final = ( + cold_storage_object_key + if cold_storage_object_key is not None + and cold_storage_object_key in (object_key("day"), object_key("hour")) + else object_key(self.resolve_partition_granularity()) ) verbose_logger.debug("s3_object_key=%s", s3_object_key) diff --git a/litellm/litellm_core_utils/litellm_logging.py b/litellm/litellm_core_utils/litellm_logging.py index 06cbfd4fc04..9c7d21cee23 100644 --- a/litellm/litellm_core_utils/litellm_logging.py +++ b/litellm/litellm_core_utils/litellm_logging.py @@ -115,6 +115,7 @@ from litellm.llms.base_llm.search.transformation import SearchResponse from litellm.responses.utils import ResponseAPILoggingUtils from litellm.types.agents import LiteLLMSendMessageResponse from litellm.types.containers.main import ContainerObject +from litellm.types.integrations.s3_v2 import S3PartitionGranularity from litellm.types.interactions import ( InteractionsAPIResponse, InteractionsAPIStreamingResponse, @@ -6032,6 +6033,7 @@ class StandardLoggingPayloadSetup: # Get the actual s3_path from the configured cold storage logger instance s3_path = "" # default value + partition_granularity: S3PartitionGranularity = "day" # Try to get the actual logger instance from the logger name try: @@ -6040,6 +6042,8 @@ class StandardLoggingPayloadSetup: ) if custom_logger and hasattr(custom_logger, "s3_path") and getattr(custom_logger, "s3_path"): s3_path = getattr(custom_logger, "s3_path") + if isinstance(custom_logger, S3V2Logger): + partition_granularity = custom_logger.resolve_partition_granularity() except Exception: # If any error occurs in getting the logger instance, use default empty s3_path pass @@ -6049,6 +6053,7 @@ class StandardLoggingPayloadSetup: prefix="", # Don't split by team alias for cold storage start_time=start_time, s3_file_name=s3_file_name, + partition_granularity=partition_granularity, ) return s3_object_key diff --git a/litellm/proxy/_types.py b/litellm/proxy/_types.py index d9fb053035b..ce6b1c991ea 100644 --- a/litellm/proxy/_types.py +++ b/litellm/proxy/_types.py @@ -3916,6 +3916,7 @@ class AllCallbacks(LiteLLMPydanticObjectBase): "AWS_SECRET_ACCESS_KEY", "AWS_REGION_NAME", "S3_LOG_PROMPTS_ONLY", + "S3_PARTITION_GRANULARITY", ], ) diff --git a/litellm/types/integrations/s3_v2.py b/litellm/types/integrations/s3_v2.py index 3b0dad97e8c..e8ad28f1a3b 100644 --- a/litellm/types/integrations/s3_v2.py +++ b/litellm/types/integrations/s3_v2.py @@ -1,5 +1,9 @@ +from typing import Literal + from pydantic import BaseModel +S3PartitionGranularity = Literal["day", "hour"] + class s3BatchLoggingElement(BaseModel): """ diff --git a/tests/unit/integrations/test_s3.py b/tests/unit/integrations/test_s3.py index fd677b9dfdf..c9a53a43d34 100644 --- a/tests/unit/integrations/test_s3.py +++ b/tests/unit/integrations/test_s3.py @@ -312,3 +312,10 @@ def test_prompts_only_payload_returns_copy_with_response_cleared(): assert stripped["messages"] == TEST_MESSAGES assert stripped is not payload assert payload == snapshot + + +def test_legacy_s3_logger_ignores_partition_granularity_and_keeps_daily_folder(): + mock_s3_client = _run_log_event({"s3_bucket_name": "b", "s3_path": "logs", "s3_partition_granularity": "hour"}) + + key = mock_s3_client.put_object.call_args.kwargs["Key"] + assert key.startswith("logs/2026-07-30/time-12-00-00-") diff --git a/tests/unit/integrations/test_s3_v2.py b/tests/unit/integrations/test_s3_v2.py index caab4ff561d..a50ca31d367 100644 --- a/tests/unit/integrations/test_s3_v2.py +++ b/tests/unit/integrations/test_s3_v2.py @@ -2522,6 +2522,230 @@ def test_prompts_only_toggle_is_exposed_to_admin_ui_for_both_s3_callbacks(callba assert "S3_LOG_PROMPTS_ONLY" in CustomLogger.get_callback_env_vars(callback_name) +_PARTITION_START: Final = datetime(2026, 9, 29, 14, 5, 9, 123456) +_PARTITION_ID: Final = "chatcmpl-partition" + + +def _partition_payload(response_id: str = _PARTITION_ID) -> StandardLoggingPayload: + return StandardLoggingPayload( + id=response_id, + metadata={"user_api_key_team_alias": "team-a", "user_api_key_alias": "key-a"}, + messages=[], + ) + + +def _partition_logger( + monkeypatch: pytest.MonkeyPatch, callback_params: dict[str, object], **kwargs: object +) -> S3Logger: + import litellm + + monkeypatch.setattr( + litellm, + "s3_callback_params", + {"s3_bucket_name": "test-bucket", "s3_region_name": "us-east-1", "s3_path": "logs", **callback_params}, + ) + return S3Logger( + s3_aws_access_key_id="test-key", + s3_aws_secret_access_key="test-secret", + s3_use_team_prefix=True, + s3_use_key_prefix=True, + **kwargs, + ) + + +_DAILY_KEY: Final = f"logs/team-a/key-a/2026-09-29/time-14-05-09-123456_{_PARTITION_ID}.json" +_HOURLY_KEY: Final = f"logs/team-a/key-a/2026-09-29/14/time-14-05-09-123456_{_PARTITION_ID}.json" + + +@pytest.mark.parametrize( + ("callback_params", "expected_key"), + [ + ({}, _DAILY_KEY), + ({"s3_partition_granularity": None}, _DAILY_KEY), + ({"s3_partition_granularity": "day"}, _DAILY_KEY), + ({"s3_partition_granularity": "hour"}, _HOURLY_KEY), + ], +) +def test_partition_granularity_sets_request_log_folder( + monkeypatch: pytest.MonkeyPatch, callback_params: dict[str, object], expected_key: str +) -> None: + monkeypatch.delenv("S3_PARTITION_GRANULARITY", raising=False) + logger = _partition_logger(monkeypatch, callback_params) + + element = logger.create_s3_batch_logging_element(_PARTITION_START, _partition_payload()) + + assert element is not None + assert element.s3_object_key == expected_key + + +@pytest.mark.parametrize("invalid", ["hourly", "HOUR", "1", 1, True]) +def test_invalid_partition_granularity_warns_and_keeps_daily_folder( + monkeypatch: pytest.MonkeyPatch, invalid: object +) -> None: + monkeypatch.delenv("S3_PARTITION_GRANULARITY", raising=False) + with patch("litellm.integrations.s3.verbose_logger") as mock_logger: + logger = _partition_logger(monkeypatch, {"s3_partition_granularity": invalid}) + element = logger.create_s3_batch_logging_element(_PARTITION_START, _partition_payload()) + second = logger.create_s3_batch_logging_element(_PARTITION_START, _partition_payload()) + + assert element is not None + assert second is not None + assert element.s3_object_key == second.s3_object_key == _DAILY_KEY + mock_logger.warning.assert_called_once() + assert mock_logger.warning.call_args.args[1:] == (invalid,) + + +def test_partition_granularity_reads_admin_ui_env_var_below_callback_params(monkeypatch: pytest.MonkeyPatch) -> None: + monkeypatch.setenv("S3_PARTITION_GRANULARITY", "hour") + + from_env = _partition_logger(monkeypatch, {}).create_s3_batch_logging_element( + _PARTITION_START, _partition_payload() + ) + from_params = _partition_logger(monkeypatch, {"s3_partition_granularity": "day"}).create_s3_batch_logging_element( + _PARTITION_START, _partition_payload() + ) + + assert from_env is not None and from_env.s3_object_key == _HOURLY_KEY + assert from_params is not None and from_params.s3_object_key == _DAILY_KEY + + +def test_partition_granularity_constructor_argument_and_os_environ_reference(monkeypatch: pytest.MonkeyPatch) -> None: + monkeypatch.delenv("S3_PARTITION_GRANULARITY", raising=False) + monkeypatch.setenv("MY_S3_PARTITION", "hour") + + from_ctor = _partition_logger(monkeypatch, {}, s3_partition_granularity="hour") + from_secret = _partition_logger(monkeypatch, {"s3_partition_granularity": "os.environ/MY_S3_PARTITION"}) + + for logger in (from_ctor, from_secret): + element = logger.create_s3_batch_logging_element(_PARTITION_START, _partition_payload()) + assert element is not None and element.s3_object_key == _HOURLY_KEY + + +def test_hourly_partition_long_key_keeps_hour_folder_within_s3_limit(monkeypatch: pytest.MonkeyPatch) -> None: + from litellm.constants import MAX_S3_OBJECT_KEY_BYTES + + monkeypatch.delenv("S3_PARTITION_GRANULARITY", raising=False) + logger = _partition_logger(monkeypatch, {"s3_partition_granularity": "hour", "s3_path": "p" * 1100}) + + element = logger.create_s3_batch_logging_element(_PARTITION_START, _partition_payload("r" * 600)) + + assert element is not None + assert len(element.s3_object_key.encode("utf-8")) <= MAX_S3_OBJECT_KEY_BYTES + assert re.search(r"/2026-09-29/14/[0-9a-f]{64}\.json$", element.s3_object_key) + + +@pytest.mark.asyncio +@pytest.mark.parametrize(("granularity", "hour_folder"), [("hour", True), ("day", False), (None, False)]) +async def test_audit_log_key_follows_audit_callback_params_partition_granularity( + monkeypatch: pytest.MonkeyPatch, granularity: str | None, hour_folder: bool +) -> None: + monkeypatch.delenv("S3_PARTITION_GRANULARITY", raising=False) + logger = S3Logger( + s3_callback_params_override={ + "s3_bucket_name": "audit-bucket", + "s3_path": "audit", + "s3_partition_granularity": granularity, + } + ) + + await logger.async_log_audit_log_event({"id": "audit-1"}) + + (element,) = logger.log_queue + match = re.fullmatch( + r"audit/audit_logs/\d{4}-\d{2}-\d{2}/(?:(\d{2})/)?(\d{2})-\d{2}-\d{2}_audit-1\.json", element.s3_object_key + ) + assert match is not None, element.s3_object_key + assert (match.group(1) is not None) is hour_folder + if hour_folder: + assert match.group(1) == match.group(2) + + +@pytest.mark.asyncio +async def test_hourly_batch_file_upload_writes_one_file_per_hour_folder(monkeypatch: pytest.MonkeyPatch) -> None: + monkeypatch.delenv("S3_PARTITION_GRANULARITY", raising=False) + logger = _partition_logger(monkeypatch, {"s3_partition_granularity": "hour"}, s3_batch_file_upload=True) + put = _RecordingPut() + logger.async_httpx_client = AsyncMock() + logger.async_httpx_client.put = put + before = logger.create_s3_batch_logging_element(datetime(2026, 9, 29, 13, 59, 59), _partition_payload("before")) + after = logger.create_s3_batch_logging_element(datetime(2026, 9, 29, 14, 0, 1), _partition_payload("after")) + assert before is not None and after is not None + logger.log_queue = [before, after] + + await logger.async_send_batch() + + by_folder = { + re.sub(r"/batch_\d{2}-\d{2}-\d{2}_[0-9a-f]{32}\.jsonl$", "", url.split(".com/", 1)[-1]): data + for url, data, _headers in put.calls + } + assert sorted(by_folder) == ["logs/team-a/key-a/2026-09-29/13", "logs/team-a/key-a/2026-09-29/14"] + assert [json.loads(line)["id"] for line in (by_folder["logs/team-a/key-a/2026-09-29/13"] or "").splitlines()] == [ + "before" + ] + assert [json.loads(line)["id"] for line in (by_folder["logs/team-a/key-a/2026-09-29/14"] or "").splitlines()] == [ + "after" + ] + + +@pytest.mark.parametrize("granularity", [None, "day", "hour"]) +def test_cold_storage_object_key_matches_the_uploaded_request_log_key( + monkeypatch: pytest.MonkeyPatch, granularity: str | None +) -> None: + import litellm + from litellm.litellm_core_utils.litellm_logging import StandardLoggingPayloadSetup + + monkeypatch.delenv("S3_PARTITION_GRANULARITY", raising=False) + monkeypatch.setattr( + litellm, + "s3_callback_params", + {"s3_bucket_name": "test-bucket", "s3_path": "coldlogs", "s3_partition_granularity": granularity}, + ) + monkeypatch.setattr(litellm, "cold_storage_custom_logger", "s3_v2") + logger = S3Logger() + uploaded = logger.create_s3_batch_logging_element( + _PARTITION_START, StandardLoggingPayload(id=_PARTITION_ID, metadata={}, messages=[]) + ) + + monkeypatch.setattr(litellm, "callbacks", [logger]) + cold_key = StandardLoggingPayloadSetup._generate_cold_storage_object_key( + start_time=_PARTITION_START, response_id=_PARTITION_ID + ) + + assert uploaded is not None + assert cold_key == uploaded.s3_object_key + assert ("/2026-09-29/14/" in cold_key) is (granularity == "hour") + + +def test_cold_storage_key_matches_upload_when_env_var_changes_mid_request(monkeypatch: pytest.MonkeyPatch) -> None: + import litellm + from litellm.litellm_core_utils.litellm_logging import StandardLoggingPayloadSetup + + monkeypatch.delenv("S3_PARTITION_GRANULARITY", raising=False) + monkeypatch.setattr(litellm, "s3_callback_params", {"s3_bucket_name": "test-bucket", "s3_path": "coldlogs"}) + monkeypatch.setattr(litellm, "cold_storage_custom_logger", "s3_v2") + logger = S3Logger() + monkeypatch.setattr(litellm, "callbacks", [logger]) + + cold_key = StandardLoggingPayloadSetup._generate_cold_storage_object_key( + start_time=_PARTITION_START, response_id=_PARTITION_ID + ) + monkeypatch.setenv("S3_PARTITION_GRANULARITY", "hour") + uploaded = logger.create_s3_batch_logging_element( + _PARTITION_START, + StandardLoggingPayload(id=_PARTITION_ID, metadata={"cold_storage_object_key": cold_key}, messages=[]), + ) + + assert uploaded is not None + assert cold_key == uploaded.s3_object_key == f"coldlogs/2026-09-29/time-14-05-09-123456_{_PARTITION_ID}.json" + + +@pytest.mark.parametrize("callback_name", ["s3", "s3_v2"]) +def test_partition_granularity_is_exposed_to_admin_ui(callback_name: str) -> None: + from litellm.integrations.custom_logger import CustomLogger + + assert "S3_PARTITION_GRANULARITY" in CustomLogger.get_callback_env_vars(callback_name) + + def _element(payload: dict[str, object], key_suffix: str) -> s3BatchLoggingElement: return s3BatchLoggingElement( s3_object_key=f"2025-09-14/test-{key_suffix}.json", diff --git a/tests/unit/litellm_core_utils/test_litellm_logging.py b/tests/unit/litellm_core_utils/test_litellm_logging.py index 2fc747e1b48..777736420a1 100644 --- a/tests/unit/litellm_core_utils/test_litellm_logging.py +++ b/tests/unit/litellm_core_utils/test_litellm_logging.py @@ -3227,6 +3227,7 @@ async def test_e2e_generate_cold_storage_object_key_successful(): prefix="", # No prefix for cold storage start_time=start_time, s3_file_name="time-10-30-45-123456_chatcmpl-test-12345", + partition_granularity="day", ) # Verify the result @@ -3276,6 +3277,7 @@ async def test_e2e_generate_cold_storage_object_key_with_custom_logger_s3_path() prefix="", start_time=start_time, s3_file_name="time-10-30-45-123456_chatcmpl-test-12345", + partition_granularity="day", ) # Verify the result @@ -3320,6 +3322,7 @@ async def test_e2e_generate_cold_storage_object_key_with_logger_no_s3_path(): prefix="", start_time=start_time, s3_file_name="time-10-30-45-123456_chatcmpl-test-12345", + partition_granularity="day", ) # Verify the result diff --git a/ui/litellm-dashboard/src/components/settings.test.tsx b/ui/litellm-dashboard/src/components/settings.test.tsx index 4ba5dd23fd1..9b67657dbca 100644 --- a/ui/litellm-dashboard/src/components/settings.test.tsx +++ b/ui/litellm-dashboard/src/components/settings.test.tsx @@ -313,6 +313,7 @@ describe("Settings", () => { "AWS_SECRET_ACCESS_KEY", "AWS_REGION_NAME", "S3_LOG_PROMPTS_ONLY", + "S3_PARTITION_GRANULARITY", ], ui_callback_name: "s3 Bucket (AWS)", }, @@ -326,6 +327,12 @@ describe("Settings", () => { dynamic_params: { s3_bucket_name: { type: "text", ui_name: "S3 Bucket Name", required: false }, s3_log_prompts_only: { type: "boolean", ui_name: "Log Prompts Only", required: false }, + s3_partition_granularity: { + type: "select", + ui_name: "Folder Partitioning", + options: ["day", "hour"], + required: false, + }, }, }, ]); @@ -409,6 +416,46 @@ describe("Settings", () => { }); }); + it("should show the saved s3_v2 folder partitioning and post the newly selected value", async () => { + mockS3Callback({ S3_LOG_PROMPTS_ONLY: null, S3_PARTITION_GRANULARITY: "hour" }, "s3_v2"); + const user = await openS3EditModal("s3_v2"); + + const dialog = screen.getByRole("dialog"); + const partitioning = await within(dialog).findByRole("combobox", { name: "Folder Partitioning" }); + expect(partitioning).toHaveTextContent("hour"); + + await user.click(partitioning); + await user.click(await screen.findByRole("option", { name: "day" })); + await user.click(within(dialog).getByRole("button", { name: "Save Changes" })); + + await waitFor(() => { + expect(vi.mocked(setCallbacksCall)).toHaveBeenCalledWith( + "token", + expect.objectContaining({ + environment_variables: expect.objectContaining({ callback: "s3_v2", s3_partition_granularity: "day" }), + litellm_settings: { success_callback: ["s3_v2"] }, + }), + ); + }); + }); + + it("should not offer folder partitioning for the legacy s3 callback, which cannot honour it", async () => { + mockS3Callback({ S3_LOG_PROMPTS_ONLY: null, S3_PARTITION_GRANULARITY: null }); + const user = await openS3EditModal(); + + const dialog = screen.getByRole("dialog"); + await within(dialog).findByRole("switch", { name: "Log Prompts Only" }); + expect(within(dialog).queryByRole("combobox", { name: "Folder Partitioning" })).not.toBeInTheDocument(); + + await user.click(within(dialog).getByRole("button", { name: "Save Changes" })); + await waitFor(() => { + expect(vi.mocked(setCallbacksCall)).toHaveBeenCalledTimes(1); + }); + const [, payload] = vi.mocked(setCallbacksCall).mock.calls[0]; + expect(payload.environment_variables).not.toHaveProperty("s3_partition_granularity"); + expect(payload.environment_variables).not.toHaveProperty("S3_PARTITION_GRANULARITY"); + }); + it("should send the typed webhook url for an alert type when the alerting tab is saved", async () => { const user = userEvent.setup(); render(); diff --git a/ui/litellm-dashboard/src/components/settings.tsx b/ui/litellm-dashboard/src/components/settings.tsx index 9247f22ec28..e376d858df8 100644 --- a/ui/litellm-dashboard/src/components/settings.tsx +++ b/ui/litellm-dashboard/src/components/settings.tsx @@ -238,6 +238,7 @@ export const CallbackSelector: React.FC = ({ }; const CALLBACK_CONFIG_ALIASES: Record = { s3_v2: "s3" }; +const CALLBACK_UNSUPPORTED_PARAMS: Record = { s3: ["s3_partition_granularity"] }; interface DynamicParamConfig { type?: string; @@ -274,7 +275,8 @@ const getDynamicParamsForCallback = ( const callbackConfig = findCallbackConfig(callbackConfigs, callbackName); if (callbackConfig?.dynamic_params) { - return Object.keys(callbackConfig.dynamic_params); + const unsupportedParams = CALLBACK_UNSUPPORTED_PARAMS[callbackName] ?? []; + return Object.keys(callbackConfig.dynamic_params).filter((param) => !unsupportedParams.includes(param)); } return fallbackVariables ? Object.keys(fallbackVariables) : [];