Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
6 changes: 3 additions & 3 deletions .repository-projection.json
Original file line number Diff line number Diff line change
Expand Up @@ -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
}
38 changes: 38 additions & 0 deletions vendor/dex-loop/src/compaction.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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<Output = CompactionPlan> + 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<Output = CompactionPlan> + Send {
async { CompactionPlan::default() }
}
}

/// An attempted summary may consume model usage even when no summary commits.
Expand Down Expand Up @@ -149,6 +155,38 @@ impl<S: Summarize> Compactor for Threshold<S> {
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`].
Expand Down
54 changes: 50 additions & 4 deletions vendor/dex-loop/src/engine.rs
Original file line number Diff line number Diff line change
Expand Up @@ -42,6 +42,14 @@
};
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.
Expand Down Expand Up @@ -413,6 +421,7 @@
) -> Result<Exit, Fenced> {
let started = Instant::now();
let mut prefetch = Prefetch::new();
let mut capacity_recovery = CapacityRecovery::Available;
loop {
self.read_control(ctx).await?;
match ctx.status() {
Expand Down Expand Up @@ -465,13 +474,29 @@
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));
Expand Down Expand Up @@ -511,15 +536,21 @@
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?;
return Ok(Exit::Failed);
}
// 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);
}
}
Expand All @@ -543,6 +574,7 @@
cancel: &'e CancellationToken,
started: Instant,
prefetch: &mut Prefetch<'e>,
capacity_recovery: &mut CapacityRecovery,
) -> Result<Option<Exit>, Fenced> {
let step = ctx.step().saturating_add(1);
self.emit(
Expand Down Expand Up @@ -844,6 +876,20 @@
}
let mut events = pending_usage;
events.push(Event::ModelAttemptAbandoned { step });
if error.class() == crate::ErrorClass::ContextCapacity
&& matches!(capacity_recovery, CapacityRecovery::Available)

Check failure on line 880 in vendor/dex-loop/src/engine.rs

View workflow job for this annotation

GitHub Actions / unresolved-review-threads / unresolved-review-threads

unresolved P1 review thread

Check failure on line 880 in vendor/dex-loop/src/engine.rs

View workflow job for this annotation

GitHub Actions / unresolved-review-threads / unresolved-review-threads

unresolved P1 review thread

Comment on lines +879 to +880

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P1 Badge Propagate context-capacity errors from the production adapter

This recovery branch is unreachable for actual context-window failures from the production AiRsModel: packages/ai-rs/src/openai.rs reports both HTTP and streamed context-overflow failures as ProviderStreamErrorKind::ProviderDeclaredFailure, while packages/dex-host-rs/src/model.rs converts that kind to ErrorClass::Unknown. The newly extended ErrorClass::of_gateway_code has no production callers (repo-wide search finds only its unit tests), so only the synthetic FakeModel scenarios supply ContextCapacity; real oversized requests still skip compaction and fail immediately. Propagate a typed context-capacity kind/code through the adapter before relying on this equality check.

Useful? React with 👍 / 👎.

&& 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,
Expand Down
1 change: 1 addition & 0 deletions vendor/dex-loop/src/event.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand Down
161 changes: 129 additions & 32 deletions vendor/dex-loop/src/summary.rs
Original file line number Diff line number Diff line change
Expand Up @@ -175,45 +175,71 @@ impl<M: Model> Summarize for ModelSummarizer<M> {
);
// 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";
Expand Down Expand Up @@ -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<usize>,
class: crate::ErrorClass,
succeed: bool,
}
impl Model for Attempts {
fn stream<'a>(
&'a self,
_: &'a Context,
tools: &'a [&'a ToolSpec],
) -> impl Stream<Item = Result<ModelChunk, ModelError>> + 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 {
Expand Down
Loading
Loading