Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
3 changes: 3 additions & 0 deletions CHANGELOG.d/observed-health-authority.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,3 @@
- Removed the uncalibrated observed-health model-admission switch; explicit or
KV activation now has no supported runtime surface. Provider HTTP 429 still
records quota cooldown evidence without charging breaker health.
21 changes: 18 additions & 3 deletions contextual_orchestrator/orchestrator.py
Original file line number Diff line number Diff line change
Expand Up @@ -8264,7 +8264,17 @@ def stream_route(
raise
last_error = upstream
decision = classify_provider_transport_failure(upstream.retryable)
if decision.circuit_failure:
rate_limit_signal = self._rate_limited_provider_signal(upstream)
if rate_limit_signal is not None:
signal_status, signal_http_error = rate_limit_signal
self._record_rate_limit(
agent.id,
resolve_retry_after_seconds(signal_http_error)
if signal_http_error is not None
else upstream.extra_detail.get("retry_after_seconds"),
status=signal_status,
)
if decision.circuit_failure and self._charges_breaker(upstream):
self._record_failure(agent.id)
if decision.action is ToolFallbackAction.FAIL_CLOSED:
raise upstream from None
Expand Down Expand Up @@ -11926,7 +11936,7 @@ def attempt(agent: ModelAgent) -> EndpointAttempt[Any]:
):
retry_attempt += 1
self._record_tool_fallback(agent.id, decision, retry_attempt)
if decision.circuit_failure and not quota_rejection:
if decision.circuit_failure and self._charges_breaker(exc):
self._record_failure(agent.id)
if isinstance(exc, ProviderUpstreamError):
_append_typed_route_failure(
Expand Down Expand Up @@ -11958,7 +11968,7 @@ def attempt(agent: ModelAgent) -> EndpointAttempt[Any]:
decision = downgrade_to_failover(decision)
action = decision.action
self._record_tool_fallback(agent.id, decision, retry_attempt)
if decision.circuit_failure and not quota_rejection:
if decision.circuit_failure and self._charges_breaker(exc):
self._record_failure(agent.id)
if action is ToolFallbackAction.FAIL_CLOSED:
_append_tool_stop_route_attempt(
Expand Down Expand Up @@ -13001,6 +13011,11 @@ def _rate_limited_provider_signal(
current = current.__context__
return None

def _charges_breaker(self, exc: BaseException) -> bool:
"""Return false when provider quota, rather than member health, failed."""
signal = self._rate_limited_provider_signal(exc)
return signal is None or signal[0] != 429

def _agent(self, agent_id: str) -> ModelAgent:
for agent in self.candidates:
if agent.id == agent_id and _agent_matches_request_endpoint(agent):
Expand Down
37 changes: 37 additions & 0 deletions docs/doctoring/document-diff-data-uri-redos-20260930.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,37 @@
# Document diff data-URI ReDoS RCA

Status: Proposed until protected integration and exact-head hosted Checks.

## Incident

The Python CodeQL dispatch for
`ContextualWisdomLab/contextual-orchestrator#1221@4dcf9e32b057cde83bca67bfd45975fc6deda458`
reported `py/polynomial-redos` with security severity 7.5 at
`contextual_orchestrator/document_diff_review.py:129`. Central run
`36447487525`, job `109084173022`, preserved SARIF artifact `11026927998`
with digest
`sha256:0145d9b03e8c0064c2a57be79bc3f8168bcd45dc93b488aa4a4b1f81c23898a5`.

## Root cause

The data-URI boundary searched caller-controlled extracted document text with
`data:[^,\s]*,`. Repeated `data:` prefixes let the unanchored expression retry
over overlapping suffixes. The per-object byte ceiling bounded the input, but
did not make polynomial work an acceptable security boundary.

## Repair and invariant

PR #1349 is the canonical owner repair. Its fixed-token scan recognizes
case-insensitive `data:`, comma, and whitespace in one pass, preserving the
original `data:[^,\s]*,` language without a header-length cutoff. Binary
signatures, long base64 runs, credentials, resident registration numbers, byte
budgets, and provider-call admission are unchanged.

The #1221 ancestor retains the original CodeQL and SARIF identity above. The
stacked merge replaces its parallel scanner and duplicate regression with
#1349's broader contract: ordinary and mixed-case data URIs, repeated prefixes,
headers beyond the rejected 256-character workaround, whitespace termination,
and deterministic varied-text equivalence against the original language.
Hosted CodeQL on the unchanged merged head remains the authoritative acceptance
gate; this record does not convert queued, skipped, or stale results into
success.
29 changes: 28 additions & 1 deletion docs/product-technical-gap-baseline.md
Original file line number Diff line number Diff line change
@@ -1,5 +1,11 @@
# Contextual Orchestrator: Product & Technical Gap Baseline

## 2026-09-30 document diff data-URI scan — Proposed

| Gap ID | Status | Exact-head evidence | Repair / next gate |
|---|---|---|---|
| CO-DOCUMENT-DIFF-REDOS-01 | **Proposed — source repaired; hosted acceptance pending** | `contextual-orchestrator#1221@4dcf9e32b057cde83bca67bfd45975fc6deda458`; central CodeQL run `36447487525`, Python job `109084173022`, rule `py/polynomial-redos`, security severity 7.5, `contextual_orchestrator/document_diff_review.py:129`; SARIF artifact `11026927998`, digest `sha256:0145d9b03e8c0064c2a57be79bc3f8168bcd45dc93b488aa4a4b1f81c23898a5`. | Replace the unanchored data-URI regular expression with a disjoint-segment linear scanner while preserving fail-closed inline-media rejection. RED imports the absent scanner; GREEN covers ordinary, case-insensitive, whitespace-terminated, later-valid, and 1,600-prefix adversarial inputs. Re-run exact-head CodeQL only after this cause change; protected Checks and independent review remain required. |

## 2026-09-19 free multimodal review routing — Proposed

Canonical owner PR
Expand Down Expand Up @@ -7443,7 +7449,6 @@ keeps the value administrator-owned through `OrchestrationPolicy`, and adds
(prompt and parser follow the policy value; default stays 6). Not established:
an ablation of the bound itself, which belongs to the #568 equal-budget lane.


## 2026-09-20 Protected-main CI regression RCA

PR #995 exact head `d4f7e135719b16a0f0d72377d1bbe9b3614a8600`
Expand Down Expand Up @@ -7492,6 +7497,28 @@ test-owned listening socket. The fixture now calls `server_close()` in
`finally` after the serving thread joins. This changes no production server
lifecycle policy and keeps ResourceWarning visible as a failure signal.

## 2026-09-27 Observed-health admission authority — Proposed

PR #1221 exposed an operator-activatable routing policy whose fixed weight,
failure-rate threshold, observation window, cooldown, exponential escalation,
and all-open ordering changed candidate admission without a released
mathematical, statistical, psychometric, standards-based, or experimentally
validated authority. An opt-in Boolean or KV value records intent, not evidence;
the failure-only sidecar replay is explicitly biased and cannot calibrate those
decisions.

RED `ee58cb09` proves both supported activation surfaces remained reachable:
the constructor accepted `observed_health_quarantine=True`, and KV `enabled`
activated the same heuristic policy. The forward repair removes that
decision-affecting runtime surface and its replay artifacts instead of replacing
one arbitrary policy with another. The independent provider boundary remains:
HTTP 429 records quota cooldown evidence but never breaker health, while 503
remains an availability failure. Focused authority and three-path HTTP status
contracts must be GREEN on the exact successor before review. A future health
policy requires immutable owner identity, executable calibration provenance,
validated sampling/failure denominators, and a versioned released contract;
until then activation fails closed.

## 2026-09-30 Data-URI leak scan without a heuristic cutoff

PR [#1340](https://github.com/ContextualWisdomLab/contextual-orchestrator/pull/1340)
Expand Down
2 changes: 1 addition & 1 deletion tests/test_document_diff_review.py
Original file line number Diff line number Diff line change
Expand Up @@ -30,8 +30,8 @@

from contextual_orchestrator import ModelAgent, TaskOrchestrator # noqa: E402
from contextual_orchestrator.document_diff_review import ( # noqa: E402
_contains_data_uri,
DocumentDiffReviewError,
_contains_data_uri,
_scan_for_leaks,
validate_document_diff_envelope,
validate_document_diff_findings,
Expand Down
37 changes: 37 additions & 0 deletions tests/test_observed_health_policy_authority.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,37 @@
"""Fail closed while observed-health admission has no released authority."""

from __future__ import annotations

import inspect

import pytest

from contextual_orchestrator import ModelAgent, TaskOrchestrator
from contextual_orchestrator.credentials import (
InMemoryCredentialBackend,
register_credential,
set_backend,
)


def test_uncalibrated_constructor_activation_is_not_supported() -> None:
"""A boolean opt-in cannot authorize heuristic model admission."""
assert "observed_health_quarantine" not in inspect.signature(TaskOrchestrator).parameters
with pytest.raises(TypeError, match="observed_health_quarantine"):
TaskOrchestrator(
[ModelAgent("worker_one", "mock/worker")],
observed_health_quarantine=True,
)


def test_uncalibrated_kv_activation_has_no_runtime_surface() -> None:
"""A KV flag cannot substitute for released calibration evidence."""
set_backend(InMemoryCredentialBackend())
try:
register_credential(
"CONTEXTUAL_ORCHESTRATOR_OBSERVED_HEALTH_QUARANTINE", "enabled"
)
orchestrator = TaskOrchestrator([ModelAgent("worker_one", "mock/worker")])
assert not hasattr(orchestrator, "observed_health_quarantine")
finally:
set_backend(None)
170 changes: 170 additions & 0 deletions tests/test_rate_limit_breaker_asymmetry.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,170 @@
"""Pin provider 429 quota evidence across each chat path."""

from __future__ import annotations

import http.client
import json
import urllib.error
from dataclasses import replace

import pytest

from contextual_orchestrator import ModelAgent, TaskOrchestrator
from contextual_orchestrator.credentials import (
InMemoryCredentialBackend,
register_credential,
set_backend,
)
from contextual_orchestrator.orchestrator import ModelClient

_MESSAGES = [{"role": "user", "content": "summarize the incident"}]


class _Response:
"""A provider 200 response for chat, passthrough, or streaming."""

def __init__(self, stream: bool) -> None:
if stream:
chunk = {
"model": "second-model",
"choices": [{"index": 0, "delta": {"content": "ok"}}],
}
self._lines = [f"data: {json.dumps(chunk)}\n".encode(), b"data: [DONE]\n"]
else:
body = {
"id": "chatcmpl_ok",
"object": "chat.completion",
"created": 0,
"model": "second-model",
"choices": [
{
"index": 0,
"message": {"role": "assistant", "content": "ok"},
"finish_reason": "stop",
}
],
}
self._lines = [json.dumps(body).encode()]
self.status = 200
self.headers: dict[str, str] = {}

def __iter__(self):
return iter(self._lines)

def read(self, amount: int | None = None) -> bytes:
del amount
body, self._lines = b"".join(self._lines), []
return body

def getheader(self, name: str, default: str | None = None) -> str | None:
del name
return default

def close(self) -> None:
self._lines = []

def __enter__(self) -> "_Response":
return self

def __exit__(self, *exc_info: object) -> bool:
del exc_info
self.close()
return False


def _provider_error(status: int) -> urllib.error.HTTPError:
headers = http.client.HTTPMessage()
if status == 429:
headers["Retry-After"] = "7"
return urllib.error.HTTPError(
"https://provider.invalid/v1/chat/completions",
status,
"provider error",
headers,
None,
)


@pytest.fixture
def pool_for(monkeypatch: pytest.MonkeyPatch):
"""Return a two-member pool whose first provider yields one HTTP error."""
set_backend(InMemoryCredentialBackend())
register_credential("SYNTHETIC_PROVIDER_KEY", "synthetic")
raised: list[urllib.error.HTTPError] = []

def build(status: int) -> TaskOrchestrator:
def open_provider(self, request, destination=None, *, timeout=None):
del self, destination, timeout
payload = json.loads(request.data.decode())
if payload["model"] == "first-model":
raised.append(_provider_error(status))
raise raised[-1]
return _Response(stream=bool(payload.get("stream")))

monkeypatch.setattr(ModelClient, "_validate_provider", lambda self, agent: None)
monkeypatch.setattr(ModelClient, "_open_provider", open_provider)
orchestrator = TaskOrchestrator(
[
ModelAgent(
"first_agent",
"first-model",
base_url="https://provider.invalid/v1",
api_key_env="SYNTHETIC_PROVIDER_KEY",
tags=("reasoning", "writing"),
priority=10,
),
ModelAgent(
"second_agent",
"second-model",
base_url="https://provider.invalid/v2",
api_key_env="SYNTHETIC_PROVIDER_KEY",
tags=("reasoning", "writing"),
priority=1,
),
],
client=ModelClient(max_retries=0),
tool_retry_attempts=0,
tool_retry_backoff_seconds=0.0,
rate_limit_wait_seconds=0.0,
)
orchestrator.policy = replace(orchestrator.policy, realtime_judge=False)
orchestrator._triage_fn = lambda text: False
return orchestrator

try:
yield build
finally:
for error in raised:
error.close()
set_backend(None)


def _serve(orchestrator: TaskOrchestrator, path: str) -> str:
if path == "route":
result = orchestrator.route_once(list(_MESSAGES), model_name=TaskOrchestrator.AUTO_MODEL)
return result["answer"]
if path == "stream":
return "".join(
orchestrator.stream_route(list(_MESSAGES), model_name=TaskOrchestrator.AUTO_MODEL)
)
response = orchestrator.proxy_completion(
{"model": TaskOrchestrator.AUTO_MODEL, "messages": list(_MESSAGES)}
)
return response["choices"][0]["message"]["content"]


@pytest.mark.parametrize("path", ["route", "stream", "passthrough"])
def test_provider_429_records_quota_but_not_breaker_health(pool_for, path: str) -> None:
"""A quota response cannot become member-health evidence."""
orchestrator = pool_for(429)
assert _serve(orchestrator, path) == "ok"
assert orchestrator._rate_limit_remaining("first_agent") is not None, path
assert "first_agent" not in orchestrator._circuit, (path, orchestrator._circuit)


@pytest.mark.parametrize("path", ["route", "stream", "passthrough"])
def test_provider_503_still_counts_as_one_breaker_failure(pool_for, path: str) -> None:
"""A provider availability failure remains breaker evidence."""
orchestrator = pool_for(503)
assert _serve(orchestrator, path) == "ok"
assert orchestrator._circuit["first_agent"]["failures"] == 1.0, path
Loading