diff --git a/AGENTS.md b/AGENTS.md index dd055a14..8e29de69 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -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. diff --git a/docs/context/architecture/catalog.md b/docs/context/architecture/catalog.md index 8e29f2c4..6a0b4c5f 100644 --- a/docs/context/architecture/catalog.md +++ b/docs/context/architecture/catalog.md @@ -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. diff --git a/docs/context/architecture/money.md b/docs/context/architecture/money.md index 4d360a6d..1cc6141b 100644 --- a/docs/context/architecture/money.md +++ b/docs/context/architecture/money.md @@ -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 diff --git a/docs/context/architecture/proxy-model.md b/docs/context/architecture/proxy-model.md index 0105b3eb..c77da2e7 100644 --- a/docs/context/architecture/proxy-model.md +++ b/docs/context/architecture/proxy-model.md @@ -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 diff --git a/scripts/catalog_validate.py b/scripts/catalog_validate.py index 21513c66..453c6113 100644 --- a/scripts/catalog_validate.py +++ b/scripts/catalog_validate.py @@ -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, diff --git a/src/treg/application/call/resolve.py b/src/treg/application/call/resolve.py index 2c38669f..9bec6f80 100644 --- a/src/treg/application/call/resolve.py +++ b/src/treg/application/call/resolve.py @@ -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) diff --git a/src/treg/application/call/settle.py b/src/treg/application/call/settle.py index d4250286..386f893a 100644 --- a/src/treg/application/call/settle.py +++ b/src/treg/application/call/settle.py @@ -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( diff --git a/tests/test_call_spool.py b/tests/test_call_spool.py index 095d6d7e..9bc0d5ac 100644 --- a/tests/test_call_spool.py +++ b/tests/test_call_spool.py @@ -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) == () diff --git a/tests/test_catalog_validate.py b/tests/test_catalog_validate.py index 33d7583e..4326d5c8 100644 --- a/tests/test_catalog_validate.py +++ b/tests/test_catalog_validate.py @@ -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 == []