返回 CodeWhale
large_output_router.rs
根目录 / crates / tui / src / tools / large_output_router.rs
1 //! Adaptive evidence routing for tool results (#4619) — explicit opt-in.
2 //!
3 //! Classic bounded spillover is the default: `tools/truncate.rs` keeps results
4 //! at or under its byte threshold fully inline and gives larger ones a
5 //! head/tail preview plus a session artifact. Set
6 //! `CODEWHALE_ADAPTIVE_OUTPUT_ROUTING` to enable the adaptive lane, which
7 //! classifies results as inline, hybrid, or handle-only by estimated tokens
8 //! and publishes non-inline results exactly once under their origin session
9 //! with immutable evidence metadata for bounded retrieval.
10
11 use std::collections::HashMap;
12 use std::io;
13 use std::path::PathBuf;
14 use std::sync::{Mutex, OnceLock};
15
16 use serde::{Deserialize, Serialize};
17
18 use crate::tools::spec::ToolResult;
19
20 // ── Constants ──────────────────────────────────────────────────────────────────
21
22 /// Default token threshold separating hybrid from handle-only evidence.
23 ///
24 /// 32K tokens (≈96 KiB of text at the 3 chars/token estimate) keeps ordinary
25 /// tool results — file reads, test runs, build logs up to a few thousand
26 /// lines — fully inline. Only genuinely large outputs spill to evidence
27 /// artifacts, where the model-facing preview names the artifact path and how
28 /// to recover the omitted range.
29 pub const DEFAULT_LARGE_OUTPUT_THRESHOLD_TOKENS: usize = 32_768;
30
31 /// Approximate characters-per-token ratio used for the heuristic estimate.
32 /// We intentionally choose a conservative value (3 chars/token) so we err
33 /// on the side of routing rather than dumping raw data into the parent.
34 const CHARS_PER_TOKEN_ESTIMATE: usize = 3;
35
36 static ACTIVE_WORKSHOP: OnceLock<Mutex<WorkshopConfig>> = OnceLock::new();
37
38 #[cfg(test)]
39 static ACTIVE_WORKSHOP_TEST_SERIAL: OnceLock<Mutex<()>> = OnceLock::new();
40
41 #[cfg(test)]
42 std::thread_local! {
43 static ACTIVE_WORKSHOP_TEST_SERIAL_HELD: std::cell::Cell<bool> = const {
44 std::cell::Cell::new(false)
45 };
46 }
47
48 /// Holds every test-side workshop activation behind one process-wide gate.
49 /// The thread-local marker lets the owning current-thread test call
50 /// `install_active` without trying to acquire its own non-reentrant lock.
51 ///
52 /// The guard also owns the process-wide test env barrier, because holding this
53 /// gate across a `reload_config` is holding it across `Settings::load` and
54 /// `Config`'s env reads, which take that barrier (#6306). Field order is the
55 /// release order: the serial gate first, then the barrier.
56 #[cfg(test)]
57 pub(crate) struct ActiveWorkshopTestGuard {
58 _serial: std::sync::MutexGuard<'static, ()>,
59 _env: Option<crate::test_support::TestEnvLock>,
60 }
61
62 #[cfg(test)]
63 impl Drop for ActiveWorkshopTestGuard {
64 fn drop(&mut self) {
65 ACTIVE_WORKSHOP_TEST_SERIAL_HELD.with(|held| held.set(false));
66 }
67 }
68
69 /// Take the workshop gate, and the env barrier under it, in that order (#6306).
70 ///
71 /// A holder of this gate keeps it across whole product calls — `reload_config`
72 /// installs the budgets and then loads settings — and those calls read the
73 /// environment through `test_support::with_test_env_lock`. A test that sealed
74 /// the environment first and then reaches `install_active` through the very
75 /// same `reload_config` takes the two locks in the opposite order, and the
76 /// pair wedges the whole binary: `cargo test -p codewhale-tui --lib config`
77 /// never returned, with dozens of unrelated tests queued on the barrier.
78 /// Acquiring the barrier here, before the gate, makes that inversion
79 /// impossible for every caller at once — the same shape as
80 /// `provider_lake::lock_live_snapshot`.
81 ///
82 /// Known limitation: this orders these two locks and nothing else. A future
83 /// process-wide test gate held across product code has to join the same order
84 /// rather than invent a third one.
85 #[cfg(test)]
86 pub(crate) fn active_workshop_test_guard() -> ActiveWorkshopTestGuard {
87 assert!(
88 !ACTIVE_WORKSHOP_TEST_SERIAL_HELD.with(std::cell::Cell::get),
89 "active workshop test guard is not reentrant"
90 );
91 let env = if crate::test_support::current_thread_holds_test_env_lock() {
92 None
93 } else {
94 Some(crate::test_support::lock_test_env())
95 };
96 let serial = ACTIVE_WORKSHOP_TEST_SERIAL
97 .get_or_init(|| Mutex::new(()))
98 .lock()
99 .unwrap_or_else(std::sync::PoisonError::into_inner);
100 ACTIVE_WORKSHOP_TEST_SERIAL_HELD.with(|held| held.set(true));
101 ActiveWorkshopTestGuard {
102 _serial: serial,
103 _env: env,
104 }
105 }
106
107 fn active_workshop_slot() -> &'static Mutex<WorkshopConfig> {
108 ACTIVE_WORKSHOP.get_or_init(|| Mutex::new(WorkshopConfig::default()))
109 }
110
111 // ── Configuration ─────────────────────────────────────────────────────────────
112
113 /// Existing `[workshop]` threshold configuration, retained for compatibility.
114 #[derive(Debug, Clone, Deserialize, Default)]
115 pub struct WorkshopConfig {
116 /// Token threshold above which results become handle-only evidence.
117 #[serde(default)]
118 pub large_output_threshold_tokens: Option<usize>,
119
120 /// Per-tool threshold overrides (tool name → token limit). A tool whose
121 /// name appears here uses this limit instead of
122 /// `large_output_threshold_tokens`.
123 #[serde(default)]
124 pub per_tool_thresholds: Option<HashMap<String, usize>>,
125
126 /// Optional model-visible byte budget for a single `read` / `read_file`
127 /// result. Absent keeps the compile-time default (#5367).
128 #[serde(default)]
129 pub read_result_max_bytes: Option<usize>,
130
131 /// Optional model-visible byte budget for a generic tool result after
132 /// spillover. Absent keeps the compile-time default (#5367).
133 #[serde(default)]
134 pub tool_result_max_bytes: Option<usize>,
135 }
136
137 impl WorkshopConfig {
138 /// Install the process-wide workshop budgets used by read/tool compactors.
139 ///
140 /// The returned immutable receipt is the snapshot written while the
141 /// singleton lock was held. Callers that need evidence of their own
142 /// activation can inspect it without racing a later process-wide update.
143 pub fn install_active(config: Option<&Self>) -> Self {
144 #[cfg(test)]
145 let _test_serial = if ACTIVE_WORKSHOP_TEST_SERIAL_HELD.with(std::cell::Cell::get) {
146 None
147 } else {
148 Some(
149 ACTIVE_WORKSHOP_TEST_SERIAL
150 .get_or_init(|| Mutex::new(()))
151 .lock()
152 .unwrap_or_else(std::sync::PoisonError::into_inner),
153 )
154 };
155 let snapshot = config.cloned().unwrap_or_default();
156 let mut slot = active_workshop_slot()
157 .lock()
158 .unwrap_or_else(std::sync::PoisonError::into_inner);
159 *slot = snapshot;
160 slot.clone()
161 }
162
163 /// Optional model-visible read budget, when the user opted in (#5367).
164 #[must_use]
165 pub fn active_read_result_max_bytes() -> Option<usize> {
166 active_workshop_slot()
167 .lock()
168 .ok()
169 .and_then(|cfg| cfg.read_result_max_bytes.filter(|n| *n > 0))
170 }
171
172 /// Optional model-visible tool-result budget, when the user opted in (#5367).
173 #[must_use]
174 pub fn active_tool_result_max_bytes() -> Option<usize> {
175 active_workshop_slot()
176 .lock()
177 .ok()
178 .and_then(|cfg| cfg.tool_result_max_bytes.filter(|n| *n > 0))
179 }
180
181 /// Resolve the effective threshold for the given tool name.
182 #[must_use]
183 pub fn threshold_for(&self, tool_name: &str) -> usize {
184 if let Some(per_tool) = self.per_tool_thresholds.as_ref()
185 && let Some(&limit) = per_tool.get(tool_name)
186 {
187 return limit;
188 }
189 self.large_output_threshold_tokens
190 .unwrap_or(DEFAULT_LARGE_OUTPUT_THRESHOLD_TOKENS)
191 }
192 }
193
194 // ── Token estimation ──────────────────────────────────────────────────────────
195
196 /// Estimate the number of tokens in `text` using a character-count heuristic.
197 ///
198 /// This avoids a real tokeniser dependency; the estimate is deliberately
199 /// conservative (under-counts tokens) so we route aggressively rather than
200 /// letting a 5K-token blob slip through.
201 #[must_use]
202 pub fn estimate_tokens(text: &str) -> usize {
203 let chars = text.chars().count();
204 // Round up: partial last token still costs a token.
205 chars.div_ceil(CHARS_PER_TOKEN_ESTIMATE)
206 }
207
208 // ── Router ────────────────────────────────────────────────────────────────────
209
210 /// Classifies tool results for adaptive evidence routing.
211 ///
212 /// This type is intentionally `Clone` and `Default` so it can be embedded
213 /// cheaply in [`ToolContext`](crate::tools::spec::ToolContext) without
214 /// requiring `Arc` wrappers.
215 #[derive(Debug, Clone, Default)]
216 pub struct LargeOutputRouter {
217 config: WorkshopConfig,
218 }
219
220 impl LargeOutputRouter {
221 /// Construct a router from the resolved workshop config.
222 #[must_use]
223 pub fn new(config: WorkshopConfig) -> Self {
224 Self { config }
225 }
226
227 #[must_use]
228 pub fn evidence_routing(
229 &self,
230 tool_name: &str,
231 result: &ToolResult,
232 ) -> (EvidenceRouting, usize, usize) {
233 let threshold = self.config.threshold_for(tool_name);
234 let estimated_tokens = estimate_tokens(&result.content);
235 // There is no per-call bypass of the context bound: exact bytes stay
236 // available through the artifact handle.
237 let routing = EvidenceRouting::from_token_estimate(estimated_tokens, threshold);
238 (routing, estimated_tokens, threshold)
239 }
240 }
241
242 // ── Adaptive evidence routing (#4619) ─────────────────────────────────────────
243
244 /// Routing policy for tool results: how much of the output stays inline in the
245 /// conversation vs. being stored as an external artifact.
246 #[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
247 #[serde(rename_all = "snake_case")]
248 pub enum EvidenceRouting {
249 /// Full result stays inline in the conversation context.
250 Inline,
251 /// A bounded observation (head/tail/summary) stays inline; the exact bytes
252 /// are stored as an artifact recoverable via handle.
253 Hybrid,
254 /// Only a handle/reference stays inline; the full result is artifact-only.
255 HandleOnly,
256 }
257
258 impl EvidenceRouting {
259 /// Determine routing from estimated token count and threshold.
260 #[must_use]
261 pub fn from_token_estimate(estimated_tokens: usize, threshold: usize) -> Self {
262 if estimated_tokens <= threshold / 4 {
263 Self::Inline
264 } else if estimated_tokens <= threshold {
265 Self::Hybrid
266 } else {
267 Self::HandleOnly
268 }
269 }
270 }
271
272 /// Immutable metadata for a stored evidence artifact (#4619).
273 #[derive(Debug, Clone, Serialize, Deserialize)]
274 pub struct EvidenceArtifact {
275 pub handle: String,
276 pub digest: String,
277 pub size_bytes: u64,
278 pub content_type: String,
279 pub tool_name: String,
280 pub call_id: String,
281 pub origin_session: String,
282 pub generation: u32,
283 pub redacted: bool,
284 pub encoding: String,
285 pub retention_state: EvidenceRetentionState,
286 pub created_at_unix_ms: u64,
287 pub retain_until_unix_ms: u64,
288 pub storage_path: PathBuf,
289 }
290
291 #[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
292 #[serde(rename_all = "snake_case")]
293 pub enum EvidenceRetentionState {
294 Live,
295 Expired,
296 }
297
298 pub const EVIDENCE_RETENTION_SECS: u64 = 7 * 24 * 60 * 60;
299
300 /// Whether adaptive evidence routing (#4619) is enabled for this process.
301 ///
302 /// Off by default — classic bounded spillover owns large results. Set
303 /// `CODEWHALE_ADAPTIVE_OUTPUT_ROUTING` to opt in. The retired rollback
304 /// variable is still honored in the negative, so
305 /// `CODEWHALE_CLASSIC_OUTPUT_ROUTING=0` also selects the adaptive lane.
306 #[must_use]
307 pub fn adaptive_output_routing_enabled() -> bool {
308 if let Ok(value) = std::env::var("CODEWHALE_ADAPTIVE_OUTPUT_ROUTING") {
309 return matches!(value.trim(), "1" | "true" | "yes" | "on");
310 }
311 matches!(
312 std::env::var("CODEWHALE_CLASSIC_OUTPUT_ROUTING")
313 .ok()
314 .as_deref()
315 .map(str::trim),
316 Some("0" | "false" | "no" | "off")
317 )
318 }
319
320 #[must_use]
321 pub fn evidence_metadata_relative_path(handle: &str) -> PathBuf {
322 PathBuf::from(crate::artifacts::ARTIFACTS_DIR_NAME).join(format!("{handle}.evidence.json"))
323 }
324
325 pub fn publish_evidence_metadata(
326 session_id: &str,
327 artifact: &EvidenceArtifact,
328 ) -> io::Result<PathBuf> {
329 let bytes = serde_json::to_vec_pretty(artifact)
330 .map_err(|err| io::Error::new(io::ErrorKind::InvalidData, err))?;
331 crate::artifacts::write_session_relative_immutable(
332 session_id,
333 &evidence_metadata_relative_path(&artifact.handle),
334 &bytes,
335 )
336 }
337
338 pub fn read_evidence_metadata(session_id: &str, handle: &str) -> io::Result<EvidenceArtifact> {
339 let relative = evidence_metadata_relative_path(handle);
340 let file = crate::artifacts::open_session_relative(session_id, &relative, false)?;
341 read_evidence_metadata_file(&file)
342 }
343
344 /// Bounded, no-follow read shared by publication/replay and authenticated HTTP
345 /// retrieval. The caller chooses the existing session-root authority.
346 pub(crate) fn read_evidence_metadata_file(
347 file: &crate::fleet::files::WorkspaceFile,
348 ) -> io::Result<EvidenceArtifact> {
349 use std::io::Read;
350 const MAX_MANIFEST_BYTES: u64 = 64 * 1024;
351 let mut raw = Vec::new();
352 file.open_file()?
353 .take(MAX_MANIFEST_BYTES + 1)
354 .read_to_end(&mut raw)?;
355 if raw.len() as u64 > MAX_MANIFEST_BYTES {
356 return Err(io::Error::new(
357 io::ErrorKind::InvalidData,
358 "evidence metadata exceeds limit",
359 ));
360 }
361 serde_json::from_slice(&raw).map_err(|err| io::Error::new(io::ErrorKind::InvalidData, err))
362 }
363
364 #[must_use]
365 pub fn unix_millis_now() -> u64 {
366 std::time::SystemTime::now()
367 .duration_since(std::time::UNIX_EPOCH)
368 .unwrap_or_default()
369 .as_millis()
370 .try_into()
371 .unwrap_or(u64::MAX)
372 }
373
374 #[must_use]
375 pub fn evidence_is_expired(artifact: &EvidenceArtifact, now_ms: u64) -> bool {
376 artifact.retention_state == EvidenceRetentionState::Expired
377 || now_ms > artifact.retain_until_unix_ms
378 }
379
380 // ── Unit tests ────────────────────────────────────────────────────────────────
381
382 #[cfg(test)]
383 mod tests {
384 use super::*;
385
386 fn make_result(content: &str) -> ToolResult {
387 ToolResult::success(content.to_string())
388 }
389
390 #[test]
391 fn default_threshold_is_32k_tokens() {
392 assert_eq!(DEFAULT_LARGE_OUTPUT_THRESHOLD_TOKENS, 32_768);
393 }
394
395 #[test]
396 fn oversized_result_becomes_handle_only_evidence() {
397 let router = LargeOutputRouter::default();
398 let big = make_result(&"a".repeat(100_000));
399 let (routing, _, _) = router.evidence_routing("bash", &big);
400 assert_eq!(routing, EvidenceRouting::HandleOnly);
401 }
402
403 #[test]
404 fn per_tool_threshold_override() {
405 let mut per_tool = HashMap::new();
406 per_tool.insert("grep_files".to_string(), 100); // very low
407 let config = WorkshopConfig {
408 large_output_threshold_tokens: Some(4096),
409 per_tool_thresholds: Some(per_tool),
410 read_result_max_bytes: None,
411 tool_result_max_bytes: None,
412 };
413 assert_eq!(config.threshold_for("grep_files"), 100);
414 assert_eq!(config.threshold_for("read_file"), 4096);
415 let default_config = WorkshopConfig::default();
416 assert_eq!(
417 default_config.threshold_for("read_file"),
418 DEFAULT_LARGE_OUTPUT_THRESHOLD_TOKENS
419 );
420 }
421
422 #[test]
423 fn workshop_byte_budgets_raise_floor_only() {
424 let _guard = active_workshop_test_guard();
425 let installed = WorkshopConfig::install_active(Some(&WorkshopConfig {
426 large_output_threshold_tokens: None,
427 per_tool_thresholds: None,
428 read_result_max_bytes: Some(102_400),
429 tool_result_max_bytes: Some(80_000),
430 }));
431 assert_eq!(installed.read_result_max_bytes, Some(102_400));
432 assert_eq!(installed.tool_result_max_bytes, Some(80_000));
433 assert_eq!(
434 WorkshopConfig::active_read_result_max_bytes(),
435 Some(102_400)
436 );
437 assert_eq!(WorkshopConfig::active_tool_result_max_bytes(), Some(80_000));
438 let cleared = WorkshopConfig::install_active(None);
439 assert_eq!(cleared.read_result_max_bytes, None);
440 assert_eq!(cleared.tool_result_max_bytes, None);
441 assert_eq!(WorkshopConfig::active_read_result_max_bytes(), None);
442 assert_eq!(WorkshopConfig::active_tool_result_max_bytes(), None);
443 }
444
445 #[test]
446 fn estimate_tokens_conservative() {
447 // 9 chars → ceil(9/3) = 3 tokens
448 assert_eq!(estimate_tokens("123456789"), 3);
449 // 10 chars → ceil(10/3) = 4 tokens
450 assert_eq!(estimate_tokens("1234567890"), 4);
451 // Empty string
452 assert_eq!(estimate_tokens(""), 0);
453 }
454 }
455
455 lines RUST