diff --git a/.repository-projection.json b/.repository-projection.json index 2f896049d..c9ed14b4e 100644 --- a/.repository-projection.json +++ b/.repository-projection.json @@ -3,11 +3,11 @@ "projection": "deixic-code", "projectionSchemaVersion": 1, "sourceRepository": "dx-corp/mono", - "sourceSha": "360c907931868953603eb6cba0a5dd0a7ad0f377", + "sourceSha": "f0184c6753b3ec240c0e7da694a16886cbecfae7", "destinationRepository": "dx-corp/code", - "priorProjectedBase": "1b0b37e6966c7f42e87a5ce2f0d9ff9c930b2704", + "priorProjectedBase": "1761784b70fe4375102cfe54aa6029e946699ffc", "definitionDigest": "cb9d429542ebb0a2de9b42a7aad60d9d8696a648ceba47c30f05c0b285ca0db7", "toolDigest": "c9112982c70b80737d8faac67d58b06d55bde0dc297eb1009c1b02296f34f4a9", - "contentDigest": "10fd59ed421dfe0ff91080c6280dd89dfca5cd14a8f76f6dc48f518cd3f06727", + "contentDigest": "cc056cd6ac865d7b0d8e7eae1032e70006f3788e20adc2871696aa27010e94e0", "publicationEligible": true } diff --git a/vendor/dex-loop/src/compaction.rs b/vendor/dex-loop/src/compaction.rs index 002424fd6..9f279ee18 100644 --- a/vendor/dex-loop/src/compaction.rs +++ b/vendor/dex-loop/src/compaction.rs @@ -17,6 +17,12 @@ pub struct Compaction { /// Decides, before each model call, whether to compact. pub trait Compactor: Send + Sync { fn plan(&self, ctx: &Context) -> impl Future + Send; + + /// A complete request exceeded capacity despite the proactive threshold. + /// Hosts without a recovery strategy leave the original failure intact. + fn recover(&self, _ctx: &Context) -> impl Future + Send { + async { CompactionPlan::default() } + } } /// An attempted summary may consume model usage even when no summary commits. @@ -149,6 +155,38 @@ impl Compactor for Threshold { async fn plan(&self, ctx: &Context) -> CompactionPlan { self.plan_at(ctx, self.max_bytes).await } + + async fn recover(&self, ctx: &Context) -> CompactionPlan { + if self.max_bytes == usize::MAX { + return CompactionPlan::default(); + } + // Request framing, tools and resolved references are not all counted + // by the history threshold. Use the same safe cuts even below it. + let mut plan = self.plan_at(ctx, 0).await; + if let Some(compaction) = &plan.compaction { + let before: usize = ctx.history().iter().map(|entry| entry.message.size()).sum(); + let mut projected = ctx.clone(); + projected.observe( + compaction.covers_to, + &crate::Event::Compaction { + covers_to_cursor: compaction.covers_to, + summary: compaction.summary.clone(), + }, + ); + let after: usize = projected + .history() + .iter() + .map(|entry| entry.message.size()) + .sum(); + if after >= before { + // A successful summarizer is not necessarily a reduction. + // Retain billed usage, but never retry an unchanged payload. + plan.compaction = None; + plan.declined = true; + } + } + plan + } } /// Where to cut history for a compaction, chosen by [`plan_cut`]. diff --git a/vendor/dex-loop/src/engine.rs b/vendor/dex-loop/src/engine.rs index 84875b62e..87006c27f 100644 --- a/vendor/dex-loop/src/engine.rs +++ b/vendor/dex-loop/src/engine.rs @@ -42,6 +42,14 @@ use crate::ports::{ }; use crate::sanitize::{DeltaFilter, Sanitizer}; +/// One capacity recovery per run; restart remains bounded by the durable +/// step and usage budgets. No failed answer or tool call is replayed. +enum CapacityRecovery { + Available, + Pending(ModelError), + Attempted, +} + /// The engine-owned discovery tool. Always offered to the model. pub const TOOLS_SEARCH: &str = "tools.search"; /// The engine-owned bounded script tool. @@ -413,6 +421,7 @@ where ) -> Result { let started = Instant::now(); let mut prefetch = Prefetch::new(); + let mut capacity_recovery = CapacityRecovery::Available; loop { self.read_control(ctx).await?; match ctx.status() { @@ -465,13 +474,29 @@ where return Ok(Exit::Failed); } let remaining_wall = self.budget.wall.saturating_sub(started.elapsed()); + let recovery_error = if matches!(capacity_recovery, CapacityRecovery::Pending(_)) { + match std::mem::replace(&mut capacity_recovery, CapacityRecovery::Attempted) { + CapacityRecovery::Pending(error) => Some(error), + _ => unreachable!(), + } + } else { + None + }; + let planning = async { + if recovery_error.is_some() { + self.compactor.recover(ctx).await + } else { + self.compactor.plan(ctx).await + } + }; let plan = tokio::select! { _ = cancel.cancelled() => return self.interrupt(ctx, &prefetch).await, - result = tokio::time::timeout(remaining_wall, self.compactor.plan(ctx)) => result, + result = tokio::time::timeout(remaining_wall, planning) => result, }; let mut compaction_declined = false; if let Ok(plan) = plan { - compaction_declined = plan.declined; + compaction_declined = + plan.declined || (recovery_error.is_some() && plan.compaction.is_none()); let mut events = Vec::new(); if plan.usage != Default::default() { events.push(Event::Usage(plan.usage)); @@ -511,7 +536,10 @@ where vec![Event::Error { class: Some(crate::ErrorClass::ContextCapacity), code: ErrorCode::ModelFailed, - message: "Conversation context could not be reduced safely. History is preserved; start a new conversation with the relevant details if recovery keeps failing.".into(), + message: recovery_error.map_or_else( + || "Conversation context could not be reduced safely. History is preserved; start a new conversation with the relevant details if recovery keeps failing.".into(), + |error| error.message, + ), }], ) .await?; @@ -519,7 +547,10 @@ where } // A step never inherits another step's reads. prefetch = Prefetch::new(); - if let Some(exit) = self.model_step(ctx, cancel, started, &mut prefetch).await? { + if let Some(exit) = self + .model_step(ctx, cancel, started, &mut prefetch, &mut capacity_recovery) + .await? + { return Ok(exit); } } @@ -543,6 +574,7 @@ where cancel: &'e CancellationToken, started: Instant, prefetch: &mut Prefetch<'e>, + capacity_recovery: &mut CapacityRecovery, ) -> Result, Fenced> { let step = ctx.step().saturating_add(1); self.emit( @@ -844,6 +876,20 @@ where } let mut events = pending_usage; events.push(Event::ModelAttemptAbandoned { step }); + if error.class() == crate::ErrorClass::ContextCapacity + && matches!(capacity_recovery, CapacityRecovery::Available) + && text.is_empty() + && calls.is_empty() + && reasoning.is_none() + && !cancel.is_cancelled() + { + // Close this attempt before planning a safe cut. The next + // loop checks control and budgets, persists a reducing summary, + // then rebuilds the complete request through normal admission. + self.emit(ctx, events).await?; + *capacity_recovery = CapacityRecovery::Pending(error); + return Ok(None); + } events.push(Event::Error { class: Some(error.class()), code: ErrorCode::ModelFailed, diff --git a/vendor/dex-loop/src/event.rs b/vendor/dex-loop/src/event.rs index b7af9395d..e170dad92 100644 --- a/vendor/dex-loop/src/event.rs +++ b/vendor/dex-loop/src/event.rs @@ -462,6 +462,7 @@ impl ErrorClass { /// Explicit wire-code boundary; human-readable details are never classified. pub fn of_gateway_code(code: &str) -> Self { match code { + "context_length_exceeded" | "context_capacity" => Self::ContextCapacity, "rate_limit_error" | "rate_limit_exceeded" | "resource_exhausted" diff --git a/vendor/dex-loop/src/summary.rs b/vendor/dex-loop/src/summary.rs index 9902d2f8f..c22a71967 100644 --- a/vendor/dex-loop/src/summary.rs +++ b/vendor/dex-loop/src/summary.rs @@ -175,45 +175,71 @@ impl Summarize for ModelSummarizer { ); // No tools or executable capability are offered to this call. Reusing // Model retains tenant authentication, provider routing, and metering. - let stream = self.model.stream(&input, &[]); - futures_util::pin_mut!(stream); let mut result = Summary::default(); let mut summary = String::new(); let mut valid = true; let mut failure = None; let deadline = tokio::time::sleep(self.timeout); tokio::pin!(deadline); - loop { - let chunk = tokio::select! { - chunk = stream.next() => match chunk { Some(chunk) => chunk, None => break }, - () = &mut deadline => { valid = false; failure = Some("timeout"); break; }, - }; - match chunk { - Ok(ModelChunk::Usage(usage)) => result.usage += usage, - Ok(ModelChunk::Text(delta)) - if summary.len().saturating_add(delta.len()) <= SUMMARY_BYTES => - { - summary.push_str(&delta) - } - Ok(ModelChunk::Text(_)) => { - valid = false; - failure = Some("oversized_summary"); - } - Ok(ModelChunk::ToolCall { .. }) => { - valid = false; - failure = Some("unexpected_tool_call"); - } - Err(_) => { - valid = false; - failure = Some("model_error"); + for attempt in 0..2 { + // Failed partial summaries are private and disposable; usage is + // accumulated across attempts. Both share the original deadline. + summary.clear(); + valid = true; + let mut retry = false; + let stream = self.model.stream(&input, &[]); + futures_util::pin_mut!(stream); + loop { + let chunk = tokio::select! { + biased; + () = &mut deadline => { valid = false; retry = false; failure = Some("timeout"); break; }, + chunk = stream.next() => match chunk { Some(chunk) => chunk, None => break }, + }; + match chunk { + Ok(ModelChunk::Usage(usage)) => result.usage += usage, + Ok(ModelChunk::Text(delta)) + if summary.len().saturating_add(delta.len()) <= SUMMARY_BYTES => + { + summary.push_str(&delta) + } + Ok(ModelChunk::Text(_)) => { + valid = false; + failure = Some("oversized_summary"); + } + Ok(ModelChunk::ToolCall { .. }) => { + valid = false; + failure = Some("unexpected_tool_call"); + } + Err(error) => { + // Only classes whose transport contract is always + // retryable. Provider rejections carry retry advice + // below this port; do not override it from prose/class. + retry = valid + && matches!( + error.class(), + crate::ErrorClass::Transport | crate::ErrorClass::Truncated + ); + valid = false; + failure = Some("model_error"); + // Drain terminal metering even after an error. No + // partial output from this attempt can be installed. + } + Ok( + ModelChunk::Reasoning(_) + | ModelChunk::Thinking(_) + | ModelChunk::Served(_) + | ModelChunk::Timing(_) + | ModelChunk::AttemptFailed { .. }, + ) => {} } - Ok( - ModelChunk::Reasoning(_) - | ModelChunk::Thinking(_) - | ModelChunk::Served(_) - | ModelChunk::Timing(_) - | ModelChunk::AttemptFailed { .. }, - ) => {} + } + if !retry || attempt == 1 { + break; + } + tokio::select! { + biased; + () = &mut deadline => { failure = Some("timeout"); break; }, + () = tokio::time::sleep(Duration::from_millis(250)) => {}, } } let mut tier = "summarize"; @@ -1113,6 +1139,77 @@ mod tests { assert!(text.len() <= SUMMARY_BYTES); assert_eq!(summary.usage.cost_micros, 5); } + + #[tokio::test(start_paused = true)] + async fn summary_retries_only_transport_failures_with_one_deadline_and_exact_usage() { + struct Attempts { + calls: Mutex, + class: crate::ErrorClass, + succeed: bool, + } + impl Model for Attempts { + fn stream<'a>( + &'a self, + _: &'a Context, + tools: &'a [&'a ToolSpec], + ) -> impl Stream> + Send + 'a { + assert!(tools.is_empty()); + let mut calls = self.calls.lock().unwrap(); + *calls += 1; + let success = *calls == 2 && self.succeed; + let mut chunks = vec![ + Ok(ModelChunk::Text( + if success { + "complete summary" + } else { + "discard this partial" + } + .into(), + )), + Ok(ModelChunk::Usage(Usage { + cost_micros: if success { 7 } else { 3 }, + ..Default::default() + })), + ]; + if !success { + chunks.push(Err(ModelError { + class: self.class, + message: "transport: temporary failure".into(), + })); + } + stream::iter(chunks) + } + } + for (class, timeout_ms, succeed, calls, cost, recovered) in [ + (crate::ErrorClass::Transport, 1000, true, 2, 10, true), + (crate::ErrorClass::Truncated, 1000, true, 2, 10, true), + (crate::ErrorClass::Transport, 1000, false, 2, 6, false), + (crate::ErrorClass::Transport, 100, true, 1, 3, false), + (crate::ErrorClass::Auth, 1000, true, 1, 3, false), + (crate::ErrorClass::ContextCapacity, 1000, true, 1, 3, false), + (crate::ErrorClass::Unavailable, 1000, true, 1, 3, false), + ] { + let summarizer = ModelSummarizer::new(Attempts { + calls: Mutex::new(0), + class, + succeed, + }) + .with_timeout(Duration::from_millis(timeout_ms)); + let ctx = context(); + let original = ctx.clone(); + let result = summarizer.summarize(&ctx, &inference_history()).await; + assert_eq!(*summarizer.model.calls.lock().unwrap(), calls, "{class:?}"); + assert_eq!(result.usage.cost_micros, cost, "{class:?}"); + let text = result + .text + .expect("successful summary or faithful fallback"); + assert!(!text.contains("discard this partial")); + assert_eq!(text.contains("complete summary"), recovered); + assert!(text.contains("attachment@v1") && text.contains("result@v1")); + assert_eq!(ctx, original); + } + } + #[tokio::test] async fn empty_summary_falls_back_to_exact_constraints_and_references_with_usage() { let usage = Usage { diff --git a/vendor/dex-loop/tests/scenarios.rs b/vendor/dex-loop/tests/scenarios.rs index b736b8af2..835d10cca 100644 --- a/vendor/dex-loop/tests/scenarios.rs +++ b/vendor/dex-loop/tests/scenarios.rs @@ -2542,6 +2542,228 @@ async fn declined_compaction_preserves_history_and_a_later_turn_can_recover() { assert_eq!(resumed, log.rehydrate()); } +#[tokio::test] +async fn capacity_recovery_reduces_history_once_before_retrying_the_same_turn() { + for second_fails in [false, true] { + let log = FakeLog::default(); + log.start_turn("old", "Never publish without approval. Budget is $10."); + log.host_append(Event::ModelStepCompleted { + step: 1, + text: "historical analysis ".repeat(100), + calls: vec![], + reasoning: None, + served: None, + timing: None, + }); + log.host_append(Event::Final { + text: "unfinished".into(), + }); + let mut ctx = log.start_turn("current", "continue investigating, do not publish"); + let capacity = || { + Err(ModelError { + class: dex_loop::ErrorClass::ContextCapacity, + message: "request too large".into(), + }) + }; + let model = FakeModel::new(vec![ + vec![usage(3, 0, 5), capacity()], + if second_fails { + vec![capacity()] + } else { + vec![text("recovered answer"), usage(7, 2, 11)] + }, + ]); + let tools = FakeTools::new(vec![]); + // History is below the proactive threshold: the complete request's + // capacity error, not its text, must trigger recovery. + let engine = support::engine(&log, &model, &tools, budget()) + .with_compactor(Threshold::for_turns(48 * 1024, FakeSummarizer)); + assert_eq!( + engine.run(&mut ctx, &CancellationToken::new()).await, + Ok(if second_fails { + Exit::Failed + } else { + Exit::Done + }) + ); + assert_eq!(model.calls(), 2); + assert_eq!(ctx.usage().cost_micros, if second_fails { 5 } else { 16 }); + assert_eq!( + log.events() + .iter() + .filter(|e| matches!(e, Event::Compaction { .. })) + .count(), + 1 + ); + assert_eq!( + log.events() + .iter() + .filter(|e| matches!(e, Event::Error { .. })) + .count(), + usize::from(second_fails) + ); + assert!(matches!( + &model.seen()[1][0], + dex_loop::Message::Summary { .. } + )); + assert!(model.seen()[1].iter().any(|m| matches!(m, + dex_loop::Message::User { text, .. } if text == "continue investigating, do not publish"))); + assert_eq!(ctx, log.rehydrate()); + assert!(tools.runs().is_empty()); + } +} + +#[tokio::test] +async fn capacity_recovery_does_not_retry_nonreducing_or_uncompactable_history() { + struct Expanding; + impl dex_loop::Summarize for Expanding { + async fn summarize( + &self, + _: &dex_loop::Context, + _: &[dex_loop::Entry], + ) -> dex_loop::Summary { + dex_loop::Summary { + text: Some("x".repeat(10_000)), + usage: dex_loop::Usage { + cost_micros: 13, + ..Default::default() + }, + } + } + } + for history in [false, true] { + let log = FakeLog::default(); + if history { + log.start_turn("old", "exact constraint"); + log.host_append(Event::ModelStepCompleted { + step: 1, + text: "done".into(), + calls: vec![], + reasoning: None, + served: None, + timing: None, + }); + log.host_append(Event::Final { + text: "done".into(), + }); + } + let mut ctx = log.start_turn("current", "keep this exact"); + let original = ctx.history().to_vec(); + let model = FakeModel::new(vec![vec![Err(ModelError { + class: dex_loop::ErrorClass::ContextCapacity, + message: "request too large".into(), + })]]); + let tools = FakeTools::new(vec![]); + let engine = support::engine(&log, &model, &tools, budget()) + .with_compactor(Threshold::for_turns(48 * 1024, Expanding)); + assert_eq!( + engine.run(&mut ctx, &CancellationToken::new()).await, + Ok(Exit::Failed) + ); + assert_eq!(model.calls(), 1); + assert_eq!(ctx.history(), original); + assert_eq!(ctx.usage().cost_micros, if history { 13 } else { 0 }); + assert!( + !log.events() + .iter() + .any(|e| matches!(e, Event::Compaction { .. })) + ); + assert_eq!(ctx, log.rehydrate()); + } +} + +#[tokio::test] +async fn capacity_recovery_respects_failure_class_visible_output_and_budget() { + for (class, partial, max_cost, exit, error_class) in [ + ( + dex_loop::ErrorClass::Unknown, + false, + 100, + Exit::Failed, + Some(dex_loop::ErrorClass::Unknown), + ), + ( + dex_loop::ErrorClass::ContextCapacity, + true, + 100, + Exit::Done, + None, + ), + ( + dex_loop::ErrorClass::ContextCapacity, + false, + 5, + Exit::Failed, + None, + ), + ] { + let log = FakeLog::default(); + log.start_turn("old", "preserve approval constraint"); + log.host_append(Event::ModelStepCompleted { + step: 1, + text: "old analysis".repeat(100), + calls: vec![], + reasoning: None, + served: None, + timing: None, + }); + log.host_append(Event::Final { + text: "unfinished".into(), + }); + let mut ctx = log.start_turn("current", "continue"); + let mut chunks = vec![usage(3, 1, 5)]; + if partial { + chunks.push(text("already visible answer")); + } + chunks.push(Err(ModelError { + class, + message: "context_length_exceeded".into(), + })); + let model = FakeModel::new(vec![chunks]); + let tools = FakeTools::new(vec![]); + let engine = support::engine( + &log, + &model, + &tools, + Budget { + max_cost_micros: max_cost, + ..budget() + }, + ) + .with_compactor(Threshold::for_turns(48 * 1024, FakeSummarizer)); + assert_eq!( + engine.run(&mut ctx, &CancellationToken::new()).await, + Ok(exit) + ); + assert_eq!(model.calls(), 1); + assert!( + !log.events() + .iter() + .any(|e| matches!(e, Event::Compaction { .. })) + ); + assert_eq!( + log.events().iter().find_map(|e| match e { + Event::Error { class, .. } => *class, + _ => None, + }), + error_class + ); + if max_cost == 5 { + assert!(log.events().iter().any(|e| matches!( + e, + Event::Error { + code: dex_loop::ErrorCode::BudgetExhausted, + .. + } + ))); + } + if partial { + assert!(log.events().iter().any(|e| matches!(e, Event::Final { text } if text.starts_with("already visible answer") && text.ends_with(CUT_OFF_NOTICE)))); + } + assert_eq!(ctx, log.rehydrate()); + } +} + #[tokio::test] async fn summary_calls_obey_the_same_wall_budget_as_normal_model_calls() { struct Slow;