Synthesis follow-ups from the combined-tree review; register the grown stream suite in the size ratchet

Adversarial review of the combined tree c18c681bd (sprint plan §9.4):
- provider_contract_ci: a complete streamed reply without a usage frame is a warning
  (streamed_reply_without_usage_frame), never RED; the recorded corpus opts out of the
  whitespace check (-whitespace) as it does out of EOL conversion.
- llm_stream: a native rejection settles only on message_delta's final counters
  (usage_final), never on message_start's lower bound; a non-JSON native frame is the
  typed unknown outcome; the tool-call kind is inferred for either payload; a clean close
  is terminal framing only when the whole expected choice set finished; four more
  forgiveness rows pinned (object/list/text delta onto another shape, data after [DONE]
  arriving in one network read).
- loop_transport: the grant→transport_unavailable transition tells the owner the new
  attempt never reached the provider; the redial-crossed-dispatch branch guards on
  msg_present locally.
- ARCHITECTURE: unknown-outcome continuation is reserved for a stream that never reached
  its terminal frame (a lost socket, or a malformed mid-stream native frame).
- size_ratchet_manifest: tests/test_transport_b_stream_deadlines.py entered the 1001-1500
  band with its rationale (the regenerator refuses this tree because of the base's own debt).
- loop-level pin: a code-less SSE error read before any usage frame keeps its unresolved
  upper bound and still walks the cross-model chain (disclosed residual).

Co-authored-by: Ouroboros <311266734+ouroboros-agent@users.noreply.github.com>
This commit is contained in:
Ouroboros 2026-09-14 08:40:11 +03:00
parent 212746bffa
commit c75aa759e9
10 changed files with 186 additions and 23 deletions

2
.gitattributes vendored
View file

@ -15,4 +15,4 @@ tests/fixtures/skill_publish_scanner/*.fixture text eol=lf
# Recorded provider wire (SSE / JSON bodies) is byte-exact evidence: no line-ending
# conversion at commit or checkout, so the replay corpus always sees the recorded bytes.
tests/fixtures/llm_wire/** -text
tests/fixtures/llm_wire/** -text -whitespace

View file

@ -1264,7 +1264,7 @@ The Claude Agent SDK gateway is retired (owner-consented). `review_native_episod
Late review completion is retained in the existing operation-addressed prompt/response CAS. The response manifest stamps a complete producer outcome and the original task/root, attempt, slot/route, operation, subject, contract and roster/epoch binding. Exact reconciliation reads and verifies that complete source and runs the existing surface reducer; a new context needs no second paid attempt or surface writer. Pending, unreadable, mismatched and unknown outcomes retain custody and their full source references. Plan waves carry their original dispatched set, health epoch and operation IDs through the same reconciliation path; free replay spends no new cycle. Skill Review reserves every chunk digest/retry key and slot operation ID in the existing `review_job.review_wave` before paid dispatch, then retains that binding in terminal history. The next authorized public call may record `review_resume_of` only for the immediately preceding unsuperseded wave with identical task/root/attempt, skill/state roots, group, content, contract, rebuttal and chunks. Its new lifecycle reaggregates complete CAS at the original wave ID without another paid stamp, even at the cycle cap; old timeout/history rows remain unchanged. Missing or partial producers stay pending and unstarted chunks do not acquire verdict authority from their reserved IDs. A new task/root, changed material/contract/rebuttal, superseding lifecycle or explicit cancellation cannot inherit the old wave. Review-state persistence, current grants/dependencies and enablement still follow the current lifecycle guards; late completion alone never revives a job or enables a skill.
Main remote completions use transport streaming, assembled inside the physical send closure before accounting settles. Compatible choices/tool-call fragments and native blocks, reasoning/signatures, usage snapshots and terminal framing produce the same normalized response shape as JSON. Partial streams never yield usable tool calls or answers; their exact wire bytes and partial assembly remain in private CAS through the existing `physical_stream` manifest and attempt ID. Wire form is tolerated (identity scalars and metadata keep their first value, index gaps are forgiven, every forgiven fact is disclosed in the stream receipt) and only completeness is enforced; a complete-but-unusable stream settles with its usage and classifies as a provider error, so unknown-outcome continuation is reserved for real socket loss. A complete response with absent final usage still has unknown money. Comments/pings do not define cognitive deadlines. Every recovery candidate checks inherited calendar and quota-adjusted execution bounds before reservation, after preparation and at dispatch; HTTP phase bounds remain distinct from an overall logical wait.
Main remote completions use transport streaming, assembled inside the physical send closure before accounting settles. Compatible choices/tool-call fragments and native blocks, reasoning/signatures, usage snapshots and terminal framing produce the same normalized response shape as JSON. Partial streams never yield usable tool calls or answers; their exact wire bytes and partial assembly remain in private CAS through the existing `physical_stream` manifest and attempt ID. Wire form is tolerated (identity scalars and metadata keep their first value, index gaps are forgiven, every forgiven fact is disclosed in the stream receipt) and only completeness is enforced; a complete-but-unusable stream settles with its usage and classifies as a provider error, so unknown-outcome continuation is reserved for a stream that never reached its terminal frame (a lost socket, or a malformed mid-stream native frame). A complete response with absent final usage still has unknown money. Comments/pings do not define cognitive deadlines. Every recovery candidate checks inherited calendar and quota-adjusted execution bounds before reservation, after preparation and at dispatch; HTTP phase bounds remain distinct from an overall logical wait.
An ordinary managed task whose provider outcome becomes unknown now stays in the existing transport-wait episode. The old attempt and unreported cost remain unknown. After non-generating upstream observation, a user-role `[SYSTEM NOTICE]` supplies explicit recovery input for one new physical attempt; existing budget, cancellation, owner deadline and absolute ceiling still apply. A granted attempt that fails again, unknown after dispatch or released before it, returns to the same episode, and so does a free redial that crosses dispatch and dies unknown: the backoff keeps growing (4→60 s), one redial is counted per wait iteration, an unknown repeat re-arms the probe's freshness bound at that latest failure and refreshes the custody, one `continued` row names the transition (`continuation_outcome_unknown`, `continuation_transport_unavailable`, `redial_outcome_unknown`), and there is no cap on continuations; budget, deadline and Stop remain the only stops. Finished tools are retained. Direct turns and configured session-nanny custody retain their own contracts; manual Restart/Panic gains no resume authority. Direct remote endpoints are observed through HEAD with the same no-proxy policy. The metadata HEAD reuses the ordinary connection allowance for every socket phase, narrowed by the owner remainder; it holds no cognitive in-flight lease. Subscription metadata qualifies only when the existing catalog reports generic `provenance="provider_http"` and an original `observedAt` after wait entry for the selected source/model and effective profile/account fingerprint. Cached reuse never advances that timestamp; a local handshake, static or pre-outage catalog, or timestamp without provider provenance cannot prove recovery. The pinned raw-model adapter is Codex; capability discovery remains authoritative rather than a new core provider table. Loss of the Claudexor control connection first keeps reading the same accepted model operation, including across endpoint rediscovery, without creating another operation.

View file

@ -208,7 +208,7 @@ def _tool_call_problem(call: Any) -> str:
return "id missing"
# Structural completeness is id + name + arguments; a fragment set that never
# carried ``type`` is the same shape the non-stream path executes.
kind = call.get("type") or ("function" if isinstance(call.get("function"), dict) else None)
kind = call.get("type") or next((k for k in ("function", "custom") if isinstance(call.get(k), dict)), None)
if kind not in {"function", "custom"}:
return f"type {kind!r}"
payload = call.get(kind)
@ -325,7 +325,8 @@ class ChatAccumulator(_Accumulator):
# A body the provider closed cleanly after every choice finished is
# terminal framing too ([DONE] is the other witness); a close before
# that is the one unknown outcome this assembler still reports.
if not self.choices or any(not choice.get("finish_reason") for choice in self.choices.values()):
if (set(self.choices) != set(range(self.expected_choices))
or any(not choice.get("finish_reason") for choice in self.choices.values())):
raise IncompleteProviderStream("Stream ended without complete terminal framing")
self._note("stream closed without [DONE] after every choice finished")
self.done = True
@ -374,13 +375,17 @@ class AnthropicAccumulator(_Accumulator):
self.body: dict = {}
self.blocks: dict[int, dict] = {}
self.open_blocks: set[int] = set()
self.usage_final = False # message_delta folded its final counters into body["usage"]
self.inputs: dict[int, str] = {}
def accept(self, event: str, data: str) -> None:
if self.done:
self._note("data after message_stop; ignored")
return
chunk = json.loads(data)
try:
chunk = json.loads(data)
except ValueError:
raise IncompleteProviderStream("Native stream chunk is not JSON") from None
if not isinstance(chunk, dict):
raise IncompleteProviderStream("Native stream chunk is not an object")
kind = chunk.get("type")
@ -440,6 +445,7 @@ class AnthropicAccumulator(_Accumulator):
raise IncompleteProviderStream("Native message delta before blocks finished")
_snapshot(self.body, chunk.get("delta") or {})
_snapshot(self.body.setdefault("usage", {}), chunk.get("usage") or {})
self.usage_final = True
elif kind == "message_stop":
self.done = True
# Future non-content events are retained in the exact wire evidence.
@ -469,8 +475,11 @@ class AnthropicAccumulator(_Accumulator):
return self.partial()
def _reject(self, detail: str) -> None:
# message_start's usage is a lower-bound snapshot: without message_delta the
# attempt keeps its unresolved upper bound (same policy as the error branch).
raise RejectedProviderStream(f"Native stream rejected after message_stop: {detail}",
usage=self.body.get("usage"), anomalies=self.anomalies)
usage=self.body.get("usage") if self.usage_final else None,
anomalies=self.anomalies)
def partial(self) -> dict:
return {**copy.deepcopy(self.body), "content": [copy.deepcopy(self.blocks[key]) for key in sorted(self.blocks)]}

View file

@ -306,9 +306,12 @@ def reconcile_transport_wait(
else "continuation_transport_unavailable"),
outcome_custody=episode.outcome_custody,
)
if error_kind == "transport_unavailable": # the grant note said "continuing"; the owner must hear otherwise
emit_progress("🌐 The new attempt could not reach the provider — waiting and redialing "
"automatically (that attempt was $0).", incident=None)
return episode
if (episode is not None and not after_local_pass and episode.wait_cause != "provider_outcome_unknown"
and unknown_again):
if (episode is not None and not msg_present and not after_local_pass
and episode.wait_cause != "provider_outcome_unknown" and unknown_again):
# A formerly free redial crossed dispatch and died unknown: the same
# episode now needs upstream proof before its next attempt; the clock,
# backoff and redial count carry over instead of a fresh 4s episode.

View file

@ -213,6 +213,7 @@ BAND_PATHS = {
"tests/test_terminal_durability_v664.py": "Entered the band from 974 lines: terminal durability coverage now pins retry-admission failure custody so an unpersisted terminal row cannot publish task_done or lose the retry marker.",
"tests/test_timeout_policy.py": "Adaptive timeout and custody regression suite covers raw-deadline admission, explicit finalization reserve, transport bounds, and late-result reconciliation.",
"tests/test_tool_result.py": "F3.1 typed-organ pins carried with the D02 organ (D04 entry 9): the closed code table, the one legacy-text adapter, the publish/sidecar seam and the meta-boundary contracts pin one organ in one suite; sibling suites (meta_boundaries, t46, classification differential) already hold the spill-over families.",
"tests/test_transport_b_stream_deadlines.py": "Stream-assembler contract suite: the issue #856 doctrine pins (terminal-vs-form verdicts, recorded-wire driver replays, native Messages assembly, the forgiveness table) share one wire/driver/ledger fixture; keeping them in one focused suite below the 1500-line band cap.",
"tests/test_transport_death_retry.py": "Physical transport-repeat and exceptional-loop evidence tests share the same scripted provider and real tool-execution fixture; retain this coherent contract suite below the module cap.",
"tests/test_tree_cost_ceiling.py": "Budget-rail coverage entered the band with the cache-split ownership, candidate-predicate, soft-landing, probe-confirmed stop and one-row-per-delegated-run regressions; one focused suite for the tree cost ceiling.",
"tests/test_ui_smoke_project_continuity.py": "Playwright smoke of the Project continuity contracts (panel/Main re-homing, lifecycle rows, the Main-root project pointer): each test drives one end-to-end owner flow across both surfaces, so the cross-surface assertions cannot be split into smaller files without losing what they prove.",

View file

@ -491,8 +491,18 @@ def assert_canary_usage(usage, canary: ProviderCanary, *, forced_tool_choice: bo
assert isinstance(usage, dict), failure("usage_not_mapping")
assert usage.get("provider") == canary.expected_provider, failure("unexpected_accounting_provider")
assert usage.get("resolved_model") == normalize_model_identity(canary.model), failure("unexpected_accounting_model")
assert _safe_nonnegative_int(usage.get("prompt_tokens")) > 0, failure("missing_prompt_tokens")
assert _safe_nonnegative_int(usage.get("completion_tokens")) > 0, failure("missing_completion_tokens")
if (isinstance(usage.get("stream_receipt"), dict) and not _safe_nonnegative_int(usage.get("prompt_tokens"))
and not _safe_nonnegative_int(usage.get("completion_tokens"))):
# A complete streamed reply whose route sent no final usage frame (a provider
# that drops stream_options under wire recovery, or ignores include_usage):
# money is unknown, the contract is not broken — weather, recorded as a
# warning through the existing telemetry, never RED.
warnings = usage.setdefault("canary_warnings", [])
if isinstance(warnings, list) and len(warnings) < 4:
warnings.append({**_canary_evidence(canary, None, usage), "code": "streamed_reply_without_usage_frame"})
else:
assert _safe_nonnegative_int(usage.get("prompt_tokens")) > 0, failure("missing_prompt_tokens")
assert _safe_nonnegative_int(usage.get("completion_tokens")) > 0, failure("missing_completion_tokens")
if canary.reasoning_effort == "medium":
expected_effort = canary.reasoning_effort
if canary.expected_provider == "deepseek":

View file

@ -174,6 +174,24 @@ def test_stream_loss_before_the_terminal_frame_is_inconclusive_but_a_rejected_bo
ProviderFailureKind.RED if code == 400 else ProviderFailureKind.INCONCLUSIVE, expected)
def test_streamed_reply_without_usage_frame_is_a_warning_not_a_contract_violation():
"""Main streams and so does the canary's first turn; a route may answer completely without a
final usage frame (stream_options dropped by wire recovery, include_usage ignored): money is
unknown, the contract is not broken — a warning, never RED. The non-stream shape still demands tokens."""
import dataclasses
from ouroboros.provider_models import normalize_model_identity
from tests.provider_contract_ci import assert_canary_usage
canary = dataclasses.replace(next(iter(provider_canary_matrix())), expected_provider="openrouter", reasoning_effort="high")
usage = {"provider": "openrouter", "resolved_model": normalize_model_identity(canary.model),
"stream_receipt": {"complete": True, "anomalies": {"count": 0, "first": []}}, "canary_warnings": []}
assert_canary_usage(usage, canary)
assert [warning["code"] for warning in usage["canary_warnings"]] == ["streamed_reply_without_usage_frame"]
with pytest.raises(AssertionError):
assert_canary_usage({**usage, "stream_receipt": None, "canary_warnings": []}, canary)
def test_provider_alarm_output_sanitizes_token_shaped_evidence(capsys):
sentinel = "sk-proj-" + ("A" * 40)
exc = _http_error(429, f'{{"error":{{"token":"{sentinel}"}}}}')

View file

@ -69,17 +69,17 @@ class WireResponse:
reason = "OK"
url = "https://provider.invalid/v1/messages"
def __init__(self, wire, *, failure=None, step=None):
self.wire, self.failure, self.step = wire, failure, step
def __init__(self, wire, *, failure=None, step=None, chunk_size=13):
self.wire, self.failure, self.step, self.chunk_size = wire, failure, step, chunk_size
self.headers = {"x-generation-id": "header-generation"}
self.closed = False
def iter_bytes(self):
# Fragment through UTF-8, CRLF and JSON boundaries.
for offset in range(0, len(self.wire), 13):
# Fragment through UTF-8, CRLF and JSON boundaries (13 bytes by default).
for offset in range(0, len(self.wire), self.chunk_size):
if self.step:
self.step()
yield self.wire[offset:offset + 13]
yield self.wire[offset:offset + self.chunk_size]
if self.failure:
raise self.failure
@ -333,6 +333,20 @@ _FORGIVEN_SHAPES = {
"model_conflict": (
lambda: sse(chunk({"content": "do"}), chunk({"content": "ne"}, "stop", usage=completion()["usage"], model="vendor/other")),
["model: 'vendor/test-stream' then 'vendor/other'; kept first"]),
"object_delta_onto_text": (
lambda: sse(chunk({"content": "done"}), chunk({"content": {"x": 1}}, "stop", usage=completion()["usage"])),
["content: object delta onto str; kept first shape"]),
"list_delta_onto_text": (
lambda: sse(chunk({"content": "done"}), chunk({"content": [1]}, "stop", usage=completion()["usage"])),
["content: list delta onto str; kept first shape"]),
"text_delta_onto_scalar": (
lambda: sse(chunk({"content": 7}), chunk({"content": "done"}, "stop", usage=completion()["usage"])),
["content: text delta onto int; kept first shape"]),
# [DONE] and a later frame in ONE network read: the assembler drains the whole chunk.
"data_after_done": (
lambda: WireResponse(sse(chunk({"content": "done"}, "stop", usage=completion()["usage"]))
+ b"data: {\"id\": \"late\"}\r\n\r\n", chunk_size=1 << 16),
["data after [DONE]; ignored"]),
}
@ -342,7 +356,9 @@ def test_form_irregularities_never_raise_before_the_terminal_verdict(isolated, s
Every forgiven shape assembles a complete reply, settles, and names the forgiven fact in the
receipt; none of them is an unknown outcome."""
build, expected = _FORGIVEN_SHAPES[shape]
result = run_driver(lambda **kw: WireResponse(build()), payload(stream=True), target()).model_dump()
wire = build()
response = wire if isinstance(wire, WireResponse) else WireResponse(wire)
result = run_driver(lambda **kw: response, payload(stream=True), target()).model_dump()
receipt = result["_stream_receipt"]
assert receipt["complete"] is True
assert all(note in receipt["anomalies"]["first"] for note in expected), receipt["anomalies"]
@ -363,17 +379,26 @@ def test_clean_close_after_every_choice_finished_is_terminal_framing(isolated):
assert rows(isolated)[-1]["state"] == "settled"
def test_tool_call_without_type_is_structurally_complete(isolated):
@pytest.mark.parametrize("payload_key,field", [("function", "arguments"), ("custom", "input")])
def test_tool_call_without_type_is_structurally_complete(isolated, payload_key, field):
"""``type`` is not part of structural completeness (id, name, arguments are): a call whose
fragments never carried it returns exactly as the non-stream path returns it."""
wire = sse(chunk({"tool_calls": [{"index": 0, "id": "t", "function": {"name": "lookup", "arguments": '{"q":"x"}'}}]},
"tool_calls", usage=completion()["usage"]))
fragments never carried it returns exactly as the non-stream path returns it, for either payload."""
call = {"index": 0, "id": "t", payload_key: {"name": "lookup", field: '{"q":"x"}'}}
wire = sse(chunk({"tool_calls": [call]}, "tool_calls", usage=completion()["usage"]))
result = run_driver(lambda **kw: WireResponse(wire), payload(stream=True), target()).model_dump()
assert result["choices"][0]["message"]["tool_calls"] == [
{"id": "t", "function": {"name": "lookup", "arguments": '{"q":"x"}'}}]
assert result["choices"][0]["message"]["tool_calls"] == [{"id": "t", payload_key: {"name": "lookup", field: '{"q":"x"}'}}]
assert result["_stream_receipt"]["anomalies"]["count"] == 0 and rows(isolated)[-1]["state"] == "settled"
def test_clean_close_with_a_missing_expected_choice_stays_unknown(isolated):
"""``n=2``, one choice finished, the body closes without ``[DONE]`` and without the second
choice: the wire is genuinely incomplete, so this is the unknown outcome, not a rejection."""
wire = sse(chunk({"role": "assistant", "content": "only one"}, "stop", usage=completion()["usage"]), done=False)
with pytest.raises(IncompleteProviderStream):
run_driver(lambda **kw: WireResponse(wire), payload(stream=True, n=2), target())
assert rows(isolated)[-1]["state"] == "unresolved"
def test_later_usage_snapshot_overrides_an_earlier_one(isolated):
"""Usage frames are cumulative snapshots: the last one read is the one settled, so an early
partial snapshot never becomes the attempt's cost."""
@ -804,6 +829,35 @@ def test_native_unusable_body_after_message_stop_is_rejected_and_settled(isolate
assert (verdict.kind, verdict.retry_same_request) == ("provider_error", False)
def test_native_rejection_without_message_delta_keeps_money_unknown(isolated, monkeypatch):
"""``message_stop`` arrived but no ``message_delta`` did (no stop_reason, no final counters): the
body is rejected, and the ledger keeps its unresolved upper bound instead of settling on
``message_start``'s lower-bound snapshot — the same policy as the mid-stream error frame."""
import requests
events, _expected = native_events()
events = [(kind, body) for kind, body in events if kind != "message_delta"]
monkeypatch.setattr(requests, "post", lambda *a, **k: WireResponse(sse(*events, done=False)))
with pytest.raises(RejectedProviderStream) as caught:
LLMClient()._chat_anthropic(target("anthropic"), MESSAGES, TOOLS, "high", 1024, "auto", stream=True)
assert "no stop_reason after message_stop" in str(caught.value) and caught.value.stream_usage is None
assert rows(isolated)[-1]["state"] == "unresolved"
assert classify_llm_exception(caught.value).kind == "provider_error"
def test_native_non_json_frame_is_the_typed_unknown_outcome(isolated, monkeypatch):
"""A truncated ``data:`` line in a native stream is ``IncompleteProviderStream`` (typed
``model_outcome_unknown``), never a raw ``JSONDecodeError`` whose classification would rest
on custody state alone."""
import requests
events, _expected = native_events()
wire = sse(*events[:3], done=False) + b"data: {broken\r\n\r\n"
monkeypatch.setattr(requests, "post", lambda *a, **k: WireResponse(wire))
with pytest.raises(IncompleteProviderStream) as caught:
LLMClient()._chat_anthropic(target("anthropic"), MESSAGES, TOOLS, "high", 1024, "auto", stream=True)
assert "not JSON" in str(caught.value) and caught.value.code == "model_outcome_unknown"
assert rows(isolated)[-1]["state"] == "unresolved"
def test_native_max_tokens_with_tool_blocks_returns_like_non_stream(isolated, monkeypatch):
"""``stop_reason == max_tokens`` beside a complete tool block is a finished reply the loop reads
(parity with the non-stream native path, which surfaces ``stop_reason``); it is not a rejection."""

View file

@ -1152,3 +1152,69 @@ def test_rejected_stream_settles_and_walks_the_fallback_chain_without_a_wait_epi
for transcript in (messages, kwargs["tools"]._ctx.messages):
assert not any("[SYSTEM NOTICE]" in str(m.get("content") or "") for m in transcript)
assert any("⚡ Fallback: test-model → other/model" in note and "provider_error" in note for note in notes)
def test_code_less_stream_error_without_usage_keeps_the_upper_bound_and_still_walks_the_chain(
data_root, tmp_path, monkeypatch, no_sleep,
):
"""An explicit SSE error frame without an HTTP-shaped code (the native overload shape) read
before any usage frame is the provider's own terminal verdict: ``provider_error`` with no
same-model repeat, the cross-model fallback dialed, no wait episode — while the first attempt
keeps its unresolved upper bound (message_start's snapshot is deliberately not settled), so the
budget fence still counts that money."""
from ouroboros import fallback_cooldown
from ouroboros.llm_stream import ProviderStreamError
fallback_cooldown.reset_for_tests()
monkeypatch.setenv("OUROBOROS_TASK_REVIEW_MODE", "off")
monkeypatch.setenv("OUROBOROS_MODEL_FALLBACKS", "other/model")
monkeypatch.delenv("USE_LOCAL_FALLBACK", raising=False)
monkeypatch.setattr(loop_mod, "_rebind_context_fit_plan", lambda *a, **k: (None, "max"))
class _OverloadedPrimaryLLM:
def __init__(self, root):
self.root = root
self.models = []
def default_model(self):
return "test-model"
def chat(self, **kwargs):
model = kwargs["model"]
self.models.append(model)
def send():
if model == "test-model":
raise ProviderStreamError({"type": "error", "error": {"type": "overloaded_error", "message": "Overloaded"}})
return {"content": "done"}
request = ua.AttemptRequest(
model=model, provider="anthropic", reservation_usd=1.0,
drive_root=self.root, task_id="t-death", root_task_id="t-death", source="test.overloaded",
)
ua.execute_physical_attempt(
request, send, extractor=lambda _resp: ({"prompt_tokens": 1, "completion_tokens": 1}, 0.01, True),
)
return OK_RESPONSE
llm = _OverloadedPrimaryLLM(data_root)
notes = []
kwargs = _loop_kwargs(tmp_path, llm, notes)
result, usage, _trace = run_llm_loop(**kwargs)
assert result == "done"
assert llm.models == ["test-model", "other/model"]
by_attempt = {}
for row in _ledger(data_root):
by_attempt.setdefault(row["attempt_id"], []).append(row)
assert [[row["state"] for row in rows] for rows in by_attempt.values()] == [
["reserved", "dispatched", "unresolved"],
["reserved", "dispatched", "settled"],
]
assert ua.usage_projection(data_root)["unresolved_upper_bound_usd"] == 1.0
api_errors = _events(tmp_path, "llm_api_error")
assert [(row["error_kind"], row["retry_same_request"], row["attempt_custody_state"]) for row in api_errors] == [
("provider_error", False, "unresolved"),
]
assert _events(tmp_path, "network_wait") == []
assert "transport_recovery" not in usage and "_pending_transport_outcome" not in usage

View file

@ -140,8 +140,10 @@ def test_alternating_unknown_and_transport_failures_keep_one_episode(tmp_path, m
notices = [row for row in sends[-1] if "NEW physical model attempt" in str(row.get("content"))]
assert len(notices) == 2 and "paid-attempt-3" in str(notices[-1]["content"])
# The owner is told that money became unknown when the free redial crossed dispatch (the
# entry note said "$0"), not only at the next grant.
# entry note said "$0"), not only at the next grant; and that the granted attempt never
# reached the provider (the grant note said "continuing").
assert sum("another charge is possible" in text for text in notes) >= 2
assert any("could not reach the provider" in text for text in notes)
def test_grant_without_a_new_attempt_is_not_a_phantom_repeat(tmp_path):