litellm/tests/integration/observability/test_s3_v2_upload_fanout.py
devin-ai-integration[bot] e47b1f2a3f
fix(s3_v2): upload fresh events first, drop terminal failures and hour-old retries by default, opt-in adaptive concurrency (#43022)
* 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>
2026-09-26 14:58:28 -07:00

1093 lines
57 KiB
Python

import json
import re
import threading
import time
import uuid
from collections.abc import Mapping
from concurrent.futures import ThreadPoolExecutor
from dataclasses import dataclass, field
from pathlib import Path
from typing import Final
import httpx
import pytest
import yaml
from _s3_v2_support import RecordingS3Sink, collect_payloads
from _s3_v2_support import s3_config as _recording_s3_config
from integration._support.client import Gateway, JsonValue, eventually
from integration._support.process import group_members, owned_proxy, owned_proxy_process
from integration._support.wire import Reply, Request, Wire, wire_server
BUCKET: Final = "integration-bucket"
PREFIX: Final = "integration-logs"
REQUESTS: Final = 64
PUT_DELAY_SECONDS: Final = 0.5
@dataclass(slots=True)
class S3Sink:
"""Accepts every PUT after a fixed delay and records the peak number of PUTs in flight."""
lock: threading.Lock = field(default_factory=threading.Lock)
in_flight: int = 0
peak: int = 0
def respond(self, request: Request) -> Reply:
assert request.method == "PUT", request.method
assert request.target.startswith(f"/{BUCKET}/{PREFIX}/"), request.target
with self.lock:
self.in_flight += 1
self.peak = max(self.peak, self.in_flight)
time.sleep(PUT_DELAY_SECONDS)
with self.lock:
self.in_flight -= 1
return Reply()
def _chat_reply(request: Request) -> Reply:
if request.method != "POST" or not request.body:
return Reply(status=404)
text: Final = json.loads(request.body)["messages"][0]["content"]
return Reply(
body=json.dumps(
{
"id": text,
"object": "chat.completion",
"created": 1,
"model": "gpt-4o-mini",
"choices": [{"index": 0, "message": {"role": "assistant", "content": text}, "finish_reason": "stop"}],
"usage": {"prompt_tokens": 11, "completion_tokens": 4, "total_tokens": 15},
}
).encode()
)
def _s3_config(path: Path, sink_url: str, extra: Mapping[str, JsonValue]) -> Path:
config: Final = yaml.safe_load(Path("tests/integration/proxy_config.yaml").read_text())
config["litellm_settings"].update(
{
"callbacks": ["s3_v2"],
"s3_callback_params": {
"s3_bucket_name": BUCKET,
"s3_region_name": "us-east-1",
"s3_endpoint_url": sink_url,
"s3_path": PREFIX,
"s3_aws_access_key_id": "AKIAIOSFODNN7EXAMPLE",
"s3_aws_secret_access_key": "wJalrXUtnFEMI/K7MDENG/bPxRfiCYEXAMPLEKEY",
**extra,
},
}
)
target: Final = path / "s3_v2.yaml"
target.write_text(yaml.safe_dump(config))
return target
def _burst(candidate: Gateway, model: str, key: str, marker: str) -> frozenset[str]:
ids: Final = tuple(f"{marker}-{index}" for index in range(REQUESTS))
def request(identity: str) -> str:
response: Final = candidate.request(
"POST",
"/v1/chat/completions",
{"model": model, "messages": [{"role": "user", "content": identity}], "cache": {"no-cache": True}},
key=key,
)
assert response.status_code == 200, response.text
return response.json()["id"]
with ThreadPoolExecutor(max_workers=32) as pool:
returned: Final = frozenset(pool.map(request, ids))
assert returned == frozenset(ids)
return returned
def _collect(bucket: Wire, count_lines: bool, expected: int) -> tuple[Request, ...]:
puts: Final[list[Request]] = [] # mutable-ok: drain() consumes the queue, later polls must keep earlier PUTs
def delivered() -> int:
puts.extend(bucket.drain())
return sum(len(put.body.splitlines()) if count_lines else 1 for put in puts)
eventually(delivered, lambda total: total >= expected, seconds=30)
return tuple(puts)
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.flush_bounds_concurrent_puts_to_default_ceiling_and_keeps_every_log")
def test_s3_v2_flush_bounds_concurrent_puts_to_the_default_ceiling(gateway: Gateway, tmp_path: Path) -> None:
marker: Final = "s3fan" + uuid.uuid4().hex[:8]
sink: Final = S3Sink()
with wire_server(_chat_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,
):
model: Final = scenario.model(api_base=provider.url + "/v1", api_key="synthetic-provider-key")
key: Final = scenario.key(models=[model])
ids: Final = _burst(candidate, model, key, marker)
puts: Final = _collect(bucket, count_lines=False, expected=REQUESTS)
assert sum(1 for r in provider.drain() if r.method == "POST") == REQUESTS
assert sink.peak <= 16, (
f"peak concurrent PUTs {sink.peak} exceeded the default width of 16 for {REQUESTS} queued logs"
)
assert all(PER_REQUEST_KEY.match(put.target) for put in puts), [put.target for put in puts]
assert frozenset(json.loads(put.body)["id"] for put in puts) == ids
assert len({put.target for put in puts}) == REQUESTS
@pytest.mark.covers("other.observability.s3_v2.configured_bound_and_env_backed_false_keeps_per_request_objects")
def test_s3_v2_honors_configured_bound_and_env_backed_false_batch_flag(gateway: Gateway, tmp_path: Path) -> None:
marker: Final = "s3cap" + uuid.uuid4().hex[:8]
sink: Final = S3Sink()
with wire_server(_chat_reply) as provider, wire_server(sink.respond) as bucket:
config: Final = _s3_config(
tmp_path,
bucket.url,
{"s3_max_concurrent_uploads": 4, "s3_batch_file_upload": "os.environ/INTEGRATION_S3_BATCH_FILE_UPLOAD"},
)
with (
owned_proxy(
gateway,
tmp_path,
{"DEFAULT_S3_FLUSH_INTERVAL_SECONDS": "3", "INTEGRATION_S3_BATCH_FILE_UPLOAD": "false"},
config=config,
) as candidate,
candidate.scenario() as scenario,
):
model: Final = scenario.model(api_base=provider.url + "/v1", api_key="synthetic-provider-key")
key: Final = scenario.key(models=[model])
ids: Final = _burst(candidate, model, key, marker)
puts: Final = _collect(bucket, count_lines=False, expected=REQUESTS)
assert sum(1 for r in provider.drain() if r.method == "POST") == REQUESTS
assert sink.peak <= 4, f"peak concurrent PUTs {sink.peak} exceeded s3_max_concurrent_uploads=4"
assert all(PER_REQUEST_KEY.match(put.target) for put in puts), [put.target for put in puts]
assert frozenset(json.loads(put.body)["id"] for put in puts) == ids
@pytest.mark.covers("other.observability.s3_v2.batch_file_upload_writes_one_ndjson_object_per_flush")
def test_s3_v2_batch_file_upload_writes_one_jsonl_object_per_flush(gateway: Gateway, tmp_path: Path) -> None:
marker: Final = "s3jsonl" + uuid.uuid4().hex[:8]
sink: Final = S3Sink()
with wire_server(_chat_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,
):
model: Final = scenario.model(api_base=provider.url + "/v1", api_key="synthetic-provider-key")
key: Final = scenario.key(models=[model])
ids: Final = _burst(candidate, model, key, marker)
puts: Final = _collect(bucket, count_lines=True, expected=REQUESTS)
assert sum(1 for r in provider.drain() if r.method == "POST") == REQUESTS
assert len(puts) <= 2, f"{len(puts)} PUTs for {REQUESTS} logs; batch mode must write one object per flush"
assert all(BATCH_KEY.match(put.target) for put in puts), [put.target for put in puts]
assert all(put.headers["content-type"] == "application/x-ndjson" for put in puts), [put.headers for put in puts]
lines: Final = tuple(line for put in puts for line in put.body.decode().splitlines())
assert frozenset(json.loads(line)["id"] for line in lines) == ids
assert len(lines) == REQUESTS
@pytest.mark.covers("other.observability.s3_v2.batch_file_upload_keeps_team_prefix_in_object_key")
def test_s3_v2_batch_file_upload_keeps_team_alias_prefix(gateway: Gateway, tmp_path: Path) -> None:
marker: Final = "s3team" + uuid.uuid4().hex[:8]
team_alias: Final = f"alpha-{uuid.uuid4().hex[:8]}"
team_batch_key: Final = re.compile(
rf"^/{BUCKET}/{PREFIX}/{team_alias}/\d{{4}}-\d{{2}}-\d{{2}}/batch_\d{{2}}-\d{{2}}-\d{{2}}_[0-9a-f]{{32}}\.jsonl$"
)
sink: Final = S3Sink()
with wire_server(_chat_reply) as provider, wire_server(sink.respond) as bucket:
config: Final = _s3_config(tmp_path, bucket.url, {"s3_batch_file_upload": True, "s3_use_team_prefix": True})
with (
owned_proxy(gateway, tmp_path, {"DEFAULT_S3_FLUSH_INTERVAL_SECONDS": "3"}, config=config) as candidate,
candidate.scenario() as scenario,
):
model: Final = scenario.model(api_base=provider.url + "/v1", api_key="synthetic-provider-key")
team: Final = scenario.team(team_alias=team_alias, models=[model])
key: Final = scenario.key(team_id=team, models=[model])
ids: Final = _burst(candidate, model, key, marker)
puts: Final = _collect(bucket, count_lines=True, expected=REQUESTS)
assert sum(1 for r in provider.drain() if r.method == "POST") == REQUESTS
assert len(puts) >= 1
assert all(team_batch_key.match(put.target) for put in puts), [put.target for put in puts]
lines: Final = tuple(line for put in puts for line in put.body.decode().splitlines())
assert frozenset(json.loads(line)["id"] for line in lines) == ids
assert len(lines) == REQUESTS
@pytest.mark.covers("other.observability.s3_v2.upstream_failure_events_land_alongside_successes")
def test_s3_v2_upstream_failure_events_land_alongside_successes(gateway: Gateway, tmp_path: Path) -> None:
marker: Final = "s3fail" + uuid.uuid4().hex[:8]
sink: Final = RecordingS3Sink()
def provider(request: Request) -> Reply:
text: Final = json.loads(request.body)["messages"][0]["content"]
if text.endswith("-fail"):
return Reply(
status=401,
body=b'{"error": {"message": "synthetic upstream rejection", "code": "synthetic_401"}}',
)
return _chat_reply(request)
with wire_server(provider) as upstream, 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,
):
model: Final = scenario.model(api_base=upstream.url + "/v1", api_key="synthetic-provider-key")
key: Final = scenario.key(models=[model])
success_ids: Final = tuple(f"{marker}-{index}" for index in range(8))
failure_ids: Final = tuple(f"{marker}-{index}-fail" for index in range(4))
def send(identity: str) -> httpx.Response:
return candidate.request(
"POST",
"/v1/chat/completions",
{"model": model, "messages": [{"role": "user", "content": identity}], "cache": {"no-cache": True}},
key=key,
)
with ThreadPoolExecutor(max_workers=12) as pool:
responses: Final = tuple(pool.map(send, (*success_ids, *failure_ids)))
ok: Final = responses[:8]
rejected: Final = responses[8:]
assert all(response.status_code == 200 for response in ok), [r.text for r in ok]
assert tuple(response.json()["id"] for response in ok) == success_ids
for response in rejected:
assert response.status_code in (400, 401), response.status_code
assert "synthetic upstream rejection" in response.text, response.text
failure_call_ids: Final = frozenset(response.headers["x-litellm-call-id"] for response in rejected)
payloads: Final = collect_payloads(sink, len(success_ids) + len(failure_ids))
assert len(upstream.drain()) == len(success_ids) + len(failure_ids)
delivered: Final = frozenset(payload["id"] for payload in payloads if payload["status"] == "success")
assert delivered == frozenset(success_ids)
failures: Final = tuple(payload for payload in payloads if payload["status"] == "failure")
assert len(failures) == len(failure_ids)
assert frozenset(payload["litellm_call_id"] for payload in failures) == failure_call_ids
assert all("synthetic upstream rejection" in json.dumps(payload["error_information"]) for payload in failures)
@pytest.mark.covers("other.observability.s3_v2.invalid_or_empty_bound_falls_back_to_default_ceiling")
@pytest.mark.parametrize(
("bad", "warns"),
[
pytest.param("abc", True, id="non_integer"),
pytest.param(0, True, id="below_one"),
pytest.param("", False, id="empty"),
],
)
def test_s3_v2_invalid_or_empty_bound_falls_back_to_default_ceiling(
gateway: Gateway, tmp_path: Path, bad: JsonValue, warns: bool
) -> None:
marker: Final = "s3bound" + uuid.uuid4().hex[:8]
sink: Final = RecordingS3Sink()
with wire_server(_chat_reply) as provider, wire_server(sink.respond) as bucket:
config: Final = _s3_config(tmp_path, bucket.url, {"s3_max_concurrent_uploads": bad})
with (
owned_proxy_process(gateway, tmp_path, {"DEFAULT_S3_FLUSH_INTERVAL_SECONDS": "3"}, config=config) as owned,
owned.gateway.scenario() as scenario,
):
model: Final = scenario.model(api_base=provider.url + "/v1", api_key="synthetic-provider-key")
key: Final = scenario.key(models=[model])
ids: Final = _burst(owned.gateway, model, key, marker)
payloads: Final = collect_payloads(sink, REQUESTS)
if warns:
eventually(
lambda: owned.log.read_text(),
lambda text: "s3_max_concurrent_uploads" in text,
seconds=15,
)
else:
assert "s3_max_concurrent_uploads" not in owned.log.read_text()
assert sum(1 for r in provider.drain() if r.method == "POST") == REQUESTS
assert sink.peak <= 16, f"peak concurrent PUTs {sink.peak} exceeded the fallback width of 16"
assert frozenset(payload["id"] for payload in payloads) == ids
@pytest.mark.covers("other.observability.s3_v2.sink_rejection_requeues_and_delivers_every_id_once")
def test_s3_v2_sink_rejection_requeues_and_delivers_every_id_once(gateway: Gateway, tmp_path: Path) -> None:
marker: Final = "s3deny" + uuid.uuid4().hex[:8]
sink: Final = RecordingS3Sink(fail_status=503, delay_seconds=0.2)
with wire_server(_chat_reply) as provider, wire_server(sink.respond) as bucket:
config: Final = _s3_config(tmp_path, bucket.url, {})
with (
owned_proxy_process(gateway, tmp_path, {"DEFAULT_S3_FLUSH_INTERVAL_SECONDS": "3"}, config=config) as owned,
owned.gateway.scenario() as scenario,
):
model: Final = scenario.model(api_base=provider.url + "/v1", api_key="synthetic-provider-key")
key: Final = scenario.key(models=[model])
sink.fail_until = time.time() + 10
ids: Final = _burst(owned.gateway, model, key, marker)
payloads: Final = collect_payloads(sink, REQUESTS, seconds=90)
eventually(
lambda: owned.log.read_text(),
lambda text: "S3BatchUploadError" in text,
seconds=15,
)
readiness: Final = owned.gateway.client.get("/health/readiness")
assert readiness.status_code == 200, readiness.text
assert sum(1 for r in provider.drain() if r.method == "POST") == REQUESTS
assert len(sink.objects()) == REQUESTS
assert frozenset(payload["id"] for payload in payloads) == ids
@dataclass(slots=True)
class RejectingS3Sink:
"""Answers every PUT whose body carries `reject_marker` with `reject_status`, accepts the rest,
and counts the rejected attempts so a test can see whether the proxy keeps re-sending them."""
reject_marker: str
reject_status: int
reject_code: str = "AccessDenied"
reject_until: float = float("inf")
lock: threading.Lock = field(default_factory=threading.Lock)
rejected_attempts: int = 0
rejected_times: list[float] = field(default_factory=list) # mutable-ok: appended under lock per rejected PUT
store: dict[str, bytes] = field(default_factory=dict) # mutable-ok: later PUTs must be visible to earlier polls
def respond(self, request: Request) -> Reply:
assert request.method == "PUT", request.method
with self.lock:
if self.reject_marker.encode() in request.body and time.time() < self.reject_until:
self.rejected_attempts += 1
self.rejected_times.append(time.time())
return Reply(status=self.reject_status, body=f"<Error><Code>{self.reject_code}</Code></Error>".encode())
self.store[request.target] = request.body
return Reply()
def landed_ids(self) -> frozenset[str]:
with self.lock:
bodies: Final = tuple(self.store.values())
return frozenset(json.loads(line)["id"] for body in bodies for line in body.splitlines())
def _send(candidate: Gateway, model: str, key: str, identity: str) -> None:
response: Final = candidate.request(
"POST",
"/v1/chat/completions",
{"model": model, "messages": [{"role": "user", "content": identity}], "cache": {"no-cache": True}},
key=key,
)
assert response.status_code == 200, response.text
def _send_and_wait_until_landed(candidate: Gateway, model: str, key: str, sink: RejectingS3Sink, identity: str) -> None:
_send(candidate, model, key, identity)
eventually(sink.landed_ids, lambda landed: identity in landed, seconds=60)
@pytest.mark.parametrize(
("status", "code"),
[
pytest.param(403, "AccessDenied", id="access_denied"),
pytest.param(404, "NoSuchBucket", id="no_such_bucket"),
pytest.param(400, "KMS.DisabledException", id="kms_disabled"),
],
)
def test_s3_v2_object_rejected_with_a_bucket_wide_code_is_delivered_once_the_fault_clears(
gateway: Gateway, tmp_path: Path, status: int, code: str
) -> None:
marker: Final = "s3fault" + uuid.uuid4().hex[:8]
sink: Final = RejectingS3Sink(reject_marker=f"{marker}-denied", reject_status=status, reject_code=code)
with wire_server(_chat_reply) as provider, wire_server(sink.respond) as bucket:
config: Final = _s3_config(tmp_path, bucket.url, {"s3_batch_file_upload": False})
with (
owned_proxy(gateway, tmp_path, {"DEFAULT_S3_FLUSH_INTERVAL_SECONDS": "2"}, config=config) as candidate,
candidate.scenario() as scenario,
):
model: Final = scenario.model(api_base=provider.url + "/v1", api_key="synthetic-provider-key")
key: Final = scenario.key(models=[model])
_send(candidate, model, key, f"{marker}-denied")
_send_and_wait_until_landed(candidate, model, key, sink, f"{marker}-first-flush")
eventually(lambda: sink.rejected_attempts, lambda attempts: attempts >= 2, seconds=30)
sink.reject_until = time.time()
eventually(sink.landed_ids, lambda landed: f"{marker}-denied" in landed, seconds=60)
readiness: Final = candidate.client.get("/health/readiness")
assert readiness.status_code == 200, readiness.text
assert sum(1 for r in provider.drain() if r.method == "POST") == 2
assert sink.landed_ids() == {f"{marker}-denied", f"{marker}-first-flush"}
def test_s3_v2_terminal_object_is_put_once_and_dropped_by_default(gateway: Gateway, tmp_path: Path) -> None:
marker: Final = "s3toolarge" + uuid.uuid4().hex[:8]
sink: Final = RejectingS3Sink(reject_marker=f"{marker}-huge", reject_status=400, reject_code="EntityTooLarge")
with wire_server(_chat_reply) as provider, wire_server(sink.respond) as bucket:
config: Final = _s3_config(tmp_path, bucket.url, {"s3_batch_file_upload": False})
with (
owned_proxy(gateway, tmp_path, {"DEFAULT_S3_FLUSH_INTERVAL_SECONDS": "2"}, config=config) as candidate,
candidate.scenario() as scenario,
):
model: Final = scenario.model(api_base=provider.url + "/v1", api_key="synthetic-provider-key")
key: Final = scenario.key(models=[model])
_send(candidate, model, key, f"{marker}-huge")
_send_and_wait_until_landed(candidate, model, key, sink, f"{marker}-sibling")
_send_and_wait_until_landed(candidate, model, key, sink, f"{marker}-second-flush")
_send_and_wait_until_landed(candidate, model, key, sink, f"{marker}-third-flush")
assert sum(1 for r in provider.drain() if r.method == "POST") == 4
assert sink.rejected_attempts == 1, (
f"an EntityTooLarge object was PUT {sink.rejected_attempts} times next to delivered siblings; "
"with the default s3_drop_on_terminal_error it must be attempted once and dropped"
)
def test_s3_v2_terminal_object_keeps_retrying_when_opted_out(gateway: Gateway, tmp_path: Path) -> None:
marker: Final = "s3keep" + uuid.uuid4().hex[:8]
sink: Final = RejectingS3Sink(reject_marker=f"{marker}-huge", reject_status=400, reject_code="EntityTooLarge")
with wire_server(_chat_reply) as provider, wire_server(sink.respond) as bucket:
config: Final = _s3_config(
tmp_path, bucket.url, {"s3_batch_file_upload": False, "s3_drop_on_terminal_error": False}
)
with (
owned_proxy(gateway, tmp_path, {"DEFAULT_S3_FLUSH_INTERVAL_SECONDS": "2"}, config=config) as candidate,
candidate.scenario() as scenario,
):
model: Final = scenario.model(api_base=provider.url + "/v1", api_key="synthetic-provider-key")
key: Final = scenario.key(models=[model])
_send(candidate, model, key, f"{marker}-huge")
_send_and_wait_until_landed(candidate, model, key, sink, f"{marker}-sibling")
_send_and_wait_until_landed(candidate, model, key, sink, f"{marker}-second-flush")
eventually(lambda: sink.rejected_attempts, lambda attempts: attempts >= 2, seconds=30)
assert sum(1 for r in provider.drain() if r.method == "POST") == 3
def test_s3_v2_aged_out_object_is_dropped_next_to_delivered_siblings(gateway: Gateway, tmp_path: Path) -> None:
marker: Final = "s3aged" + uuid.uuid4().hex[:8]
sink: Final = RejectingS3Sink(reject_marker=f"{marker}-doomed", reject_status=503, reject_code="InternalError")
with wire_server(_chat_reply) as provider, wire_server(sink.respond) as bucket:
config: Final = _s3_config(tmp_path, bucket.url, {"s3_batch_file_upload": False, "s3_max_retry_age_seconds": 1})
with (
owned_proxy_process(gateway, tmp_path, {"DEFAULT_S3_FLUSH_INTERVAL_SECONDS": "2"}, config=config) as owned,
owned.gateway.scenario() as scenario,
):
model: Final = scenario.model(api_base=provider.url + "/v1", api_key="synthetic-provider-key")
key: Final = scenario.key(models=[model])
_send(owned.gateway, model, key, f"{marker}-doomed")
_send_and_wait_until_landed(owned.gateway, model, key, sink, f"{marker}-sibling")
_send(owned.gateway, model, key, f"{marker}-trigger")
eventually(
lambda: owned.log.read_text(),
lambda text: "retrying longer than s3_max_retry_age_seconds=1)" in text,
seconds=60,
)
exhausted: Final = sink.rejected_attempts
_send_and_wait_until_landed(owned.gateway, model, key, sink, f"{marker}-one-flush-later")
_send_and_wait_until_landed(owned.gateway, model, key, sink, f"{marker}-two-flushes-later")
assert sum(1 for r in provider.drain() if r.method == "POST") == 5
assert 3 <= exhausted <= 3 * 3, f"{exhausted} PUTs for an object that aged out after its second flush"
assert sink.rejected_attempts == exhausted, (
f"a 503 object kept being PUT after ageing out: {exhausted} -> {sink.rejected_attempts}"
)
def test_s3_v2_aged_out_object_stays_queued_while_the_whole_sink_is_down(gateway: Gateway, tmp_path: Path) -> None:
marker: Final = "s3down" + uuid.uuid4().hex[:8]
sink: Final = RecordingS3Sink(fail_status=503, delay_seconds=0.1)
with wire_server(_chat_reply) as provider, wire_server(sink.respond) as bucket:
config: Final = _s3_config(tmp_path, bucket.url, {"s3_batch_file_upload": False, "s3_max_retry_age_seconds": 1})
with (
owned_proxy(gateway, tmp_path, {"DEFAULT_S3_FLUSH_INTERVAL_SECONDS": "2"}, config=config) as candidate,
candidate.scenario() as scenario,
):
model: Final = scenario.model(api_base=provider.url + "/v1", api_key="synthetic-provider-key")
key: Final = scenario.key(models=[model])
sink.fail_until = time.time() + 12
ids: Final = _push(candidate, model, key, marker, 4)
payloads: Final = collect_payloads(sink, 4, seconds=90)
assert sum(1 for r in provider.drain() if r.method == "POST") == 4
assert frozenset(payload["id"] for payload in payloads) == ids, (
"a bucket-wide outage longer than the age budget lost events"
)
def test_s3_v2_failing_sink_trims_the_oldest_failed_events_past_the_queue_cap(gateway: Gateway, tmp_path: Path) -> None:
marker: Final = "s3cap" + uuid.uuid4().hex[:8]
sink: Final = RecordingS3Sink(fail_status=503, delay_seconds=0.05)
with wire_server(_chat_reply) as provider, wire_server(sink.respond) as bucket:
config: Final = _s3_config(tmp_path, bucket.url, {"s3_batch_file_upload": False, "s3_max_queue_size": 4})
with (
owned_proxy_process(gateway, tmp_path, {"DEFAULT_S3_FLUSH_INTERVAL_SECONDS": "2"}, config=config) as owned,
owned.gateway.scenario() as scenario,
):
model: Final = scenario.model(api_base=provider.url + "/v1", api_key="synthetic-provider-key")
key: Final = scenario.key(models=[model])
sink.fail_until = time.time() + 15
_send(owned.gateway, model, key, f"{marker}-probe")
eventually(lambda: owned.log.read_text(), lambda text: "S3BatchUploadError" in text, seconds=30)
for index in range(24):
_send(owned.gateway, model, key, f"{marker}-{index}")
eventually(
lambda: owned.log.read_text(),
lambda text: "after a failed flush, dropped" in text,
seconds=30,
)
payloads: Final = collect_payloads(sink, 4, seconds=90)
landed: Final = frozenset(payload["id"] for payload in payloads)
assert sum(1 for r in provider.drain() if r.method == "POST") == 25
assert len(landed) == 4, f"{len(landed)} objects landed with s3_max_queue_size=4"
assert f"{marker}-probe" not in landed and f"{marker}-0" not in landed, (
f"the oldest events survived the cap: {landed}"
)
assert f"{marker}-23" in landed, f"the newest event was dropped: {landed}"
def test_s3_v2_retry_age_zero_keeps_aged_object_queued(gateway: Gateway, tmp_path: Path) -> None:
marker: Final = "s3agezero" + uuid.uuid4().hex[:8]
sink: Final = RejectingS3Sink(reject_marker=f"{marker}-doomed", reject_status=503, reject_code="SlowDown")
with wire_server(_chat_reply) as provider, wire_server(sink.respond) as bucket:
config: Final = _s3_config(tmp_path, bucket.url, {"s3_batch_file_upload": False, "s3_max_retry_age_seconds": 0})
with (
owned_proxy_process(gateway, tmp_path, {"DEFAULT_S3_FLUSH_INTERVAL_SECONDS": "2"}, config=config) as owned,
owned.gateway.scenario() as scenario,
):
model: Final = scenario.model(api_base=provider.url + "/v1", api_key="synthetic-provider-key")
key: Final = scenario.key(models=[model])
_send(owned.gateway, model, key, f"{marker}-doomed")
_send_and_wait_until_landed(owned.gateway, model, key, sink, f"{marker}-sibling")
_send_and_wait_until_landed(owned.gateway, model, key, sink, f"{marker}-second-flush")
_send_and_wait_until_landed(owned.gateway, model, key, sink, f"{marker}-third-flush")
attempts_before_clear: Final = sink.rejected_attempts
log_text: Final = owned.log.read_text()
assert "uploads dropped" not in log_text, log_text
assert "retrying longer than" not in log_text, log_text
sink.reject_until = time.time()
eventually(sink.landed_ids, lambda landed: f"{marker}-doomed" in landed, seconds=60)
assert sum(1 for r in provider.drain() if r.method == "POST") == 4
assert attempts_before_clear >= 3, (
f"only {attempts_before_clear} PUTs for an object that stayed queued through three delivered flushes; "
"with s3_max_retry_age_seconds=0 it must keep retrying longer than any enabled budget"
)
def test_s3_v2_throttled_429_object_is_put_once_per_flush(gateway: Gateway, tmp_path: Path) -> None:
marker: Final = "s3throttle" + uuid.uuid4().hex[:8]
sink: Final = RejectingS3Sink(reject_marker=f"{marker}-throttled", reject_status=429, reject_code="TooManyRequests")
with wire_server(_chat_reply) as provider, wire_server(sink.respond) as bucket:
config: Final = _s3_config(tmp_path, bucket.url, {"s3_batch_file_upload": False})
with (
owned_proxy(gateway, tmp_path, {"DEFAULT_S3_FLUSH_INTERVAL_SECONDS": "2"}, config=config) as candidate,
candidate.scenario() as scenario,
):
model: Final = scenario.model(api_base=provider.url + "/v1", api_key="synthetic-provider-key")
key: Final = scenario.key(models=[model])
_send(candidate, model, key, f"{marker}-throttled")
_send_and_wait_until_landed(candidate, model, key, sink, f"{marker}-sibling")
_send_and_wait_until_landed(candidate, model, key, sink, f"{marker}-second-flush")
_send_and_wait_until_landed(candidate, model, key, sink, f"{marker}-third-flush")
sink.reject_until = time.time()
eventually(sink.landed_ids, lambda landed: f"{marker}-throttled" in landed, seconds=60)
assert sum(1 for r in provider.drain() if r.method == "POST") == 4
times: Final = tuple(sink.rejected_times)
gaps: Final = tuple(round(later - earlier, 3) for earlier, later in zip(times, times[1:]))
assert len(times) >= 3 and min(gaps) >= 1.5, (
f"PUTs for a 429 object ran {gaps} apart; the 2 s flush interval allows exactly one attempt per flush "
"because 429 is not an in-call retry status"
)
def test_s3_v2_default_config_retries_access_denied_and_every_event_lands(gateway: Gateway, tmp_path: Path) -> None:
marker: Final = "s3denied" + uuid.uuid4().hex[:8]
sink: Final = RecordingS3Sink(fail_status=403, fail_code="AccessDenied", delay_seconds=0.05)
with wire_server(_chat_reply) as provider, wire_server(sink.respond) as bucket:
config: Final = _s3_config(tmp_path, bucket.url, {"s3_batch_file_upload": False})
with (
owned_proxy(gateway, tmp_path, {"DEFAULT_S3_FLUSH_INTERVAL_SECONDS": "2"}, config=config) as candidate,
candidate.scenario() as scenario,
):
model: Final = scenario.model(api_base=provider.url + "/v1", api_key="synthetic-provider-key")
key: Final = scenario.key(models=[model])
sink.fail_until = time.time() + 20
ids: Final = _push(candidate, model, key, marker, 8)
payloads: Final = collect_payloads(sink, 8, seconds=120)
assert sum(1 for r in provider.drain() if r.method == "POST") == 8
assert frozenset(payload["id"] for payload in payloads) == ids, (
f"a default-config run lost events through a 20s AccessDenied outage: {len(payloads)} landed"
)
attempt_totals: Final = tuple(sorted(sink.attempt_counts.values()))
assert len(attempt_totals) == 8 and all(count >= 4 for count in attempt_totals), (
f"each object must see at least one full 3-PUT retry burst before landing: {attempt_totals}"
)
@pytest.mark.covers("other.observability.s3_v2.batch_retry_resends_identical_key_and_body")
def test_s3_v2_batch_retry_resends_identical_key_and_body(gateway: Gateway, tmp_path: Path) -> None:
marker: Final = "s3retry" + uuid.uuid4().hex[:8]
sink: Final = RecordingS3Sink(fail_status=500, delay_seconds=0.2)
with wire_server(_chat_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,
):
model: Final = scenario.model(api_base=provider.url + "/v1", api_key="synthetic-provider-key")
key: Final = scenario.key(models=[model])
sink.fail_until = time.time() + 8
ids: Final = _burst(candidate, model, key, marker)
payloads: Final = collect_payloads(sink, REQUESTS, seconds=90)
puts: Final = bucket.drain()
assert sum(1 for r in provider.drain() if r.method == "POST") == REQUESTS
by_target: Final = {}
for put in puts:
by_target.setdefault(put.target, set()).add(put.body) # mutable-ok: grouping attempts seen so far per target
assert all(len(bodies) == 1 for bodies in by_target.values()), "a retried batch PUT changed key or body"
assert max(sum(1 for put in puts if put.target == target) for target in by_target) >= 2, "no retried PUT observed"
assert frozenset(payload["id"] for payload in payloads) == ids
assert len(payloads) == REQUESTS
@pytest.mark.covers("other.observability.s3_v2.unknown_model_rejection_keeps_other_requests_logging")
def test_s3_v2_unknown_model_rejection_keeps_other_requests_logging(gateway: Gateway, tmp_path: Path) -> None:
marker: Final = "s3ghost" + uuid.uuid4().hex[:8]
sink: Final = RecordingS3Sink()
with wire_server(_chat_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,
):
model: Final = scenario.model(api_base=provider.url + "/v1", api_key="synthetic-provider-key")
key: Final = scenario.key(models=[model])
ghost: Final = candidate.request(
"POST",
"/v1/chat/completions",
{"model": f"ghost-{uuid.uuid4().hex}", "messages": [{"role": "user", "content": "hi"}]},
key=key,
)
assert ghost.status_code in (400, 403, 404), ghost.text
ids: Final = _burst(candidate, model, key, marker)
eventually(
lambda: frozenset(payload["id"] for payload in sink.payloads()),
lambda landed: ids <= landed,
seconds=90,
)
payloads: Final = sink.payloads()
assert sum(1 for r in provider.drain() if r.method == "POST") == REQUESTS
assert ids <= frozenset(payload["id"] for payload in payloads)
extras: Final = tuple(payload for payload in payloads if payload["id"] not in ids)
assert all(payload["status"] == "failure" for payload in extras), extras
@pytest.mark.covers("other.observability.s3_v2.batch_flag_ignored_when_s3_v2_is_cold_storage_logger")
def test_s3_v2_batch_flag_ignored_when_s3_v2_is_cold_storage_logger(gateway: Gateway, tmp_path: Path) -> None:
marker: Final = "s3cold" + uuid.uuid4().hex[:8]
sink: Final = RecordingS3Sink()
with wire_server(_chat_reply) as provider, wire_server(sink.respond) as bucket:
config: Final = _recording_s3_config(
tmp_path,
bucket.url,
{"s3_batch_file_upload": True},
{"cold_storage_custom_logger": "s3_v2"},
)
with (
owned_proxy_process(gateway, tmp_path, {"DEFAULT_S3_FLUSH_INTERVAL_SECONDS": "3"}, config=config) as owned,
owned.gateway.scenario() as scenario,
):
model: Final = scenario.model(api_base=provider.url + "/v1", api_key="synthetic-provider-key")
key: Final = scenario.key(models=[model])
response: Final = owned.gateway.request(
"POST",
"/v1/chat/completions",
{"model": model, "messages": [{"role": "user", "content": marker}], "cache": {"no-cache": True}},
key=key,
)
assert response.status_code == 200, response.text
request_id: Final = str(response.json()["id"])
payloads: Final = collect_payloads(sink, 1)
assert all(PER_REQUEST_KEY.match(target) for target in sink.objects()), list(sink.objects())
eventually(
lambda: owned.log.read_text(),
lambda text: "s3_batch_file_upload is ignored because s3_v2 is the cold storage logger" in text,
seconds=15,
)
spend: Final = eventually(
lambda: owned.gateway.request("GET", f"/spend/logs/ui/{request_id}"),
lambda reply: reply.status_code == 200 and bool((reply.json() or {}).get("messages")),
seconds=60,
)
assert spend.status_code == 200, spend.text
body: Final = spend.json()
assert body["messages"], spend.text
assert body["response"], spend.text
assert payloads[0]["id"] == request_id
@pytest.mark.covers("other.observability.s3_v2.identical_requests_land_distinct_objects")
def test_s3_v2_identical_requests_land_distinct_objects(gateway: Gateway, tmp_path: Path) -> None:
marker: Final = "s3same" + uuid.uuid4().hex[:8]
sink: Final = RecordingS3Sink()
with wire_server(_chat_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,
):
model: Final = scenario.model(api_base=provider.url + "/v1", api_key="synthetic-provider-key")
key: Final = scenario.key(models=[model])
def send(_: int) -> str:
response: Final = candidate.request(
"POST",
"/v1/chat/completions",
{"model": model, "messages": [{"role": "user", "content": marker}], "cache": {"no-cache": True}},
key=key,
)
assert response.status_code == 200, response.text
return str(response.json()["id"])
with ThreadPoolExecutor(max_workers=16) as pool:
returned: Final = frozenset(pool.map(send, range(16)))
payloads: Final = collect_payloads(sink, 16)
assert sum(1 for r in provider.drain() if r.method == "POST") == 16
assert returned == {marker}, "the upstream echo keeps the same id for identical requests"
assert len(sink.objects()) == 16, "identical requests must still land as distinct objects"
assert all(payload["id"] == marker for payload in payloads)
@pytest.mark.covers("other.observability.s3_v2.two_workers_bound_and_deliver_every_id")
def test_s3_v2_two_workers_bound_and_deliver_every_id(gateway: Gateway, tmp_path: Path) -> None:
marker: Final = "s3work" + uuid.uuid4().hex[:8]
sink: Final = RecordingS3Sink()
with wire_server(_chat_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, workers=2
) as candidate,
candidate.scenario() as scenario,
):
model: Final = scenario.model(api_base=provider.url + "/v1", api_key="synthetic-provider-key")
key: Final = scenario.key(models=[model])
ids: Final = _burst(candidate, model, key, marker)
payloads: Final = collect_payloads(sink, REQUESTS)
assert sum(1 for r in provider.drain() if r.method == "POST") == REQUESTS
assert sink.peak <= 32, f"peak concurrent PUTs {sink.peak} exceeded two workers at the default bound"
assert len(sink.objects()) == REQUESTS
assert frozenset(payload["id"] for payload in payloads) == ids
@pytest.mark.covers("other.observability.s3_v2.slow_sink_never_duplicates_or_stalls_readiness")
def test_s3_v2_slow_sink_never_duplicates_or_stalls_readiness(gateway: Gateway, tmp_path: Path) -> None:
marker: Final = "s3slow" + uuid.uuid4().hex[:8]
sink: Final = RecordingS3Sink(delay_seconds=1.5)
with wire_server(_chat_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": "1"}, config=config) as candidate,
candidate.scenario() as scenario,
):
model: Final = scenario.model(api_base=provider.url + "/v1", api_key="synthetic-provider-key")
key: Final = scenario.key(models=[model])
ids: Final = _burst(candidate, model, key, marker)
def delivered() -> int:
readiness: Final = candidate.client.get("/health/readiness")
assert readiness.status_code == 200, readiness.text
return sum(len(body.splitlines()) for body in sink.objects().values())
eventually(delivered, lambda total: total >= REQUESTS, seconds=90)
payloads: Final = sink.payloads()
puts: Final = bucket.drain()
targets: Final = tuple(put.target for put in puts)
assert sum(1 for r in provider.drain() if r.method == "POST") == REQUESTS
assert len(set(targets)) == len(targets), "the same object was PUT more than once"
assert frozenset(payload["id"] for payload in payloads) == ids
assert len(payloads) == REQUESTS
@pytest.mark.covers("other.observability.s3_v2.worker_kill_mid_burst_keeps_surviving_deliveries")
def test_s3_v2_worker_kill_mid_burst_keeps_surviving_deliveries(gateway: Gateway, tmp_path: Path) -> None:
marker: Final = "s3kill" + uuid.uuid4().hex[:8]
sink: Final = RecordingS3Sink()
with wire_server(_chat_reply) as provider, wire_server(sink.respond) as bucket:
config: Final = _s3_config(tmp_path, bucket.url, {})
with (
owned_proxy_process(
gateway, tmp_path, {"DEFAULT_S3_FLUSH_INTERVAL_SECONDS": "3"}, config=config, workers=2
) as owned,
owned.gateway.scenario() as scenario,
):
model: Final = scenario.model(api_base=provider.url + "/v1", api_key="synthetic-provider-key")
key: Final = scenario.key(models=[model])
sent: Final = tuple(f"{marker}-{index}" for index in range(REQUESTS))
def send(identity: str) -> tuple[str, bool]:
try:
response: Final = owned.gateway.request(
"POST",
"/v1/chat/completions",
{
"model": model,
"messages": [{"role": "user", "content": identity}],
"cache": {"no-cache": True},
},
key=key,
)
except Exception:
return identity, False
return identity, response.status_code == 200
with ThreadPoolExecutor(max_workers=32) as pool:
futures: Final = tuple(pool.submit(send, identity) for identity in sent)
time.sleep(0.5)
children: Final = tuple(
process for process in group_members(owned.process.pid) if process.pid != owned.process.pid
)
assert children, "no worker children found to kill"
children[0].kill()
results: Final = tuple(future.result() for future in futures)
survivors: Final = frozenset(identity for identity, ok in results if ok)
assert survivors, "no request survived the worker kill"
readiness: Final = owned.gateway.client.get("/health/readiness")
assert readiness.status_code == 200, readiness.text
payloads: Final = collect_payloads(sink, len(survivors), seconds=90)
landed: Final = frozenset(payload["id"] for payload in payloads)
assert survivors <= landed, "an id whose response succeeded never landed"
assert landed <= frozenset(sent), "an id that was never sent landed"
@pytest.mark.covers("other.observability.s3_v2.sigterm_mid_burst_loses_only_inflight_without_duplicates")
def test_s3_v2_sigterm_mid_burst_loses_only_inflight_without_duplicates(gateway: Gateway, tmp_path: Path) -> None:
marker: Final = "s3term" + uuid.uuid4().hex[:8]
sink: Final = RecordingS3Sink()
with wire_server(_chat_reply) as provider, wire_server(sink.respond) as bucket:
config: Final = _s3_config(tmp_path, bucket.url, {})
owned: Final = owned_proxy_process(gateway, tmp_path, {"DEFAULT_S3_FLUSH_INTERVAL_SECONDS": "3"}, config=config)
candidate_owned: Final = owned.__enter__()
try:
created: Final = candidate_owned.gateway.post(
"/model/new",
{
"model_name": f"integration-{marker}",
"litellm_params": {
"model": "openai/gpt-4o-mini",
"api_key": "synthetic-provider-key",
"api_base": provider.url + "/v1",
},
"model_info": {},
},
)
model: Final = str(created["model_name"])
key: Final = str(candidate_owned.gateway.post("/key/generate", {"models": [model]})["key"])
sent: Final = tuple(f"{marker}-{index}" for index in range(REQUESTS))
def send(identity: str) -> tuple[str, bool]:
try:
response: Final = candidate_owned.gateway.request(
"POST",
"/v1/chat/completions",
{
"model": model,
"messages": [{"role": "user", "content": identity}],
"cache": {"no-cache": True},
},
key=key,
)
except Exception:
return identity, False
return identity, response.status_code == 200
with ThreadPoolExecutor(max_workers=32) as pool:
futures: Final = tuple(pool.submit(send, identity) for identity in sent)
time.sleep(0.5)
candidate_owned.process.terminate()
results: Final = tuple(future.result() for future in futures)
candidate_owned.process.wait(timeout=30)
finally:
owned.__exit__(None, None, None)
answered: Final = frozenset(identity for identity, ok in results if ok)
landed: Final = frozenset(payload["id"] for payload in sink.payloads())
assert landed <= answered, (
"a delivered object has no matching answered request; lost in-flight ids are expected, extras are not"
)
targets: Final = tuple(sink.objects())
assert len(set(targets)) == len(targets), "the same object was PUT more than once"
RAMP_REQUESTS: Final = 400
RAMP_PUT_DELAY_SECONDS: Final = 1.0
def _push(candidate: Gateway, model: str, key: str, marker: str, count: int) -> frozenset[str]:
ids: Final = tuple(f"{marker}-{index}" for index in range(count))
def request(identity: str) -> str:
response: Final = candidate.request(
"POST",
"/v1/chat/completions",
{"model": model, "messages": [{"role": "user", "content": identity}], "cache": {"no-cache": True}},
key=key,
)
assert response.status_code == 200, response.text
return response.json()["id"]
with ThreadPoolExecutor(max_workers=64) as pool:
returned: Final = frozenset(pool.map(request, ids))
assert returned == frozenset(ids)
return returned
def test_s3_v2_slow_sink_ramps_concurrency_and_drains_the_backlog(gateway: Gateway, tmp_path: Path) -> None:
marker: Final = "s3ramp" + uuid.uuid4().hex[:8]
sink: Final = RecordingS3Sink(delay_seconds=RAMP_PUT_DELAY_SECONDS)
with wire_server(_chat_reply) as provider, wire_server(sink.respond) as bucket:
config: Final = _s3_config(
tmp_path, bucket.url, {"s3_batch_file_upload": False, "s3_adaptive_concurrency": True}
)
with (
owned_proxy(
gateway,
tmp_path,
{"DEFAULT_S3_FLUSH_INTERVAL_SECONDS": "1", "DEFAULT_S3_BATCH_SIZE": "5000"},
config=config,
) as candidate,
candidate.scenario() as scenario,
):
model: Final = scenario.model(api_base=provider.url + "/v1", api_key="synthetic-provider-key")
key: Final = scenario.key(models=[model])
ids: Final = _push(candidate, model, key, marker, RAMP_REQUESTS)
drain_started: Final = time.monotonic()
payloads: Final = collect_payloads(sink, RAMP_REQUESTS, seconds=180)
drained_seconds: Final = time.monotonic() - drain_started
fixed_sixteen_estimate: Final = RAMP_REQUESTS * RAMP_PUT_DELAY_SECONDS / 16
assert sum(1 for r in provider.drain() if r.method == "POST") == RAMP_REQUESTS
assert frozenset(payload["id"] for payload in payloads) == ids
assert sink.peak > 16, f"adaptive limiter never ramped past the old fixed bound: peak {sink.peak}"
assert drained_seconds < 2 * fixed_sixteen_estimate, (
f"backlog of {RAMP_REQUESTS} drained in {drained_seconds:.1f}s with peak concurrency {sink.peak}; "
f"even a fixed bound of 16 would need only ~{fixed_sixteen_estimate:.1f}s, so the uploads stalled"
)
def test_s3_v2_throttled_sink_halves_in_flight_puts(gateway: Gateway, tmp_path: Path) -> None:
marker: Final = "s3throt" + uuid.uuid4().hex[:8]
sink: Final = RecordingS3Sink(fail_status=503, fail_code="SlowDown", delay_seconds=0.3)
with wire_server(_chat_reply) as provider, wire_server(sink.respond) as bucket:
config: Final = _s3_config(
tmp_path,
bucket.url,
{"s3_batch_file_upload": False, "s3_adaptive_concurrency": True, "s3_max_concurrent_uploads": 4},
)
with (
owned_proxy(gateway, tmp_path, {"DEFAULT_S3_FLUSH_INTERVAL_SECONDS": "2"}, config=config) as candidate,
candidate.scenario() as scenario,
):
model: Final = scenario.model(api_base=provider.url + "/v1", api_key="synthetic-provider-key")
key: Final = scenario.key(models=[model])
healthy_ids: Final = _push(candidate, model, key, f"{marker}-healthy", REQUESTS)
collect_payloads(sink, REQUESTS)
healthy_peak: Final = sink.peak
window_start: Final = time.time()
sink.fail_until = window_start + 60
throttled_ids: Final = _push(candidate, model, key, f"{marker}-throttled", REQUESTS)
first_fail_at: Final = eventually(
lambda: next((when for when, _ in sink.attempt_log if when >= window_start), None),
lambda when: when is not None,
seconds=30,
)
window_end: Final = first_fail_at + 8.0
sink.fail_until = window_end
payloads: Final = collect_payloads(sink, 2 * REQUESTS, seconds=120)
throttled_peak: Final = sink.peak_between(first_fail_at + 5.0, window_end)
throttled_attempts: Final = sum(
1 for when, _ in sink.attempt_log if first_fail_at + 5.0 <= when < window_end
)
assert sum(1 for r in provider.drain() if r.method == "POST") == 2 * REQUESTS
assert healthy_peak > 4, (
f"healthy peak {healthy_peak} never rose above the configured width 4; nothing to back off from"
)
assert throttled_attempts > 0, (
"no PUTs observed in the measured SlowDown window; the back-off assertion would be vacuous"
)
assert throttled_peak < healthy_peak, (
f"in-flight PUTs during the SlowDown window peaked at {throttled_peak}, not below the healthy peak "
f"{healthy_peak}; the limiter did not back off"
)
assert frozenset(payload["id"] for payload in payloads) == healthy_ids | throttled_ids
assert len(sink.objects()) == 2 * REQUESTS
def test_s3_v2_coded_403_is_transient_and_every_id_lands_once(gateway: Gateway, tmp_path: Path) -> None:
marker: Final = "s3coded" + uuid.uuid4().hex[:8]
sink: Final = RecordingS3Sink(fail_attempts=3, fail_status=403, fail_code="RequestTimeout", delay_seconds=0.1)
with wire_server(_chat_reply) as provider, wire_server(sink.respond) as bucket:
config: Final = _s3_config(tmp_path, bucket.url, {"s3_batch_file_upload": False})
with (
owned_proxy(gateway, tmp_path, {"DEFAULT_S3_FLUSH_INTERVAL_SECONDS": "2"}, config=config) as candidate,
candidate.scenario() as scenario,
):
model: Final = scenario.model(api_base=provider.url + "/v1", api_key="synthetic-provider-key")
key: Final = scenario.key(models=[model])
ids: Final = _push(candidate, model, key, marker, 4)
payloads: Final = collect_payloads(sink, 4)
assert frozenset(payload["id"] for payload in payloads) == ids
assert sink.attempts >= 7, (
f"only {sink.attempts} PUT attempts for 4 objects whose first 3 uploads 403 RequestTimeout; "
"coded 403s must be retried"
)
def test_s3_v2_success_callback_mode_logs_only_successes(gateway: Gateway, tmp_path: Path) -> None:
marker: Final = "s3succ" + uuid.uuid4().hex[:8]
sink: Final = RecordingS3Sink()
with wire_server(_chat_reply) as provider, wire_server(sink.respond) as bucket:
config: Final = _recording_s3_config(
tmp_path,
bucket.url,
{},
{"callbacks": [], "success_callback": ["s3_v2"]},
)
with (
owned_proxy(gateway, tmp_path, {"DEFAULT_S3_FLUSH_INTERVAL_SECONDS": "2"}, config=config) as candidate,
candidate.scenario() as scenario,
):
model: Final = scenario.model(api_base=provider.url + "/v1", api_key="synthetic-provider-key")
key: Final = scenario.key(models=[model])
ghost: Final = candidate.request(
"POST",
"/v1/chat/completions",
{"model": f"ghost-{uuid.uuid4().hex}", "messages": [{"role": "user", "content": "hi"}]},
key=key,
)
assert ghost.status_code in (400, 403, 404), ghost.text
_send(candidate, model, key, marker)
payloads: Final = collect_payloads(sink, 1)
assert len(payloads) == 1
assert payloads[0]["id"] == marker
assert payloads[0]["status"] == "success"
def test_s3_v2_failure_callback_mode_logs_only_failures(gateway: Gateway, tmp_path: Path) -> None:
marker: Final = "s3failcb" + uuid.uuid4().hex[:8]
sink: Final = RecordingS3Sink()
with wire_server(_chat_reply) as provider, wire_server(sink.respond) as bucket:
config: Final = _recording_s3_config(
tmp_path,
bucket.url,
{},
{"callbacks": [], "failure_callback": ["s3_v2"]},
)
with (
owned_proxy(gateway, tmp_path, {"DEFAULT_S3_FLUSH_INTERVAL_SECONDS": "2"}, config=config) as candidate,
candidate.scenario() as scenario,
):
model: Final = scenario.model(api_base=provider.url + "/v1", api_key="synthetic-provider-key")
key: Final = scenario.key(models=[model])
_send(candidate, model, key, marker)
ghost: Final = candidate.request(
"POST",
"/v1/chat/completions",
{"model": f"ghost-{uuid.uuid4().hex}", "messages": [{"role": "user", "content": "hi"}]},
key=key,
)
assert ghost.status_code in (400, 403, 404), ghost.text
payloads: Final = collect_payloads(sink, 1)
assert len(payloads) == 1
assert payloads[0]["status"] == "failure"
assert payloads[0]["id"] != marker
assert isinstance(payloads[0]["litellm_call_id"], str) and payloads[0]["litellm_call_id"]