From 21a629810dd08d1b0c5c3190e67f87f26bb3baf5 Mon Sep 17 00:00:00 2001 From: yucheng Date: Wed, 30 Sep 2026 09:07:23 +0000 Subject: [PATCH] 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> --- pyproject.toml | 1 + tests/e2e/coverage_registry/logging.yaml | 1 + tests/e2e/e2e_config.py | 1 + tests/e2e/logging/test_s3_log_e2e.py | 43 +- tests/integration/_support/proxy.py | 14 +- .../test_s3_v2_partition_granularity.py | 371 +++++++++++++++++- uv.lock | 15 + 7 files changed, 438 insertions(+), 8 deletions(-) diff --git a/pyproject.toml b/pyproject.toml index fb21d8fa23b..bb01ff9f040 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -236,6 +236,7 @@ dev = [ "pytest-timeout==2.4.0", "vcrpy==8.2.1", "pytest-recording==0.13.4", + "libfaketime==3.0.1", ] e2e-dev = [ "playwright==1.61.0", diff --git a/tests/e2e/coverage_registry/logging.yaml b/tests/e2e/coverage_registry/logging.yaml index 7c83e4d3aea..5424fb2d12f 100644 --- a/tests/e2e/coverage_registry/logging.yaml +++ b/tests/e2e/coverage_registry/logging.yaml @@ -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"} diff --git a/tests/e2e/e2e_config.py b/tests/e2e/e2e_config.py index b2682c04841..3fa9f534ffd 100644 --- a/tests/e2e/e2e_config.py +++ b/tests/e2e/e2e_config.py @@ -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", "") diff --git a/tests/e2e/logging/test_s3_log_e2e.py b/tests/e2e/logging/test_s3_log_e2e.py index 7a1ee1e6536..1612fa315a6 100644 --- a/tests/e2e/logging/test_s3_log_e2e.py +++ b/tests/e2e/logging/test_s3_log_e2e.py @@ -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\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\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 diff --git a/tests/integration/_support/proxy.py b/tests/integration/_support/proxy.py index a444b93757d..621ced7e114 100644 --- a/tests/integration/_support/proxy.py +++ b/tests/integration/_support/proxy.py @@ -1,12 +1,20 @@ -"""Run the normal single-process CLI with the existing behavior-suite test entitlement.""" +"""Run the normal CLI with the existing behavior-suite test entitlement in the parent and every spawned worker.""" import signal import sys from types import FrameType +from typing import Final from unittest.mock import patch from litellm import run_server +_ENTITLEMENT: Final = patch( # test-quality-ok: route entitlement only; license checks are outside these contracts + "litellm.proxy.auth.litellm_license.LicenseCheck.is_premium", return_value=True +) + +if __name__ == "__mp_main__": + _ENTITLEMENT.start() + def _exit_on_reraised_term(signum: int, frame: FrameType | None) -> None: sys.exit(0) @@ -14,9 +22,7 @@ def _exit_on_reraised_term(signum: int, frame: FrameType | None) -> None: def main() -> None: signal.signal(signal.SIGTERM, _exit_on_reraised_term) - with patch( # test-quality-ok: route entitlement only; license validation is outside these HTTP/DB contracts - "litellm.proxy.auth.litellm_license.LicenseCheck.is_premium", return_value=True - ): + with _ENTITLEMENT: run_server() diff --git a/tests/integration/observability/test_s3_v2_partition_granularity.py b/tests/integration/observability/test_s3_v2_partition_granularity.py index ee9a2623116..6a840249aa0 100644 --- a/tests/integration/observability/test_s3_v2_partition_granularity.py +++ b/tests/integration/observability/test_s3_v2_partition_granularity.py @@ -7,6 +7,7 @@ from concurrent.futures import ThreadPoolExecutor from contextlib import contextmanager from dataclasses import dataclass, field from datetime import datetime, timedelta +from importlib.resources import files from pathlib import Path from typing import Final from urllib.parse import quote, unquote @@ -29,11 +30,13 @@ from _s3_v2_support import ( ) from integration._support.client import Gateway, JsonValue, Scenario, eventually, object_value from integration._support.database import read_rows, scratch_database +from integration._support.database_relay import database_relay from integration._support.process import OwnedProxy, group_members, owned_proxy_process from integration._support.wire import Reply, Request, wire_server FLUSH: Final = {"DEFAULT_S3_FLUSH_INTERVAL_SECONDS": "1"} HOUR: Final = {"s3_partition_granularity": "hour"} +FAKETIME_LIBRARY: Final = files("libfaketime").joinpath("vendor", "libfaketime", "src", "libfaketime.so.1") ANTHROPIC_MODEL: Final = "anthropic/claude-sonnet-4-5-20250929" WARNING: Final = "s3 logging: s3_partition_granularity=" SINK_CREDENTIALS: Final = { @@ -177,6 +180,49 @@ def _update_environment(candidate: Gateway, values: Mapping[str, JsonValue]) -> ) +def _keys_on_fresh_connections(candidate: Gateway, aliases: tuple[str, ...]) -> tuple[tuple[str, str], ...]: + def generate(alias: str) -> tuple[str, str]: + with httpx.Client(base_url=candidate.client.base_url, timeout=30, trust_env=False) as fresh: + response: Final = fresh.post( + "/key/generate", + json={"key_alias": alias}, + headers={"Authorization": f"Bearer {candidate.key}", "Connection": "close"}, + ) + assert response.status_code == 200, response.text + return str(response.json()["key"]), str(response.json()["token_id"]) + + with ThreadPoolExecutor(max_workers=len(aliases)) as pool: + return tuple(pool.map(generate, aliases)) + + +SETUP_AUDITS: Final = ( + ("created", "LiteLLM_ProxyModelTable"), + ("created", "LiteLLM_ProxyModelTable"), + ("created", "LiteLLM_VerificationToken"), +) + + +def _audit_changes(sink: RecordingS3Sink, audit_prefix: str) -> tuple[tuple[str, str], ...]: + bodies: Final = (body for target, body in sink.objects().items() if target.startswith(audit_prefix)) + lines: Final = b"\n".join(bodies).splitlines() + return tuple(sorted((str(audit["action"]), str(audit["table_name"])) for audit in map(_audit_line, lines))) + + +def _audit_line(line: bytes) -> Mapping[str, JsonValue]: + return object_value(json.loads(line)) + + +def _created_key_hashes(sink: RecordingS3Sink, audit_prefix: str) -> frozenset[str]: + created: Final = ( + object_value(json.loads(body)) for target, body in sink.objects().items() if target.startswith(audit_prefix) + ) + return frozenset( + str(audit["object_id"]) + for audit in created + if audit["action"] == "created" and audit["table_name"] == "LiteLLM_VerificationToken" + ) + + def test_s3_v2_hour_granularity_files_every_surface_under_its_hour_folder(gateway: Gateway, tmp_path: Path) -> None: marker: Final = "s3hour" + uuid.uuid4().hex[:8] upstream: Final = CountingUpstream() @@ -576,17 +622,20 @@ def test_s3_v2_audit_logs_follow_the_audit_params_granularity_not_the_request_lo "s3_audit_callback_params": {**SINK_CREDENTIALS, "s3_endpoint_url": bucket.url, **HOUR}, } with ( - _s3_proxy(gateway, tmp_path, bucket.url, {}, settings, workers=1) as owned, + _s3_proxy(gateway, tmp_path, bucket.url, {}, settings) as owned, owned.gateway.scenario() as scenario, ): openai_model, _, key = _models(scenario, provider.url, key_alias=marker) returned: Final = _sdk_chats(owned.gateway, openai_model, key, (marker,)) + aliases: Final = tuple(f"{marker}-fresh{index}" for index in range(16)) + fresh_keys: Final = _keys_on_fresh_connections(owned.gateway, aliases) audit_prefix: Final = f"/{BUCKET}/{PREFIX}/audit_logs/" eventually( - lambda: tuple(target for target in sink.objects() if target.startswith(audit_prefix)), - lambda targets: len(targets) >= 1, + lambda: _created_key_hashes(sink, audit_prefix), + lambda created: frozenset(token for _, token in fresh_keys) <= created, seconds=30, ) + owned.gateway.post("/key/delete", {"keys": [key for key, _ in fresh_keys]}) collect_payloads(sink, 2) objects: Final = sink.objects() audits: Final = { @@ -613,6 +662,40 @@ def test_s3_v2_audit_logs_follow_the_audit_params_granularity_not_the_request_lo ) +@pytest.mark.parametrize("level", ["key", "team"]) +def test_s3_v2_key_and_team_logging_callback_vars_cannot_change_the_proxy_hour_layout( + gateway: Gateway, tmp_path: Path, level: str +) -> None: + marker: Final = f"s3h{level}vars" + uuid.uuid4().hex[:8] + logging: Final[list[JsonValue]] = [ + {"callback_name": "s3_v2", "callback_type": "success", "callback_vars": {"s3_partition_granularity": "day"}} + ] + upstream: Final = CountingUpstream() + sink: Final = RecordingS3Sink(delay_seconds=0.05) + with ( + wire_server(upstream.respond) as provider, + wire_server(sink.respond) as bucket, + _s3_proxy(gateway, tmp_path, bucket.url, HOUR) as owned, + owned.gateway.scenario() as scenario, + ): + openai_model, _, _ = _models(scenario, provider.url) + key: Final = ( + scenario.key(models=[openai_model], metadata={"logging": logging}) + if level == "key" + else scenario.key(models=[openai_model], team_id=scenario.team(metadata={"logging": logging})) + ) + prompts: Final = tuple(f"{marker}-{index}" for index in range(8)) + returned: Final = _sdk_chats(owned.gateway, openai_model, key, prompts) + collect_payloads(sink, len(prompts)) + objects: Final = sink.objects() + assert returned == prompts + assert sorted(upstream.received()) == sorted(prompts) + assert sorted(str(object_value(json.loads(body))["id"]) for body in objects.values()) == sorted(prompts), ( + f"{level}-level s3_v2 logging must land exactly one object per request" + ) + assert _outside_layout(objects, "hour") == (), f"{level}-level callback_vars must not change the proxy granularity" + + def test_s3_v2_admin_ui_granularity_update_moves_live_traffic_on_both_workers(gateway: Gateway, tmp_path: Path) -> None: marker: Final = "s3hui" + uuid.uuid4().hex[:8] upstream: Final = CountingUpstream() @@ -713,6 +796,195 @@ def test_s3_v2_granularity_toggles_mid_burst_keep_every_cold_storage_key_on_its_ ) +def test_s3_v2_in_flight_request_keeps_its_cold_storage_key_on_its_object_across_owner_and_granularity_switches( + gateway: Gateway, tmp_path: Path +) -> None: + marker: Final = "s3hflight" + uuid.uuid4().hex[:8] + held_prompt: Final = f"{marker}-held" + upstream: Final = CountingUpstream() + arrived: Final = threading.Event() + release: Final = threading.Event() + + def held(request: Request) -> Reply: + if held_prompt.encode() in request.body: + arrived.set() + assert release.wait(90), "held request was never released" + return upstream.respond(request) + + sink: Final = RecordingS3Sink(delay_seconds=0.05) + with ( + scratch_database() as database_url, + wire_server(held) as provider, + wire_server(sink.respond) as bucket, + _s3_proxy( + gateway, + tmp_path, + bucket.url, + {}, + {"cold_storage_custom_logger": "s3_v2"}, + environment={"DATABASE_URL": database_url}, + ) as owned, + owned.gateway.scenario() as scenario, + ): + openai_model, _, key = _models(scenario, provider.url) + with ThreadPoolExecutor(max_workers=1) as flight: + pending: Final = flight.submit(_sdk_chats, owned.gateway, openai_model, key, (held_prompt,)) + assert arrived.wait(60), "held request never reached the upstream" + owner_switch: Final = owned.gateway.request( + "POST", "/config/update", {"litellm_settings": {"cold_storage_custom_logger": "gcs_bucket"}} + ) + _update_environment(owned.gateway, HOUR) + probe_round: Final = iter(range(1000)) + + def probe() -> Mapping[str, bytes]: + round_id: Final = next(probe_round) + prompts: Final = tuple(f"{marker}-probe{round_id}-{index}" for index in range(8)) + _sdk_chats(owned.gateway, openai_model, key, prompts) + eventually( + lambda: frozenset(str(payload["id"]) for payload in sink.payloads()), + lambda landed: frozenset(prompts) <= landed, + seconds=20, + ) + return {target: body for target, body in sink.objects().items() if f"-probe{round_id}-" in target} + + eventually(probe, lambda probed: len(probed) == 8 and _outside_layout(probed, "hour") == (), seconds=60) + release.set() + returned: Final = pending.result() + eventually( + lambda: frozenset(str(payload["id"]) for payload in sink.payloads()), + lambda landed: held_prompt in landed, + seconds=30, + ) + held_objects: Final = {target: body for target, body in sink.objects().items() if held_prompt in target} + cold_key: Final = _cold_storage_key(held_prompt, database_url) + assert owner_switch.status_code == 400, owner_switch.text + assert "cold_storage_custom_logger" in owner_switch.text and "config file" in owner_switch.text, owner_switch.text + assert returned == (held_prompt,) + assert upstream.received().count(held_prompt) == 1 + assert frozenset(held_objects) == frozenset({f"/{BUCKET}/{quote(cold_key, safe='/')}"}), ( + "the in-flight request's cold_storage_object_key must name the one object the logger uploaded", + cold_key, + tuple(held_objects), + ) + assert _outside_layout(held_objects, "hour") == () + + +def test_s3_v2_cold_storage_owner_saved_through_config_update_is_not_applied_to_a_running_proxy( + gateway: Gateway, tmp_path: Path +) -> None: + marker: Final = "s3howner" + uuid.uuid4().hex[:8] + upstream: Final = CountingUpstream() + sink: Final = RecordingS3Sink() + with ( + scratch_database() as database_url, + wire_server(upstream.respond) as provider, + wire_server(sink.respond) as bucket, + _s3_proxy(gateway, tmp_path, bucket.url, HOUR, environment={"DATABASE_URL": database_url}) as owned, + owned.gateway.scenario() as scenario, + ): + openai_model, _, key = _models(scenario, provider.url) + saved: Final = owned.gateway.request( + "POST", "/config/update", {"litellm_settings": {"cold_storage_custom_logger": "s3_v2"}} + ) + prompts: Final = tuple(f"{marker}-{index}" for index in range(8)) + answered: Final = _sdk_chats(owned.gateway, openai_model, key, prompts) + landed: Final = collect_payloads(sink, len(prompts)) + rows: Final = eventually( + lambda: read_rows( + 'SELECT request_id, metadata FROM "LiteLLM_SpendLogs" WHERE request_id = ANY(%s)', + (list(answered),), + database_url=database_url, + ), + lambda values: len(values) == len(prompts), + seconds=60, + ) + objects: Final = sink.objects() + stored: Final = read_rows( + 'SELECT param_value FROM "LiteLLM_Config" WHERE param_name = %s', + ("litellm_settings",), + database_url=database_url, + ) + cold_keys: Final = { + str(row["request_id"]): object_value( + json.loads(row["metadata"]) if isinstance(row["metadata"], str) else row["metadata"] + ).get("cold_storage_object_key") + for row in rows + } + assert saved.status_code == 200, saved.text + assert [ + object_value(json.loads(row["param_value"]) if isinstance(row["param_value"], str) else row["param_value"]).get( + "cold_storage_custom_logger" + ) + for row in stored + ] == ["s3_v2"], "the owner switch must be persisted, so the unchanged live keys are not a rejected write" + assert sorted(upstream.received()) == sorted(prompts) + assert sorted(_prompt(payload) for payload in landed) == sorted(prompts) + assert cold_keys == dict.fromkeys(answered), "a DB-saved cold storage owner must not change a live request" + assert _outside_layout(objects, "hour") == () + + +def test_s3_v2_hour_postgres_outage_mid_mixed_burst_lands_every_id_exactly_once_and_recovers( + gateway: Gateway, tmp_path: Path +) -> None: + marker: Final = "s3hpg" + uuid.uuid4().hex[:8] + upstream: Final = CountingUpstream() + sink: Final = RecordingS3Sink(delay_seconds=0.05) + sent: Final = _surface_prompts(marker, 5) + with ( + scratch_database() as database_url, + database_relay(database_url, b'"LiteLLM_SpendLogs"') as (relay, relayed_url), + wire_server(upstream.respond) as provider, + wire_server(sink.respond) as bucket, + _s3_proxy( + gateway, + tmp_path, + bucket.url, + HOUR, + {"cold_storage_custom_logger": "s3_v2"}, + environment={"DATABASE_URL": relayed_url}, + ) as owned, + owned.gateway.scenario() as scenario, + ): + openai_model, anthropic_model, key = _models(scenario, provider.url) + warm: Final = mixed_burst(owned.gateway, openai_model, anthropic_model, key, f"{marker}warm", per_surface=2) + eventually( + lambda: frozenset(_prompt(payload) for payload in sink.payloads()), + lambda landed: _surface_prompts(f"{marker}warm", 2) <= landed, + seconds=60, + ) + relay.arm() + answered: Final = mixed_burst(owned.gateway, openai_model, anthropic_model, key, marker, per_surface=5) + assert relay.tripped.wait(90), "no spend log write reached the database during the burst" + eventually(lambda: relay.refused, lambda count: count >= 1, seconds=30) + burst_payloads: Final = eventually( + lambda: tuple(payload for payload in sink.payloads() if _prompt(payload) in sent), + lambda landed: frozenset(_prompt(payload) for payload in landed) == frozenset(sent), + seconds=60, + ) + recovered_prompt: Final = f"{marker}-recovered" + recovered: Final = _sdk_chats(owned.gateway, openai_model, key, (recovered_prompt,)) + recovered_key: Final = _cold_storage_key(recovered_prompt, database_url) + eventually( + lambda: frozenset(str(payload["id"]) for payload in sink.payloads()), + lambda landed: recovered_prompt in landed, + seconds=30, + ) + objects: Final = sink.objects() + uploads: Final = sink.attempts + burst: Final = burst_payloads + assert len(warm) == len(_surface_prompts(f"{marker}warm", 2)) + assert len(answered) == len(sent) == 30 + assert sorted(prompt for prompt in upstream.received() if prompt.startswith(f"{marker}-")) == sorted( + (*sent, recovered_prompt) + ) + assert matched_ids(burst, answered) == frozenset(str(payload["id"]) for payload in burst) + assert sorted(_prompt(payload) for payload in burst) == sorted(sent), "every burst id lands exactly once" + assert uploads == len(objects), "no object is uploaded twice" + assert _outside_layout(objects, "hour") == () + assert recovered == (recovered_prompt,) + assert f"/{BUCKET}/{quote(recovered_key, safe='/')}" in objects, "cold key written after recovery names its object" + + def test_legacy_s3_callback_ignores_hour_granularity(gateway: Gateway, tmp_path: Path) -> None: marker: Final = "s3v1hour" + uuid.uuid4().hex[:8] upstream: Final = CountingUpstream() @@ -785,6 +1057,99 @@ def test_s3_v2_hour_coded_403_retries_reuse_the_same_hour_key(gateway: Gateway, assert _outside_layout(objects, "hour") == () +def test_s3_v2_hour_rollover_inside_one_batch_flush_splits_files_by_hour_folder( + gateway: Gateway, tmp_path: Path +) -> None: + assert FAKETIME_LIBRARY.is_file(), "the libfaketime dev dependency drives the proxy clock" + marker: Final = "s3hroll" + uuid.uuid4().hex[:8] + clock: Final = tmp_path / "proxy-clock" + clock.write_text("@2026-09-29 10:59:30\n") + upstream: Final = CountingUpstream() + sink: Final = RecordingS3Sink(delay_seconds=0.05) + audit_prefix: Final = f"/{BUCKET}/{PREFIX}/audit_logs/" + with ( + wire_server(upstream.respond) as provider, + wire_server(sink.respond) as bucket, + _s3_proxy( + gateway, + tmp_path, + bucket.url, + {**HOUR, "s3_batch_file_upload": True}, + {"store_audit_logs": True, "audit_log_callbacks": ["s3_v2"]}, + environment={ + "DEFAULT_S3_FLUSH_INTERVAL_SECONDS": "3600", + "DEFAULT_S3_BATCH_SIZE": "1", + "LD_PRELOAD": str(FAKETIME_LIBRARY), + "FAKETIME_TIMESTAMP_FILE": str(clock), + "FAKETIME_CACHE_DURATION": "1", + "FAKETIME_DONT_FAKE_MONOTONIC": "1", + }, + workers=1, + ) as owned, + owned.gateway.scenario() as scenario, + ): + openai_model, _, key = _models(scenario, provider.url) + setup_audits: Final = eventually( + lambda: _audit_changes(sink, audit_prefix), + lambda changes: changes == SETUP_AUDITS, + seconds=30, + ) + + def advance(stamp: str, shown: str) -> None: + clock.write_text(f"@2026-09-29 {stamp}\n") + eventually( + lambda: owned.gateway.client.get("/health/liveliness").headers["date"], + lambda date: f" {shown}:" in date, + seconds=10, + ) + + advance("10:59:40", "10:59") + flushed_before: Final = frozenset(sink.objects()) + before: Final = tuple(f"{marker}-before-{index}" for index in range(4)) + after: Final = tuple(f"{marker}-after-{index}" for index in range(4)) + answered_before: Final = _sdk_chats(owned.gateway, openai_model, key, before) + advance("11:00:05", "11:00") + answered_after: Final = _sdk_chats(owned.gateway, openai_model, key, after) + spent: Final = eventually( + lambda: read_rows( + 'SELECT request_id FROM "LiteLLM_SpendLogs" WHERE request_id = ANY(%s)', + (list(answered_before + answered_after),), + ), + lambda rows: len(rows) == len(before) + len(after), + seconds=70, + ) + pending: Final = frozenset(sink.objects()) - flushed_before + scenario.key(models=[openai_model]) + batches: Final = eventually( + lambda: { + target: body + for target, body in sink.objects().items() + if target not in flushed_before and not target.startswith(audit_prefix) + }, + lambda landed: sum(len(body.splitlines()) for body in landed.values()) >= len(before) + len(after), + seconds=30, + ) + batch_file: Final = re.compile( + rf"/{BUCKET}/{PREFIX}/2026-09-29/(\d{{2}})/batch_(\d{{2}}-\d{{2}}-\d{{2}})_[0-9a-f]{{32}}\.jsonl" + ) + layout: Final = { + (match.group(1), match.group(2)): sorted(_prompt(object_value(json.loads(line))) for line in body.splitlines()) + for target, body in batches.items() + if (match := batch_file.fullmatch(unquote(target))) + } + assert setup_audits == SETUP_AUDITS + assert answered_before == before and answered_after == after + assert sorted(str(row["request_id"]) for row in spent) == sorted(before + after) + assert sorted(upstream.received()) == sorted(before + after), "every prompt must reach the upstream exactly once" + assert pending == frozenset(), ( + f"request logs must stay queued until the audit log fills the batch: {sorted(pending)}" + ) + assert len(layout) == len(batches) == 2, tuple(batches) + assert sorted(hour for hour, _ in layout) == ["10", "11"], layout + assert len({stamp for _, stamp in layout}) == 1, f"both hour files must come from one flush: {layout}" + assert {hour: prompts for (hour, _), prompts in layout.items()} == {"10": sorted(before), "11": sorted(after)} + + def test_s3_v2_hour_slow_sink_batches_never_duplicate_an_upload(gateway: Gateway, tmp_path: Path) -> None: marker: Final = "s3hslow" + uuid.uuid4().hex[:8] upstream: Final = CountingUpstream() diff --git a/uv.lock b/uv.lock index 527f53bd372..d39c331e85e 100644 --- a/uv.lock +++ b/uv.lock @@ -4399,6 +4399,19 @@ wheels = [ { url = "https://files.pythonhosted.org/packages/41/a0/b91504515c1f9a299fc157967ffbd2f0321bce0516a3d5b89f6f4cad0355/lazy_object_proxy-1.12.0-pp39.pp310.pp311.graalpy311-none-any.whl", hash = "sha256:c3b2e0af1f7f77c4263759c4824316ce458fabe0fceadcd24ef8ca08b2d1e402", size = 15072, upload-time = "2025-08-22T13:50:05.498Z" }, ] +[[package]] +name = "libfaketime" +version = "3.0.1" +source = { registry = "https://pypi.org/simple" } +dependencies = [ + { name = "python-dateutil" }, + { name = "pytz" }, +] +sdist = { url = "https://files.pythonhosted.org/packages/df/85/b46c502837430823420be8ba856479819acc5a5fa4e46de167b0188f4ba3/libfaketime-3.0.1.tar.gz", hash = "sha256:49c2b06250bd5a1206efe5ecfdc575850517081b50b42c6c58718be1a4a09fbd", size = 75878, upload-time = "2026-08-19T16:44:13.598Z" } +wheels = [ + { url = "https://files.pythonhosted.org/packages/87/a3/af0d453324bff5ee0c1ae51be5604df56c5c6d1d5f61e430a1b96a80e2ac/libfaketime-3.0.1-py3-none-any.whl", hash = "sha256:5eed7c17f3d4ba0f6f828de4975f69a622b06b1e792674e6c1b069c0689fab79", size = 92987, upload-time = "2026-08-19T16:44:12.085Z" }, +] + [[package]] name = "librt" version = "0.15.0" @@ -4687,6 +4700,7 @@ dev = [ { name = "hypothesis" }, { name = "keyring" }, { name = "langfuse" }, + { name = "libfaketime" }, { name = "mypy" }, { name = "numpy", version = "1.26.4", source = { registry = "https://pypi.org/simple" }, marker = "python_full_version < '3.12'" }, { name = "numpy", version = "2.4.4", source = { registry = "https://pypi.org/simple" }, marker = "python_full_version >= '3.12'" }, @@ -4902,6 +4916,7 @@ dev = [ { name = "hypothesis", specifier = "==6.165.10" }, { name = "keyring", specifier = "==25.7.0" }, { name = "langfuse", specifier = ">=4.7,<5.0" }, + { name = "libfaketime", specifier = "==3.0.1" }, { name = "mypy", specifier = "==1.20.1" }, { name = "numpy", specifier = ">=1.26.0,<3.0" }, { name = "openapi-core", specifier = "==0.22.0" },