diff --git a/.repository-projection.json b/.repository-projection.json index b80ec3d2e..a23b41ae0 100644 --- a/.repository-projection.json +++ b/.repository-projection.json @@ -3,11 +3,11 @@ "projection": "deixic-code", "projectionSchemaVersion": 1, "sourceRepository": "dx-corp/mono", - "sourceSha": "ac693f34cdb782d83afd46151bf3a25a2df324be", + "sourceSha": "4c3e830640c4066c99303b7582fa2bdb74095f0d", "destinationRepository": "dx-corp/code", - "priorProjectedBase": "1f00b9e589e080257b8366fc0884453e2596ce90", + "priorProjectedBase": "ca90e530cac35a9d63eac51f920d58891a511b21", "definitionDigest": "82936441c776e3e8edb5d215a75007ec9714a233f489d460075d79d5ef5ba32f", "toolDigest": "c244d99199a7ae3eb8ff644a99462163c23b0bb6a83ef50af01efbdca0b81d04", - "contentDigest": "42d4d938700f91af78522c50de4f281b3d06804428a6744df1f37d0cd3ddc839", + "contentDigest": "f684b3dd9cbf2e4f672beaa7d0c536d87035881e1cdb7f5d28afeebe1e49b530", "publicationEligible": true } diff --git a/packages/dex-host-rs/src/model.rs b/packages/dex-host-rs/src/model.rs index b2f708737..5852189eb 100644 --- a/packages/dex-host-rs/src/model.rs +++ b/packages/dex-host-rs/src/model.rs @@ -177,6 +177,7 @@ impl ChunkTranslator { Ok(value) => value, Err(error) => { return vec![Err(ModelError { + class: dex_loop::ErrorClass::Protocol, message: format!( "provider returned invalid tool-call JSON for {name}: {error}" ), @@ -202,8 +203,25 @@ impl ChunkTranslator { // Not attributed here; see the module doc comment. cost_micros: 0, }))], - StreamEvent::ProviderError { message, .. } | StreamEvent::Error { message } => { - vec![Err(ModelError { message })] + StreamEvent::ProviderError { kind, message } => { + use maestro_ai::ProviderStreamErrorKind; + let class = match kind { + ProviderStreamErrorKind::TransientProtocol => dex_loop::ErrorClass::Truncated, + ProviderStreamErrorKind::OutputTokenExhaustion + | ProviderStreamErrorKind::IncompleteResponse => { + dex_loop::ErrorClass::Incomplete + } + ProviderStreamErrorKind::ProviderDeclaredFailure => { + dex_loop::ErrorClass::Unknown + } + }; + vec![Err(ModelError { class, message })] + } + StreamEvent::Error { message } => { + vec![Err(ModelError { + class: dex_loop::ErrorClass::Unknown, + message, + })] } _ => Vec::new(), } @@ -234,6 +252,7 @@ impl dex_loop::Model for AiRsModel { } Err(error) => { let _ = tx.send(Err(ModelError { + class: dex_loop::ErrorClass::Unknown, message: format!("{error:#}"), })); } @@ -258,6 +277,42 @@ mod tests { } } + #[test] + fn failure_class_comes_from_the_provider_kind_not_display_words() { + use maestro_ai::ProviderStreamErrorKind; + for (kind, expected) in [ + ( + ProviderStreamErrorKind::TransientProtocol, + dex_loop::ErrorClass::Truncated, + ), + ( + ProviderStreamErrorKind::OutputTokenExhaustion, + dex_loop::ErrorClass::Incomplete, + ), + ( + ProviderStreamErrorKind::IncompleteResponse, + dex_loop::ErrorClass::Incomplete, + ), + ( + ProviderStreamErrorKind::ProviderDeclaredFailure, + dex_loop::ErrorClass::Unknown, + ), + ] { + let errors = ChunkTranslator::default().translate(StreamEvent::ProviderError { + kind, + message: "budget auth 429 refusal".into(), + }); + assert_eq!(errors[0].as_ref().unwrap_err().class, expected); + } + let errors = ChunkTranslator::default().translate(StreamEvent::Error { + message: "budget auth 429 refusal".into(), + }); + assert_eq!( + errors[0].as_ref().unwrap_err().class, + dex_loop::ErrorClass::Unknown + ); + } + #[test] fn retains_provider_reported_cache_usage() { let chunks = ChunkTranslator::default().translate(StreamEvent::Usage { diff --git a/packages/dex-host-rs/src/turn.rs b/packages/dex-host-rs/src/turn.rs index bb35cdf5a..fa1ed2838 100644 --- a/packages/dex-host-rs/src/turn.rs +++ b/packages/dex-host-rs/src/turn.rs @@ -224,6 +224,7 @@ mod tests { .pop_front() .unwrap_or_else(|| { vec![Err(ModelError { + class: dex_loop::ErrorClass::Unknown, message: "no script left".into(), })] }); diff --git a/packages/dex-host-rs/tests/turn.rs b/packages/dex-host-rs/tests/turn.rs index 57f4f243d..10ade3289 100644 --- a/packages/dex-host-rs/tests/turn.rs +++ b/packages/dex-host-rs/tests/turn.rs @@ -53,6 +53,7 @@ impl Model for ScriptedModel { .pop_front() .unwrap_or_else(|| { vec![Err(ModelError { + class: dex_loop::ErrorClass::Unknown, message: "no script left".into(), })] }); diff --git a/vendor/dex-loop/src/engine.rs b/vendor/dex-loop/src/engine.rs index 1407f0070..485d10838 100644 --- a/vendor/dex-loop/src/engine.rs +++ b/vendor/dex-loop/src/engine.rs @@ -23,7 +23,6 @@ use std::task::Poll; use std::time::{Duration, Instant, SystemTime, UNIX_EPOCH}; use futures_util::StreamExt; -use futures_util::stream::FuturesUnordered; use tokio_util::sync::CancellationToken; use crate::budget::{Budget, BudgetAxis}; @@ -46,6 +45,7 @@ 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"; +const NOT_RUN_WALL: &str = "not executed: the time limit was reached before this call started"; /// A call that already started (its tool vanished from the catalog, or the /// ledger still shows it running) must never be told "unknown tool" or shown /// a `Running` outcome that nothing will update: both invite the model to @@ -144,7 +144,7 @@ type PrefetchFuture<'e> = Pin + Sen struct Reads<'e> { /// Still running. Each carries its own deadline, the one a wave would /// have given it. - pending: FuturesUnordered>, + pending: HashMap>, /// Returned before the step committed. done: HashMap, completion_order: VecDeque, @@ -170,7 +170,7 @@ struct ReadsHandle<'e>(Arc>>); impl<'e> ReadsHandle<'e> { fn new() -> Self { Self(Arc::new(Mutex::new(Reads { - pending: FuturesUnordered::new(), + pending: HashMap::new(), done: HashMap::new(), completion_order: VecDeque::new(), }))) @@ -181,8 +181,14 @@ impl<'e> ReadsHandle<'e> { self.0.lock().unwrap_or_else(PoisonError::into_inner) } - fn push(&self, read: PrefetchFuture<'e>) { - self.lock().pending.push(read); + fn push(&self, index: usize, read: PrefetchFuture<'e>) { + self.lock().pending.insert(index, read); + } + + fn discard(&self, index: usize) { + let mut reads = self.lock(); + reads.pending.remove(&index); + reads.done.remove(&index); } fn has_pending(&self) -> bool { @@ -193,28 +199,52 @@ impl<'e> ReadsHandle<'e> { self.lock().done.remove(&index) } - fn take_pending(&self) -> FuturesUnordered> { - std::mem::take(&mut self.lock().pending) + fn take_pending(&self, wave: &[usize]) -> Vec<(usize, PrefetchFuture<'e>)> { + let mut reads = self.lock(); + wave.iter() + .filter_map(|index| reads.pending.remove(index).map(|future| (*index, future))) + .collect() } /// The next read to return, if any is running. Cancel-safe: dropping it /// loses nothing, the read stays in the set. async fn next_done(&self) -> Option<(usize, ToolResult)> { - std::future::poll_fn(|cx| self.lock().pending.poll_next_unpin(cx)).await + std::future::poll_fn(|cx| { + let mut reads = self.lock(); + let ready = + reads + .pending + .values_mut() + .find_map(|future| match future.as_mut().poll(cx) { + Poll::Ready(result) => Some(result), + Poll::Pending => None, + }); + if let Some((index, result)) = ready { + reads.pending.remove(&index); + return Poll::Ready(Some((index, result))); + } + if reads.pending.is_empty() { + Poll::Ready(None) + } else { + Poll::Pending + } + }) + .await } - /// Drain peers completed during a log append before polling remaining reads. + /// Completed results first, in completion order, so peers that finished + /// during a log append are drained before polling the reads still running. + /// A wave consumes each completed result once. async fn next_completed(&self) -> Option<(usize, ToolResult)> { - std::future::poll_fn(|cx| { + { let mut reads = self.lock(); while let Some(index) = reads.completion_order.pop_front() { if let Some(result) = reads.done.remove(&index) { - return Poll::Ready(Some((index, result))); + return Some((index, result)); } } - reads.pending.poll_next_unpin(cx) - }) - .await + } + self.next_done().await } /// Awaits `work` while polling the running reads, so a read suspended in @@ -227,7 +257,16 @@ impl<'e> ReadsHandle<'e> { return Poll::Ready(output); } let mut reads = self.lock(); - while let Poll::Ready(Some((index, result))) = reads.pending.poll_next_unpin(cx) { + let ready: Vec<_> = reads + .pending + .values_mut() + .filter_map(|future| match future.as_mut().poll(cx) { + Poll::Ready(result) => Some(result), + Poll::Pending => None, + }) + .collect(); + for (index, result) in ready { + reads.pending.remove(&index); reads.record_done(index, result); } Poll::Pending @@ -388,11 +427,7 @@ where return self.interrupt(ctx, &prefetch).await; } if ctx.open_step().is_some() { - let reads = prefetch.reads.clone(); - if let Some(exit) = reads - .drive(self.dispatch(ctx, cancel, started, &mut prefetch)) - .await? - { + if let Some(exit) = self.dispatch(ctx, cancel, started, &mut prefetch).await? { return Ok(exit); } continue; @@ -405,6 +440,7 @@ where self.emit( ctx, vec![Event::Error { + class: None, code: ErrorCode::BudgetExhausted, message, }], @@ -503,22 +539,21 @@ where // not only against `cancel`. loop { let remaining = self.budget.wall.saturating_sub(started.elapsed()); - // `FuturesUnordered::next` on an empty set resolves at once - // with `None`; the guard keeps it out of the race until a - // read is actually running. + // Only pending reads participate in the stream race; + // completed results remain available for authorized adoption. let has_pending = prefetch.reads.has_pending(); let outcome = tokio::select! { biased; () = cancel.cancelled() => StreamStep::Cancelled, () = tokio::time::sleep(remaining) => StreamStep::WallExceeded, - item = stream.next() => match item { - Some(chunk) => StreamStep::Chunk(chunk), - None => StreamStep::Ended, - }, done = prefetch.reads.next_done(), if has_pending => match done { Some((index, result)) => StreamStep::Prefetched(index, result), None => continue, }, + item = stream.next() => match item { + Some(chunk) => StreamStep::Chunk(chunk), + None => StreamStep::Ended, + }, }; match outcome { StreamStep::Chunk(Ok(ModelChunk::Thinking(delta))) => { @@ -583,7 +618,7 @@ where prefetch.reads.drive(self.log.append(&event)).await?; } StreamStep::Chunk(Err(error)) => { - failure = Some(error.message); + failure = Some(error); break; } StreamStep::Ended | StreamStep::Cancelled => break, @@ -620,6 +655,7 @@ where let mut events = pending_usage; events.push(Event::ModelAttemptAbandoned { step }); events.push(Event::Error { + class: Some(crate::ErrorClass::BudgetExhausted), code: ErrorCode::BudgetExhausted, message, }); @@ -652,12 +688,13 @@ where self.emit(ctx, events).await?; return Ok(Some(Exit::Done)); } - if let Some(message) = failure { + if let Some(error) = failure { let mut events = pending_usage; events.push(Event::ModelAttemptAbandoned { step }); events.push(Event::Error { + class: Some(error.class()), code: ErrorCode::ModelFailed, - message, + message: error.message, }); self.emit(ctx, events).await?; return Ok(Some(Exit::Failed)); @@ -685,6 +722,7 @@ where let mut events = pending_usage; events.push(Event::ModelAttemptAbandoned { step }); events.push(Event::Error { + class: Some(crate::ErrorClass::BudgetExhausted), code: ErrorCode::BudgetExhausted, message, }); @@ -750,10 +788,24 @@ where return Ok(()); }; let eligible = spec.read_only + && !ctx.client_tools().iter().any(|tool| tool.name == call.tool) && !matches!(spec.executor, ExecutorKind::User | ExecutorKind::Client) && validate_args(&spec, &call.args).is_ok() && !ctx.has_uncertain_call(call); - if !eligible || self.tools.policy(ctx, call).await != Verdict::Allow { + if !eligible { + return Ok(()); + } + let deadline = tokio::time::Instant::now() + self.call_deadline(run_started); + let verdict = tokio::select! { + biased; + () = cancel.cancelled() => return Ok(()), + () = tokio::time::sleep_until(deadline) => return Ok(()), + verdict = self.tools.policy(ctx, call) => verdict, + }; + if verdict != Verdict::Allow + || cancel.is_cancelled() + || tokio::time::Instant::now() >= deadline + { return Ok(()); } // Appended without `ctx.observe`: the stream still borrows `ctx`. @@ -770,17 +822,22 @@ where }; prefetch.unobserved.push((cursor, event)); prefetch.started.insert(index); - let deadline = self.call_deadline(run_started); let thread = ctx.thread().clone(); let call = call.clone(); - prefetch.reads.push(Box::pin(async move { - let run = self.tools.run(&thread, &call, cancel); - let result = match tokio::time::timeout(deadline, run).await { - Ok(result) => result, - Err(_elapsed) => ToolResult::error(DEADLINE_READ), - }; - (index, result) - })); + prefetch.reads.push( + index, + Box::pin(async move { + if cancel.is_cancelled() || tokio::time::Instant::now() >= deadline { + return (index, ToolResult::error(DEADLINE_READ)); + } + let run = self.tools.run(&thread, &call, cancel); + let result = match tokio::time::timeout_at(deadline, run).await { + Ok(result) => result, + Err(_elapsed) => ToolResult::error(DEADLINE_READ), + }; + (index, result) + }), + ); Ok(()) } @@ -820,6 +877,9 @@ where if cancel.is_cancelled() { break; } + if let Some(exit) = self.stop_expired_dispatch(ctx, run_started).await? { + return Ok(Some(exit)); + } let state = ctx .open_step() .and_then(|step| step.states.get(index)) @@ -888,15 +948,12 @@ where .. }) => Some(decision), Some(CallState::Started) => { - // Started before a restart and never finished. Reads run - // again; mutations resolve through the ledger, which - // never dispatches a claimed call a second time. if call.tool.as_str() == TOOLS_SEARCH { self.search_tools(ctx, call).await?; continue; } match self.offered_spec(ctx, &call.tool) { - Some(spec) if spec.read_only => wave.push(index), + Some(spec) if spec.read_only => None, Some(_) => { if self .flush(ctx, &calls, &mut wave, cancel, run_started, prefetch) @@ -905,16 +962,14 @@ where break; } self.run_mutation(ctx, call, cancel, run_started).await?; + continue; + } + None => { + prefetch.reads.discard(index); + self.resolve_started(ctx, call).await?; + continue; } - // The tool is no longer offered (deploy, grant - // revoke), but this call already started: it may - // have run. Resolve through the ledger instead of - // telling the model to retry with a different tool, - // which would dispatch a new call id for the same - // mutation. - None => self.resolve_started(ctx, call).await?, } - continue; } Some(CallState::Todo) => None, }; @@ -931,6 +986,7 @@ where // executor: the model sees why and can call again. (A call the // stream already started passed this check before it ran.) if let Err(reason) = validate_args(&spec, &call.args) { + prefetch.reads.discard(index); self.finish(ctx, call, ToolResult::error(reason)).await?; continue; } @@ -963,8 +1019,20 @@ where } // Current policy first, even for a decided call: a revoked // grant or changed policy denies it. - let verdict = match self.tools.policy(ctx, call).await { + let remaining = self.budget.wall.saturating_sub(run_started.elapsed()); + let current_policy = tokio::select! { + biased; + () = cancel.cancelled() => break, + () = tokio::time::sleep(remaining) => { + prefetch.reads.discard(index); + self.finish(ctx, call, ToolResult::error(DEADLINE_READ)).await?; + break; + } + verdict = self.tools.policy(ctx, call) => verdict, + }; + let verdict = match current_policy { Verdict::Deny(reason) => { + prefetch.reads.discard(index); self.finish(ctx, call, ToolResult::error(format!("denied: {reason}"))) .await?; continue; @@ -1059,7 +1127,73 @@ where if cancel.is_cancelled() { return self.interrupt(ctx, prefetch).await.map(Some); } - Ok(None) + self.stop_expired_dispatch(ctx, run_started).await + } + + /// Seal the open step after its wall budget: never dispatch a fresh call + /// through a zero-duration timeout, which polls its inner future first. + /// Previously started effects still resolve through their durable owner. + async fn stop_expired_dispatch( + &self, + ctx: &mut Context, + run_started: Instant, + ) -> Result, Fenced> { + if run_started.elapsed() < self.budget.wall { + return Ok(None); + } + let pending: Vec<_> = ctx + .open_step() + .map(|step| { + step.calls + .iter() + .cloned() + .zip(step.states.iter().cloned()) + .collect() + }) + .unwrap_or_default(); + for (call, state) in pending { + match state { + CallState::Done(_) => {} + CallState::Started + if self + .offered_spec(ctx, &call.tool) + .is_some_and(|spec| spec.read_only) => + { + self.finish(ctx, &call, ToolResult::error(DEADLINE_READ)) + .await?; + } + CallState::Started => self.resolve_started(ctx, &call).await?, + CallState::AwaitingClient { + result: Some(result), + .. + } => { + self.finish_client_result(ctx, &call, result).await?; + } + CallState::AwaitingClient { result: None, .. } => { + self.timeout_client_call(ctx, &call).await?; + } + CallState::Asked { + answer: Some(answer), + } => { + self.finish(ctx, &call, ToolResult::text(answer)).await?; + } + _ => { + self.finish(ctx, &call, ToolResult::error(NOT_RUN_WALL)) + .await? + } + } + } + let message = self.budget_message(ctx, BudgetAxis::Wall); + self.emit( + ctx, + vec![Event::Error { + code: ErrorCode::BudgetExhausted, + message, + class: Some(crate::ErrorClass::BudgetExhausted), + }], + ) + .await?; + Ok(Some(Exit::Failed)) } /// Refusal does not claim or dispatch a second effect. Reads remain safe @@ -1172,7 +1306,7 @@ where prefetch, ) .await?; - Ok(cancel.is_cancelled()) + Ok(cancel.is_cancelled() || run_started.elapsed() >= self.budget.wall) } /// Runs read-only calls concurrently. Each `ToolFinished` is appended as @@ -1189,7 +1323,7 @@ where run_started: Instant, prefetch: &mut Prefetch<'e>, ) -> Result<(), Fenced> { - if wave.is_empty() || cancel.is_cancelled() { + if wave.is_empty() || cancel.is_cancelled() || run_started.elapsed() >= self.budget.wall { return Ok(()); } let starts: Vec = wave @@ -1201,36 +1335,44 @@ where Some(started(call, &spec)) }) .collect(); + let deadline = tokio::time::Instant::now() + self.call_deadline(run_started); self.emit(ctx, starts).await?; let thread = ctx.thread().clone(); // One deadline for the wave: its reads run concurrently, so each // gets the full time. A read that overruns is dropped and finished // `Failed`; a read has no effect to wait for, so retrying is safe. - let deadline = self.call_deadline(run_started); let running = ReadsHandle::new(); for &index in &wave { if let Some(result) = prefetch.reads.take_done(index) { - self.finish(ctx, &calls[index], result).await?; + running + .drive(self.finish(ctx, &calls[index], result)) + .await?; continue; } if prefetch.started.contains(&index) { // Still running from the stream; it arrives through - // `prefetch.pending` below. + // the indexed pending wave below. continue; } let call = calls[index].clone(); let thread = thread.clone(); - running.push(Box::pin(async move { - let run = self.tools.run(&thread, &call, cancel); - let result = match tokio::time::timeout(deadline, run).await { - Ok(result) => result, - Err(_elapsed) => ToolResult::error(DEADLINE_READ), - }; - (index, result) - })); + running.push( + index, + Box::pin(async move { + if cancel.is_cancelled() || tokio::time::Instant::now() >= deadline { + return (index, ToolResult::error(DEADLINE_READ)); + } + let run = self.tools.run(&thread, &call, cancel); + let result = match tokio::time::timeout_at(deadline, run).await { + Ok(result) => result, + Err(_elapsed) => ToolResult::error(DEADLINE_READ), + }; + (index, result) + }), + ); } - for future in prefetch.reads.take_pending() { - running.push(future); + for (index, future) in prefetch.reads.take_pending(&wave) { + running.push(index, future); } // On `Fenced` the remaining reads are dropped: a stale owner must not // append, and the new owner runs them again. @@ -1241,7 +1383,7 @@ where .is_some_and(|state| !matches!(state, CallState::Done(_))); // A prefetched read whose call `dispatch` already finished // (policy denied it on re-check) has nothing left to report. - if !still_open || !(wave.contains(&index) || prefetch.started.contains(&index)) { + if !still_open || !wave.contains(&index) { continue; } running @@ -1273,23 +1415,59 @@ where // never as "unknown tool". return self.resolve_started(ctx, call).await; }; + let already_started = ctx.open_step().is_some_and(|step| { + step.calls + .iter() + .zip(&step.states) + .any(|(proposal, state)| { + proposal.id == call.id && matches!(state, CallState::Started) + }) + }); + if self.call_deadline(run_started).is_zero() { + return if already_started { + self.resolve_started(ctx, call).await + } else { + self.finish(ctx, call, ToolResult::error(NOT_RUN_WALL)) + .await + }; + } match self.effects.claim(call).await? { Claim::Existing(result) => { let result = self.settle_claim(&call.id, result).await?; self.finish(ctx, call, result).await } Claim::Granted => { + // Claim persistence can consume the remaining wall. A fresh + // claim is known not executed; a historical Started call may + // already have taken effect and must remain Unknown. + if self.call_deadline(run_started).is_zero() { + let result = if already_started { + ToolResult::unknown(UNKNOWN_NO_RETRY) + } else { + ToolResult::error(NOT_RUN_WALL) + }; + self.effects.record(&call.id, &result).await?; + return self.finish(ctx, call, result).await; + } self.emit(ctx, vec![started(call, &spec)]).await?; // A mutation that overruns its deadline is dropped, not // cancelled: the effect may still land. `Unknown` is recorded // under the claim, so a later resume of this call adopts it - // instead of dispatching the mutation a second time. An - // interrupt only fires `cancel`; the run is still awaited. - let run = self.tools.run(ctx.thread(), call, cancel); - let result = match tokio::time::timeout(self.call_deadline(run_started), run).await - { - Ok(result) => result, - Err(_elapsed) => ToolResult::unknown(DEADLINE_MUTATION), + // instead of dispatching the mutation a second time. Interrupts + // reach cancellable process tools; settlement is still awaited. + let deadline = self.call_deadline(run_started); + let result = if deadline.is_zero() { + if already_started { + ToolResult::unknown(UNKNOWN_NO_RETRY) + } else { + ToolResult::error(NOT_RUN_WALL) + } + } else { + let run = self.tools.run(ctx.thread(), call, cancel); + match tokio::time::timeout(deadline, run).await { + Ok(result) => result, + Err(_elapsed) => ToolResult::unknown(DEADLINE_MUTATION), + } }; self.effects.record(&call.id, &result).await?; self.finish(ctx, call, result).await @@ -1711,17 +1889,10 @@ fn non_empty_str<'a>(call: &'a ProposedCall, key: &str) -> Option<&'a str> { .filter(|value| !value.trim().is_empty()) } -/// `args` against `spec.schema`. A schema that is absent, not an object, or -/// does not compile validates nothing (the executor still checks what it -/// needs); a schema violation names the first error so the model can call -/// again with arguments that fit. +/// Validate boolean and object JSON schemas; malformed schemas fail closed. fn validate_args(spec: &ToolSpec, args: &serde_json::Value) -> Result<(), String> { - if !spec.schema.is_object() { - return Ok(()); - } - let Ok(validator) = jsonschema::validator_for(&spec.schema) else { - return Ok(()); - }; + let validator = jsonschema::validator_for(&spec.schema) + .map_err(|error| format!("invalid tool schema: {error}"))?; match validator.iter_errors(args).next() { None => Ok(()), Some(error) => Err(format!("invalid arguments: {error}")), diff --git a/vendor/dex-loop/src/event.rs b/vendor/dex-loop/src/event.rs index e52cef9f8..195dbddbe 100644 --- a/vendor/dex-loop/src/event.rs +++ b/vendor/dex-loop/src/event.rs @@ -171,13 +171,21 @@ pub struct Usage { /// The part of `input_tokens` served from the prompt cache. Zero on rows /// written before the field existed, and for providers that do not /// report it. - #[serde(default)] + #[serde(default, alias = "cache_read_tokens", skip_serializing_if = "is_zero")] pub cache_read_input_tokens: u64, /// The part of `input_tokens` written to the prompt cache. Same default. - #[serde(default)] + #[serde( + default, + alias = "cache_creation_tokens", + skip_serializing_if = "is_zero" + )] pub cache_creation_input_tokens: u64, } +fn is_zero(value: &u64) -> bool { + *value == 0 +} + impl Usage { pub fn tokens(&self) -> u64 { self.input_tokens.saturating_add(self.output_tokens) @@ -363,6 +371,100 @@ impl ErrorCode { } } +/// A content-free failure class supplied by the typed model adapter. +#[derive(Clone, Copy, Debug, Default, PartialEq, Eq, Serialize, Deserialize)] +#[serde(rename_all = "snake_case")] +pub enum ErrorClass { + /// The turn hit its step, token, cost or wall budget. + BudgetExhausted, + /// The stream ended before the response completed. + Truncated, + /// The provider ended the response early (for example the output token + /// limit). + Incomplete, + /// The model proposed a tool call whose arguments were not a JSON object. + MalformedToolCall, + /// The response held neither text nor tool calls. + EmptyCompletion, + /// The gateway could not be reached or the connection broke. + Transport, + /// The provider or gateway throttled the request. + RateLimited, + /// The provider or gateway reported itself unavailable or overloaded. + Unavailable, + /// The provider or a safety filter refused to answer. + Refusal, + /// The service token for the gateway could not be minted. + Auth, + /// The wire broke the SSE or Responses event contract. + Protocol, + /// Stored history could not be loaded for the request. + Resolve, + /// The host could not prepare the request. + Host, + /// Any other gateway or provider rejection. + Rejected, + /// An error the classifier does not recognise. Old and new codes land + /// here rather than failing. + #[default] + #[serde(other)] + Unknown, +} + +impl ErrorClass { + /// The stable label stored in `dex_turn_outcomes.error_class`. + pub fn as_str(self) -> &'static str { + match self { + ErrorClass::BudgetExhausted => "budget_exhausted", + ErrorClass::Truncated => "truncated", + ErrorClass::Incomplete => "incomplete", + ErrorClass::MalformedToolCall => "malformed_tool_call", + ErrorClass::EmptyCompletion => "empty_completion", + ErrorClass::Transport => "transport", + ErrorClass::RateLimited => "rate_limited", + ErrorClass::Unavailable => "unavailable", + ErrorClass::Refusal => "refusal", + ErrorClass::Auth => "auth", + ErrorClass::Protocol => "protocol", + ErrorClass::Resolve => "resolve", + ErrorClass::Host => "host", + ErrorClass::Rejected => "rejected", + ErrorClass::Unknown => "unknown", + } + } + + /// Explicit wire-code boundary; human-readable details are never classified. + pub fn of_gateway_code(code: &str) -> Self { + match code { + "rate_limit_error" + | "rate_limit_exceeded" + | "resource_exhausted" + | "quota_exceeded" + | "too_many_requests" => Self::RateLimited, + "refusal" + | "content_filter" + | "content_filter_error" + | "safety" + | "policy_violation" => Self::Refusal, + "upstream_unavailable" + | "overloaded_error" + | "unavailable" + | "api_error" + | "internal_error" + | "server_error" + | "timeout" => Self::Unavailable, + "authentication_error" | "unauthorized" => Self::Auth, + "invalid_request" + | "invalid_request_error" + | "bad_request" + | "permission_error" + | "forbidden" + | "not_found" => Self::Rejected, + _ => Self::Unknown, + } + } +} + /// Who can answer a turn's approval requests. /// /// `Headless` turns come from callers with no human to click Approve (service @@ -628,6 +730,9 @@ pub enum Event { Error { code: ErrorCode, message: String, + /// Typed adapter classification; old rows and old readers remain compatible. + #[serde(default, skip_serializing_if = "Option::is_none")] + class: Option, }, Interrupted, } @@ -866,6 +971,7 @@ mod tests { receipt: None, }, Event::Error { + class: None, code: ErrorCode::BudgetExhausted, message: "steps".into(), }, @@ -1109,4 +1215,53 @@ mod tests { usage ); } + #[test] + fn usage_rows_written_before_the_cache_split_still_decode() { + let old: Usage = + serde_json::from_str(r#"{"input_tokens":5,"output_tokens":2,"cost_micros":9}"#) + .expect("old row decodes"); + assert_eq!( + old, + Usage { + input_tokens: 5, + output_tokens: 2, + cost_micros: 9, + ..Usage::default() + } + ); + // A zero split is not written, so unsplit providers keep their shape. + let json = serde_json::to_value(old).expect("encodes"); + assert_eq!( + json, + serde_json::json!({"input_tokens":5,"output_tokens":2,"cost_micros":9}) + ); + let split = Usage { + cache_read_input_tokens: 3, + ..old + }; + let back: Usage = + serde_json::from_value(serde_json::to_value(split).expect("encodes")).expect("decodes"); + assert_eq!(back, split); + } + + #[test] + fn old_error_rows_do_not_classify_human_readable_prose() { + let old = serde_json::json!({"type":"error", "code":"model_failed", "message":"rate_limit_error: arbitrary prose"}); + assert!(matches!( + serde_json::from_value::(old).unwrap(), + Event::Error { class: None, .. } + )); + assert_eq!( + serde_json::to_string(&ErrorCode::ModelFailed).unwrap(), + r#""model_failed""# + ); + assert_eq!( + ErrorClass::of_gateway_code("not_a_rate_limit_error"), + ErrorClass::Unknown + ); + assert_eq!( + ErrorClass::of_gateway_code("rate_limit_error"), + ErrorClass::RateLimited + ); + } } diff --git a/vendor/dex-loop/src/lib.rs b/vendor/dex-loop/src/lib.rs index c347a81ed..06fda88e6 100644 --- a/vendor/dex-loop/src/lib.rs +++ b/vendor/dex-loop/src/lib.rs @@ -34,9 +34,9 @@ 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, 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, + Cursor, ErrorClass, 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, diff --git a/vendor/dex-loop/src/ports.rs b/vendor/dex-loop/src/ports.rs index 8fa59b139..9599ee377 100644 --- a/vendor/dex-loop/src/ports.rs +++ b/vendor/dex-loop/src/ports.rs @@ -9,8 +9,8 @@ use tokio_util::sync::CancellationToken; use crate::context::Context; use crate::event::{ - ApprovalId, AttemptNext, CallId, Cursor, Event, PrincipalId, ProposedCall, ProviderReasoning, - ServedBy, StepTiming, ThreadId, ToolName, ToolResult, Usage, + ApprovalId, AttemptNext, CallId, Cursor, ErrorClass, 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 @@ -94,6 +94,14 @@ pub enum ModelChunk { #[error("model call failed: {message}")] pub struct ModelError { pub message: String, + pub class: ErrorClass, +} + +impl ModelError { + /// The adapter-supplied class; message text cannot change it. + pub fn class(&self) -> ErrorClass { + self.class + } } /// The model, streaming. diff --git a/vendor/dex-loop/tests/prefetch.rs b/vendor/dex-loop/tests/prefetch.rs index 611731211..30e933435 100644 --- a/vendor/dex-loop/tests/prefetch.rs +++ b/vendor/dex-loop/tests/prefetch.rs @@ -317,6 +317,7 @@ async fn a_failed_stream_closes_the_reads_it_started_and_records_no_result() { let model = FakeModel::new(vec![vec![ call("search", json!({"key": "a"})), Err(ModelError { + class: dex_loop::ErrorClass::Unknown, message: "connection reset".into(), }), ]]) @@ -354,6 +355,7 @@ async fn a_new_turn_after_a_failed_stream_starts_its_own_calls_once() { vec![ call("search", json!({"key": "a"})), Err(ModelError { + class: dex_loop::ErrorClass::Unknown, message: "connection reset".into(), }), ], diff --git a/vendor/dex-loop/tests/prefetch_authority.rs b/vendor/dex-loop/tests/prefetch_authority.rs new file mode 100644 index 000000000..fcb53aec5 --- /dev/null +++ b/vendor/dex-loop/tests/prefetch_authority.rs @@ -0,0 +1,481 @@ +//! Policy, schema and original deadline boundaries for speculative reads. +#[allow(dead_code)] +mod support; +use dex_loop::{ + Budget, CancellationToken, Context, Engine, Event, Exit, Lexicon, PrincipalId, ProposedCall, + ThreadId, ToolName, ToolResult, ToolSpec, Tools, Verdict, +}; +use serde_json::json; +use std::sync::{ + Arc, + atomic::{AtomicUsize, Ordering}, +}; +use std::time::Duration; +use support::*; + +#[derive(Clone)] +struct PolicyTools { + inner: FakeTools, + checks: Arc, + allow_checks: usize, + delay: Duration, +} +impl Tools for PolicyTools { + fn catalog(&self) -> &[ToolSpec] { + self.inner.catalog() + } + async fn search(&self, p: &PrincipalId, q: &str) -> Vec { + self.inner.search(p, q).await + } + async fn policy(&self, _: &Context, _: &ProposedCall) -> Verdict { + let index = self.checks.fetch_add(1, Ordering::SeqCst); + tokio::time::sleep(self.delay).await; + if index < self.allow_checks { + Verdict::Allow + } else { + Verdict::Deny("grant revoked".into()) + } + } + async fn run(&self, t: &ThreadId, c: &ProposedCall, cancel: &CancellationToken) -> ToolResult { + self.inner.run(t, c, cancel).await + } +} +fn limits(wall: Duration) -> Budget { + Budget { + max_steps: 10, + max_tokens: 1_000_000, + max_cost_micros: 1_000_000, + wall, + } +} +fn governed( + log: &FakeLog, + model: &FakeModel, + tools: PolicyTools, + wall: Duration, +) -> Engine { + Engine::new( + log.clone(), + model.clone(), + tools, + FakeEffects::default(), + Lexicon::default(), + limits(wall), + ) +} +fn tools(inner: FakeTools, allow_checks: usize, delay: Duration) -> PolicyTools { + PolicyTools { + inner, + checks: Arc::default(), + allow_checks, + delay, + } +} +fn finishes(log: &FakeLog) -> Vec { + log.events() + .into_iter() + .filter_map(|e| match e { + Event::ToolFinished { output, .. } => Some(output), + _ => None, + }) + .collect() +} +#[tokio::test] +async fn completed_prefetch_is_discarded_when_current_policy_revokes_it() { + let log = FakeLog::default(); + let model = FakeModel::new(vec![ + vec![call("read", json!({})), text("tail")], + vec![text("done")], + ]) + .with_chunk_delay(Duration::from_millis(30)); + let inner = FakeTools::new(vec![read_tool("read")]); + let tools = tools(inner.clone(), 1, Duration::ZERO); + let checks = tools.checks.clone(); + let engine = governed(&log, &model, tools, Duration::from_secs(2)); + let mut ctx = log.start_turn("t1", "read it"); + assert_eq!( + engine.run(&mut ctx, &CancellationToken::new()).await, + Ok(Exit::Done) + ); + assert_eq!(checks.load(Ordering::SeqCst), 2); + assert_eq!(inner.run_ids(), strings(&["t1-1-0"])); + assert_eq!( + finishes(&log), + vec![dex_loop::Output::Text("denied: grant revoked".into())] + ); + assert!(!history(&ctx).join("\n").contains("out/t1-1-0")); + assert_eq!(log.rehydrate(), ctx); +} +#[tokio::test] +async fn a_restarted_started_read_rechecks_authority_before_execution() { + let log = FakeLog::default(); + log.start_turn("t1", "read it"); + let proposal = dex_loop::ProposedCall::new( + call_id("t1", 1, 0), + ToolName::new("read"), + json!({}), + alice(), + ); + log.host_append(Event::StepStarted { + step: 1, + control_through: dex_loop::Cursor::START, + }); + log.host_append(Event::ModelStepCompleted { + served: None, + step: 1, + text: String::new(), + calls: vec![proposal.clone()], + reasoning: None, + timing: None, + }); + log.host_append(Event::ToolStarted { + call: proposal.id, + tool: proposal.tool, + label: "read".into(), + principal: alice(), + }); + let model = FakeModel::new(vec![vec![text("done")]]); + let inner = FakeTools::new(vec![read_tool("read")]); + let engine = governed( + &log, + &model, + tools(inner.clone(), 0, Duration::ZERO), + Duration::from_secs(2), + ); + let mut ctx = log.rehydrate(); + assert_eq!( + engine.run(&mut ctx, &CancellationToken::new()).await, + Ok(Exit::Done) + ); + assert!(inner.run_ids().is_empty()); + assert_eq!( + finishes(&log), + vec![dex_loop::Output::Text("denied: grant revoked".into())] + ); +} +#[tokio::test] +async fn malformed_and_false_schemas_never_reach_a_read_or_mutation_executor() { + for schema in [ + json!(false), + json!(null), + json!({"type":7}), + json!({"$ref":"#/$defs/missing"}), + ] { + for read_only in [true, false] { + let mut spec = if read_only { + read_tool("check") + } else { + write_tool("check") + }; + spec.schema = schema.clone(); + let log = FakeLog::default(); + let model = FakeModel::new(vec![vec![call("check", json!({}))], vec![text("done")]]); + let inner = FakeTools::new(vec![spec]); + let engine = engine(&log, &model, &inner, limits(Duration::from_secs(2))); + let mut ctx = log.start_turn("t1", "check"); + assert_eq!( + engine.run(&mut ctx, &CancellationToken::new()).await, + Ok(Exit::Done) + ); + assert!(inner.run_ids().is_empty(), "{schema}"); + assert!( + !log.events() + .iter() + .any(|e| matches!(e, Event::ToolStarted { .. })), + "{schema}" + ); + assert_eq!(finishes(&log).len(), 1); + } + } +} +#[tokio::test] +async fn boolean_true_and_valid_object_schemas_remain_callable() { + for schema in [json!(true), json!({"type":"object"})] { + let mut spec = read_tool("check"); + spec.schema = schema; + let log = FakeLog::default(); + let model = FakeModel::new(vec![vec![call("check", json!({}))], vec![text("done")]]); + let inner = FakeTools::new(vec![spec]); + let engine = engine(&log, &model, &inner, limits(Duration::from_secs(2))); + let mut ctx = log.start_turn("t1", "check"); + assert_eq!( + engine.run(&mut ctx, &CancellationToken::new()).await, + Ok(Exit::Done) + ); + assert_eq!(inner.run_ids(), strings(&["t1-1-0"])); + } +} +#[tokio::test] +async fn a_pending_prefetch_policy_cannot_outlive_the_wall_or_interrupt() { + for interrupt in [false, true] { + let log = FakeLog::default(); + let model = FakeModel::new(vec![vec![call("read", json!({}))]]); + let inner = FakeTools::new(vec![read_tool("read")]); + let engine = governed( + &log, + &model, + tools(inner.clone(), usize::MAX, Duration::from_secs(60)), + Duration::from_millis(150), + ); + let mut ctx = log.start_turn("t1", "read"); + let cancel = CancellationToken::new(); + let host = async { + if interrupt { + tokio::time::sleep(Duration::from_millis(30)).await; + cancel.cancel(); + } + }; + let (exit, ()) = tokio::time::timeout(Duration::from_secs(1), async { + tokio::join!(engine.run(&mut ctx, &cancel), host) + }) + .await + .expect("policy must be bounded"); + assert_eq!( + exit, + Ok(if interrupt { + Exit::Interrupted + } else { + Exit::Failed + }) + ); + assert!(inner.run_ids().is_empty()); + assert!( + !log.events() + .iter() + .any(|e| matches!(e, Event::ToolStarted { .. })) + ); + } +} + +#[derive(Clone)] +struct SlowStartLog(FakeLog); +impl dex_loop::Log for SlowStartLog { + async fn append(&self, events: &[Event]) -> Result, dex_loop::Fenced> { + let cursors = self.0.append(events).await?; + if events + .iter() + .any(|e| matches!(e, Event::ToolStarted { .. })) + { + tokio::time::sleep(Duration::from_millis(80)).await; + } + Ok(cursors) + } + async fn append_text(&self, text: String) -> Result<(), dex_loop::Fenced> { + self.0.append_text(text).await + } + async fn control_since( + &self, + c: dex_loop::Cursor, + ) -> Result, dex_loop::Fenced> { + self.0.control_since(c).await + } +} + +#[tokio::test] +async fn an_expired_queued_read_is_never_polled_into_the_executor() { + let log = FakeLog::default(); + let model = FakeModel::new(vec![vec![call("read", json!({}))], vec![text("done")]]); + let first_polls = Arc::new(AtomicUsize::new(0)); + let count = first_polls.clone(); + let inner = FakeTools::new(vec![read_tool("read")]).on_run(move |_| { + count.fetch_add(1, Ordering::SeqCst); + }); + let engine = Engine::new( + SlowStartLog(log.clone()), + model, + inner, + FakeEffects::default(), + Lexicon::default(), + limits(Duration::from_secs(2)), + ) + .with_tool_call_deadline(Duration::from_millis(40)); + let mut ctx = log.start_turn("t1", "read it"); + assert_eq!( + engine.run(&mut ctx, &CancellationToken::new()).await, + Ok(Exit::Done) + ); + assert_eq!(first_polls.load(Ordering::SeqCst), 0); + assert_eq!(finishes(&log).len(), 1); + assert!( + matches!(&finishes(&log)[0], dex_loop::Output::Text(text) if text.contains("time limit")) + ); + assert_eq!(log.rehydrate(), ctx); +} + +#[tokio::test] +async fn a_running_revoked_prefetch_is_dropped_without_waiting_or_exposing_output() { + let log = FakeLog::default(); + let model = FakeModel::new(vec![ + vec![call("read", json!({"key":"slow"})), text("tail")], + vec![text("done")], + ]) + .with_chunk_delay(Duration::from_millis(20)); + let first_polls = Arc::new(AtomicUsize::new(0)); + let count = first_polls.clone(); + let inner = FakeTools::new(vec![read_tool("read")]) + .delay("slow", Duration::from_secs(60)) + .on_run(move |_| { + count.fetch_add(1, Ordering::SeqCst); + }); + let engine = governed( + &log, + &model, + tools(inner, 1, Duration::ZERO), + Duration::from_secs(2), + ); + let mut ctx = log.start_turn("t1", "read"); + assert_eq!( + tokio::time::timeout( + Duration::from_secs(1), + engine.run(&mut ctx, &CancellationToken::new()) + ) + .await + .expect("revoked future must be dropped"), + Ok(Exit::Done) + ); + assert_eq!(first_polls.load(Ordering::SeqCst), 1); + assert_eq!( + finishes(&log), + vec![dex_loop::Output::Text("denied: grant revoked".into())] + ); +} + +#[tokio::test] +async fn an_in_process_client_bridge_waits_for_the_committed_model_step() { + let log = FakeLog::default(); + let model = FakeModel::new(vec![ + vec![call("page.read", json!({})), text("tail")], + vec![text("done")], + ]) + .with_chunk_delay(Duration::from_millis(20)); + let committed = log.clone(); + let inner = FakeTools::new(vec![read_tool("page.read")]).on_run(move |_| { + assert!( + committed + .events() + .iter() + .any(|event| matches!(event, Event::ModelStepCompleted { step: 1, .. })), + "a client bridge must not wait for a response before its step commits" + ); + }); + let engine = engine(&log, &model, &inner, limits(Duration::from_secs(2))); + let mut ctx = log.start_turn_with_client_tools( + "t1", + "read my page", + vec![dex_loop::ClientToolSpec { + name: ToolName::new("page.read"), + schema: json!({}), + read_only: true, + label: "Read page".into(), + }], + ); + assert_eq!( + engine.run(&mut ctx, &CancellationToken::new()).await, + Ok(Exit::Done) + ); + assert_eq!(inner.run_ids(), strings(&["t1-1-0"])); + assert_eq!(log.rehydrate(), ctx); +} + +/// A host read may suspend while owning the same row lock as a progress/log +/// write. The engine must keep polling it through commit and wave completion. +#[tokio::test] +async fn host_read_log_contention_does_not_deadlock_commit_or_wave_completion() { + use dex_loop::{Cursor, Fenced, Log}; + #[derive(Clone)] + struct ContendedLog { + inner: FakeLog, + row: Arc>, + waits: Arc, + } + impl Log for ContendedLog { + async fn append(&self, events: &[Event]) -> Result, Fenced> { + if self.row.try_lock().is_err() { + self.waits.fetch_add(1, Ordering::SeqCst); + } + let _row = self.row.lock().await; + self.inner.append(events).await + } + async fn append_text(&self, text: String) -> Result<(), Fenced> { + let _row = self.row.lock().await; + self.inner.append_text(text).await + } + async fn control_since(&self, after: Cursor) -> Result, Fenced> { + self.inner.control_since(after).await + } + } + #[derive(Clone)] + struct ContendedTools { + inner: FakeTools, + row: Arc>, + } + impl Tools for ContendedTools { + fn catalog(&self) -> &[ToolSpec] { + self.inner.catalog() + } + async fn search(&self, principal: &PrincipalId, query: &str) -> Vec { + self.inner.search(principal, query).await + } + async fn policy(&self, ctx: &Context, call: &ProposedCall) -> Verdict { + self.inner.policy(ctx, call).await + } + async fn run( + &self, + thread: &ThreadId, + call: &ProposedCall, + cancel: &CancellationToken, + ) -> ToolResult { + let _row = self.row.lock().await; + tokio::time::sleep(Duration::from_millis(30)).await; + self.inner.run(thread, call, cancel).await + } + } + let inner = FakeLog::default(); + let row = Arc::new(tokio::sync::Mutex::new(())); + let waits = Arc::new(AtomicUsize::new(0)); + let log = ContendedLog { + inner: inner.clone(), + row: row.clone(), + waits: waits.clone(), + }; + let reads = FakeTools::new(vec![read_tool("read")]); + let tools = ContendedTools { + inner: reads.clone(), + row, + }; + let model = FakeModel::new(vec![ + vec![ + call("read", json!({})), + call("read", json!({})), + text("tail"), + ], + vec![text("done")], + ]) + .with_chunk_delay(Duration::from_millis(1)); + let engine = Engine::new( + log, + model, + tools, + FakeEffects::default(), + Lexicon::default(), + limits(Duration::from_secs(2)), + ); + let mut ctx = inner.start_turn("t1", "read both"); + assert_eq!( + tokio::time::timeout( + Duration::from_secs(2), + engine.run(&mut ctx, &CancellationToken::new()) + ) + .await + .expect("row lock released while its host read is polled"), + Ok(Exit::Done) + ); + assert!( + waits.load(Ordering::SeqCst) > 0, + "fixture must contend on the real log path" + ); + let mut runs = reads.run_ids(); + runs.sort(); + assert_eq!(runs, strings(&["t1-1-0", "t1-1-1"])); + assert_eq!(inner.rehydrate(), ctx); +} diff --git a/vendor/dex-loop/tests/scenarios.rs b/vendor/dex-loop/tests/scenarios.rs index e491213af..39bc6b963 100644 --- a/vendor/dex-loop/tests/scenarios.rs +++ b/vendor/dex-loop/tests/scenarios.rs @@ -257,13 +257,14 @@ async fn an_approval_class_call_is_granted_at_once_recorded_and_runs_in_order() tools.run_of(&call_id("t1", 1, 1)).args, json!({"key": "w", "to": "ops@example.com"}) ); - // Policy ran once per call: nothing resumed, so nothing re-checked. + // Policy ran once per call, plus once more for the read prefetched while + // the model streamed: adoption re-checks it under current authority. let checks: Vec = tools .policy_checks() .into_iter() .map(|(call, _)| call) .collect(); - assert_eq!(checks, strings(&["t1-1-0", "t1-1-1", "t1-1-2"])); + assert_eq!(checks, strings(&["t1-1-0", "t1-1-0", "t1-1-1", "t1-1-2"])); assert_eq!( view(&model.seen()[1])[2..], strings(&[ @@ -495,8 +496,10 @@ async fn steer_from_another_principal_is_checked_under_that_principal() { "user:check prod", "step:1", "started:t1-1-0", - "completed::[t1-1-0]", + // The read is prefetched while the model streams, so the steer it + // triggers lands before the model step completes. "steer:also update staging", + "completed::[t1-1-0]", "finished:t1-1-0:ok", "step:2", "completed::[t1-2-0]", @@ -519,6 +522,8 @@ async fn steer_from_another_principal_is_checked_under_that_principal() { assert_eq!( tools.policy_checks(), vec![ + // Prefetched while streaming, then re-checked at adoption. + ("t1-1-0".to_owned(), "alice".to_owned()), ("t1-1-0".to_owned(), "alice".to_owned()), ("t1-2-0".to_owned(), "bob".to_owned()), ] @@ -1002,6 +1007,7 @@ async fn usage_reported_before_a_failed_attempt_still_lands_with_the_abandon() { FakeModel::new(vec![vec![ usage(10, 5, 100), Err(ModelError { + class: dex_loop::ErrorClass::Unknown, message: "boom".into(), }), ]]), @@ -1027,6 +1033,7 @@ async fn a_stream_that_fails_after_text_keeps_the_answer_marked_cut_off() { text("nearly done"), usage(10, 5, 100), Err(ModelError { + class: dex_loop::ErrorClass::Unknown, message: "provider_stream_timeout: provider stream timed out".into(), }), ]]), diff --git a/vendor/dex-loop/tests/sim/fakes.rs b/vendor/dex-loop/tests/sim/fakes.rs index 65821e83a..7269ceb07 100644 --- a/vendor/dex-loop/tests/sim/fakes.rs +++ b/vendor/dex-loop/tests/sim/fakes.rs @@ -449,6 +449,7 @@ impl Model for SimModel { StepScript::Abandoned => vec![ Ok(ModelChunk::Text("about to fail".into())), Err(ModelError { + class: dex_loop::ErrorClass::Unknown, message: format!("upstream reset ({key})"), }), ], diff --git a/vendor/dex-loop/tests/support/mod.rs b/vendor/dex-loop/tests/support/mod.rs index e483e9fca..beb4a2218 100644 --- a/vendor/dex-loop/tests/support/mod.rs +++ b/vendor/dex-loop/tests/support/mod.rs @@ -331,6 +331,7 @@ impl Model for FakeModel { .push(tools.iter().map(|spec| spec.name.to_string()).collect()); let script = state.scripts.pop_front().unwrap_or_else(|| { vec![Err(ModelError { + class: dex_loop::ErrorClass::Unknown, message: "no script left".into(), })] }); @@ -785,7 +786,7 @@ pub fn shape(event: &Event) -> String { } => format!("client_tool_result:{call}:{}", outcome(*result)), Event::Compaction { summary, .. } => format!("compaction:{summary}"), Event::Final { text } => format!("final:{text}"), - Event::Error { code, message } => format!("error:{}:{message}", code.as_str()), + Event::Error { code, message, .. } => format!("error:{}:{message}", code.as_str()), Event::Interrupted => "interrupted".into(), } } diff --git a/vendor/dex-loop/tests/tool_deadline.rs b/vendor/dex-loop/tests/tool_deadline.rs index 84836f90f..2673723ff 100644 --- a/vendor/dex-loop/tests/tool_deadline.rs +++ b/vendor/dex-loop/tests/tool_deadline.rs @@ -17,14 +17,46 @@ #[allow(dead_code)] mod support; +use std::sync::{ + Arc, + atomic::{AtomicUsize, Ordering}, +}; use std::time::Duration; -use dex_loop::{Budget, CancellationToken, Exit, Outcome, Output}; +use dex_loop::{ + Budget, CallId, CancellationToken, Claim, Context, Cursor, Effects, Engine, ErrorCode, Event, + Exit, Fenced, Lexicon, Log, Message, Outcome, Output, ProposedCall, ToolName, ToolResult, +}; use serde_json::json; use support::*; const NEVER: Duration = Duration::from_secs(60 * 60); +fn assert_complete_tool_history(ctx: &Context, expected: &[CallId]) { + let proposed: Vec<_> = ctx + .history() + .iter() + .filter_map(|entry| match &entry.message { + Message::Assistant { calls, .. } => Some(calls.iter().map(|call| call.id.clone())), + _ => None, + }) + .flatten() + .collect(); + let completed: Vec<_> = ctx + .history() + .iter() + .filter_map(|entry| match &entry.message { + Message::Tool { call, .. } => Some(call.clone()), + _ => None, + }) + .collect(); + assert_eq!(proposed, expected, "all proposed calls remain in history"); + assert_eq!( + completed, expected, + "every call has exactly one result in history" + ); +} + fn budget_with_wall(wall: Duration) -> Budget { Budget { max_steps: 10, @@ -93,7 +125,9 @@ async fn a_mutation_that_overruns_the_deadline_is_recorded_unknown() { let tools = FakeTools::new(vec![write_tool("update")]).delay("slow", NEVER); let effects = FakeEffects::default(); let engine = engine_with(&log, &model, &tools, &effects, budget_with_wall(NEVER)) - .with_tool_call_deadline(Duration::from_millis(50)); + .with_tool_call_deadline(Duration::from_millis(50)) + // Composing a compactor must preserve the configured per-call bound. + .with_compactor(dex_loop::Threshold::new(usize::MAX, 1, FakeSummarizer)); let mut ctx = log.start_turn("t1", "change it"); let exit = tokio::time::timeout( @@ -179,6 +213,268 @@ async fn a_call_never_outlives_the_wall_budget() { assert_eq!(model.seen().len(), 1, "no model step after the wall budget"); } +#[tokio::test(flavor = "current_thread")] +async fn a_read_that_consumes_the_wall_does_not_start_the_next_mutation() { + let log = FakeLog::default(); + let model = FakeModel::new(vec![vec![ + call("search", json!({"key": "slow"})), + call("update", json!({"key": "fresh"})), + ]]); + let mutation_hits = Arc::new(AtomicUsize::new(0)); + let hits = mutation_hits.clone(); + let tools = FakeTools::new(vec![read_tool("search"), write_tool("update")]) + .delay("slow", NEVER) + .on_run(move |call| { + if call.tool.as_str() == "update" { + // Runs on the future's first poll, including the first poll + // of an otherwise expired Tokio timeout. + hits.fetch_add(1, Ordering::SeqCst); + } + }); + let effects = FakeEffects::default(); + let engine = engine_with( + &log, + &model, + &tools, + &effects, + budget_with_wall(Duration::from_millis(100)), + ) + .with_tool_call_deadline(NEVER); + let mut ctx = log.start_turn("t1", "read then change it"); + + assert_eq!( + engine.run(&mut ctx, &CancellationToken::new()).await, + Ok(Exit::Failed) + ); + assert_eq!(mutation_hits.load(Ordering::SeqCst), 0); + assert_eq!( + effects.recorded(&call_id("t1", 1, 1)), + None, + "fresh mutation must not even claim its effect" + ); + assert_eq!( + log.shapes_after(1), + strings(&[ + "step:1", + "started:t1-1-0", + "completed::[t1-1-0,t1-1-1]", + "finished:t1-1-0:err", + "finished:t1-1-1:err", + "error:budget_exhausted:wall budget exhausted: 100ms", + ]) + ); + assert_complete_tool_history(&ctx, &[call_id("t1", 1, 0), call_id("t1", 1, 1)]); + assert_eq!(ctx, log.rehydrate()); +} + +#[tokio::test(flavor = "current_thread")] +async fn a_read_that_consumes_the_wall_does_not_request_client_execution_or_approval() { + for read_only in [true, false] { + let log = FakeLog::default(); + let model = FakeModel::new(vec![vec![ + call("search", json!({"key": "slow"})), + call("browser.act", json!({})), + ]]); + let tools = FakeTools::new(vec![ + read_tool("search"), + client_executed_tool("browser.act", read_only), + ]) + .delay("slow", NEVER); + let engine = engine( + &log, + &model, + &tools, + budget_with_wall(Duration::from_millis(100)), + ) + .with_tool_call_deadline(NEVER); + let mut ctx = log.start_turn("t1", "read then use my browser"); + + assert_eq!( + engine.run(&mut ctx, &CancellationToken::new()).await, + Ok(Exit::Failed) + ); + assert!(log.events().iter().all(|event| !matches!( + event, + Event::ClientToolRequested { .. } | Event::ApprovalRequested { .. } + ))); + assert!(log.events().iter().any(|event| matches!(event, + Event::ToolFinished { call, outcome: Outcome::Failed, .. } + if call == &call_id("t1", 1, 1) + ))); + assert!(log.events().iter().any(|event| matches!( + event, + Event::Error { + code: ErrorCode::BudgetExhausted, + .. + } + ))); + assert_complete_tool_history(&ctx, &[call_id("t1", 1, 0), call_id("t1", 1, 1)]); + assert_eq!(ctx, log.rehydrate()); + } +} + +#[tokio::test(flavor = "current_thread")] +async fn exhausted_wall_still_adopts_an_already_started_mutations_ledger_result() { + let log = FakeLog::default(); + log.start_turn("t1", "resume the calls"); + let read = ProposedCall::new( + call_id("t1", 1, 0), + ToolName::new("search"), + json!({"key": "slow"}), + alice(), + ); + let write = ProposedCall::new( + call_id("t1", 1, 1), + ToolName::new("update"), + json!({"key": "existing"}), + alice(), + ); + log.host_append(Event::StepStarted { + step: 1, + control_through: Cursor::START, + }); + log.host_append(Event::ModelStepCompleted { + served: None, + step: 1, + text: String::new(), + calls: vec![read, write.clone()], + reasoning: None, + timing: None, + }); + log.host_append(Event::ToolStarted { + call: write.id.clone(), + tool: write.tool.clone(), + label: "Update".into(), + principal: write.principal.clone(), + }); + let recorded = ToolResult::unknown("executor lost the result; check the effect"); + let effects = FakeEffects::default().seed(write.id.clone(), Some(recorded.clone())); + let model = FakeModel::new(vec![]); + let tools = + FakeTools::new(vec![read_tool("search"), write_tool("update")]).delay("slow", NEVER); + let engine = engine_with( + &log, + &model, + &tools, + &effects, + budget_with_wall(Duration::from_millis(100)), + ) + .with_tool_call_deadline(NEVER); + let mut ctx = log.rehydrate(); + + assert_eq!( + engine.run(&mut ctx, &CancellationToken::new()).await, + Ok(Exit::Failed) + ); + assert!( + tools.runs().is_empty(), + "historical mutation must not execute again" + ); + assert_eq!(effects.recorded(&write.id), Some(Some(recorded.clone()))); + assert!(log.events().iter().any(|event| matches!(event, + Event::ToolFinished { call, outcome, output, .. } + if call == &write.id && outcome == &recorded.outcome && output == &recorded.output + ))); + assert!(log.events().iter().any(|event| matches!( + event, + Event::Error { + code: ErrorCode::BudgetExhausted, + .. + } + ))); + assert_complete_tool_history(&ctx, &[call_id("t1", 1, 0), call_id("t1", 1, 1)]); + assert_eq!(ctx, log.rehydrate()); +} + +#[tokio::test(flavor = "current_thread")] +async fn mutation_claim_or_start_persistence_consuming_wall_never_polls_the_effect() { + struct SlowClaim { + effects: FakeEffects, + delay: bool, + } + impl Effects for SlowClaim { + async fn claim(&self, call: &ProposedCall) -> Result { + if self.delay { + tokio::time::sleep(Duration::from_millis(150)).await; + } + self.effects.claim(call).await + } + async fn record(&self, call: &CallId, result: &ToolResult) -> Result<(), Fenced> { + self.effects.record(call, result).await + } + } + struct SlowStart { + log: FakeLog, + delay: bool, + } + impl Log for SlowStart { + async fn append(&self, events: &[Event]) -> Result, Fenced> { + if self.delay + && events + .iter() + .any(|event| matches!(event, Event::ToolStarted { .. })) + { + tokio::time::sleep(Duration::from_millis(150)).await; + } + self.log.append(events).await + } + async fn append_text(&self, text: String) -> Result<(), Fenced> { + self.log.append_text(text).await + } + async fn control_since(&self, after: Cursor) -> Result, Fenced> { + self.log.control_since(after).await + } + } + for delay_claim in [true, false] { + let log = FakeLog::default(); + let model = FakeModel::new(vec![vec![call("update", json!({"key": "fresh"}))]]); + let hits = Arc::new(AtomicUsize::new(0)); + let effect_hits = hits.clone(); + let tools = FakeTools::new(vec![write_tool("update")]).on_run(move |_| { + effect_hits.fetch_add(1, Ordering::SeqCst); + }); + let effects = FakeEffects::default(); + let engine = Engine::new( + SlowStart { + log: log.clone(), + delay: !delay_claim, + }, + model, + tools, + SlowClaim { + effects: effects.clone(), + delay: delay_claim, + }, + Lexicon::default(), + budget_with_wall(Duration::from_millis(100)), + ) + .with_tool_call_deadline(NEVER); + let mut ctx = log.start_turn("t1", "change it"); + + assert_eq!( + engine.run(&mut ctx, &CancellationToken::new()).await, + Ok(Exit::Failed) + ); + assert_eq!(hits.load(Ordering::SeqCst), 0); + assert_eq!( + log.events() + .iter() + .any(|event| matches!(event, Event::ToolStarted { .. })), + !delay_claim + ); + let recorded = effects + .recorded(&call_id("t1", 1, 0)) + .flatten() + .expect("claimed but undispatched result must be durable"); + assert_eq!(recorded.outcome, Outcome::Failed); + assert!( + matches!(recorded.output, Output::Text(ref text) if text.starts_with("not executed:")) + ); + assert_complete_tool_history(&ctx, &[call_id("t1", 1, 0)]); + assert_eq!(ctx, log.rehydrate()); + } +} + #[tokio::test(flavor = "current_thread", start_paused = true)] async fn a_call_that_finishes_in_time_is_unaffected() { let log = FakeLog::default(); @@ -209,3 +505,49 @@ async fn a_call_that_finishes_in_time_is_unaffected() { ]) ); } + +#[tokio::test(flavor = "current_thread")] +async fn compaction_time_consumes_the_tool_calls_remaining_wall_budget() { + struct SlowSummary; + impl dex_loop::Summarize for SlowSummary { + async fn summarize(&self, _entries: &[dex_loop::Entry]) -> Option { + tokio::time::sleep(Duration::from_millis(40)).await; + Some("historical summary".into()) + } + } + let log = FakeLog::default(); + log.start_turn("old", "old history to summarize"); + log.host_append(dex_loop::Event::Final { + text: "old done".into(), + }); + let mut ctx = log.start_turn("current", "read"); + let model = FakeModel::new(vec![ + vec![call("search", json!({"key": "slow"}))], + vec![text("never reached")], + ]); + let tools = FakeTools::new(vec![read_tool("search")]).delay("slow", Duration::from_millis(80)); + let engine = engine( + &log, + &model, + &tools, + budget_with_wall(Duration::from_millis(100)), + ) + .with_tool_call_deadline(NEVER) + .with_compactor(dex_loop::Threshold::new(1, 1, SlowSummary)); + assert_eq!( + engine.run(&mut ctx, &CancellationToken::new()).await, + Ok(Exit::Failed) + ); + assert!( + tools.runs().is_empty(), + "a read must not finish using a wall budget restarted after summarization" + ); + assert!(log.events().iter().all(|event| !matches!( + event, + dex_loop::Event::ToolFinished { + outcome: Outcome::Succeeded, + .. + } + ))); + assert_eq!(ctx, log.rehydrate()); +}