返回 CodeWhale
superfast.rs
根目录 / crates / tui / src / superfast.rs
1 //! Superfast Decision Gate — a small, off-by-default "System One" front-door
2 //! classifier for the agent turn loop (#6603).
3 //!
4 //! Every user message normally wakes a large, slow, expensive model just to
5 //! decide intent and whether a tool is needed. The Decision Gate asks a small,
6 //! fast decision model (a Jev-compatible System One endpoint) those routine
7 //! questions in one non-generating pass, and derives a conservative routing
8 //! recommendation. A later, validated step could send each turn to the cheapest
9 //! correct path; this first increment only measures and logs.
10 //!
11 //! Design contract:
12 //! - Off by default. Nothing runs unless `SUPERFAST_ENABLED` is truthy, and
13 //! an enabled gate also needs `SUPERFAST_PROVIDER` naming the route it may
14 //! call, so turning the gate on never picks a paid endpoint by itself.
15 //! - Shadow mode. The gate classifies the turn and logs a typed
16 //! [`ShadowOutcome`] through `tracing` (target `superfast`), but never
17 //! changes routing, never skips or delays the model call, and never alters
18 //! any user-visible behavior.
19 //! - Fail open. Any error, timeout, non-2xx response, unreachable backend,
20 //! or malformed body is a typed failure class and the turn continues
21 //! exactly as if the gate were off. The call runs on a detached task.
22 //! - One transport. The call goes through the existing System One client
23 //! (`client::system_one`: `CodewhaleClient::for_decision_route` and
24 //! `system_one_decide`) that serves the `[auto.router] kind = "decision"`
25 //! router — same auth, TLS, secret redaction and one-attempt policy. There
26 //! is no second HTTP client.
27 //!
28 //! Configuration (environment, read when a turn starts):
29 //! - `SUPERFAST_ENABLED` — `1` / `true` / `yes` / `on` turns the gate on.
30 //! - `SUPERFAST_PROVIDER` — `typesafe` or `openrouter` (required when on).
31 //! The key comes from the same place the decision router reads it.
32 //! - `SUPERFAST_BASE_URL` — optional TypeSafe-route base override, e.g. a
33 //! self-hosted Jev server at `http://localhost:8000/v1`; `/systemone` is
34 //! appended.
35 //! - `SUPERFAST_MODEL` — decision model id (default `jev-latest` for
36 //! TypeSafe, `~typesafe/jev-latest` for OpenRouter).
37 //! - `SUPERFAST_TIMEOUT_MS` — per-call deadline, 1..=10000 (default 150).
38 //!
39 //! Known limits (written down so nobody assumes them):
40 //! - Shadow only: the recommendation is logged, never acted on.
41 //! - Only the latest user message's text is sent, truncated to 4,000
42 //! characters and redacted of configured secrets. No prompt text is
43 //! logged; the log carries the route, failure class and latency.
44 //! - The TypeSafe route needs a TypeSafe key even for a self-hosted server
45 //! (the transport always authenticates); a server that ignores auth can
46 //! be given any placeholder key.
47 //! - Usage settles through the originating turn's shared ledger. Missing
48 //! usage or cancellation after dispatch records a coverage gap. TypeSafe
49 //! is an unpriced Custom route until a billing basis is reviewed.
50 //! - Misconfiguration while enabled (missing or unknown provider, bad
51 //! timeout, missing key) is logged at `warn` for each turn and nothing is
52 //! sent.
53 //!
54 //! The Decision Gate concept and the reference implementation are by Andrea
55 //! Bruno, released under Creative Commons Attribution 4.0 (CC BY 4.0); see
56 //! <https://github.com/Andrea-Bruno/harness-superfast>. The decision models
57 //! (Von, OpenJev, Laya) are third-party open models; only the integration
58 //! architecture and the routing method here are covered by that attribution.
59
60 use std::sync::atomic::{AtomicBool, Ordering};
61 use std::time::{Duration, Instant};
62
63 use serde_json::{Value, json};
64
65 use codewhale_core::request::{ContentBlock, Message};
66 use codewhale_core::role::Role;
67
68 use crate::client::CodewhaleClient;
69 use crate::client::system_one::{DecisionRouterRoute, SystemOneResponse};
70 use crate::config::Config;
71 use crate::model_routing::{
72 AutoRouterFailure, auto_route_usage_source_id, decision_usage_batch, truncate_for_auto_router,
73 };
74 use tokio_util::sync::CancellationToken;
75
76 #[derive(Clone)]
77 struct ShadowUsageContext {
78 scope: crate::cost_status::CostScopeToken,
79 runtime_owner: Option<String>,
80 // Retains the origin's existing durable ledger until settlement finishes.
81 _lease: Option<crate::cost_status::RuntimeUsageLease>,
82 #[cfg(test)]
83 test_origin: std::thread::ThreadId,
84 }
85
86 impl ShadowUsageContext {
87 fn capture(owner: Option<&str>) -> Self {
88 let owner = owner.map(str::trim).filter(|owner| !owner.is_empty());
89 Self {
90 scope: crate::cost_status::scope_token(),
91 runtime_owner: owner.map(str::to_owned),
92 _lease: owner.and_then(crate::cost_status::acquire_runtime_usage_lease),
93 #[cfg(test)]
94 test_origin: crate::cost_status::test_cost_scope_id(),
95 }
96 }
97
98 async fn report(&self, batch: crate::cost_status::RuntimeUsageBatch) {
99 let context = self.clone();
100 // Existing sinks may write the origin-session ledger. Keep their
101 // filesystem/SQLite work off Tokio workers as well.
102 let settled = tokio::task::spawn_blocking(move || {
103 #[cfg(test)]
104 let _origin = crate::cost_status::bind_test_cost_scope(context.test_origin);
105 crate::cost_status::report_runtime_usage_batch(
106 context.scope,
107 context.runtime_owner.as_deref(),
108 &batch,
109 );
110 // Release the lease while its captured test binding is still live.
111 drop(context);
112 })
113 .await;
114 if settled.is_err() {
115 tracing::warn!(target: "superfast", "decision usage settlement worker failed");
116 }
117 }
118 }
119
120 /// Master switch. The gate never runs unless this env var is truthy.
121 const ENABLED_VAR: &str = "SUPERFAST_ENABLED";
122 /// Which System One route the gate may call: `typesafe` or `openrouter`.
123 const PROVIDER_VAR: &str = "SUPERFAST_PROVIDER";
124 /// Optional TypeSafe-route base URL (`/systemone` is appended).
125 const BASE_URL_VAR: &str = "SUPERFAST_BASE_URL";
126 /// Decision model id sent in the request body.
127 const MODEL_VAR: &str = "SUPERFAST_MODEL";
128 /// Per-call deadline in milliseconds.
129 const TIMEOUT_VAR: &str = "SUPERFAST_TIMEOUT_MS";
130
131 const DEFAULT_TIMEOUT_MS: u64 = 150;
132 const MAX_TIMEOUT_MS: u64 = 10_000;
133 /// Characters of the latest user message sent as decision state.
134 const MAX_STATE_CHARS: usize = 4_000;
135
136 /// Conservative routing recommendation derived from a turn's answers. Only a
137 /// decisive set of numbers produces a fast route; anything else is `Unknown`,
138 /// which means "fall back to the full model exactly as today".
139 #[derive(Debug, Clone, Copy, PartialEq, Eq)]
140 pub(crate) enum Route {
141 NeedsTool,
142 AnswerFromContext,
143 PlainChat,
144 Unknown,
145 }
146
147 impl Route {
148 fn as_str(self) -> &'static str {
149 match self {
150 Route::NeedsTool => "needs_tool",
151 Route::AnswerFromContext => "answer_from_context",
152 Route::PlainChat => "plain_chat",
153 Route::Unknown => "unknown",
154 }
155 }
156 }
157
158 /// Everything one shadow evaluation can end in. Bounded by construction: a
159 /// route or a non-secret failure class, plus the measured latency. Provider
160 /// bodies and prompt text never enter this type.
161 #[derive(Debug, Clone, Copy, PartialEq, Eq)]
162 pub(crate) enum ShadowOutcome {
163 /// The decision model answered; `route` is the derived recommendation.
164 Recommendation { route: Route, latency_ms: u64 },
165 /// Fail open: nothing reached the model request. `NotRunnable` means the
166 /// route could not be built (no key, bad URL) and nothing was sent.
167 Failed {
168 failure: AutoRouterFailure,
169 latency_ms: u64,
170 },
171 }
172
173 /// The route an enabled gate calls, resolved from the environment.
174 #[derive(Debug, Clone, PartialEq, Eq)]
175 pub(crate) struct ShadowSettings {
176 route: DecisionRouterRoute,
177 base_url: Option<String>,
178 model: String,
179 timeout: Duration,
180 }
181
182 impl ShadowSettings {
183 /// `None` when the gate is off; `Some(Err)` when it is on but
184 /// misconfigured (the message names the variable to fix).
185 fn from_env() -> Option<Result<Self, String>> {
186 Self::from_lookup(|name| std::env::var(name).ok())
187 }
188
189 fn from_lookup(lookup: impl Fn(&str) -> Option<String>) -> Option<Result<Self, String>> {
190 let value = |name: &str| {
191 lookup(name)
192 .map(|value| value.trim().to_string())
193 .filter(|value| !value.is_empty())
194 };
195 let enabled = value(ENABLED_VAR).is_some_and(|flag| {
196 matches!(
197 flag.to_ascii_lowercase().as_str(),
198 "1" | "true" | "yes" | "on"
199 )
200 });
201 enabled.then(|| Self::parse(&value))
202 }
203
204 fn parse(value: &impl Fn(&str) -> Option<String>) -> Result<Self, String> {
205 let provider = value(PROVIDER_VAR).ok_or_else(|| {
206 format!(
207 "{ENABLED_VAR} is set but {PROVIDER_VAR} is not; set it to typesafe or openrouter"
208 )
209 })?;
210 let route = DecisionRouterRoute::parse(&provider).ok_or_else(|| {
211 format!(
212 "{PROVIDER_VAR}={provider:?} is not a decision route; use typesafe or openrouter"
213 )
214 })?;
215 let timeout_ms = match value(TIMEOUT_VAR) {
216 None => DEFAULT_TIMEOUT_MS,
217 Some(raw) => raw
218 .parse::<u64>()
219 .ok()
220 .filter(|ms| (1..=MAX_TIMEOUT_MS).contains(ms))
221 .ok_or_else(|| {
222 format!("{TIMEOUT_VAR}={raw:?} must be 1..={MAX_TIMEOUT_MS} milliseconds")
223 })?,
224 };
225 let model = value(MODEL_VAR).unwrap_or_else(|| {
226 match route {
227 DecisionRouterRoute::Typesafe => "jev-latest",
228 DecisionRouterRoute::Openrouter => "~typesafe/jev-latest",
229 }
230 .to_string()
231 });
232 Ok(Self {
233 route,
234 base_url: value(BASE_URL_VAR),
235 model,
236 timeout: Duration::from_millis(timeout_ms),
237 })
238 }
239 }
240
241 /// Text of the last user message, or `None` when there is none or it is blank.
242 fn last_user_text(messages: &[Message]) -> Option<String> {
243 let message = messages.iter().rev().find(|m| m.role == Role::User)?;
244 let text = message
245 .content
246 .iter()
247 .filter_map(|block| match block {
248 ContentBlock::Text { text, .. } => Some(text.as_str()),
249 _ => None,
250 })
251 .collect::<Vec<_>>()
252 .join("\n");
253 let trimmed = text.trim();
254 (!trimmed.is_empty()).then(|| trimmed.to_string())
255 }
256
257 /// Fire the shadow gate for a turn's first model request. Returns at once:
258 /// the detached evaluation uses the turn's cancellation and accounting owner. No task is started
259 /// (and `None` is returned) when the gate is off or misconfigured, when there
260 /// is no user text, or when no Tokio runtime is present. The handle exists
261 /// for tests; the turn loop drops it.
262 pub(crate) fn spawn_shadow_gate(
263 config: &Config,
264 messages: &[Message],
265 runtime_owner: Option<&str>,
266 cancel_token: &CancellationToken,
267 ) -> Option<tokio::task::JoinHandle<ShadowOutcome>> {
268 let settings = match ShadowSettings::from_env()? {
269 Ok(settings) => settings,
270 Err(message) => {
271 tracing::warn!(target: "superfast", "decision gate (shadow) not run: {message}");
272 return None;
273 }
274 };
275 let latest_request = last_user_text(messages)?;
276 let runtime = tokio::runtime::Handle::try_current().ok()?;
277 let config = config.clone();
278 // Capture before detaching; the next session/turn must never acquire it.
279 let usage_context = ShadowUsageContext::capture(runtime_owner);
280 let cancel_token = cancel_token.clone();
281 Some(runtime.spawn(async move {
282 let outcome = evaluate(
283 &config,
284 &settings,
285 &latest_request,
286 &usage_context,
287 &cancel_token,
288 )
289 .await;
290 log_outcome(outcome);
291 outcome
292 }))
293 }
294
295 fn log_outcome(outcome: ShadowOutcome) {
296 match outcome {
297 ShadowOutcome::Recommendation { route, latency_ms } => tracing::info!(
298 target: "superfast",
299 route = route.as_str(),
300 latency_ms,
301 "decision gate (shadow) recommendation"
302 ),
303 ShadowOutcome::Failed {
304 failure: AutoRouterFailure::NotRunnable,
305 ..
306 } => tracing::warn!(
307 target: "superfast",
308 "decision gate (shadow) not run: route not runnable (check the {PROVIDER_VAR} key and {BASE_URL_VAR})"
309 ),
310 ShadowOutcome::Failed {
311 failure,
312 latency_ms,
313 } => tracing::debug!(
314 target: "superfast",
315 failure = %failure.label(),
316 latency_ms,
317 "decision gate (shadow) no opinion (fail-open)"
318 ),
319 }
320 }
321
322 /// One shadow evaluation over the existing System One transport.
323 async fn evaluate(
324 config: &Config,
325 settings: &ShadowSettings,
326 latest_request: &str,
327 usage_context: &ShadowUsageContext,
328 cancel_token: &CancellationToken,
329 ) -> ShadowOutcome {
330 if cancel_token.is_cancelled() {
331 return ShadowOutcome::Failed {
332 failure: AutoRouterFailure::Cancelled,
333 latency_ms: 0,
334 };
335 }
336 // Client construction resolves keys (environment, secret store), so it
337 // runs off the async worker (#6149).
338 let built = {
339 let config = config.clone();
340 let route = settings.route;
341 let base_url = settings.base_url.clone();
342 #[cfg(test)]
343 let ticket = crate::test_support::env_scope_ticket();
344 tokio::task::spawn_blocking(move || {
345 #[cfg(test)]
346 let _membership = crate::test_support::join_env_scope(ticket);
347 CodewhaleClient::for_decision_route(&config, route, base_url.as_deref())
348 })
349 .await
350 };
351 let Ok(Ok(client)) = built else {
352 return ShadowOutcome::Failed {
353 failure: AutoRouterFailure::NotRunnable,
354 latency_ms: 0,
355 };
356 };
357 let body = decision_body(&client, &settings.model, latest_request);
358 let request_route = client.effective_route_envelope(&settings.model, chrono::Utc::now());
359 let dispatched = AtomicBool::new(false);
360 let started = Instant::now();
361 let answer = tokio::select! {
362 biased;
363 () = cancel_token.cancelled() => Err(AutoRouterFailure::Cancelled),
364 result = tokio::time::timeout(settings.timeout, client.system_one_decide(&body, &dispatched)) => {
365 result.unwrap_or(Err(AutoRouterFailure::Timeout))
366 }
367 };
368 let latency_ms = u64::try_from(started.elapsed().as_millis()).unwrap_or(u64::MAX);
369 match answer {
370 Err(failure) => {
371 if dispatched.load(Ordering::Acquire) {
372 usage_context
373 .report(crate::cost_status::RuntimeUsageBatch {
374 decisions: Vec::new(),
375 drop_records: vec![crate::cost_status::RuntimeUsageDropRecord {
376 reason: crate::cost_status::RuntimeUsageMissingReason::default(),
377 source_id: auto_route_usage_source_id(
378 &request_route,
379 "shadow:dispatched-unreceipted",
380 ),
381 route: request_route.sanitized_for_persistence(),
382 }],
383 dropped_records: 1,
384 ..Default::default()
385 })
386 .await;
387 }
388 ShadowOutcome::Failed {
389 failure,
390 latency_ms,
391 }
392 }
393 Ok(response) => {
394 // Billing evidence settles even when strict answer validation
395 // prevents the policy from producing a recommendation.
396 let mut batch = decision_usage_batch(&request_route, &response);
397 let intent = crate::model_routing::validated_choice(
398 response.answers.get("intent"),
399 &["code_change", "code_question", "command", "chat", "other"],
400 );
401 batch
402 .decisions
403 .push(crate::cost_status::RuntimeDecisionReceipt {
404 source_id: auto_route_usage_source_id(
405 &request_route,
406 response.id.as_deref().unwrap_or("systemone"),
407 ),
408 route: request_route.sanitized_for_persistence(),
409 usage: response.usage.as_ref().map(|u| codewhale_models::Usage {
410 input_tokens: u.input_tokens,
411 output_tokens: u.output_tokens,
412 ..Default::default()
413 }),
414 usage_complete: response
415 .usage
416 .as_ref()
417 .is_some_and(|u| u.complete && (u.input_tokens > 0 || u.output_tokens > 0)),
418 shadow: true,
419 valid_answers: response.answers_validated != Some(false),
420 evidence: crate::model_routing::AutoRouteDecisionEvidence {
421 choice: intent
422 .as_ref()
423 .map_or_else(|| "invalid".to_string(), |v| v.choice.clone()),
424 probabilities_bp: intent
425 .as_ref()
426 .map_or_else(Default::default, |v| v.probabilities_bp.clone()),
427 confidence_bp: intent.as_ref().map_or(0, |v| v.confidence_bp),
428 min_confidence_bp: 5_000,
429 cost_saving_kept_fast: false,
430 thinking: None,
431 provider_reported_cost_usd: response
432 .usage
433 .as_ref()
434 .and_then(|u| u.reported_cost()),
435 latency_ms,
436 response_model: response.model.clone(),
437 },
438 });
439 usage_context.report(batch).await;
440 if response.answers_validated == Some(false) {
441 ShadowOutcome::Failed {
442 failure: AutoRouterFailure::InvalidAnswer,
443 latency_ms,
444 }
445 } else {
446 ShadowOutcome::Recommendation {
447 route: derive_route(&response),
448 latency_ms,
449 }
450 }
451 }
452 }
453 }
454
455 /// The System One request: two `noul` questions and one `choice` over the
456 /// redacted, bounded latest request.
457 fn decision_body(client: &CodewhaleClient, model: &str, latest_request: &str) -> Value {
458 json!({
459 "model": model,
460 "state": {
461 "latest_request": client
462 .redact_model_bound_text(&truncate_for_auto_router(latest_request, MAX_STATE_CHARS)),
463 },
464 "questions": {
465 "needs_tool": {
466 "type": "noul",
467 "instructions": "Does answering this request require taking an action with a tool (reading, writing, running, searching), rather than replying from what is already known?"
468 },
469 "answerable_from_context": {
470 "type": "noul",
471 "instructions": "Can this request be answered from information already present in the conversation, without any new investigation?"
472 },
473 "intent": {
474 "type": "choice",
475 "instructions": "Classify the primary intent of the user request.",
476 "criteria": {
477 "code_change": "Create, edit, or delete code or files.",
478 "code_question": "Explain or reason about code without changing it.",
479 "command": "Run a command or operation.",
480 "chat": "Casual conversation or a question needing no tools.",
481 "other": "None of the above."
482 }
483 }
484 }
485 })
486 }
487
488 /// A finite probability in [0, 1], or `None` ("no evidence").
489 fn unit_interval(value: Option<f64>) -> Option<f64> {
490 value.filter(|value| value.is_finite() && (0.0..=1.0).contains(value))
491 }
492
493 /// Read a `noul` answer only when it is typed as one and its value is a real
494 /// probability. Anything else (absent, wrong type, out of range) is "no
495 /// evidence", so a mis-scaled or missing answer can never produce a decisive
496 /// fast route.
497 fn read_noul(response: &SystemOneResponse, key: &str) -> Option<f64> {
498 let answer = response.answers.get(key)?;
499 (answer.kind == "noul")
500 .then_some(answer.noul)
501 .and_then(unit_interval)
502 }
503
504 /// Derive a conservative route. The gate only recommends a fast route when the
505 /// relevant numbers are decisive; otherwise it says `Unknown` so the caller
506 /// falls back to the normal path.
507 fn derive_route(response: &SystemOneResponse) -> Route {
508 let needs_tool = read_noul(response, "needs_tool");
509 let from_context = read_noul(response, "answerable_from_context");
510
511 // Decisive "needs a tool" wins first — the harness must not skip work.
512 if needs_tool.is_some_and(|nt| nt >= 0.85) {
513 return Route::NeedsTool;
514 }
515
516 // Strongly answerable from context, with a present and low tool-need signal.
517 if from_context.is_some_and(|fc| fc >= 0.85) && needs_tool.is_some_and(|nt| nt <= 0.3) {
518 return Route::AnswerFromContext;
519 }
520
521 // Clearly chat, with a calibrated intent and a present, low tool-need signal.
522 let intent_is_calibrated_chat = crate::model_routing::validated_choice(
523 response.answers.get("intent"),
524 &["code_change", "code_question", "command", "chat", "other"],
525 )
526 .is_some_and(|intent| intent.choice == "chat" && intent.confidence_bp >= 5_000);
527 if intent_is_calibrated_chat && needs_tool.is_some_and(|nt| nt <= 0.2) {
528 return Route::PlainChat;
529 }
530
531 Route::Unknown
532 }
533
534 #[cfg(test)]
535 mod tests {
536 use super::*;
537 use wiremock::matchers::{header, method, path};
538 use wiremock::{Mock, MockServer, Request, ResponseTemplate};
539
540 const TYPESAFE_TEST_KEY: &str = "sf-typesafe-test-key-0123456789";
541
542 /// Fields drop in declaration order: restore the environment before the
543 /// lock is released, or another test observes our overrides.
544 struct Env {
545 _guards: Vec<crate::test_support::EnvVarGuard>,
546 _home: tempfile::TempDir,
547 _lock: crate::test_support::TestEnvLock,
548 }
549
550 /// A hermetic home with a TypeSafe key. `gate` switches the gate on
551 /// against the TypeSafe route at that server with that deadline (ms).
552 fn hermetic_env(gate: Option<(&MockServer, u64)>) -> Env {
553 use crate::test_support::EnvVarGuard;
554 let lock = crate::test_support::lock_test_env();
555 let home = tempfile::tempdir().expect("test home");
556 let mut guards = vec![
557 EnvVarGuard::set("CODEWHALE_HOME", home.path()),
558 EnvVarGuard::remove("OPENROUTER_API_KEY"),
559 EnvVarGuard::set("TYPESAFE_API_KEY", TYPESAFE_TEST_KEY),
560 EnvVarGuard::remove(MODEL_VAR),
561 ];
562 match gate {
563 Some((server, timeout_ms)) => guards.extend([
564 EnvVarGuard::set(ENABLED_VAR, "1"),
565 EnvVarGuard::set(PROVIDER_VAR, "typesafe"),
566 EnvVarGuard::set(BASE_URL_VAR, format!("{}/v1", server.uri())),
567 EnvVarGuard::set(TIMEOUT_VAR, timeout_ms.to_string()),
568 ]),
569 None => guards.extend([
570 EnvVarGuard::remove(ENABLED_VAR),
571 EnvVarGuard::remove(PROVIDER_VAR),
572 EnvVarGuard::remove(BASE_URL_VAR),
573 EnvVarGuard::remove(TIMEOUT_VAR),
574 ]),
575 }
576 Env {
577 _guards: guards,
578 _home: home,
579 _lock: lock,
580 }
581 }
582
583 /// A chat provider the client can be built from; the gate never calls it.
584 fn config() -> Config {
585 Config {
586 provider: Some("deepseek".to_string()),
587 default_text_model: Some("deepseek-v4-pro".to_string()),
588 providers: Some(crate::config::ProvidersConfig {
589 deepseek: crate::config::ProviderConfig {
590 api_key: Some("ds-test-key".to_string()),
591 ..Default::default()
592 },
593 ..Default::default()
594 }),
595 ..Default::default()
596 }
597 }
598
599 fn text(role: Role, text: &str) -> Message {
600 Message {
601 role,
602 content: vec![ContentBlock::Text {
603 text: text.to_string(),
604 cache_control: None,
605 }],
606 }
607 }
608
609 fn turn(latest: &str) -> Vec<Message> {
610 vec![
611 text(Role::User, "earlier question"),
612 text(Role::Assistant, "earlier answer"),
613 text(Role::User, latest),
614 ]
615 }
616
617 fn noul_body(needs_tool: f64, from_context: f64) -> Value {
618 json!({
619 "id": "sf-1",
620 "model": "jev-latest",
621 "usage": { "input_tokens": 121, "output_tokens": 8 },
622 "answers": {
623 "needs_tool": { "type": "noul", "noul": needs_tool },
624 "answerable_from_context": { "type": "noul", "noul": from_context },
625 "intent": {
626 "type": "choice",
627 "choice": "code_change",
628 "probabilities": {
629 "code_change": 0.9, "code_question": 0.04, "command": 0.03,
630 "chat": 0.02, "other": 0.01
631 },
632 "confidence": 0.8
633 }
634 }
635 })
636 }
637
638 async fn requests(server: &MockServer) -> Vec<Request> {
639 server.received_requests().await.expect("recorded")
640 }
641
642 async fn run(messages: &[Message]) -> Option<ShadowOutcome> {
643 let handle = spawn_shadow_gate(&config(), messages, None, &CancellationToken::new())?;
644 Some(handle.await.expect("shadow task"))
645 }
646
647 #[tokio::test]
648 async fn disabled_gate_sends_nothing_and_starts_no_task() {
649 let server = MockServer::start().await;
650 Mock::given(method("POST"))
651 .respond_with(ResponseTemplate::new(200).set_body_json(noul_body(0.9, 0.1)))
652 .mount(&server)
653 .await;
654 let _env = hermetic_env(None);
655
656 assert_eq!(run(&turn("Refactor the parser")).await, None);
657 assert!(
658 requests(&server).await.is_empty(),
659 "a disabled gate must not call"
660 );
661 }
662
663 #[tokio::test]
664 async fn enabled_gate_posts_one_redacted_bounded_decision_and_recommends() {
665 let server = MockServer::start().await;
666 Mock::given(method("POST"))
667 .and(path("/v1/systemone"))
668 .and(header(
669 "authorization",
670 format!("Bearer {TYPESAFE_TEST_KEY}").as_str(),
671 ))
672 .respond_with(ResponseTemplate::new(200).set_body_json(noul_body(0.93, 0.1)))
673 .expect(1)
674 .mount(&server)
675 .await;
676 let _env = hermetic_env(Some((&server, 2_000)));
677
678 // The latest message echoes the key and is longer than the state cap.
679 let latest = format!(
680 "Fix the build. token={TYPESAFE_TEST_KEY} {}",
681 "x".repeat(6_000)
682 );
683 let outcome = run(&turn(&latest)).await.expect("enabled gate runs");
684 assert!(
685 matches!(
686 outcome,
687 ShadowOutcome::Recommendation {
688 route: Route::NeedsTool,
689 ..
690 }
691 ),
692 "{outcome:?}"
693 );
694
695 let requests = requests(&server).await;
696 assert_eq!(requests.len(), 1);
697 let raw = String::from_utf8(requests[0].body.clone()).expect("utf8 body");
698 assert!(
699 !raw.contains(TYPESAFE_TEST_KEY),
700 "key leaked into decision state"
701 );
702 assert!(
703 !raw.contains("earlier question"),
704 "only the latest message is sent"
705 );
706 let body: Value = serde_json::from_str(&raw).expect("json body");
707 assert_eq!(body["model"], "jev-latest");
708 let state = body["state"]["latest_request"]
709 .as_str()
710 .expect("state text");
711 assert!(state.starts_with("Fix the build."));
712 assert!(
713 state.chars().count() <= MAX_STATE_CHARS + 3,
714 "state is bounded"
715 );
716 let questions: Vec<&str> = body["questions"]
717 .as_object()
718 .expect("questions")
719 .keys()
720 .map(String::as_str)
721 .collect();
722 assert_eq!(
723 questions,
724 ["needs_tool", "answerable_from_context", "intent"]
725 );
726 }
727
728 #[tokio::test]
729 async fn http_error_and_malformed_answers_fail_open_as_typed_classes() {
730 for (response, expected) in [
731 (
732 ResponseTemplate::new(500).set_body_string("upstream exploded: secret body"),
733 AutoRouterFailure::Http { status: 500 },
734 ),
735 (
736 ResponseTemplate::new(200).set_body_string("not json"),
737 AutoRouterFailure::InvalidAnswer,
738 ),
739 ] {
740 let server = MockServer::start().await;
741 Mock::given(method("POST"))
742 .and(path("/v1/systemone"))
743 .respond_with(response)
744 .expect(1)
745 .mount(&server)
746 .await;
747 let _env = hermetic_env(Some((&server, 2_000)));
748
749 let outcome = run(&turn("Explain the diff"))
750 .await
751 .expect("enabled gate runs");
752 assert!(
753 matches!(outcome, ShadowOutcome::Failed { failure, .. } if failure == expected),
754 "{outcome:?}"
755 );
756 }
757 }
758
759 #[tokio::test]
760 async fn slow_endpoint_times_out_without_holding_up_the_caller() {
761 let server = MockServer::start().await;
762 Mock::given(method("POST"))
763 .and(path("/v1/systemone"))
764 .respond_with(
765 ResponseTemplate::new(200)
766 .set_body_json(noul_body(0.93, 0.1))
767 .set_delay(Duration::from_secs(3)),
768 )
769 .mount(&server)
770 .await;
771 let _env = hermetic_env(Some((&server, 100)));
772
773 let started = Instant::now();
774 let handle = spawn_shadow_gate(
775 &config(),
776 &turn("Run the tests"),
777 None,
778 &CancellationToken::new(),
779 )
780 .expect("task");
781 assert!(
782 started.elapsed() < Duration::from_millis(100),
783 "spawning must not wait on the decision call"
784 );
785 let outcome = handle.await.expect("shadow task");
786 assert!(
787 matches!(
788 outcome,
789 ShadowOutcome::Failed {
790 failure: AutoRouterFailure::Timeout,
791 ..
792 }
793 ),
794 "{outcome:?}"
795 );
796 assert!(
797 started.elapsed() < Duration::from_secs(2),
798 "the deadline bounds the call"
799 );
800 }
801
802 #[tokio::test]
803 async fn no_op_and_unrunnable_turns_send_nothing() {
804 let server = MockServer::start().await;
805 Mock::given(method("POST"))
806 .respond_with(ResponseTemplate::new(200).set_body_json(noul_body(0.9, 0.1)))
807 .mount(&server)
808 .await;
809 let _env = hermetic_env(Some((&server, 2_000)));
810
811 // No user text: no task.
812 assert_eq!(run(&[text(Role::Assistant, "only assistant")]).await, None);
813 assert_eq!(run(&turn(" ")).await, None);
814
815 // No TypeSafe key: the route is not runnable and nothing is sent.
816 let _no_key = crate::test_support::EnvVarGuard::remove("TYPESAFE_API_KEY");
817 let outcome = run(&turn("Refactor the parser"))
818 .await
819 .expect("enabled gate runs");
820 assert_eq!(
821 outcome,
822 ShadowOutcome::Failed {
823 failure: AutoRouterFailure::NotRunnable,
824 latency_ms: 0,
825 }
826 );
827 assert!(requests(&server).await.is_empty());
828 }
829
830 fn lookup(pairs: &'static [(&'static str, &'static str)]) -> impl Fn(&str) -> Option<String> {
831 move |name| {
832 pairs
833 .iter()
834 .find_map(|(key, value)| (*key == name).then(|| (*value).to_string()))
835 }
836 }
837
838 #[test]
839 fn settings_are_off_by_default_and_misconfiguration_is_named() {
840 assert_eq!(ShadowSettings::from_lookup(lookup(&[])), None);
841 assert_eq!(
842 ShadowSettings::from_lookup(lookup(&[(ENABLED_VAR, "0"), (PROVIDER_VAR, "typesafe")])),
843 None
844 );
845 let missing = ShadowSettings::from_lookup(lookup(&[(ENABLED_VAR, "on")]))
846 .expect("enabled")
847 .expect_err("no provider");
848 assert!(missing.contains(PROVIDER_VAR), "{missing}");
849 let unknown =
850 ShadowSettings::from_lookup(lookup(&[(ENABLED_VAR, "1"), (PROVIDER_VAR, "deepseek")]))
851 .expect("enabled")
852 .expect_err("not a decision route");
853 assert!(unknown.contains("deepseek"), "{unknown}");
854 let timeout = ShadowSettings::from_lookup(lookup(&[
855 (ENABLED_VAR, "1"),
856 (PROVIDER_VAR, "openrouter"),
857 (TIMEOUT_VAR, "0"),
858 ]))
859 .expect("enabled")
860 .expect_err("zero timeout");
861 assert!(timeout.contains(TIMEOUT_VAR), "{timeout}");
862 let settings = ShadowSettings::from_lookup(lookup(&[
863 (ENABLED_VAR, "yes"),
864 (PROVIDER_VAR, "openrouter"),
865 ]))
866 .expect("enabled")
867 .expect("valid");
868 assert_eq!(settings.route, DecisionRouterRoute::Openrouter);
869 assert_eq!(settings.model, "~typesafe/jev-latest");
870 assert_eq!(settings.timeout, Duration::from_millis(DEFAULT_TIMEOUT_MS));
871 assert_eq!(settings.base_url, None);
872 }
873
874 fn response(answers: Value) -> SystemOneResponse {
875 serde_json::from_value(json!({ "answers": answers })).expect("response")
876 }
877
878 #[test]
879 fn only_decisive_typed_answers_route_fast() {
880 let cases = [
881 (
882 json!({ "needs_tool": { "type": "noul", "noul": 0.9 } }),
883 Route::NeedsTool,
884 ),
885 (
886 json!({
887 "needs_tool": { "type": "noul", "noul": 0.1 },
888 "answerable_from_context": { "type": "noul", "noul": 0.9 }
889 }),
890 Route::AnswerFromContext,
891 ),
892 (
893 json!({
894 "needs_tool": { "type": "noul", "noul": 0.05 },
895 "intent": { "type": "choice", "choice": "chat", "confidence": 0.8,
896 "probabilities": { "chat": 0.9, "code_change": 0.04, "code_question": 0.03, "command": 0.02, "other": 0.01 } }
897 }),
898 Route::PlainChat,
899 ),
900 // Out of range, wrong type, or untyped answers are no evidence.
901 (
902 json!({
903 "needs_tool": { "type": "noul", "noul": 5.0 },
904 "answerable_from_context": { "type": "choice", "noul": 0.95 }
905 }),
906 Route::Unknown,
907 ),
908 (json!({ "needs_tool": { "noul": 0.95 } }), Route::Unknown),
909 // Chat without a present, low tool-need signal is not decisive.
910 (
911 json!({ "intent": { "type": "choice", "choice": "chat", "confidence": 0.9 } }),
912 Route::Unknown,
913 ),
914 (json!({}), Route::Unknown),
915 ];
916 for (answers, expected) in cases {
917 assert_eq!(
918 derive_route(&response(answers.clone())),
919 expected,
920 "{answers}"
921 );
922 }
923 }
924
925 #[test]
926 fn last_user_text_picks_the_latest_user_text() {
927 assert_eq!(
928 last_user_text(&turn(" latest question ")).as_deref(),
929 Some("latest question")
930 );
931 assert_eq!(
932 last_user_text(&[text(Role::Assistant, "only assistant")]),
933 None
934 );
935 }
936 #[tokio::test]
937 async fn cancelled_before_dispatch_sends_nothing() {
938 let server = MockServer::start().await;
939 let _env = hermetic_env(Some((&server, 2_000)));
940 let cancel = CancellationToken::new();
941 cancel.cancel();
942 let outcome = spawn_shadow_gate(&config(), &turn("hi"), None, &cancel)
943 .expect("task")
944 .await
945 .expect("joined");
946 assert!(matches!(
947 outcome,
948 ShadowOutcome::Failed {
949 failure: AutoRouterFailure::Cancelled,
950 ..
951 }
952 ));
953 assert!(requests(&server).await.is_empty());
954 }
955
956 #[tokio::test]
957 async fn cancelled_dispatched_call_records_gap_for_its_owner() {
958 let server = MockServer::start().await;
959 Mock::given(method("POST"))
960 .and(path("/v1/systemone"))
961 .respond_with(
962 ResponseTemplate::new(200)
963 .set_delay(Duration::from_secs(2))
964 .set_body_json(noul_body(0.9, 0.1)),
965 )
966 .expect(1)
967 .mount(&server)
968 .await;
969 let _env = hermetic_env(Some((&server, 3_000)));
970 let owner = "superfast-cancel-owner";
971 let dropped = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
972 let sink = dropped.clone();
973 crate::cost_status::register_runtime_usage_sink_with_drop(
974 owner,
975 std::sync::Arc::new(|_| true),
976 Some(std::sync::Arc::new(move |record| {
977 sink.lock().expect("drop sink").push(record);
978 true
979 })),
980 );
981
982 let cancel = CancellationToken::new();
983 let handle = spawn_shadow_gate(&config(), &turn("hi"), Some(owner), &cancel).expect("task");
984 tokio::time::timeout(Duration::from_secs(1), async {
985 while requests(&server).await.is_empty() {
986 tokio::task::yield_now().await;
987 }
988 })
989 .await
990 .expect("request admitted");
991 let start = Instant::now();
992 cancel.cancel();
993 let outcome = tokio::time::timeout(Duration::from_millis(300), handle)
994 .await
995 .expect("cancellation stops the pending request")
996 .expect("joined");
997 assert!(matches!(
998 outcome,
999 ShadowOutcome::Failed {
1000 failure: AutoRouterFailure::Cancelled,
1001 ..
1002 }
1003 ));
1004 assert!(start.elapsed() < Duration::from_millis(300));
1005 let records = dropped.lock().expect("drop records");
1006 assert_eq!(records.len(), 1);
1007 assert_eq!(records[0].route.provider_identity, "typesafe");
1008 let batch = crate::cost_status::take_runtime_usage(owner);
1009 assert!(batch.records.is_empty());
1010 assert!(batch.drop_records.is_empty());
1011 crate::cost_status::finish_runtime_usage_owner(owner);
1012 }
1013
1014 #[tokio::test]
1015 async fn owner_lease_retains_late_usage_and_invalid_answer_cost() {
1016 let server = MockServer::start().await;
1017 let mut body = noul_body(0.9, 0.1);
1018 body["answers"]["needs_tool"]["noul"] = 2.0.into();
1019 body["usage"]["cost"] = 0.000012054.into();
1020 Mock::given(method("POST"))
1021 .and(path("/v1/systemone"))
1022 .respond_with(
1023 ResponseTemplate::new(200)
1024 .set_delay(Duration::from_millis(100))
1025 .set_body_json(body),
1026 )
1027 .expect(1)
1028 .mount(&server)
1029 .await;
1030 let _env = hermetic_env(Some((&server, 2_000)));
1031 let owner = "superfast-late-owner";
1032 let observed = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
1033 let sink = observed.clone();
1034 crate::cost_status::register_runtime_usage_sink(
1035 owner,
1036 std::sync::Arc::new(move |record| {
1037 sink.lock().expect("sink").push(record);
1038 true
1039 }),
1040 );
1041 let decisions = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
1042 let decision_sink = decisions.clone();
1043 crate::cost_status::register_runtime_decision_sink(
1044 owner,
1045 std::sync::Arc::new(move |receipt| {
1046 decision_sink.lock().expect("decision sink").push(receipt);
1047 true
1048 }),
1049 );
1050 let handle = spawn_shadow_gate(
1051 &config(),
1052 &turn("hi"),
1053 Some(owner),
1054 &CancellationToken::new(),
1055 )
1056 .expect("task");
1057 crate::cost_status::finish_runtime_usage_owner(owner);
1058 let outcome = handle.await.expect("joined");
1059 assert!(matches!(
1060 outcome,
1061 ShadowOutcome::Failed {
1062 failure: AutoRouterFailure::InvalidAnswer,
1063 ..
1064 }
1065 ));
1066 let records = observed.lock().expect("records");
1067 assert_eq!(
1068 records.len(),
1069 1,
1070 "retired origin still owns its late response"
1071 );
1072 assert_eq!(records[0].usage.usage.input_tokens, 121);
1073 assert_eq!(records[0].usage.usage.output_tokens, 8);
1074 assert_eq!(records[0].usage.route.provider_identity, "typesafe");
1075 assert_eq!(
1076 records[0].usage.route.billing_mode,
1077 crate::cost_status::RouteBillingMode::Unknown
1078 );
1079 let receipts = decisions.lock().expect("decision receipts");
1080 assert_eq!(receipts.len(), 1);
1081 assert!(!receipts[0].valid_answers);
1082 assert!(receipts[0].shadow);
1083 assert_eq!(
1084 receipts[0].evidence.provider_reported_cost_usd.as_deref(),
1085 Some("0.000012054")
1086 );
1087 assert!(
1088 crate::cost_status::take_runtime_usage(owner)
1089 .records
1090 .is_empty(),
1091 "settled response must not enter fallback journal"
1092 );
1093 }
1094
1095 #[tokio::test]
1096 async fn oversized_decision_body_fails_closed_with_dispatched_coverage_gap() {
1097 let server = MockServer::start().await;
1098 Mock::given(method("POST"))
1099 .and(path("/v1/systemone"))
1100 .respond_with(ResponseTemplate::new(200).set_body_string("x".repeat(256 * 1024 + 1)))
1101 .expect(1)
1102 .mount(&server)
1103 .await;
1104 let _env = hermetic_env(Some((&server, 2_000)));
1105 let owner = "superfast-oversized-owner";
1106 let drops = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
1107 let sink = drops.clone();
1108 crate::cost_status::register_runtime_usage_sink_with_drop(
1109 owner,
1110 std::sync::Arc::new(|_| true),
1111 Some(std::sync::Arc::new(move |record| {
1112 sink.lock().expect("drops").push(record);
1113 true
1114 })),
1115 );
1116 let handle = spawn_shadow_gate(
1117 &config(),
1118 &turn("hi"),
1119 Some(owner),
1120 &CancellationToken::new(),
1121 )
1122 .expect("task");
1123 assert!(matches!(
1124 handle.await.expect("joined"),
1125 ShadowOutcome::Failed {
1126 failure: AutoRouterFailure::InvalidAnswer,
1127 ..
1128 }
1129 ));
1130 assert_eq!(drops.lock().expect("drop receipts").len(), 1);
1131 crate::cost_status::finish_runtime_usage_owner(owner);
1132 }
1133 #[tokio::test]
1134 async fn empty_origin_and_receipt_only_batch_use_the_existing_interactive_pool() {
1135 let _scope = crate::cost_status::test_scope();
1136 let context = ShadowUsageContext::capture(Some(" "));
1137 assert!(context.runtime_owner.is_none());
1138 context
1139 .report(crate::cost_status::RuntimeUsageBatch {
1140 decisions: vec![crate::cost_status::decision_receipt_fixture(
1141 "ownerless-fixture",
1142 )],
1143 ..Default::default()
1144 })
1145 .await;
1146 let projected = crate::cost_status::drain();
1147 assert!(
1148 projected
1149 .route_receipts
1150 .iter()
1151 .any(|r| r.contains("0.000012054"))
1152 );
1153 assert_eq!(
1154 projected.unpriced_turns, 0,
1155 "diagnostic receipts do not mint a new charge"
1156 );
1157 assert!(
1158 crate::cost_status::take_runtime_usage(" ")
1159 .decisions
1160 .is_empty()
1161 );
1162 }
1163 }
1164
1164 lines RUST