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": "c30dc2f64dce3300144e0e0d0dc97eb97b8744ee",
"sourceSha": "ac693f34cdb782d83afd46151bf3a25a2df324be",
"destinationRepository": "dx-corp/code",
"priorProjectedBase": "358d47f82eb9a64f88ecfa2865002aa704e0baa5",
"priorProjectedBase": "1f00b9e589e080257b8366fc0884453e2596ce90",
"definitionDigest": "82936441c776e3e8edb5d215a75007ec9714a233f489d460075d79d5ef5ba32f",
"toolDigest": "c244d99199a7ae3eb8ff644a99462163c23b0bb6a83ef50af01efbdca0b81d04",
"contentDigest": "fc3edcce20cf8d94ea1cd95fce135752740dc77c7bd79715e242599e385cb806",
"contentDigest": "42d4d938700f91af78522c50de4f281b3d06804428a6744df1f37d0cd3ddc839",
"publicationEligible": true
}
5 changes: 4 additions & 1 deletion vendor/dex-loop/src/context.rs
Original file line number Diff line number Diff line change
Expand Up @@ -467,7 +467,10 @@ impl Context {
self.flush_steers(cursor, *control_through);
}
// Text reaches history through `ModelStepCompleted`.
Event::TextDelta { .. } | Event::ThinkingDelta { .. } | Event::ToolProgress { .. } => {}
Event::TextDelta { .. }
| Event::ThinkingDelta { .. }
| Event::ModelAttemptFailed { .. }
| Event::ToolProgress { .. } => {}
Event::Usage(usage) => self.usage += *usage,
Event::ModelStepCompleted {
text,
Expand Down
36 changes: 30 additions & 6 deletions vendor/dex-loop/src/engine.rs
Original file line number Diff line number Diff line change
Expand Up @@ -43,6 +43,7 @@ use crate::sanitize::{DeltaFilter, Sanitizer};
pub const TOOLS_SEARCH: &str = "tools.search";

const NOT_RUN_INTERRUPTED: &str = "not run: the turn was interrupted";
const READ_INTERRUPTED: &str = "not completed: the read was interrupted; it is safe to try again";
const UNKNOWN_INTERRUPTED: &str =
"outcome unknown: the turn was interrupted before the result was recorded";
/// A call that already started (its tool vanished from the catalog, or the
Expand Down Expand Up @@ -384,7 +385,7 @@ where
self.emit(ctx, events).await?;
}
if ctx.interrupt_requested() || cancel.is_cancelled() {
return self.interrupt(ctx).await;
return self.interrupt(ctx, &prefetch).await;
}
if ctx.open_step().is_some() {
let reads = prefetch.reads.clone();
Expand Down Expand Up @@ -564,6 +565,23 @@ where
StreamStep::Chunk(Ok(ModelChunk::Timing(t))) => {
timing = Some(t);
}
StreamStep::Chunk(Ok(ModelChunk::AttemptFailed {
provider,
model,
code,
elapsed_ms,
then,
})) => {
let event = [Event::ModelAttemptFailed {
step,
provider,
model,
code,
elapsed_ms,
then,
}];
prefetch.reads.drive(self.log.append(&event)).await?;
}
StreamStep::Chunk(Err(error)) => {
failure = Some(error.message);
break;
Expand Down Expand Up @@ -657,7 +675,7 @@ where
timing: timing.take(),
});
self.emit(ctx, events).await?;
return self.interrupt(ctx).await.map(Some);
return self.interrupt(ctx, prefetch).await.map(Some);
}
if !calls.is_empty() && self.budget.answer_only(step.saturating_sub(1)) {
// Asked for a tool on the call that offered none. Nothing can run
Expand Down Expand Up @@ -1039,7 +1057,7 @@ where
self.run_wave(ctx, &calls, wave, cancel, run_started, prefetch)
.await?;
if cancel.is_cancelled() {
return self.interrupt(ctx).await.map(Some);
return self.interrupt(ctx, prefetch).await.map(Some);
}
Ok(None)
}
Expand Down Expand Up @@ -1352,15 +1370,21 @@ where
}

