fix(e2e): clean up batch files reliably and expire Azure inputs

This commit is contained in:
Yuneng Jiang 2026-09-07 12:15:06 -07:00
parent 55fe4a7894
commit 6c21be6194
No known key found for this signature in database
9 changed files with 381 additions and 50 deletions

View file

@ -120,6 +120,24 @@ create traverse gateway -> gateway -> OpenAI (LIT-5347, PR #36240). The pin:
nested managed ids round-trip retrieve. This self-chaining only needs the proxy to
reach its own `PROXY_BASE_URL`, which holds both locally and on the e2e stage.
## Cleanup
Batch teardown cancels active batches before deleting their input files and keys.
Raw file IDs from both `model_param` and `provider_fallback` uploads use the upload
provider when deleted. Model-encoded and managed file IDs route themselves
File deletion and batch cancellation check their responses and retry transient
failures up to three times. Teardown attempts every registered cleanup before
reporting failures as test errors. Already deleted files and batches that are
terminal are safe to clean up again. Cancellation polls for up to ten minutes
before input deletion, because accepting cancellation does not finish it
Azure input uploads request `expires_after` anchored to `created_at` with
`seconds=1209600`, and the lifecycle tests check the returned expiry. This is a
fallback for interrupted runs: immediate deletion remains the normal cleanup.
Azure's minimum supported native expiry is 14 days, so a three-day expiry cannot
be requested through its Files API
## Terminal state + cost write-back (cross-run marker baton)
The 24h completion window rules out submit-and-wait inside one run, so

View file

@ -0,0 +1,90 @@
from collections.abc import Callable
from time import monotonic, sleep
from typing import Final, Protocol
from pydantic import BaseModel
from batch_client import BatchObject, FileDeleteResponse
from e2e_http import NetworkError, RateLimitedError, Result, Success, UnknownApiError
CLEANUP_DELAYS: Final = (1.0, 2.0, 4.0)
BATCH_TERMINAL_STATUSES: Final = frozenset({"completed", "failed", "expired", "cancelled"})
BATCH_CANCEL_TIMEOUT_SECONDS: Final = 600.0
BATCH_CANCEL_POLL_SECONDS: Final = 10.0
class BatchCleanupClient(Protocol):
def delete_file(self, file_id: str, *, key: str, provider: str | None = None) -> Result[FileDeleteResponse]: ...
def retrieve_batch(self, batch_id: str, *, key: str, provider: str | None = None) -> Result[BatchObject]: ...
def cancel_batch(self, batch_id: str, *, key: str, provider: str | None = None) -> Result[BatchObject]: ...
def cleanup_result[R: BaseModel](
action: Callable[[], Result[R]], *, wait: Callable[[float], None] = sleep
) -> Result[R]:
for delay, result in ((delay, action()) for delay in CLEANUP_DELAYS):
match result:
case NetworkError() | RateLimitedError():
wait(delay)
case UnknownApiError(status_code=code) if code in {408, 429, 500, 502, 503, 504}:
wait(delay)
case _:
return result
return action()
def _require_cleanup_success[R: BaseModel](result: Result[R], operation: str) -> R:
match result:
case Success(data=data):
return data
case UnknownApiError(status_code=code):
raise AssertionError(f"{operation} failed: HTTP {code}")
case _:
raise AssertionError(f"{operation} failed: {result.kind}")
def cleanup_file(client: BatchCleanupClient, file_id: str, *, key: str, provider: str | None = None) -> None:
result: Final = cleanup_result(lambda: client.delete_file(file_id, key=key, provider=provider))
if isinstance(result, UnknownApiError) and result.status_code == 404:
return
deleted: Final = _require_cleanup_success(result, f"Delete file {file_id}")
assert deleted.deleted, f"Delete file {file_id} did not confirm deletion"
def cleanup_batch(
client: BatchCleanupClient,
batch_id: str,
*,
key: str,
provider: str | None = None,
wait: Callable[[float], None] = sleep,
clock: Callable[[], float] = monotonic,
) -> None:
fetched: Final = _require_cleanup_success(
cleanup_result(lambda: client.retrieve_batch(batch_id, key=key, provider=provider)),
f"Retrieve batch {batch_id} for cleanup",
)
if fetched.status in BATCH_TERMINAL_STATUSES:
return
if fetched.status != "cancelling":
result: Final = cleanup_result(lambda: client.cancel_batch(batch_id, key=key, provider=provider))
if not (isinstance(result, UnknownApiError) and result.status_code in {400, 409}):
cancelled: Final = _require_cleanup_success(result, f"Cancel batch {batch_id}")
assert cancelled.status in BATCH_TERMINAL_STATUSES | {"cancelling"}, (
f"Cancel batch {batch_id} left status {cancelled.status}"
)
deadline: Final = clock() + BATCH_CANCEL_TIMEOUT_SECONDS
while True:
current = _require_cleanup_success(
cleanup_result(lambda: client.retrieve_batch(batch_id, key=key, provider=provider)),
f"Retrieve batch {batch_id} after cancellation",
)
if current.status in BATCH_TERMINAL_STATUSES:
return
assert current.status == "cancelling", f"Cancel batch {batch_id} left status {current.status}"
assert clock() < deadline, (
f"Batch {batch_id} cancellation did not finish within {BATCH_CANCEL_TIMEOUT_SECONDS}s"
)
wait(BATCH_CANCEL_POLL_SECONDS)

View file

@ -13,8 +13,9 @@ co-located here because only this suite uses them.
from __future__ import annotations
from dataclasses import dataclass
from typing import Final, Literal
from pydantic import BaseModel
from pydantic import BaseModel, Field
from proxy_client import ProxyClient
from e2e_http import (
@ -27,6 +28,18 @@ from e2e_http import (
from models import LiteLLMParamsBody
UPLOAD_FILENAME = "batch_input.jsonl"
AZURE_FILE_EXPIRY_SECONDS: Final = 14 * 24 * 60 * 60
class ExpiringFileUploadForm(FileUploadForm):
expires_after_anchor: Literal["created_at"] = Field(default="created_at", alias="expires_after[anchor]")
expires_after_seconds: int = Field(default=AZURE_FILE_EXPIRY_SECONDS, alias="expires_after[seconds]")
def batch_upload_form(provider: str, *, target_model_names: str | None = None) -> FileUploadForm:
if provider == "azure":
return ExpiringFileUploadForm(target_model_names=target_model_names)
return FileUploadForm(target_model_names=target_model_names)
class FileObject(BaseModel):
@ -37,6 +50,7 @@ class FileObject(BaseModel):
bytes: int | None = None
status: str | None = None
created_at: int | None = None
expires_at: int | None = None
class FileList(BaseModel):

View file

@ -108,6 +108,10 @@ class Capability:
def id(self) -> str:
return f"{self.provider}-{self.scenario}"
@property
def file_provider(self) -> str | None:
return self.provider if self.scenario in {"model_param", "provider_fallback"} else None
@property
def jsonl_model(self) -> str:
# Always the provider deployment name. Unified routes via

View file

@ -13,7 +13,7 @@ the proxy config.
from __future__ import annotations
import os
from typing import Iterator
from typing import Final, Iterator
import pytest
@ -21,6 +21,7 @@ from batch_client import BatchClient, build_client
from capabilities import PROVIDERS
from e2e_config import MANAGED_FILES_OPT_IN_ENV
from e2e_http import NoBody
from lifecycle import ResourceManager
from proxy_client import ProxyClient
@ -52,6 +53,13 @@ def client(proxy: ProxyClient) -> BatchClient:
return build_client(proxy)
@pytest.fixture
def resources(client: BatchClient) -> Iterator[ResourceManager]:
manager: Final = ResourceManager(client=client.proxy, strict_cleanup=True)
yield manager
manager.teardown()
@pytest.fixture(scope="session")
def batch_deployments(client: BatchClient) -> Iterator[None]:
probe = client.proxy.probe("/health/liveliness", params=NoBody())

View file

@ -0,0 +1,187 @@
from builtins import ExceptionGroup
from collections.abc import Iterator
from dataclasses import dataclass, field
from typing import Final
import pytest
from batch_cleanup import BATCH_CANCEL_TIMEOUT_SECONDS, CLEANUP_DELAYS, cleanup_batch, cleanup_file, cleanup_result
from batch_client import AZURE_FILE_EXPIRY_SECONDS, BatchObject, FileDeleteResponse, batch_upload_form
from capabilities import CAPABILITIES, Capability
from e2e_http import NetworkError, RateLimitedError, Result, Success, UnknownApiError
from lifecycle import ResourceManager
from models import KeyGenerateBody
@dataclass
class CleanupClient:
files: Iterator[Result[FileDeleteResponse]] = field(default_factory=lambda: iter(()))
batches: Iterator[Result[BatchObject]] = field(default_factory=lambda: iter(()))
cancellations: Iterator[Result[BatchObject]] = field(default_factory=lambda: iter(()))
calls: list[str] = field(default_factory=list)
def delete_file(self, file_id: str, *, key: str, provider: str | None = None) -> Result[FileDeleteResponse]:
self.calls.append(f"delete {provider} {file_id}")
return next(self.files)
def retrieve_batch(self, batch_id: str, *, key: str, provider: str | None = None) -> Result[BatchObject]:
self.calls.append(f"retrieve {provider} {batch_id}")
return next(self.batches)
def cancel_batch(self, batch_id: str, *, key: str, provider: str | None = None) -> Result[BatchObject]:
self.calls.append(f"cancel {provider} {batch_id}")
return next(self.cancellations)
def generate_key(self, body: KeyGenerateBody) -> str:
return "test-key"
def delete_key(self, key: str) -> None:
self.calls.append(f"delete key {key}")
def delete_customers(self, user_ids: list[str]) -> None:
self.calls.append(f"delete customers {user_ids}")
def batch(status: str) -> Success[BatchObject]:
return Success(status_code=200, data=BatchObject(id="batch-1", status=status))
def deleted_file(*, deleted: bool = True) -> Success[FileDeleteResponse]:
return Success(status_code=200, data=FileDeleteResponse(id="file-1", deleted=deleted))
class TestFileCleanup:
@pytest.mark.parametrize("cap", CAPABILITIES, ids=[cap.id for cap in CAPABILITIES])
def test_deletes_raw_files_through_the_upload_provider(self, cap: Capability) -> None:
client: Final = CleanupClient(files=iter((deleted_file(),)))
cleanup_file(client, "file-1", key="test-key", provider=cap.file_provider)
expected_provider: Final = cap.provider if cap.scenario in {"model_param", "provider_fallback"} else None
assert client.calls == [f"delete {expected_provider} file-1"]
def test_failed_delete_is_reported_after_remaining_resources_are_cleaned(self) -> None:
client: Final = CleanupClient(files=iter((UnknownApiError(status_code=403, body="secret response"),)))
manager: Final = ResourceManager(client=client, strict_cleanup=True)
key: Final = manager.key()
manager.defer(lambda: cleanup_file(client, "file-1", key=key, provider="azure"))
with pytest.raises(ExceptionGroup) as caught:
manager.teardown()
assert client.calls == ["delete azure file-1", "delete key test-key"]
assert len(caught.value.exceptions) == 1
assert str(caught.value.exceptions[0]) == "Delete file file-1 failed: HTTP 403"
def test_success_response_must_confirm_deletion(self) -> None:
client: Final = CleanupClient(files=iter((deleted_file(deleted=False),)))
with pytest.raises(AssertionError, match="did not confirm deletion"):
cleanup_file(client, "file-1", key="test-key")
def test_cleanup_is_idempotent_when_file_is_already_deleted(self) -> None:
client: Final = CleanupClient(files=iter((UnknownApiError(status_code=404, body="missing"),)))
cleanup_file(client, "file-1", key="test-key", provider="azure")
assert client.calls == ["delete azure file-1"]
def test_default_resource_cleanup_keeps_existing_best_effort_behavior(self) -> None:
client: Final = CleanupClient(files=iter((UnknownApiError(status_code=403, body="forbidden"),)))
manager: Final = ResourceManager(client=client)
key: Final = manager.key()
manager.defer(lambda: cleanup_file(client, "file-1", key=key))
manager.teardown()
assert client.calls == ["delete None file-1", "delete key test-key"]
class TestCleanupRetries:
@pytest.mark.parametrize(
"failure",
[NetworkError(message="offline"), RateLimitedError(), UnknownApiError(status_code=503, body="unavailable")],
)
def test_transient_error_retries_and_returns_success(self, failure: Result[FileDeleteResponse]) -> None:
outcomes: Final = iter((failure, deleted_file()))
delays: Final[list[float]] = []
result: Final[Result[FileDeleteResponse]] = cleanup_result(lambda: next(outcomes), wait=delays.append)
assert isinstance(result, Success) and result.data.deleted
assert delays == [1.0]
def test_persistent_error_has_bounded_retries(self) -> None:
failure: Final = UnknownApiError(status_code=503, body="unavailable")
outcomes: Final[Iterator[Result[FileDeleteResponse]]] = iter((failure,) * (len(CLEANUP_DELAYS) + 1))
delays: Final[list[float]] = []
result: Final[Result[FileDeleteResponse]] = cleanup_result(lambda: next(outcomes), wait=delays.append)
assert result is failure
assert tuple(delays) == CLEANUP_DELAYS
assert next(outcomes, None) is None
def test_permanent_error_is_not_retried(self) -> None:
failure: Final = UnknownApiError(status_code=403, body="forbidden")
outcomes: Final = iter((failure, deleted_file()))
delays: Final[list[float]] = []
assert cleanup_result(lambda: next(outcomes), wait=delays.append) is failure
assert delays == []
assert isinstance(next(outcomes), Success)
class TestBatchCancellation:
def test_cancelling_batch_is_polled_until_terminal_without_cancelling_again(self) -> None:
client: Final = CleanupClient(batches=iter((batch("cancelling"), batch("cancelling"), batch("cancelled"))))
delays: Final[list[float]] = []
cleanup_batch(client, "batch-1", key="test-key", wait=delays.append)
assert client.calls == ["retrieve None batch-1"] * 3
assert delays == [10.0]
def test_cancellation_timeout_is_reported_but_file_and_key_cleanup_still_run(self) -> None:
client: Final = CleanupClient(
batches=iter((batch("cancelling"), batch("cancelling"))), files=iter((deleted_file(),))
)
ticks: Final = iter((0.0, BATCH_CANCEL_TIMEOUT_SECONDS))
manager: Final = ResourceManager(client=client, strict_cleanup=True)
key: Final = manager.key()
manager.defer(lambda: cleanup_file(client, "file-1", key=key))
manager.defer(lambda: cleanup_batch(client, "batch-1", key=key, clock=lambda: next(ticks)))
with pytest.raises(ExceptionGroup) as caught:
manager.teardown()
assert "cancellation did not finish" in str(caught.value.exceptions[0])
assert client.calls == [
"retrieve None batch-1",
"retrieve None batch-1",
"delete None file-1",
"delete key test-key",
]
@pytest.mark.parametrize("status", ["completed", "failed", "expired", "cancelled"])
def test_inactive_batch_needs_no_cancellation(self, status: str) -> None:
client: Final = CleanupClient(batches=iter((batch(status),)))
cleanup_batch(client, "batch-1", key="test-key")
assert client.calls == ["retrieve None batch-1"]
def test_active_batch_is_cancelled_through_its_provider(self) -> None:
client: Final = CleanupClient(
batches=iter((batch("in_progress"), batch("cancelled"))), cancellations=iter((batch("cancelling"),))
)
cleanup_batch(client, "batch-1", key="test-key", provider="azure")
assert client.calls == ["retrieve azure batch-1", "cancel azure batch-1", "retrieve azure batch-1"]
@pytest.mark.parametrize("status", ["completed", "in_progress"])
def test_cancellation_conflict_is_accepted_only_when_batch_became_inactive(self, status: str) -> None:
client: Final = CleanupClient(
batches=iter((batch("in_progress"), batch(status))),
cancellations=iter((UnknownApiError(status_code=409, body="conflict"),)),
)
if status == "completed":
cleanup_batch(client, "batch-1", key="test-key")
else:
with pytest.raises(AssertionError, match="Cancel batch batch-1 left status in_progress"):
cleanup_batch(client, "batch-1", key="test-key")
assert client.calls == ["retrieve None batch-1", "cancel None batch-1", "retrieve None batch-1"]
class TestAzureFileExpiry:
def test_azure_form_serializes_native_expiry_for_the_proxy(self) -> None:
form: Final = batch_upload_form("azure", target_model_names="azure-test")
assert form.model_dump(by_alias=True, exclude_none=True) == {
"purpose": "batch",
"target_model_names": "azure-test",
"expires_after[anchor]": "created_at",
"expires_after[seconds]": AZURE_FILE_EXPIRY_SECONDS,
}
@pytest.mark.parametrize("provider", ["openai", "vertex_ai", "bedrock"])
def test_other_providers_keep_their_existing_upload_fields(self, provider: str) -> None:
assert batch_upload_form(provider).model_dump(by_alias=True, exclude_none=True) == {"purpose": "batch"}

View file

@ -21,14 +21,16 @@ import os
import re
import time
from datetime import datetime, timedelta, timezone
from typing import Callable
import pytest
from pydantic import BaseModel
from e2e_config import PROXY_BASE_URL, unique_marker
from batch_cleanup import cleanup_batch, cleanup_file
from batch_client import (
AZURE_FILE_EXPIRY_SECONDS,
batch_upload_form,
UPLOAD_FILENAME,
BatchClient,
BatchCreateBody,
@ -155,19 +157,19 @@ def upload_for_scenario(
if cap.scenario == "encoded":
return client.upload_file(
content=content,
form=FileUploadForm(purpose="batch"),
form=batch_upload_form(cap.provider),
model=cap.model,
key=key,
)
if cap.scenario == "unified":
return client.upload_file(
content=content,
form=FileUploadForm(purpose="batch", target_model_names=cap.model),
form=batch_upload_form(cap.provider, target_model_names=cap.model),
key=key,
)
return client.upload_file(
content=content,
form=FileUploadForm(purpose="batch"),
form=batch_upload_form(cap.provider),
key=key,
provider=cap.provider,
)
@ -188,20 +190,11 @@ def create_for_scenario(
def op_provider(cap: Capability) -> str | None:
"""provider_fallback ids are raw, so retrieve/cancel/list/delete need the provider
"""provider_fallback batch ids are raw, so retrieve/cancel/list need the provider
hint; the other scenarios encode it into the id and route automatically."""
return cap.provider if cap.scenario == "provider_fallback" else None
def quietly(action: Callable[[], object]) -> Callable[[], None]:
"""Adapt a value-returning call into a best-effort cleanup the teardown can run."""
def run() -> None:
action()
return run
def assert_file_object(file: FileObject, *, provider: str) -> None:
assert file.object == "file", f"file.object={file.object!r}"
assert file.purpose == "batch", f"file.purpose={file.purpose!r}"
@ -209,6 +202,10 @@ def assert_file_object(file: FileObject, *, provider: str) -> None:
if provider != "bedrock":
assert file.bytes > 0, f"file.bytes={file.bytes!r}"
assert file.status, "file.status missing"
if provider == "azure":
assert file.expires_at is not None, "Azure batch input has no automatic expiry"
assert file.created_at is not None
assert file.expires_at - file.created_at == AZURE_FILE_EXPIRY_SECONDS
assert (
file.created_at is not None and file.created_at > 0
), "file.created_at missing"
@ -249,7 +246,7 @@ def test_batch_lifecycle(
file = unwrap(upload_for_scenario(client, cap, render_jsonl(cap.jsonl_model), key))
resources.defer(
quietly(lambda: client.delete_file(file.id, key=key, provider=provider))
lambda: cleanup_file(client, file.id, key=key, provider=cap.file_provider)
)
assert_file_object(file, provider=cap.provider)
assert matches_id_shape(
@ -260,7 +257,7 @@ def test_batch_lifecycle(
require_successful_call(created)
batch = BatchObject.model_validate_json(created.body)
resources.defer(
quietly(lambda: client.cancel_batch(batch.id, key=key, provider=provider))
lambda: cleanup_batch(client, batch.id, key=key, provider=provider)
)
assert batch.id, f"create returned no batch id (body={created.body[:200]})"
@ -339,7 +336,7 @@ def test_batch_key_model_access_denied(
denied_upload = client.upload_file(
content=render_jsonl(AZURE_BATCH_MODEL),
form=FileUploadForm(purpose="batch"),
form=batch_upload_form("azure"),
model=AZURE_BATCH_MODEL,
key=key,
)
@ -356,7 +353,7 @@ def test_batch_key_model_access_denied(
)
).id
resources.defer(
quietly(lambda: client.delete_file(raw_file, key=key, provider="openai"))
lambda: cleanup_file(client, raw_file, key=key, provider="openai")
)
denied_create = client.create_batch(
@ -383,6 +380,7 @@ def test_file_upload_and_delete_outputs(
key=key,
)
)
resources.defer(lambda: cleanup_file(client, file.id, key=key))
assert_file_object(file, provider="openai")
deleted = unwrap(client.delete_file(file.id, key=key))
@ -458,12 +456,12 @@ def test_rate_limited_batch_create_leaves_no_unattributed_spend_row(
key=key,
)
)
resources.defer(quietly(lambda: client.delete_file(file.id, key=key)))
resources.defer(lambda: cleanup_file(client, file.id, key=key))
created = client.create_batch(body=BatchCreateBody(input_file_id=file.id), key=key)
require_successful_call(created)
batch = BatchObject.model_validate_json(created.body)
resources.defer(quietly(lambda: client.cancel_batch(batch.id, key=key)))
resources.defer(lambda: cleanup_batch(client, batch.id, key=key))
_ = client.proxy.poll_logs_for_key(key, min_rows=1)
@ -517,7 +515,7 @@ class TestBatchFileContent:
key=key,
)
)
resources.defer(quietly(lambda: client.delete_file(file.id, key=key)))
resources.defer(lambda: cleanup_file(client, file.id, key=key))
assert file.id
downloaded = client.proxy.transport.download(
@ -559,11 +557,11 @@ class TestBatchFileContent:
file = unwrap(
client.upload_file(
content=payload,
form=FileUploadForm(purpose="batch", target_model_names=provider.model),
form=batch_upload_form(provider.name, target_model_names=provider.model),
key=key,
)
)
resources.defer(quietly(lambda: client.delete_file(file.id, key=key)))
resources.defer(lambda: cleanup_file(client, file.id, key=key))
assert_file_object(file, provider=provider.name)
assert is_managed_id(file.id), (
f"{provider.name}: unified upload must return a managed file id, got {file.id!r}"
@ -626,7 +624,7 @@ class TestOpenAIFiles:
)
)
resources.defer(
quietly(lambda: client.delete_file(file.id, key=key, provider="openai"))
lambda: cleanup_file(client, file.id, key=key, provider="openai")
)
listed = unwrap(client.list_files(key=key))
@ -690,7 +688,7 @@ class TestOpenAIFiles:
key=key,
)
)
resources.defer(quietly(lambda: client.delete_file(file.id, key=key)))
resources.defer(lambda: cleanup_file(client, file.id, key=key))
fetched = unwrap(client.retrieve_file(file.id, key=key))
assert fetched.id == file.id, "retrieve must echo the uploaded file id"
@ -760,7 +758,7 @@ class TestBatchRateLimitErrorMapping:
key=key,
)
)
resources.defer(quietly(lambda: client.delete_file(file.id, key=key)))
resources.defer(lambda: cleanup_file(client, file.id, key=key))
created = client.create_batch(body=BatchCreateBody(input_file_id=file.id), key=key)
@ -813,7 +811,7 @@ class TestBatchEnqueuedTokenLimit:
key=key,
)
)
resources.defer(quietly(lambda: client.delete_file(file.id, key=key)))
resources.defer(lambda: cleanup_file(client, file.id, key=key))
return file
def _generate_enqueued_key(
@ -861,7 +859,7 @@ class TestBatchEnqueuedTokenLimit:
)
require_successful_call(created)
batch = BatchObject.model_validate_json(created.body)
resources.defer(quietly(lambda: client.cancel_batch(batch.id, key=key)))
resources.defer(lambda: cleanup_batch(client, batch.id, key=key))
@pytest.mark.covers(
"quota_management.ratelimit.batch_enqueued_tokens.blocks_when_exhausted",
@ -904,7 +902,7 @@ class TestBatchEnqueuedTokenLimit:
first = client.create_batch(body=BatchCreateBody(input_file_id=file.id), key=key)
require_successful_call(first)
first_batch = BatchObject.model_validate_json(first.body)
resources.defer(quietly(lambda: client.cancel_batch(first_batch.id, key=key)))
resources.defer(lambda: cleanup_batch(client, first_batch.id, key=key))
blocked = client.create_batch(body=BatchCreateBody(input_file_id=file.id), key=key)
assert blocked.status_code == 429, (
@ -928,7 +926,7 @@ class TestBatchEnqueuedTokenLimit:
)
require_successful_call(retried)
retry_batch = BatchObject.model_validate_json(retried.body)
resources.defer(quietly(lambda: client.cancel_batch(retry_batch.id, key=key)))
resources.defer(lambda: cleanup_batch(client, retry_batch.id, key=key))
ASSUME_ROLE_RAW_MODEL = "bedrock/us.anthropic.claude-haiku-4-5-20251001-v1:0"
@ -984,13 +982,13 @@ class TestBedrockBatchAssumeRole:
key=key,
)
)
resources.defer(quietly(lambda: client.delete_file(file.id, key=key)))
resources.defer(lambda: cleanup_file(client, file.id, key=key))
assert_file_object(file, provider="bedrock")
created = client.create_batch(body=BatchCreateBody(input_file_id=file.id), key=key)
require_successful_call(created)
batch = BatchObject.model_validate_json(created.body)
resources.defer(quietly(lambda: client.cancel_batch(batch.id, key=key)))
resources.defer(lambda: cleanup_batch(client, batch.id, key=key))
assert batch.id, f"assume-role create returned no batch id: {created.body[:200]}"
assert is_managed_id(batch.id), (
@ -1044,7 +1042,7 @@ class TestGeminiFiles:
key=key,
)
)
resources.defer(quietly(lambda: client.delete_file(file.id, key=key)))
resources.defer(lambda: cleanup_file(client, file.id, key=key))
assert_file_object(file, provider="gemini")
assert file.id, "gemini file upload returned no id"
@ -1099,13 +1097,13 @@ class TestHostedVllmBatch:
key=key,
)
)
resources.defer(quietly(lambda: client.delete_file(file.id, key=key)))
resources.defer(lambda: cleanup_file(client, file.id, key=key))
assert_file_object(file, provider="hosted_vllm")
created = client.create_batch(body=BatchCreateBody(input_file_id=file.id), key=key)
require_successful_call(created)
batch = BatchObject.model_validate_json(created.body)
resources.defer(quietly(lambda: client.cancel_batch(batch.id, key=key)))
resources.defer(lambda: cleanup_batch(client, batch.id, key=key))
assert batch.id, f"hosted_vllm create returned no batch id: {created.body[:200]}"
assert batch.status in CREATED_BATCH_STATUSES, (
@ -1192,7 +1190,7 @@ class TestBatchFailurePaths:
key=key,
)
)
resources.defer(quietly(lambda: client.delete_file(file.id, key=key)))
resources.defer(lambda: cleanup_file(client, file.id, key=key))
created = client.create_batch(body=BatchCreateBody(input_file_id=file.id), key=key)
require_successful_call(created)
@ -1243,12 +1241,12 @@ class TestBatchFailurePaths:
file = unwrap(
client.upload_file(
content=render_jsonl(AZURE_BATCH_RAW_MODEL),
form=FileUploadForm(purpose="batch"),
form=batch_upload_form("azure"),
model=AZURE_BATCH_MODEL,
key=key,
)
)
resources.defer(quietly(lambda: client.delete_file(file.id, key=key)))
resources.defer(lambda: cleanup_file(client, file.id, key=key))
assert decoded_model_from_id(file.id) == AZURE_BATCH_MODEL, (
f"upload did not encode the azure deployment into the file id: {file.id!r}"
)
@ -1258,7 +1256,7 @@ class TestBatchFailurePaths:
)
require_successful_call(created)
batch = BatchObject.model_validate_json(created.body)
resources.defer(quietly(lambda: client.cancel_batch(batch.id, key=key)))
resources.defer(lambda: cleanup_batch(client, batch.id, key=key))
assert decoded_model_from_id(batch.id) == AZURE_BATCH_MODEL, (
"create with a foreign encoded file id must route by the file's embedded model, "
@ -1307,7 +1305,7 @@ class TestBatchSecondHop:
key=key,
)
)
resources.defer(quietly(lambda: client.delete_file(file.id, key=key)))
resources.defer(lambda: cleanup_file(client, file.id, key=key))
assert is_managed_id(file.id), (
f"second-hop unified upload must return a managed file id, got {file.id!r}"
)
@ -1315,7 +1313,7 @@ class TestBatchSecondHop:
created = client.create_batch(body=BatchCreateBody(input_file_id=file.id), key=key)
require_successful_call(created)
batch = BatchObject.model_validate_json(created.body)
resources.defer(quietly(lambda: client.cancel_batch(batch.id, key=key)))
resources.defer(lambda: cleanup_batch(client, batch.id, key=key))
assert is_managed_id(batch.id), (
f"second-hop create must return a managed batch id, got {batch.id!r}"

View file

@ -21,6 +21,7 @@ from typing import Iterator
import pytest
from batch_client import BatchClient, FileObject
from batch_cleanup import cleanup_file
from capabilities import batch_model_name, is_managed_id, openai_batch_params
from e2e_config import unique_marker
from e2e_http import FileUploadForm, Result, UnknownApiError, unwrap
@ -108,7 +109,7 @@ def test_cross_user_managed_id_denied_owner_allowed(
key=owner_key,
)
)
resources.defer(lambda: client.delete_file(uploaded.id, key=owner_key))
resources.defer(lambda: cleanup_file(client, uploaded.id, key=owner_key))
assert is_managed_id(uploaded.id), f"expected a managed unified file id, got {uploaded.id}"
denied = client.retrieve_file(uploaded.id, key=other_key)

View file

@ -8,8 +8,9 @@ ResourceManager; the test registers a cleanup for every resource it creates, and
the fixture's teardown releases them all even when the test body raises.
"""
from builtins import ExceptionGroup
from dataclasses import dataclass, field
from typing import Callable, List, Protocol, runtime_checkable
from typing import Callable, Final, List, Protocol, runtime_checkable
from proxy_client import ProxyClient
from models import KeyGenerateBody
@ -52,6 +53,7 @@ class ResourceManager:
"""
client: ResourceClient
strict_cleanup: bool = False
_cleanups: List[Callable[[], object]] = field(
default_factory=list
) # mutable-ok: append-only teardown registry
@ -82,8 +84,17 @@ class ResourceManager:
return customer_id
def teardown(self) -> None:
for cleanup in reversed(self._cleanups):
try:
cleanup()
except Exception:
pass # best-effort: a failed cleanup must not block the rest
failures: Final = tuple(
failure for cleanup in reversed(self._cleanups)
if (failure := _run_cleanup(cleanup)) is not None
)
if failures and self.strict_cleanup:
raise ExceptionGroup("Resource cleanup failed", failures)
def _run_cleanup(cleanup: Callable[[], object]) -> Exception | None:
try:
cleanup()
except Exception as exc:
return exc
return None