diff --git a/CHANGELOG.d/observed-health-authority.md b/CHANGELOG.d/observed-health-authority.md new file mode 100644 index 000000000..50589d9d3 --- /dev/null +++ b/CHANGELOG.d/observed-health-authority.md @@ -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. diff --git a/contextual_orchestrator/orchestrator.py b/contextual_orchestrator/orchestrator.py index fac6367f6..4e15e89b6 100644 --- a/contextual_orchestrator/orchestrator.py +++ b/contextual_orchestrator/orchestrator.py @@ -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 @@ -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( @@ -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( @@ -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): diff --git a/docs/doctoring/document-diff-data-uri-redos-20260930.md b/docs/doctoring/document-diff-data-uri-redos-20260930.md new file mode 100644 index 000000000..6ef251ef2 --- /dev/null +++ b/docs/doctoring/document-diff-data-uri-redos-20260930.md @@ -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. diff --git a/docs/product-technical-gap-baseline.md b/docs/product-technical-gap-baseline.md index 55498da2f..c90440a0a 100644 --- a/docs/product-technical-gap-baseline.md +++ b/docs/product-technical-gap-baseline.md @@ -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 @@ -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` @@ -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) diff --git a/tests/test_document_diff_review.py b/tests/test_document_diff_review.py index ca37fea91..f65ef3e81 100644 --- a/tests/test_document_diff_review.py +++ b/tests/test_document_diff_review.py @@ -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, diff --git a/tests/test_observed_health_policy_authority.py b/tests/test_observed_health_policy_authority.py new file mode 100644 index 000000000..5316aab7c --- /dev/null +++ b/tests/test_observed_health_policy_authority.py @@ -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) diff --git a/tests/test_rate_limit_breaker_asymmetry.py b/tests/test_rate_limit_breaker_asymmetry.py new file mode 100644 index 000000000..8d5a5b748 --- /dev/null +++ b/tests/test_rate_limit_breaker_asymmetry.py @@ -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