diff --git a/INSTALLATION.md b/INSTALLATION.md index 6354f120f..15755f0c1 100644 --- a/INSTALLATION.md +++ b/INSTALLATION.md @@ -36,6 +36,47 @@ See [Getting Started](docs/getting_started.md#server-path) for a complete TOML deployment and [`switchyard-server`](crates/switchyard-server/README.md) for the configuration reference. +## Linux Codex Service + +This setup is for single-user Linux machines. It requires a systemd user +session, the Rust toolchain listed above, and Codex CLI 0.134.0 or newer +logged in with a ChatGPT account. Another user's process could take the local +port while the service is stopped and receive your login and prompts. + +From a checkout, preview or install the service: + +```bash +make install-linux-dry-run +make install-linux +systemctl --user status switchyard +codex -p sy +``` + +The installer builds and installs `~/.switchyard/bin/switchyard-server`. +It creates `~/.switchyard/composite.toml` if missing and keeps existing edits. +It replaces `~/.config/systemd/user/switchyard.service` and restarts the service. +It writes `~/.codex/sy.config.toml`, backing up a changed profile first. +It leaves shell rc files unchanged. Use `codex -p sy` to select the profile; +`codex login` and other management commands still work as usual. +The profile format requires [Codex 0.134.0 or newer](https://developers.openai.com/codex/config-advanced#profiles). +Remove any old `[profiles.sy]` table from `~/.codex/config.toml` before using it. + +Set `SY_HOME` or `SY_PORT` to change the server directory or port (default 4123). +`SY_HOME` must not contain whitespace, control characters, or a trailing +backslash; `SY_PORT` must contain only digits. `XDG_CONFIG_HOME` and `CODEX_HOME` +set the systemd and Codex config directories. + +Read server logs with `journalctl --user -u switchyard`. Routing records are +stored in `~/.switchyard/routing.jsonl`. Use `systemctl --user edit switchyard` +for service changes that survive reinstalling. + +Remove the service and profile with `make uninstall-linux`. This also removes +marked Codex aliases left by older installs from existing `.bashrc` and +`.zshrc` files. If upgrading an older install, uninstall first and run +`unalias codex` in any open shell. The server directory, routing records, +profile backups, and systemd drop-in files stay in place. Use the same path +overrides when installing and uninstalling. + ## Rust Libraries Add the crates needed by an embedded application: diff --git a/Makefile b/Makefile new file mode 100644 index 000000000..eadc87cff --- /dev/null +++ b/Makefile @@ -0,0 +1,22 @@ +# SPDX-FileCopyrightText: Copyright (c) 2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved. +# SPDX-License-Identifier: Apache-2.0 + +.DEFAULT_GOAL := help +.PHONY: help install-linux install-linux-dry-run uninstall-linux + +help: + @echo "install-linux Install the systemd user service and Codex profile" + @echo "install-linux-dry-run Preview installation without changes" + @echo "uninstall-linux Remove the service and Codex profile" + +## Install the Switchyard background server as a systemd user service. +install-linux: + @scripts/linux/install.sh + +## Print what install-linux would do, without changing anything. +install-linux-dry-run: + @scripts/linux/install.sh --dry-run + +## Remove the systemd user service, the sy Codex profile, and the codex alias. +uninstall-linux: + @scripts/linux/uninstall.sh diff --git a/README.md b/README.md index 460ddb6a3..59963b09c 100644 --- a/README.md +++ b/README.md @@ -83,7 +83,7 @@ releases. Pin the version you integrate. | `switchyard-libsy` | **Beta** | Routing embedded in your own gateway or harness. You own model calls, credentials, and retries. | Trial integrations. API will change before v1.0. | | `switchyard-llm-client` | **Alpha** | HTTP model calls and protocol translation alongside libsy. | Experiments and pilots. | | `switchyard-runner` | **Alpha** | Running configured routes inside another runtime, such as NeMo Relay. | Integration work and supervised pilots. | -| `switchyard-server` | **Demo** | A standalone OpenAI- and Anthropic-compatible proxy. | Demos and evaluation only. Not for production. | +| `switchyard-server` | **Demo** | A standalone OpenAI- and Anthropic-compatible proxy, including a local Codex service. | Demos, evaluation, and personal use on single-user machines. Not for production. | ## Community and license diff --git a/benchmark/SWE_ATLAS_ESCALATION_REPORT.md b/benchmark/SWE_ATLAS_ESCALATION_REPORT.md new file mode 100644 index 000000000..92751180b --- /dev/null +++ b/benchmark/SWE_ATLAS_ESCALATION_REPORT.md @@ -0,0 +1,224 @@ +# SWE-Atlas escalation-router findings + +## Summary + +The escalation router could interpret duplicate command serialization in a NeMo Gym Terminus +trajectory as repeated failed work. A single assistant turn may contain the same command both as +raw JSON text and as a structured tool call. The judge prompt previously treated a command shown +two or more times as loop evidence without requiring those attempts to occur in distinct turns. + +This change makes escalation depend on fresh, consistent failure evidence across turns. It keeps +the existing one-way switch and latch: a session can move from the efficient model to the capable +model at most once and remains there afterward. + +On a deterministic 40-task paired sample (20 SWE-Atlas RF and 20 SWE-Atlas TW), the patched router +matched the two single-model arms at 15/40 correct and improved over the original router's 10/40. +It switched 5 times instead of 9. Under the synthetic token prices defined below, it cost +$0.889/task versus $1.171/task for the original router. GLM-only also scored 15/40 and remained the +least expensive arm at $0.501/task, so the result supports the patch over the original router but +does not show that routing beats the best fixed model on this sample. + +## Changes + +- Normalize only the transcript sent to the escalation judge. When a raw Terminus command batch + has a one-to-one match with structured tool calls in the same assistant turn, the redundant raw + representation is omitted. The worker model still receives the original history. +- Require the judge to classify positive evidence as `repetition`, `false_progress`, `drift`, + `desperation`, or `capability_gap`, and report whether the evidence is new. +- Advance the confirmation streak only for fresh positive evidence in the same category. A + category change starts a new streak; a decline or stale evidence resets it. An unavailable judge + preserves the prior streak rather than converting an infrastructure failure into a routing + decision. +- Clarify in the judge prompt that duplicate serialization within one turn and an immediate, + adaptive change of terminal strategy are not repeated failed attempts. +- Preserve the existing maximum-one-switch and post-switch latch behavior. + +## Focused SWE-Atlas reproducer observations + +The live checks used GLM 5.2 as the efficient model and Opus 4.8 as the capable model. The reward +column is the task's binary verifier result. "GLM only" means that the patched router did not +switch models during the trajectory. + +| Dataset | Task | Patched behavior | Reward | Observation | +| --- | --- | --- | ---: | --- | +| RF | `697e7458be1623d850a88838` | 28 GLM calls, then 36 Opus calls; one switch | 1 | Two fresh `repetition` findings confirmed the escalation path and latch. | +| RF | `69391d8d1ce51c407be1e531` | GLM only; 93 logged episodes | 0 | Findings changed from `repetition` to `false_progress`, so they did not form one confirmation streak. The earlier router switched without improving the reward. | +| RF | `694b4b99829f00e24fd11889` | GLM only; 34 logged episodes | 0 | No positive escalation finding. | +| TW | `6902ef3ab97fe23e2ad271f3` | GLM only; 58 logged episodes | 0 | One isolated `repetition` finding did not trigger a switch. The earlier router switched without improving the reward. | +| RF | `696719205599a51110d4b45f` | Three GLM-only runs: 72, 93, and 70 logged episodes | 0, 0, 0 | One isolated positive finding across the replicas; no latch. Direct GLM and direct Opus also scored 0. | +| RF | `696719205599a51110d4b455` | Three GLM-only runs: 51, 60, and 48 logged episodes | 0, 0, 0 | No transition in any replica. Direct GLM and direct Opus also scored 0. | + +A separate passing run of `697e7458be1623d850a88838` produced a positive finding followed by a +decline, stayed on GLM, and scored 1. This checks that non-consecutive findings do not accumulate. + +The positive switched task also scored 1 in a direct-GLM run. It therefore demonstrates correct +switching, confirmation, and latching, but does not establish a causal accuracy improvement. +Similarly, the unsuccessful controls show improved escalation precision and avoided Opus calls; +they do not show that the efficient model could solve those tasks. + +## Paired benchmark + +### Design + +The paired benchmark used the following fixed setup: + +- 20 tasks from `scale-ai/swe-atlas-rf@1` and 20 from `scale-ai/swe-atlas-tw@1`. + Tasks were selected deterministically by ranking the hash of + `20260904::`; they were not selected by outcome. +- Four arms per task: direct GLM, direct Opus, the original escalation router, and the patched + escalation router. This produced 160 task-arm runs. +- GLM was `nvidia/zai-org/glm-5.2`; Opus was + `aws/anthropic/bedrock-claude-opus-4-8` with medium effort. GLM also served as the escalation + judge. The two escalation arms used two confirmations, maximum one switch, and a post-switch + latch. +- The original router was built from upstream commit + `9a743e89223a0d5b14011f1226d5b068f730a3b8`; the patched router was built from + `507376b1d851b825378cf147d09a8b5594773298`. +- There was no agent-turn limit. Worker requests allowed up to 128,000 output tokens and a + 900-second model-server timeout. Five initial runs reached Harbor's separate 3,600-second agent + timeout. Only those five task-arm runs were repeated with a 3x agent-timeout multiplier; all five + retries completed and replaced the invalid attempts. The final effective matrix is 160/160 valid. +- Each table entry is one trial per task and arm. Binary verifier reward is reported as correct. + Because model sampling is stochastic, differences between unswitched GLM router runs and direct + GLM runs are not necessarily routing effects. + +### Accuracy and routing + +| Dataset | Arm | Correct | Accuracy | Switches | Correct after switch | +| --- | --- | ---: | ---: | ---: | ---: | +| RF (20) | Direct GLM | 7 | 0.350 | -- | -- | +| RF (20) | Always Opus | 9 | 0.450 | -- | -- | +| RF (20) | Original escalation | 6 | 0.300 | 7 | 2 | +| RF (20) | **Patched escalation** | **9** | **0.450** | **4** | **2** | +| TW (20) | Direct GLM | 8 | 0.400 | -- | -- | +| TW (20) | Always Opus | 6 | 0.300 | -- | -- | +| TW (20) | Original escalation | 4 | 0.200 | 2 | 0 | +| TW (20) | **Patched escalation** | **6** | **0.300** | **1** | **0** | +| Combined (40) | Direct GLM | 15 | 0.375 | -- | -- | +| Combined (40) | Always Opus | 15 | 0.375 | -- | -- | +| Combined (40) | Original escalation | 10 | 0.250 | 9 | 2 | +| Combined (40) | **Patched escalation** | **15** | **0.375** | **5** | **2** | + +Against the original router on the same 40 tasks, the patched router won 7 outcomes, lost 2, and +tied 31. Against direct GLM it won 5 and lost 5; against direct Opus it won 6 and lost 6. All 14 +switched trajectories across the two router arms had exactly one GLM-to-Opus transition and no +hand-back, confirming the switch and latch invariants in live execution. + +The patch traded earlier sensitivity for precision. Original-router switches occurred after a +median of 23 GLM calls; patched-router switches occurred after a median of 51. One patched RF retry +did not switch until 141 GLM calls, then made 57 latched Opus calls and still scored 0. Direct Opus +and the original router also scored 0 on that task. This avoided a premature switch but exposed a +cost problem: judging every turn can be expensive when the destination model is unlikely to help. + +### Synthetic cost calculation + +The inference service did not provide billable cost. The values below are synthetic estimates for +comparing routing behavior, not NVIDIA prices or charges. The assumed rates are: + +| Model | Uncached input / 1M tokens | Cached input / 1M tokens | Output / 1M tokens | +| --- | ---: | ---: | ---: | +| GLM 5.2 | $0.50 | $0.05 | $2.00 | +| Opus 4.8 | $5.00 | $0.50 | $25.00 | + +For each model in each trajectory: + +```text +uncached_input_tokens = max(prompt_tokens - cached_tokens, 0) + +model_cost = ( + uncached_input_tokens * uncached_input_rate + + cached_tokens * cached_input_rate + + output_tokens * output_rate +) / 1,000,000 + +trajectory_cost = sum(worker_model_costs) + sum(escalation_judge_costs) +``` + +Cache-creation tokens are part of the non-cached prompt-token remainder and therefore receive the +uncached-input rate. The calculation includes GLM judge calls for the two router arms. It excludes +the Harbor verifier, cluster resources, and other infrastructure because those calls are outside +Switchyard's per-trajectory statistics. + +| Dataset | Arm | Synthetic cost | Cost/task | Cost/correct | +| --- | --- | ---: | ---: | ---: | +| RF (20) | Direct GLM | $12.953 | $0.648 | $1.850 | +| RF (20) | Always Opus | $32.196 | $1.610 | $3.577 | +| RF (20) | Original escalation | $35.169 | $1.758 | $5.862 | +| RF (20) | **Patched escalation** | **$27.200** | **$1.360** | **$3.022** | +| TW (20) | Direct GLM | $7.093 | $0.355 | $0.887 | +| TW (20) | Always Opus | $14.829 | $0.741 | $2.471 | +| TW (20) | Original escalation | $11.657 | $0.583 | $2.914 | +| TW (20) | **Patched escalation** | **$8.356** | **$0.418** | **$1.393** | +| Combined (40) | Direct GLM | $20.045 | $0.501 | $1.336 | +| Combined (40) | Always Opus | $47.025 | $1.176 | $3.135 | +| Combined (40) | Original escalation | $46.826 | $1.171 | $4.683 | +| Combined (40) | **Patched escalation** | **$35.556** | **$0.889** | **$2.370** | + +The patched router was 24.1% cheaper than the original router and produced five additional correct +outcomes. It was 24.4% cheaper than always-Opus at the same aggregate accuracy. It was 77.4% more +expensive than direct GLM at the same aggregate accuracy. Of the patched router's $35.556 synthetic +total, $3.698 (10.4%) came from judge calls; this motivates reducing judge cadence after repeated +negative verdicts. + +### TW model-order reversal + +TW demonstrates a limitation that this patch does not solve: the escalation route assumes its +configured Opus target has positive expected gain over GLM. On this sample, direct GLM scored 8/20 +while direct Opus scored 6/20, so that assumption is not valid at the workload level. + +The patched router switched on only one TW task. Direct GLM, direct Opus, and both router arms all +scored 0 on that task, so the switch added cost without changing its outcome. The patched router's +two-correct deficit versus direct GLM was not caused by switching: it lost three independently +sampled GLM-only outcomes and gained one different GLM-only outcome. The original router provides +a clearer harmful-switch example on task `6902ef3ab97fe23e2ad2727e`: direct GLM scored 1, direct Opus +scored 0, and the original router switched and scored 0. The patch avoided that switch, although its +independent GLM-only run also scored 0. + +The next routing change should therefore be separate from the trajectory-evidence fix: gate +escalation on a workload- or task-conditioned estimate of expected accuracy gain and cost. When +recent paired evidence says the nominal strong model is worse, the route should disable or reverse +that transition rather than asking only whether the current trajectory looks stuck. This requires +held-out calibration or online exploration; the router cannot infer relative model quality from one +model's failing trajectory alone. + +## Validation evidence + +The implementation was validated with Rust 1.96.1 on a Slurm compute node using: + +```text +cargo fmt --all --check +cargo clippy --workspace --all-targets --all-features -- -D warnings +cargo test --workspace --all-features +uv run ruff check . +uv run mypy switchyard +uv run maturin develop --uv +uv run pytest tests/ -v +``` + +The validation completed successfully: the Rust workspace passed, including 289 `libsy` tests; +ruff and mypy passed; the native wheel built and installed; and pytest reported 115 passed, +2 deselected, and 2 subtests passed. +Focused regression coverage includes exact and multi-command transcript normalization, partial and +multiplicity mismatches, category changes, stale evidence, unavailable-judge behavior, and the +one-way latch. + +After rebasing onto current `main` on 2026-09-22, the review fixes were revalidated with Rust +1.96.1 and Python 3.12 using `cargo fmt --all --check`, workspace clippy with all targets and +features, and workspace tests with all features. The full Rust gate passed, including 323 `libsy` +tests and all 56 server integration tests. This post-rebase run caught and fixed stale call +signatures and mock verdicts that still used the pre-rebase runtime interfaces and escalation +schema. + +No live provider calls are made by the repository validation commands. The live SWE-Atlas checks +described above were separate, explicitly configured benchmark runs. + +## Conclusion + +The patch fixes the concrete duplicate-serialization false positive and materially improves the +tested router: it matches the single-model aggregate accuracy, gains 5 correct outcomes over the +original router, cuts switches from 9 to 5, and reduces synthetic cost by 24.1%. It does not make +the router the best cost/accuracy arm; direct GLM matches its aggregate accuracy at substantially +lower cost. The remaining work is not to loosen the new evidence rule globally. It is to add +model-order calibration for workloads such as TW and reduce judge overhead or stop judging when +escalation has low expected value. Replicated trials on held-out tasks are needed before treating +the point estimates as stable production gains. diff --git a/crates/libsy-llm-client/README.md b/crates/libsy-llm-client/README.md index 34ec8a989..8b0ef430d 100644 --- a/crates/libsy-llm-client/README.md +++ b/crates/libsy-llm-client/README.md @@ -248,6 +248,12 @@ fn build_multi_format_client( and `anthropic-version`. Header names are case-insensitive. - Per-target top-level request defaults go in `HttpBackendConfig::extra_body`. The merge is shallow and fields already present in the request take precedence. +- Anthropic reconstruction preserves the caller's native `thinking` settings, including + disabled thinking and manual budgets. `output_config.effort` is kept separately. + Effort implies adaptive thinking only when native `thinking` is absent. + To replace either top-level field for a target, list it in `omit_body_fields` and + supply its replacement in `extra_body`. Omission runs before defaults are merged. + Target `reasoning_effort` overrides apply to OpenAI backends. - `HttpBackendConfig::max_retries` controls additional attempts after retryable transport failures, timeouts, HTTP 408/429, and 5xx responses. Buffered body transport failures are retried; streaming body failures are not replayed after diff --git a/crates/libsy-llm-client/tests/observability.rs b/crates/libsy-llm-client/tests/observability.rs index 32aae6feb..93c12adb5 100644 --- a/crates/libsy-llm-client/tests/observability.rs +++ b/crates/libsy-llm-client/tests/observability.rs @@ -772,7 +772,9 @@ async fn stateful_escalation_warns_once_without_a_session_id() -> switchyard_lib })?) as Arc; let client = Arc::new(JudgeClient { judge_model: "warning-judge".into(), - outcome: JudgeOutcome::Reply(r#"{"escalate":false,"reason":"progressing"}"#), + outcome: JudgeOutcome::Reply( + r#"{"escalate":false,"category":"none","new_evidence":false,"reason":"progressing"}"#, + ), }) as Arc; for _ in 0..2 { @@ -828,7 +830,9 @@ async fn deescalation_evidence_stays_pending_until_confirmed() -> switchyard_lib switchyard_llm_client::run( router.clone(), - ClientRouter::single(client(r#"{"escalate":true,"reason":"stuck"}"#)), + ClientRouter::single(client( + r#"{"escalate":true,"category":"repetition","new_evidence":true,"reason":"stuck"}"#, + )), request.clone(), classifier_models("evidence-judge", "evidence-efficient", "evidence-capable"), None, @@ -836,7 +840,9 @@ async fn deescalation_evidence_stays_pending_until_confirmed() -> switchyard_lib .await?; let outcome = switchyard_llm_client::decide( router, - ClientRouter::single(client(r#"{"escalate":false,"reason":"recovered"}"#)), + ClientRouter::single(client( + r#"{"escalate":false,"category":"none","new_evidence":false,"reason":"recovered"}"#, + )), request, classifier_models("evidence-judge", "evidence-efficient", "evidence-capable"), ) diff --git a/crates/libsy/src/algorithms/advisor_gate.rs b/crates/libsy/src/algorithms/advisor_gate.rs index d638864f8..4551ac703 100644 --- a/crates/libsy/src/algorithms/advisor_gate.rs +++ b/crates/libsy/src/algorithms/advisor_gate.rs @@ -52,6 +52,7 @@ mod trigger; mod turn; use super::util::buffered_response::{BufferedResponse, buffer_response}; +use super::util::robustness::safe_error_summary; use budget::{ReviewBudget, ScopeKey, budget_scope, stall_key}; use signals::{GateSignalProcessor, GateSignals}; use telemetry::{ @@ -440,14 +441,17 @@ impl AdvisorGate { "advisor consult failed (fail_open = false): {error}" ))); } + // An upstream error's Display can quote request content back, + // so only its redacted summary reaches logs or the audit trail. + let summary = safe_error_summary(&error); tracing::warn!( target: "libsy", - error = %error, + error = %summary, "advisor gate: consult failed; passing the turn through (fail open)" ); emit_review_audit(ReviewAudit { verdict: "APPROVE", - error: Some(error.to_string()), + error: Some(summary), latency_ms, reply_head: None, usage: None, diff --git a/crates/libsy/src/algorithms/advisor_gate/tests.rs b/crates/libsy/src/algorithms/advisor_gate/tests.rs index 72d39615b..10bb968d4 100644 --- a/crates/libsy/src/algorithms/advisor_gate/tests.rs +++ b/crates/libsy/src/algorithms/advisor_gate/tests.rs @@ -676,6 +676,85 @@ async fn fail_closed_propagates_refunds_and_counts() { assert_eq!(completion_text(&agg_of(response).await), "recovered"); } +/// Collects everything the `fmt` subscriber renders, so assertions can run +/// against the final log sink rather than a field mid-pipeline. +#[derive(Clone, Default)] +struct LogCapture(Arc>>); + +impl std::io::Write for LogCapture { + fn write(&mut self, buf: &[u8]) -> std::io::Result { + self.0.lock().unwrap().extend_from_slice(buf); + Ok(buf.len()) + } + + fn flush(&mut self) -> std::io::Result<()> { + Ok(()) + } +} + +#[tokio::test] +async fn fail_open_logs_redact_the_upstream_error_body() { + // A target name unique to this test keys the captured warn line and audit + // record to this drive: the process-wide capture also sees parallel + // advisor-gate tests' events. + const ADVISOR_TARGET: &str = "advisor-redaction-probe"; + const MARKER: &str = "SECRET-UPSTREAM-QUOTE-the-user-prompt"; + let gate = gate(AdvisorGateConfig::default()); + let models: HashMap> = [ + (Category::Efficient, vec![target(EXECUTOR)]), + (Category::Judge, vec![target(ADVISOR_TARGET)]), + ] + .into(); + let serve = move |model: ModelId, _request: Request| { + Box::pin(async move { + let model = model.to_string(); + if model == ADVISOR_TARGET { + // UpstreamHttp's Display interpolates the raw body, which is + // exactly the content that must not reach a sink. + Err(LlmClientError::UpstreamHttp { + status: http::StatusCode::INTERNAL_SERVER_ERROR, + body: format!("{MARKER}: validation failed"), + }) + } else { + Ok(reply("done")) + } + }) + }; + + // Installed as the process default so the capture sees events from every + // thread the drive touches: the suite runs tests in parallel, and a + // thread-local default does not cover them all. + let capture = LogCapture::default(); + let writer = capture.clone(); + let subscriber = tracing_subscriber::fmt() + .with_ansi(false) + .with_max_level(tracing_subscriber::filter::LevelFilter::INFO) + .with_writer(move || writer.clone()) + .finish(); + tracing::subscriber::set_global_default(subscriber).expect("unused global subscriber"); + + let (_, response) = test_drive_with_models(gate, task_request(), models, serve) + .await + .expect("fail-open run"); + assert_eq!(completion_text(&agg_of(response).await), "done"); + + let logs = String::from_utf8(capture.0.lock().unwrap().clone()).expect("logs are utf-8"); + assert!( + !logs.contains(MARKER), + "the upstream error body reached a log sink: {logs}" + ); + // Only this drive's decision target names the redacted summary, so its + // presence proves this drive's warn rendered with the safe summary. + let summary = r#"client call to target "advisor-redaction-probe" failed: upstream HTTP 500"#; + assert!(logs.contains(summary), "{logs}"); + // The audit record still renders, keyed to this drive's target, with the + // redacted summary rather than the body. + let audit = logs + .lines() + .find(|line| line.contains("advisor_review=") && line.contains(ADVISOR_TARGET)); + assert!(audit.is_some(), "{logs}"); +} + #[tokio::test] async fn unparseable_verdict_refunds_and_approves() { let script = Script::new(); diff --git a/crates/libsy/src/algorithms/escalation.rs b/crates/libsy/src/algorithms/escalation.rs index 534ef34a0..76c77d76b 100644 --- a/crates/libsy/src/algorithms/escalation.rs +++ b/crates/libsy/src/algorithms/escalation.rs @@ -15,8 +15,8 @@ use super::util::buffered_response::buffer_response; use super::util::classifier_contract::ClassifierContractConfig; use super::util::decisive; use super::util::escalation::{ - self, DeescalationConfig, EscalationJudge, EscalationJudgeConfig, EscalationPolicy, - EvaluationPhase, + self, DeescalationConfig, EscalationCategory, EscalationJudge, EscalationJudgeConfig, + EscalationPolicy, EvaluationPhase, }; use super::util::llm_judge::JudgeClassifier; use crate::core::algorithm::Driver; @@ -25,6 +25,8 @@ use crate::core::state::{State, StateValue}; use crate::{LibsyError, Result}; const STREAK_KEY: &str = "escalation_streak"; +/// Session-state key holding the category currently being confirmed. +const CATEGORY_KEY: &str = "escalation_category"; const STRONG_CALLS_KEY: &str = "escalation_strong_calls"; const RELEASE_STREAK_KEY: &str = "escalation_release_streak"; const WEAK_COOLDOWN_KEY: &str = "escalation_weak_cooldown"; @@ -42,6 +44,20 @@ fn set_count(state: &mut State, key: &str, value: u32) { .insert(key.to_string(), StateValue::Count(value)); } +/// Reads the failure category whose confirmation streak is in progress, if any. +fn category(state: &State) -> Option<&str> { + match state.extra.get(CATEGORY_KEY) { + Some(StateValue::String(category)) => Some(category), + _ => None, + } +} + +/// Clears the confirmation streak together with the category it was confirming. +fn clear_streak(state: &mut State) { + set_count(state, STREAK_KEY, 0); + state.extra.remove(CATEGORY_KEY); +} + fn assistant_message(response: &AggLlmResponse) -> Message { Message { role: Role::Assistant, @@ -171,7 +187,7 @@ impl EscalationClassifier { set_count(state, RELEASE_STREAK_KEY, release_streak); if release_streak >= deescalation.config.confirmations { - set_count(state, STREAK_KEY, 0); + clear_streak(state); set_count(state, STRONG_CALLS_KEY, 0); set_count(state, RELEASE_STREAK_KEY, 0); tracing::debug!( @@ -225,7 +241,7 @@ impl Classifier for EscalationClassifier { .strong_max_calls .is_some_and(|strong_max_calls| strong_calls >= strong_max_calls) { - set_count(state, STREAK_KEY, 0); + clear_streak(state); set_count(state, STRONG_CALLS_KEY, 0); set_count(state, RELEASE_STREAK_KEY, 0); set_count( @@ -336,19 +352,52 @@ impl Classifier for EscalationClassifier { )); } - let (classification, _) = self + let judge_models = driver.models_for(&Category::Judge); + if judge_models.is_empty() { + return Err(LibsyError::AlgorithmError { + message: "no models available for category Judge".to_string(), + }); + } + let verdict = self .escalation_judge - .score(state, &mut judge_request, driver) - .await?; + .verdict(state, &judge_request, driver, judge_models) + .await; let held = count(state, STREAK_KEY); - let best = classification.argmax(false)?; - let (escalate, pending) = match &best { - Some(score) if score.target == capable => (true, held.saturating_add(1)), - Some(_) => (false, 0), - None => (false, held), + let held_category = category(state).map(str::to_string); + let (escalate, pending, pending_category) = match verdict.as_ref() { + Some(verdict) => { + let category = verdict.category.label(); + tracing::info!( + escalate = verdict.escalate, + category, + new_evidence = verdict.new_evidence, + "escalation judge verdict" + ); + if verdict.escalate + && verdict.new_evidence + && verdict.category != EscalationCategory::None + { + let next = if held_category.as_deref() == Some(category) { + held.saturating_add(1) + } else { + 1 + }; + (true, next, Some(category.to_string())) + } else { + (false, 0, None) + } + } + None => (false, held, held_category), }; set_count(state, STREAK_KEY, pending); + if let Some(category) = pending_category { + state + .extra + .insert(CATEGORY_KEY.to_string(), StateValue::String(category)); + } else { + state.extra.remove(CATEGORY_KEY); + } if escalate && pending >= self.confirmations { driver.set_evidence(serde_json::json!({ @@ -367,6 +416,11 @@ impl Classifier for EscalationClassifier { "source": "escalation", "verdict": "pending", })); + } else if verdict.is_some() { + driver.set_evidence(serde_json::json!({ + "source": "escalation", + "verdict": "continue", + })); } Ok(( @@ -475,13 +529,13 @@ mod tests { } } - /// Builds a router with escalation enabled (`confirmations=1` latches immediately). - fn escalation_router() -> Result> { + /// Builds a router with escalation enabled. + fn escalation_router_with_confirmations(confirmations: u32) -> Result> { Ok(Arc::new(LlmTaskClassifier::new( LlmClassifierConfig::Escalation { contract: ClassifierContractConfig::default(), config: EscalationJudgeConfig { - confirmations: 1, + confirmations, ..EscalationJudgeConfig::default() }, max_output_tokens: DEFAULT_JUDGE_MAX_OUTPUT_TOKENS, @@ -526,9 +580,16 @@ mod tests { Ok(selected) } + /// Builds a router that latches on its first supported escalation verdict. + fn escalation_router() -> Result> { + escalation_router_with_confirmations(1) + } + #[tokio::test] async fn serves_efficient_when_judge_declines() -> Result<()> { - let judge = Queue::new([r#"{"escalate":false,"reason":"progressing"}"#]); + let judge = Queue::new([ + r#"{"escalate":false,"category":"none","new_evidence":false,"reason":"progressing"}"#, + ]); let model = Queue::new(["efficient answer"]); let (selected_model, response) = test_drive_with_models( @@ -547,6 +608,44 @@ mod tests { Ok(()) } + /// A parsed decline records continue evidence on the routing outcome. + #[tokio::test] + async fn records_continue_evidence_when_judge_declines() -> Result<()> { + let judge = Queue::new([ + r#"{"escalate":false,"category":"none","new_evidence":false,"reason":"progressing"}"#, + ]); + let model = Queue::new(["efficient answer"]); + let serve = Arc::new(queued(model, judge)); + let routing_serve = Arc::clone(&serve); + let outcome = crate::drive( + escalation_router()?, + classify_request(), + Arc::new(crate::core::algorithm::RuntimeModels::new(runtime_models())), + move |call| { + let serve = Arc::clone(&routing_serve); + async move { + let target = call.models.first().cloned().ok_or(LibsyError::NoTargets)?; + let request = call.request.clone(); + let response = serve + .serve(target.clone(), request) + .await + .map_err(|source| LibsyError::client_call(target, source)); + call.respond(response) + } + }, + ) + .await?; + + assert_eq!( + outcome.metadata.and_then(|metadata| metadata.evidence), + Some(serde_json::json!({ + "source": "escalation", + "verdict": "continue", + })) + ); + Ok(()) + } + #[tokio::test] async fn config_overrides_the_packaged_prompt() -> Result<()> { let prompts = Arc::new(Mutex::new(Vec::new())); @@ -564,7 +663,9 @@ mod tests { }) }); recorded.lock().extend(prompt); - std::future::ready(Ok(reply(r#"{"escalate":false,"reason":"progressing"}"#))) + std::future::ready(Ok(reply( + r#"{"escalate":false,"category":"none","new_evidence":false,"reason":"progressing"}"#, + ))) } else { std::future::ready(Ok(reply("efficient answer"))) } @@ -586,7 +687,9 @@ mod tests { #[tokio::test] async fn upgrades_to_capable_when_judge_escalates() -> Result<()> { - let judge = Queue::new([r#"{"escalate":true,"reason":"stuck in a loop"}"#]); + let judge = Queue::new([ + r#"{"escalate":true,"category":"repetition","new_evidence":true,"reason":"stuck in a loop"}"#, + ]); let model = Queue::new(["efficient draft", "capable answer"]); let (selected_model, response) = test_drive_with_models( @@ -605,9 +708,106 @@ mod tests { Ok(()) } + /// A verdict in a different category restarts the confirmation streak at one. + #[tokio::test] + async fn confirmation_streak_requires_the_same_category() -> Result<()> { + let judge = Queue::new([ + r#"{"escalate":true,"category":"repetition","new_evidence":true,"reason":"repeated command"}"#, + r#"{"escalate":true,"category":"drift","new_evidence":true,"reason":"off task"}"#, + r#"{"escalate":true,"category":"drift","new_evidence":true,"reason":"still off task"}"#, + ]); + let model = Queue::new(["efficient t1", "efficient t2", "efficient t3", "capable t3"]); + let router = escalation_router_with_confirmations(2)?; + let request = classify_session_request(); + + let (first, _) = test_drive_with_models( + router.clone(), + request.clone(), + runtime_models(), + queued(Arc::clone(&model), Arc::clone(&judge)), + ) + .await?; + let (second, _) = test_drive_with_models( + router.clone(), + request.clone(), + runtime_models(), + queued(Arc::clone(&model), Arc::clone(&judge)), + ) + .await?; + let (third, _) = + test_drive_with_models(router, request, runtime_models(), queued(model, judge)).await?; + + assert_eq!(first, "efficient"); + assert_eq!(second, "efficient"); + assert_eq!(third, "capable"); + Ok(()) + } + + /// An escalate verdict without fresh evidence resets the streak instead of extending it. + #[tokio::test] + async fn verdict_without_new_evidence_resets_the_streak() -> Result<()> { + let judge = Queue::new([ + r#"{"escalate":true,"category":"repetition","new_evidence":true,"reason":"repeated command"}"#, + r#"{"escalate":true,"category":"repetition","new_evidence":false,"reason":"only old evidence remains"}"#, + r#"{"escalate":true,"category":"repetition","new_evidence":true,"reason":"new repeated command"}"#, + ]); + let model = Queue::new(["efficient t1", "efficient t2", "efficient t3"]); + let router = escalation_router_with_confirmations(2)?; + let request = classify_session_request(); + + for _ in 0..3 { + let (selected, _) = test_drive_with_models( + router.clone(), + request.clone(), + runtime_models(), + queued(Arc::clone(&model), Arc::clone(&judge)), + ) + .await?; + assert_eq!(selected, "efficient"); + } + Ok(()) + } + + /// An unparseable verdict keeps the streak and its category for the next turn. + #[tokio::test] + async fn unavailable_judge_preserves_the_category_streak() -> Result<()> { + let judge = Queue::new([ + r#"{"escalate":true,"category":"repetition","new_evidence":true,"reason":"repeated command"}"#, + "not json", + r#"{"escalate":true,"category":"repetition","new_evidence":true,"reason":"another repeated command"}"#, + ]); + let model = Queue::new(["efficient t1", "efficient t2", "efficient t3", "capable t3"]); + let router = escalation_router_with_confirmations(2)?; + let request = classify_session_request(); + + let (first, _) = test_drive_with_models( + router.clone(), + request.clone(), + runtime_models(), + queued(Arc::clone(&model), Arc::clone(&judge)), + ) + .await?; + let (second, _) = test_drive_with_models( + router.clone(), + request.clone(), + runtime_models(), + queued(Arc::clone(&model), Arc::clone(&judge)), + ) + .await?; + let (third, _) = + test_drive_with_models(router, request, runtime_models(), queued(model, judge)).await?; + + assert_eq!(first, "efficient"); + assert_eq!(second, "efficient"); + assert_eq!(third, "capable"); + Ok(()) + } + #[tokio::test] async fn stays_capable_after_latch() -> Result<()> { - let judge = Queue::new([r#"{"escalate":true,"reason":"stuck"}"#]); + let judge = Queue::new([ + r#"{"escalate":true,"category":"repetition","new_evidence":true,"reason":"stuck"}"#, + ]); let model = Queue::new(["efficient draft", "capable t1", "capable t2"]); let router = escalation_router()?; let request = classify_session_request(); @@ -629,10 +829,10 @@ mod tests { #[tokio::test] async fn deescalation_holds_then_returns_to_efficient() -> Result<()> { let judge = Queue::new([ - r#"{"escalate":true,"reason":"stuck"}"#, - r#"{"escalate":false,"reason":"recovered"}"#, - r#"{"escalate":false,"reason":"routine"}"#, - r#"{"escalate":false,"reason":"progressing"}"#, + r#"{"escalate":true,"category":"repetition","new_evidence":true,"reason":"stuck"}"#, + r#"{"escalate":false,"category":"none","new_evidence":false,"reason":"recovered"}"#, + r#"{"escalate":false,"category":"none","new_evidence":false,"reason":"routine"}"#, + r#"{"escalate":false,"category":"none","new_evidence":false,"reason":"progressing"}"#, ]); let model = Queue::new([ "efficient draft", @@ -658,10 +858,10 @@ mod tests { #[tokio::test] async fn deescalation_hard_limit_forces_a_weak_cooldown() -> Result<()> { let judge = Queue::new([ - r#"{"escalate":true,"reason":"stuck"}"#, - r#"{"escalate":true,"reason":"still hard"}"#, - r#"{"escalate":true,"reason":"still hard"}"#, - r#"{"escalate":false,"reason":"progressing"}"#, + r#"{"escalate":true,"category":"repetition","new_evidence":true,"reason":"stuck"}"#, + r#"{"escalate":true,"category":"repetition","new_evidence":true,"reason":"still hard"}"#, + r#"{"escalate":true,"category":"repetition","new_evidence":true,"reason":"still hard"}"#, + r#"{"escalate":false,"category":"none","new_evidence":false,"reason":"progressing"}"#, ]); let model = Queue::new([ "efficient draft", @@ -695,7 +895,9 @@ mod tests { #[tokio::test] async fn strong_review_returns_an_efficient_fallback_without_judging_it() -> Result<()> { - let judge = Queue::new([r#"{"escalate":true,"reason":"stuck"}"#]); + let judge = Queue::new([ + r#"{"escalate":true,"category":"repetition","new_evidence":true,"reason":"stuck"}"#, + ]); let model = Queue::new(["efficient draft", "capable t1", "efficient fallback"]); let router = deescalation_router(DeescalationConfig { strong_min_calls: 1, diff --git a/crates/libsy/src/algorithms/llm_class.rs b/crates/libsy/src/algorithms/llm_class.rs index 488bc6585..06e5449cc 100644 --- a/crates/libsy/src/algorithms/llm_class.rs +++ b/crates/libsy/src/algorithms/llm_class.rs @@ -454,7 +454,8 @@ impl CustomClassifierPolicy { pub struct CustomClassifierConfig { /// System prompt sent to the classifier judge. pub prompt: String, - /// Inner JSON Schema placed inside the provider's structured-output wrapper. + /// Inner JSON Schema for the verdict. Sent in the provider's structured-output wrapper in + /// JSON Schema mode, or appended to the prompt in JSON Object mode. pub response_schema: Value, /// Deterministic policy applied after the verdict passes schema validation. pub policy: CustomClassifierPolicy, @@ -466,6 +467,8 @@ pub struct CustomClassifierConfig { pub recent_turn_window: Option, /// Maximum completion tokens available to the classifier verdict. pub max_output_tokens: u64, + /// Structured-output mode requested from the classifier judge. + pub response_format_type: ClassifierResponseFormat, } impl CustomClassifierConfig { @@ -483,6 +486,7 @@ impl CustomClassifierConfig { message_hash_fallback: false, recent_turn_window: None, max_output_tokens: DEFAULT_JUDGE_MAX_OUTPUT_TOKENS, + response_format_type: ClassifierResponseFormat::default(), } } @@ -666,8 +670,10 @@ impl LlmTaskClassifier { message_hash_fallback, recent_turn_window, max_output_tokens, + response_format_type, } = config; - let contract = ClassifierContract::from_inner_schema(&prompt, response_schema)?; + let contract = + ClassifierContract::from_inner_schema(&prompt, response_schema, response_format_type)?; let policy = match policy { CustomClassifierPolicy::TargetSelector { selector } => { CustomPolicyRuntime::TargetSelector(TargetSelectorPolicy::new(selector)?) diff --git a/crates/libsy/src/algorithms/util/classifier_contract.rs b/crates/libsy/src/algorithms/util/classifier_contract.rs index 6cdc42551..defe4b53d 100644 --- a/crates/libsy/src/algorithms/util/classifier_contract.rs +++ b/crates/libsy/src/algorithms/util/classifier_contract.rs @@ -78,6 +78,7 @@ impl ClassifierContract { response_format_json: &str, ) -> Result { let prompt_template = config.prompt().unwrap_or(default_prompt); + validate_prompt(prompt_template)?; let response_format: Value = serde_json::from_str(response_format_json).map_err(|error| { LibsyError::AlgorithmError { @@ -90,29 +91,33 @@ impl ClassifierContract { message: "response schema has no json_schema.schema".to_string(), })?; match config.response_format_type() { - ClassifierResponseFormat::JsonSchema => { - Self::from_response_format(prompt_template, response_format, None) - } + ClassifierResponseFormat::JsonSchema => Ok(Self::from_response_format( + prompt_template.to_string(), + response_format, + None, + )), ClassifierResponseFormat::JsonObject => { - validate_prompt(prompt_template)?; let validator = compile_schema(schema)?; - let rendered_schema = serde_json::to_string_pretty(schema).map_err(|error| { - algorithm_error(format!("response schema could not be rendered: {error}")) - })?; - let system_prompt = format!( - "{prompt_template}\n\nReturn exactly one JSON object matching this JSON Schema:\n{rendered_schema}" - ); - Self::from_response_format( - &system_prompt, + Ok(Self::from_response_format( + schema_in_prompt(prompt_template, schema)?, json!({"type": "json_object"}), Some(validator), - ) + )) } } } - /// Builds a provider response format around a user-supplied inner JSON Schema. - pub(crate) fn from_inner_schema(prompt_template: &str, schema: Value) -> Result { + /// Builds a contract around a user-supplied inner JSON Schema. + /// + /// The verdict is validated against the schema locally in both modes. JSON Schema mode sends it through the + /// provider's strict wrapper; JSON Object mode appends it to the prompt, because the + /// provider then guarantees only that the reply is JSON. + pub(crate) fn from_inner_schema( + prompt_template: &str, + schema: Value, + response_format_type: ClassifierResponseFormat, + ) -> Result { + validate_prompt(prompt_template)?; if schema.get("json_schema").is_some() { return Err(LibsyError::AlgorithmError { message: @@ -121,32 +126,53 @@ impl ClassifierContract { }); } let validator = compile_schema(&schema)?; - Self::from_response_format( - prompt_template, - json!({ - "type": "json_schema", - "json_schema": { - "name": "switchyard_classifier_response", - "strict": true, - "schema": schema, + match response_format_type { + ClassifierResponseFormat::JsonSchema => Ok(Self::from_response_format( + prompt_template.to_string(), + json!({ + "type": "json_schema", + "json_schema": { + "name": "switchyard_classifier_response", + "strict": true, + "schema": schema, + } + }), + Some(validator), + )), + ClassifierResponseFormat::JsonObject => { + // Catches the common case of a root `type` that excludes objects: the provider + // only returns objects, so every verdict would fail and route to the default. + // A schema with no `type` (e.g. `{"enum": [1, 2]}`) is not checked and can + // still never match; that is left to the author. + let is_object_allowed = match schema.get("type") { + Some(Value::String(root)) => root == "object", + Some(Value::Array(roots)) => roots.iter().any(|root| root == "object"), + _ => true, + }; + if !is_object_allowed { + return Err(algorithm_error( + "response_schema must describe a JSON object in json_object mode", + )); } - }), - Some(validator), - ) + Ok(Self::from_response_format( + schema_in_prompt(prompt_template, &schema)?, + json!({"type": "json_object"}), + Some(validator), + )) + } + } } fn from_response_format( - prompt_template: &str, + system_prompt: String, response_format: Value, validator: Option, - ) -> Result { - validate_prompt(prompt_template)?; - - Ok(Self { - system_prompt: prompt_template.to_string(), + ) -> Self { + Self { + system_prompt, response_format, validator, - }) + } } pub(crate) fn system_prompt(&self) -> &str { @@ -187,6 +213,19 @@ pub(super) fn validate_prompt(prompt_template: &str) -> Result<()> { Ok(()) } +/// Appends the verdict schema to a prompt for JSON Object mode. +/// +/// Callers validate the template, not this result: the appended schema makes the prompt +/// non-empty and may itself contain `{{RESPONSE_SCHEMA}}` text. +fn schema_in_prompt(prompt_template: &str, schema: &Value) -> Result { + let rendered_schema = serde_json::to_string_pretty(schema).map_err(|error| { + algorithm_error(format!("response schema could not be rendered: {error}")) + })?; + Ok(format!( + "{prompt_template}\n\nReturn exactly one JSON object matching this JSON Schema:\n{rendered_schema}" + )) +} + fn compile_schema(schema: &Value) -> Result { if !schema.is_object() { return Err(algorithm_error("response_schema must be a JSON object")); @@ -307,6 +346,7 @@ mod tests { "required": ["decision"], "additionalProperties": false }), + ClassifierResponseFormat::JsonSchema, )?; assert_eq!( @@ -335,12 +375,110 @@ mod tests { #[test] fn a_provider_wrapper_is_rejected_as_an_inner_schema() { + for response_format_type in [ + ClassifierResponseFormat::JsonSchema, + ClassifierResponseFormat::JsonObject, + ] { + let error = ClassifierContract::from_inner_schema( + "classify", + json!({"json_schema": {"schema": {"type": "object"}}}), + response_format_type, + ) + .expect_err("a wrapped schema should error"); + + assert!( + error.to_string().contains("inner JSON Schema"), + "{response_format_type:?}: {error}" + ); + } + } + + #[test] + fn a_custom_contract_can_request_a_json_object() -> Result<()> { + let contract = ClassifierContract::from_inner_schema( + "Choose a target.", + json!({ + "type": "object", + "properties": {"target": {"type": "string", "enum": ["sonnet", "opus"]}}, + "required": ["target"], + "additionalProperties": false + }), + ClassifierResponseFormat::JsonObject, + )?; + + assert_eq!(contract.response_format(), &json!({"type": "json_object"})); + assert!(contract.system_prompt().starts_with("Choose a target.")); + assert!(contract.system_prompt().contains("JSON Schema")); + assert!(contract.system_prompt().contains("\"target\"")); + // The provider no longer enforces the schema, so the local validator must. + contract.validate_verdict(&json!({"target": "sonnet"}))?; + assert!( + contract + .validate_verdict(&json!({"target": "unknown"})) + .is_err() + ); + Ok(()) + } + + #[test] + fn json_object_mode_accepts_only_schemas_whose_root_allows_an_object() { + for (schema, accepted) in [ + (json!({"type": "object"}), true), + (json!({"type": ["object"]}), true), + (json!({"type": ["object", "null"]}), true), + (json!({"type": "string"}), false), + (json!({"type": "array"}), false), + (json!({"type": ["string", "null"]}), false), + ] { + let result = ClassifierContract::from_inner_schema( + "classify", + schema.clone(), + ClassifierResponseFormat::JsonObject, + ); + + match result { + Ok(_) => assert!(accepted, "{schema} should be rejected"), + Err(error) => { + assert!(!accepted, "{schema}: {error}"); + assert!( + error.to_string().contains("must describe a JSON object"), + "{schema}: {error}" + ); + } + } + } + } + + #[test] + fn a_schema_that_mentions_the_placeholder_is_accepted_in_both_modes() { + let schema = json!({ + "type": "object", + "properties": {"target": {"type": "string"}}, + "required": ["target"], + "description": "{{RESPONSE_SCHEMA}}" + }); + for response_format_type in [ + ClassifierResponseFormat::JsonSchema, + ClassifierResponseFormat::JsonObject, + ] { + ClassifierContract::from_inner_schema("pick", schema.clone(), response_format_type) + .unwrap_or_else(|error| panic!("{response_format_type:?}: {error}")); + } + } + + #[test] + fn an_empty_prompt_is_rejected_in_json_object_mode() { let error = ClassifierContract::from_inner_schema( - "classify", - json!({"json_schema": {"schema": {"type": "object"}}}), + " ", + json!({"type": "object"}), + ClassifierResponseFormat::JsonObject, ) - .expect_err("provider wrapper should be rejected"); + .expect_err("an empty prompt should be rejected before the schema is appended"); - assert!(error.to_string().contains("inner JSON Schema")); + assert!( + error + .to_string() + .contains("classifier prompt must not be empty") + ); } } diff --git a/crates/libsy/src/algorithms/util/escalation.rs b/crates/libsy/src/algorithms/util/escalation.rs index 0335b4d3f..2160cb7ec 100644 --- a/crates/libsy/src/algorithms/util/escalation.rs +++ b/crates/libsy/src/algorithms/util/escalation.rs @@ -94,8 +94,9 @@ impl DeescalationConfig { #[derive(Clone, Debug, Deserialize)] #[serde(default, deny_unknown_fields)] pub struct EscalationJudgeConfig { - /// Consecutive escalate verdicts required before a turn moves to the capable tier, which - /// is also the turn that latches the session. Any decline clears the streak. + /// Consecutive fresh-evidence verdicts in the same category required before a turn moves to + /// the capable tier, which is also the turn that latches the session. Any decline or stale + /// evidence clears the streak. /// `1` escalates on the first verdict; the router's main cost dial. /// `2` or higher needs a session id, since the streak is retained per session. pub confirmations: u32, @@ -157,15 +158,39 @@ impl EvaluationPhase { } } -/// The judge's verdict. The schema also requires a `reason`, which makes the judge state its -/// case and measurably sharpens the verdict. Routing reads only the boolean; the reason is -/// kept solely so an operator can see why the judge held or escalated when the -/// `switchyard_libsy::algorithms::util::escalation` target is enabled at `debug`. +/// Bounded trouble pattern used to correlate escalation confirmations across turns. +#[derive(Clone, Copy, Debug, Deserialize, Eq, PartialEq)] +#[serde(rename_all = "snake_case")] +pub(crate) enum EscalationCategory { + None, + Repetition, + FalseProgress, + Drift, + Desperation, + CapabilityGap, +} + +impl EscalationCategory { + /// Stable state and telemetry label. + pub(crate) const fn label(self) -> &'static str { + match self { + Self::None => "none", + Self::Repetition => "repetition", + Self::FalseProgress => "false_progress", + Self::Drift => "drift", + Self::Desperation => "desperation", + Self::CapabilityGap => "capability_gap", + } + } +} + +/// The judge's typed verdict, including the evidence needed to confirm a stable pattern. #[derive(Deserialize)] pub(crate) struct EscalationVerdict { - escalate: bool, - #[serde(default)] - reason: String, + pub(crate) escalate: bool, + pub(crate) category: EscalationCategory, + pub(crate) new_evidence: bool, + pub(crate) reason: String, } /// Builds the condensed trajectory presented to the escalation judge. @@ -314,21 +339,102 @@ pub(crate) fn conversation_turn(request: &Request) -> usize { /// relies on. fn message_text(message: &Message) -> String { let mut parts = Vec::new(); - collect_text(&message.content, &mut parts); + let terminus_commands = if message.role == Role::Assistant { + message + .content + .iter() + .filter_map(|block| match block { + ContentBlock::ToolCall(call) if call.name == "bash_command" => call + .arguments + .get("keystrokes") + .and_then(|value| value.as_str()), + _ => None, + }) + .collect::>() + } else { + Vec::new() + }; + collect_text(&message.content, &mut parts, &terminus_commands); parts.join(" ") } +/// Removes a Terminus command batch when a structured bash call carries the same action. +/// +/// The model-facing request remains untouched. Only the judge's plain-text view is normalized, +/// so one action cannot look like two attempts while the agent still sees its native history. +fn without_duplicated_terminus_commands(text: &str, tool_commands: &[&str]) -> String { + let mut normalized_text = String::with_capacity(text.len()); + let mut unmatched_tool_commands = tool_commands.to_vec(); + let mut copied_through = 0; + let mut scan_from = 0; + + while let Some(relative_start) = text[scan_from..].find('{') { + let start = scan_from + relative_start; + + let mut values = + serde_json::Deserializer::from_str(&text[start..]).into_iter::(); + let Some(Ok(mut value)) = values.next() else { + scan_from = start + 1; + continue; + }; + let end = start + values.byte_offset(); + let Some(commands) = value + .get("commands") + .and_then(|commands| commands.as_array()) + else { + scan_from = start + 1; + continue; + }; + let Some(command_batch) = commands + .iter() + .map(|command| command.get("keystrokes").and_then(|value| value.as_str())) + .collect::>>() + else { + scan_from = start + 1; + continue; + }; + let mut remaining_tool_commands = unmatched_tool_commands.clone(); + let fully_encoded = command_batch.iter().all(|command| { + let Some(index) = remaining_tool_commands + .iter() + .position(|candidate| candidate == command) + else { + return false; + }; + remaining_tool_commands.swap_remove(index); + true + }); + if command_batch.is_empty() || !fully_encoded { + scan_from = start + 1; + continue; + } + + value["commands"] = serde_json::Value::Array(Vec::new()); + let Ok(normalized) = serde_json::to_string(&value) else { + scan_from = start + 1; + continue; + }; + normalized_text.push_str(&text[copied_through..start]); + normalized_text.push_str(&normalized); + copied_through = end; + scan_from = end; + unmatched_tool_commands = remaining_tool_commands; + } + normalized_text.push_str(&text[copied_through..]); + normalized_text +} + /// Appends the judge-relevant text of each block, descending into tool results. -fn collect_text(content: &[ContentBlock], parts: &mut Vec) { +fn collect_text(content: &[ContentBlock], parts: &mut Vec, tool_commands: &[&str]) { for block in content { match block { ContentBlock::Text { text } | ContentBlock::Refusal { text } => { - parts.push(text.clone()); + parts.push(without_duplicated_terminus_commands(text, tool_commands)); } ContentBlock::ToolCall(call) => { parts.push(format!("tool_call {}({})", call.name, call.arguments)); } - ContentBlock::ToolResult(result) => collect_text(&result.content, parts), + ContentBlock::ToolResult(result) => collect_text(&result.content, parts, &[]), _ => {} } } @@ -378,7 +484,7 @@ fn summarize_for_judge( for instruction in instructions { let mut parts = Vec::new(); - collect_text(&instruction.content, &mut parts); + collect_text(&instruction.content, &mut parts, &[]); instruction_anchors.push(format!( "[{}] {}", role_label(instruction.role), @@ -648,6 +754,194 @@ mod tests { assert_eq!(message_text(&result), "no such file"); } + /// A raw command batch fully mirrored by structured tool calls is emptied in the judge view. + #[test] + fn message_text_deduplicates_terminus_commands_for_the_judge() { + let first_command = "grep -n bug app.py\n"; + let second_command = "sed -n '1,80p' app.py\n"; + let message = Message { + role: Role::Assistant, + content: vec![ + ContentBlock::Text { + text: format!( + "Before\n```json\n{}\n```\nAfter", + json!({ + "analysis": "inspect the reported file", + "commands": [ + {"keystrokes": first_command, "duration": 0.1}, + {"keystrokes": second_command, "duration": 0.1}, + ], + }) + ), + }, + ContentBlock::ToolCall(ToolCall { + id: "call-1".to_string(), + name: "bash_command".to_string(), + arguments: json!({"keystrokes": first_command, "duration": 0.1}), + }), + ContentBlock::ToolCall(ToolCall { + id: "call-2".to_string(), + name: "bash_command".to_string(), + arguments: json!({"keystrokes": second_command, "duration": 0.1}), + }), + ], + }; + + let text = message_text(&message); + + assert!(text.contains("inspect the reported file"), "{text}"); + assert!(text.contains("Before"), "{text}"); + assert!(text.contains("After"), "{text}"); + assert!(text.contains(r#""commands":[]"#), "{text}"); + assert_eq!(text.matches("grep -n bug app.py").count(), 1, "{text}"); + assert_eq!(text.matches("sed -n '1,80p' app.py").count(), 1, "{text}"); + assert_eq!(text.matches("tool_call bash_command(").count(), 2, "{text}"); + } + + /// Each structured tool call can absorb only one rendered batch, so a repeated batch stays. + #[test] + fn message_text_deduplicates_multiple_batches_once_per_tool_call() { + let first_command = "grep -n bug app.py\n"; + let second_command = "sed -n '1,80p' app.py\n"; + let batch = |command| { + json!({ + "analysis": "inspect", + "commands": [{"keystrokes": command}], + }) + .to_string() + }; + let message = Message { + role: Role::Assistant, + content: vec![ + ContentBlock::Text { + text: format!( + "First {} second {} repeated {}", + batch(first_command), + batch(second_command), + batch(first_command) + ), + }, + ContentBlock::ToolCall(ToolCall { + id: "call-1".to_string(), + name: "bash_command".to_string(), + arguments: json!({"keystrokes": first_command}), + }), + ContentBlock::ToolCall(ToolCall { + id: "call-2".to_string(), + name: "bash_command".to_string(), + arguments: json!({"keystrokes": second_command}), + }), + ], + }; + + let text = message_text(&message); + + assert_eq!(text.matches(r#""commands":[]"#).count(), 2, "{text}"); + assert_eq!(text.matches("grep -n bug app.py").count(), 2, "{text}"); + assert_eq!(text.matches("sed -n '1,80p' app.py").count(), 1, "{text}"); + } + + /// A rendered batch stays when no structured tool call carries the same command. + #[test] + fn message_text_keeps_terminus_commands_without_matching_tool_call() { + let command = "grep -n bug app.py\n"; + let text = json!({ + "analysis": "inspect the reported file", + "commands": [{"keystrokes": command}], + }) + .to_string(); + let without_tool_call = Message { + role: Role::Assistant, + content: vec![ContentBlock::Text { text: text.clone() }], + }; + let mismatched_tool_call = Message { + role: Role::Assistant, + content: vec![ + ContentBlock::Text { text }, + ContentBlock::ToolCall(ToolCall { + id: "call-1".to_string(), + name: "bash_command".to_string(), + arguments: json!({"keystrokes": "sed -n '1,20p' app.py\n", "duration": 0.1}), + }), + ], + }; + + assert_eq!( + message_text(&without_tool_call) + .matches("grep -n bug app.py") + .count(), + 1 + ); + assert_eq!( + message_text(&mismatched_tool_call) + .matches("grep -n bug app.py") + .count(), + 1 + ); + } + + /// A batch stays intact when only some of its commands have matching tool calls. + #[test] + fn message_text_keeps_a_partially_encoded_terminus_batch() { + let first_command = "grep -n bug app.py\n"; + let second_command = "sed -n '1,80p' app.py\n"; + let message = Message { + role: Role::Assistant, + content: vec![ + ContentBlock::Text { + text: json!({ + "commands": [ + {"keystrokes": first_command}, + {"keystrokes": second_command}, + ], + }) + .to_string(), + }, + ContentBlock::ToolCall(ToolCall { + id: "call-1".to_string(), + name: "bash_command".to_string(), + arguments: json!({"keystrokes": first_command}), + }), + ], + }; + + let text = message_text(&message); + + assert_eq!(text.matches("grep -n bug app.py").count(), 2, "{text}"); + assert_eq!(text.matches("sed -n '1,80p' app.py").count(), 1, "{text}"); + assert!(!text.contains(r#""commands":[]"#), "{text}"); + } + + /// Duplicate commands in one batch need one tool call each before the batch is removed. + #[test] + fn message_text_keeps_duplicate_commands_without_one_tool_call_each() { + let command = "grep -n bug app.py\n"; + let message = Message { + role: Role::Assistant, + content: vec![ + ContentBlock::Text { + text: json!({ + "commands": [ + {"keystrokes": command}, + {"keystrokes": command}, + ], + }) + .to_string(), + }, + ContentBlock::ToolCall(ToolCall { + id: "call-1".to_string(), + name: "bash_command".to_string(), + arguments: json!({"keystrokes": command}), + }), + ], + }; + + let text = message_text(&message); + + assert_eq!(text.matches("grep -n bug app.py").count(), 3, "{text}"); + assert!(!text.contains(r#""commands":[]"#), "{text}"); + } + #[test] fn truncate_middle_keeps_head_and_tail() { let text = "a".repeat(40) + &"z".repeat(40); diff --git a/crates/libsy/src/algorithms/util/llm_judge.rs b/crates/libsy/src/algorithms/util/llm_judge.rs index 4798178e9..cb8277dae 100644 --- a/crates/libsy/src/algorithms/util/llm_judge.rs +++ b/crates/libsy/src/algorithms/util/llm_judge.rs @@ -261,7 +261,7 @@ where /// mid-stream, or unparseable reply — is logged and folded into `None` for the policy's /// fallback branch. A closed driver stream is folded too; the algorithm's next driver /// call surfaces it, so nothing is masked. - async fn verdict( + pub(crate) async fn verdict( &self, state: &mut State, request: &Request, diff --git a/crates/libsy/src/prompts/escalation/deescalation.md b/crates/libsy/src/prompts/escalation/deescalation.md index 3d1e3fe9c..4f9870829 100644 --- a/crates/libsy/src/prompts/escalation/deescalation.md +++ b/crates/libsy/src/prompts/escalation/deescalation.md @@ -27,5 +27,12 @@ The routing input begins with one of these router-generated markers: success, or a context compaction. A strong-tier turn that merely reads files or plans is not by itself evidence that the trouble is resolved. + In this phase the router reads only `escalate`. Fill the other fields + consistently anyway: when retaining, set `category` to the trouble + pattern that caused the escalation and `new_evidence` to true only when + the newest strong-tier turn or its tool output shows that trouble is + still active; when releasing, return `category: "none"` and + `new_evidence: false`. + The router, not the judge, applies confirmation counts and decides when to change tiers. Judge only the phase named in the routing input. diff --git a/crates/libsy/src/prompts/escalation/prompt.md b/crates/libsy/src/prompts/escalation/prompt.md index 6d57369bb..d29a4017f 100644 --- a/crates/libsy/src/prompts/escalation/prompt.md +++ b/crates/libsy/src/prompts/escalation/prompt.md @@ -9,11 +9,20 @@ messages and tool results). Judge the *trajectory* — is the agent making real progress toward the stated task — not the difficulty of the task itself. Return exactly one JSON object: -{"escalate": boolean, "reason": "one short sentence naming the pattern"} +{"escalate": boolean, "category": "none|repetition|false_progress|drift|desperation|capability_gap", "new_evidence": boolean, "reason": "one short sentence naming the pattern"} + +Use `category: "none"` whenever `escalate` is false. `new_evidence` +is true only when the NEWEST assistant turn or its resulting tool output +adds evidence for the named pattern. Older turns may establish context, +but do not repeat an escalation vote solely because old trouble remains +visible in the rolling transcript. If the newest turn recovered, adapted, +or made progress, return `escalate: false`, `category: "none"`, and +`new_evidence: false`. Escalation is one-way for the rest of the task and expensive. Escalate only on a clear PATTERN of trouble, never on a single failed command. -When the evidence is thin or ambiguous, return {"escalate": false}. +When the evidence is thin or ambiguous, decline with `category: "none"` +and `new_evidence: false`. The bar is not "is there friction" — agentic coding is full of friction the efficient tier works through on its own. The bar is "is this run @@ -58,8 +67,10 @@ Hold weak — no model can fix these, so escalation is pure waste: # Trouble patterns — escalate when you see these Repetition and loops (the most common way agent runs die): -- The same command or edit failing 2+ times with materially the same - error, especially with unrelated changes in between. +- The same command or edit failing across 2+ DISTINCT assistant turns + with materially the same error, especially with unrelated changes in + between. Count executed attempts across turns, not repeated renderings + inside one message. - Near-identical tool calls repeated, or the same files re-read, without new information gained — including longer cycles (A -> B -> C -> A). - Fighting the environment: repeatedly invoking a missing executable, @@ -105,6 +116,15 @@ Desperation: # Expected friction — do NOT escalate on these Agentic coding is full of failures that are part of healthy work: +- A command appearing both in the assistant's JSON/text and as one or + more structured tool-call blocks in that SAME turn. Agent harnesses + commonly serialize one intended action more than once; this is one + attempt, not a loop. Repetition is evidence only when separate turns + show separate executions with materially the same failed result. +- Terminal-input serialization trouble (for example tabs triggering + completion or a heredoc being mangled) when the next turn changes the + write mechanism, quoting, or transport. That is adaptation, not + repeated failure. - A test written to fail first (TDD) or a bug being reproduced on purpose. - A compile, lint, or test error fixed or meaningfully acted on in the @@ -130,49 +150,56 @@ Agentic coding is full of failures that are part of healthy work: - A long-running command (build, install, test suite) that simply has not finished, or the agent waiting on information it asked for. -The distinguishing question: is each failure producing new information -that changes the next action? Failing forward is fine; failing in place -is trouble. Also weigh the session's own recovery record: if this same -session already shows friction the agent subsequently cleared (a failure -followed by a verified fix or passing check), lean toward holding — a -session that has recovered before will usually recover again. +The distinguishing question: across DISTINCT assistant turns, is each +failure producing new information that changes the next action? Failing +forward is fine; failing in place is trouble. Never infer a multi-turn +pattern from duplicated representations inside one turn. Also weigh the +session's own recovery record: if this same session already shows +friction the agent subsequently cleared (a failure followed by a verified +fix or passing check), lean toward holding — a session that has recovered +before will usually recover again. # Worked examples (none drawn from any benchmark task set) * Turn 3; the agent ran the test suite, 4 tests fail, and it is now - reading the first failing test. -> {"escalate": false} — reproducing - failures is the job. + reading the first failing test. -> {"escalate": false, "category": + "none", "new_evidence": false, "reason": "working through the first + reproduced failure"} * The agent has run `pytest tests/test_api.py` 4 times with the same ImportError, editing an unrelated config file between attempts. -> - {"escalate": true, "reason": "same ImportError 4 times while editing - unrelated files"} + {"escalate": true, "category": "repetition", "new_evidence": true, + "reason": "same ImportError 4 times while editing unrelated files"} * `conda` is not installed; the agent has tried `conda install` five ways instead of using the `pip` that earlier output showed present. - -> {"escalate": true, "reason": "fighting missing executable instead - of adapting"} + -> {"escalate": true, "category": "repetition", "new_evidence": + true, "reason": "fighting missing executable instead of adapting"} * Task: "make the provided integration tests pass." Recent turns: renaming variables and reformatting docstrings; tests not run in 8 - turns. -> {"escalate": true, "reason": "drifted to cosmetic edits, - verification abandoned"} + turns. -> {"escalate": true, "category": "drift", "new_evidence": + true, "reason": "drifted to cosmetic edits, verification abandoned"} * The agent says "All tests pass, task complete" but the last visible test output shows "2 failed, 11 passed". -> {"escalate": true, - "reason": "claims success contradicted by latest test output"} + "category": "false_progress", "new_evidence": true, "reason": + "claims success contradicted by latest test output"} * The agent wrote a reproduction script that exits 0 without invoking the code path the issue describes, concluded "bug not reproducible", - and is wrapping up. -> {"escalate": true, "reason": "reproduction - never exercised the reported code path"} + and is wrapping up. -> {"escalate": true, "category": + "false_progress", "new_evidence": true, "reason": "reproduction never + exercised the reported code path"} * Two turns of edits, one failed build, then a fixed build and a - passing test. -> {"escalate": false} + passing test. -> {"escalate": false, "category": "none", + "new_evidence": false, "reason": "latest build recovered and passed"} * `npm install` has been running for one turn with no output yet. -> - {"escalate": false} — slow command, not a stall. + {"escalate": false, "category": "none", "new_evidence": false, + "reason": "command is still running"} * Four different serialization libraries failed to import; the agent is now writing the converter with a fifth approach it has not tried - before. -> {"escalate": false} — sequential alternatives are - adaptation, even when none has succeeded yet. + before. -> {"escalate": false, "category": "none", "new_evidence": + false, "reason": "latest turn changed approach"} * Task: tune a slow batch pipeline. The agent is investigating why the message broker fails to start, since the pipeline cannot be - measured without it. -> {"escalate": false} — unblocking - verification serves the task. + measured without it. -> {"escalate": false, "category": "none", + "new_evidence": false, "reason": "working to unblock verification"} Do not emit markdown, commentary, or chain-of-thought — only the JSON object. diff --git a/crates/libsy/src/prompts/escalation/schema.json b/crates/libsy/src/prompts/escalation/schema.json index 977c41882..e98b3c5fd 100644 --- a/crates/libsy/src/prompts/escalation/schema.json +++ b/crates/libsy/src/prompts/escalation/schema.json @@ -10,12 +10,21 @@ "type": "boolean", "description": "True when the run is likely doomed without escalation to the strong tier." }, + "category": { + "type": "string", + "enum": ["none", "repetition", "false_progress", "drift", "desperation", "capability_gap"], + "description": "The single trouble pattern supporting escalation, or none when escalation is false." + }, + "new_evidence": { + "type": "boolean", + "description": "True only when the newest assistant turn or its resulting tool output adds evidence for the named trouble pattern." + }, "reason": { "type": "string", "description": "One short sentence naming the trouble pattern, or stating why the run is progressing." } }, - "required": ["escalate", "reason"], + "required": ["escalate", "category", "new_evidence", "reason"], "additionalProperties": false } } diff --git a/crates/protocol/src/decision.rs b/crates/protocol/src/decision.rs new file mode 100644 index 000000000..7d1ea28a8 --- /dev/null +++ b/crates/protocol/src/decision.rs @@ -0,0 +1,136 @@ +// SPDX-FileCopyrightText: Copyright (c) 2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved. +// SPDX-License-Identifier: Apache-2.0 + +//! Provider-neutral decision questions and answers, separate from LLM messages. +//! +//! Enums use snake-case `type` tags and a `data` payload in serialized form. +//! Fields are public; providers and callers are responsible for valid values. + +use std::collections::BTreeMap; + +use serde::{Deserialize, Serialize}; +use serde_json::Value; + +use crate::{ModelId, Usage}; + +/// Answer probability on a `[0, 1]` scale. +#[derive(Clone, Copy, Debug, PartialEq, Serialize, Deserialize)] +#[serde(transparent)] +pub struct Probability(pub f64); + +/// Position in the request's rubric, including fractional positions. +#[derive(Clone, Copy, Debug, PartialEq, Serialize, Deserialize)] +#[serde(transparent)] +pub struct ScoreValue(pub f64); + +/// Provider confidence; its scale and meaning are provider-specific. +#[derive(Clone, Copy, Debug, PartialEq, Serialize, Deserialize)] +#[serde(transparent)] +pub struct ProviderConfidence(pub f64); + +/// Shared context evaluated against independent, named questions. +#[derive(Clone, Debug, PartialEq, Serialize, Deserialize)] +pub struct DecisionRequest { + /// Optional until a target is selected. + pub model: Option, + /// Conversation, application state, or other material to evaluate. + pub context: Value, + /// Independent questions keyed by their IDs. + pub questions: BTreeMap, +} + +/// Instructions and the expected answer shape for one question. +#[derive(Clone, Debug, PartialEq, Serialize, Deserialize)] +pub struct DecisionQuestion { + /// Structured or textual instructions shared with the provider. + pub instructions: Value, + /// Expected answer shape. + pub kind: DecisionKind, +} + +/// The answer shape and any options or ordered rubric levels. +#[derive(Clone, Debug, PartialEq, Serialize, Deserialize)] +#[serde(tag = "type", content = "data", rename_all = "snake_case")] +pub enum DecisionKind { + /// A Boolean judgment or probability of true. + Boolean { + /// Meaning of a true answer, when needed. + true_description: Option, + /// Meaning of a false answer, when needed. + false_description: Option, + }, + /// Select one of the declared options. + Choice { + /// Nonempty options with unique IDs, preserving caller order. + options: Vec, + }, + /// A position on an ordered rubric, not an arbitrary numeric measurement. + Score { + /// At least two levels, ordered low to high and indexed from zero. + levels: Vec, + }, +} + +/// An identified choice with an optional structured description. +#[derive(Clone, Debug, PartialEq, Serialize, Deserialize)] +pub struct ChoiceOption { + /// Stable identifier used by choice answers and distributions. + pub id: String, + /// Meaning of this option, when its ID alone is insufficient. + pub description: Option, +} + +/// Answers keyed by the matching request's question IDs. +#[derive(Clone, Debug, PartialEq, Serialize, Deserialize)] +pub struct DecisionResponse { + /// Provider-reported response identifier. + pub id: Option, + /// Provider-reported model identifier. + pub model: Option, + /// Typed answers corresponding to the request's questions. + pub answers: BTreeMap, + /// Available token counts; absent counts remain unknown. + #[serde(default)] + pub usage: Usage, +} + +/// An answer and separate, optional provider confidence. +#[derive(Clone, Debug, PartialEq, Serialize, Deserialize)] +pub struct DecisionAnswer { + /// The estimate for the matching question. + pub value: DecisionValue, + /// Provider-specific confidence, distinct from answer probabilities. + pub provider_confidence: Option, +} + +/// A typed estimate; missing distributions remain unknown. +#[derive(Clone, Debug, PartialEq, Serialize, Deserialize)] +#[serde(tag = "type", content = "data", rename_all = "snake_case")] +pub enum DecisionValue { + /// A Boolean judgment or probability, without an implicit threshold. + Boolean(BooleanEstimate), + /// One selected option with an optional complete distribution. + Choice { + /// Must name an option in the matching question. + selected: String, + /// Maps every declared option ID to its probability when available. + probabilities: Option>, + }, + /// A fractional position in the matching request's rubric. + Score { + /// Must lie in `0..=N-1` for the request's N levels. + value: ScoreValue, + /// Follows the request's level order. Retain that request to interpret it. + probabilities: Option>, + }, +} + +/// Preserves Boolean-only answers without inventing probability or certainty. +#[derive(Clone, Copy, Debug, PartialEq, Serialize, Deserialize)] +#[serde(tag = "type", content = "data", rename_all = "snake_case")] +pub enum BooleanEstimate { + /// A Boolean judgment with no probability supplied. + Value(bool), + /// Probability of true; algorithms choose their own thresholds. + ProbabilityTrue(Probability), +} diff --git a/crates/protocol/src/lib.rs b/crates/protocol/src/lib.rs index 0d76b4d17..b32dc6c43 100644 --- a/crates/protocol/src/lib.rs +++ b/crates/protocol/src/lib.rs @@ -7,6 +7,7 @@ pub mod category; pub mod client; pub mod codex_namespaces; +pub mod decision; pub mod envelope; pub mod format; pub mod llm; @@ -16,6 +17,7 @@ pub mod stream; pub use category::*; pub use client::*; +pub use decision::*; pub use envelope::*; pub use format::*; pub use llm::*; diff --git a/crates/switchyard-py/src/libsy_bindings.rs b/crates/switchyard-py/src/libsy_bindings.rs index 23c90d7ad..b049c5429 100644 --- a/crates/switchyard-py/src/libsy_bindings.rs +++ b/crates/switchyard-py/src/libsy_bindings.rs @@ -218,7 +218,8 @@ impl PyCustomClassifierConfig { session_affinity=false, message_hash_fallback=false, recent_turn_window=None, - max_output_tokens=4096 + max_output_tokens=4096, + response_format_type="json_schema" ))] #[allow(clippy::too_many_arguments)] fn new( @@ -229,6 +230,7 @@ impl PyCustomClassifierConfig { message_hash_fallback: bool, recent_turn_window: Option, max_output_tokens: u64, + response_format_type: &str, ) -> PyResult { // Convert the Python schema into serde JSON and pair it with the target-selector policy; // conversion failures propagate to Python through `PyResult`. @@ -241,6 +243,7 @@ impl PyCustomClassifierConfig { inner.message_hash_fallback = message_hash_fallback; inner.recent_turn_window = recent_turn_window; inner.max_output_tokens = max_output_tokens; + inner.response_format_type = parse_response_format_type(response_format_type)?; Ok(Self { inner }) } } @@ -348,16 +351,17 @@ fn classifier_contract( if let Some(prompt) = prompt { contract = contract.with_prompt(prompt); } - let response_format_type = match response_format_type { - "json_schema" => ClassifierResponseFormat::JsonSchema, - "json_object" => ClassifierResponseFormat::JsonObject, - other => { - return Err(PyValueError::new_err(format!( - "response_format_type must be 'json_schema' or 'json_object', got {other:?}" - ))); - } - }; - Ok(contract.with_response_format_type(response_format_type)) + Ok(contract.with_response_format_type(parse_response_format_type(response_format_type)?)) +} + +fn parse_response_format_type(response_format_type: &str) -> PyResult { + match response_format_type { + "json_schema" => Ok(ClassifierResponseFormat::JsonSchema), + "json_object" => Ok(ClassifierResponseFormat::JsonObject), + other => Err(PyValueError::new_err(format!( + "response_format_type must be 'json_schema' or 'json_object', got {other:?}" + ))), + } } /// Judge target and policy used when stage-router signals are inconclusive. diff --git a/crates/switchyard-runner/src/algorithm.rs b/crates/switchyard-runner/src/algorithm.rs index e742d48ba..187030a48 100644 --- a/crates/switchyard-runner/src/algorithm.rs +++ b/crates/switchyard-runner/src/algorithm.rs @@ -133,6 +133,7 @@ struct CustomClassifierRouteConfig { message_hash_fallback: bool, recent_turn_window: Option, max_output_tokens: u64, + response_format_type: ClassifierResponseFormat, } /// Runtime model groups for a custom classifier, keyed by group name. @@ -993,10 +994,9 @@ impl LlmClassifierRouteConfig { || base_threshold.is_some() || threshold_step.is_some() || escalation.is_some() - || *response_format_type != ClassifierResponseFormat::JsonSchema { return Err(AlgorithmConfigError::new(format!( - "llm_classifier route {route_name} mode custom cannot use capability or escalation fields and response_format_type must be 'json_schema'" + "llm_classifier route {route_name} mode custom cannot use capability or escalation fields" ))); } let models = required_classifier_field(route_name, "models", models)?; @@ -1032,6 +1032,7 @@ impl LlmClassifierRouteConfig { message_hash_fallback: *message_hash_fallback, recent_turn_window: *recent_turn_window, max_output_tokens: *max_output_tokens, + response_format_type: *response_format_type, }, )) } @@ -1118,6 +1119,7 @@ fn build_subagent_router_config( ); classifier_config.recent_turn_window = config.recent_turn_window; classifier_config.max_output_tokens = config.max_output_tokens; + classifier_config.response_format_type = config.response_format_type; let classifier = Arc::new( LlmTaskClassifier::new(LlmClassifierConfig::Custom { default_target: config.default_target.clone(), @@ -1271,6 +1273,7 @@ fn build_algorithm( classifier_config.message_hash_fallback = config.message_hash_fallback; classifier_config.recent_turn_window = config.recent_turn_window; classifier_config.max_output_tokens = config.max_output_tokens; + classifier_config.response_format_type = config.response_format_type; LlmTaskClassifier::new(LlmClassifierConfig::Custom { default_target: config.default_target, config: classifier_config, diff --git a/crates/switchyard-runner/src/config.rs b/crates/switchyard-runner/src/config.rs index f37b4cda9..2a9e76753 100644 --- a/crates/switchyard-runner/src/config.rs +++ b/crates/switchyard-runner/src/config.rs @@ -1289,6 +1289,35 @@ new = ["send_message"] ); } + #[test] + fn mode_custom_accepts_json_object_output() -> RunnerResult<()> { + let top_level = format!( + r#"{VALID_CONFIG} +[routes.custom] +id = "switchyard/custom" +type = "llm_classifier" +mode = "custom" +response_format_type = "json_object" +models = {{ judge = ["classifier"], capable = ["strong"], efficient = ["weak"], any = ["strong", "weak"] }} +default_target = "efficient" +prompt = "Select a target for this task." +response_schema = '{{"type":"object","properties":{{"target":{{"type":"string","enum":["capable","efficient"]}}}},"required":["target"],"additionalProperties":false}}' +policy = {{ type = "target_selector", selector = "/target" }} +"# + ); + let subagent = with_subagent_llm_classifier( + VALID_CONFIG, + "passthrough", + "\nresponse_format_type = \"json_object\"", + ); + + for (name, config) in [("top-level", top_level), ("subagent", subagent)] { + runner_from_toml(&config) + .map_err(|error| RunnerError::configuration(format!("{name}: {error}")))?; + } + Ok(()) + } + #[test] fn stage_router_rejects_an_unknown_field() { let config = stage_config().replace( diff --git a/crates/switchyard-server/tests/server.rs b/crates/switchyard-server/tests/server.rs index cd21d3c91..727916a79 100644 --- a/crates/switchyard-server/tests/server.rs +++ b/crates/switchyard-server/tests/server.rs @@ -371,9 +371,15 @@ async fn upstream_chat( .into_response(); } + // JSON Object mode carries the custom schema in the judge prompt instead of + // `response_format`, so recognize it there too. let custom_target_schema = body .pointer("/response_format/json_schema/schema/properties/decision/properties/target") - .is_some(); + .is_some() + || (body["response_format"]["type"] == "json_object" + && body["messages"][0]["content"] + .as_str() + .is_some_and(|prompt| prompt.contains("\"decision\""))); let requests_invalid_verdict = body["messages"].as_array().is_some_and(|messages| { messages.iter().any(|message| { message["content"] @@ -399,7 +405,10 @@ async fn upstream_chat( }) }); let content = if model == "model/classifier" && custom_target_schema { - if requests_invalid_verdict { + if requests_schema_invalid_verdict { + // The selector would accept "fast"; only the schema's additionalProperties rejects it. + r#"{"decision":{"target":"fast"},"extra":1}"#.to_string() + } else if requests_invalid_verdict { r#"{"decision":{"target":"unknown"}}"#.to_string() } else { let group = requested_group.unwrap_or_else(|| "efficient".to_string()); @@ -409,7 +418,8 @@ async fn upstream_chat( .pointer("/response_format/json_schema/schema/properties/escalate") .is_some() { - r#"{"escalate":false,"reason":"making progress"}"#.to_string() + r#"{"escalate":false,"category":"none","new_evidence":false,"reason":"making progress"}"# + .to_string() } else if model == "model/classifier" && requests_schema_invalid_verdict { r#"{"crux":"bounded task","primary_rule":"SUP-1","capability_boundary":"supported","p_solve":0.1,"unexpected":true}"#.to_string() } else if model == "model/classifier" { @@ -598,7 +608,12 @@ async fn upstream_responses_silo( .get("escalate") .is_some() { - json!({"escalate": strong, "reason": "state probe"}) + json!({ + "escalate": strong, + "category": if strong { "capability_gap" } else { "none" }, + "new_evidence": strong, + "reason": "state probe", + }) } else { json!({ "crux": "state probe", "primary_rule": if strong { "LIM-1" } else { "SUP-1" }, @@ -2486,6 +2501,172 @@ selector = "/decision/target" Ok(()) } +#[tokio::test] +async fn custom_classifier_can_request_json_object_output() -> TestResult { + let upstream = MockUpstream::start().await?; + let state = load_test_config(&format!( + r#" +schema_version = 1 + +[llm_clients.upstream] +format = "openai_chat" +base_url = "{base_url}" + +[targets.classifier] +id = "model/classifier" +llm_client = "upstream" + +[targets.strong] +id = "model/strong" +llm_client = "upstream" + +[targets.weak] +id = "model/weak" +llm_client = "upstream" + +[routes.custom] +id = "switchyard/custom" +type = "llm_classifier" +mode = "custom" +response_format_type = "json_object" +models = {{ judge = ["classifier"], fast = ["weak"], reasoning = ["strong"], any = ["weak", "strong"] }} +default_target = "reasoning" +prompt = "CUSTOM JSON OBJECT" +response_schema = ''' +{{ + "type": "object", + "properties": {{ + "decision": {{ + "type": "object", + "properties": {{ + "target": {{"type": "string", "enum": ["fast", "reasoning"]}} + }}, + "required": ["target"], + "additionalProperties": false + }} + }}, + "required": ["decision"], + "additionalProperties": false +}} +''' + +[routes.custom.policy] +type = "target_selector" +selector = "/decision/target" +"#, + base_url = upstream.base_url + ))?; + let app = build_switchyard_router(state); + + // The provider only guarantees JSON here. The second verdict names a valid target, so the + // selector alone would pick `fast`; only local schema validation rejects it and falls + // back to `default_target`. + for (task, selected) in [ + ("route to fast", "model/weak"), + ("return a schema-invalid verdict", "model/strong"), + ] { + upstream.calls.lock().await.clear(); + let response = send( + &app, + "POST", + "/v1/chat/completions", + Some(json!({ + "model": "switchyard/custom", + "messages": [{"role": "user", "content": task}] + })), + ) + .await?; + assert_eq!(response.status, StatusCode::OK, "{task}"); + assert_eq!( + response.headers["x-model-router-selected-model"], selected, + "{task}" + ); + + let calls = upstream.calls.lock().await; + let judge_call = calls + .iter() + .find(|call| call["model"] == "model/classifier") + .ok_or("custom classifier target was not called")?; + assert_eq!( + judge_call["response_format"], + json!({"type": "json_object"}) + ); + let prompt = judge_call["messages"][0]["content"] + .as_str() + .ok_or("custom classifier prompt was not text")?; + assert!(prompt.starts_with("CUSTOM JSON OBJECT"), "{prompt}"); + assert!(prompt.contains("JSON Schema"), "{prompt}"); + assert!(prompt.contains("\"decision\""), "{prompt}"); + } + Ok(()) +} + +#[tokio::test] +async fn subagent_custom_classifier_can_request_json_object_output() -> TestResult { + let upstream = MockUpstream::start().await?; + let app = build_switchyard_router(load_test_config(&format!( + r#" +schema_version = 1 +[llm_clients.upstream] +format = "openai_chat" +base_url = "{base_url}" +[targets] +classifier = {{ id = "model/classifier", llm_client = "upstream" }} +strong = {{ id = "model/strong", llm_client = "upstream" }} +weak = {{ id = "model/weak", llm_client = "upstream" }} +[routes.agent] +id = "agent" +type = "passthrough" +target = "weak" +[routes.agent.subagents] +type = "llm_classifier" +mode = "custom" +response_format_type = "json_object" +models = {{ judge = ["classifier"], capable = ["strong"], efficient = ["weak"], any = ["strong", "weak"] }} +default_target = "efficient" +prompt = "classify the delegated task" +response_schema = '''{{"type":"object","properties":{{"decision":{{"type":"object","properties":{{"target":{{"type":"string","enum":["capable","efficient"]}}}},"required":["target"],"additionalProperties":false}}}},"required":["decision"],"additionalProperties":false}}''' +[routes.agent.subagents.policy] +type = "target_selector" +selector = "/decision/target" +"#, + base_url = upstream.base_url + ))?); + + let response = send_with_headers( + &app, + "POST", + "/v1/chat/completions", + Some(json!({ + "model": "agent", + "messages": [{"role": "user", "content": "route to capable"}] + })), + &[ + ("x-claude-code-session-id", "root-session"), + ("x-claude-code-agent-id", "child-agent"), + ], + ) + .await?; + + assert_eq!(response.status, StatusCode::OK); + assert_eq!( + response.headers["x-model-router-selected-model"], + "model/strong" + ); + let calls = upstream.calls.lock().await; + let judge_call = calls + .iter() + .find(|call| call["model"] == "model/classifier") + .ok_or("subagent classifier target was not called")?; + // The sub-agent path builds its own classifier config; this pins that the + // setting reaches the judge instead of silently staying on JSON Schema. + assert_eq!( + judge_call["response_format"], + json!({"type": "json_object"}) + ); + Ok(()) +} + #[tokio::test] async fn classifier_contract_overrides_reach_every_server_mode() -> TestResult { let upstream = MockUpstream::start().await?; diff --git a/crates/switchyard-translation/src/codecs/anthropic/buffered.rs b/crates/switchyard-translation/src/codecs/anthropic/buffered.rs index c6cdcfba5..231d97b84 100644 --- a/crates/switchyard-translation/src/codecs/anthropic/buffered.rs +++ b/crates/switchyard-translation/src/codecs/anthropic/buffered.rs @@ -291,8 +291,16 @@ impl FormatCodec for AnthropicMessagesCodec { if request.stream { body.insert("stream".to_string(), Value::Bool(true)); } + // Native thinking controls the mode independently of output effort. + // Other codecs store their own provider's reasoning object in `raw`. + if is_anthropic_request(request) + && let Some(thinking) = &request.reasoning.raw + { + body.insert("thinking".to_string(), thinking.clone()); + } if let Some(effort) = &request.reasoning.effort { - body.insert("thinking".to_string(), json!({"type": "adaptive"})); + body.entry("thinking".to_string()) + .or_insert_with(|| json!({"type": "adaptive"})); body.insert("output_config".to_string(), json!({"effort": effort})); } if let Some(response_format) = &request.output.response_format diff --git a/crates/switchyard-translation/tests/request_translation.rs b/crates/switchyard-translation/tests/request_translation.rs index 2af80f9d3..eceb54693 100644 --- a/crates/switchyard-translation/tests/request_translation.rs +++ b/crates/switchyard-translation/tests/request_translation.rs @@ -484,6 +484,38 @@ fn anthropic_thinking_to_responses_uses_normalized_effort() -> TestResult { Ok(()) } +#[test] +fn anthropic_reconstruction_preserves_thinking() -> TestResult { + let engine = TranslationEngine::default(); + let policy = normalized_policy(); + for (thinking, effort) in [ + (json!({"type": "disabled"}), Some("high")), + (json!({"type": "enabled", "budget_tokens": 2048}), None), + (json!({"type": "adaptive"}), Some("high")), + ] { + let mut body = json!({ + "model": "caller", "max_tokens": 4096, + "messages": [{"role": "user", "content": "hi"}], + "thinking": thinking + }); + if let Some(effort) = effort { + body["output_config"] = json!({"effort": effort}); + } + let mut request = engine + .decode_request(WireFormat::AnthropicMessages, &body, &policy)? + .request; + prepare_request_for_target(&mut request, &"target/model".into(), Some("target prompt")); + let output = engine + .encode_request(WireFormat::AnthropicMessages, &request, &policy)? + .body; + + body["model"] = json!("target/model"); + body["system"] = json!("target prompt"); + assert_eq!(output, body); + } + Ok(()) +} + #[test] fn anthropic_target_prompt_preserves_native_request_fields() -> TestResult { let engine = TranslationEngine::default(); diff --git a/docs/reference/toml_schema.md b/docs/reference/toml_schema.md index 8bf382a08..031f405a6 100644 --- a/docs/reference/toml_schema.md +++ b/docs/reference/toml_schema.md @@ -229,7 +229,7 @@ Runs one of three judge-backed modes: `capability`, `escalation`, or `custom`. | `mode` | No | `capability` | Classifier behavior. Set it explicitly for new configurations. | | `classifier_target` | Capability, escalation | — | Target the judge is called through. Not a routing destination. Custom mode uses `models.judge`. | | `max_output_tokens` | No | `4096` | Maximum completion tokens for the judge verdict. Must be at least `1`. | -| `response_format_type` | No | `json_schema` | Structured-output mode for capability and escalation judges. Use `json_object` when the provider does not support JSON Schema; Switchyard adds the schema to the prompt and validates the verdict locally. Custom mode always uses its configured JSON Schema. | +| `response_format_type` | No | `json_schema` | Structured-output mode for the judge in every mode. Use `json_object` when the provider does not support JSON Schema; Switchyard adds the schema to the prompt and validates the verdict locally. Custom mode appends `response_schema`; do not copy it into the prompt. | Capability mode classifies before serving. See [LLM Classifier Routing](../routing_algorithms/llm_classifier_routing.md). @@ -253,7 +253,7 @@ Escalation mode serves the weak target first and judges the completed turn. See | `strong_target` | Yes | — | Target used after the session latches. | | `weak_target` | Yes | — | Target served before the latch. | | `prompt` | No | packaged prompt | Replaces the trajectory-judge prompt. | -| `escalation.confirmations` | No | `2` | Consecutive escalate verdicts required to latch. Above `1` needs a stable session ID. | +| `escalation.confirmations` | No | `2` | Consecutive fresh-evidence verdicts for the same failure category required to latch. Above `1` needs a stable session ID. | | `escalation.recent_turn_window` | No | `28` | Trailing messages shown to the judge. | | `escalation.window_message_chars` | No | `500` | Per-message cap inside that window. | | `escalation.deescalation` | No | unset | Enables phase-aware de-escalation. Requires a stable session ID. | diff --git a/docs/routing_algorithms/escalation_router_routing.md b/docs/routing_algorithms/escalation_router_routing.md index 411446478..994359a08 100644 --- a/docs/routing_algorithms/escalation_router_routing.md +++ b/docs/routing_algorithms/escalation_router_routing.md @@ -61,8 +61,10 @@ The route-level `prompt` key replaces the packaged trajectory-judge prompt. It uses the escalation verdict schema rather than the capability verdict schema. Switchyard supplies that schema according to the route's `response_format_type`: through the structured-output request in the default `json_schema` mode, or in -the prompt in `json_object` mode. When de-escalation is enabled, Switchyard -appends the phase-specific verdict contract to packaged and custom prompts. +the prompt in `json_object` mode. The verdict includes an `escalate` decision, a +bounded failure `category`, whether the newest turn adds `new_evidence`, and a +short `reason`. When de-escalation is enabled, Switchyard appends the +phase-specific verdict contract to packaged and custom prompts. ## How the decision works @@ -72,8 +74,10 @@ For each turn on an unlatched session, Switchyard: 2. Appends that reply to the transcript and asks the judge to rule on the completed turn. The judge therefore rates work the weak model actually did, not a prediction about work it might do. -3. Increments a consecutive-escalate streak on an escalate verdict, and resets it - to zero on a decline. +3. Increments a confirmation streak only when an escalate verdict cites fresh + evidence for the same failure category as the preceding vote. A different + category starts a new streak at one; a decline or stale evidence resets it to + zero. 4. Serves the buffered weak reply when the streak has not yet reached `confirmations` — so a judged turn that does not escalate costs one weak call plus one judge call, and no strong call. @@ -119,7 +123,7 @@ configuration, so a bare `escalation = {}` is a valid, tuned route: | Key | Default | Meaning | |---|---|---| -| `confirmations` | `2` | Consecutive escalate verdicts required before the session latches to strong. Must be at least `1`. | +| `confirmations` | `2` | Consecutive fresh-evidence verdicts for the same failure category required before the session latches to strong. Must be at least `1`. | | `recent_turn_window` | `28` | Trailing messages shown to the judge on top of the anchors. Must be at least `1`. | | `window_message_chars` | `500` | Per-message truncation cap inside that trailing window. Must be at least `50`. | @@ -146,11 +150,13 @@ weak_cooldown_calls = 8 With this table present, Switchyard marks judge input as either `EFFICIENT_EVALUATION` or `STRONG_EVALUATION`. In the strong phase, -`escalate: true` keeps the strong tier. An `escalate: false` verdict can release -the next request only after `strong_min_calls` is reached and the configured -confirmation streak is complete. A timeout, error, or unparseable verdict -retains the strong tier. Omitting the table preserves the permanent latch and -does not add phase markers to judge input. +`escalate: true` keeps the strong tier; the router reads only `escalate` there, +so `category` and `new_evidence` are reported but do not affect release. An +`escalate: false` verdict can release the next request only after +`strong_min_calls` is reached and the configured confirmation streak is +complete. A timeout, error, or unparseable verdict retains the strong tier. +Omitting the table preserves the permanent latch and does not add phase markers +to judge input. The strong-phase verdict is judged against the trouble that caused the escalation: the packaged rules release only once the failure that triggered the @@ -179,8 +185,8 @@ route. The request still succeeds, but its temporary state cannot carry into the next request. Anchor and transcript caps remain fixed. Set the route-level -`max_output_tokens` key to change the judge's reply budget. Any decline still -resets the streak to zero. +`max_output_tokens` key to change the judge's reply budget. Any decline or +verdict without new evidence resets the streak to zero. ## Run the route @@ -226,6 +232,10 @@ per-session routing stats under the judge's model id, tagged with the `classifier` tier — so per-session token accounting includes judge overhead alongside the tiers the session was served by. +The server log records each parsed escalation verdict's `escalate` decision, +category, and `new_evidence` flag. The judge's free-form reason is neither +retained nor logged. + ## When not to use escalation routing - **One-shot requests.** No trajectory to judge. Use diff --git a/docs/routing_algorithms/llm_classifier_routing.md b/docs/routing_algorithms/llm_classifier_routing.md index de35015ff..537ca426d 100644 --- a/docs/routing_algorithms/llm_classifier_routing.md +++ b/docs/routing_algorithms/llm_classifier_routing.md @@ -92,7 +92,7 @@ Switchyard does not parse provider-specific reasoning fields such as to `strong_target` even when the judge request returned HTTP 200. With session affinity, that fallback can be reused without another judge call. -Capability and escalation routes use JSON Schema structured output by default. +Every classifier mode uses JSON Schema structured output by default. For a provider that supports JSON Object mode but not JSON Schema, set `response_format_type = "json_object"` on the route. Switchyard then adds the verdict schema to the judge prompt and validates the returned object locally. @@ -123,7 +123,7 @@ for the server merge behavior. | `classify_trigger` | `every_request` | When the judge runs. `every_request` judges every request, tool continuations included. `user_turn` judges each new user message and holds that target across the tool calls between. `new_session` judges once and reuses that target for the session. | | `message_hash_fallback` | `false` | When session metadata is absent, keys affinity from the first user-message text. Requires `classify_trigger = "new_session"` or `"user_turn"`. | | `prompt` | packaged capability prompt | Replaces the classifier's system prompt. The packaged verdict schema and routing policy remain active. | -| `response_format_type` | `json_schema` | Structured-output mode for capability and escalation judges. Use `json_object` for providers without JSON Schema support. | +| `response_format_type` | `json_schema` | Structured-output mode for the judge in every mode. Use `json_object` for providers without JSON Schema support. | | `max_output_tokens` | `4096` | Maximum completion tokens available to the classifier verdict. Must be at least `1`. | ### Override the classifier prompt @@ -200,9 +200,11 @@ type = "target_selector" selector = "/decision/target" ``` -The names in `models` reference existing target tables. Switchyard passes the -schema to the provider in a strict structured-output wrapper and validates the -returned JSON again. `jsonptr` resolves the selector against that verdict. A +The names in `models` reference existing target tables. By default Switchyard +passes the schema to the provider in a strict structured-output wrapper; with +`response_format_type = "json_object"` it appends the schema to the prompt +instead (see below). Either way it validates the returned JSON against the +schema. `jsonptr` resolves the selector against that verdict. A missing, non-string, or unconfigured label falls back to `default_target`, and `judge` is never routable. @@ -218,6 +220,13 @@ order and is not a completion destination. carry the tier meaning the stage and composite routers give it; otherwise any name works. +If the judge's provider supports JSON Object mode but not JSON Schema, add +`response_format_type = "json_object"` to the route. `response_schema` is still +required. Switchyard appends it to your prompt, asks for a JSON object, and +validates the reply against it. The configured schema is the source of truth, so +do not paste a copy into the prompt. A reply that fails the schema falls back to +`default_target`. + This separation applies to every classifier mode. Prompts containing the legacy `{{RESPONSE_SCHEMA}}` placeholder are rejected during configuration validation. diff --git a/scripts/linux/common.sh b/scripts/linux/common.sh new file mode 100644 index 000000000..25b16dd83 --- /dev/null +++ b/scripts/linux/common.sh @@ -0,0 +1,47 @@ +# SPDX-FileCopyrightText: Copyright (c) 2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved. +# SPDX-License-Identifier: Apache-2.0 +# +# Paths, markers, and helpers shared by install.sh and uninstall.sh. +# The markers must match between the two, which is why they live here. + +SY_HOME="${SY_HOME:-$HOME/.switchyard}" +SY_PORT="${SY_PORT:-4123}" +SERVICE_NAME="switchyard.service" +SYSTEMD_USER_DIR="${XDG_CONFIG_HOME:-$HOME/.config}/systemd/user" +CODEX_DIR="${CODEX_HOME:-$HOME/.codex}" +# Since Codex 0.134.0, `--profile sy` reads sy.config.toml and fails if +# config.toml still has [profiles.sy]. +CODEX_PROFILE_CONFIG="$CODEX_DIR/sy.config.toml" +ALIAS_START="# >>> switchyard codex alias >>>" +ALIAS_END="# <<< switchyard codex alias <<<" + +say() { printf '%s\n' "$*"; } +step() { printf '\n==> %s\n' "$*"; } + +# Deletes the marked block, inclusive, leaving the rest of the file alone. +strip_block() { + local path="$1" start="$2" end="$3" label="${4:-the switchyard block}" + if [[ ! -f "$path" ]] || ! grep -qF "$start" "$path"; then + return 1 + fi + if ! grep -qF "$end" "$path"; then + say " $path has no end marker for $label; not editing it" >&2 + return 1 + fi + if (( DRY_RUN )); then + say " would remove $label from $path" + return 0 + fi + local temp + temp="$(mktemp)" + awk -v start="$start" -v end="$end" ' + index($0, start) { skipping = 1 } + !skipping { print } + index($0, end) { skipping = 0 } + END { if (skipping) exit 1 } + ' "$path" > "$temp" + cat "$temp" > "$path" + rm -f "$temp" + say " removed $label from $path" + return 0 +} diff --git a/scripts/linux/install.sh b/scripts/linux/install.sh new file mode 100755 index 000000000..7bf19e422 --- /dev/null +++ b/scripts/linux/install.sh @@ -0,0 +1,216 @@ +#!/usr/bin/env bash +# SPDX-FileCopyrightText: Copyright (c) 2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved. +# SPDX-License-Identifier: Apache-2.0 +# +# Installs the Switchyard background server as a systemd --user service and sets +# up a `sy` Codex profile. +# +# Keeps existing composite.toml, replaces the service unit, and backs up +# sy.config.toml before replacing it. Use systemctl --user edit switchyard +# for service changes that survive reinstalling. +# Run with --dry-run to print what would happen. + +set -euo pipefail + +SCRIPT_DIR="$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd)" +REPO_ROOT="$(cd "$SCRIPT_DIR/../.." && pwd)" + +DRY_RUN=0 +case "$#:${1:-}" in + 0:) ;; + 1:--dry-run) DRY_RUN=1 ;; + *) printf 'Usage: %s [--dry-run]\n' "$0" >&2; exit 2 ;; +esac + +# shellcheck source=scripts/linux/common.sh +source "$SCRIPT_DIR/common.sh" + +if [[ "$SY_HOME" =~ [[:space:][:cntrl:]] || "$SY_HOME" == *\\ ]]; then + say "SY_HOME must not contain whitespace, control characters, or a trailing backslash." >&2 + exit 1 +fi +if [[ ! "$SY_PORT" =~ ^[0-9]+$ ]]; then + say "SY_PORT must contain only digits." >&2 + exit 1 +fi + +# Runs a command, or prints it when dry running. +run() { + if (( DRY_RUN )); then + say " would run: $*" + else + "$@" + fi +} + +# Writes stdin to a file, leaving an existing file untouched. +write_once() { + local path="$1" + if [[ -f "$path" ]]; then + say " keeping existing $path" + cat >/dev/null + return + fi + if (( DRY_RUN )); then + say " would create $path" + cat >/dev/null + else + mkdir -p "$(dirname "$path")" + cat > "$path" + say " created $path" + fi +} + +# Writes stdin to a file, keeping any existing version as a timestamped backup. +write_with_backup() { + local path="$1" + if (( DRY_RUN )); then + [[ -f "$path" ]] && say " would back up $path" + say " would write $path" + cat >/dev/null + return + fi + mkdir -p "$(dirname "$path")" + local incoming + incoming="$(mktemp)" + cat > "$incoming" + if [[ -f "$path" ]] && cmp -s "$incoming" "$path"; then + say " $path is already up to date" + rm -f "$incoming" + return + fi + if [[ -f "$path" ]]; then + local backup + backup="$path.switchyard-backup.$(date +%Y%m%d%H%M%S)" + cp "$path" "$backup" + say " backed up $path to $backup" + fi + cat "$incoming" > "$path" + rm -f "$incoming" + say " wrote $path" +} + +# Writes stdin to a file, replacing it. Used only for files this script owns. +write_always() { + local path="$1" + if (( DRY_RUN )); then + say " would write $path" + cat >/dev/null + else + mkdir -p "$(dirname "$path")" + cat > "$path" + say " wrote $path" + fi +} + +if [[ "$(uname -s)" != "Linux" ]]; then + say "This installer is for Linux only." >&2 + exit 1 +fi + +step "Building release binary" +run cargo build --release --manifest-path "$REPO_ROOT/Cargo.toml" -p switchyard-server + +step "Installing binary into $SY_HOME/bin" +run mkdir -p "$SY_HOME/bin" +run install -m 755 "$REPO_ROOT/target/release/switchyard-server" "$SY_HOME/bin/switchyard-server" + +step "Writing server config" +# The composite router from examples/run_codex.sh: Terra classifies each user +# turn and sets the tier, Stage drives the tool loop underneath it. +write_once "$SY_HOME/composite.toml" <<'EOF' +schema_version = 1 + +[llm_clients.chatgpt_backend] +format = "openai_responses" +base_url = "https://chatgpt.com/backend-api/codex" +forward_auth = true + +[targets.capable] +id = "gpt-5.6-sol" +llm_client = "chatgpt_backend" + +[targets.efficient] +id = "gpt-5.6-luna" +llm_client = "chatgpt_backend" + +# chatgpt.com/backend-api/codex is Codex CLI's own private endpoint, not the +# public OpenAI Responses API. It 400s unless store=false and stream=true are +# set explicitly, and it rejects max_output_tokens outright, so the +# classifier's own token cap has to be dropped before the request goes out. +[targets.terra] +id = "gpt-5.6-terra" +llm_client = "chatgpt_backend" +extra_body = { store = false, stream = true } +omit_body_fields = ["max_output_tokens"] + +[routes.switchyard] +id = "switchyard" +type = "composite" + +[routes.switchyard.classifier] +target = "terra" +base_threshold = 0.5 +classify_trigger = "user_turn" + +[routes.switchyard.stage] +capable_target = "capable" +efficient_target = "efficient" +confidence_threshold = 0.5 +EOF + +step "Validating the server config" +if (( DRY_RUN )); then + say " would run: $SY_HOME/bin/switchyard-server --config $SY_HOME/composite.toml --dry-run" +else + "$SY_HOME/bin/switchyard-server" --config "$SY_HOME/composite.toml" --dry-run +fi + +step "Writing the systemd user service" +write_always "$SYSTEMD_USER_DIR/$SERVICE_NAME" <&2 + exit 1 +fi + +step "Adding the sy Codex profile" +# This profile only changes which router answers. Approval and sandbox +# settings are deliberately left out, so the profile cannot loosen how Codex +# asks before it acts. Set those yourself if you want them. +write_with_backup "$CODEX_PROFILE_CONFIG" <&2; exit 2 ;; +esac + +# shellcheck source=scripts/linux/common.sh +source "$SCRIPT_DIR/common.sh" + +# Deletes a file this script owns. +remove_file() { + local path="$1" + if [[ ! -e "$path" ]]; then + say " nothing to remove at $path" + elif (( DRY_RUN )); then + say " would delete $path" + else + rm -f "$path" + say " deleted $path" + fi +} + +step "Removing the sy Codex profile" +remove_file "$CODEX_PROFILE_CONFIG" + +step "Removing the codex alias" +for rc in "$HOME/.zshrc" "$HOME/.bashrc"; do + if [[ -f "$rc" ]] && grep -qF "$ALIAS_START" "$rc"; then + strip_block "$rc" "$ALIAS_START" "$ALIAS_END" "the codex alias" + else + say " no alias in $rc" + fi +done + +step "Stopping the systemd user service" +if (( DRY_RUN )); then + say " would run: systemctl --user disable --now $SERVICE_NAME" + say " would delete $SYSTEMD_USER_DIR/$SERVICE_NAME" +else + systemctl --user disable --now "$SERVICE_NAME" 2>/dev/null || true + rm -f "$SYSTEMD_USER_DIR/$SERVICE_NAME" + systemctl --user daemon-reload + say " stopped and removed $SERVICE_NAME" +fi + +step "Done" +say "Left in place, delete them if you want:" +say " $SY_HOME (binary, config, routing log)" +say " $CODEX_PROFILE_CONFIG.switchyard-backup.* (backups taken at install time)" diff --git a/switchyard_rust/libsy.py b/switchyard_rust/libsy.py index 7bdf2cb6c..491bc4e7a 100644 --- a/switchyard_rust/libsy.py +++ b/switchyard_rust/libsy.py @@ -77,6 +77,7 @@ def __init__( message_hash_fallback: bool = False, recent_turn_window: int | None = None, max_output_tokens: int = 4096, + response_format_type: Literal["json_schema", "json_object"] = "json_schema", ) -> None: ... @final diff --git a/tests/test_libsy_minimal_bindings.py b/tests/test_libsy_minimal_bindings.py index e28ffc126..d9badd2ae 100644 --- a/tests/test_libsy_minimal_bindings.py +++ b/tests/test_libsy_minimal_bindings.py @@ -358,6 +358,64 @@ async def call(self, request: dict[str, Any]) -> dict[str, Any]: assert response["model"] == "weak" +async def test_custom_classifier_config_accepts_json_object_output() -> None: + """Verify that Python can select JSON Object mode for a custom classifier judge.""" + + class JudgeClient(EchoClient): + async def call(self, request: dict[str, Any]) -> dict[str, Any]: + self.calls.append(request) + return { + "model": self.model, + "outputs": [ + { + "role": "assistant", + "content": [{"type": "text", "text": '{"target":"efficient"}'}], + "stop_reason": "end_turn", + } + ], + } + + schema = { + "type": "object", + "additionalProperties": False, + "required": ["target"], + "properties": {"target": {"type": "string", "enum": ["capable", "efficient"]}}, + } + judge = JudgeClient("judge") + algorithm = algorithms.llm_classifier( + LlmClassifierConfig.custom( + default_target="capable", + config=CustomClassifierConfig( + "Choose a target.", + schema, + "/target", + response_format_type="json_object", + ), + ) + ) + + _, response = await run_algorithm( + algorithm, + { + "judge": judge, + "model-a": EchoClient("model-a"), + "model-b": EchoClient("model-b"), + }, + models={ + "judge": ["judge"], + "capable": ["model-a"], + "efficient": ["model-b"], + "any": ["model-a", "model-b"], + }, + ) + + assert judge.calls[0]["output"]["response_format"] == {"type": "json_object"} + prompt = judge.calls[0]["instructions"][0]["content"][0]["text"] + assert prompt.startswith("Choose a target.") + assert '"efficient"' in prompt + assert response["model"] == "model-b" + + def test_classifier_config_rejects_unknown_response_format() -> None: invalid_response_format: Any = "yaml" @@ -368,6 +426,22 @@ def test_classifier_config_rejects_unknown_response_format() -> None: TaskClassifierConfig(0.5, response_format_type=invalid_response_format) +def test_custom_classifier_config_rejects_unknown_response_format() -> None: + invalid_response_format: Any = "yaml" + schema = {"type": "object"} + + with pytest.raises( + ValueError, + match="response_format_type must be 'json_schema' or 'json_object'", + ): + CustomClassifierConfig( + "Choose a target.", + schema, + "/target", + response_format_type=invalid_response_format, + ) + + async def test_random_weights_and_seed_are_reproducible() -> None: def algorithm(): return algorithms.random( diff --git a/tests/test_linux_install.py b/tests/test_linux_install.py new file mode 100644 index 000000000..e41dde9f7 --- /dev/null +++ b/tests/test_linux_install.py @@ -0,0 +1,195 @@ +# SPDX-FileCopyrightText: Copyright (c) 2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved. +# SPDX-License-Identifier: Apache-2.0 + +import os +import shutil +import subprocess +from pathlib import Path + +import pytest + +REPO = Path(__file__).resolve().parents[1] +START = "# >>> switchyard codex alias >>>" +END = "# <<< switchyard codex alias <<<" +RC = f"before\n{START}\nalias codex='codex -p sy'\n{END}\nafter\n" + + +@pytest.fixture +def setup(tmp_path): + scripts = tmp_path / "repo" / "scripts" / "linux" + shutil.copytree(REPO / "scripts" / "linux", scripts) + home = tmp_path / "home" + home.mkdir() + bin_dir = tmp_path / "bin" + bin_dir.mkdir() + env = { + **os.environ, + "HOME": str(home), + "SY_HOME": str(home / ".switchyard"), + "SY_PORT": "4123", + "XDG_CONFIG_HOME": str(home / ".config"), + "CODEX_HOME": str(home / ".codex"), + "TMPDIR": str(tmp_path), + "PATH": f"{bin_dir}:/usr/bin:/bin", + "SYSTEMCTL_LOG": str(tmp_path / "systemctl.log"), + "FAIL_SYSTEMCTL": "", + } + stubs = { + "cargo": "exit 0\n", + "uname": "echo Linux\n", + "sleep": "exit 0\n", + "install": 'printf "#!/bin/sh\\nexit 0\\n" > "$4"\nchmod +x "$4"\n', + "systemctl": 'echo "$*" >> "$SYSTEMCTL_LOG"\n' + '[[ "$*" != *"$FAIL_SYSTEMCTL"* || -z "$FAIL_SYSTEMCTL" ]]\n', + } + for name, body in stubs.items(): + stub = bin_dir / name + stub.write_text("#!/bin/bash\nset -eu\n" + body) + stub.chmod(0o755) + return scripts, home, bin_dir, env + + +def run(setup, script, *args): + scripts, _, _, env = setup + return subprocess.run( + ["bash", str(scripts / script), *args], + env=env, + capture_output=True, + text=True, + timeout=10, + ) + + +@pytest.mark.parametrize("script", ["install.sh", "uninstall.sh"]) +@pytest.mark.parametrize("args", [["--dryrun"], ["--help"], ["-n"], ["--dry-run", "extra"], [""]]) +def test_bad_arguments_leave_files_unchanged(setup, script, args): + _, home, _, env = setup + rc = home / ".bashrc" + rc.write_text(RC) + result = run(setup, script, *args) + assert result.returncode == 2 + assert "Usage:" in result.stderr + assert rc.read_text() == RC + assert not Path(env["SYSTEMCTL_LOG"]).exists() + + +@pytest.mark.parametrize("failure", ["mktemp", "awk", "missing-end", "second-missing-end"]) +def test_uninstall_failure_preserves_shell_files(setup, failure): + _, home, bin_dir, _ = setup + contents = RC + if failure == "missing-end": + contents = RC.replace(END + "\n", "") + elif failure == "second-missing-end": + contents += f"{START}\nalias codex='codex -p sy'\ntail\n" + else: + (bin_dir / failure).write_text("#!/bin/bash\nexit 1\n") + (bin_dir / failure).chmod(0o755) + for name in [".zshrc", ".bashrc"]: + (home / name).write_text(contents) + result = run(setup, "uninstall.sh") + assert result.returncode != 0 + for name in [".zshrc", ".bashrc"]: + assert (home / name).read_text() == contents + + +def test_uninstall_cleans_profile_and_aliases_before_systemctl_failure(setup): + _, home, _, env = setup + env["FAIL_SYSTEMCTL"] = "--user" + profile = Path(env["CODEX_HOME"]) / "sy.config.toml" + profile.parent.mkdir() + profile.write_text("old profile\n") + for name in [".zshrc", ".bashrc"]: + (home / name).write_text(RC) + result = run(setup, "uninstall.sh") + assert result.returncode != 0 + assert not profile.exists() + for name in [".zshrc", ".bashrc"]: + assert (home / name).read_text() == "before\nafter\n" + + +@pytest.mark.parametrize("script", ["install.sh", "uninstall.sh"]) +def test_dry_run_leaves_files_unchanged(setup, script): + _, home, _, env = setup + (home / ".bashrc").write_text(RC) + result = run(setup, script, "--dry-run") + assert result.returncode == 0, result.stderr + assert (home / ".bashrc").read_text() == RC + assert sorted(path.name for path in home.iterdir()) == [".bashrc"] + assert not Path(env["SYSTEMCTL_LOG"]).exists() + + +@pytest.mark.parametrize( + "key,value", + [ + ("SY_HOME", "/tmp/with space"), + ("SY_HOME", "/tmp/line\nbreak"), + ("SY_HOME", "/tmp/control\x01"), + ("SY_HOME", "/tmp/backslash\\"), + ("SY_PORT", "bad"), + ("SY_PORT", "4123\n"), + ], +) +def test_invalid_unit_values_fail_before_install(setup, key, value): + _, home, _, env = setup + env[key] = value + result = run(setup, "install.sh") + assert result.returncode != 0 + assert key in result.stderr + assert not list(home.iterdir()) + + +def test_reinstall_restarts_service_and_keeps_user_config(setup): + _, home, _, env = setup + (home / ".bashrc").write_text("user settings\n") + assert run(setup, "install.sh").returncode == 0 + config = Path(env["SY_HOME"]) / "composite.toml" + config.write_text("user config\n") + env["SY_PORT"] = "5000" + result = run(setup, "install.sh") + assert result.returncode == 0, result.stderr + assert config.read_text() == "user config\n" + assert (home / ".bashrc").read_text() == "user settings\n" + assert not (home / ".zshrc").exists() + profile = Path(env["CODEX_HOME"]) / "sy.config.toml" + assert "127.0.0.1:5000/v1" in profile.read_text() + backups = list(profile.parent.glob("sy.config.toml.switchyard-backup.*")) + assert len(backups) == 1 + assert "127.0.0.1:4123/v1" in backups[0].read_text() + calls = Path(env["SYSTEMCTL_LOG"]).read_text().splitlines() + assert ( + calls + == [ + "--user daemon-reload", + "--user enable switchyard.service", + "--user restart switchyard.service", + "--user is-active --quiet switchyard.service", + ] + * 2 + ) + + +def test_failed_start_leaves_existing_profile_unchanged(setup): + _, _, _, env = setup + env["FAIL_SYSTEMCTL"] = "is-active" + profile = Path(env["CODEX_HOME"]) / "sy.config.toml" + profile.parent.mkdir() + profile.write_text("user profile\n") + result = run(setup, "install.sh") + assert result.returncode != 0 + assert "journalctl" in result.stderr + assert profile.read_text() == "user profile\n" + + +def test_default_make_only_prints_help(setup): + _, home, _, env = setup + result = subprocess.run( + ["make", "-f", str(REPO / "Makefile")], + cwd=REPO, + env=env, + capture_output=True, + text=True, + timeout=10, + ) + assert result.returncode == 0, result.stderr + assert "install-linux-dry-run" in result.stdout + assert not list(home.iterdir())