mirror of
https://github.com/BerriAI/litellm.git
synced 2026-10-02 02:11:58 +00:00
feat(s3_v2): add s3_partition_granularity option for hourly S3 folders (#43748)
* 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> * test(integration): cover s3 v2 partition granularity across surfaces, settings and chaos Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> * test(integration): cover previous_response_id history rebuilt from an hourly cold storage object Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> * fix(s3_v2): reuse the cold storage key only when s3_v2 owns cold storage Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> * test(s3_v2): cover hour rollover, postgres outage, in-flight switches, key/team vars and real S3 layout Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> * chore(liccheck): authorize libfaketime, the GPLv2 dev-only clock the s3 rollover integration test preloads Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> * test(s3_v2): wait for the rejected-request cell's payloads by id, not by line count Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> * test(s3_v2): declare the postgres outage cell's models in config and trip the relay on burst ids Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> * test(s3_v2): drop the libfaketime hour rollover cell and its dev dependency Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> * test(e2e): deselect the s3_v2 live e2e on the stage-mirror stack The stage-mirror config enables no s3_v2 callback, so every test in test_s3_log_e2e.py fails its readiness check there. The file keeps running in the Buildkite e2e lane, which configures s3_v2 Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> * test(s3_v2): declare the sink outage burst models in config Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> * refactor(s3_v2): read cold storage metadata without an empty dict default Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> * test(integration): wait for the proxy to reconnect before the postgres outage recovery request Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> --------- Co-authored-by: mrinal <mrinal@berri.ai> Co-authored-by: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> Co-authored-by: yucheng <yucheng@berri.ai>
This commit is contained in:
parent
d24240c014
commit
bfd3f39dca
18 changed files with 1679 additions and 9 deletions
1
.github/e2e-stack/select_tests.py
vendored
1
.github/e2e-stack/select_tests.py
vendored
|
|
@ -11,6 +11,7 @@ UNSUPPORTED: Final = re.compile(
|
|||
r"|^tests/e2e/guardrails/test_presidio_masking_e2e\.py$"
|
||||
r"|^tests/e2e/logging/test_otel_v2_langfuse_generation_output_e2e\.py$"
|
||||
r"|^tests/e2e/logging/test_langsmith_batch_serialization_e2e\.py$"
|
||||
r"|^tests/e2e/logging/test_s3_log_e2e\.py$"
|
||||
r"|^tests/e2e/secret_manager/"
|
||||
)
|
||||
HARNESS: Final = re.compile(
|
||||
|
|
|
|||
|
|
@ -69,6 +69,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))
|
||||
|
|
|
|||
|
|
@ -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"
|
||||
|
|
|
|||
|
|
@ -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:
|
||||
|
|
|
|||
|
|
@ -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,27 @@ 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,
|
||||
)
|
||||
|
||||
metadata: Final = standard_logging_payload.get("metadata")
|
||||
cold_storage_object_key: Final = (
|
||||
metadata.get("cold_storage_object_key")
|
||||
if metadata is not None and litellm.cold_storage_custom_logger == "s3_v2"
|
||||
else None
|
||||
)
|
||||
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)
|
||||
|
||||
|
|
|
|||
|
|
@ -117,6 +117,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,
|
||||
|
|
@ -6059,6 +6060,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:
|
||||
|
|
@ -6067,6 +6069,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
|
||||
|
|
@ -6076,6 +6080,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
|
||||
|
|
|
|||
|
|
@ -3957,6 +3957,7 @@ class AllCallbacks(LiteLLMPydanticObjectBase):
|
|||
"AWS_SECRET_ACCESS_KEY",
|
||||
"AWS_REGION_NAME",
|
||||
"S3_LOG_PROMPTS_ONLY",
|
||||
"S3_PARTITION_GRANULARITY",
|
||||
],
|
||||
)
|
||||
|
||||
|
|
|
|||
|
|
@ -1,5 +1,9 @@
|
|||
from typing import Literal
|
||||
|
||||
from pydantic import BaseModel
|
||||
|
||||
S3PartitionGranularity = Literal["day", "hour"]
|
||||
|
||||
|
||||
class s3BatchLoggingElement(BaseModel):
|
||||
"""
|
||||
|
|
|
|||
|
|
@ -1,6 +1,7 @@
|
|||
# Logging integration delivery (behavior features). Grounded in litellm/integrations/.
|
||||
- {id: logging.s3.success.writes_object, module: logging, tier: P0, event: success, assertions: [writes_object], exercised_on: [chat_completions, messages, embeddings], source: "integrations/s3_v2.py", rationale: "Primary audit trail; batch flush no-drop"}
|
||||
- {id: logging.s3.failure.writes_object, module: logging, tier: P0, event: failure, assertions: [writes_object], exercised_on: [chat_completions, messages], source: "integrations/s3_v2.py", rationale: "Failed calls persisted for compliance"}
|
||||
- {id: logging.s3.success.partition_layout, module: logging, tier: P1, event: success, assertions: [object_key_layout], exercised_on: [chat_completions], source: "integrations/s3_v2.py / LIT-8985", rationale: "s3_partition_granularity picks the date or date/hour folder every downstream query and lifecycle rule reads"}
|
||||
- {id: logging.gcs_bucket.success.writes_object, module: logging, tier: P0, event: success, assertions: [writes_object], exercised_on: [chat_completions, messages, embeddings], source: "integrations/gcs_bucket/gcs_bucket.py", rationale: "GCS parallel to S3"}
|
||||
- {id: logging.datadog.success.exports_metric, module: logging, tier: P0, event: success, assertions: [exports_metric], exercised_on: [chat_completions, messages, responses, embeddings], source: "integrations/datadog/datadog.py", rationale: "Powers dashboards/alerts; cardinality regressions common"}
|
||||
- {id: logging.datadog.stream.exports_metric, module: logging, tier: P0, event: stream, assertions: [exports_metric], exercised_on: [chat_completions, messages, responses], source: "integrations/datadog/datadog.py", rationale: "Streaming aggregates usage after the last chunk; delivery and cost must survive that path"}
|
||||
|
|
|
|||
|
|
@ -106,6 +106,7 @@ UI_BASE_URL = os.environ.get("E2E_UI_BASE_URL", PROXY_BASE_URL).rstrip("/")
|
|||
|
||||
CHEAP_ANTHROPIC_MODEL = os.environ.get("E2E_CHEAP_ANTHROPIC_MODEL", "claude-haiku-4-5")
|
||||
CHEAP_OPENAI_MODEL = os.environ.get("E2E_CHEAP_OPENAI_MODEL", "gpt-5.5")
|
||||
S3_PARTITION_GRANULARITY = os.environ.get("E2E_S3_PARTITION_GRANULARITY", "day")
|
||||
|
||||
LINEAR_MCP_URL = os.environ.get("E2E_LINEAR_MCP_URL", "https://mcp.linear.app/mcp")
|
||||
LINEAR_STORAGE_STATE = os.environ.get("E2E_LINEAR_STORAGE_STATE", "")
|
||||
|
|
|
|||
|
|
@ -22,11 +22,12 @@ alias per test turns the poll into a cheap prefix listing.
|
|||
from __future__ import annotations
|
||||
|
||||
import math
|
||||
import re
|
||||
import time
|
||||
|
||||
import pytest
|
||||
|
||||
from e2e_config import CHEAP_ANTHROPIC_MODEL, unique_marker
|
||||
from e2e_config import CHEAP_ANTHROPIC_MODEL, S3_PARTITION_GRANULARITY, unique_marker
|
||||
from lifecycle import ResourceManager
|
||||
from logging_client import (
|
||||
INVALID_UPSTREAM_API_KEY,
|
||||
|
|
@ -106,6 +107,46 @@ class TestS3LogDelivery:
|
|||
record.response_cost, outcome.response_cost, rel_tol=1e-9
|
||||
), f"payload response_cost {record.response_cost!r} must equal the header cost {outcome.response_cost}"
|
||||
|
||||
@pytest.mark.covers("logging.s3.success.partition_layout", exercised_on=["chat_completions"])
|
||||
def test_chat_completions_object_key_follows_the_partition_granularity(
|
||||
self, client: LoggingClient, s3_logs: S3LogReader, resources: ResourceManager
|
||||
) -> None:
|
||||
"""The one object a call writes must sit in the folder layout the proxy's
|
||||
s3_partition_granularity names: {alias}/{date}/ for day and
|
||||
{alias}/{date}/{HH}/ for hour, where HH is the hour the object's own
|
||||
time- file name records. E2E_S3_PARTITION_GRANULARITY tells the test
|
||||
which one the proxy under test runs."""
|
||||
_assert_s3_configured(client)
|
||||
|
||||
alias = f"s3-layout-{unique_marker()}"
|
||||
key = client.key_with_alias(alias, models=[CHEAP_ANTHROPIC_MODEL])
|
||||
resources.defer(lambda: client.delete_key(key))
|
||||
|
||||
outcome = first_ok(
|
||||
client,
|
||||
lambda: client.chat_raw(
|
||||
key, CHEAP_ANTHROPIC_MODEL, f"reply with one word {unique_marker()}", max_tokens=16
|
||||
),
|
||||
)
|
||||
body_id = completion_response_id(outcome.body)
|
||||
assert body_id is not None, "the completion body must carry an id (it names the s3 object)"
|
||||
records = s3_logs.poll_records(prefix=f"{alias}/", predicate=lambda r: r.id == body_id)
|
||||
assert len(records) == 1, f"expected exactly ONE s3 object for response {body_id}, got {len(records)}"
|
||||
|
||||
file_id = body_id.replace("/", "_").replace(":", "_")
|
||||
hour_folder = r"(?P<folder_hour>\d{2})/" if S3_PARTITION_GRANULARITY == "hour" else ""
|
||||
layout = re.compile(
|
||||
rf"{re.escape(alias)}/\d{{4}}-\d{{2}}-\d{{2}}/{hour_folder}"
|
||||
rf"time-(?P<file_hour>\d{{2}})-\d{{2}}-\d{{2}}-\d{{6}}_{re.escape(file_id)}\.json"
|
||||
)
|
||||
keys = [object_key for object_key in s3_logs.list_keys(f"{alias}/") if file_id in object_key]
|
||||
assert len(keys) == 1, f"expected one object key for response {body_id}, got {keys}"
|
||||
match = layout.fullmatch(keys[0])
|
||||
assert match is not None, f"{keys[0]!r} is outside the {S3_PARTITION_GRANULARITY} layout {layout.pattern!r}"
|
||||
assert S3_PARTITION_GRANULARITY != "hour" or match.group("folder_hour") == match.group("file_hour"), (
|
||||
f"the hour folder must be the hour the object's file name records: {keys[0]!r}"
|
||||
)
|
||||
|
||||
@pytest.mark.covers("logging.s3.failure.writes_object", exercised_on=["chat_completions"])
|
||||
def test_chat_completions_failure_writes_one_object(
|
||||
self, client: LoggingClient, s3_logs: S3LogReader, resources: ResourceManager
|
||||
|
|
|
|||
|
|
@ -29,6 +29,7 @@ class DatabaseRelay:
|
|||
self._armed: Final = threading.Event()
|
||||
self.tripped: Final = threading.Event()
|
||||
self.refused = 0
|
||||
self.reconnected: Final = threading.Event()
|
||||
self._tripped_at = 0.0
|
||||
self._writers: tuple[asyncio.StreamWriter, ...] = ()
|
||||
self._ready: Final = threading.Event()
|
||||
|
|
@ -61,6 +62,8 @@ class DatabaseRelay:
|
|||
self.refused += 1
|
||||
client_writer.close()
|
||||
return
|
||||
if self.tripped.is_set():
|
||||
self.reconnected.set()
|
||||
server_reader, server_writer = await asyncio.open_connection(self._upstream_host, self._upstream_port)
|
||||
self._writers = (*self._writers, client_writer, server_writer)
|
||||
|
||||
|
|
|
|||
1241
tests/integration/observability/test_s3_v2_partition_granularity.py
Normal file
1241
tests/integration/observability/test_s3_v2_partition_granularity.py
Normal file
File diff suppressed because it is too large
Load diff
|
|
@ -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-")
|
||||
|
|
|
|||
|
|
@ -2552,6 +2552,253 @@ 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"
|
||||
|
||||
|
||||
def test_hour_upload_ignores_a_cold_storage_key_owned_by_another_logger(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_partition_granularity": "hour"}
|
||||
)
|
||||
monkeypatch.setattr(litellm, "cold_storage_custom_logger", "gcs_bucket")
|
||||
logger = S3Logger()
|
||||
cold_key = StandardLoggingPayloadSetup._generate_cold_storage_object_key(
|
||||
start_time=_PARTITION_START, response_id=_PARTITION_ID
|
||||
)
|
||||
uploaded = logger.create_s3_batch_logging_element(
|
||||
_PARTITION_START,
|
||||
StandardLoggingPayload(id=_PARTITION_ID, metadata={"cold_storage_object_key": cold_key}, messages=[]),
|
||||
)
|
||||
|
||||
assert cold_key == f"2026-09-29/time-14-05-09-123456_{_PARTITION_ID}.json"
|
||||
assert uploaded is not None
|
||||
assert uploaded.s3_object_key == f"2026-09-29/14/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",
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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(<Settings {...defaultProps} />);
|
||||
|
|
|
|||
|
|
@ -238,6 +238,7 @@ export const CallbackSelector: React.FC<CallbackSelectorProps> = ({
|
|||
};
|
||||
|
||||
const CALLBACK_CONFIG_ALIASES: Record<string, string> = { s3_v2: "s3" };
|
||||
const CALLBACK_UNSUPPORTED_PARAMS: Record<string, readonly string[]> = { 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) : [];
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue