feat(call): spool large metered answers and settle from their usage

A metered answer is buffered whole before headers so settlement sees complete evidence, and
anything over 8 MiB fails uncharged. Providers that inline generated media in JSON (Gemini
returns a base64 image plus a large thoughtSignature: ~9 MB at 2K, ~23 MB at 4K) can never fit.
Streaming instead would settle after the response starts, losing the exact X-Treg-Cost-Micro
and adding a second close-once path for disconnects, routed children and overflow; raising the
buffer would hold tens of megabytes per call in RAM.

An endpoint now declares `spooled_response: true`. Its metered 2xx is written to an unlinked
temp file under `spool_max_bytes` (64 MiB) and a per-process `spool_budget_bytes` (512 MiB),
claimed whole up front when the provider declares a Content-Length;
crossing either is the existing uncharged `response_buffer_limit`. The file is parsed once with
stdlib json in a worker thread (a 23 MB answer: ~30 ms, ~45 MB peak) behind
`spool_parse_concurrency`, and only the top-level keys its usage paths start at
(`settlement.usage_roots`) become the body every settlement consumer reads, so the evidence
cannot drift from the price. Settlement runs before headers as before; the router then relays
the file byte for byte, and its close (or garbage collection) returns the budget once. Spooled
bodies are not archived or kept for idempotent replay. Errors, own-key calls and routed
children keep their existing paths. `_buffer_response` and the spool share the refusal text and
the content-length rewrite.

The validator requires settle: usage and refuses the field beside async, resource_ownership or
managed_resource. AGENTS.md non-negotiable 4 records the new path.

