Skip to main content

switchyard_libsy/algorithms/
advisor_gate.rs

1// SPDX-FileCopyrightText: Copyright (c) 2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved.
2// SPDX-License-Identifier: Apache-2.0
3
4//! Executor gated by a once-per-session advisor review.
5//!
6//! The executor answers every client-visible turn. Turns with tool calls pass
7//! through unreviewed; the first *terminal* turn — no tool calls (or a text
8//! match under the `pattern` trigger) — is buffered and shown to a stronger
9//! advisor model together with the full transcript. `APPROVE` releases the
10//! buffered turn unchanged; `REDO` appends the discarded turn's text and the
11//! advisor's plan as feedback, then re-invokes the executor so it keeps
12//! working. Each budget scope (one benchmark evaluation, one session, or the
13//! whole instance — see [`budget::budget_scope`]) is reviewed at most
14//! `max_reviews` times; afterwards every call is a pure passthrough.
15//!
16//! This design is a near-superset of solo executor behavior: identical until
17//! the executor first claims to be done, plus one quality gate that catches
18//! premature convergence. Front-loading advice was measured to suppress the
19//! executor's own test-and-iterate loop, so no advice is injected up front.
20//!
21//! Failure posture: executor errors always propagate (including
22//! `ContextWindowExceeded`, which hosts map to a client-visible 400 so agent
23//! harnesses can compact). Advisor errors honor `fail_open` — the buffered
24//! turn passes through as an implicit APPROVE — refund the consumed review,
25//! and count toward a per-scope failure cap that stops consulting a down
26//! advisor entirely.
27//!
28//! Structure: [`AdvisorGate`] is a thin orchestrator — the [`signals`]
29//! processor folds each event's facts into per-turn state, the [`trigger`]
30//! classifier reads them after the executor call, and the [`budget`] ledger
31//! holds the only mutable state.
32
33use std::sync::Arc;
34use std::time::Instant;
35
36use switchyard_protocol::{
37    Category, ContentBlock, InstructionBlock, LlmRequest, Message, ModelId, OutputParams, Request,
38    Role, SamplingParams,
39};
40
41use crate::core::algorithm::{Algorithm, Driver, RoutingOutcome};
42use crate::core::processor::{Event, Processor};
43use crate::{LibsyError, Result};
44
45mod budget;
46mod signals;
47mod telemetry;
48#[cfg(test)]
49mod tests;
50mod transcript;
51mod trigger;
52mod turn;
53
54use super::util::buffered_response::{BufferedResponse, buffer_response};
55use super::util::robustness::safe_error_summary;
56use budget::{ReviewBudget, ScopeKey, budget_scope, stall_key};
57use signals::{GateSignalProcessor, GateSignals};
58use telemetry::{
59    ReviewAudit, emit_discarded_audit, emit_review_audit, record_consult_failure, record_discarded,
60    record_review,
61};
62use transcript::{VERDICT_PATTERN, Verdict, advisor_reply_text, parse_verdict, review_transcript};
63use trigger::TriggerClassifier;
64#[cfg(test)]
65use turn::has_tool_use;
66use turn::{reasoning_text, visible_text};
67
68/// APPROVE/REDO reviewer contract sent as the advisor's system prompt.
69pub const REVIEWER_SYSTEM_PROMPT: &str =
70    include_str!("../prompts/advisor-gate/reviewer-system-prompt.md");
71
72/// Prepended to the advisor's REDO plan when it is fed back as a user turn,
73/// instructing the executor to continue rather than stop.
74pub const REDO_FEEDBACK_PREFIX: &str = concat!(
75    include_str!("../prompts/advisor-gate/redo-feedback-prefix.md"),
76    "\n"
77);
78
79/// Labels the executor's internal reasoning when a turn has no visible text,
80/// so the advisor still has evidence to review (reasoning models on vLLM/NIM
81/// can emit turns whose only output is reasoning).
82const REASONING_TAIL_LABEL: &str =
83    "(the executor produced no visible text this turn; its internal reasoning follows)\n";
84/// REDO echo when the discarded turn had neither text nor reasoning; strict
85/// endpoints (Anthropic) reject empty text blocks, so never echo "".
86const EMPTY_ECHO_PLACEHOLDER: &str = "(the executor produced no output this turn)";
87/// Benchmark harnesses stamp every request of one evaluation — sub-agents
88/// included — with this header, so it is the review budget's first-choice
89/// scope: "reviews for *this* task" survives gateways shared by many tasks.
90const BENCH_SESSION_HEADER: &str = "proxy_x_session_id";
91
92/// How the gate decides a buffered executor turn is terminal.
93#[derive(Clone, Debug, PartialEq)]
94pub enum GateTrigger {
95    /// First turn without tool calls (subject to `gate_min_tool_results`).
96    NoToolCall,
97    /// First turn whose visible text matches this regex (searched, not anchored) —
98    /// for text-protocol harnesses where every turn lacks tool calls and
99    /// completion is declared with a textual marker instead.
100    Pattern(String),
101}
102
103/// Gate knobs; defaults mirror the benchmarked Python advisor configuration.
104#[derive(Clone, Debug)]
105pub struct AdvisorGateConfig {
106    /// System prompt for the advisor's review call; states the APPROVE/REDO contract.
107    pub reviewer_system_prompt: String,
108    /// Prepended to the advisor's REDO plan when fed back to the executor.
109    pub redo_feedback_prefix: String,
110    /// What fires the review.
111    pub gate_trigger: GateTrigger,
112    /// Reviews allowed per budget scope. 1 keeps the original once-per-task
113    /// gate; higher values re-review later terminal turns, making the gate a
114    /// sequential best-of-(N+1) with the advisor as judge.
115    pub max_reviews: u32,
116    /// When > 0, additionally review (once per conversation, consuming budget)
117    /// the first request already carrying at least this many assistant turns —
118    /// a mid-task checkpoint for executors that grind without declaring
119    /// completion. 0 disables.
120    pub gate_stall_turns: u32,
121    /// For the `no_tool_call` trigger: only review once the conversation
122    /// carries at least this many tool results, skipping early commentary
123    /// turns on chatty harnesses. 0 reviews from the first terminal turn.
124    pub gate_min_tool_results: u32,
125    /// Cap on the advisor's output per consult.
126    pub advisor_max_tokens: u64,
127    /// Sampling temperature for the consult; `None` omits the field on the wire.
128    pub advisor_temperature: Option<f64>,
129    /// Cap on the serialized transcript handed to the advisor; the middle of
130    /// an over-cap conversation is dropped (task head + recent tail survive).
131    pub transcript_max_chars: usize,
132    /// When true (default), an advisor failure degrades to APPROVE; when
133    /// false, it propagates as the turn's error.
134    pub fail_open: bool,
135}
136
137impl Default for AdvisorGateConfig {
138    fn default() -> Self {
139        Self {
140            reviewer_system_prompt: REVIEWER_SYSTEM_PROMPT.to_string(),
141            redo_feedback_prefix: REDO_FEEDBACK_PREFIX.to_string(),
142            gate_trigger: GateTrigger::NoToolCall,
143            max_reviews: 1,
144            gate_stall_turns: 0,
145            gate_min_tool_results: 0,
146            advisor_max_tokens: 2048,
147            advisor_temperature: None,
148            transcript_max_chars: 200_000,
149            fail_open: true,
150        }
151    }
152}
153
154/// Advisor review gate: executor turns pass through until the first terminal
155/// turn, which a stronger advisor reviews once per scope budget (APPROVE
156/// releases it, REDO feeds the plan back and re-invokes the executor).
157pub struct AdvisorGate {
158    config: AdvisorGateConfig,
159    /// Folds request- and response-side facts into the per-turn [`GateSignals`].
160    signals: GateSignalProcessor,
161    /// Decides, from the signals, whether the buffered turn warrants review.
162    trigger: TriggerClassifier,
163    /// Reserve/refund review ledger and stall latch — the gate's only mutable state.
164    budget: ReviewBudget,
165    verdict_re: regex::Regex,
166}
167
168impl AdvisorGate {
169    /// Validates the config. Models are supplied when each request runs.
170    pub fn new(config: AdvisorGateConfig) -> Result<Self> {
171        if config.max_reviews < 1 {
172            return Err(algorithm_error("max_reviews must be at least 1"));
173        }
174        if config.advisor_max_tokens < 1 {
175            return Err(algorithm_error("advisor_max_tokens must be at least 1"));
176        }
177        if config.transcript_max_chars < 256 {
178            return Err(algorithm_error("transcript_max_chars must be at least 256"));
179        }
180        let trigger = TriggerClassifier::new(&config)?;
181        let verdict_re = regex::Regex::new(VERDICT_PATTERN).map_err(|error| {
182            algorithm_error(format!("verdict pattern failed to compile: {error}"))
183        })?;
184        let budget = ReviewBudget::new(config.max_reviews);
185        Ok(Self {
186            config,
187            signals: GateSignalProcessor,
188            trigger,
189            budget,
190            verdict_re,
191        })
192    }
193
194    // ── Gate flow ───────────────────────────────────────────────────────────
195
196    async fn route_inner(
197        &self,
198        driver: &Driver,
199        request: Request,
200        scope: &ScopeKey,
201    ) -> Result<RoutingOutcome> {
202        let executor_models = driver.models_for(&Category::Efficient).to_vec();
203        let executor = executor_models
204            .first()
205            .ok_or_else(|| LibsyError::AlgorithmError {
206                message: "no models available for category Efficient".to_string(),
207            })?;
208
209        // Spent budget (or failure cap): pure passthrough — live stream,
210        // verbatim preserved-body replay, zero buffering. Executor errors
211        // (including ContextWindowExceeded) propagate for the host's
212        // client-visible mapping.
213        if self.budget.check_exhausted(scope) {
214            return Ok(RoutingOutcome::route_to(
215                executor.clone(),
216                executor_models[1..].to_vec(),
217                request,
218            ));
219        }
220
221        // Request-side signals fold in before the executor runs.
222        let mut request = request;
223        let mut signals = GateSignals::default();
224        self.signals
225            .process(
226                &mut signals,
227                Event::Request {
228                    request: &mut request,
229                    driver,
230                },
231            )
232            .await?;
233
234        // Gated phase: generate the turn once, fully buffered, so the gate
235        // can inspect it before the client sees anything.
236        let response = driver
237            .call_model(request.clone(), executor_models.clone())
238            .await?;
239        let served_executor = response
240            .served_model()
241            .cloned()
242            .unwrap_or_else(|| executor.clone());
243        let turn = buffer_response(served_executor.as_str(), response).await?;
244
245        // Response-side signals fold in after it: the terminal turn never
246        // appears on a later request, so the trigger runs on this event.
247        self.signals
248            .process(&mut signals, Event::ModelResponse(&turn.agg))
249            .await?;
250
251        let decision = self.trigger.classify(&signals);
252        // The stall checkpoint fires once per conversation regardless of the
253        // turn's shape. Only a stall with no simultaneous trigger latches
254        // (atomically — one winner per conversation). The latch is provisional
255        // until a review completes: a refunded consult or a spent budget
256        // re-arms it so a later eligible turn is reviewed instead of silently
257        // passing through.
258        let stall_key = stall_key(&request);
259        let stall = decision.fired.is_none()
260            && decision.stalled
261            && self.budget.try_mark_stall_fired(stall_key);
262        if decision.fired.is_none() && !stall {
263            return Ok(RoutingOutcome::answered(
264                served_executor.clone(),
265                request,
266                turn.into_response(),
267            ));
268        }
269        if !self.budget.try_reserve(scope) {
270            if stall {
271                self.budget.clear_stall_fired(stall_key);
272            }
273            return Ok(RoutingOutcome::answered(
274                served_executor.clone(),
275                request,
276                turn.into_response(),
277            ));
278        }
279
280        let trigger_label = decision.fired.unwrap_or("stall");
281        let review_tail = visible_text(&turn.agg).or_else(|| {
282            reasoning_text(&turn.agg).map(|reasoning| format!("{REASONING_TAIL_LABEL}{reasoning}"))
283        });
284        match self
285            .consult(driver, &request, review_tail.as_deref(), trigger_label)
286            .await
287        {
288            Ok(ConsultOutcome::Approve) => {
289                driver.set_evidence(serde_json::json!({
290                    "source": "advisor",
291                    "verdict": "approve",
292                    "trigger": trigger_label,
293                }));
294                Ok(RoutingOutcome::answered(
295                    served_executor.clone(),
296                    request,
297                    turn.into_response(),
298                ))
299            }
300            Ok(ConsultOutcome::Redo { plan }) => {
301                driver.set_evidence(serde_json::json!({
302                    "source": "advisor",
303                    "verdict": "redo",
304                    "trigger": trigger_label,
305                }));
306                Ok(self.redo(
307                    executor,
308                    &executor_models[1..],
309                    &served_executor,
310                    request,
311                    turn,
312                    &plan,
313                ))
314            }
315            Ok(ConsultOutcome::Failed { reason }) => {
316                self.budget.refund_failure(scope);
317                if stall {
318                    self.budget.clear_stall_fired(stall_key);
319                }
320                driver.set_evidence(serde_json::json!({
321                    "source": "advisor",
322                    "verdict": "fail_open",
323                    "trigger": trigger_label,
324                    "reason_code": reason,
325                }));
326                Ok(RoutingOutcome::answered(
327                    served_executor,
328                    request,
329                    turn.into_response(),
330                ))
331            }
332            Err(error) => {
333                self.budget.refund_failure(scope);
334                if stall {
335                    self.budget.clear_stall_fired(stall_key);
336                }
337                Err(error)
338            }
339        }
340    }
341
342    /// REDO: the client never sees the gated turn. Its text (or reasoning) is
343    /// echoed as an assistant message, the advisor's plan follows as user
344    /// feedback, and the executor continues as a pure passthrough call.
345    fn redo(
346        &self,
347        executor: &ModelId,
348        executor_fallbacks: &[ModelId],
349        served_executor: &ModelId,
350        request: Request,
351        turn: BufferedResponse,
352        plan: &str,
353    ) -> RoutingOutcome {
354        record_discarded(&turn.agg.usage);
355        emit_discarded_audit(served_executor.as_str(), &turn.agg.usage);
356        let echo = visible_text(&turn.agg)
357            .or_else(|| reasoning_text(&turn.agg))
358            .unwrap_or_else(|| EMPTY_ECHO_PLACEHOLDER.to_string());
359        let mut redo = request;
360        redo.llm_request
361            .messages
362            .push(Message::text(Role::Assistant, echo));
363        redo.llm_request.messages.push(Message::text(
364            Role::User,
365            format!("{}{}", self.config.redo_feedback_prefix, plan),
366        ));
367        // Mandatory after any message mutation: codecs otherwise replay the
368        // preserved pre-surgery body verbatim and the feedback never reaches
369        // the executor.
370        crate::algorithms::util::prompts::drop_exact_replay(&mut redo);
371        RoutingOutcome::route_to(executor.clone(), executor_fallbacks.to_vec(), redo)
372    }
373
374    /// Consults the advisor over the buffered transcript and parses the
375    /// verdict. `Ok(Failed)` covers fail-open errors and unparseable replies
376    /// (the caller refunds); fail-closed errors return `Err`.
377    async fn consult(
378        &self,
379        driver: &Driver,
380        base: &Request,
381        review_tail: Option<&str>,
382        trigger: &'static str,
383    ) -> Result<ConsultOutcome> {
384        // The advisor reviews the FULL transcript: system/developer content is
385        // normalized out of `messages` into `instructions`, so prepend it back
386        // as leading messages (identical {role, content} shape) — the task
387        // constraints the verdict must check against usually live there.
388        let transcript_messages: Vec<Message> = base
389            .llm_request
390            .instructions
391            .iter()
392            .map(|block| Message {
393                role: block.role,
394                content: block.content.clone(),
395            })
396            .chain(base.llm_request.messages.iter().cloned())
397            .collect();
398        let transcript = review_transcript(
399            &transcript_messages,
400            review_tail,
401            self.config.transcript_max_chars,
402        );
403        let consult_request = self.build_consult_request(base, transcript);
404        let started = Instant::now();
405        // An unresolvable advisor is treated like any other consult failure, so
406        // fail_open still returns the buffered executor turn to the client.
407        let reply = match driver.first_model_for(&Category::Judge) {
408            Ok(advisor) => {
409                let advisor = advisor.clone();
410                let advisor_models = driver.models_for(&Category::Judge).to_vec();
411                match driver
412                    .call_model_with_error_recovery(
413                        consult_request,
414                        advisor_models,
415                        self.config.fail_open,
416                    )
417                    .await
418                {
419                    Ok(response) => {
420                        let served_advisor = response
421                            .served_model()
422                            .cloned()
423                            .unwrap_or_else(|| advisor.clone());
424                        response
425                            .llm_response
426                            .into_agg()
427                            .await
428                            .map_err(|source| LibsyError::client_call(served_advisor, source))
429                    }
430                    Err(error) => Err(error),
431                }
432            }
433            Err(error) => Err(error),
434        };
435        let latency_ms = started.elapsed().as_secs_f64() * 1000.0;
436        let agg = match reply {
437            Ok(agg) => agg,
438            Err(error) => {
439                let reason = crate::algorithms::util::llm_judge::libsy_error_reason(&error);
440                record_consult_failure(reason);
441                if !self.config.fail_open {
442                    // Surface as an algorithm failure (5xx), never as the
443                    // advisor's own client error: a typed ContextWindowExceeded
444                    // from the consult would otherwise reach the client as 400
445                    // context_length_exceeded and trigger compaction of a
446                    // healthy conversation.
447                    return Err(algorithm_error(format!(
448                        "advisor consult failed (fail_open = false): {error}"
449                    )));
450                }
451                // An upstream error's Display can quote request content back,
452                // so only its redacted summary reaches logs or the audit trail.
453                let summary = safe_error_summary(&error);
454                tracing::warn!(
455                    target: "libsy",
456                    error = %summary,
457                    "advisor gate: consult failed; passing the turn through (fail open)"
458                );
459                emit_review_audit(ReviewAudit {
460                    verdict: "APPROVE",
461                    error: Some(summary),
462                    latency_ms,
463                    reply_head: None,
464                    usage: None,
465                });
466                return Ok(ConsultOutcome::Failed { reason });
467            }
468        };
469        let reply_text = advisor_reply_text(&agg);
470        let reply_head: String = reply_text.chars().take(160).collect();
471        match parse_verdict(&self.verdict_re, &reply_text) {
472            Some(Verdict::Approve) => {
473                record_review("approve", trigger);
474                emit_review_audit(ReviewAudit {
475                    verdict: "APPROVE",
476                    error: None,
477                    latency_ms,
478                    reply_head: Some(reply_head),
479                    usage: Some(&agg.usage),
480                });
481                Ok(ConsultOutcome::Approve)
482            }
483            Some(Verdict::Redo { plan }) => {
484                record_review("redo", trigger);
485                emit_review_audit(ReviewAudit {
486                    verdict: "REDO",
487                    error: None,
488                    latency_ms,
489                    reply_head: Some(reply_head),
490                    usage: Some(&agg.usage),
491                });
492                Ok(ConsultOutcome::Redo { plan })
493            }
494            None => {
495                // The advisor spent real tokens on a reply the gate cannot
496                // act on; the observer already recorded them. Refunded by
497                // the caller so a flaky advisor cannot burn the budget.
498                record_review("unparseable", trigger);
499                emit_review_audit(ReviewAudit {
500                    verdict: "UNPARSEABLE",
501                    error: None,
502                    latency_ms,
503                    reply_head: Some(reply_head),
504                    usage: Some(&agg.usage),
505                });
506                Ok(ConsultOutcome::Failed {
507                    reason: "parse_error",
508                })
509            }
510        }
511    }
512
513    /// A fresh, buffered, tool-free request carrying the reviewer contract and
514    /// the serialized transcript; metadata is kept for session correlation.
515    fn build_consult_request(&self, base: &Request, transcript: String) -> Request {
516        Request {
517            llm_request: LlmRequest {
518                model: base.llm_request.model.clone(),
519                instructions: vec![InstructionBlock {
520                    role: Role::System,
521                    content: vec![ContentBlock::Text {
522                        text: self.config.reviewer_system_prompt.clone(),
523                    }],
524                }],
525                messages: vec![Message::text(Role::User, transcript)],
526                sampling: SamplingParams {
527                    temperature: self.config.advisor_temperature,
528                    ..SamplingParams::default()
529                },
530                output: OutputParams {
531                    max_output_tokens: Some(self.config.advisor_max_tokens),
532                    ..OutputParams::default()
533                },
534                ..LlmRequest::default()
535            },
536            raw_request: None,
537            metadata: base.metadata.clone(),
538        }
539    }
540}
541
542#[async_trait::async_trait]
543impl Algorithm for AdvisorGate {
544    fn name(&self) -> &str {
545        "advisor_gate"
546    }
547
548    async fn route(self: Arc<Self>, driver: Driver, request: Request) -> Result<RoutingOutcome> {
549        let scope = budget_scope(&request);
550        let session_final = request
551            .metadata
552            .as_ref()
553            .and_then(|metadata| metadata.session_final)
554            == Some(true);
555        let result = self.route_inner(&driver, request, &scope).await;
556        if session_final {
557            self.budget.evict_scope(&scope);
558        }
559        result
560    }
561}
562
563/// Outcome of one consult; `Failed` = fail-open error or unparseable reply.
564enum ConsultOutcome {
565    Approve,
566    Redo { plan: String },
567    Failed { reason: &'static str },
568}
569
570fn algorithm_error(message: impl Into<String>) -> LibsyError {
571    LibsyError::AlgorithmError {
572        message: message.into(),
573    }
574}