feat(e2e): say why a response was not recorded

Build 226 routed Bedrock streaming for the first time and rejected 62 of
220 misses on that mount, and the counters could not say why. A flat
rejected count covers three unrelated things with opposite fixes: the
consumer walking away mid-capture, a body that arrived whole and failed
its endpoint's rule, and a provider that could not be reached. Each now
also counts its own reason.

A consumer that walks away was counting nothing at all. Abandoning the
capture generator raises GeneratorExit at its yield, so neither branch of
the old accounting ran and the miss simply vanished from the report, which
is also why misses could exceed writes plus rejected with nothing to
explain the gap. The decision moves into settle() so the generator's
finally owns the accounting and an abandoned capture is counted like any
other rejection.
This commit is contained in:
Yuneng Jiang 2026-09-16 07:47:31 -07:00
parent 8a553ceb58
commit 7006da9cde
No known key found for this signature in database
3 changed files with 78 additions and 16 deletions

View file

@ -516,6 +516,45 @@ def test_counters_attribute_every_outcome_to_its_mount(
assert counts["mount:anthropic:rejected"] == 1 and "mount:openai:rejected" not in counts
def test_a_rejection_says_whether_the_body_was_cut_short_or_simply_unfinished(
store: RedisResponseStore, provider: Provider,
) -> None:
"""One `rejected` count cannot tell a connection that dropped from a body the
provider finished sending and the rules turned down, and those have opposite
fixes: the first is the client going away mid-capture, the second is a grammar
the cache does not accept. A mount whose rejections are mostly one or the other
is a different problem, so the report has to be able to say which."""
upstream: Final = f"http://127.0.0.1:{provider.server_port}"
cut_short: Final = cache_edge(store)
provider.stream = True
provider.truncated = True
provider.response = b'data: {"choices":[{"index":0,"delta":{"content":"hi"},"finish_reason":"stop"}]}\n\ndata: [DONE]\n\n'
running: Final = start_provider_edge(cut_short, mounts={"openai": upstream})
try:
forward("POST", running.edge.api_base("openai") + "/v1/chat/completions",
headers=HEADERS, body=MARKED, timeout=5)
finally:
running.shutdown()
unfinished: Final = cache_edge(store)
provider.stream = False
provider.truncated = False
provider.response = b'{"choices":[{"index":0,"message":{"content":"hi"}}]}'
second: Final = start_provider_edge(unfinished, mounts={"openai": upstream})
try:
call(second.edge.api_base("openai") + "/v1/chat/completions", MARKED)
finally:
second.shutdown()
cut: Final = dict(cut_short.counters.counts)
turned_down: Final = dict(unfinished.counters.counts)
assert cut["mount:openai:rejected"] == 1 and turned_down["mount:openai:rejected"] == 1
assert cut["mount:openai:rejected_cut_short"] == 1
assert "mount:openai:rejected_incomplete" not in cut
assert turned_down["mount:openai:rejected_incomplete"] == 1
assert "mount:openai:rejected_cut_short" not in turned_down
EMBEDDING_SUCCESS: Final = (
b'{"object":"list","data":[{"object":"embedding","index":0,"embedding":[0.1,0.2]}],'
b'"model":"text-embedding-3-small","usage":{"prompt_tokens":2,"total_tokens":2}}'

View file

@ -54,7 +54,7 @@ The trusted runner receives:
- `E2E_PROVIDER_CACHE_NAMESPACE`: shared environment namespace, independent of build and candidate revision
- `E2E_PROVIDER_CACHE_METRICS_DIR`: optional per-process counter artifact directory
Do not give cache credentials to candidate deployments. Counter artifacts contain no recorded payloads or credentials. Hits count shared-cache responses; upstream attempts count actual forwards from the edge. Every counter is emitted twice, once as a flat total and once under `mount:{mount}:`, so a hit rate can be read per provider rather than only in aggregate. Existing application-cache observations still count requests arriving at the edge, including shared-cache hits
Do not give cache credentials to candidate deployments. Counter artifacts contain no recorded payloads or credentials. Hits count shared-cache responses; upstream attempts count actual forwards from the edge. A rejection also counts its reason, one of `rejected_cut_short` (the consumer walked away mid-capture), `rejected_incomplete` (the body arrived whole and failed its endpoint's rule) or `rejected_unreachable` (the provider could not be reached). A mount whose rejections are nearly all one or the other is a different problem, and the flat count cannot tell them apart. Every counter is emitted twice, once as a flat total and once under `mount:{mount}:`, so a hit rate can be read per provider rather than only in aggregate. Existing application-cache observations still count requests arriving at the edge, including shared-cache hits
Tests that require real provider timing, limits or state use `@pytest.mark.provider_live`. The marker keeps newly registered models on live routes without weakening their assertions. The provider prompt-caching tests carry it because a replayed priming response reports cache creation rather than a cache read.

View file

@ -52,6 +52,9 @@ BEDROCK_SUFFIXES: Final = (
BEDROCK_INVOKE_STREAM_SUFFIX,
)
EVENTSTREAM_PRELUDE_BYTES: Final = 4
CUT_SHORT: Final = "cut_short"
INCOMPLETE: Final = "incomplete"
UNREACHABLE: Final = "unreachable"
EVENT_TYPE_HEADER: Final = ":event-type"
EVENTSTREAM_HEADERS: Final[TypeAdapter[dict[str, str]]] = TypeAdapter(dict[str, str])
OPENAI_JSON_PATHS: Final = frozenset({"/v1/chat/completions", "/v1/messages", "/v1/embeddings", "/v1/responses"})
@ -551,7 +554,7 @@ class CacheEdge:
)
prepared: Final = prepare_forward(method, url, self.outbound(mount, method, url, headers, body), body)
if isinstance(prepared, NetworkError):
self.count(mount, "rejected")
self.reject(mount, UNREACHABLE)
return prepared
identity: Final = request_identity(
self.secret, test_key, method, url, self.keyed(mount, prepared.headers), body,
@ -575,7 +578,7 @@ class CacheEdge:
return head
if isinstance(head, NetworkError):
self.store.release(key, capture_slot)
self.count(mount, "rejected")
self.reject(mount, UNREACHABLE)
return head
return StreamHead(
head.status_code, head.headers, primed_steps(self.capture(mount, key, capture_slot, url, head)),
@ -585,25 +588,45 @@ class CacheEdge:
self, mount: str, key: str, lease: CaptureLease, url: str, head: StreamHead,
) -> Generator[StreamStep, None, None]:
capture: Final = ResponseCapture()
reason = CUT_SHORT # rebind-ok: a consumer that walks away never reaches the settle call below
try:
with closing(head.steps):
yield StreamChunk(b"")
for step in head.steps:
yield step
capture.observe(step)
chunks: Final = capture.chunks() if capture.eligible else ()
headers: Final = {
name: value for name, value in head.headers.items() if name.lower() not in UNRECORDED_RESPONSE_HEADERS
}
if not capture.eligible or not successful_response(mount, url, head.status_code, headers, b"".join(chunks)):
self.count(mount, "rejected")
return
response: Final = CachedResponse(
request_key=key, status_code=head.status_code, headers=headers,
chunks=tuple(base64.b64encode(chunk).decode("ascii") for chunk in chunks),
)
published: Final = self.store.publish(key, lease, encode_response(self.secret, response))
self.count(mount, "writes" if published else "write_failures")
reason = self.settle(mount, key, lease, url, head, capture)
finally:
self.reject(mount, reason)
self.store.release(key, lease)
capture.buffer.close()
def settle(
self, mount: str, key: str, lease: CaptureLease, url: str, head: StreamHead, capture: ResponseCapture,
) -> str | None:
"""None once the response is stored, otherwise the reason it was not."""
if not capture.eligible:
return CUT_SHORT
headers: Final = {
name: value for name, value in head.headers.items() if name.lower() not in UNRECORDED_RESPONSE_HEADERS
}
chunks: Final = capture.chunks()
if not successful_response(mount, url, head.status_code, headers, b"".join(chunks)):
return INCOMPLETE
response: Final = CachedResponse(
request_key=key, status_code=head.status_code, headers=headers,
chunks=tuple(base64.b64encode(chunk).decode("ascii") for chunk in chunks),
)
published: Final = self.store.publish(key, lease, encode_response(self.secret, response))
self.count(mount, "writes" if published else "write_failures")
return None
def reject(self, mount: str, reason: str | None) -> None:
"""A flat rejection count cannot separate a connection that went away from
a body the provider finished sending and the rules turned down, and the two
have opposite fixes. A mount whose rejections are nearly all one or the
other is a different problem, so the report has to be able to say which."""
if reason is None:
return
self.count(mount, "rejected")
self.count(mount, f"rejected_{reason}")