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>
This commit is contained in:
yucheng 2026-09-30 09:07:23 +00:00
parent b51762556d
commit 21a629810d
7 changed files with 438 additions and 8 deletions

View file

@ -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",

View file

@ -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"}

View file

@ -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", "")

View file

@ -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

View file

@ -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()

View file

@ -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()

15
uv.lock generated
View file

@ -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" },