mirror of
https://github.com/superdesigndev/treg.git
synced 2026-10-02 03:24:35 +08:00
feat(call): keep a spooled answer's evidence by path, including a success rule
A spooled answer kept only the top-level objects its usage terms read, so only token-metered
rows could spool. A fixed-price answer (Lyria bills per song) settles on its `expect` success
leaf instead, and that leaf, `candidates.0.finishReason`, shares its top-level key with the
multi-megabyte audio.
Evidence is now a set of paths (`_spool_evidence_paths`): each usage term's top-level object
plus the `expect` leaf, projected at their original places (`{"candidates": [{"finishReason":
"STOP"}]}`) so the existing settlement readers work unchanged. The validator accepts
`spooled_response` with either `settle: usage` or an `expect` rule.
Fragments: architecture/proxy-model.md, money.md, catalog.md. AGENTS.md non-negotiable 4 names
the evidence precisely.
This commit is contained in:
@@ -49,8 +49,9 @@ Everything else in this file is guidance; these are the contract, and they win o
|
||||
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.
|
||||
unlinked temp file under a per-body cap and per-process budget, settles from only the paths its
|
||||
row reads (usage meters, `expect` success leaf), 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.
|
||||
|
||||
@@ -1032,10 +1032,12 @@ catalog metadata, not a provider-specific billing branch, and cannot be combined
|
||||
|
||||
`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.
|
||||
4K; Lyria songs arrive as base64 MP3). Its metered 2xx is read to disk and settled from exactly
|
||||
the paths its row reads: the top-level objects its usage terms start at, and its `expect` success
|
||||
leaf (`resolve._spool_evidence_paths`; proxy-model.md), so the evidence cannot drift from the
|
||||
price. A token-metered row settles on `usage`; a fixed price (Lyria's per song) needs an `expect`
|
||||
rule so a refused generation is not billed. The validator requires one of the two and refuses the
|
||||
field beside `async`, `resource_ownership` or `managed_resource`, which need the whole body.
|
||||
|
||||
A `pathParams` field that declares an `enum` is enforced on treg's key: the value names what the
|
||||
shared credential is spent on (Google AI's `model`), so any other value is a 400 before reserve.
|
||||
|
||||
@@ -1020,7 +1020,7 @@ does not observe the original generation task or persist a response for idempote
|
||||
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
|
||||
the body on disk: `_spool_response` settles from the paths its row reads (usage meters, `expect` leaf) 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
|
||||
|
||||
@@ -610,8 +610,10 @@ a web process. The endpoint instead declares `spooled_response: true` and
|
||||
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
|
||||
`spool_parse_concurrency` (2) per process. Only the paths the row's settlement reads survive
|
||||
(`_spool_evidence_paths`: each usage term's top-level object and the `expect` leaf, kept at its
|
||||
original place, e.g. `{"candidates": [{"finishReason": "STOP"}]}`), 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
|
||||
|
||||
@@ -608,18 +608,18 @@ def check_async_descriptor(descriptor: object, where: str, provider: str,
|
||||
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."""
|
||||
read to disk and settled from what the row declares it reads: its usage meters (`settle:
|
||||
usage`) or, for a fixed price, its `expect` success rule. 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")
|
||||
if (ep.get("cost") or {}).get("settle") != "usage" and not (ep.get("expect") or {}).get("json_path"):
|
||||
fail(errors, where, "spooled_response keeps only the evidence the row declares: it needs "
|
||||
"settle: usage or an expect success rule")
|
||||
|
||||
|
||||
def check_resource_ownership(rule: object, where: str, input_schema: object,
|
||||
|
||||
@@ -340,8 +340,8 @@ 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.
|
||||
# On an endpoint declaring `spooled_response`: the response paths its settlement reads (the
|
||||
# usage meters, the `expect` success rule). A metered 2xx goes to disk; only these are kept.
|
||||
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.
|
||||
@@ -1511,6 +1511,15 @@ def _enforce_apify_run_options(ep: dict, query) -> None:
|
||||
)
|
||||
|
||||
|
||||
def _spool_evidence_paths(ep: dict, cost: dict) -> tuple[str, ...]:
|
||||
"""What a spooled answer must keep for settlement: the top-level objects its usage terms read,
|
||||
and the leaf its `expect` success rule compares. Everything else stays on disk."""
|
||||
if not ep.get("spooled_response"):
|
||||
return ()
|
||||
rule = (ep.get("expect") or {}).get("json_path")
|
||||
return settlement_basis.usage_roots(cost) + ((str(rule),) if rule else ())
|
||||
|
||||
|
||||
def _enforce_platform_request(ep: dict, body: bytes, headers=None, query=None) -> None:
|
||||
"""Check explicit platform constraints and fixed pricing selectors before reserve/relay.
|
||||
|
||||
@@ -2085,8 +2094,7 @@ 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 ()),
|
||||
spooled_evidence=_spool_evidence_paths(ep, raw_cost),
|
||||
)
|
||||
if chosen_tool is not None:
|
||||
return MarketplaceCall(tool=chosen_tool, tier="tool", **common)
|
||||
|
||||
@@ -859,12 +859,38 @@ def _release_spool(file, claimed: int) -> 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`."""
|
||||
def _project(document, paths: tuple[str, ...]) -> dict:
|
||||
"""A document holding only the values at `paths`, each at its original place: `usageMetadata`
|
||||
keeps that whole object, `candidates.0.finishReason` keeps one leaf inside a one-item list.
|
||||
A path the document does not carry is simply absent, as it was in the original."""
|
||||
projection: dict = {}
|
||||
for path in paths:
|
||||
value = json_path(document, path)
|
||||
if value is None:
|
||||
continue
|
||||
parts = path.split(".")
|
||||
node: dict | list = projection
|
||||
for depth, part in enumerate(parts):
|
||||
key: int | str = int(part) if isinstance(node, list) else part
|
||||
if isinstance(node, list):
|
||||
node.extend({} for _ in range(key + 1 - len(node)))
|
||||
if depth == len(parts) - 1:
|
||||
node[key] = value
|
||||
break
|
||||
want = list if parts[depth + 1].isdigit() else dict
|
||||
current = node[key] if isinstance(node, list) else node.get(key)
|
||||
if not isinstance(current, want):
|
||||
node[key] = want()
|
||||
node = node[key]
|
||||
return projection
|
||||
|
||||
|
||||
def _spool_evidence(file, paths: tuple[str, ...]) -> bytes:
|
||||
"""The values at `paths` in a spooled JSON object, re-serialized in place; 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)
|
||||
@@ -872,7 +898,7 @@ def _spool_evidence(file, keys: tuple[str, ...]) -> bytes:
|
||||
return b""
|
||||
if not isinstance(document, dict):
|
||||
return b""
|
||||
return json.dumps({key: document[key] for key in keys if key in document}).encode()
|
||||
return json.dumps(_project(document, paths)).encode()
|
||||
|
||||
|
||||
async def _spool_response(
|
||||
|
||||
@@ -199,3 +199,26 @@ async def test_a_temp_file_that_cannot_be_created_still_closes_the_upstream(monk
|
||||
with pytest.raises(OSError):
|
||||
await _spool_response(upstream, ("usageMetadata",))
|
||||
assert state["closes"] == 1
|
||||
|
||||
|
||||
async def test_evidence_keeps_a_success_leaf_without_its_siblings():
|
||||
"""A fixed-price answer settles on its `expect` leaf: `candidates.0.finishReason` survives, the
|
||||
inline media beside it does not."""
|
||||
upstream, _ = _upstream(BODY, declared=len(BODY))
|
||||
response, evidence, _ = await _spool_response(
|
||||
upstream, ("usageMetadata", "candidates.0.finishReason"))
|
||||
assert json.loads(evidence) == {"usageMetadata": USAGE,
|
||||
"candidates": [{"finishReason": "STOP"}]}
|
||||
await response.close()
|
||||
|
||||
|
||||
def test_spool_evidence_paths_come_from_usage_and_the_success_rule():
|
||||
from treg.application.call.resolve import _spool_evidence_paths
|
||||
usage = {"settle": "usage", "usage": {"unit": "usd", "terms": [
|
||||
{"path": "usageMetadata.promptTokenCount", "rate": 1}]}}
|
||||
assert _spool_evidence_paths({"spooled_response": True}, usage) == ("usageMetadata",)
|
||||
rule = {"spooled_response": True,
|
||||
"expect": {"json_path": "candidates.0.finishReason", "equals": "STOP"}}
|
||||
assert _spool_evidence_paths(rule, {"type": "per_success", "value": 0.08}) == (
|
||||
"candidates.0.finishReason",)
|
||||
assert _spool_evidence_paths({}, usage) == ()
|
||||
|
||||
@@ -932,7 +932,7 @@ def test_spooled_response_settles_from_reported_usage():
|
||||
(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"),
|
||||
(lambda ep: ep["cost"].update(settle="base"), False, "settle: usage or an expect success rule"),
|
||||
])
|
||||
def test_spooled_response_rejects_bodies_something_else_must_read(mutate, is_async, message):
|
||||
ep = _spooled_endpoint()
|
||||
@@ -940,3 +940,14 @@ def test_spooled_response_rejects_bodies_something_else_must_read(mutate, is_asy
|
||||
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
|
||||
|
||||
|
||||
def test_a_fixed_price_spools_when_its_success_rule_is_declared():
|
||||
"""A per-song price has no meters: the `expect` leaf is the evidence a spooled answer keeps."""
|
||||
ep = _spooled_endpoint()
|
||||
ep["cost"].update(settle="base")
|
||||
ep["cost"].pop("usage")
|
||||
ep["expect"] = {"json_path": "candidates.0.finishReason", "equals": "STOP"}
|
||||
errors: list[str] = []
|
||||
validator.check_spooled_response(ep, None, "x", errors)
|
||||
assert errors == []
|
||||
|
||||
Reference in New Issue
Block a user