mirror of
https://github.com/razzant/ouroboros.git
synced 2026-10-03 04:07:04 +00:00
P3: transport honesty and route evidence for the subscription lane
Implements sprint phase P3 (steps P3-1..P3-7) from the Claudexor codex run run-126c36cb8dcc (commits could not be created inside the run sandbox, so the delivered patch is committed here as one phase commit): - P3-1: keep the cache-friendly prompt prefix stable across rounds (loop_llm_call.py, net 0 lines at the 1600 cap); golden payload untouched. - P3-2: normalize the all-text send copy for the subscription transport and preserve non-text parts (llm_claudexor.py). - P3-3: record requested/applied options and options_honored as durable evidence; the owner mismatch notice is NOT emitted (no lawful owner-message seam beside the _model_route consumer; reported as BLOCKED, see plan §8.4 A16). - P3-4: the fallback notice names the account, the typed reason and the pin disclosure and carries a model_lane_switch incident (loop_model_call.py). - P3-5: drop the unsupported "after retries and same-model reroute" clause (loop_transport.py); the unknown-outcome branch keeps provider_recovery_hint. - P3-6: observability failures warn non-fatally; both retained execution-drive roots are counted once (agent_startup_checks.py, context_budget.py). - P3-7: same-execution, same-route refusal evidence suppresses exactly the next preferredProfileId preference while pins stay strict (llm_claudexor.py). Docs: ARCHITECTURE.md (data layout, send copy/options, Auto account behaviour, canonical-versus-execution continuity); DEVELOPMENT.md (hot-store enrollment, owner-visible terminal facts, account-selection contract, cache affinity, lane-switch incidents). Tests: new tests/test_review_prompt_caching_p3.py plus focused additions; 401 focused tests passed in the run, the one failure (test_safety_supervisor_prompt_text_is_unchanged_by_the_cache_block_shape) is pre-existing on the base and untouched by this phase (to be re-verified).
This commit is contained in:
parent
57634d75b3
commit
cb08adc4a3
17 changed files with 388 additions and 19 deletions
|
|
@ -563,7 +563,7 @@ Bundled resources use the §1 CLI/headless lookup order rather than assuming the
|
|||
│ ├── settings.json ← user settings (API keys, models, budget)
|
||||
│ ├── task_results/ ← durable task results (task_results/<id>.json, every write stamped `_schema_version: 1`; an unstamped, future, malformed or retired-key row is QUARANTINED with log-only visibility and keeps its id occupied rather than being re-minted); artifacts/<task_id>/ holds .artifact_manifest.json (private metadata) + artifact files; .scratch_manifest.json declares ephemeral scratch {abs_path: sha256} excluded from patch capture only while content matches
|
||||
│ ├── artifact_versions/<task_id>/ ← artifact recovery history, last 5 versions per name
|
||||
│ ├── task_drives/<task_id>/ ← task-scoped scratch; startup prunes terminal tasks after the headless retention window
|
||||
│ ├── task_drives/<task_id>/ ← task-scoped scratch, including live per-call manifests; startup prunes terminal tasks after the headless retention window
|
||||
│ ├── task_trees/<root>/blackboard.jsonl ← append-only swarm blackboard + beacons; tree-scoped and ephemeral (task_tree_ledger.py), pruned on root terminal
|
||||
│ ├── state/
|
||||
│ │ ├── state.json ← runtime state + compatibility cost projection; never the monetary authority
|
||||
|
|
@ -599,7 +599,7 @@ Bundled resources use the §1 CLI/headless lookup order rather than assuming the
|
|||
│ │ ├── review_continuations/ ← durable blocked-review continuations (+ corrupt/ quarantine; archived/ holds settled un-resumed rows ≥7 days, never deleted)
|
||||
│ │ ├── workspace_executor_processes/ ← durable local/docker executor cleanup records
|
||||
│ │ ├── consciousness_observations.jsonl ← append-only inbox; rows retained until a settled successful cycle appends an ACK; malformed rows stay visible as source gaps
|
||||
│ │ ├── headless_tasks/<task_id>/data ← forked/empty child memory drives (CLI / Headless Boundary above)
|
||||
│ │ ├── headless_tasks/<task_id>/data ← forked/empty child execution drives whose live per-call manifests are promoted at terminal (CLI / Headless Boundary above)
|
||||
│ │ ├── pycache/ ← embedded-interpreter bytecode (packaged builds; CLI / Headless Boundary above)
|
||||
│ │ ├── python-userbase/ ← embedded-interpreter user installs (packaged builds)
|
||||
│ │ ├── betterleaks/ ← versioned scanner runtime + archive cache, created only by the explicit source-checkout installer
|
||||
|
|
@ -617,7 +617,7 @@ Bundled resources use the §1 CLI/headless lookup order rather than assuming the
|
|||
│ │ ├── registry.md ← memory awareness map
|
||||
│ │ └── owner_mailbox/ ← per-task user message files
|
||||
│ ├── projects/<id>/knowledge/ ← per-project facts + provenance sidecars; logs/task_reflections.jsonl holds full reflections with a bounded pointer row in the canonical log
|
||||
│ ├── observability/ ← private forensic ledger: blobs/<sha256>.json.gz compressed CAS payloads (0600) + calls/<task_id>/<call_id>.json manifests
|
||||
│ ├── observability/ ← canonical private forensic ledger: blobs/<sha256>.json.gz compressed CAS payloads (0600) + terminal-promoted calls/<task_id>/<call_id>.json manifests
|
||||
│ ├── services/<task_id>/<service>.log ← service runner logs; public tool output is bounded redacted tails + private blob refs
|
||||
│ ├── logs/
|
||||
│ │ ├── chat.jsonl ← canonical chat: one logical message stored once, projected into Main/Project lenses
|
||||
|
|
@ -1368,6 +1368,8 @@ canonical history retains it. A nonempty `refusal` is preserved as assistant
|
|||
text, including when the prior provider supplied no ordinary content. The exact
|
||||
`model_request_invalid` create refusal proves validation failed before command
|
||||
admission; a generic HTTP error or failed status read cannot prove non-dispatch.
|
||||
The send copy also normalizes all-text tool results, so moving the message-side
|
||||
cache boundary does not rewrite bytes already sent; non-text blocks remain intact.
|
||||
|
||||
`ClaudexorModelError.display_message` adds sanitized typed `vendorCode` and
|
||||
`parameter` details to the existing error event and terminal preview only after
|
||||
|
|
@ -1427,6 +1429,8 @@ never imply identical roles or accounts. Auto uses the largest advertised window
|
|||
of the exact account route, not CLI compaction thresholds. A manual value is a
|
||||
sizing assertion, not a provider unlock or a scope-review acknowledgement. Unknown
|
||||
capacity stays unknown; input size, response reservation and capacity are separate.
|
||||
Each model operation records its submitted options beside the engine's applied
|
||||
options; an absent applied-options report remains explicitly unknown.
|
||||
An account change rebinds preparation before another physical send.
|
||||
Ordinary sends, prospective wrap-up payloads and forced final replies share the
|
||||
same acting-role/account binding. Prospective subscription accounting uses the
|
||||
|
|
@ -1451,7 +1455,13 @@ helper result or verdict.
|
|||
actions use the existing mailbox and `/api/decisions` family
|
||||
`model_wait:<task>:<wait>`; task-result rows are projections, not restartable stack
|
||||
checkpoints or a second attempt ledger. Metadata polling makes no generation.
|
||||
Auto exhausts suitable same-model accounts before waiting; Pin never rotates.
|
||||
On Auto, the host may send the last successful same-route account as a preference,
|
||||
but omits that preference for the next request after a status-null or typed
|
||||
per-subject refusal on that account in the same execution. The engine remains
|
||||
the account chooser: without engine refusal evidence it may select the same
|
||||
top-headroom account again, so rotation is possible rather than guaranteed.
|
||||
Pin never rotates. The existing fallback budget remains unchanged; its defaults
|
||||
are one attempt per model and a 120-second cross-model cooldown.
|
||||
Even a configured API fallback waits for an explicit owner switch after a quota
|
||||
refusal. A temporary switch changes only the waiting role; optional persistence
|
||||
uses the ordinary settings writer. Pending acceptance, settings saved and worker
|
||||
|
|
@ -2098,7 +2108,7 @@ otherwise the view is partial and the consumer remains non-final or abstains.
|
|||
| Background Consciousness observations | `data/state/consciousness_observations.jsonl`, append-only enqueue/ACK rows owned by `BackgroundConsciousness` | Pending count/oldest metadata and a bounded recent observation rendering | `read_file(root='runtime_data', path='state/consciousness_observations.jsonl')` | Unacknowledged rows survive restart/overflow/error. Gaps block ACK and the existing direct identity rewrite; only a settled successful cycle appends ACK. |
|
||||
| Plan/review authority | Exact task-artifact/observability wave bodies, evidence selectors, reviewer route/thread receipts, and the bounded review hot index | Review status, latest wave, obligations, and compact findings; a predecessor's inherited `plan_review_state` is first projected to a compact authority core ordered around the newest wave's identity, acceptance claims, findings, and dispositions, with reviewer transport removed and `need_evidence_seen` last-priority. Every bounded collection names its total and omitted count; the projection discloses `full_chars` plus `source_ref`, and the named `include_authority` source stays complete | Exact artifact/source handle plus SHA/range/thread selectors | Missing or partial evidence is `DEGRADED`/`NOT_RUN`, never PASS. Exact artifacts remain bound to the reviewed candidate SHA; hot indexes may rotate only after the source is retained. |
|
||||
| Task acceptance (three deliveries) | The FULL host packet (`review_evidence.build_task_acceptance_evidence` under the host ladder, with its `__provenance__` table), the applied host run retained through canonical task source handles, and the paid-identity wallet ledger | The per-delivery work order: the api pack for a packet row; the FULL packet plus absolute pointers and the access disclosure for an agent-session row; the packet without its freely degradable tail plus the real data root for a native inspection row (`loop_acceptance_review.acceptance_retrieving_work_order`) | Exact `evidence_refs` from the packet's enumerable exhibit vocabulary; absolute pointers to the task's active workspace, task result record, artifact directory, verification receipts and tool-trajectory log; `review_projection.panels[].applied_source_ref` for the complete redacted applied review | Refs resolve against the FULL packet only, never the rendered projection; a session's reads are unobserved by the host (disclosed), a native episode's are `host_observed`; the immutable-core overflow refuses every delivery, a partial tool-result projection only packet rows; one strict wallet claim per panel whatever the rows' deliveries (owner R11). |
|
||||
| Canonical versus execution roots | Canonical budget/data root owns identity, authority, biography, results, and promoted observability; execution drives own tools, workspace, and transient trajectory | Project/fork/task lenses and status projections | Existing canonical-root resolver, task-result pointers, and source handles | A fork is an execution lens, not a second mind. Copy-back/promotion precedes GC for anything referenced by a canonical result; missing legacy bytes become an explicit gap. |
|
||||
| Canonical versus execution roots | Canonical budget/data root owns identity, authority, biography, results, and promoted observability; execution drives own tools, workspace, transient trajectory, and per-call manifests while a task runs | Project/fork/task lenses and status projections | Existing canonical-root resolver, task-result pointers, and source handles | A fork is an execution lens, not a second mind. Copy-back/promotion precedes GC for anything referenced by a canonical result; before terminal promotion the canonical reader cannot resolve a ref bound to a child drive (tracked as #805), and missing legacy bytes become an explicit gap. |
|
||||
|
||||
---
|
||||
|
||||
|
|
|
|||
|
|
@ -568,7 +568,10 @@ the hot-store growth health invariant
|
|||
that introduces a new append-only store read on an interactive path must
|
||||
enroll that store in the `ouroboros/context_budget.py` threshold table (with a
|
||||
justified constant) in the same commit — an unenrolled hot store is invisible
|
||||
to the tripwire.
|
||||
to the tripwire. Retained execution drives under both `state/headless_tasks`
|
||||
and `task_drives` are enrolled by direct-child count at
|
||||
`context_budget.RETAINED_EXECUTION_DRIVES_WARN_COUNT`; startup never recursively
|
||||
sizes those trees.
|
||||
|
||||
### Invariant: Source-complete decision pipeline
|
||||
|
||||
|
|
@ -951,6 +954,8 @@ what the owner reads:
|
|||
to a durable full copy (e.g. an observability `response_ref`). Reviewer
|
||||
rationale is a cognitive artifact (BIBLE P1): projecting it truncated while
|
||||
the full copy sits unreferenced in private blobs is partial memory loss.
|
||||
Terminal text asserts only recovery facts carried by the round record, never
|
||||
a route mechanism that the selected transport cannot perform.
|
||||
- **Model-bound projections** (review packs, context sections, tool-result
|
||||
transport) keep their disclosed-truncation budgets — those are real context
|
||||
economics.
|
||||
|
|
@ -2028,6 +2033,13 @@ owner, owed terminal delivery, cascade postconditions — lives in ARCHITECTURE
|
|||
saved-but-undiscovered choices stay visible and editable; a compound effort
|
||||
slug plus a conflicting separate effort is a validation error, never two
|
||||
applied efforts.
|
||||
On the Auto lane, the host may prefer the last successful same-route account.
|
||||
After a status-null or typed per-subject refusal, only the next request in the
|
||||
same execution omits that preference and lets the engine choose; without an
|
||||
engine refusal fact, selecting a sibling is possible, not guaranteed. Pin
|
||||
remains exact and never rotates. `OUROBOROS_FALLBACK_ATTEMPTS_PER_MODEL=1`
|
||||
and `OUROBOROS_FALLBACK_COOLDOWN_SEC=120` keep their existing escalation
|
||||
budget and do not turn preference suppression into a retry or cooldown.
|
||||
- Saved intent, generated drafts, and live status are different axes: a
|
||||
status/catalog failure annotates a loaded row and never erases it; GET may
|
||||
return an unsaved candidate but only explicit Save or onboarding completion
|
||||
|
|
@ -2135,7 +2147,9 @@ Focused regressions: `test_review_late_cas_recovery.py`, `test_delivery_control_
|
|||
the same operation ID; record unknown outcomes as unknown. ACK only after the
|
||||
existing private CAS owns the exact result. Optional host hints must be chosen
|
||||
by their caller according to transport capability; explicit unsupported options
|
||||
refuse, rather than being silently removed and retried.
|
||||
refuse, rather than being silently removed and retried. Record submitted model
|
||||
options beside the engine's applied options on the usage row; an absent report
|
||||
stays unknown and a mismatch is disclosure, never a dispatch gate.
|
||||
- The engine's active-turn token is one of those transport facts, so the CALLER
|
||||
owns its slot (`llm_claudexor.ModelTurnState` on the loop context, a wake-scoped
|
||||
one in Background Consciousness) and the engine boundary is its only writer.
|
||||
|
|
@ -2262,7 +2276,9 @@ by "Provider Independence" above. Call-site imperatives:
|
|||
qualifies, migrated in the same call -- preserved on the
|
||||
direct-Anthropic lane by `_anthropic_blocks_from_content` and on
|
||||
OpenRouter by `supports_message_cache_control`, and pinned by
|
||||
`tests/test_review_prompt_caching.py`;
|
||||
`tests/test_review_prompt_caching.py`. The main loop declares an
|
||||
execution-scoped cache affinity only for subscription transport; API-compatible
|
||||
lanes retain their prefix-derived session identity;
|
||||
`review_substrate.assert_cache_breakpoint_cap` covers only the review
|
||||
builders. Review gate: CHECKLISTS item 22 (`cache_friendliness`).
|
||||
- Provider fallback is disabled only when the transcript carries a SEALED
|
||||
|
|
@ -2389,7 +2405,8 @@ by "Provider Independence" above. Call-site imperatives:
|
|||
waiting tunes the passive wait only.
|
||||
- Preserve raw terminal model/salvage bytes separately from the host-authored
|
||||
`terminal_provider_notice`. Existing receipts and secondary notices consume
|
||||
those same facts; a retained answer must not hide wait or unknown-attempt
|
||||
those same facts: attempted repeats, the last provider error, and an unknown
|
||||
dispatched outcome. A retained answer must not hide wait or unknown-attempt
|
||||
evidence or invite a blind rerun. Ephemeral and message/deferred Presence
|
||||
responses render one host-labelled status section; cached Presence output
|
||||
is already rendered. Preserve silent/tool-delivered authority and never
|
||||
|
|
@ -2405,7 +2422,9 @@ by "Provider Independence" above. Call-site imperatives:
|
|||
reads the tone through `normalizeTone` and keeps the alarm tone for a
|
||||
frame without one — never parse `toast_once` or the text for it; `OuroborosAgent._emit_progress` is the
|
||||
production implementation and a test fake mirrors it
|
||||
(`lambda text, *, incident=None: ...`).
|
||||
(`lambda text, *, incident=None: ...`). A cross-model lane switch is the
|
||||
second owner note carrying this pair; it names both models, the selected
|
||||
account, and the typed failure reason when the round record has one.
|
||||
- Timeout contract classes differ; keep the axes separate. A transport
|
||||
timeout only bounds a dead socket
|
||||
(`OUROBOROS_LLM_TRANSPORT_READ_TIMEOUT_SEC`) — it is not a reasoning cutoff
|
||||
|
|
|
|||
|
|
@ -860,6 +860,24 @@ def hot_store_growth_notes(env: Any) -> list:
|
|||
"replay scans this chain on ownership questions. Investigate chain "
|
||||
"indexing/compaction; archives are durable history and are never deleted."
|
||||
)
|
||||
from ouroboros.context_budget import RETAINED_EXECUTION_DRIVES_WARN_COUNT
|
||||
from ouroboros.headless import HEADLESS_TASKS_DIR, TASK_DRIVES_DIR
|
||||
retained_drive_count = 0
|
||||
for retained_root in (
|
||||
drive_root / HEADLESS_TASKS_DIR,
|
||||
drive_root / TASK_DRIVES_DIR,
|
||||
):
|
||||
try:
|
||||
retained_drive_count += sum(path.is_dir() for path in retained_root.iterdir())
|
||||
except OSError:
|
||||
pass
|
||||
if retained_drive_count > RETAINED_EXECUTION_DRIVES_WARN_COUNT:
|
||||
notes.append(
|
||||
"WARNING: HOT STORE GROWTH — retained execution drives under "
|
||||
f"state/headless_tasks and task_drives total {retained_drive_count} "
|
||||
f"(threshold {RETAINED_EXECUTION_DRIVES_WARN_COUNT}). Terminal-task retention "
|
||||
"or pruning is lagging; inspect lifecycle GC without recursively sizing drives."
|
||||
)
|
||||
return notes
|
||||
|
||||
|
||||
|
|
|
|||
|
|
@ -283,6 +283,9 @@ CHAT_ARCHIVE_SCAN_WARN_BYTES = 100_000_000
|
|||
# archives stay durable history (never GC'd), so the remediation is chain
|
||||
# indexing/compaction, never deletion.
|
||||
EVENTS_ARCHIVE_SCAN_WARN_BYTES = 100_000_000
|
||||
# Warn before the observed 242-of-253 retained-drive corpus becomes routine;
|
||||
# count only direct children because startup health is an interactive path.
|
||||
RETAINED_EXECUTION_DRIVES_WARN_COUNT = 200
|
||||
|
||||
|
||||
def estimate_message_chars(messages: Any) -> int:
|
||||
|
|
|
|||
|
|
@ -40,6 +40,7 @@ from __future__ import annotations
|
|||
|
||||
import asyncio
|
||||
import copy
|
||||
import contextvars
|
||||
from dataclasses import replace
|
||||
import json
|
||||
import logging
|
||||
|
|
@ -48,6 +49,7 @@ import time
|
|||
from typing import Any
|
||||
|
||||
from ouroboros import config
|
||||
from ouroboros import context_fit
|
||||
from ouroboros._usage_response import provider_cost_value
|
||||
from ouroboros.anthropic_native_custody import scrub_native_custody
|
||||
from ouroboros.claudexor_daemon import ensure_owned_gateway, owned_engine_version, read_owned_gateway
|
||||
|
|
@ -66,6 +68,14 @@ from ouroboros.usage_accounting import (
|
|||
from ouroboros.utils import append_jsonl, sanitize_tool_result_for_log, utc_now_iso
|
||||
|
||||
log = logging.getLogger(__name__)
|
||||
_FAILED_PROFILE = contextvars.ContextVar("claudexor_failed_profile", default=())
|
||||
_PER_SUBJECT_REFUSALS = frozenset({
|
||||
"auth_required", "auth_refresh_failed", "credential_unusable", "provider_refused",
|
||||
"rate_limited", "subscription_window_exhausted",
|
||||
})
|
||||
_NON_PROVIDER_FAILURES = frozenset({
|
||||
"model_operation_cancelled", "model_operation_interrupted", "model_outcome_unknown",
|
||||
})
|
||||
|
||||
|
||||
def model_catalog(source: str, credential_profile_id: str | None = None, *,
|
||||
|
|
@ -229,6 +239,16 @@ def adopt_turn_state(slot: ModelTurnState | None, payload: dict, result: dict) -
|
|||
slot.envelope = copy.deepcopy(envelope) if isinstance(envelope, dict) else None
|
||||
|
||||
|
||||
def _remember_failed_profile(target: dict, parameters: dict, error: ClaudexorModelError) -> None:
|
||||
route = error.route or {}
|
||||
key = (parameters.get("cache_affinity"), route.get("source"), route.get("model"))
|
||||
if (key == (parameters.get("cache_affinity"), target["source"], target["resolved_model"])
|
||||
and key[0] and route.get("credentialProfileId")
|
||||
and ((error.status_code == 0 and error.code not in _NON_PROVIDER_FAILURES)
|
||||
or error.code in _PER_SUBJECT_REFUSALS)):
|
||||
_FAILED_PROFILE.set((*key, route["credentialProfileId"]))
|
||||
|
||||
|
||||
def _request(target: dict, messages: list, tools: list | None, parameters: dict) -> dict:
|
||||
from ouroboros.llm_messages import _MessageShapingMixin
|
||||
|
||||
|
|
@ -260,12 +280,22 @@ def _request(target: dict, messages: list, tools: list | None, parameters: dict)
|
|||
if isinstance(block, dict):
|
||||
for name in ("_caption", "_source_path", "_context_capsule", "cache_control"):
|
||||
block.pop(name, None)
|
||||
if message.get("role") == "tool" and isinstance(content, list) and content and all(
|
||||
isinstance(block, dict) and block.get("type") == "text" for block in content
|
||||
):
|
||||
message["content"] = context_fit.extract_plain_text_from_content(content)
|
||||
role = parameters.get("model_role", "")
|
||||
override = parameters.get("model_account_override")
|
||||
if override is not None and not isinstance(override, str):
|
||||
raise ValueError("model_account_override must be a profile name, empty Auto, or None")
|
||||
pin = override.strip() if override is not None else model_role_option(MODEL_ACCOUNTS_KEY, role)
|
||||
account = {"mode": "pin", "profileId": pin} if pin else {"mode": "auto"}
|
||||
failed = _FAILED_PROFILE.get()
|
||||
failed_key = (parameters.get("cache_affinity"), target["source"], target["resolved_model"])
|
||||
same_execution = len(failed) == 4 and failed[0] == failed_key[0]
|
||||
failed_profile = failed[3] if same_execution and failed[:3] == failed_key else ""
|
||||
if same_execution:
|
||||
_FAILED_PROFILE.set(()) # one request only; Pin still consumes the failure fact
|
||||
if not pin:
|
||||
# Carry the conversation's last account as a preference, not admission.
|
||||
# The engine is still the only actor choosing an eligible account.
|
||||
|
|
@ -273,7 +303,7 @@ def _request(target: dict, messages: list, tools: list | None, parameters: dict)
|
|||
native = message.get("nativeContinuation") or {}
|
||||
route = native.get("route") or {}
|
||||
if route.get("source") == target["source"] and route.get("model") == target["resolved_model"]:
|
||||
if route.get("credentialProfileId"):
|
||||
if route.get("credentialProfileId") and route["credentialProfileId"] != failed_profile:
|
||||
account["preferredProfileId"] = route["credentialProfileId"]
|
||||
break
|
||||
options = {wire: parameters[key] for key, wire in (
|
||||
|
|
@ -522,12 +552,17 @@ class _ModelInvocation:
|
|||
def finish(self, result: dict) -> tuple[dict, dict]:
|
||||
usage, cost, final = _usage(result)
|
||||
route = result.get("route") or {}
|
||||
requested_options = copy.deepcopy(self.payload.get("options") or {})
|
||||
applied_options = copy.deepcopy(result.get("appliedOptions"))
|
||||
options_honored = "unknown" if applied_options is None else (
|
||||
"mismatch" if any(applied_options[key] != value for key, value in requested_options.items() if key in applied_options) else "confirmed")
|
||||
usage.update(provider="claudexor", resolved_model=self.target["usage_model"], cost=cost, cost_final=final,
|
||||
cost_estimated=cost is not None and not final,
|
||||
claudexor={"operation_id": self.operation_id, "model_role": self.role,
|
||||
"route": copy.deepcopy(route), "cost_evidence": copy.deepcopy(result.get("cost")),
|
||||
"outcome": result.get("outcome"), "problem": copy.deepcopy(result.get("problem")),
|
||||
"applied_options": copy.deepcopy(result.get("appliedOptions")),
|
||||
"requested_options": requested_options, "applied_options": applied_options,
|
||||
"options_honored": options_honored,
|
||||
"output_reserve_tokens": self.output_reserve, "output_cap_applied": False,
|
||||
"result_custody": {"state": "pending", "operation_id": self.operation_id,
|
||||
"response_ref": self.response_ref,
|
||||
|
|
@ -682,6 +717,9 @@ def chat_claudexor(target: dict, messages: list, tools: list | None, **parameter
|
|||
raise
|
||||
payload = updated
|
||||
retry_preparation = _native_retry_preparation(target, payload, parameters, error)
|
||||
except ClaudexorModelError as error:
|
||||
_remember_failed_profile(target, parameters, error)
|
||||
raise
|
||||
except PhysicalAttemptPreparationFailed as error:
|
||||
cause = error.__cause__
|
||||
if isinstance(cause, ClaudexorModelError):
|
||||
|
|
@ -725,6 +763,9 @@ async def chat_claudexor_async(target: dict, messages: list, tools: list | None,
|
|||
raise
|
||||
payload = updated
|
||||
retry_preparation = _native_retry_preparation(target, payload, parameters, error)
|
||||
except ClaudexorModelError as error:
|
||||
_remember_failed_profile(target, parameters, error)
|
||||
raise
|
||||
except PhysicalAttemptPreparationFailed as error:
|
||||
cause = error.__cause__
|
||||
if isinstance(cause, (ClaudexorModelError, asyncio.CancelledError)):
|
||||
|
|
|
|||
|
|
@ -21,7 +21,7 @@ def persist_observed_call(root: Any, *, payload: Any, writer: Any = None, **iden
|
|||
try:
|
||||
return (writer or persist_call)(root, payload=public_custody_projection(payload), **identity)
|
||||
except Exception:
|
||||
logging.getLogger(__name__).debug("Failed to persist LLM observability payload", exc_info=True)
|
||||
logging.getLogger(__name__).warning("Failed to persist LLM observability payload", exc_info=True)
|
||||
return {}
|
||||
|
||||
|
||||
|
|
|
|||
|
|
@ -1379,7 +1379,7 @@ def call_llm_with_retry(
|
|||
"max_tokens": MAIN_LOOP_MAX_TOKENS,
|
||||
"stream": True, "caller_deadline_ts": (None if deadline_ts is None
|
||||
else float(deadline_ts) - float(transport_reserve_sec or 0.0)),
|
||||
"use_local": use_local,
|
||||
"use_local": use_local, "cache_affinity": execution_id if provider_for_model(model) == "claudexor" else "",
|
||||
# These are optional host hints, not required tools. This
|
||||
# transport has neither provider-owned web tools nor a bypass
|
||||
# knob; ordinary Ouroboros web tools stay in the schema.
|
||||
|
|
|
|||
|
|
@ -99,7 +99,7 @@ def _run_cross_model_fallback_chain(
|
|||
"""Try fallbacks; unknown dispatch stops the chain."""
|
||||
from ouroboros import fallback_cooldown as _fcd
|
||||
from ouroboros.config import fallback_candidate_targets
|
||||
from ouroboros.model_slots import parse_fallback_chain
|
||||
from ouroboros.model_slots import MODEL_ACCOUNTS_KEY, model_role_option, parse_fallback_chain
|
||||
from ouroboros.loop_llm_call import _COOLDOWN_ERROR_KINDS as _cooldown_kinds
|
||||
|
||||
def _cooled(model: str, use_local: bool) -> None:
|
||||
|
|
@ -130,7 +130,11 @@ def _run_cross_model_fallback_chain(
|
|||
break
|
||||
ptag = " (local)" if active_use_local else ""
|
||||
ftag = " (local)" if fallback_use_local else ""
|
||||
emit_progress(f"⚡ Fallback: {active_model}{ptag} → {fallback_model}{ftag}")
|
||||
fallback_account = str(model_role_option(MODEL_ACCOUNTS_KEY, fallback_role) or "")
|
||||
reason = str(accumulated_usage.get("_last_llm_error_kind") or "")
|
||||
emit_progress(f"⚡ Fallback: {active_model}{ptag} → {fallback_model}{ftag}; account: {fallback_account or 'Auto'}"
|
||||
f"{f'; reason: {reason}' if reason else ''}{'; pinned account: siblings were not tried' if fallback_account else ''}",
|
||||
incident={"task_incident": "model_lane_switch", "toast_once": f"{task_id}:model_lane_switch:{round_idx}:{fallback_model}"})
|
||||
# Cross-FAMILY fallback must not replay the primary's
|
||||
# provider-private reasoning to a different family (the GLM->Claude
|
||||
# 400 "Invalid signature" death); the SSOT sanitizer no-ops same-family.
|
||||
|
|
|
|||
|
|
@ -762,7 +762,7 @@ def provider_terminal_fallback_text(
|
|||
text += provider_recovery_hint(accumulated_usage)
|
||||
return text
|
||||
return (
|
||||
"⚠️ The model provider returned no usable response after retries and same-model reroute."
|
||||
"⚠️ The model provider returned no usable response."
|
||||
f"{provider_failure_hint(accumulated_usage)}{provider_recovery_hint(accumulated_usage)} "
|
||||
"Any files written so far are preserved in the workspace."
|
||||
)
|
||||
|
|
|
|||
|
|
@ -6,6 +6,12 @@ import pytest
|
|||
import yaml
|
||||
|
||||
from ouroboros.gateways.claudexor import final_attempt_facts
|
||||
from ouroboros.llm_claudexor import (
|
||||
ClaudexorModelError,
|
||||
_ModelInvocation,
|
||||
_remember_failed_profile,
|
||||
_request,
|
||||
)
|
||||
|
||||
|
||||
def _write_telemetry(tmp_path, attempts, *, final_id="a02", run_id="run-fixture"):
|
||||
|
|
@ -108,3 +114,81 @@ def test_missing_engine_run_directory_never_reads_the_working_directory(detail,
|
|||
|
||||
monkeypatch.setattr(Path, "read_text", unexpected_read)
|
||||
assert final_attempt_facts(detail, "run-fixture") == {}
|
||||
|
||||
|
||||
@pytest.mark.parametrize("applied,expected", [
|
||||
({"reasoningEffort": "xhigh"}, "confirmed"),
|
||||
({"reasoningEffort": "medium"}, "mismatch"),
|
||||
(None, "unknown"),
|
||||
])
|
||||
def test_model_invocation_records_requested_and_applied_options(applied, expected):
|
||||
requested = {"reasoningEffort": "xhigh"}
|
||||
invocation = _ModelInvocation(
|
||||
{"usage_model": "claudexor::codex=model"}, {"options": requested}, {}
|
||||
)
|
||||
result = {
|
||||
"outcome": "completed",
|
||||
"message": {"role": "assistant", "content": "done"},
|
||||
}
|
||||
if applied is not None:
|
||||
result["appliedOptions"] = applied
|
||||
|
||||
_message, usage = invocation.finish(result)
|
||||
observed = usage["claudexor"]
|
||||
assert observed["requested_options"] == requested
|
||||
assert observed["applied_options"] == applied
|
||||
assert observed["options_honored"] == expected
|
||||
|
||||
|
||||
def _continuation(profile="profile-a"):
|
||||
return [{
|
||||
"role": "assistant",
|
||||
"content": "prior answer",
|
||||
"nativeContinuation": {"route": {
|
||||
"source": "codex", "model": "gpt-6", "credentialProfileId": profile,
|
||||
}},
|
||||
}]
|
||||
|
||||
|
||||
def test_successful_profile_remains_an_auto_lane_preference():
|
||||
target = {"source": "codex", "resolved_model": "gpt-6"}
|
||||
|
||||
payload = _request(target, _continuation(), None, {"cache_affinity": "execution-success"})
|
||||
|
||||
assert payload["account"] == {"mode": "auto", "preferredProfileId": "profile-a"}
|
||||
|
||||
|
||||
def test_status_null_failure_suppresses_only_the_next_same_route_preference():
|
||||
target = {"source": "codex", "resolved_model": "gpt-6"}
|
||||
parameters = {"cache_affinity": "execution-failed"}
|
||||
error = ClaudexorModelError(
|
||||
{"code": "server_error", "message": "stream ended"},
|
||||
route={"source": "codex", "model": "gpt-6", "credentialProfileId": "profile-a"},
|
||||
)
|
||||
_remember_failed_profile(target, parameters, error)
|
||||
|
||||
assert _request(target, _continuation(), None, parameters)["account"] == {"mode": "auto"}
|
||||
assert _request(target, _continuation(), None, parameters)["account"] == {
|
||||
"mode": "auto", "preferredProfileId": "profile-a",
|
||||
}
|
||||
|
||||
|
||||
def test_failure_fact_does_not_change_pin_or_single_account_auto_mode():
|
||||
target = {"source": "codex", "resolved_model": "gpt-6"}
|
||||
parameters = {"cache_affinity": "execution-pin"}
|
||||
error = ClaudexorModelError(
|
||||
{"code": "subscription_window_exhausted", "message": "window spent",
|
||||
"context": {"httpStatus": 429}},
|
||||
route={"source": "codex", "model": "gpt-6", "credentialProfileId": "only-profile"},
|
||||
)
|
||||
_remember_failed_profile(target, parameters, error)
|
||||
|
||||
pinned = _request(target, _continuation("only-profile"), None, {
|
||||
**parameters, "model_account_override": "only-profile",
|
||||
})
|
||||
assert pinned["account"] == {"mode": "pin", "profileId": "only-profile"}
|
||||
|
||||
_remember_failed_profile(target, parameters, error)
|
||||
assert _request(target, _continuation("only-profile"), None, parameters)["account"] == {
|
||||
"mode": "auto",
|
||||
}
|
||||
|
|
|
|||
|
|
@ -666,3 +666,11 @@ def test_affordance_map_carries_label_path_pairs(tmp_path):
|
|||
# The read-only orchestrator root is resolvable when visible to the profile.
|
||||
if "subagent_projects" in result.get("readonly_roots", []):
|
||||
assert paths.get("subagent_projects")
|
||||
|
||||
|
||||
def test_provider_terminal_notice_never_claims_unsupported_reroute():
|
||||
import inspect
|
||||
from ouroboros.loop_transport import provider_terminal_fallback_text
|
||||
|
||||
source = inspect.getsource(provider_terminal_fallback_text)
|
||||
assert "same-model reroute" not in source
|
||||
|
|
|
|||
|
|
@ -20,6 +20,7 @@ from __future__ import annotations
|
|||
|
||||
import asyncio
|
||||
import json
|
||||
import logging
|
||||
import types
|
||||
|
||||
from starlette.requests import Request
|
||||
|
|
@ -28,6 +29,58 @@ from ouroboros import usage_accounting as ua
|
|||
from ouroboros import usage_ledger
|
||||
|
||||
|
||||
def test_retained_execution_drive_tripwire_counts_both_roots(tmp_path, monkeypatch):
|
||||
from ouroboros import context_budget
|
||||
from ouroboros.agent_startup_checks import hot_store_growth_notes
|
||||
from ouroboros.headless import HEADLESS_TASKS_DIR, TASK_DRIVES_DIR
|
||||
from supervisor.state import ISOLATED_BENCHMARK_SENTINEL
|
||||
|
||||
monkeypatch.setattr(context_budget, "RETAINED_EXECUTION_DRIVES_WARN_COUNT", 2)
|
||||
env = types.SimpleNamespace(
|
||||
drive_root=tmp_path,
|
||||
drive_path=lambda rel: tmp_path / rel,
|
||||
)
|
||||
headless = tmp_path / HEADLESS_TASKS_DIR
|
||||
task_drives = tmp_path / TASK_DRIVES_DIR
|
||||
(headless / "headless-1").mkdir(parents=True)
|
||||
(task_drives / "drive-1").mkdir(parents=True)
|
||||
|
||||
assert hot_store_growth_notes(env) == []
|
||||
|
||||
(task_drives / "drive-2").mkdir()
|
||||
notes = hot_store_growth_notes(env)
|
||||
assert len(notes) == 1
|
||||
assert "retained execution drives" in notes[0]
|
||||
assert "total 3 (threshold 2)" in notes[0]
|
||||
|
||||
(tmp_path / ISOLATED_BENCHMARK_SENTINEL).write_text("isolated\n", encoding="utf-8")
|
||||
assert hot_store_growth_notes(env) == []
|
||||
|
||||
|
||||
def test_observability_write_failure_warns_once_and_stays_nonfatal(tmp_path, caplog):
|
||||
from ouroboros.llm_observability import persist_observed_call
|
||||
|
||||
def fail_write(*args, **kwargs):
|
||||
raise OSError("test write failure")
|
||||
|
||||
with caplog.at_level(logging.WARNING, logger="ouroboros.llm_observability"):
|
||||
result = persist_observed_call(
|
||||
tmp_path,
|
||||
payload={"model": "claudexor/test"},
|
||||
writer=fail_write,
|
||||
task_id="task-1",
|
||||
)
|
||||
|
||||
assert result == {}
|
||||
records = [
|
||||
record for record in caplog.records
|
||||
if record.name == "ouroboros.llm_observability"
|
||||
]
|
||||
assert len(records) == 1
|
||||
assert records[0].levelno == logging.WARNING
|
||||
assert records[0].getMessage() == "Failed to persist LLM observability payload"
|
||||
|
||||
|
||||
def _seeded_accounting_root(tmp_path, monkeypatch):
|
||||
"""Temp drive root with a settled + reserved attempt in the usage ledger."""
|
||||
root = tmp_path / "data"
|
||||
|
|
|
|||
|
|
@ -160,6 +160,24 @@ def test_deadline_text_does_not_hide_an_existing_unknown_attempt():
|
|||
assert "no retry or paid fallback was sent" in text
|
||||
|
||||
|
||||
def test_provider_terminal_text_claims_only_recorded_recovery_facts():
|
||||
ordinary = loop_transport.provider_terminal_fallback_text(
|
||||
{"_last_llm_error_kind": "provider_transient", "_last_llm_error": "HTTP 503"},
|
||||
is_context_overflow=False, is_transport_wait=False, waited_sec=0.0,
|
||||
interactive=False, is_deadline_exhausted=False,
|
||||
)
|
||||
assert "provider returned no usable response" in ordinary
|
||||
assert "same-model reroute" not in ordinary
|
||||
|
||||
unknown = loop_transport.provider_terminal_fallback_text(
|
||||
{"_last_llm_error_kind": "provider_outcome_unknown"},
|
||||
is_context_overflow=False, is_transport_wait=False, waited_sec=0.0,
|
||||
interactive=False, is_deadline_exhausted=False,
|
||||
)
|
||||
assert unknown.count("dispatched request has no terminal provider outcome") == 1
|
||||
assert "same-model reroute" not in unknown
|
||||
|
||||
|
||||
def test_body_error_diagnostic_is_masked_before_terminal_publication(tmp_path, monkeypatch):
|
||||
from ouroboros.utils import sanitize_tool_result_for_log
|
||||
from tests.test_transport_death_retry import _ScriptedLLM, _death, _primary_call
|
||||
|
|
|
|||
|
|
@ -130,11 +130,46 @@ def test_fallback_dispatch_lane_stays_the_global_flag(tmp_path, monkeypatch, cap
|
|||
messages=[], active_model="primary", active_use_local=False, tool_schemas=[],
|
||||
active_effort="high", max_retries=1, drive_logs=tmp_path / "logs", task_id="t",
|
||||
round_idx=1, event_queue=None, accumulated_usage={}, task_type="task",
|
||||
emit_progress=lambda _: None, context_fit_plan=None, active_context_mode="max",
|
||||
emit_progress=lambda _text, *, incident=None: None,
|
||||
context_fit_plan=None, active_context_mode="max",
|
||||
)
|
||||
assert dispatched == [("remote-model", captured_local), ("other (local)", captured_local)]
|
||||
|
||||
|
||||
def test_fallback_notice_carries_lane_switch_incident_reason_and_pin(tmp_path, monkeypatch):
|
||||
from types import SimpleNamespace
|
||||
from ouroboros import fallback_cooldown, loop, loop_model_call
|
||||
|
||||
fallback = "claudexor::codex=fallback"
|
||||
monkeypatch.setenv("OUROBOROS_MODEL_FALLBACKS", fallback)
|
||||
monkeypatch.setenv("OUROBOROS_MODEL_ACCOUNTS", '{"fallback":["account-a"]}')
|
||||
monkeypatch.setattr(fallback_cooldown, "is_cooling_down", lambda *_: False)
|
||||
monkeypatch.setattr(loop, "_task_deadline_epoch", lambda _: None)
|
||||
monkeypatch.setattr(loop, "_rebind_context_fit_plan", lambda *a, **k: (None, "max"))
|
||||
monkeypatch.setattr(loop, "_call_round_model", lambda _ctx: ({"role": "assistant"}, 0, "max"))
|
||||
progress = []
|
||||
ctx = SimpleNamespace(active_model="primary", active_use_local=False)
|
||||
tools = SimpleNamespace(_ctx=SimpleNamespace())
|
||||
|
||||
loop_model_call._run_cross_model_fallback_chain(
|
||||
llm=None, ctx=ctx, tools=tools, messages=[], active_model="primary",
|
||||
active_use_local=False, tool_schemas=[], active_effort="high", max_retries=1,
|
||||
drive_logs=tmp_path / "logs", task_id="task-7", round_idx=3, event_queue=None,
|
||||
accumulated_usage={"_last_llm_error_kind": "provider_transient"}, task_type="task",
|
||||
emit_progress=lambda text, *, incident=None: progress.append((text, incident)),
|
||||
context_fit_plan=None, active_context_mode="max",
|
||||
)
|
||||
|
||||
assert len(progress) == 1
|
||||
text, incident = progress[0]
|
||||
assert "account: account-a" in text and "reason: provider_transient" in text
|
||||
assert "pinned account: siblings were not tried" in text
|
||||
assert incident == {
|
||||
"task_incident": "model_lane_switch",
|
||||
"toast_once": f"task-7:model_lane_switch:3:{fallback}",
|
||||
}
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# Seam 2: the reviewer model lists (review_model_routes / reviewer slots).
|
||||
# ---------------------------------------------------------------------------
|
||||
|
|
|
|||
41
tests/test_review_prompt_caching_p3.py
Normal file
41
tests/test_review_prompt_caching_p3.py
Normal file
|
|
@ -0,0 +1,41 @@
|
|||
"""Main-loop subscription cache-affinity regressions."""
|
||||
|
||||
import queue
|
||||
|
||||
from ouroboros.llm_claudexor import _request
|
||||
from ouroboros.loop_llm_call import call_llm_with_retry
|
||||
|
||||
|
||||
def test_claudexor_request_projects_explicit_affinity_to_cache_key():
|
||||
payload = _request(
|
||||
{"source": "codex", "resolved_model": "model"},
|
||||
[{"role": "user", "content": "work"}],
|
||||
None,
|
||||
{"cache_affinity": "execution-7", "model_account_override": ""},
|
||||
)
|
||||
assert payload["options"]["cacheKey"] == "execution-7"
|
||||
|
||||
|
||||
def test_main_loop_scopes_execution_affinity_to_claudexor(tmp_path):
|
||||
captured = []
|
||||
|
||||
class LLM:
|
||||
def chat(self, **kwargs):
|
||||
captured.append(kwargs)
|
||||
return ({"content": "done", "tool_calls": [], "finish_reason": "stop"}, {
|
||||
"provider": "fixture", "resolved_model": kwargs["model"], "cost": 0.0,
|
||||
"prompt_tokens": 1, "completion_tokens": 1,
|
||||
})
|
||||
|
||||
logs = tmp_path / "logs"
|
||||
logs.mkdir()
|
||||
for model in ("claudexor::codex=model", "openrouter::openai/model"):
|
||||
usage = {"execution_id": "execution-7"}
|
||||
message, _cost = call_llm_with_retry(
|
||||
LLM(), [{"role": "user", "content": "work"}], model, None, "medium", 1,
|
||||
logs, "task", 1, queue.Queue(), usage,
|
||||
)
|
||||
assert message["content"] == "done"
|
||||
|
||||
assert captured[0]["cache_affinity"] == "execution-7"
|
||||
assert captured[1]["cache_affinity"] == ""
|
||||
|
|
@ -239,7 +239,8 @@ def test_original_fallback_ordinal_keeps_account_after_filters(tmp_path, monkeyp
|
|||
llm=ctx.llm, ctx=ctx.tools._ctx, tools=ctx.tools, messages=ctx.messages,
|
||||
active_model=ctx.active_model, active_use_local=False, tool_schemas=[], active_effort="medium",
|
||||
max_retries=1, drive_logs=ctx.drive_logs, task_id=ctx.task_id, round_idx=1,
|
||||
event_queue=None, accumulated_usage={}, task_type="task", emit_progress=lambda _text: None,
|
||||
event_queue=None, accumulated_usage={}, task_type="task",
|
||||
emit_progress=lambda _text, *, incident=None: None,
|
||||
context_fit_plan=ctx.context_fit_plan, active_context_mode="max")
|
||||
assert seen == ["fallback:1", "fallback:3"]
|
||||
|
||||
|
|
|
|||
|
|
@ -218,6 +218,40 @@ def test_anthropic_messages_pass_through_list_content_for_tool_result():
|
|||
assert tool_result_block["content"] == sealed_content
|
||||
|
||||
|
||||
def test_claudexor_send_copy_is_stable_when_rolling_seal_moves():
|
||||
"""The subscription projection normalizes both sides of a seal move."""
|
||||
from ouroboros.llm_claudexor import _request
|
||||
|
||||
target = {"source": "codex", "resolved_model": "model"}
|
||||
parameters = {"model_account_override": ""}
|
||||
messages = _make_messages(n_tool_rounds=7, prefix_per_tool=3000)
|
||||
seal_task_transcript(messages, keep_active=5, min_prefix_tokens=100)
|
||||
first = _request(target, messages, None, parameters)["messages"]
|
||||
|
||||
messages.append(_assistant_msg("new round"))
|
||||
messages.append(_tool_msg("new_result" * 500, call_id="tc_new"))
|
||||
seal_task_transcript(messages, keep_active=5, min_prefix_tokens=100)
|
||||
second = _request(target, messages, None, parameters)["messages"]
|
||||
|
||||
assert second[:len(first)] == first
|
||||
|
||||
|
||||
def test_claudexor_send_copy_preserves_nontext_tool_blocks():
|
||||
from ouroboros.llm_claudexor import _request
|
||||
|
||||
content = [
|
||||
{"type": "text", "text": "caption"},
|
||||
{"type": "image_url", "image_url": {"url": "data:image/png;base64,AA=="}},
|
||||
]
|
||||
payload = _request(
|
||||
{"source": "codex", "resolved_model": "model"},
|
||||
[{"role": "tool", "tool_call_id": "tc-image", "content": content}],
|
||||
None,
|
||||
{"model_account_override": ""},
|
||||
)
|
||||
assert payload["messages"][0]["content"] == content
|
||||
|
||||
|
||||
def test_compaction_receives_plain_strings():
|
||||
"""After seal + revert cycle, all tool messages are plain strings (safe for compaction)."""
|
||||
msgs = _make_messages(n_tool_rounds=8, prefix_per_tool=3000)
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue