mirror of
https://github.com/BerriAI/litellm.git
synced 2026-10-06 02:48:13 +00:00
* fix(s3_v2): drop terminal upload failures, bound retries per flush and enforce the queue cap at enqueue Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> * fix(s3_v2): keep retrying credential-rotation 403s, only AccessDenied-style errors are terminal Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> * test(env_keys): exclude DEFAULT_S3_MAX_FLUSH_ATTEMPTS as an internal tuning var Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> * fix(s3_v2): retry every 5xx, warn on first queue overflow, validate the flush budget Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> * fix(s3_v2): read the queue cap defensively so un-initialized loggers still enqueue Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> * fix(s3_v2): drop the getattr in _enqueue and tighten the retry tests Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> * fix(s3_v2): keep the constructor flush budget when the callback override is invalid Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> * feat(s3_v2): adapt per-object upload concurrency to sink latency and throttling Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> * refactor(s3_v2): tidy adaptive limiter Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> * perf(s3_v2): wake one waiter per released upload slot Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> * feat(s3_v2): make the enqueue queue cap configurable with s3_max_queue_size Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> * test(s3_v2): audit cells for cache hits, coded 403 and callback modes Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> * fix(s3_v2): retry bucket-wide failures by default, age-budget requeues and make terminal drops and adaptive concurrency opt-in Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> * refactor(s3_v2): suppress the missing-waiter ValueError explicitly in the adaptive limiter Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> * fix(s3_v2): count oldest events trimmed after a failed flush as callback failures Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> * style(s3_v2): move the mutable-ok marker onto the list literal it suppresses Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> * fix(s3_v2): report post-flush overflow drops once and grow adaptive concurrency above the floor before asserting back-off Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> * test(s3_v2): fail the SlowDown back-off test when the measured window sees no PUTs Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> * fix(s3_v2): restore base retry defaults, opt-in age budget, no enqueue cap, back off outside the limiter slot Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> * fix(s3_v2): hoist the default no-op upload slot to a module constant Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> * fix(s3_v2): drop unused mutable-ok suppressions on queue appends Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> * fix(s3_v2): keep the retry queue oldest-first and prioritise fresh events at upload time Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> * style(s3_v2): keep the mutable-ok marker on the queue list literal Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> * fix(s3_v2): build request bodies inside the upload slot and keep the sync retry set at base parity Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> * test(s3_v2): drop wall-clock sleeps from the unit tests Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> * fix(s3_v2): rebuild the request body inside the slot on every retry attempt Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> * test(s3_v2): anchor the backoff window on the first observed failure and tighten shard assertions Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> * fix(s3_v2): default the upload slot to the logger limiter so monkeypatched doubles keep working Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> * test(s3_v2): clear ambient AWS env credentials so the rotating profile signs the sync retry test Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> * test(s3_v2): shrink the linear send-batch perf test to 2k/8k elements Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> * test(s3_v2): fix stale batch sizes in the perf test assert message Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> * fix(s3_v2): keep async in-call retries on the base 403/500/503 set Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> * fix(s3_v2): hold the upload slot across retries, restore the bool upload contract, and fail safe on bool config Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> * fix(s3_v2): mark dropped uploads by element identity so a shared key cannot mask a retryable sibling Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> * fix(s3_v2): make the per-flush drop lookup constant time Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> * fix(s3_v2): take the upload slot in the caller like base, build the body once per attempt loop Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> * fix(s3_v2): match base retry, logging and hook behaviour unless the new options are opted in Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> * test(s3_v2): assert the wire key in the init-bypassed sync upload test Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> * test(s3_v2): assert the signed headers and wire key in the init-bypassed sync upload test Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> * fix(s3_v2): default the upload limiter at class level instead of reading it with getattr Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> * fix(s3_v2): drop the duplicate annotations that redeclare the class-level counters Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> * feat(s3_v2): drop terminal-failed uploads by default and bound retry age to one hour Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> * fix(s3_v2): fall back to the configured retry age on invalid values Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> * test(s3_v2): drive retry-age tests from a fixed clock Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> * test(s3_v2): audit cells for retry-age opt-out and 429 single-put parity Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> --------- Co-authored-by: yucheng <yucheng@berri.ai> Co-authored-by: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com>
136 lines
7.3 KiB
Python
136 lines
7.3 KiB
Python
import os
|
|
import re
|
|
import uuid
|
|
from pathlib import Path
|
|
from typing import Final
|
|
|
|
import pytest
|
|
from redis import Redis
|
|
from _s3_v2_support import (
|
|
BUCKET,
|
|
PREFIX,
|
|
SURFACES,
|
|
RecordingS3Sink,
|
|
call_surface,
|
|
collect_payloads,
|
|
matched_ids,
|
|
mixed_burst,
|
|
s3_config,
|
|
surface_reply,
|
|
)
|
|
from integration._support.client import Gateway, eventually
|
|
from integration._support.process import owned_proxy
|
|
from integration._support.wire import wire_server
|
|
|
|
PER_REQUEST_KEY: Final = re.compile(rf"^/{BUCKET}/{PREFIX}/\d{{4}}-\d{{2}}-\d{{2}}/.+\.json$")
|
|
BATCH_KEY: Final = re.compile(
|
|
rf"^/{BUCKET}/{PREFIX}/\d{{4}}-\d{{2}}-\d{{2}}/batch_\d{{2}}-\d{{2}}-\d{{2}}_[0-9a-f]{{32}}\.jsonl$"
|
|
)
|
|
|
|
|
|
@pytest.mark.covers("other.observability.s3_v2.mixed_surface_burst_bounds_puts_one_object_per_response_id")
|
|
def test_s3_v2_mixed_surface_burst_bounds_puts_one_object_per_response_id(gateway: Gateway, tmp_path: Path) -> None:
|
|
marker: Final = "s3mix" + uuid.uuid4().hex[:8]
|
|
sink: Final = RecordingS3Sink()
|
|
with wire_server(surface_reply) as provider, wire_server(sink.respond) as bucket:
|
|
config: Final = s3_config(tmp_path, bucket.url, {})
|
|
with (
|
|
owned_proxy(gateway, tmp_path, {"DEFAULT_S3_FLUSH_INTERVAL_SECONDS": "3"}, config=config) as candidate,
|
|
candidate.scenario() as scenario,
|
|
):
|
|
openai_model: Final = scenario.model(api_base=provider.url + "/v1", api_key="synthetic-provider-key")
|
|
anthropic_model: Final = scenario.model(
|
|
model="anthropic/claude-sonnet-4-5-20250929", api_base=provider.url, api_key="synthetic-provider-key"
|
|
)
|
|
key: Final = scenario.key(models=[openai_model, anthropic_model])
|
|
answered: Final = mixed_burst(candidate, openai_model, anthropic_model, key, marker)
|
|
payloads: Final = collect_payloads(sink, len(answered))
|
|
targets: Final = tuple(sink.objects())
|
|
assert sum(1 for r in provider.drain() if r.method == "POST") == 48
|
|
assert sink.peak <= 16, f"peak concurrent PUTs {sink.peak} exceeded the default bound"
|
|
assert all(PER_REQUEST_KEY.match(target) for target in targets), list(targets)
|
|
assert len(targets) == 48
|
|
assert matched_ids(payloads, answered)
|
|
|
|
|
|
@pytest.mark.covers("other.observability.s3_v2.mixed_surface_batch_writes_ndjson_lines_per_response_id")
|
|
def test_s3_v2_mixed_surface_batch_writes_ndjson_lines_per_response_id(gateway: Gateway, tmp_path: Path) -> None:
|
|
marker: Final = "s3mixb" + uuid.uuid4().hex[:8]
|
|
sink: Final = RecordingS3Sink()
|
|
with wire_server(surface_reply) as provider, wire_server(sink.respond) as bucket:
|
|
config: Final = s3_config(tmp_path, bucket.url, {"s3_batch_file_upload": True})
|
|
with (
|
|
owned_proxy(gateway, tmp_path, {"DEFAULT_S3_FLUSH_INTERVAL_SECONDS": "3"}, config=config) as candidate,
|
|
candidate.scenario() as scenario,
|
|
):
|
|
openai_model: Final = scenario.model(api_base=provider.url + "/v1", api_key="synthetic-provider-key")
|
|
anthropic_model: Final = scenario.model(
|
|
model="anthropic/claude-sonnet-4-5-20250929", api_base=provider.url, api_key="synthetic-provider-key"
|
|
)
|
|
key: Final = scenario.key(models=[openai_model, anthropic_model])
|
|
answered: Final = mixed_burst(candidate, openai_model, anthropic_model, key, marker)
|
|
payloads: Final = collect_payloads(sink, len(answered))
|
|
targets: Final = tuple(sink.objects())
|
|
puts: Final = bucket.drain()
|
|
assert sum(1 for r in provider.drain() if r.method == "POST") == 48
|
|
assert all(BATCH_KEY.match(target) for target in targets), list(targets)
|
|
assert all(put.headers["content-type"] == "application/x-ndjson" for put in puts), [put.headers for put in puts]
|
|
assert matched_ids(payloads, answered)
|
|
assert len(payloads) == 48
|
|
|
|
|
|
@pytest.mark.covers("other.observability.s3_v2.sink_outage_mid_mixed_burst_recovers_every_response_id")
|
|
def test_s3_v2_sink_outage_mid_mixed_burst_recovers_every_response_id(gateway: Gateway, tmp_path: Path) -> None:
|
|
marker: Final = "s3mixo" + uuid.uuid4().hex[:8]
|
|
sink: Final = RecordingS3Sink(fail_attempts=30, fail_status=503, delay_seconds=0.2)
|
|
with wire_server(surface_reply) as provider, wire_server(sink.respond) as bucket:
|
|
config: Final = s3_config(tmp_path, bucket.url, {})
|
|
with (
|
|
owned_proxy(gateway, tmp_path, {"DEFAULT_S3_FLUSH_INTERVAL_SECONDS": "3"}, config=config) as candidate,
|
|
candidate.scenario() as scenario,
|
|
):
|
|
openai_model: Final = scenario.model(api_base=provider.url + "/v1", api_key="synthetic-provider-key")
|
|
anthropic_model: Final = scenario.model(
|
|
model="anthropic/claude-sonnet-4-5-20250929", api_base=provider.url, api_key="synthetic-provider-key"
|
|
)
|
|
key: Final = scenario.key(models=[openai_model, anthropic_model])
|
|
answered: Final = mixed_burst(candidate, openai_model, anthropic_model, key, marker)
|
|
payloads: Final = collect_payloads(sink, len(answered), seconds=90)
|
|
assert sum(1 for r in provider.drain() if r.method == "POST") == 48
|
|
assert matched_ids(payloads, answered)
|
|
assert len(payloads) == 48, "a stored id was overwritten or duplicated"
|
|
|
|
|
|
def test_s3_v2_cache_hit_twins_log_one_object_per_request(gateway: Gateway, tmp_path: Path) -> None:
|
|
marker: Final = "s3cache" + uuid.uuid4().hex[:8]
|
|
sink: Final = RecordingS3Sink(delay_seconds=0.1)
|
|
with wire_server(surface_reply) as provider, wire_server(sink.respond) as bucket:
|
|
config: Final = s3_config(tmp_path, bucket.url, {})
|
|
with (
|
|
owned_proxy(gateway, tmp_path, {"DEFAULT_S3_FLUSH_INTERVAL_SECONDS": "3"}, config=config) as candidate,
|
|
candidate.scenario() as scenario,
|
|
):
|
|
openai_model: Final = scenario.model(api_base=provider.url + "/v1", api_key="synthetic-provider-key")
|
|
anthropic_model: Final = scenario.model(
|
|
model="anthropic/claude-sonnet-4-5-20250929", api_base=provider.url, api_key="synthetic-provider-key"
|
|
)
|
|
key: Final = scenario.key(models=[openai_model, anthropic_model])
|
|
cache: Final = Redis(host=os.environ["REDIS_HOST"], port=int(os.environ["REDIS_PORT"]))
|
|
keys_before: Final = cache.dbsize()
|
|
warmed: Final = tuple(
|
|
call_surface(candidate, surface, openai_model, anthropic_model, key, f"{marker}-{surface}")
|
|
for surface in SURFACES
|
|
)
|
|
eventually(cache.dbsize, lambda size: size >= keys_before + len(SURFACES), seconds=30)
|
|
repeated: Final = tuple(
|
|
call_surface(candidate, surface, openai_model, anthropic_model, key, f"{marker}-{surface}", False)
|
|
for surface in SURFACES
|
|
)
|
|
payloads: Final = collect_payloads(sink, 2 * len(SURFACES))
|
|
assert sum(1 for r in provider.drain() if r.method == "POST") == len(SURFACES), (
|
|
"a repeated request reached the upstream; the six repeats must all be served from cache"
|
|
)
|
|
assert len(payloads) == 12
|
|
assert sum(1 for payload in payloads if payload["cache_hit"] is True) == 6
|
|
assert sum(1 for payload in payloads if payload["cache_hit"] is not True) == 6
|
|
assert matched_ids(payloads, warmed + repeated)
|