Fragments: architecture/proxy-model.md, money.md, archive.md, catalog.md, interface/api.md,
ops/deploy.md.
This commit is contained in:
SToneX
2026-09-29 23:27:00 +08:00
parent 15cd572533
commit 2e3097cb1c
18 changed files with 514 additions and 23 deletions
+3
View File
@@ -48,6 +48,9 @@ Everything else in this file is guidance; these are the contract, and they win o
omits credential injection but does not strip or rewrite caller headers.
Routed endpoints and overflow wrap the child's answer and say so; they never alter it. Responses needing settlement or ownership evidence are buffered by the application
up to 8 MiB; exceeding that limit fails without charging, never returns a successful prefix.
An endpoint declaring `spooled_response` (inline media) instead reads its metered 2xx to an
unlinked temp file under a per-body cap and per-process budget, settles from the reported
usage it reads, then relays the file byte for byte; it is never archived or replayed.
Authorized free final fetches needing no body evidence stream in full.
5. Balances change only through money's five entries: grant, topup, reserve, settle, release.
There is deliberately no refund or adjustment entry; an ops correction is a grant.
+3 -1
View File
@@ -380,7 +380,9 @@ No `.env` is needed for local dev — every setting has a working default (ephem
decrypt its secret(s) → apply each binding's injector → stream to the upstream → fire-and-forget
audit record. The infra relay streams bytes without business logic. The call application buffers
responses needing settlement or ownership evidence up to 8 MiB; larger responses return a 502
without charging instead of a truncated success. Authorized free final downloads needing no body
without charging instead of a truncated success. Endpoints that inline media (Gemini images)
declare `spooled_response`: their metered answer is read to a temp file instead, settled from its
reported usage, and relayed whole. Authorized free final downloads needing no body
evidence stream in full, as do own-key and own-tool responses.
**Module map** (`src/treg/`):
+2
View File
@@ -586,6 +586,8 @@ Successful free final fetches that qualify for `MarketplaceCall.streamable_free_
lookup and recording even though their zero-amount money lifecycle remains metered. They have no
buffered body, so no empty or partial body/hash is recorded. Other calls exceeding the settlement
buffer's 8 MiB limit fail before recording and cannot populate a cache or idempotent success.
A spooled answer (`spooled_response`, inline media) is settled from evidence on disk and is never
recorded: its body is not in memory, and multi-megabyte generated media is not a reusable answer.
## The cache key
+7
View File
@@ -1030,6 +1030,13 @@ ceiling. A finite nonnegative response value settles the call at that amount; mi
non-finite evidence falls back to the normal estimate/miss rules. `reported_charge` is generic
catalog metadata, not a provider-specific billing branch, and cannot be combined with `cost.settle`.
`spooled_response: true` marks a synchronous endpoint whose answer inlines media too large for
the 8 MiB settlement buffer (Gemini returns images as base64 in its JSON: ~9 MB at 2K, ~23 MB at
4K). Its metered 2xx is read to disk and settled from the top-level keys its usage paths start at
(`settlement.usage_roots`; proxy-model.md), so the evidence cannot drift from the price. The
validator requires `settle: usage` and refuses the field beside `async`, `resource_ownership` or
`managed_resource`, which need the whole body.
`platform_request` fixes exact body, header or query values needed only on the shared credential.
A `queryParams.*` pin must appear exactly once and is read as the pinned value's type, so a run
option such as a spend cap or memory size can bound what one call costs. An Apify `per_result` price may add
+7
View File
@@ -1019,6 +1019,13 @@ reserve/settle gates and settles with an explicit zero override before returning
does not observe the original generation task or persist a response for idempotent replay; the
label is released and retrying performs another free read. MIME type never decides billability.
A metered 2xx on an endpoint declaring `spooled_response` keeps the same completeness rule with
the body on disk: `_spool_response` settles from the top-level keys its usage paths read before headers go
out, so `X-Treg-Cost-Micro` stays exact and the hold closes once on the normal path. An answer the
provider completed but whose evidence is missing or unparseable settles at the reserve (the image
exists); an oversized or budget-refused body is a `response_buffer_limit` release. The label is
released rather than stored, so an idempotent retry is a new, separately billed generation.
## HarvestAPI integration
+33
View File
@@ -591,6 +591,39 @@ metered JSON. The fault is attributed to treg's buffer limit, not to the provide
remain outside this limit. `tests/test_call_response_limits.py` exercises both real HTTP hops,
CLI output, boundaries, Range, disconnects, settlement evidence, archive and replay behavior.
### Spooled evidence for inline media
Some providers return generated media inside their JSON answer: Gemini's `generateContent` puts
a base64 image and a multi-megabyte `thoughtSignature` in the body, about 9 MB at 2K and 23 MB at
4K, with its token meters (`usageMetadata`) after them. Streaming that answer would mean settling
after the response is sent (no exact `X-Treg-Cost-Micro`, a second close-once path for disconnects,
routed children and overflow rebuilt); raising the buffer would put tens of megabytes per call in
a web process. The endpoint instead declares `spooled_response: true` and
`_read_evidence` hands its metered 2xx to `_spool_response`:
- The body streams into `tempfile.TemporaryFile` (unlinked from creation, so a crashed worker
leaves nothing on disk) under `spool_max_bytes` (64 MiB) and a per-process
`spool_budget_bytes` (512 MiB) shared by concurrent spools. A declared `Content-Length` is
claimed whole before the first read, so concurrent answers are admitted or refused whole rather
than all stalling half-read; without one, the claim grows chunk by chunk. Crossing either limit
raises the same `response_buffer_limit` before headers (the budget refusal says it is temporary);
the upstream closes in `finally` as in `_buffer_response`, temp-file creation included.
- The file is parsed once with the stdlib `json` (measured on a 23 MB Gemini answer: ~30 ms and
~45 MB peak; about three times the body at worst) in a worker thread, at most
`spool_parse_concurrency` (2) per process. Only the top-level keys the endpoint's usage paths
start at (`settlement.usage_roots`) survive, re-serialized as the `body` every later consumer sees: usage
settlement, result classification and capacity signatures. A body that is not a JSON object
yields empty evidence; a `usage` basis then settles at its reserve.
- Settlement then runs exactly as for a buffered body, before the response starts, and the router
relays the file in 256 KiB reads with the provider's headers and a recomputed `content-length`.
The replay's close returns the budget; `weakref.finalize` does so too for a response dropped
without closing.
- A spooled body is not archived and not stored for idempotent replay; its audit row records the
full size. Non-2xx answers, own-key calls and routed children (whose parent reads the child's
body) keep their existing paths. `tests/test_call_spool.py` covers the lifecycle (oversize,
budget, reset, cancellation, garbage collection, parse gate) and `test_call_response_limits.py`
a 20 MB Gemini-shaped answer over two real HTTP hops.
## HarvestAPI integration
+6 -1
View File
@@ -208,7 +208,12 @@ response. The response remains `Cache-Control: no-store` and precedes any login
`response_buffer_limit` is a treg-attributed 502 with a structured `detail.error` of the same
name. It means response evidence exceeded the 8 MiB settlement buffer before delivery; the new
call is not charged and its hold/idempotency claim is released. Authorized free final GET fetches
call is not charged and its hold/idempotency claim is released. On an endpoint declaring
`spooled_response` the same error means the body passed `spool_max_bytes` (64 MiB) or the
process's concurrent spool budget; the latter is temporary (its message says so) and clears as other
calls finish, so unlike an oversized answer it is worth retrying. A successful spooled
answer carries the usual exact `X-Treg-Cost-Micro`, but it is not stored for idempotent replay: a
retry with the same key calls the provider again and is charged again. Authorized free final GET fetches
needing no body evidence stream without that limit and return zero cost; their retries read the
provider again. See `proxy-model.md` for the eligibility and close-once lifecycle.
+8 -1
View File
@@ -311,7 +311,14 @@ whole record.
Archive recording has independent task-count and byte limits. `_MAX_PENDING_BYTES` bounds bodies
retained for database writes, and `archive_r2_max_pending_bytes` separately bounds object-store work.
These are admission budgets, not an RSS ceiling: SDK buffers, compression and mandatory terminal
evidence require additional memory headroom. Production observations and incident history live in
evidence require additional memory headroom.
Spooled settlement evidence (`spooled_response`, proxy-model.md) writes large metered answers to
the process temp directory, not RAM. `TREG_SPOOL_BUDGET_BYTES` (default 512 MiB per process) bounds
that disk use, `TREG_SPOOL_MAX_BYTES` (64 MiB) one answer, `TREG_SPOOL_PARSE_CONCURRENCY` (2) the
transient parse memory (up to about three times the body each), and `TREG_SPOOL_DIR` overrides the directory.
Size the budget against the host's ephemeral disk: a platform that evicts an instance over its
local-storage allowance turns an oversized budget into a restart. Production observations and incident history live in
the private [database-capacity runbook](https://github.com/superdesigndev/treg-internal/blob/main/docs/production/database-capacity.md).
## Hosted feature switches
+20 -1
View File
@@ -443,8 +443,8 @@ def check_usage_block(cost: dict, settle: object, where: str, errors: list[str],
"""`settle: usage` names the dotted path and unit of the provider's own reported charge."""
usage = cost.get("usage")
if settle == "usage":
terms = usage.get("terms") if isinstance(usage, dict) else None
if isinstance(usage, dict) and set(usage) == {"terms", "unit"}:
terms = usage["terms"]
if not isinstance(terms, list) or not terms or not all(
isinstance(term, dict) and set(term) == {"path", "rate"}
and isinstance(term.get("path"), str) and USAGE_TERM_PATH.fullmatch(term["path"])
@@ -605,6 +605,23 @@ def check_async_descriptor(descriptor: object, where: str, provider: str,
fail(errors, where, "an endpoint with async must have cost.type per_success")
def check_spooled_response(ep: dict, effective_async: object, where: str,
errors: list[str]) -> None:
"""`spooled_response: true` - a synchronous metered answer too large to buffer (inline media),
read to disk and settled from the top-level keys its usage paths start at. So only
`settle: usage`, and nothing else may need the body: no async task, resource ownership or
managed resource."""
if ep["spooled_response"] is not True:
fail(errors, where, "spooled_response must be true when present")
return
if effective_async is not None or ep.get("resource_ownership") or ep.get("managed_resource"):
fail(errors, where, "spooled_response cannot be combined with async, resource_ownership "
"or managed_resource: those read the whole body")
if (ep.get("cost") or {}).get("settle") != "usage":
fail(errors, where, "spooled_response settles from the answer's reported usage: it needs "
"settle: usage")
def check_resource_ownership(rule: object, where: str, input_schema: object,
errors: list[str]) -> None:
"""Validate declarative ownership of opaque ids on treg's shared provider account."""
@@ -1228,6 +1245,8 @@ def main(argv: list[str]) -> int:
check_resource_ownership(ep["resource_ownership"], where, inp, errors)
if ep.get("managed_resource") is not None:
check_managed_resource(ep["managed_resource"], where, inp, errors)
if "spooled_response" in ep:
check_spooled_response(ep, effective_async, where, errors)
if ep.get("verified"):
ex = ep.get("example_response")
if not ex:
+5
View File
@@ -340,6 +340,9 @@ class MarketplaceCall:
resource_ownership: dict | None = None
managed_resource: dict | None = None
public_resource_ids: tuple[str, ...] = ()
# On an endpoint declaring `spooled_response`: the top-level response keys its usage
# settlement reads. A metered 2xx goes to disk and only these keys are kept as evidence.
spooled_evidence: tuple[str, ...] = ()
# A platform-key utility poll was authorized against this org-owned submission. The buffered
# response may teach the same row its provider result/file id before the background worker runs.
async_owner_call_id: str | None = None
@@ -2069,6 +2072,8 @@ async def _resolve_marketplace_call(
settlement_basis=basis, request_data=request_data,
async_descriptor=ep.get("async"), resource_ownership=ep.get("resource_ownership"),
managed_resource=ep.get("managed_resource"),
spooled_evidence=(settlement_basis.usage_roots(raw_cost)
if ep.get("spooled_response") else ()),
)
if chosen_tool is not None:
return MarketplaceCall(tool=chosen_tool, tier="tool", **common)
+17 -6
View File
@@ -62,7 +62,6 @@ from . import overflow as overflow_cycle
from . import route as routed
from .settle import _dig
from .settle import (
_buffer_response,
_finish_cancelled_call as finish_cancelled_call,
_note_capacity_recovery,
_note_capacity_signal,
@@ -70,6 +69,7 @@ from .settle import (
_platform_settle,
_read_whole_if_small,
_record_first_call,
_read_evidence,
)
from .types import (
AuthorizationFailed,
@@ -596,6 +596,9 @@ async def _execute_call(request: _ApplicationRequest, upstream_client: httpx.Asy
drop_params: set[str] = set()
streaming_free_result = False
# Size of a metered 2xx relayed from disk (`spooled_response`); None when the body was read
# into memory. A spooled body is settled from its usage evidence and never retained.
spooled_bytes: int | None = None
served_hit = False # a cached hit — set where the archive answers instead of the vendor
served_repeat = False # …and this team had already paid for the question: the repeat price
# The archive identities of this call's answer (question key + exact bytes), set where the
@@ -1230,8 +1233,10 @@ async def _execute_call(request: _ApplicationRequest, upstream_client: httpx.Asy
if (served is None and mk is not None and (mk.metered or mk.free_owned_poll)
and not streaming_free_result):
# Settlement reads the body; owned free polls also need it to learn result ownership.
# A failure while draining remains an upstream failure on either path.
response, body = await _buffer_response(response)
# A failure while draining remains an upstream failure on either path. An endpoint
# that inlines media declares `spooled_response`: its 2xx goes to disk and `body`
# is only the usage evidence from here on.
response, body, spooled_bytes = await _read_evidence(mk, response)
if (platform_tier and response.status == 429 and _idempotent_read(request)
and (retry_s := _burst_retry_after(mk.provider, response, body)) is not None):
# Half two: ONE bounded wait on the provider's own `retry-after`, then the identical
@@ -1242,7 +1247,7 @@ async def _execute_call(request: _ApplicationRequest, upstream_client: httpx.Asy
response = await relay(
upstream_request, upstream_url, tool, secrets, upstream_client,
drop_params=drop_params or None, force_identity=True)
response, body = await _buffer_response(response)
response, body, spooled_bytes = await _read_evidence(mk, response)
smoothed.append("retry=1")
if mk.async_owner_call_id and 200 <= response.status < 300:
try:
@@ -1274,6 +1279,7 @@ async def _execute_call(request: _ApplicationRequest, upstream_client: httpx.Asy
# TREG_ARCHIVE_MODE says otherwise; record() is fire-and-forget and never raises.
# `own_credential` here means billed OAuth: the org's token, treg's bill.
if (mk.metered and archive.recording() and 200 <= response.status < 300
and spooled_bytes is None
and not (own_credential and _echoes_own_credential(tool, secrets, body))
and not _account_out_2xx(mk, response, body)):
_ct = next((v.decode("latin-1") for k, v in response.raw_headers
@@ -1500,7 +1506,9 @@ async def _execute_call(request: _ApplicationRequest, upstream_client: httpx.Asy
pending = _audit(response.status, observed_micro=observed,
charged_micro=None if deferred else charged,
duration_ms=duration_ms,
response_bytes=None if streaming_free_result else len(body), hit=result.hit,
response_bytes=(None if streaming_free_result else spooled_bytes
if spooled_bytes is not None else len(body)),
hit=result.hit,
capacity_signal=capacity_signal, error_request=err_request, error_response=err_response,
defer_analytics=may_overflow)
served_via = ""
@@ -1547,7 +1555,10 @@ async def _execute_call(request: _ApplicationRequest, upstream_client: httpx.Asy
try:
await _store_idempotent(idem_key, caller, status_code=response.status, body=body,
media_type=_response_header(response, "content-type"),
charged_micro=charged, metered=not streaming_free_result,
charged_micro=charged,
# A streamed or spooled body is not kept for replay: a
# retry calls the provider again.
metered=not streaming_free_result and spooled_bytes is None,
call_ref=call_ref)
except asyncio.CancelledError:
await _finish_cancelled_call(request, mk, call_ref, response)
+150 -13
View File
@@ -6,6 +6,8 @@ import asyncio
import json
import logging
import math
import tempfile
import weakref
from decimal import Decimal, InvalidOperation, ROUND_HALF_UP
from collections.abc import Callable
from dataclasses import dataclass
@@ -786,6 +788,24 @@ def _dig(doc, dotted: str):
return json_path(doc, dotted)
def _buffer_limit(message: str) -> GatewayFailed:
"""The uncharged refusal for evidence treg will not hold: raised before any header is sent."""
return GatewayFailed("response_buffer_limit", status_code=502, detail={
"error": "response_buffer_limit",
"message": f"{message}; no response was delivered and this call was not charged"})
def _with_length(raw_headers, size: int) -> tuple[tuple[bytes, bytes], ...]:
"""The upstream's headers verbatim (the relay already dropped hop-by-hop + our own), with a
content-length that matches what we are actually about to send."""
return tuple([(k, v) for k, v in raw_headers if k.lower() != b"content-length"]
+ [(b"content-length", str(size).encode())])
def _as_bytes(chunk) -> bytes:
return chunk if isinstance(chunk, bytes) else str(chunk).encode("utf-8", "replace")
async def _buffer_response(response: UpstreamResponse) -> tuple[UpstreamResponse, bytes]:
"""Read complete settlement evidence before sending headers; never return a prefix.
@@ -795,14 +815,10 @@ async def _buffer_response(response: UpstreamResponse) -> tuple[UpstreamResponse
chunks, size = [], 0
try:
async for chunk in response.body_stream:
raw = chunk if isinstance(chunk, bytes) else str(chunk).encode("utf-8", "replace")
raw = _as_bytes(chunk)
size += len(raw)
if size > _PLATFORM_BODY_MAX:
raise GatewayFailed(
"response_buffer_limit", status_code=502,
detail={"error": "response_buffer_limit",
"message": "upstream response exceeds treg's 8 MiB settlement buffer; "
"no response was delivered and this call was not charged"})
raise _buffer_limit("upstream response exceeds treg's 8 MiB settlement buffer")
chunks.append(raw)
body = b"".join(chunks)
finally:
@@ -814,16 +830,137 @@ async def _buffer_response(response: UpstreamResponse) -> tuple[UpstreamResponse
async def closed() -> None:
return None
# Carry the upstream's headers verbatim (the relay already dropped hop-by-hop + our own), with a
# content-length that matches what we are actually about to send.
raw_headers = tuple(
[(k, v) for k, v in response.raw_headers if k.lower() != b"content-length"]
+ [(b"content-length", str(len(body)).encode())]
)
out = UpstreamResponse(response.status, raw_headers, buffered_body(), closed)
out = UpstreamResponse(response.status, _with_length(response.raw_headers, len(body)),
buffered_body(), closed)
return out, body
_SPOOL_CHUNK = 256 * 1024
_spool_in_use = 0 # bytes held by this process's live spools, against `spool_budget_bytes`
_spool_gate: asyncio.Semaphore | None = None
_spool_gate_loop: asyncio.AbstractEventLoop | None = None
def _spool_parse_gate() -> asyncio.Semaphore:
"""`spool_parse_concurrency` parses per process, on the CURRENT loop (recreated if the loop
changed, like `audit._get_sem`)."""
global _spool_gate, _spool_gate_loop
loop = asyncio.get_running_loop()
if _spool_gate is None or _spool_gate_loop is not loop:
_spool_gate = asyncio.Semaphore(get_settings().spool_parse_concurrency)
_spool_gate_loop = loop
return _spool_gate
def _release_spool(file, claimed: int) -> None:
global _spool_in_use
_spool_in_use -= claimed
if file is not None:
file.close()
def _spool_evidence(file, keys: tuple[str, ...]) -> bytes:
"""The named top-level keys of a spooled JSON object, re-serialized; b"" when the body is not
a JSON object. Parsing the whole document is deliberate: a correct parser must scan every byte
anyway, and the stdlib one does it at C speed (a 23 MB answer in ~30 ms). Peak memory is about
three times the body (bytes, decoded text, parsed strings), bounded by `spool_max_bytes` and
`spool_parse_concurrency`."""
file.seek(0)
try:
document = json.load(file)
except (ValueError, RecursionError):
return b""
if not isinstance(document, dict):
return b""
return json.dumps({key: document[key] for key in keys if key in document}).encode()
async def _spool_response(
response: UpstreamResponse, keys: tuple[str, ...],
) -> tuple[UpstreamResponse, bytes, int]:
"""Read a large metered 2xx to disk, settle from its evidence keys, relay it verbatim.
The `_buffer_response` contract - complete evidence before headers, never a prefix, an
oversized body fails uncharged - with the body in an anonymous temp file instead of RAM, so a
provider that inlines media in its JSON (Gemini's base64 images) can be metered. Returns the
replay response, the evidence (a small JSON object of `keys`, b"" when the body is not a JSON
object) and the body's size. The file is unlinked from creation, so a crash leaks nothing on
disk; the replay's close (or its garbage collection) closes it and returns its budget, once."""
global _spool_in_use
settings = get_settings()
file = None
size = claimed = 0
def claim(total: int) -> None:
"""Hold budget for `total` bytes of this body; refuse (uncharged) when the process has no
room. Admission is all at once when the provider declares a length, so concurrent large
answers are admitted or refused whole instead of all stalling half-read."""
nonlocal claimed
global _spool_in_use
if total > settings.spool_max_bytes:
raise _buffer_limit(
f"upstream response exceeds treg's {settings.spool_max_bytes // (1024 * 1024)} "
"MiB settlement limit")
if total > claimed:
if _spool_in_use + total - claimed > settings.spool_budget_bytes:
raise _buffer_limit("treg is settling too many large responses right now; this "
"is temporary, retry shortly")
_spool_in_use += total - claimed
claimed = total
try:
try:
declared = next((v for k, v in response.raw_headers
if k.lower() == b"content-length"), b"")
if declared.isdigit():
claim(int(declared))
file = tempfile.TemporaryFile(dir=settings.spool_dir or None)
pending = bytearray()
async for chunk in response.body_stream:
raw = _as_bytes(chunk)
claim(size + len(raw))
size += len(raw)
pending += raw
if len(pending) >= _SPOOL_CHUNK:
await asyncio.to_thread(file.write, pending)
pending.clear()
if pending:
await asyncio.to_thread(file.write, pending)
finally:
await response.close()
async with _spool_parse_gate():
evidence = await asyncio.to_thread(_spool_evidence, file, keys)
except BaseException:
_release_spool(file, claimed)
raise
async def replay():
file.seek(0)
while chunk := await asyncio.to_thread(file.read, _SPOOL_CHUNK):
yield chunk
async def close() -> None:
release()
out = UpstreamResponse(response.status, _with_length(response.raw_headers, size), replay(), close)
release = weakref.finalize(out, _release_spool, file, claimed)
return out, evidence, size
async def _read_evidence(
mk: MarketplaceCall, response: UpstreamResponse,
) -> tuple[UpstreamResponse, bytes, int | None]:
"""Complete settlement evidence for a metered or owned-poll answer: the whole body in memory,
or, for a metered 2xx on an endpoint declaring `spooled_response`, the body on disk and the
keys its usage settlement reads. The third value is the spooled size, None when the body is in memory.
A routed child always reads into memory: its parent builds the answer from the child's body."""
if (mk.spooled_evidence and mk.metered and mk.deferred is None
and 200 <= response.status < 300):
return await _spool_response(response, mk.spooled_evidence)
response, body = await _buffer_response(response)
return response, body, None
async def _peek_stream_head(response: UpstreamResponse, limit: int) -> tuple[UpstreamResponse, bytes]:
"""Read at most ``limit`` response bytes for evidence, then replay every byte to the caller.
+10
View File
@@ -156,6 +156,16 @@ class Settings(BaseSettings):
# (a 30s ceiling 502'd one live), and merchant routes have been observed at 9s→105s.
call_timeout_s: int = 180
hold_grace_s: int = 60
# Spooled settlement evidence (application/call/settle.py `_spool_response`): a catalog
# endpoint declaring `spooled_response` writes its metered 2xx body to an anonymous temp file
# instead of RAM, reads the top-level keys its usage settles from, settles, then relays the file.
# Per-body cap, per-process budget across concurrent spools (over it a call fails uncharged,
# like the 8 MiB buffer), and how many bodies one process parses for evidence at once.
# `spool_dir` empty = the system temp dir; the files are unlinked from birth.
spool_dir: str = ""
spool_max_bytes: int = Field(default=64 * 1024 * 1024, ge=1)
spool_budget_bytes: int = Field(default=512 * 1024 * 1024, ge=1)
spool_parse_concurrency: int = Field(default=2, ge=1, le=16)
# ---- referral program (referrals.py) --------------------------------------------------------
# Flat bounties, not a percentage of top-ups. At 0% platform margin a percentage would be a
+3
View File
@@ -708,6 +708,9 @@ def _normalize(raw: dict, provider: str, directory: Path) -> dict:
"resource_ownership": raw.get("resource_ownership") or None,
# User-visible long-lived objects created on a shared provider account.
"managed_resource": raw.get("managed_resource") or None,
# A metered 2xx too large to buffer (inline media): read to disk and settled from the
# keys its usage terms read (application/call/settle.py `_spool_response`).
"spooled_response": raw.get("spooled_response") is True,
"platform_request": raw.get("platform_request") or None,
# How treg serves the catalog fallback after the team's own tool/credential ladder misses.
# Absent means the provider credential is required. `anonymous` means the verified public
+8
View File
@@ -188,6 +188,14 @@ def usage_term_path(document: object, dotted: str) -> object:
return current
def usage_roots(cost: dict) -> tuple[str, ...]:
"""The top-level response keys a `usage` settlement reads - all a spooled answer keeps."""
usage = cost.get("usage") or {}
paths = [term.get("path") for term in usage.get("terms") or []] or [usage.get("path")]
roots = (str(path).split(".", 1)[0].split("[", 1)[0] for path in paths if path)
return tuple(dict.fromkeys(roots))
def usage_evidence(basis: dict, evidence: dict[str, Any]) -> float | None:
"""The provider-reported usage figure a `usage` basis settles on, or None when the terminal
response does not carry a usable one.
+3
View File
@@ -1073,6 +1073,9 @@ def test_usage_terms_sum_every_reported_meter_at_its_rate():
# A malformed meter poisons the figure rather than silently billing less.
bad = {"usageMetadata": {**usage, "thoughtsTokenCount": -1}}
assert settlement.usage_evidence(basis, {"terminal": bad}) is None
# A spooled answer keeps only the keys these meters start at.
assert settlement.usage_roots(cost) == ("usageMetadata",)
assert settlement.usage_roots({"usage": {"path": "usage.cost", "unit": "usd"}}) == ("usage",)
def test_basis_derivation_and_settlement_table_vs_usage():
+201
View File
@@ -0,0 +1,201 @@
"""Spooled settlement evidence: a large metered 2xx is read to disk, never into RAM.
`_spool_response` keeps `_buffer_response`'s contract (complete evidence before headers, never a
prefix, an oversized body fails uncharged) for endpoints that declare `spooled_response`. These
tests pin the byte-exact replay, the evidence projection, and that every path returns the file
and its budget: success, oversize, a busy process, an upstream failure, cancellation and a
response that is dropped without being closed.
"""
import asyncio
import gc
import json
import pytest
from treg.application.call import settle
from treg.application.call.settle import _spool_response
from treg.application.call.types import GatewayFailed, UpstreamResponse
from treg.config import get_settings
MB = 1024 * 1024
IMAGE = b"A" * (3 * MB) # base64 is quote- and escape-free, like a real inlineData payload
USAGE = {"promptTokenCount": 17, "candidatesTokenCount": 1229,
"candidatesTokensDetails": [{"modality": "IMAGE", "tokenCount": 1120}]}
BODY = (b'{"candidates":[{"content":{"parts":[{"inlineData":{"mimeType":"image/jpeg","data":"'
+ IMAGE + b'"}}]},"finishReason":"STOP"}],"usageMetadata":'
+ json.dumps(USAGE).encode() + b',"modelVersion":"gemini-3-pro-image"}')
@pytest.fixture(autouse=True)
def fresh_budget(monkeypatch):
monkeypatch.setattr(settle, "_spool_in_use", 0)
monkeypatch.setattr(settle, "_spool_gate", None)
yield
gc.collect()
assert settle._spool_in_use == 0
def _upstream(body: bytes, *, chunk: int = 64 * 1024, fail_after: int | None = None,
declared: int | None = None):
state = {"closes": 0, "read": 0}
async def stream():
for start in range(0, len(body), chunk):
if fail_after is not None and start >= fail_after:
raise ConnectionResetError("synthetic upstream reset")
state["read"] += 1
yield body[start:start + chunk]
async def close():
state["closes"] += 1
headers = [(b"content-type", b"application/json")]
if declared is not None:
headers.append((b"content-length", str(declared).encode()))
return UpstreamResponse(200, tuple(headers), stream(), close), state
async def _drain(response: UpstreamResponse) -> bytes:
return b"".join([chunk async for chunk in response.body_stream])
async def test_replay_is_byte_exact_and_evidence_is_the_declared_keys():
upstream, state = _upstream(BODY, declared=len(BODY))
response, evidence, size = await _spool_response(upstream, ("usageMetadata", "missing"))
assert state["closes"] == 1 # the upstream is done before headers go out
assert size == len(BODY)
assert json.loads(evidence) == {"usageMetadata": USAGE}
assert dict(response.raw_headers)[b"content-length"] == str(len(BODY)).encode()
assert settle._spool_in_use == len(BODY)
assert await _drain(response) == BODY
await response.close()
await response.close() # close is idempotent: the router may close twice
assert settle._spool_in_use == 0
@pytest.mark.parametrize("body", [b"not json", b'["a list"]', b'{"truncated": "'])
async def test_a_body_that_is_not_a_json_object_yields_no_evidence(body):
upstream, _ = _upstream(body)
response, evidence, size = await _spool_response(upstream, ("usageMetadata",))
assert evidence == b"" and size == len(body)
assert await _drain(response) == body
await response.close()
async def test_an_oversized_body_fails_uncharged_and_returns_its_budget(monkeypatch):
monkeypatch.setattr(get_settings(), "spool_max_bytes", 2 * MB)
upstream, state = _upstream(BODY)
with pytest.raises(GatewayFailed) as failure:
await _spool_response(upstream, ("usageMetadata",))
assert failure.value.kind == "response_buffer_limit"
assert "not charged" in failure.value.detail["message"]
assert state["closes"] == 1 and settle._spool_in_use == 0
async def test_a_busy_process_refuses_instead_of_filling_the_disk(monkeypatch):
monkeypatch.setattr(get_settings(), "spool_budget_bytes", 4 * MB)
first, _ = _upstream(BODY)
held, _, _ = await _spool_response(first, ("usageMetadata",))
second, state = _upstream(BODY)
with pytest.raises(GatewayFailed) as failure:
await _spool_response(second, ("usageMetadata",))
assert "too many large responses" in failure.value.detail["message"]
assert state["closes"] == 1 and settle._spool_in_use == len(BODY)
await held.close()
third, _ = _upstream(BODY) # the budget came back with the first close
again, _, _ = await _spool_response(third, ("usageMetadata",))
await again.close()
async def test_an_upstream_reset_mid_body_releases_everything():
upstream, state = _upstream(BODY, fail_after=MB)
with pytest.raises(ConnectionResetError):
await _spool_response(upstream, ("usageMetadata",))
assert state["closes"] == 1 and settle._spool_in_use == 0
async def test_cancellation_while_spooling_releases_everything():
started = asyncio.Event()
state = {"closes": 0}
async def stream():
yield b"x" * MB
started.set()
await asyncio.sleep(3600)
yield b""
async def close():
state["closes"] += 1
task = asyncio.create_task(_spool_response(UpstreamResponse(200, (), stream(), close), ()))
await started.wait()
task.cancel()
with pytest.raises(asyncio.CancelledError):
await task
assert state["closes"] == 1 and settle._spool_in_use == 0
async def test_a_dropped_response_still_returns_its_budget():
upstream, _ = _upstream(BODY)
response, _, _ = await _spool_response(upstream, ("usageMetadata",))
assert settle._spool_in_use == len(BODY)
del response
gc.collect()
assert settle._spool_in_use == 0
async def test_parsing_is_gated_per_process(monkeypatch):
monkeypatch.setattr(get_settings(), "spool_parse_concurrency", 1)
active = peak = 0
original = settle._spool_evidence
def tracked(file, keys):
nonlocal active, peak
active += 1
peak = max(peak, active)
try:
return original(file, keys)
finally:
active -= 1
monkeypatch.setattr(settle, "_spool_evidence", tracked)
results = await asyncio.gather(*[
_spool_response(_upstream(BODY)[0], ("usageMetadata",)) for _ in range(3)])
assert peak == 1
for response, evidence, _ in results:
assert json.loads(evidence) == {"usageMetadata": USAGE}
await response.close()
async def test_a_declared_length_is_admitted_or_refused_whole(monkeypatch):
"""Concurrent large answers must not all claim budget piecemeal and stall half-read: with a
Content-Length the whole body is admitted up front, or refused before a byte is read."""
monkeypatch.setattr(get_settings(), "spool_budget_bytes", len(BODY) + MB)
first, _ = _upstream(BODY, declared=len(BODY))
second, second_state = _upstream(BODY, declared=len(BODY))
held, _, _ = await _spool_response(first, ("usageMetadata",))
with pytest.raises(GatewayFailed) as failure:
await _spool_response(second, ("usageMetadata",))
assert "temporary" in failure.value.detail["message"]
assert second_state["read"] == 0 and second_state["closes"] == 1
await held.close()
async def test_a_declared_length_over_the_cap_is_refused_before_reading(monkeypatch):
monkeypatch.setattr(get_settings(), "spool_max_bytes", 2 * MB)
upstream, state = _upstream(BODY, declared=len(BODY))
with pytest.raises(GatewayFailed):
await _spool_response(upstream, ("usageMetadata",))
assert state["read"] == 0 and state["closes"] == 1
async def test_a_temp_file_that_cannot_be_created_still_closes_the_upstream(monkeypatch):
def no_disk(*args, **kwargs):
raise OSError("synthetic: no space left on device")
monkeypatch.setattr(settle.tempfile, "TemporaryFile", no_disk)
upstream, state = _upstream(BODY, declared=len(BODY))
with pytest.raises(OSError):
await _spool_response(upstream, ("usageMetadata",))
assert state["closes"] == 1
+28
View File
@@ -912,3 +912,31 @@ def test_usage_terms_reject_bad_shapes(usage, message):
errors: list[str] = []
validator.check_cost(cost, "x", errors, [], provider="google-ai")
assert any(message in e for e in errors), errors
def _spooled_endpoint() -> dict:
cost = _flat_usage_cost()
cost["usage"] = {"terms": _TERMS, "unit": "usd"}
return {"id": "google-ai.image-gen.demo", "cost": cost, "spooled_response": True}
def test_spooled_response_settles_from_reported_usage():
errors: list[str] = []
validator.check_spooled_response(_spooled_endpoint(), None, "x", errors)
assert errors == []
@pytest.mark.parametrize(("mutate", "is_async", "message"), [
(lambda ep: ep.update(spooled_response={"evidence": ["usageMetadata"]}), False, "must be true"),
(lambda ep: ep.update(spooled_response=False), False, "must be true"),
(lambda ep: ep.update(resource_ownership={"requires": {}}), False, "cannot be combined"),
(lambda ep: None, True, "cannot be combined"),
(lambda ep: ep["cost"].update(settle="base"), False, "needs settle: usage"),
])
def test_spooled_response_rejects_bodies_something_else_must_read(mutate, is_async, message):
ep = _spooled_endpoint()
mutate(ep)
errors: list[str] = []
validator.check_spooled_response(ep, {"id_from": "id"} if is_async else None, "x", errors)
assert any(message in e for e in errors), errors