/// Every call without a result gets one, then `Interrupted`.
async fn interrupt(&self, ctx: &mut Context) -> Result<Exit, Fenced> {
async fn interrupt(&self, ctx: &mut Context, prefetch: &Prefetch<'_>) -> Result<Exit, Fenced> {
let mut events: Vec<Event> = ctx
.open_step()
.map(|step| {
step.calls
.iter()
.zip(&step.states)
.filter_map(|(call, state)| match state {
.enumerate()
.filter_map(|(index, (call, state))| match state {
CallState::Done(_) => None,
// Only this attempt's admitted local reads enter prefetch.
// Do not reinterpret an older started call using a changed catalog.
CallState::Started if prefetch.started.contains(&index) => {
Some(finished(call, ToolResult::error(READ_INTERRUPTED)))
}
CallState::Started => {
Some(finished(call, ToolResult::unknown(UNKNOWN_INTERRUPTED)))
}
Expand All @@ -1370,7 +1394,7 @@ where
})
.unwrap_or_default();
events.push(Event::Interrupted);
self.emit(ctx, events).await?;
prefetch.reads.drive(self.emit(ctx, events)).await?;
Ok(Exit::Interrupted)
}

Expand Down
50 changes: 50 additions & 0 deletions vendor/dex-loop/src/event.rs
Original file line number Diff line number Diff line change
Expand Up @@ -149,6 +149,18 @@ pub struct ServedBy {
pub model: String,
}

/// What the model port did after one attempt failed.
#[derive(Clone, Copy, Debug, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum AttemptNext {
/// Tried the same route again.
Retry,
/// Moved to the next route.
Failover,
/// Gave up; the step fails.
Abandon,
}

/// Model spend reported by one model response.
#[derive(Clone, Copy, Debug, Default, PartialEq, Eq, Serialize, Deserialize)]
pub struct Usage {
Expand Down Expand Up @@ -510,6 +522,18 @@ pub enum Event {
ModelAttemptAbandoned {
step: u32,
},
/// A model attempt on one route ended without a `ModelStepCompleted`.
/// Debug data for staff: never model history, never rendered. Carries
/// only the route's configured provider/model, dex-model's fixed error
/// code and timing; never the error message.
ModelAttemptFailed {
step: u32,
provider: String,
model: String,
code: String,
elapsed_ms: u64,
then: AttemptNext,
},

ToolStarted {
call: CallId,
Expand Down Expand Up @@ -684,6 +708,32 @@ impl Event {
mod tests {
use super::*;

#[test]
fn model_attempt_failed_round_trips_as_snake_case_json() {
let event = Event::ModelAttemptFailed {
step: 2,
provider: "vertex-ai".into(),
model: "gemini-3.6-flash".into(),
code: "provider_unavailable".into(),
elapsed_ms: 1500,
then: AttemptNext::Failover,
};
let json = serde_json::to_value(&event).expect("serialize");
assert_eq!(json["type"], "model_attempt_failed");
assert_eq!(json["then"], "failover");
assert_eq!(
serde_json::from_value::<Event>(json).expect("decode"),
event
);
assert!(!event.is_control());
for (then, name) in [
(AttemptNext::Retry, "retry"),
(AttemptNext::Abandon, "abandon"),
] {
assert_eq!(serde_json::to_value(then).expect("serialize"), name);
}
}

#[test]
fn auto_approval_receipts_decode_losslessly_and_never_become_control_events() {
let event = Event::AutoApproved {
Expand Down
8 changes: 4 additions & 4 deletions vendor/dex-loop/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -33,10 +33,10 @@ pub use compaction::{Compaction, Compactor, Cut, NoCompaction, Summarize, Thresh
pub use context::{Context, Entry, Message};
pub use engine::{CUT_OFF_NOTICE, DEFAULT_TOOL_CALL_DEADLINE, Engine, Exit, TOOLS_SEARCH};
pub use event::{
AUTO_APPROVER, ApprovalId, ApprovalMode, ArtifactRef, CallId, ClientToolSpec, Cursor,
ErrorCode, Event, HEADLESS_AUTO_APPROVER, MessageId, Outcome, Output, OutputRef, PrincipalId,
ProposedCall, ProviderReasoning, ReceiptId, ServedBy, StepTiming, ThreadId, ToolName,
ToolResult, TurnId, Usage, args_digest,
AUTO_APPROVER, ApprovalId, ApprovalMode, ArtifactRef, AttemptNext, CallId, ClientToolSpec,
Cursor, ErrorCode, Event, HEADLESS_AUTO_APPROVER, MessageId, Outcome, Output, OutputRef,
PrincipalId, ProposedCall, ProviderReasoning, ReceiptId, ServedBy, StepTiming, ThreadId,
ToolName, ToolResult, TurnId, Usage, args_digest,
};
pub use ports::{
Claim, Effects, ExecutorKind, Fenced, GovernanceClass, Log, Model, ModelChunk, ModelError,
Expand Down
16 changes: 14 additions & 2 deletions vendor/dex-loop/src/ports.rs
Original file line number Diff line number Diff line change
Expand Up @@ -9,8 +9,8 @@ use tokio_util::sync::CancellationToken;

use crate::context::Context;
use crate::event::{
ApprovalId, CallId, Cursor, Event, PrincipalId, ProposedCall, ProviderReasoning, ServedBy,
StepTiming, ThreadId, ToolName, ToolResult, Usage,
ApprovalId, AttemptNext, CallId, Cursor, Event, PrincipalId, ProposedCall, ProviderReasoning,
ServedBy, StepTiming, ThreadId, ToolName, ToolResult, Usage,
};

/// The log or the effect ledger refused a write. The engine stops at once and
Expand Down Expand Up @@ -75,6 +75,18 @@ pub enum ModelChunk {
/// Where the step's time went. Sent at most once, after a clean
/// terminal; the engine stores it on `ModelStepCompleted`.
Timing(StepTiming),
/// A model attempt on one route failed. May be sent any number of times
/// before the first other chunk of the step (a final `Abandon` may follow
/// output already streamed); the engine appends it to the log at once
/// and never stores it on the step. `code` is the port's fixed error
/// vocabulary, never the error message.
AttemptFailed {
provider: String,
model: String,
code: String,
elapsed_ms: u64,
then: AttemptNext,
},
}

/// The model call failed after the `Model` port's own retries.
Expand Down
56 changes: 56 additions & 0 deletions vendor/dex-loop/tests/scenarios.rs
Original file line number Diff line number Diff line change
Expand Up @@ -2258,3 +2258,59 @@ async fn uncertain_mutations_keep_their_principal_identity() {
assert_eq!(tools.policy_checks()[0].1, "bob");
assert_eq!(ctx, log.rehydrate());
}

// Failed route attempts are debug rows: the engine appends each one at once,
// before the step's commit point, and they never change the history the next
// model request renders.
#[tokio::test]
async fn failed_attempt_rows_precede_the_commit_and_never_change_history() {
let failed = |then| {
Ok(dex_loop::ModelChunk::AttemptFailed {
provider: "vertex-anthropic".into(),
model: "claude-opus-5-5".into(),
code: "provider_unavailable".into(),
elapsed_ms: 42,
then,
})
};
let log = FakeLog::default();
let model = FakeModel::new(vec![vec![
failed(dex_loop::AttemptNext::Retry),
failed(dex_loop::AttemptNext::Failover),
text("Hello"),
usage(3, 2, 7),
]]);
let tools = FakeTools::new(vec![]);
let engine = engine(&log, &model, &tools, budget());
let mut ctx = log.start_turn("t1", "hi");
assert_eq!(
engine.run(&mut ctx, &CancellationToken::new()).await,
Ok(Exit::Done)
);
assert_eq!(
log.shapes_after(2),
strings(&[
"attempt_failed:1:provider_unavailable:Retry",
"attempt_failed:1:provider_unavailable:Failover",
"delta:Hello",
"usage:5",
"completed:Hello:[]",
"final:Hello",
])
);
let with_rows = log.entries();
let without_rows: Vec<_> = with_rows
.iter()
.filter(|(_, event)| !matches!(event, Event::ModelAttemptFailed { .. }))
.cloned()
.collect();
assert_ne!(with_rows.len(), without_rows.len());
assert_eq!(
history(&dex_loop::rehydrate(thread(), &with_rows)),
history(&dex_loop::rehydrate(thread(), &without_rows)),
);
assert_eq!(
history(&ctx),
history(&dex_loop::rehydrate(thread(), &with_rows))
);
}
1 change: 1 addition & 0 deletions vendor/dex-loop/tests/sim/invariants.rs
Original file line number Diff line number Diff line change
Expand Up @@ -69,6 +69,7 @@ fn kind_str(event: &Event) -> &'static str {
Event::Usage(_) => "usage",
Event::ModelStepCompleted { .. } => "model_step_completed",
Event::ModelAttemptAbandoned { .. } => "model_attempt_abandoned",
Event::ModelAttemptFailed { .. } => "model_attempt_failed",
Event::ToolStarted { .. } => "tool_started",
Event::ToolProgress { .. } => "tool_progress",
Event::ToolsExposed { .. } => "tools_exposed",
Expand Down
5 changes: 5 additions & 0 deletions vendor/dex-loop/tests/support/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -754,6 +754,11 @@ pub fn shape(event: &Event) -> String {
format!("completed:{text}:[{}]", ids(calls))
}
Event::ModelAttemptAbandoned { step } => format!("abandoned:{step}"),
Event::ModelAttemptFailed {
step, code, then, ..
} => {
format!("attempt_failed:{step}:{code}:{then:?}")
}
Event::ToolStarted { call, .. } => format!("started:{call}"),
Event::ToolProgress { call, label } => format!("progress:{call}:{label}"),
Event::ToolsExposed { tools, .. } => format!(
Expand Down
Loading