Oregami
Repositories/oxedyne/daimond

oxedyne/daimond/src/agent.rs

182 KiB, 1 run

created by r2519314175:933, which is this file's identity for as long as the history lasts, whatever it is later renamed to

download · who wrote it · its history

1//! Agent loop — the core Daimond agent that drives conversations.
2//!
3//! Receives a user message, sends it to the LLM with conversation
4//! history, streams the response back to the client via events,
5//! and stores the exchange in the session.
6
7use oxedyne_fe2o3_core::prelude::*;
8
9use std::cell::{Cell, RefCell};
10use std::rc::Rc;
11
12use crate::llm::{Delta, LlmClient};
13use crate::prompts::Role;
14use crate::protocol::{AgentEvent, ChatMessage, Dropped, MessageContent, Session};
15use crate::tools::ToolRegistry;
16
17/// Folding a conversation that has outgrown the model's context window.
18///
19/// Declared here, from its own file, rather than beside the other modules in `lib.rs`:
20/// compaction is part of running a turn and has no caller outside this one.
21#[path = "compact.rs"]
22pub mod compact;
23
24/// Which of a round's tool calls may run at the same time.
25///
26/// Declared here for the same reason `compact` is: deciding what a round dispatches together is
27/// part of running a turn, and nothing outside this module asks.
28#[path = "batch.rs"]
29pub mod batch;
30
31use crate::agent::compact::{Gauge, Limits};
32
33// The TLS client-config helper below is native-only; the wasm build
34// delegates TLS trust to the browser and constructs `LlmClient`
35// without a `ClientConfig`.
36#[cfg(not(target_arch = "wasm32"))]
37use std::sync::Arc;
38#[cfg(not(target_arch = "wasm32"))]
39use tokio_rustls::rustls::ClientConfig;
40
41
42// ┌───────────────────────────────────────────────────────────────┐
43// │ Speaking into a turn that is already running │
44// └───────────────────────────────────────────────────────────────┘
45
46/// Words the user has said since the turn began, waiting for a gap to be said in.
47///
48/// **What this is for.** A turn is not one request; it is a round of requests, each
49/// carrying the last round's tool results. A model twenty tool calls into the wrong
50/// approach cannot be corrected by waiting -- by the time the turn ends the work is
51/// done and the cost is spent -- and stopping it throws away everything it has
52/// learned on the way. So the user's correction is put where it can be acted on:
53/// into the message list, at the seam between one round and the next.
54///
55/// **What it is NOT.** Nothing can be added to a request that is already in flight;
56/// the message list is fixed the moment it is sent, and no provider offers otherwise.
57/// So a turn spending a minute writing prose with no tool call has no seam in it and
58/// cannot be interrupted -- for that, stopping is still the only answer. The seam
59/// exists because agentic work is made of many requests, not because a request can
60/// be reopened.
61///
62/// Shared by `Rc<RefCell<..>>` rather than a channel: the browser build is
63/// single-threaded, the queue is short, and the UI needs to READ it to draw what is
64/// waiting, which a consumed channel cannot offer.
65pub type Interjections = Rc<RefCell<Vec<String>>>;
66
67/// A fresh, empty interjection queue.
68pub fn new_interjections() -> Interjections {
69 Rc::new(RefCell::new(Vec::new()))
70}
71
72
73// ┌───────────────────────────────────────────────────────────────┐
74// │ Agent │
75// └───────────────────────────────────────────────────────────────┘
76
77/// The Daimond agent — drives a single conversation turn.
78///
79/// Holds a reference to the LLM client (shared across sessions) and
80/// the system prompt to prepend to every conversation.
81#[derive(Clone, Debug)]
82pub struct Agent {
83 pub llm: LlmClient,
84 pub system_prompt: String,
85 /// Cumulative prompt tokens for the turn in flight, updated round by round
86 /// alongside the session total. Held here, outside the session, so the
87 /// browser can read a running agent's spend without borrowing the session
88 /// the turn already holds mutably -- reading it there would panic the
89 /// `RefCell`, so a running tile could not show its cost.
90 pub live_prompt: Cell<u64>,
91 /// Cumulative completion tokens for the turn in flight; see `live_prompt`.
92 pub live_completion: Cell<u64>,
93 /// Cumulative cached prompt tokens for the turn in flight; see `live_prompt`.
94 pub live_cached: Cell<u64>,
95 /// Cumulative provider-reported USD for the turn in flight; see `live_prompt`.
96 pub live_cost: Cell<f64>,
97 /// What the user has said since this turn began (see [`Interjections`]).
98 pub interject: Interjections,
99 /// Facts about the machine this turn can reach, refreshed before each turn.
100 ///
101 /// Held behind a shared cell rather than folded into `system_prompt` because it changes
102 /// per TURN and the prompt does not: what a command may touch depends on the Diamond's
103 /// bounds and on whether this turn has read a stranger's words, neither of which is known
104 /// when the agent is built. Shared on clone, exactly as [`Interjections`] is, so a
105 /// derived agent sees the same machine -- and, unlike cloning the whole agent, this does
106 /// not silently detach the live token counters the panel is reading.
107 ///
108 /// Empty when no hand is attached, and then nothing is added to the prompt at all: an
109 /// absent capability is not worth describing on every request of every turn.
110 pub briefing: Rc<RefCell<String>>,
111 /// How long a turn may run and how big its conversation may get.
112 ///
113 /// Shared on clone, exactly as [`Interjections`] is, so a worker dispatched from a turn
114 /// inherits the model's context window rather than falling back to the default and
115 /// folding its own conversation at the wrong size.
116 pub limits: Rc<RefCell<Limits>>,
117 /// What a byte of this conversation costs in tokens, as the provider last charged it.
118 pub gauge: Rc<Gauge>,
119 /// The user's replacement for the compactor's prompt, from `prompts/compactor.md`;
120 /// empty means the shipped default.
121 ///
122 /// Held here rather than composed into `system_prompt` because it is not this
123 /// agent's prompt at all -- it is what a DIFFERENT, tool-less model is told when the
124 /// conversation is folded. Shared on clone for the same reason [`Interjections`] is.
125 pub fold_prompt: Rc<RefCell<String>>,
126 // How the last turn ended. Shared on clone rather than copied, exactly as the
127 // interjection queue is: the page holds a clone of the agent and has to be able to read
128 // the ending of the turn it just watched, which a detached cell would not carry. Written
129 // once per turn, at the single exit in `Agent::ended`.
130 ending: Rc<RefCell<Option<TurnEnding>>>,
131}
132
133
134// ┌───────────────────────────────────────────────────────────────┐
135// │ Why a fold is happening │
136// └───────────────────────────────────────────────────────────────┘
137
138/// What brought a fold about, which decides two things a boolean could not keep apart.
139///
140/// The distinction matters because [`Limits::learn_from_refusal`] moves the window DOWN
141/// and never up: it is the one occasion the provider speaks about its own size, and
142/// believing it is what lets a chat against an unpublished model recover. A fold the
143/// user asked for carries no such news -- the prompt was never refused -- so taking it as
144/// a refusal would shrink the window on every press of a button, and a user who folded
145/// three times would end up with a quarter of the context they started with.
146#[derive(Clone, Copy, Debug, PartialEq, Eq)]
147pub enum Fold {
148 /// Fold only if the estimate says the next prompt will not fit.
149 IfNeeded,
150 /// The provider refused this prompt, so fold whatever the estimate says -- and take
151 /// the refusal as the truth about the window.
152 Refused,
153 /// The user pressed Fold. Fold regardless of the estimate, and learn nothing about
154 /// the window from it.
155 ByHand,
156}
157
158impl Fold {
159 /// Should the estimate be overridden and the conversation folded anyway?
160 fn forces(self) -> bool {
161 !matches!(self, Self::IfNeeded)
162 }
163
164 /// Does this fold carry news about how big the window really is?
165 fn teaches_window(self) -> bool {
166 matches!(self, Self::Refused)
167 }
168}
169
170
171// ┌───────────────────────────────────────────────────────────────┐
172// │ How a turn ended, and whether its claims stood up │
173// └───────────────────────────────────────────────────────────────┘
174
175/// How a turn stopped running.
176///
177/// Three of these were already announced -- the output cap, the round limit and a turn that said
178/// nothing -- and every other ending was silence. A user watched a model say it would rewrite
179/// lines 43 to 49, watched the spinner clear, and was left with "I have no visibility on what
180/// occurred here": nothing had gone wrong, so nothing had been said.
181#[derive(Clone, Copy, Debug, Eq, PartialEq)]
182pub enum TurnEnd {
183 Answered, // a reply with no tool call in it, which is how a turn is meant to end
184 Stopped, // the user cancelled while the reply was streaming
185 Capped, // the tool-round budget ran out with work still going
186 Silent, // the final message carried no text at all
187 Failed, // the provider or the transport ended the turn
188}
189
190impl TurnEnd {
191
192 /// The word this ending travels as, on the wire and in a stored turn.
193 ///
194 /// Spelled here and nowhere else, for the reason [`crate::tools::CallOutcome::wire`] gives:
195 /// the browser knows these words and no others, so a second speller would not fail loudly,
196 /// it would quietly draw an ending nobody recognises.
197 pub fn wire(&self) -> &'static str {
198 match self {
199 Self::Answered => "answered",
200 Self::Stopped => "stopped",
201 Self::Capped => "capped",
202 Self::Silent => "silent",
203 Self::Failed => "failed",
204 }
205 }
206}
207
208/// What a turn came to, in figures the app measured rather than sentences the model wrote.
209///
210/// **Every field here is decidable from the tool log.** Nothing in it is read out of the model's
211/// prose, and nothing in it may be: this crate removed thirty-four prose sniffs on one night of
212/// 2026-08 and `dev/CONTRACT_OUTCOME.md` exists because four consumers were guessing a tool's
213/// outcome by reading its reply.
214#[derive(Clone, Debug, Eq, PartialEq)]
215pub struct TurnEnding {
216 pub how: TurnEnd,
217 pub offered: usize, // tools this turn was allowed to call
218 pub rounds: usize, // requests it sent
219 pub calls: usize, // tool calls it dispatched
220 pub refused: usize, // ... of which a door turned away
221 pub failed: usize, // ... of which broke
222 // Paths a completed call said it had left on the store, and which are not there.
223 pub missing: Vec<String>,
224}
225
226impl TurnEnding {
227
228 /// Is there anything here for a reader to act on?
229 ///
230 /// **This is the line between a status and a warning, and the reason the app can afford to
231 /// report every ending.** A turn that did what it was asked answers `false` and is drawn as
232 /// furniture; only a turn with a refusal, a breakage or a missing file answers `true`. An app
233 /// that appended a warning to every turn would teach its reader to skip the one that mattered,
234 /// which is the failure this whole mechanism exists to prevent.
235 pub fn unaccounted(&self) -> bool {
236 self.refused > 0 || self.failed > 0 || !self.missing.is_empty()
237 }
238}
239
240/// One dispatched tool call, as the audit sees it.
241#[derive(Clone, Debug)]
242struct Claim {
243 outcome: crate::tools::CallOutcome,
244 opaque: bool, // see [`crate::tools::Tool::opaque`]
245 paths: Vec<(String, crate::tools::PathClaim)>,
246}
247
248/// The turn's tool log, and the findings that can be read out of it.
249///
250/// Built as the turn runs and audited once at the end. It holds no prose and offers no way to
251/// look at any: a reader of this type cannot accidentally start sniffing sentences.
252#[derive(Clone, Debug, Default)]
253struct Claims {
254 calls: Vec<Claim>,
255}
256
257impl Claims {
258
259 /// Record one dispatched call.
260 ///
261 /// A call that was refused or that broke states nothing about the store, so its paths are not
262 /// taken: a write the fence stopped is not a file anybody should be looking for.
263 fn record(&mut self, name: &str, args: &str, outcome: crate::tools::CallOutcome) {
264 let tool = crate::tools::Tool::from_name(name);
265 let paths = match tool {
266 Some(t) if outcome == crate::tools::CallOutcome::Done => t.path_claims(args),
267 _ => Vec::new(),
268 };
269 // A refused shell ran nothing, so it cannot have moved a file. Anything else that
270 // reached a shell, a build or a worker did.
271 let opaque = tool.map(|t| t.opaque()).unwrap_or(false)
272 && outcome != crate::tools::CallOutcome::Refused;
273 self.calls.push(Claim { outcome, opaque, paths });
274 }
275
276 fn tally(&self, want: crate::tools::CallOutcome) -> usize {
277 self.calls.iter().filter(|c| c.outcome == want).count()
278 }
279
280 /// Can a missing file be blamed on anybody?
281 ///
282 /// False once the turn has run something whose reach is not in its arguments; see
283 /// [`crate::tools::Tool::opaque`] for why that silences the check rather than qualifying it.
284 fn accountable(&self) -> bool {
285 !self.calls.iter().any(|c| c.opaque)
286 }
287
288 /// The paths this turn's completed calls say are on the store now.
289 ///
290 /// Resolved in call order, so the turn's LAST word about a path is the turn's word about it: a
291 /// file written and then deleted is claimed by nobody, and one deleted and then written back
292 /// is claimed. Without that, every scratch file a turn tidied up after itself would be
293 /// reported as a write that did not happen.
294 fn standing(&self) -> Vec<String> {
295 let mut seen: Vec<(String, crate::tools::PathClaim)> = Vec::new();
296 for call in &self.calls {
297 for (path, claim) in &call.paths {
298 match seen.iter_mut().find(|(p, _)| p == path) {
299 Some(slot) => slot.1 = *claim,
300 None => seen.push((path.clone(), *claim)),
301 }
302 }
303 }
304 seen.into_iter()
305 .filter(|(_, c)| *c == crate::tools::PathClaim::Left)
306 .map(|(p, _)| p)
307 .collect()
308 }
309}
310
311impl Agent {
312
313 pub fn new(llm: LlmClient, system_prompt: &str) -> Self {
314 Self {
315 llm,
316 system_prompt: system_prompt.to_string(),
317 live_prompt: Cell::new(0),
318 live_completion: Cell::new(0),
319 live_cached: Cell::new(0),
320 live_cost: Cell::new(0.0),
321 interject: new_interjections(),
322 briefing: Rc::new(RefCell::new(String::new())),
323 limits: Rc::new(RefCell::new(Limits::default())),
324 gauge: Rc::new(Gauge::default()),
325 fold_prompt: Rc::new(RefCell::new(String::new())),
326 ending: Rc::new(RefCell::new(None)),
327 }
328 }
329
330 /// How the last turn on this agent ended.
331 ///
332 /// `None` before any turn has run. Read after `run_turn` returns -- including after it
333 /// returns an error, which is the ending that used to be the least visible of all.
334 pub fn ending(&self) -> Option<TurnEnding> {
335 self.ending.borrow().clone()
336 }
337
338 /// Settle the turn's ending, and say it.
339 ///
340 /// **The one exit.** Every path out of a turn comes through here, which is what makes the
341 /// promise -- every turn says how it ended -- checkable rather than aspirational: a new exit
342 /// that forgot to call this would leave `ending()` holding the turn before it, and the test
343 /// that reads the ending after each kind of finish is what catches that.
344 ///
345 /// # Arguments
346 /// * `ending` - What the turn came to; see [`TurnEnding`].
347 fn ended(&self, ending: TurnEnding, on_event: &mut impl FnMut(AgentEvent)) {
348 // THE ONE EXIT. Every path out of a turn comes through here, which is why the emit can
349 // be one line and why no second exit can grow beside it -- an ending reported from two
350 // places is an ending that can be reported from neither.
351 on_event(AgentEvent::Ended {
352 how: ending.how.wire().to_string(),
353 offered: ending.offered,
354 rounds: ending.rounds,
355 calls: ending.calls,
356 refused: ending.refused,
357 failed: ending.failed,
358 missing: ending.missing.clone(),
359 });
360 *self.ending.borrow_mut() = Some(ending);
361 }
362
363 /// Check the turn's claims against the store, and settle what it came to.
364 ///
365 /// The audit is two questions, and neither of them is asked of the model's prose:
366 ///
367 /// 1. **A refused call is not a completed step.** The tool layer already decided this and
368 /// `AgentEvent::ToolResult` already carries it, so the count is a tally rather than a
369 /// reading.
370 /// 2. **A file a completed call said it left is on the store.** The claim is in the call's
371 /// ARGUMENTS, which name a path; whether that path is there afterwards is a fact.
372 ///
373 /// # Arguments
374 /// * `how` - How the turn stopped running.
375 /// * `rounds` - Requests the turn sent.
376 /// * `claims` - The turn's tool log; empty on the pure-chat path.
377 /// * `registry` - The tools this turn held, which is also what answers for the store.
378 async fn audit(
379 &self,
380 how: TurnEnd,
381 rounds: usize,
382 claims: &Claims,
383 registry: Option<&ToolRegistry>,
384 )
385 -> TurnEnding
386 {
387 let mut missing = Vec::new();
388 if let Some(reg) = registry {
389 if claims.accountable() {
390 for path in claims.standing() {
391 if !reg.path_is_there(&path).await {
392 missing.push(path);
393 }
394 }
395 }
396 }
397 TurnEnding {
398 how,
399 offered: registry.map(|r| r.offered().len()).unwrap_or(0),
400 rounds,
401 calls: claims.calls.len(),
402 refused: claims.tally(crate::tools::CallOutcome::Refused),
403 failed: claims.tally(crate::tools::CallOutcome::Failed),
404 missing,
405 }
406 }
407
408 /// Tell this agent how big the model's context window is, so a conversation is folded
409 /// before the provider refuses it rather than after.
410 ///
411 /// Zero means nobody has said, and the default window is assumed instead; the reactive
412 /// path still catches a refusal either way.
413 ///
414 /// # Arguments
415 /// * `tokens` - The window the provider publishes for this model.
416 pub fn set_context_window(&self, tokens: u64) {
417 self.limits.borrow_mut().window = tokens;
418 }
419
420 /// Set how many tool-call rounds one turn may take.
421 ///
422 /// # Arguments
423 /// * `n` - The ceiling; zero is ignored, since a turn that may take no rounds is a turn
424 /// with no tools.
425 pub fn set_max_rounds(&self, n: usize) {
426 if n > 0 {
427 self.limits.borrow_mut().max_rounds = n;
428 }
429 }
430
431 /// Set the fraction of the window at which this agent folds.
432 ///
433 /// Held inside [`compact::FOLD_AT_MIN`]..[`compact::FOLD_AT_MAX`] here rather than only in
434 /// [`compact::Limits::budget`], so `Agent::limits` reports the figure that is actually in
435 /// force -- a control that draws itself from the getter would otherwise show a number the
436 /// arithmetic never used. Zero, or anything below it, leaves the default alone, which is
437 /// how a caller says "the user has not chosen".
438 ///
439 /// # Arguments
440 /// * `f` - The fraction, between 0 and 1; zero or less is ignored.
441 pub fn set_fold_at(&self, f: f64) {
442 if f > 0.0 {
443 self.limits.borrow_mut().fold_at = f.clamp(compact::FOLD_AT_MIN, compact::FOLD_AT_MAX);
444 }
445 }
446
447 /// Fold this agent's conversations with a different model from the one it chats with.
448 ///
449 /// Empty -- the default -- means the chat's own model. A summary becomes the session's
450 /// memory, so the cheaper model is not chosen on the user's behalf: they choose it.
451 ///
452 /// # Arguments
453 /// * `model` - The provider's id for the model, or empty for the chat's own.
454 pub fn set_fold_model(&self, model: &str) {
455 self.limits.borrow_mut().fold_model = model.to_string();
456 }
457
458 /// What bounds this agent's turns right now.
459 pub fn limits(&self) -> Limits {
460 self.limits.borrow().clone()
461 }
462
463 /// Fold by the same figures as another agent.
464 ///
465 /// [`Agent::new`] starts from [`Limits::default`], which assumes a window nobody has
466 /// published and the shipped round ceiling. A Diamond's daimon and its reducer are
467 /// each built that way, from the chat's own client -- so without this they would fold
468 /// the SAME model's conversation at a different size from the chat, and on a model
469 /// whose real window is smaller than the assumed one they would learn that the hard
470 /// way all over again. Cloning an agent shares these already; this is for the case
471 /// where a fresh one is constructed because its system prompt differs.
472 ///
473 /// # Arguments
474 /// * `from` - The agent whose limits are the right ones.
475 pub fn adopt_limits(&self, from: &Agent) {
476 *self.limits.borrow_mut() = from.limits();
477 // The fold PROMPT travels with them, for the same reason and by the same argument.
478 // It is not part of `Limits` because it is text rather than a figure, but it is the
479 // same setting: what the folding model is told. Without this line a Diamond's
480 // daimon and its reducer folded on the user's chosen model -- `fold_model` rides in
481 // `Limits` -- while ignoring the instructions the user wrote for it in
482 // `prompts/compactor.md`, which is the half of the setting that is visible on disk.
483 *self.fold_prompt.borrow_mut() = from.fold_prompt.borrow().clone();
484 }
485
486 /// Fold this conversation because the user asked, not because it had to be folded.
487 ///
488 /// The same path a turn takes when the estimate says the next prompt will not fit --
489 /// there is deliberately no second folding routine -- but entered with [`Fold::ByHand`],
490 /// so nothing is learned about the window from a prompt the provider never saw.
491 ///
492 /// Returns whether anything actually moved. It can be false: a conversation of six
493 /// messages or fewer has no tail to cut and no bulk to elide, and saying so is better
494 /// than a spinner that ends with the meter where it was.
495 ///
496 /// # Arguments
497 /// * `session` - The durable conversation, folded in place.
498 /// * `on_event` - Where the fold is announced, exactly as an automatic one is.
499 pub async fn fold_by_hand(
500 &self,
501 session: &mut Session,
502 on_event: &mut impl FnMut(AgentEvent),
503 )
504 -> bool
505 {
506 // The working list a turn would build: the system prompt, then the conversation.
507 // Rebuilt here rather than borrowed because there is no turn in flight -- this is
508 // the user at rest, between turns, pressing a button.
509 let mut working = vec![ChatMessage::system(self.system_prompt.clone())];
510 working.extend(session.messages.iter().cloned());
511 self.fold_if_needed(session, &mut working, 0, Fold::ByHand, on_event).await
512 }
513
514 /// Set what the model folding this agent's conversations is told.
515 ///
516 /// Empty -- the default -- means [`Role::Compactor`]'s shipped prompt, which is how
517 /// deleting `prompts/compactor.md` puts the original back.
518 ///
519 /// # Arguments
520 /// * `text` - What the user wrote, or empty for the default.
521 pub fn set_fold_prompt(&self, text: &str) {
522 *self.fold_prompt.borrow_mut() = text.to_string();
523 }
524
525 /// What the folding model is told, composed from the user's text or the default.
526 pub fn fold_prompt(&self) -> String {
527 Role::Compactor.compose(&self.fold_prompt.borrow())
528 }
529
530 /// Set what this turn should know about the machine it can reach.
531 ///
532 /// Called before a turn, from the caller that can await the hand. Empty clears it.
533 ///
534 /// # Arguments
535 /// * `text` - The briefing, already composed.
536 pub fn set_briefing(&self, text: &str) {
537 *self.briefing.borrow_mut() = text.to_string();
538 }
539
540 /// Say something into a turn that is already running.
541 ///
542 /// Takes effect at the next seam between rounds, which is the earliest moment a
543 /// model can act on it. Returns how many are now waiting, so the caller can draw
544 /// them without reaching into the queue itself.
545 ///
546 /// # Arguments
547 /// * `text` - What the user said. Blank input is ignored rather than queued.
548 pub fn interject(&self, text: &str) -> usize {
549 let t = text.trim();
550 if t.is_empty() {
551 return self.interject.borrow().len();
552 }
553 let mut q = self.interject.borrow_mut();
554 q.push(t.to_string());
555 q.len()
556 }
557
558 /// What is waiting to be said, for the UI to draw.
559 pub fn interjections(&self) -> Vec<String> {
560 self.interject.borrow().clone()
561 }
562
563 /// Take everything waiting, leaving the queue empty.
564 ///
565 /// Drained rather than read so that a message cannot be delivered twice: it is
566 /// pushed into the conversation the moment it is taken, and the conversation is
567 /// the record from then on.
568 fn take_interjections(&self) -> Vec<String> {
569 let mut q = self.interject.borrow_mut();
570 if q.is_empty() { return Vec::new(); }
571 std::mem::take(&mut *q)
572 }
573
574 /// Run a single agent turn.
575 ///
576 /// 1. Append the user message to the session.
577 /// 2. Build the LLM request: system prompt + conversation history.
578 /// 3. Call the LLM with streaming.
579 /// 4. Stream tokens back to the caller via `on_event`.
580 /// 5. Append the assistant response to the session.
581 /// 6. Emit `Done`.
582
583 /// The three pieces the system message is built from, in the order they are joined.
584 ///
585 /// **Read by the Wire view, and by nothing else that decides anything.** It exists so a person
586 /// can see what is actually sent -- which of it is their own role prompt, which was appended
587 /// after their edits, which is derived from the fence -- and the only way that view can be
588 /// trusted is if it is composed by the same code the request is. So `run_turn` calls this and
589 /// so does the getter; there is no second assembly to drift from the first.
590 ///
591 /// # Arguments
592 /// * `registry` - The tools this turn holds, which decide the middle sentence.
593 pub fn system_parts(&self, registry: &ToolRegistry) -> (String, String, String) {
594 let tools = if registry.is_empty() {
595 String::new()
596 } else {
597 fmt!(
598 "You have exactly these tools, all scoped to the user's \
599 workspace: {}. Use them to inspect and change the workspace \
600 when completing a task. You have no other tools; never claim \
601 to have performed an action you had no tool to perform.",
602 registry.tool_names().join(", "))
603 };
604 let brief = self.briefing.borrow().trim().to_string();
605 (self.system_prompt.clone(), tools, brief)
606 }
607
608
609 pub async fn run_turn(
610 &self,
611 session: &mut Session,
612 user_msg: String,
613 registry: &ToolRegistry,
614 on_event: &mut impl FnMut(AgentEvent),
615 ) -> Outcome<()> {
616 // THE TURN'S BYTE LEDGER STARTS HERE, and here is the only place it does. Every turn in
617 // the app arrives through this function -- the browser chat, a Diamond's daimon, a
618 // dispatched worker and `examples/devcycle_probe.rs` alike -- and a Diamond's daimon
619 // SHARES its `read_seen` with the chat that made it, so an allowance reset anywhere else
620 // would leak from one turn into the next. See `crate::tools::TurnState::spent`.
621 registry.ctx.begin_turn();
622 // Append the user message to the persisted history.
623 session.messages.push(ChatMessage::user(user_msg));
624
625 // Build the working conversation: system prompt + history.
626 let mut working = Vec::with_capacity(session.messages.len() + 1);
627 if !self.system_prompt.is_empty() {
628 // ONE COMPOSER, read here and by the Wire view. The sentence naming the tools is
629 // derived from the registry because a fixed one once promised a shell tool the
630 // browser build has not got, so a capable model called it, failed, and reported the
631 // failure as work done. The machine note goes LAST, so it sits closest to the
632 // conversation and is the most recent thing the model read before the user's words.
633 let (mut sys, tools, brief) = self.system_parts(registry);
634 if !tools.is_empty() {
635 sys.push_str("\n\n");
636 sys.push_str(&tools);
637 }
638 if !brief.is_empty() {
639 sys.push_str("\n\n");
640 sys.push_str(&brief);
641 }
642 working.push(ChatMessage::system(sys));
643 }
644 working.extend(session.messages.iter().cloned());
645
646 if registry.is_empty() {
647 return self.run_streaming(session, working, on_event).await;
648 }
649 self.run_tool_loop(session, working, registry, on_event).await
650 }
651
652 /// Pure-chat path: stream tokens as they arrive (no tools).
653 ///
654 /// Folded before the request goes out, and folded again if the provider refuses it for
655 /// being too long -- the second is the net under the first. Without it a chat whose
656 /// window nobody published dies once and then dies on every turn after, because the
657 /// same oversized history is sent again.
658 async fn run_streaming(
659 &self,
660 session: &mut Session,
661 mut working: Vec<ChatMessage>,
662 on_event: &mut impl FnMut(AgentEvent),
663 ) -> Outcome<()> {
664 let mut refolded = false;
665 let mut rounds = 0usize;
666 loop {
667 self.fold_if_needed(session, &mut working, 0, Fold::IfNeeded, on_event).await;
668 let sent = compact::conversation_bytes(&working, &self.llm.open_folds());
669 let mut full = String::new();
670 rounds += 1;
671 let result = self.llm.chat_stream(
672 &working,
673 &mut |d| match d {
674 Delta::Text(token) => {
675 full.push_str(token);
676 on_event(AgentEvent::Text(token.to_string()));
677 }
678 // THE WORKING GOES OUT WHILE IT IS STILL BEING DONE, which is the whole
679 // point of reading it: a reasoning model spends most of a round thinking,
680 // and a page that waits for the round to end shows a blank spinner for all
681 // of it. Its own event, never `Text`, so `full` -- which becomes the
682 // assistant's message -- cannot pick it up.
683 Delta::Reasoning(think) => on_event(AgentEvent::Thinking(think.to_string())),
684 },
685 ).await;
686 match result {
687 Ok(resp) => {
688 // The working is NOT emitted here. It went out delta by delta while the
689 // round was running (the `Delta::Reasoning` arm above), which is the only
690 // way it can do the job it is drawn for: a reasoning model spends most of
691 // a round thinking, and working delivered after the round is over arrives
692 // at the one moment it no longer explains the wait. `resp.thinking` still
693 // carries the whole of it for a caller that wants it in one piece.
694 let content = if resp.content.is_empty() { full } else { resp.content };
695 // A TURN MUST NEVER END IN SILENCE. Both empty means the provider
696 // returned a final message with nothing in it -- which happens on a
697 // reasoning model whose answer went entirely to a channel this app
698 // does not print, and happened to a user after the model had read
699 // five files and then appeared to do nothing at all. The spinner
700 // clears, the screen does not change, and there is no way to tell a
701 // finished turn from a hung one. Whatever the cause, saying so is
702 // strictly better than saying nothing.
703 let silent = content.trim().is_empty();
704 session.messages.push(ChatMessage::assistant(content));
705 if silent {
706 on_event(AgentEvent::Error(
707 "The model ended its turn without saying anything.".to_string()));
708 session.messages.push(compact::empty_turn_note());
709 }
710 session.prompt_tokens += resp.prompt_tokens;
711 session.completion_tokens += resp.completion_tokens;
712 session.cached_tokens += resp.cached_tokens;
713 session.cost_usd += resp.cost_usd;
714 if resp.prompt_tokens > 0 { session.last_prompt_tokens = resp.prompt_tokens; }
715 if resp.truncated { on_event(AgentEvent::Truncated); }
716 self.gauge.observe(sent, resp.prompt_tokens);
717 self.live_prompt.set(session.prompt_tokens);
718 self.live_completion.set(session.completion_tokens);
719 self.live_cached.set(session.cached_tokens);
720 self.live_cost.set(session.cost_usd);
721 // A pure chat holds no tools, so it can claim nothing about the store and
722 // its ending is the shape of the turn and nothing else.
723 let how = if silent { TurnEnd::Silent } else { TurnEnd::Answered };
724 let ending = self.audit(how, rounds, &Claims::default(), None).await;
725 self.ended(ending, on_event);
726 on_event(AgentEvent::Done);
727 return Ok(());
728 }
729 Err(e) => {
730 if !refolded && self.overflowed(&e, sent) {
731 refolded = true;
732 if self.fold_if_needed(session, &mut working, 0, Fold::Refused, on_event).await {
733 continue;
734 }
735 }
736 on_event(AgentEvent::Error(e.to_string()));
737 let ending = self.audit(TurnEnd::Failed, rounds, &Claims::default(), None).await;
738 self.ended(ending, on_event);
739 return Err(e);
740 }
741 }
742 }
743 }
744
745 /// Agentic path: streaming request/response, executing tool calls
746 /// between rounds until the model returns a final answer. Each round
747 /// streams with tools enabled, so assistant text arrives token by
748 /// token even while tools are active (via `chat_stream_tools`); tool
749 /// calls are reconstructed from the streamed fragments and fired as
750 /// before. The whole exchange -- the assistant turn that asked for the
751 /// tools, each tool result, and the final answer -- is persisted to the
752 /// session, so a later turn still sees what this agent did. Persisting
753 /// only the final text once left the model amnesiac: asked a follow-up, it
754 /// had no record of its own tool calls and could not say what it had done.
755 /// Run one tool call, retrying it while the failure is the ROAD rather than an answer.
756 ///
757 /// **THE LADDER GUARDS THE CALL TO THE MODEL; THIS GUARDS THE CALLS THE MODEL MAKES.** The
758 /// provider retry in [`crate::llm::LlmClient::stream_turn`] has always survived a laptop
759 /// waking or a phone coming back, and that was read as the problem being solved. It was not.
760 /// A `web_fetch` that dies because the page was frozen is handed back to the model AS A TOOL
761 /// RESULT SAYING IT FAILED, and the model then does the reasonable thing with a failed tool:
762 /// it apologises and answers around it. The owner met that on a real iPhone on 2026-08-28 --
763 /// "I can't get through to the web right now to look this up" -- and the sentence is now a
764 /// permanent turn in the conversation.
765 ///
766 /// **`self.llm.retry`, not a schedule of its own.** The same eight attempts, the same
767 /// jittered backoff, the same total bound, and the same object -- so a test client's fast
768 /// policy governs this ladder too and the suite does not grow two minutes per road failure.
769 ///
770 /// Only a READ is climbed: see [`crate::tools::Tool::road_retryable`]. Anything else fails
771 /// on the first attempt, which still means the model is told nothing -- it simply means the
772 /// turn ends sooner rather than after eight tries.
773 ///
774 /// # Arguments
775 /// * `name` - The tool's wire name, as the model spelled it.
776 /// * `args` - The raw argument object.
777 async fn over_the_road(
778 &self,
779 registry: &ToolRegistry,
780 name: &str,
781 args: &str,
782 on_event: &mut impl FnMut(AgentEvent),
783 )
784 -> Outcome<MessageContent>
785 {
786 let again = crate::tools::Tool::from_name(name)
787 .map(|t| t.road_retryable())
788 .unwrap_or(false);
789 let mut waited = 0u64;
790 let mut retries = 0u32;
791 loop {
792 match registry.try_dispatch(name, args).await {
793 Ok(c) => return Ok(c),
794 Err(e) => {
795 if !again {
796 return Err(e);
797 }
798 let delay = match self.llm.retry.next_delay(retries, waited, None) {
799 Some(d) => d,
800 None => return Err(e),
801 };
802 waited += delay;
803 retries += 1;
804 // SAID WHILE IT HAPPENS, and said by the APP. A tool call quietly retrying
805 // for up to two minutes looks exactly like a hung turn, which is the one
806 // thing a spinner cannot tell a user. It is a `Roading` event and not
807 // assistant text for the reason a fold is not assistant text: this is
808 // something Daimond is doing, the model neither said it nor will read it.
809 on_event(AgentEvent::Roading {
810 name: name.to_string(),
811 attempt: retries + 1,
812 of: self.llm.retry.max_attempts,
813 wait_ms: delay,
814 });
815 crate::llm::sleep_ms(delay).await;
816 // AND THEN WAIT FOR THE PAGE, which is the half a backoff cannot do. On iOS
817 // the page is frozen for as long as the user is in another app, so a ladder
818 // that only sleeps spends all eight attempts into a dead page and reports
819 // failure the moment the user comes back. Charged to the same budget, so
820 // being backgrounded cannot extend a turn without bound.
821 #[cfg(target_arch = "wasm32")]
822 {
823 waited += crate::wasm::await_restored(waited).await;
824 }
825 }
826 }
827 }
828 }
829
830 /// The event that tells the page a call has started.
831 ///
832 /// The ID travels with it. A stored conversation's `say` fold is opened and closed on the
833 /// page, and the page has to be able to name WHICH call it is talking about when it tells the
834 /// engine what is open -- see `set_open_folds`. Nothing else reads it.
835 ///
836 /// One function because a batch announces its first call before it runs and the rest as their
837 /// results are recorded, and two spellings of the same event would eventually differ.
838 fn announce(tc: &crate::protocol::ToolCall) -> AgentEvent {
839 AgentEvent::ToolCall {
840 id: tc.id.clone(),
841 name: tc.name.clone(),
842 args: tc.arguments.clone(),
843 }
844 }
845
846 /// One tool call, from the model's arguments to a result the conversation can carry.
847 ///
848 /// # Arguments
849 /// * `truncated` - Whether the reply these calls were parsed out of hit the output limit.
850 async fn one_call(
851 &self,
852 registry: &ToolRegistry,
853 tc: &crate::protocol::ToolCall,
854 truncated: bool,
855 on_event: &mut impl FnMut(AgentEvent),
856 )
857 -> Outcome<MessageContent>
858 {
859 // A call cut at the output limit is not a call. Its arguments are a JSON object that
860 // stops in the middle, and dispatching it yields a parse error -- which reads to the
861 // model as its own mistake, so it writes the same thing again and is cut in the same
862 // place. Told what actually happened, it splits the work instead.
863 if cut_short(truncated, &tc.arguments) {
864 // The cap that was SENT, not the one configured. The note says the number so the
865 // model can size its next call by it, and telling a model its reply was cut at 4,096
866 // when it was cut at 32,000 sends it splitting the work eight times finer than it
867 // needs to.
868 return Ok(MessageContent::text(truncated_call_note(self.reply_cap())));
869 }
870 self.over_the_road(registry, &tc.name, &tc.arguments, on_event).await
871 }
872
873 /// [`Self::one_call`] with its events put in a buffer instead of on the wire.
874 ///
875 /// The one thing a batch cannot share is the `&mut` event closure, so each call in a batch
876 /// gets a `Vec` of its own and the round drains them in the model's order. A method rather
877 /// than a closure at the call site because the borrow of the sink has to outlive the future,
878 /// which a closure built inside a `map` cannot arrange.
879 async fn call_buffered(
880 &self,
881 registry: &ToolRegistry,
882 tc: &crate::protocol::ToolCall,
883 truncated: bool,
884 sink: &mut Vec<AgentEvent>,
885 )
886 -> Outcome<MessageContent>
887 {
888 let mut into = |e: AgentEvent| sink.push(e);
889 self.one_call(registry, tc, truncated, &mut into).await
890 }
891
892 /// Leave a round whose tool call died on the road, with nothing false left behind.
893 ///
894 /// The assistant turn that asked for the tools is already in the session -- it has to be,
895 /// because the API requires it to precede the replies -- and abandoning the round leaves its
896 /// unanswered calls dangling. A conversation in that state is rejected WHOLE by every
897 /// provider, so it cannot simply be left.
898 ///
899 /// **[`crate::protocol::pair_up`], and deliberately not a rollback.** Dropping the whole
900 /// assistant turn would also discard the prose the model wrote before deciding to call the
901 /// tool -- text the user has already read and already paid for, and exactly what
902 /// `continueTurn` (www/js/daimond.js) works to preserve when a PROVIDER call dies. A road
903 /// failure during a tool call and a road failure during a provider call would then leave
904 /// different things behind, which reads as a bug to whoever meets it. `pair_up` keeps the
905 /// prose and the answered calls and drops only the unanswered ones, which is what the
906 /// restore path has always done to the same conversation coming back out of the store.
907 ///
908 /// **It is needed for the SITTING, not for the reload.** Press Continue, or reload, and
909 /// `restore_session` runs `pair_up` anyway. Ignore the badge and simply type again, and the
910 /// live session is the one holding the dangling call -- which is the case nothing else
911 /// covers.
912 ///
913 /// # Arguments
914 /// * `working` - This turn's own message list, repaired alongside the session's.
915 fn abandon_round(
916 &self,
917 session: &mut Session,
918 working: &mut Vec<ChatMessage>,
919 ) {
920 let msgs = std::mem::take(&mut session.messages);
921 session.messages = crate::protocol::pair_up(msgs);
922 let work = std::mem::take(working);
923 *working = crate::protocol::pair_up(work);
924 }
925
926 async fn run_tool_loop(
927 &self,
928 session: &mut Session,
929 mut working: Vec<ChatMessage>,
930 registry: &ToolRegistry,
931 on_event: &mut impl FnMut(AgentEvent),
932 ) -> Outcome<()> {
933 let tools_json = registry.definitions_json();
934 // The tool schema rides on every request and is not in the conversation, so it is
935 // counted separately -- a budget blind to it is short by however many tools are
936 // registered, which for the browser set is several thousand tokens.
937 let schema = tools_json.as_ref().map(|s| s.len() as u64).unwrap_or(0);
938 let max_rounds = self.limits.borrow().max_rounds;
939 // THE TURN'S TOOL LOG, and the only thing the audit reads. It carries an outcome and a
940 // set of argument-named paths per call, and no prose at all -- so no reader of it can
941 // slip back into working out what happened by reading what the model said about it.
942 let mut claims = Claims::default();
943 let mut rounds = 0usize;
944 for _ in 0..max_rounds {
945 rounds += 1;
946 // Fold before the request rather than after the refusal. Checked every round,
947 // because a single turn of fifty file reads can outgrow the window on its own,
948 // without any earlier turn being large at all.
949 self.fold_if_needed(session, &mut working, schema, Fold::IfNeeded, on_event).await;
950 let mut sent = compact::conversation_bytes(&working, &self.llm.open_folds());
951
952 // Stream this round's assistant text as it arrives; the tool
953 // calls (if any) are returned assembled once the round ends.
954 let mut refolded = false;
955 // What this endpoint would take BEFORE the request goes out. Only a refusal can
956 // change it and it changes it exactly once, so sampling either side of the round is
957 // the whole of the learned signal -- no counter and no flag of our own.
958 let could_see = self.llm.can_take_images();
959 let resp = loop {
960 match self.llm.chat_stream_tools(
961 &working,
962 tools_json.as_deref(),
963 &mut |d| match d {
964 Delta::Text(token) => on_event(AgentEvent::Text(token.to_string())),
965 Delta::Reasoning(t) => on_event(AgentEvent::Thinking(t.to_string())),
966 },
967 ).await {
968 Ok(r) => break r,
969 Err(e) => {
970 // The net under the proactive fold: a window nobody published, or
971 // an estimate that ran short. One retry, and only for a refusal
972 // that looks like the prompt being too long.
973 if !refolded && self.overflowed(&e, sent + schema) {
974 refolded = true;
975 if self.fold_if_needed(session, &mut working, schema, Fold::Refused, on_event).await {
976 sent = compact::conversation_bytes(&working, &self.llm.open_folds());
977 continue;
978 }
979 }
980 on_event(AgentEvent::Error(e.to_string()));
981 let ending = self.audit(
982 TurnEnd::Failed, rounds, &claims, Some(registry)).await;
983 self.ended(ending, on_event);
984 return Err(e);
985 }
986 }
987 };
988 // The working is NOT emitted here either; see the note in `run_streaming`. It
989 // streamed while the round ran, and a tool loop is where that matters most --
990 // this is the path that runs many rounds, each of which may think for a minute
991 // before it says a word.
992 // The endpoint has just been caught refusing pictures, mid-turn, having taken
993 // them a moment ago. Said out loud rather than left in the client: it is the one
994 // moment the app knows it is on the wrong model, and `stream_turn` has already
995 // stripped the pictures and answered anyway, so nothing downstream would ever
996 // find out.
997 if could_see && !self.llm.can_take_images() {
998 let images: usize = working.iter().map(|m| m.content().images().count()).sum();
999 on_event(AgentEvent::Unseeable { images, model: self.llm.model.clone() });
1000 }
1001 self.gauge.observe(sent + schema, resp.prompt_tokens);
1002 session.prompt_tokens += resp.prompt_tokens;
1003 session.completion_tokens += resp.completion_tokens;
1004 // Both are per ROUND, and a tool loop runs many: each round's prompt
1005 // is the last one's plus a little, so the cache read and the reported
1006 // cost accumulate exactly as the token counters do.
1007 session.cached_tokens += resp.cached_tokens;
1008 session.cost_usd += resp.cost_usd;
1009 if resp.prompt_tokens > 0 { session.last_prompt_tokens = resp.prompt_tokens; }
1010 self.live_prompt.set(session.prompt_tokens);
1011 self.live_completion.set(session.completion_tokens);
1012 self.live_cached.set(session.cached_tokens);
1013 self.live_cost.set(session.cost_usd);
1014
1015 // Said outright rather than left to be inferred. It is not an error and is
1016 // not retried: the request succeeded and a setting was reached, so sending
1017 // it again would cost the same money and produce the same cut.
1018 if resp.truncated {
1019 on_event(AgentEvent::Truncated);
1020 }
1021
1022 // Cancelled mid-stream: keep the partial answer already
1023 // streamed and end the turn cleanly, without an error.
1024 if resp.aborted {
1025 session.messages.push(ChatMessage::Assistant {
1026 content: MessageContent::text(crate::llm::seamed(resp.content)),
1027 tool_calls: Vec::new(),
1028 });
1029 let ending = self.audit(TurnEnd::Stopped, rounds, &claims, Some(registry)).await;
1030 self.ended(ending, on_event);
1031 on_event(AgentEvent::Done);
1032 return Ok(());
1033 }
1034
1035 if resp.tool_calls.is_empty() {
1036 // Final answer — its text has already streamed via the
1037 // token callback, so it is not re-emitted here.
1038 //
1039 // A final message with nothing in it is the ending the user met: the spinner
1040 // clears, the screen does not change, and a finished turn is indistinguishable
1041 // from a hung one. The streaming path has said so since it was found; here the
1042 // ending names it, which costs nothing and is the same fact.
1043 let how = if resp.content.trim().is_empty() {
1044 TurnEnd::Silent
1045 } else {
1046 TurnEnd::Answered
1047 };
1048 // THE SEAM, and only here. A run of prose with a tool call after it is
1049 // working rather than an answer -- `demoteToWorking` in the page draws it as
1050 // the model's own thinking -- so a `Fold:` line in one of those would build a
1051 // control over something nobody is meant to read as a reply.
1052 session.messages.push(ChatMessage::Assistant {
1053 content: MessageContent::text(crate::llm::seamed(resp.content)),
1054 tool_calls: Vec::new(),
1055 });
1056 let ending = self.audit(how, rounds, &claims, Some(registry)).await;
1057 self.ended(ending, on_event);
1058 on_event(AgentEvent::Done);
1059 return Ok(());
1060 }
1061
1062 // Any interim assistant text alongside the tool calls has already
1063 // streamed. Record the assistant turn in both the working vec, which
1064 // drives the rest of this turn, and the session, so a later turn
1065 // still carries it. The API also requires that an assistant turn
1066 // bearing tool_calls be followed by a tool reply for each of them,
1067 // which the loop below then supplies.
1068 let asked = ChatMessage::Assistant {
1069 content: MessageContent::text(resp.content.clone()),
1070 tool_calls: resp.tool_calls.clone(),
1071 };
1072 working.push(asked.clone());
1073 session.messages.push(asked);
1074
1075 // Whether this round put a question to the user, which is what ends the turn. Set
1076 // from the tool RESULT and not from the call, because a question can be refused --
1077 // see [`ends_turn`], and see what happened the last time a rule like this read a
1078 // name alone.
1079 let mut asked = false;
1080
1081 // Execute the round's tool calls, recording every result in both places for the
1082 // same reason. WHAT MAY RUN TOGETHER, AND WHY SO LITTLE MAY, IS `batch`'s RULE --
1083 // read it there rather than inferring it from here. Everything this loop does with
1084 // a result it does in the model's own order, whatever ran together underneath.
1085 let names: Vec<&str> = resp.tool_calls.iter().map(|t| t.name.as_str()).collect();
1086 for span in batch::batches(&names) {
1087 let group = &resp.tool_calls[span];
1088 // THE EVENTS STAY STRICTLY ALTERNATING: one `ToolCall`, then its `ToolResult`,
1089 // then the next. `www/js/daimond.js` keeps ONE `pendingTool` and ONE
1090 // `pendingCallId` -- in the chat, in the worker dock and in the daimon alike --
1091 // and files a result against whichever call was announced last. Announce two
1092 // calls back to back and the first result is filed under the second call while
1093 // every result after it is dropped, in the transcript and in the write-ahead
1094 // journal both. So a batch is invisible on the wire and shows only as the round
1095 // being quicker.
1096 //
1097 // The FIRST call is announced BEFORE the batch runs rather than after, so
1098 // "Running <tool>, step n…" names something that is actually running for as long
1099 // as the batch is in flight, exactly as it does for a batch of one.
1100 on_event(Self::announce(&group[0]));
1101 // A sink per call while a batch runs, drained into `on_event` in the model's
1102 // order as each result is recorded. Empty for a batch of one, which keeps its
1103 // events live. Today it is provably empty for a batch of several as well --
1104 // only `over_the_road` writes to it, and nothing in `batch::may_run_beside` is
1105 // `road_retryable` -- but a tool that later becomes both would otherwise have its
1106 // caption emitted from inside a concurrent poll, against a `&mut` closure that
1107 // cannot be shared.
1108 let mut sinks: Vec<Vec<AgentEvent>> = Vec::new();
1109 let outs: Vec<Outcome<MessageContent>> = if group.len() == 1 {
1110 vec![self.one_call(registry, &group[0], resp.truncated, on_event).await]
1111 } else {
1112 sinks = group.iter().map(|_| Vec::new()).collect();
1113 let futs: Vec<_> = group.iter().zip(sinks.iter_mut())
1114 .map(|(tc, sink)| self.call_buffered(registry, tc, resp.truncated, sink))
1115 .collect();
1116 batch::all_of(futs).await
1117 };
1118 for (n, (tc, out)) in group.iter().zip(outs).enumerate() {
1119 // Announced here for every call after the first, whose announcement went out
1120 // before the batch started.
1121 if n > 0 {
1122 on_event(Self::announce(tc));
1123 }
1124 if let Some(sink) = sinks.get_mut(n) {
1125 for e in std::mem::take(sink) {
1126 on_event(e);
1127 }
1128 }
1129 let result = match out {
1130 Ok(c) => c,
1131 // THE ROAD IS SPENT, SO THE TURN IS OVER, and the model is told nothing.
1132 //
1133 // This is the whole point of the exercise. Handing the model an
1134 // exhausted ladder produces the same apology the first failure would
1135 // have produced, ninety seconds later and for more money -- and that
1136 // apology is a durable assistant turn, re-sent on every turn after it.
1137 // So nothing is written for it to read. The app says what happened
1138 // instead, in its own voice, and offers the turn back.
1139 //
1140 // `Error` then `Err`, in that order, because that is the pair the page
1141 // is already built around: `runTurn`'s error handler recognises a road
1142 // failure and declines to write a line for it, and its catch classifies
1143 // the same failure a second time and hands the turn back badged with a
1144 // Continue. Both live in www/js/daimond.js and neither needed changing.
1145 //
1146 // A SIBLING THAT RAN BESIDE THIS ONE AND SUCCEEDED IS DROPPED WITH IT,
1147 // and that is deliberate: serially it would never have run at all, so
1148 // keeping its result would put something in the conversation that
1149 // today's behaviour does not. The calls BEFORE it in the model's order
1150 // are already recorded, exactly as they are serially, and `abandon_round`
1151 // pairs off what is left.
1152 Err(e) => {
1153 self.abandon_round(session, &mut working);
1154 on_event(AgentEvent::Error(e.plain()));
1155 let ending = self.audit(
1156 TurnEnd::Failed, rounds, &claims, Some(registry)).await;
1157 self.ended(ending, on_event);
1158 return Err(e);
1159 }
1160 };
1161 // The event carries the TEXT of the result and not the image. Everything
1162 // downstream of it -- the panel, the journal's write-ahead log, the
1163 // transcript on screen -- renders a string, and an image inside that string
1164 // would be a base64 wall in a tile. The image travels in the message instead,
1165 // where the model is the only reader of it.
1166 //
1167 // WHAT THE CALL CAME TO IS DECIDED HERE, ONCE. This is the only place the
1168 // event is built, so it is the only place the outcome is set; every reader
1169 // downstream asks rather than re-reading the prose. Four readers used to
1170 // guess it back out of the text, and one of them did not know that a refusal
1171 // opens "Refused" rather than "Error" -- so a write the fence had just
1172 // stopped was drawn as a completed step.
1173 let text = result.as_text().into_owned();
1174 let outcome = crate::tools::call_outcome(&text);
1175 // AND THE AUDIT IS KEPT FROM THE SAME VERDICT, not from a second reading of
1176 // it. The paths come from the ARGUMENTS the model sent, which name the file
1177 // it meant whatever the reply says about it.
1178 claims.record(&tc.name, &tc.arguments, outcome);
1179 // ASKED OF THE RESULT, and of every call in the round rather than of the
1180 // first: a round may carry several, and the question may not be the one that
1181 // came back first.
1182 asked |= ends_turn(&tc.name, outcome);
1183 on_event(AgentEvent::ToolResult {
1184 name: tc.name.clone(),
1185 result: text,
1186 outcome,
1187 });
1188 // A PICTURE FOR A MODEL THAT WILL NOT TAKE ONE NEVER ENTERS THE SESSION.
1189 //
1190 // `sighted()` takes it out of every request anyway, so the model sees the
1191 // same words either way -- but stored, the picture is folded around,
1192 // reloaded, and re-stripped for the life of the conversation. That is what
1193 // bricked a real Diamond on 2026-08-13. Left out here it costs one elision
1194 // that names the file, and the act is announced once instead of being
1195 // invisible.
1196 let result = if result.has_image() && !self.llm.can_take_images() {
1197 on_event(AgentEvent::Unseeable {
1198 images: result.images().count(),
1199 model: self.llm.model.clone(),
1200 });
1201 result.without_images(Dropped::Unseeable)
1202 } else {
1203 result
1204 };
1205 let reply = ChatMessage::tool(tc.id.clone(), result);
1206 working.push(reply.clone());
1207 session.messages.push(reply);
1208 }
1209 }
1210
1211 // A QUESTION IS THE ANSWER, so the turn is over. The model has just put a decision
1212 // on the user's screen and cannot say anything useful until it is answered: going
1213 // round again would buy a whole extra request whose only possible content is a
1214 // paragraph restating the question, printed under a card that already asks it.
1215 //
1216 // **The turn ending is what makes an unanswered question free.** Nothing is held --
1217 // no promise, no engine, no slot -- so a question nobody answers for an hour costs
1218 // exactly what a question nobody answers for a second does, and there is no timeout
1219 // to invent because there is nothing to time out. That is the difference between
1220 // this and `parkConsent`, which holds a worker on a `resolve` in memory and loses it
1221 // to a reload.
1222 //
1223 // Ended as `Answered` rather than under an ending of its own: this is a reply with
1224 // nothing further to say, which is what that word means, and a sixth `TurnEnd` would
1225 // be a word every locale and every reader of the ledger had to learn to draw a
1226 // distinction nothing acts on.
1227 if asked {
1228 let ending = self.audit(TurnEnd::Answered, rounds, &claims, Some(registry)).await;
1229 self.ended(ending, on_event);
1230 on_event(AgentEvent::Done);
1231 return Ok(());
1232 }
1233
1234 // THE SEAM. The tool replies are in, and the next request has not gone out,
1235 // so this is the one moment in a round where the conversation can grow by
1236 // something the model has not already been told. Anything the user said
1237 // while the tools ran goes in here, and the next round is built with it.
1238 //
1239 // After the tool replies rather than before them, because the API requires
1240 // every tool_call to be answered by a tool message: a user turn wedged
1241 // between the two is a malformed request, and a provider is entitled to
1242 // reject the whole thing.
1243 for said in self.take_interjections() {
1244 on_event(AgentEvent::Interjected(said.clone()));
1245 let msg = ChatMessage::user(said);
1246 working.push(msg.clone());
1247 session.messages.push(msg);
1248 }
1249 }
1250
1251 // Exceeded the tool-round budget.
1252 //
1253 // Recorded in the SYSTEM voice, because that is whose it is. It used to be pushed
1254 // as an assistant message reading "[Reached the tool-call round limit (25).]", so
1255 // on the next turn the model read its own surrender back as something it had said
1256 // and behaved accordingly -- a turn that was stopped from outside became, in the
1257 // record, a turn that gave up. The boundary belongs to the app, so it is said in
1258 // the app's voice, and it says the work may be unfinished rather than that it is
1259 // over.
1260 let msg = fmt!("Reached the tool-call round limit ({}).", max_rounds);
1261 on_event(AgentEvent::Error(msg.clone()));
1262 session.messages.push(compact::round_limit_note(max_rounds));
1263 let ending = self.audit(TurnEnd::Capped, rounds, &claims, Some(registry)).await;
1264 self.ended(ending, on_event);
1265 on_event(AgentEvent::Done);
1266 Ok(())
1267 }
1268
1269
1270 // ┌───────────────────────────────────────────────────────────┐
1271 // │ Folding a conversation that no longer fits │
1272 // └───────────────────────────────────────────────────────────┘
1273
1274 /// The output cap this turn's requests will ACTUALLY carry, which is not always the one the
1275 /// client was configured with.
1276 ///
1277 /// [`crate::compact::Limits::budget`] subtracts the reply from the window to decide how big a
1278 /// prompt may be, so it has to be given the figure the provider will be sent. It was being
1279 /// given `llm.max_tokens`, and on the one path that matters that is the wrong figure by a
1280 /// factor of eight: a streamed Anthropic request to a model that takes adaptive thinking is
1281 /// sent `THINKING_MIN_MAX_TOKENS`, because thinking is billed as output and counts against the
1282 /// same cap.
1283 ///
1284 /// **On a large window the error is invisible and on a small one it is fatal.** With 131,072
1285 /// tokens the fraction ceiling is the lower of the two and wins whatever the reserve says. But
1286 /// a window is not always the published one: [`Limits::learn_from_refusal`] sets it from a
1287 /// provider's refusal, and at 40,000 the old figure reserved 5,120 tokens for a reply that may
1288 /// run to 32,000, left a budget the conversation already fitted, folded nothing, sent the same
1289 /// prompt, and was refused again -- a turn with no way out, arrived at by the very mechanism
1290 /// that exists to recover from the first refusal.
1291 ///
1292 /// The agent always streams (see [`Agent::run_streaming`] and the tool loop), so the streaming
1293 /// half of the client's rule is a constant here and only the model has to be asked about.
1294 fn reply_cap(&self) -> u32 {
1295 match self.llm.dialect {
1296 crate::llm::Dialect::Anthropic
1297 if crate::llm::model_takes_adaptive_thinking(&self.llm.model) =>
1298 self.llm.max_tokens.max(crate::llm::THINKING_MIN_MAX_TOKENS),
1299 _ => self.llm.max_tokens,
1300 }
1301 }
1302
1303 /// Whether a failed round was the provider refusing an oversized prompt.
1304 ///
1305 /// # Arguments
1306 /// * `e` - The error the call returned.
1307 /// * `bytes` - Bytes of prompt that were refused.
1308 fn overflowed(&self, e: &Error<ErrTag>, bytes: u64) -> bool {
1309 let budget = self.limits.borrow().budget(self.reply_cap());
1310 compact::looks_like_overflow(&fmt!("{}", e), self.gauge.tokens(bytes), budget)
1311 }
1312
1313 /// Fold the conversation if it no longer fits, and rebuild `working` from what is left.
1314 ///
1315 /// Returns whether anything changed. Two mechanisms, tried in that order and both
1316 /// bounded:
1317 ///
1318 /// 1. **Fold.** Everything before the cut becomes one note carrying a summary and a
1319 /// ledger of what was touched; everything after it is kept exactly as it was.
1320 /// 2. **Elide.** Whatever is still too big has its older tool results shrunk in place.
1321 /// This adds and removes nothing, so it cannot orphan a tool call, and it is the only
1322 /// thing that helps when the conversation is one enormous turn with no earlier part
1323 /// to fold.
1324 ///
1325 /// It never gives up and sends the oversized history, because that is the bug: the
1326 /// provider refuses it, the turn dies, and the next turn sends the same thing.
1327 ///
1328 /// # Arguments
1329 /// * `session` - The durable conversation, folded in place.
1330 /// * `working` - This turn's message list, rebuilt from the session when anything moved.
1331 /// * `schema` - Bytes of tool definitions riding alongside the conversation.
1332 /// * `why` - What brought the fold about; see [`Fold`].
1333 /// * `on_event` - Where the user is told, since a silent fold is the one people hate.
1334 async fn fold_if_needed(
1335 &self,
1336 session: &mut Session,
1337 working: &mut Vec<ChatMessage>,
1338 schema: u64,
1339 why: Fold,
1340 on_event: &mut impl FnMut(AgentEvent),
1341 )
1342 -> bool
1343 {
1344 // THE FOLDS THE USER HAS OPEN, taken once for the whole fold. Every size below is a
1345 // size on the wire, and a closed fold's detail is not on the wire: `msg_bytes` asks the
1346 // serialiser's own `sent_args_len` what a `say` costs, and that answer depends on this
1347 // set. A copy rather than a borrow, because the page may set the folds again while the
1348 // summarising call is in flight and a live borrow at that moment would panic.
1349 let open = self.llm.open_folds();
1350
1351 // A refusal is the only occasion the provider ever says anything about the size of
1352 // its window. Believing it is what lets a chat against a model nobody published a
1353 // window for recover instead of folding to a budget that was never the real one --
1354 // and it is why a fold the USER asked for must not come through here, since that
1355 // one carries no news at all.
1356 if why.teaches_window() {
1357 let refused = self.gauge.tokens(compact::conversation_bytes(working, &open) + schema);
1358 self.limits.borrow_mut().learn_from_refusal(refused);
1359 }
1360 let (budget, tail, model) = {
1361 let l = self.limits.borrow();
1362 let cap = self.reply_cap();
1363 (l.budget(cap), l.tail_budget(cap), l.fold_model.clone())
1364 };
1365 let before = compact::conversation_bytes(&session.messages, &open);
1366 if !why.forces()
1367 && self.gauge.tokens(compact::conversation_bytes(working, &open) + schema) <= budget
1368 {
1369 return false;
1370 }
1371
1372 let mut folded = 0usize;
1373 let mut trouble = String::new();
1374 // The hard ceiling is the budget itself: however few messages that leaves, a tail
1375 // bigger than what may be sent is a fold that changed nothing.
1376 let ceiling = self.gauge.bytes(budget).saturating_sub(schema);
1377 let cut = compact::tail_start(&session.messages, self.gauge.bytes(tail).min(ceiling),
1378 compact::MIN_KEEP_MESSAGES, ceiling, &open);
1379 if cut > 0 {
1380 // Built before the summarising call, so a call that fails still leaves a
1381 // truthful record: which files were read, which were written, what ran, and
1382 // which of those failed. That is the part no model is asked for and so no model
1383 // can lose.
1384 let ledger = compact::ledger_of(&session.messages[..cut]);
1385 let rendered = compact::render_for_fold(&session.messages[..cut],
1386 compact::FOLD_INPUT_CAP);
1387 let (prose, why) = match self.summarise(&rendered, &model, session).await {
1388 Ok(s) => (s, None),
1389 Err(e) => (String::new(), Some(fmt!("{}", e))),
1390 };
1391 if let Some(ref w) = why { trouble = w.clone(); }
1392 let note = compact::notice(cut, &prose, &ledger, why.as_deref());
1393 match compact::fold(&session.messages, cut, note) {
1394 Ok(new) => {
1395 // A fold that made the conversation bigger is not a fold; it happens
1396 // when the folded part was small and the note is not.
1397 if compact::conversation_bytes(&new, &open) < before {
1398 session.messages = new;
1399 folded = cut;
1400 }
1401 },
1402 // Refused rather than allowed to orphan a tool call. Eliding below still
1403 // shrinks the same conversation, and cannot orphan anything at all.
1404 Err(e) => trouble = fmt!("{}", e),
1405 }
1406 }
1407
1408 // SHORTENED ON THE WAY OUT, AND NOWHERE ELSE. This edited `session.messages` until
1409 // 2026-08-28, and that list is the one `crate::wasm::app::DaimondApp::export_session`
1410 // hands the browser to store, to back up and to put in the sync parcel -- so a
1411 // thousand-word answer became four hundred characters in the user's own record,
1412 // permanently, with nothing said and no way back. A model's window is a property of the
1413 // REQUEST; a lossy form of the conversation is therefore the request's and is built here,
1414 // from a record that stays whole. The owner's ruling: the model gets the shortened
1415 // version, his transcript keeps every word.
1416 //
1417 // Two passes, and the second is what makes the guarantee hold: the first leaves the
1418 // newest messages alone, and if the conversation is STILL too big it is because those
1419 // are the bulky ones, so the second reaches them too. The user's own words are never
1420 // touched by either.
1421 //
1422 // IT ALSO SETTLES WHAT THE FOLD READS. `render_for_fold` above summarises
1423 // `session.messages`, and while the elision edited that list a later fold summarised
1424 // whatever an earlier elision had left of it -- a summary of four-hundred-character
1425 // stubs, which is a second silent loss standing behind the first. Nothing clips that
1426 // list now, so there is no arrangement of turns in which it can happen.
1427 let mut sent = session.messages.clone();
1428 let mut elided = compact::elide_bulk(&mut sent, ceiling,
1429 compact::MIN_KEEP_MESSAGES, &open);
1430 elided += compact::elide_bulk(&mut sent, ceiling, 1, &open);
1431 let changed = folded > 0 || elided > 0;
1432 if !changed {
1433 // And `working` is left exactly as it was, elisions and all. Rebuilding it from the
1434 // session on the way out of a fold that changed nothing would throw away an earlier
1435 // round's shortening and hand the caller a request over the budget again.
1436 return false;
1437 }
1438 // Measured on what will be SENT, since that is what the sentence below is about.
1439 let after = compact::conversation_bytes(&sent, &open);
1440
1441 // Rebuild the turn's list: the system prompt this turn was built with, then the
1442 // conversation as it now stands, shortened to fit. The two are kept in lockstep for
1443 // the whole turn, so anything else would leave the model reading a history the session
1444 // no longer holds.
1445 let sys = match working.first() {
1446 Some(m @ ChatMessage::System { .. }) => Some(m.clone()),
1447 _ => None,
1448 };
1449 working.clear();
1450 if let Some(s) = sys { working.push(s); }
1451 working.extend(sent);
1452
1453 // Told, not done quietly. A fold is lossy, and a user who is not shown one has no
1454 // way to tell a model that forgot from a model that never knew.
1455 //
1456 // "TOOL RESULTS" WAS NOT TRUE. `compact::elide_bulk` shrinks a long ASSISTANT turn on
1457 // exactly the same rule it shrinks a tool reply on, and a pure chat has no tool replies
1458 // at all -- so a user whose own answers had just been clipped to 400 characters was told
1459 // the app had shortened some tool output. That is the fold telling them the one thing
1460 // they would not go looking for.
1461 //
1462 // AND IT SAYS WHERE THE SHORTENING APPLIES, which is four words and the whole of the
1463 // second half of the ruling. The sentence was true of a record that no longer changes:
1464 // a reader who is told his answers were shortened, and is looking at them in full on
1465 // the screen above, has been handed a contradiction to resolve on his own.
1466 //
1467 // AND IT NAMES ONLY WHAT HAPPENED. One sentence covered both mechanisms and reported the
1468 // one that did not fire as a zero, so a conversation that was merely shortened opened
1469 // with "Folded 0 earlier messages and", under a heading that had just said the opposite.
1470 // A count of nothing is not information; it is a reader working out which half to ignore.
1471 let did = match (folded, elided) {
1472 (0, n) => fmt!("Shortened {} long tool results and answers on the way to the model",
1473 n),
1474 (f, 0) => fmt!("Folded {} earlier messages", f),
1475 (f, n) => fmt!("Folded {} earlier messages and shortened {} long tool results and \
1476 answers on the way to the model", f, n),
1477 };
1478 let mut said = fmt!("{}: {} tokens of conversation became about {}.",
1479 did, self.gauge.tokens(before), self.gauge.tokens(after));
1480 if !trouble.is_empty() {
1481 said.push_str(&fmt!(" The summary could not be written ({}), so only the record \
1482 of what was read and written was kept.", trouble));
1483 }
1484 // Its own event, not a borrowed tool row. A fold is something the app did to the
1485 // user's conversation, and it is lossy; the counts travel beside the sentence so a
1486 // client can draw it as the act it was rather than parse prose to find out.
1487 on_event(AgentEvent::Compacted {
1488 folded,
1489 kept: session.messages.len(),
1490 note: said,
1491 });
1492 true
1493 }
1494
1495 /// Ask a model to summarise the part of the conversation being folded.
1496 ///
1497 /// Non-streaming and tool-less: nothing here should reach the user's thread or touch
1498 /// their files. The input is already bounded by [`compact::render_for_fold`], and the
1499 /// output is capped, so the call costs a fixed amount however long the session got --
1500 /// which matters, because the alternative to folding is paying for the whole history on
1501 /// every round from now until the chat is closed.
1502 ///
1503 /// # Arguments
1504 /// * `rendered` - The folded part as a bounded transcript.
1505 /// * `model` - The model to fold with, or empty for the chat's own.
1506 /// * `session` - Charged with what the call cost, so the fold is not spent invisibly.
1507 async fn summarise(
1508 &self,
1509 rendered: &str,
1510 model: &str,
1511 session: &mut Session,
1512 )
1513 -> Outcome<String>
1514 {
1515 let mut llm = self.llm.clone();
1516 if !model.trim().is_empty() {
1517 llm.model = model.trim().to_string();
1518 }
1519 llm.max_tokens = compact::FOLD_MAX_TOKENS;
1520 let msgs = vec![
1521 ChatMessage::system(self.fold_prompt()),
1522 ChatMessage::user(fmt!(
1523 "Here is the earlier part of the conversation.\n\n{}", rendered)),
1524 ];
1525 let resp = res!(llm.chat_once(&msgs, None).await);
1526 session.prompt_tokens += resp.prompt_tokens;
1527 session.completion_tokens += resp.completion_tokens;
1528 session.cached_tokens += resp.cached_tokens;
1529 session.cost_usd += resp.cost_usd;
1530 self.live_prompt.set(session.prompt_tokens);
1531 self.live_completion.set(session.completion_tokens);
1532 self.live_cached.set(session.cached_tokens);
1533 self.live_cost.set(session.cost_usd);
1534 if resp.content.trim().is_empty() {
1535 return Err(err!("Fold: the model returned an empty summary."; Invalid, Data));
1536 }
1537 Ok(resp.content)
1538 }
1539}
1540
1541
1542// ┌───────────────────────────────────────────────────────────────┐
1543// │ The tool that ends a turn │
1544// └───────────────────────────────────────────────────────────────┘
1545
1546/// Whether this tool result is the turn's answer, so the loop stops here.
1547///
1548/// **The rule is back and the tool under it is not the old one.** It used to say `say`, which
1549/// folded a reply for the reader -- one-way presentation that a model could simply decline, and
1550/// that is now a convention in the model's own prose with no tool to call. `ask` is the opposite
1551/// shape: a round trip, whose whole point is that the model has stopped and is waiting. Prose
1552/// cannot return a tap, so there is no convention that could replace it.
1553///
1554/// **It reads the OUTCOME and never the name alone**, and that clause is the scar. A refused
1555/// `say` ended a worker's turn, so the report the refusal had just told it to write was never
1556/// written and the whole errand came back as whatever prose accompanied the call -- work done,
1557/// paid for and thrown away. `ask` has as many ways to be refused as `say` had, each of them
1558/// telling the model to put the question properly, and every one of those is advice the model
1559/// must be given a round to take. [`crate::tools::call_outcome`] is the tool layer's own
1560/// statement of what became of a call; nothing here reads the wording again.
1561///
1562/// # Arguments
1563/// * `name` - The tool the call named.
1564/// * `outcome` - What the layer said became of it.
1565fn ends_turn(name: &str, outcome: crate::tools::CallOutcome) -> bool {
1566 name == crate::tools::Tool::Ask.name()
1567 && matches!(outcome, crate::tools::CallOutcome::Done)
1568}
1569
1570
1571// ┌───────────────────────────────────────────────────────────────┐
1572// │ A reply that ran out of room │
1573// └───────────────────────────────────────────────────────────────┘
1574
1575/// Whether a tool call's arguments are the wreckage of a reply cut at the output limit.
1576///
1577/// Both halves are needed. The provider's `finish_reason` says the reply was cut, but a
1578/// turn can be cut in its trailing prose with every tool call already whole -- and
1579/// refusing a complete call because a later sentence was truncated would break a turn
1580/// that was working. So the arguments are checked too, and only a call that is BOTH
1581/// truncated and structurally incomplete is treated as one.
1582///
1583/// # Arguments
1584/// * `truncated` - What the provider said about why it stopped.
1585/// * `arguments` - The raw JSON arguments object as the model produced it.
1586pub fn cut_short(truncated: bool, arguments: &str) -> bool {
1587 truncated && !json_object_is_whole(arguments)
1588}
1589
1590/// Whether `s` is a single JSON object with every brace, bracket and quote closed.
1591///
1592/// Structural only -- it says nothing about whether the fields are the right ones -- and
1593/// that is the whole question here: an object that stops in the middle is a reply that
1594/// ran out of room, and one that closes is a call worth dispatching whatever else may be
1595/// wrong with it. String contents are skipped, so a brace inside a path or a piece of
1596/// source code is not counted.
1597///
1598/// # Arguments
1599/// * `s` - The raw arguments text.
1600pub fn json_object_is_whole(s: &str) -> bool {
1601 let t = s.trim();
1602 // The empty call, which `StreamAcc` normalises to `{}` and which is legal.
1603 if t.is_empty() {
1604 return true;
1605 }
1606 let b = t.as_bytes();
1607 if b[0] != b'{' {
1608 return false;
1609 }
1610 let mut depth = 0i32;
1611 let mut i = 0usize;
1612 while i < b.len() {
1613 match b[i] {
1614 b'{' | b'[' => depth += 1,
1615 b'}' | b']' => {
1616 depth -= 1;
1617 if depth == 0 {
1618 // Anything after the close is not this object's business, but the
1619 // object itself arrived whole.
1620 return true;
1621 }
1622 if depth < 0 {
1623 return false;
1624 }
1625 },
1626 // Skip the string's contents. Without this a brace inside a file being
1627 // written -- `fn main() {` -- would be counted as structure, and a complete
1628 // call would read as a cut one. A string the reply stopped INSIDE, which is
1629 // the commonest cut of all, runs to the end of the input and falls out of
1630 // the loop below unclosed, which is already the answer.
1631 b'"' => {
1632 i += 1;
1633 while i < b.len() {
1634 if b[i] == b'\\' { i += 2; continue; }
1635 if b[i] == b'"' { break; }
1636 i += 1;
1637 }
1638 },
1639 _ => {},
1640 }
1641 i += 1;
1642 }
1643 false
1644}
1645
1646/// What the model is told when its own tool call was cut at the output limit.
1647///
1648/// A sentence it can act on, not a parse error. Shown a parse error a model reads its
1649/// own JSON as the mistake and writes the same call again, which is cut in the same
1650/// place; told that the reply ran out of room, it splits the work. It says the number,
1651/// because "smaller" without a figure is advice rather than a constraint.
1652///
1653/// # Arguments
1654/// * `max_tokens` - The output cap this turn ran under.
1655pub fn truncated_call_note(max_tokens: u32) -> String {
1656 fmt!(
1657 "Error: your reply was cut off at the output limit of {} tokens, part-way through \
1658 the arguments of this tool call, so it could not be run. Nothing was changed. \
1659 The JSON was not wrong -- there was no room left for the rest of it. Do the same \
1660 work in smaller pieces: write or edit a part at a time, or read a range rather \
1661 than a whole file, so that no single call carries this much text.",
1662 max_tokens)
1663}
1664
1665
1666// ┌───────────────────────────────────────────────────────────────┐
1667// │ TLS config helper │
1668// └───────────────────────────────────────────────────────────────┘
1669
1670/// Build a TLS client config using the system CA bundle.
1671///
1672/// Reused from Steel's `build_outbound_tls_client` — same approach
1673/// but kept here so `daimond` can be used standalone.
1674#[cfg(not(target_arch = "wasm32"))]
1675pub fn build_tls_client_config() -> Outcome<Arc<ClientConfig>> {
1676 use tokio_rustls::rustls::{
1677 ClientConfig,
1678 RootCertStore,
1679 pki_types::CertificateDer,
1680 };
1681
1682 let ca_paths = [
1683 "/etc/ssl/certs/ca-certificates.crt",
1684 "/etc/pki/tls/certs/ca-bundle.crt",
1685 "/etc/ssl/cert.pem",
1686 ];
1687 let ca_file = match ca_paths.iter().find(|p| std::path::Path::new(p).exists()) {
1688 Some(p) => *p,
1689 None => return Err(err!(
1690 "No system CA bundle found."; Init, Missing, File)),
1691 };
1692
1693 let pem_data = match std::fs::read(ca_file) {
1694 Ok(d) => d,
1695 Err(e) => return Err(err!(e, "Failed to read CA bundle."; File, Read)),
1696 };
1697
1698 let mut roots = RootCertStore::empty();
1699 let certs: Vec<CertificateDer> = rustls_pemfile::certs(&mut pem_data.as_slice())
1700 .filter_map(|c| c.ok())
1701 .map(CertificateDer::from)
1702 .collect();
1703 for cert in certs {
1704 let _ = roots.add(cert);
1705 }
1706
1707 let mut config = ClientConfig::builder()
1708 .with_root_certificates(roots)
1709 .with_no_client_auth();
1710 // Advertise HTTP/1.1 via ALPN so CDN-fronted servers (e.g.
1711 // Fireworks.ai behind Cloudflare) don't close the connection
1712 // after the TLS handshake when no protocol is negotiated.
1713 config.alpn_protocols = vec![b"http/1.1".to_vec()];
1714
1715 Ok(Arc::new(config))
1716}
1717
1718
1719// ┌───────────────────────────────────────────────────────────────┐
1720// │ Tests │
1721// └───────────────────────────────────────────────────────────────┘
1722
1723#[cfg(test)]
1724mod tests {
1725 use super::*;
1726 use crate::llm::LlmClient;
1727 use crate::tools::CallOutcome;
1728
1729 use oxedyne_fe2o3_jdat::prelude::*;
1730
1731 fn make_test_agent() -> Agent {
1732 let tls = build_test_tls_config();
1733 let llm = LlmClient::new("api.test.com", 443, "/v1/chat", "key", "model", 4096, tls);
1734 Agent::new(llm, "You are Daimond, an AI assistant.")
1735 }
1736
1737 fn build_test_tls_config() -> Arc<ClientConfig> {
1738 use rustls::crypto::ring;
1739 let _ = ring::default_provider().install_default();
1740 ClientConfig::builder()
1741 .dangerous()
1742 .with_custom_certificate_verifier(Arc::new(crate::llm::tests::NoVerify))
1743 .with_no_client_auth().into()
1744 }
1745
1746 #[test]
1747 fn test_agent_creation() {
1748 let agent = make_test_agent();
1749 assert_eq!(agent.system_prompt, "You are Daimond, an AI assistant.");
1750 }
1751
1752 // ── Speaking into a running turn ────────────────────────────────
1753
1754 #[test]
1755 fn test_what_the_user_says_mid_turn_is_kept_until_there_is_a_seam_00() {
1756 let a = make_test_agent();
1757 assert_eq!(a.interjections().len(), 0, "nothing is waiting before anything is said");
1758 assert_eq!(a.interject("no, use the other file"), 1);
1759 assert_eq!(a.interject(" and run the tests "), 2);
1760 // Held, not lost: the UI has to be able to draw what is waiting, because a
1761 // correction that vanishes between typing it and its taking effect reads as
1762 // an app that ignored you.
1763 assert_eq!(a.interjections(), vec![
1764 fmt!("no, use the other file"), fmt!("and run the tests")]);
1765 }
1766
1767 #[test]
1768 fn test_blank_input_is_not_an_interjection_00() {
1769 let a = make_test_agent();
1770 assert_eq!(a.interject(" "), 0);
1771 assert_eq!(a.interject("\n\t"), 0);
1772 assert!(a.interjections().is_empty(), "whitespace is not a correction");
1773 }
1774
1775 #[test]
1776 fn test_taking_them_empties_the_queue_so_none_is_said_twice_00() {
1777 // Drained rather than read: once it is in the conversation, the conversation
1778 // is the record. A queue that still held it would say it again next round,
1779 // and a model told the same correction three times reasonably concludes it
1780 // has not yet complied.
1781 let a = make_test_agent();
1782 a.interject("stop and summarise");
1783 let taken = a.take_interjections();
1784 assert_eq!(taken, vec![fmt!("stop and summarise")]);
1785 assert!(a.interjections().is_empty());
1786 assert!(a.take_interjections().is_empty(), "a second take yields nothing");
1787 }
1788
1789 #[test]
1790 fn test_an_interjection_is_a_user_turn_in_the_conversation_00() {
1791 // The shape the seam pushes. It has to be a User message: an Assistant or a
1792 // Tool message would be the model reading its own words back as though it
1793 // had said them, which is the opposite of being corrected.
1794 let a = make_test_agent();
1795 a.interject("actually, target wasm");
1796 let taken = a.take_interjections();
1797 let msg = ChatMessage::user(taken[0].clone());
1798 assert_eq!(msg.role(), "user");
1799 assert_eq!(msg.text(), "actually, target wasm");
1800 }
1801
1802 // ── What bounds a turn ──────────────────────────────────────────────
1803
1804 #[test]
1805 fn test_the_round_limit_and_the_window_are_settable_00() {
1806 // Neither is a constant any more: the window is per-model and comes from the
1807 // provider's own catalogue, and a user who wants a longer leash should be able to
1808 // have one without a rebuild.
1809 let a = make_test_agent();
1810 assert_eq!(a.limits().max_rounds, compact::DEFAULT_MAX_ROUNDS);
1811 assert_eq!(a.limits().window, 0, "nobody has said what the window is yet");
1812 a.set_context_window(204_800);
1813 a.set_max_rounds(60);
1814 assert_eq!(a.limits().window, 204_800);
1815 assert_eq!(a.limits().max_rounds, 60);
1816 // A ceiling of zero is a turn with no tools, which is not what anyone means.
1817 a.set_max_rounds(0);
1818 assert_eq!(a.limits().max_rounds, 60);
1819 }
1820
1821 #[test]
1822 fn test_a_worker_inherits_the_window_of_the_turn_that_dispatched_it_00() {
1823 // Shared on clone, exactly as the interjection queue is. A worker that fell back to
1824 // the default window would fold its own conversation at the wrong size.
1825 let a = make_test_agent();
1826 a.set_context_window(32_768);
1827 let worker = a.clone();
1828 assert_eq!(worker.limits().window, 32_768);
1829 worker.set_max_rounds(9);
1830 assert_eq!(a.limits().max_rounds, 9);
1831 }
1832
1833 #[test]
1834 fn test_an_agent_built_fresh_can_adopt_the_chats_figures_00() {
1835 // A Diamond's daimon and its reducer are each built with `Agent::new` rather
1836 // than cloned, because their system prompt differs -- and `Agent::new` starts
1837 // from the default window, which is a figure nobody published. Against the same
1838 // model as the chat they would then fold at a different size, and on a model
1839 // whose real window is smaller than the assumed one they would learn it the hard
1840 // way all over again, one dead turn each.
1841 let chat = make_test_agent();
1842 chat.set_context_window(32_768);
1843 chat.set_max_rounds(40);
1844 chat.set_fold_model("cheap/fast");
1845 let derived = Agent::new(chat.llm.clone(), "a different prompt");
1846 assert_eq!(derived.limits().window, 0, "a fresh agent starts knowing nothing");
1847 derived.adopt_limits(&chat);
1848 assert_eq!(derived.limits().window, 32_768);
1849 assert_eq!(derived.limits().max_rounds, 40);
1850 assert_eq!(derived.limits().fold_model, "cheap/fast");
1851 assert_eq!(derived.limits().budget(4_096), chat.limits().budget(4_096),
1852 "the two fold at different sizes against the same model");
1853 }
1854
1855 #[test]
1856 fn test_adopting_leaves_the_agent_it_copied_from_alone_00() {
1857 // It copies one way. Writing into `from` instead -- or as well -- would let a
1858 // worker's own overflow shrink the chat's window under the user, which is a
1859 // thing they never did and cannot see.
1860 let chat = make_test_agent();
1861 chat.set_context_window(32_768);
1862 chat.set_max_rounds(40);
1863 let derived = Agent::new(chat.llm.clone(), "worker");
1864 derived.set_max_rounds(9);
1865 derived.adopt_limits(&chat);
1866 assert_eq!(chat.limits().window, 32_768, "the chat's window moved");
1867 assert_eq!(chat.limits().max_rounds, 40, "the chat's round ceiling moved");
1868 // And afterwards the two are independent, not one cell shared between them.
1869 derived.set_context_window(8_192);
1870 assert_eq!(chat.limits().window, 32_768, "the two share one figure");
1871 }
1872
1873 // ── What the fold is told ───────────────────────────────────────────
1874
1875 #[test]
1876 fn test_the_fold_runs_under_the_compactors_prompt_00() {
1877 // It used to be a private constant in `compact`, which made it the one prompt in
1878 // the app the user could neither read nor change.
1879 let a = make_test_agent();
1880 assert_eq!(a.fold_prompt(), compact_role_default());
1881 assert!(a.fold_prompt().contains("context window"), "{}", a.fold_prompt());
1882 }
1883
1884 #[test]
1885 fn test_a_rewritten_fold_prompt_reaches_the_summarising_call_00() {
1886 let a = make_test_agent();
1887 a.set_fold_prompt("Answer with the file names and nothing else.");
1888 assert_eq!(a.fold_prompt(), "Answer with the file names and nothing else.");
1889 // And emptying it puts the shipped prompt back, which is what deleting
1890 // `prompts/compactor.md` does.
1891 a.set_fold_prompt(" ");
1892 assert_eq!(a.fold_prompt(), compact_role_default());
1893 }
1894
1895 #[test]
1896 fn test_the_fold_prompt_is_not_the_reducers_00() {
1897 // A user who has rewritten `prompts/reducer.md` for their Diamonds must not
1898 // thereby change how their chats are folded.
1899 let a = make_test_agent();
1900 assert_ne!(a.fold_prompt(), crate::prompts::DEFAULT_REDUCER);
1901 assert!(!a.fold_prompt().contains("crystal"), "{}", a.fold_prompt());
1902 }
1903
1904 #[test]
1905 fn test_a_derived_agent_folds_by_the_instructions_the_user_wrote_00() {
1906 // `fold_model` rides in `Limits` and so was adopted already; the fold PROMPT is
1907 // text and did not, so a Diamond's daimon folded on the user's chosen model
1908 // while ignoring what they had written for it in `prompts/compactor.md` -- the
1909 // half of the setting that is visible on disk, and so the half whose absence
1910 // looks like the file not being read at all.
1911 let chat = make_test_agent();
1912 chat.set_fold_prompt("Keep only the file names.");
1913 let derived = Agent::new(chat.llm.clone(), "daimon");
1914 assert_eq!(derived.fold_prompt(), compact_role_default(),
1915 "a fresh agent starts on the shipped prompt");
1916 derived.adopt_limits(&chat);
1917 assert_eq!(derived.fold_prompt(), "Keep only the file names.");
1918 // One way, like the figures: a worker must not rewrite the chat's.
1919 derived.set_fold_prompt("something else");
1920 assert_eq!(chat.fold_prompt(), "Keep only the file names.",
1921 "the chat's fold prompt moved under it");
1922 }
1923
1924 // ── Why a fold is happening ─────────────────────────────────────────
1925
1926 #[test]
1927 fn test_only_a_providers_refusal_teaches_the_window_00() {
1928 // The distinction the `Fold` enum exists for. `learn_from_refusal` moves the
1929 // window DOWN and never up, so routing a user's button press through the same
1930 // arm as a refusal would shrink the window on every press: fold three times and
1931 // a third of the context is gone, with nothing on screen to say why.
1932 assert!(Fold::Refused.teaches_window());
1933 assert!(!Fold::ByHand.teaches_window(), "a hand fold shrinks the window");
1934 assert!(!Fold::IfNeeded.teaches_window());
1935 // Both of the other two override the estimate; only `IfNeeded` consults it.
1936 assert!(Fold::Refused.forces());
1937 assert!(Fold::ByHand.forces(), "a hand fold that consults the estimate is a no-op");
1938 assert!(!Fold::IfNeeded.forces());
1939 }
1940
1941 #[tokio::test]
1942 async fn test_a_hand_fold_leaves_the_window_where_it_was_00() {
1943 // The property stated as the user would see it, rather than as the enum states
1944 // it: press Fold on a chat that is nowhere near full, and the window it is
1945 // measured against must be the one it had a moment ago.
1946 let a = make_test_agent();
1947 a.set_context_window(32_768);
1948 let mut s = Session::new("s".into(), "u".into(), "m".into());
1949 for i in 0..12 {
1950 s.messages.push(ChatMessage::user(fmt!("message {}", i)));
1951 }
1952 let mut seen = Vec::new();
1953 let mut sink = |ev: AgentEvent| seen.push(ev);
1954 let _ = a.fold_by_hand(&mut s, &mut sink).await;
1955 assert_eq!(a.limits().window, 32_768,
1956 "the user's own fold was read as a provider refusal");
1957 }
1958
1959 #[tokio::test]
1960 async fn test_a_hand_fold_of_a_short_conversation_says_nothing_moved_00() {
1961 // `false` rather than a fold that changed nothing. Below `MIN_KEEP_MESSAGES`
1962 // there is no tail to cut and no bulky tool result to shorten, and the honest
1963 // answer is that this conversation cannot be made smaller -- not a spinner that
1964 // ends with the meter exactly where it started.
1965 let a = make_test_agent();
1966 a.set_context_window(32_768);
1967 let mut s = Session::new("s".into(), "u".into(), "m".into());
1968 s.messages.push(ChatMessage::user("hello"));
1969 let mut sink = |_ev: AgentEvent| {};
1970 let moved = a.fold_by_hand(&mut s, &mut sink).await;
1971 assert!(!moved, "a two-message conversation reported a fold");
1972 assert_eq!(s.messages.len(), 1, "the one message was folded away");
1973 }
1974
1975 /// What the compactor is told when the user has not said otherwise.
1976 fn compact_role_default() -> String {
1977 Role::Compactor.compose("")
1978 }
1979
1980 // ── How a fold reaches the user ─────────────────────────────────────
1981
1982 /// An agent whose provider is a port nothing is listening on.
1983 ///
1984 /// Every call fails at once, which is the point: what is under test is what the app
1985 /// does BEFORE the request goes out, and a fold has to happen and be announced
1986 /// whether or not the turn that provoked it then succeeds. Retrying is turned off
1987 /// so a refused connection costs no backoff.
1988 fn dead_agent() -> Agent {
1989 let tls = build_test_tls_config();
1990 let mut llm = LlmClient::new("127.0.0.1", 1, "/v1/chat", "key", "model", 4_096, tls);
1991 llm.retry.max_attempts = 1;
1992 Agent::new(llm, "You are Daimond.")
1993 }
1994
1995 /// A registry with no tools, so a turn takes the plain streaming path.
1996 fn no_tools() -> crate::tools::ToolRegistry {
1997 // A scratch directory under the user cache, not the tmpfs at `/tmp`, and one
1998 // per call rather than one per process: keyed on the process identifier, every
1999 // turn in the file shared a workspace.
2000 let dir = match oxedyne_fe2o3_test::scratch::scratch_dir("daimond_agent_test") {
2001 Ok(d) => d,
2002 Err(e) => panic!("a scratch directory: {}", e),
2003 };
2004 let ws = match crate::workspace::Workspace::new(dir) {
2005 Ok(w) => w,
2006 Err(e) => panic!("a scratch workspace: {}", e),
2007 };
2008 crate::tools::ToolRegistry::new(Vec::new(), crate::tools::ToolContext {
2009 workspace: ws,
2010 executor: crate::executor::Executor::local_default(),
2011 cwd: String::new(),
2012 path_prefix: String::new(),
2013 root: crate::tools::FileRoot::Workspace,
2014 read_seen: crate::tools::new_read_cache(),
2015 no_write: Vec::new(),
2016 daimon_of: String::new(),
2017 })
2018 }
2019
2020 /// A registry holding one tool, so a turn takes the agentic path.
2021 fn one_tool() -> crate::tools::ToolRegistry {
2022 let mut r = no_tools();
2023 r.tools = vec![crate::tools::Tool::FileWrite];
2024 r
2025 }
2026
2027 // ── A picture in front of a model that will not take one ────────────
2028
2029 /// The one-pixel PNG whose base64 `src/llm.rs` documents.
2030 ///
2031 /// The same bytes, so a test asserting the picture did NOT travel is naming a string a
2032 /// provider published rather than one this file invented.
2033 const COVER_PNG_B64: &str = "iVBORw0KGgoAAAANSUhEUgAAAAEAAAABCAIAAACQd1PeAAAADElEQVR4nG\
2034 P4z8AAAAMBAQDJ/pLvAAAAAElFTkSuQmCC";
2035
2036 /// A registry holding `file_read` over a scratch workspace with `cover.png` in it.
2037 fn image_tools() -> crate::tools::ToolRegistry {
2038 let dir = match oxedyne_fe2o3_test::scratch::scratch_dir("daimond_agent_vision") {
2039 Ok(d) => d,
2040 Err(e) => panic!("a scratch directory: {}", e),
2041 };
2042 let png = match oxedyne_fe2o3_text::base64::decode(COVER_PNG_B64) {
2043 Ok(b) => b,
2044 Err(e) => panic!("the documented base64 must decode: {}", e),
2045 };
2046 if let Err(e) = std::fs::write(dir.join("cover.png"), &png) {
2047 panic!("the fixture picture could not be written: {}", e);
2048 }
2049 let ws = match crate::workspace::Workspace::new(dir) {
2050 Ok(w) => w,
2051 Err(e) => panic!("a scratch workspace: {}", e),
2052 };
2053 crate::tools::ToolRegistry::new(vec![crate::tools::Tool::FileRead],
2054 crate::tools::ToolContext {
2055 workspace: ws,
2056 executor: crate::executor::Executor::local_default(),
2057 cwd: String::new(),
2058 path_prefix: String::new(),
2059 root: crate::tools::FileRoot::Workspace,
2060 read_seen: crate::tools::new_read_cache(),
2061 no_write: Vec::new(),
2062 daimon_of: String::new(),
2063 })
2064 }
2065
2066 /// One streamed round that asks to LOOK at the fixture picture.
2067 fn asks_to_look() -> crate::llm::tests::Reply {
2068 crate::llm::tests::Reply::Sse {
2069 chunks: vec![
2070 "data: {\"choices\":[{\"delta\":{\"tool_calls\":[{\"index\":0,\"id\":\"c1\",\
2071 \"type\":\"function\",\"function\":{\"name\":\"file_read\",\"arguments\":\
2072 \"{\\\"path\\\":\\\"cover.png\\\",\\\"as\\\":\\\"image\\\"}\"}}]}}]}\n\n"
2073 .to_string(),
2074 "data: {\"choices\":[{\"delta\":{},\"finish_reason\":\"tool_calls\"}],\
2075 \"usage\":{\"prompt_tokens\":10,\"completion_tokens\":5}}\n\n".to_string(),
2076 "data: [DONE]\n\n".to_string(),
2077 ],
2078 reset_after: None,
2079 }
2080 }
2081
2082 /// A plain streamed answer, ending the turn.
2083 fn plain_answer() -> crate::llm::tests::Reply {
2084 crate::llm::tests::Reply::Sse {
2085 chunks: vec![
2086 "data: {\"choices\":[{\"delta\":{\"content\":\"Done\"}}]}\n\n".to_string(),
2087 "data: {\"choices\":[],\"usage\":{\"prompt_tokens\":11,\"completion_tokens\":2}}\n\n"
2088 .to_string(),
2089 "data: [DONE]\n\n".to_string(),
2090 ],
2091 reset_after: None,
2092 }
2093 }
2094
2095 /// A refusal that says nothing about images, which is what a real one often does.
2096 fn refuses() -> crate::llm::tests::Reply {
2097 crate::llm::tests::Reply::Http {
2098 status: 404, reason: "Not Found", headers: Vec::new(),
2099 body: "{\"error\":{\"message\":\"No endpoint found\"}}".to_string(),
2100 }
2101 }
2102
2103 /// Every `Unseeable` in a run, as `(images, model)`.
2104 fn unseeable(events: &[AgentEvent]) -> Vec<(usize, String)> {
2105 events.iter().filter_map(|e| match e {
2106 AgentEvent::Unseeable { images, model } => Some((*images, model.clone())),
2107 _ => None,
2108 }).collect()
2109 }
2110
2111 #[tokio::test]
2112 async fn test_a_picture_for_a_model_known_blind_is_announced_and_left_out_00() {
2113 // The DECLARED half: `model_can_see` refuses this family before any request, so the
2114 // picture is taken out of the tool reply and never reaches the session at all. What it
2115 // leaves behind names the file, so the same read on a sighted endpoint still works.
2116 let (port, _seen) = crate::llm::tests::start_stub(vec![
2117 asks_to_look(),
2118 plain_answer(),
2119 ]).await;
2120 let mut llm = crate::llm::tests::stub_client(port);
2121 llm.model = fmt!("openai/gpt-3.5-turbo-0125");
2122 let a = Agent::new(llm, "You are Daimond.");
2123 let registry = image_tools();
2124 let mut session = Session::new(fmt!("s"), fmt!("look"), fmt!("openai/gpt-3.5-turbo-0125"));
2125 let mut events: Vec<AgentEvent> = Vec::new();
2126 let _ = a.run_turn(&mut session, fmt!("what is on the cover"), &registry,
2127 &mut |ev| events.push(ev)).await;
2128
2129 let said = unseeable(&events);
2130 assert_eq!(said.len(), 1, "the picture was left out and nothing said so: {:?}", events);
2131 assert_eq!(said[0].0, 1, "the count of pictures left out is wrong");
2132 assert!(said[0].1.contains("gpt-3.5"), "the model was not named: {}", said[0].1);
2133
2134 let tool_text = session.messages.iter()
2135 .filter(|m| m.role() == "tool")
2136 .map(|m| m.text())
2137 .collect::<Vec<_>>()
2138 .join("\n");
2139 assert!(!tool_text.is_empty(), "the tool reply never reached the session");
2140 assert!(!tool_text.contains(COVER_PNG_B64), "the bytes went into the session anyway");
2141 assert!(tool_text.contains("cannot be shown"),
2142 "the elision does not say why: {}", tool_text);
2143 assert!(tool_text.contains("cover.png"), "the file is not named in its place");
2144 assert!(!session.messages.iter().any(|m| m.content().has_image()),
2145 "an image survived into the stored conversation");
2146 }
2147
2148 #[tokio::test]
2149 async fn test_a_model_learned_blind_mid_turn_is_announced_once_00() {
2150 // The LEARNED half: nothing declares this model blind, so the picture goes out, the
2151 // endpoint refuses it, and `stream_turn` strips and retries. That retry is the only
2152 // moment the app knows it is on the wrong model, and it used to pass in silence.
2153 let (port, _seen) = crate::llm::tests::start_stub(vec![
2154 asks_to_look(),
2155 refuses(),
2156 plain_answer(),
2157 ]).await;
2158 let a = Agent::new(crate::llm::tests::stub_client(port), "You are Daimond.");
2159 let registry = image_tools();
2160 let mut session = Session::new(fmt!("s"), fmt!("look"), fmt!("anthropic/claude-opus-5"));
2161 let mut events: Vec<AgentEvent> = Vec::new();
2162 let _ = a.run_turn(&mut session, fmt!("what is on the cover"), &registry,
2163 &mut |ev| events.push(ev)).await;
2164
2165 let said = unseeable(&events);
2166 assert_eq!(said.len(), 1,
2167 "a refusal learned mid-turn was announced {} time(s): {:?}", said.len(), events);
2168 assert!(said[0].1.contains("claude-opus-5"), "the model was not named: {}", said[0].1);
2169 }
2170
2171 #[tokio::test]
2172 async fn test_a_fold_reaches_the_user_as_a_fold_00() {
2173 // It used to borrow the tool surface -- a ToolCall and a ToolResult both named
2174 // `context_compaction` -- so the browser drew the app's own lossy edit of the
2175 // user's conversation as an action row the model had taken. A fold is neither a
2176 // tool nor prose the model produced, and it now says so in its own variant.
2177 let a = dead_agent();
2178 a.set_context_window(8_192);
2179 let mut session = Session::new(fmt!("s1"), fmt!("long"), fmt!("model"));
2180 for i in 0..40 {
2181 session.messages.push(ChatMessage::user(fmt!("step {}", i)));
2182 session.messages.push(ChatMessage::Assistant {
2183 content: MessageContent::text("x".repeat(2_000)), tool_calls: Vec::new(),
2184 });
2185 }
2186 let registry = no_tools();
2187 let mut events: Vec<AgentEvent> = Vec::new();
2188 let _ = a.run_turn(&mut session, fmt!("carry on"), &registry,
2189 &mut |ev| events.push(ev)).await;
2190
2191 let folds: Vec<&AgentEvent> = events.iter()
2192 .filter(|e| matches!(e, AgentEvent::Compacted { .. })).collect();
2193 assert_eq!(folds.len(), 1, "the conversation was {} events and none was a fold",
2194 events.len());
2195 match folds[0] {
2196 AgentEvent::Compacted { folded, kept, note } => {
2197 assert!(*folded > 0, "a fold that folded nothing");
2198 assert_eq!(*kept, session.messages.len(),
2199 "the count does not match what the session now holds");
2200 assert!(note.contains("Folded"), "{}", note);
2201 }
2202 other => panic!("{:?}", other),
2203 }
2204 // And nothing on the borrowed surface, which a client draws as a collapsible
2205 // action row: a turn with no tools registered has no tool events at all.
2206 assert!(!events.iter().any(|e| matches!(e,
2207 AgentEvent::ToolCall { .. } | AgentEvent::ToolResult { .. })),
2208 "the fold is still announcing itself as a tool");
2209 }
2210
2211 // ── What a tool call came to travels with the event ─────────────────
2212
2213 /// Three real calls in one round: one that works, one the fence stops, one that is not there.
2214 ///
2215 /// Written as a script for the stub provider rather than as three hand-built strings, because
2216 /// a fixture that states the reply AND the expected reading agrees with itself whatever the
2217 /// tool layer does. What is under test is that the app carries the tool layer's own verdict,
2218 /// so the verdict has to come from the tool layer.
2219 fn three_calls() -> crate::llm::tests::Reply {
2220 crate::llm::tests::Reply::Sse {
2221 chunks: vec![
2222 // Works: a relative path inside the scratch workspace.
2223 "data: {\"choices\":[{\"delta\":{\"tool_calls\":[{\"index\":0,\"id\":\"c0\",\
2224 \"type\":\"function\",\"function\":{\"name\":\"file_write\",\"arguments\":\
2225 \"{\\\"path\\\":\\\"notes/ok.txt\\\",\\\"content\\\":\\\"hi\\\"}\"}}]}}]}\n\n"
2226 .to_string(),
2227 // Refused: an absolute path on the machine, which `guard` stops before any write.
2228 "data: {\"choices\":[{\"delta\":{\"tool_calls\":[{\"index\":1,\"id\":\"c1\",\
2229 \"type\":\"function\",\"function\":{\"name\":\"file_write\",\"arguments\":\
2230 \"{\\\"path\\\":\\\"/etc/passwd\\\",\\\"content\\\":\\\"x\\\"}\"}}]}}]}\n\n"
2231 .to_string(),
2232 // Failed: a real tool, not registered here, so `dispatch` composes an error line.
2233 "data: {\"choices\":[{\"delta\":{\"tool_calls\":[{\"index\":2,\"id\":\"c2\",\
2234 \"type\":\"function\",\"function\":{\"name\":\"spawn_agent\",\"arguments\":\
2235 \"{}\"}}]}}]}\n\n".to_string(),
2236 "data: {\"choices\":[{\"delta\":{},\"finish_reason\":\"tool_calls\"}]}\n\n"
2237 .to_string(),
2238 "data: [DONE]\n\n".to_string(),
2239 ],
2240 reset_after: None,
2241 }
2242 }
2243
2244 /// Every `ToolResult` in a run, as `(name, outcome, text)`.
2245 fn tool_results(events: &[AgentEvent]) -> Vec<(String, CallOutcome, String)> {
2246 events.iter().filter_map(|e| match e {
2247 AgentEvent::ToolResult { name, result, outcome } =>
2248 Some((name.clone(), *outcome, result.clone())),
2249 _ => None,
2250 }).collect()
2251 }
2252
2253 #[tokio::test]
2254 async fn test_a_tool_result_carries_what_the_call_came_to_00() {
2255 // The app used to flatten the tool layer's verdict into the reply's opening word and let
2256 // four browser consumers read it back out of the prose, each with its own reading. One of
2257 // them tested for "Error" alone -- so a refusal, which opens "Refused", was drawn as a
2258 // completed step, journalled as a success and reported to the Optimiser as a tool that
2259 // worked.
2260 let (port, _seen) = crate::llm::tests::start_stub(vec![
2261 three_calls(),
2262 plain_answer(),
2263 ]).await;
2264 let mut llm = crate::llm::tests::stub_client(port);
2265 llm.retry.max_attempts = 1;
2266 let a = Agent::new(llm, "You are Daimond.");
2267 a.set_max_rounds(2);
2268
2269 let registry = one_tool();
2270 let mut session = Session::new(fmt!("s1"), fmt!("three"), fmt!("model"));
2271 let mut events: Vec<AgentEvent> = Vec::new();
2272 let _ = a.run_turn(&mut session, fmt!("do three things"), &registry,
2273 &mut |ev| events.push(ev)).await;
2274
2275 let got = tool_results(&events);
2276 assert_eq!(3, got.len(), "three calls went out and {} results came back: {:?}",
2277 got.len(), events.iter().map(|e| fmt!("{:?}", e)).collect::<Vec<_>>());
2278
2279 // Each verdict is the tool layer's, on a reply the tool layer actually composed.
2280 assert_eq!(CallOutcome::Done, got[0].1, "a write inside the workspace: {}", got[0].2);
2281 assert_eq!(CallOutcome::Refused, got[1].1, "a write the fence stopped: {}", got[1].2);
2282 assert_eq!(CallOutcome::Failed, got[2].1, "a tool that is not here: {}", got[2].2);
2283
2284 // THE DEFECT, NAMED. The refused reply does not contain the word a reader looking for
2285 // failure would look for, which is why reading the prose lost it.
2286 assert!(got[1].2.trim_start().starts_with("Refused"),
2287 "the refusal no longer opens with its own word: {}", got[1].2);
2288 assert!(!got[1].2.starts_with("Error"),
2289 "if a refusal ever opens 'Error' this test stops proving anything: {}", got[1].2);
2290 assert!(got[2].2.trim_start().starts_with("Error"), "{}", got[2].2);
2291
2292 // And the refused call is a refusal and not a write with a warning on it: if the fence
2293 // ever let this path through, the reply would be `file_write`'s own success sentence and
2294 // the outcome above would be Done -- which is the whole failure, arriving one layer down.
2295 assert!(!got[1].2.contains("Wrote"),
2296 "the fence let the write through: {}", got[1].2);
2297
2298 // The wire. Exactly the three words, in the key the browser reads.
2299 let maps: Vec<DaticleMap> = events.iter()
2300 .filter(|e| matches!(e, AgentEvent::ToolResult { .. }))
2301 .map(|e| e.to_datmap()).collect();
2302 let words = ["done", "refused", "failed"];
2303 for (i, m) in maps.iter().enumerate() {
2304 assert_eq!(Some(&dat!("tool_result")), m.get(&dat!("type")));
2305 assert_eq!(Some(&dat!(words[i])), m.get(&dat!("outcome")),
2306 "call {} spelled its outcome {:?}", i, m.get(&dat!("outcome")));
2307 // The two fields that were always there are still there: this is a field added,
2308 // not a shape changed. The event carries the result TEXT, and not the image.
2309 assert_eq!(Some(&dat!(got[i].0.clone())), m.get(&dat!("name")));
2310 assert_eq!(Some(&dat!(got[i].2.clone())), m.get(&dat!("content")));
2311 }
2312 }
2313
2314
2315 // ── The round's dispatch: what runs together, and in what order ──────
2316 //
2317 // The rule itself is `agent::batch`, and its own tests cover which tools may share a batch.
2318 // What is covered HERE is the round loop around it: that batching changes nothing a reader of
2319 // the event stream or of the conversation can see, and in particular that it does not change
2320 // the one thing the page depends on.
2321
2322 /// A registry over a scratch workspace holding three small files, with the file tools a batch
2323 /// is made of and the write tool that breaks one.
2324 fn batch_tools() -> crate::tools::ToolRegistry {
2325 let dir = match oxedyne_fe2o3_test::scratch::scratch_dir("daimond_agent_batch") {
2326 Ok(d) => d,
2327 Err(e) => panic!("a scratch directory: {}", e),
2328 };
2329 for (name, body) in [("one.txt", "first"), ("two.txt", "second"), ("three.txt", "third")] {
2330 if let Err(e) = std::fs::write(dir.join(name), body) {
2331 panic!("the fixture file '{}' could not be written: {}", name, e);
2332 }
2333 }
2334 let ws = match crate::workspace::Workspace::new(dir) {
2335 Ok(w) => w,
2336 Err(e) => panic!("a scratch workspace: {}", e),
2337 };
2338 let mut r = crate::tools::ToolRegistry::new(Vec::new(), crate::tools::ToolContext {
2339 workspace: ws,
2340 executor: crate::executor::Executor::local_default(),
2341 cwd: String::new(),
2342 path_prefix: String::new(),
2343 root: crate::tools::FileRoot::Workspace,
2344 read_seen: crate::tools::new_read_cache(),
2345 no_write: Vec::new(),
2346 daimon_of: String::new(),
2347 });
2348 r.tools = vec![crate::tools::Tool::FileRead, crate::tools::Tool::FileWrite];
2349 r
2350 }
2351
2352 /// Every tool call and tool result in a run, in the order they reached the page, as
2353 /// `("call"|"result", name)`.
2354 fn call_stream(events: &[AgentEvent]) -> Vec<(&'static str, String)> {
2355 events.iter().filter_map(|e| match e {
2356 AgentEvent::ToolCall { name, .. } => Some(("call", name.clone())),
2357 AgentEvent::ToolResult { name, .. } => Some(("result", name.clone())),
2358 _ => None,
2359 }).collect()
2360 }
2361
2362 /// Run one turn of the given tool calls against a scripted provider, and hand back every
2363 /// event it produced.
2364 async fn round_of(
2365 calls: &[(&str, &str)],
2366 registry: &ToolRegistry,
2367 )
2368 -> Vec<AgentEvent>
2369 {
2370 let (port, _seen) = crate::llm::tests::start_stub(vec![
2371 tool_round(calls),
2372 plain_answer(),
2373 ]).await;
2374 let mut llm = crate::llm::tests::stub_client(port);
2375 llm.retry.max_attempts = 1;
2376 let a = Agent::new(llm, "You are Daimond.");
2377 a.set_max_rounds(2);
2378 let mut session = Session::new(fmt!("s1"), fmt!("batch"), fmt!("model"));
2379 let mut events: Vec<AgentEvent> = Vec::new();
2380 let _ = a.run_turn(&mut session, fmt!("do the work"), registry,
2381 &mut |ev| events.push(ev)).await;
2382 events
2383 }
2384
2385 /// A batch of reads answers in the order the MODEL gave, not the order the reads finished.
2386 ///
2387 /// The three files hold different words, so an answer that came back out of order is caught
2388 /// by its content and not merely by its position.
2389 #[tokio::test]
2390 async fn test_a_batch_of_reads_answers_in_the_models_order_00() {
2391 let registry = batch_tools();
2392 let events = round_of(&[
2393 ("file_read", "{\"path\":\"one.txt\"}"),
2394 ("file_read", "{\"path\":\"two.txt\"}"),
2395 ("file_read", "{\"path\":\"three.txt\"}"),
2396 ], &registry).await;
2397
2398 let got = tool_results(&events);
2399 assert_eq!(3, got.len(), "three reads went out and {} results came back", got.len());
2400 for (i, word) in ["first", "second", "third"].iter().enumerate() {
2401 assert!(got[i].2.contains(word),
2402 "result {} is not the answer to call {}: {}", i, i, got[i].2);
2403 }
2404 }
2405
2406 /// ONE CALL ANNOUNCED, THEN ITS RESULT, WHATEVER RAN TOGETHER UNDERNEATH.
2407 ///
2408 /// This is the page's contract and not a tidiness: `www/js/daimond.js` keeps one
2409 /// `pendingTool` and one `pendingCallId` for the chat, the worker dock and the daimon alike,
2410 /// and files each result against whichever call was announced last. Announce a whole batch
2411 /// up front and the first result is filed under the last call while the rest are dropped --
2412 /// out of the transcript and out of the write-ahead journal both.
2413 #[tokio::test]
2414 async fn test_a_round_announces_one_call_at_a_time_00() {
2415 let registry = batch_tools();
2416 let events = round_of(&[
2417 ("file_read", "{\"path\":\"one.txt\"}"),
2418 ("file_read", "{\"path\":\"two.txt\"}"),
2419 ("file_read", "{\"path\":\"three.txt\"}"),
2420 ], &registry).await;
2421
2422 let stream = call_stream(&events);
2423 assert_eq!(6, stream.len(), "three calls did not produce three pairs: {:?}", stream);
2424 for (i, (kind, _)) in stream.iter().enumerate() {
2425 let want = if i % 2 == 0 { "call" } else { "result" };
2426 assert_eq!(want, *kind, "the stream stopped alternating at {}: {:?}", i, stream);
2427 }
2428 }
2429
2430 /// A WRITE KEEPS ITS PLACE AMONG THE READS AROUND IT.
2431 ///
2432 /// The order the model gave is the order it intended, and a write is what separates the reads
2433 /// before it from the reads after it. Asserted on the stream the page sees rather than on the
2434 /// batching, so it holds however the batching is later rewritten.
2435 #[tokio::test]
2436 async fn test_a_write_keeps_its_place_among_the_reads_00() {
2437 let registry = batch_tools();
2438 let events = round_of(&[
2439 ("file_read", "{\"path\":\"one.txt\"}"),
2440 ("file_read", "{\"path\":\"two.txt\"}"),
2441 ("file_write", "{\"path\":\"four.txt\",\"content\":\"fourth\"}"),
2442 ("file_read", "{\"path\":\"three.txt\"}"),
2443 ], &registry).await;
2444
2445 let names: Vec<String> = call_stream(&events).into_iter()
2446 .filter(|(kind, _)| *kind == "call")
2447 .map(|(_, name)| name)
2448 .collect();
2449 assert_eq!(vec!["file_read", "file_read", "file_write", "file_read"], names,
2450 "the round did not run the calls in the order the model gave them");
2451 }
2452
2453 // ── The claims audit, and how a turn says it ended ──────────────────
2454
2455 /// A registry over its own scratch workspace, holding the tools named.
2456 fn tools_over_scratch(tools: Vec<crate::tools::Tool>) -> crate::tools::ToolRegistry {
2457 let mut r = no_tools();
2458 r.tools = tools;
2459 r
2460 }
2461
2462 /// One round of tool calls, as a provider streams them.
2463 ///
2464 /// Scripted rather than hand-built so a test states the CALL and lets the tool layer decide
2465 /// what becomes of it. What is under test is that the app carries the tool layer's own
2466 /// verdict and the model's own arguments, so both have to come from the real thing.
2467 ///
2468 /// # Arguments
2469 /// * `calls` - Each call as its wire name and its arguments JSON.
2470 fn tool_round(calls: &[(&str, &str)]) -> crate::llm::tests::Reply {
2471 let mut chunks = Vec::new();
2472 for (i, (name, args)) in calls.iter().enumerate() {
2473 // The arguments ride as a JSON string inside a JSON object, so they are escaped once
2474 // here and unescaped once by the stream accumulator.
2475 let esc = args.replace('\\', "\\\\").replace('"', "\\\"");
2476 chunks.push(fmt!(
2477 "data: {{\"choices\":[{{\"delta\":{{\"tool_calls\":[{{\"index\":{},\
2478 \"id\":\"c{}\",\"type\":\"function\",\"function\":{{\"name\":\"{}\",\
2479 \"arguments\":\"{}\"}}}}]}}}}]}}\n\n",
2480 i, i, name, esc));
2481 }
2482 chunks.push(
2483 "data: {\"choices\":[{\"delta\":{},\"finish_reason\":\"tool_calls\"}]}\n\n".to_string());
2484 chunks.push("data: [DONE]\n\n".to_string());
2485 crate::llm::tests::Reply::Sse { chunks, reset_after: None }
2486 }
2487
2488 /// Run one turn against a scripted provider, and hand back what the turn came to.
2489 async fn ran(
2490 script: Vec<crate::llm::tests::Reply>,
2491 registry: &ToolRegistry,
2492 max_rounds: usize,
2493 )
2494 -> TurnEnding
2495 {
2496 let (port, _seen) = crate::llm::tests::start_stub(script).await;
2497 let mut llm = crate::llm::tests::stub_client(port);
2498 llm.retry.max_attempts = 1;
2499 let a = Agent::new(llm, "You are Daimond.");
2500 a.set_max_rounds(max_rounds);
2501 let mut session = Session::new(fmt!("s1"), fmt!("audit"), fmt!("model"));
2502 let _ = a.run_turn(&mut session, fmt!("do the work"), registry, &mut |_| {}).await;
2503 match a.ending() {
2504 Some(e) => e,
2505 None => panic!("a turn ran and said nothing at all about how it ended"),
2506 }
2507 }
2508
2509 /// A turn's byte allowance starts AT THE TURN, wherever the turn came from.
2510 ///
2511 /// Nothing resets itself. The ledger lives on a [`ToolContext`] that OUTLIVES the turn -- the
2512 /// browser chat builds one for the life of the app, and a Diamond's daimon shares that very
2513 /// one -- so if the reset sat anywhere but at the entry to a turn, a single long read would
2514 /// narrow every turn after it for the rest of the session.
2515 #[tokio::test]
2516 async fn test_a_turn_starts_its_byte_allowance_at_the_turn() {
2517 let a = dead_agent();
2518 let registry = no_tools();
2519 registry.ctx.charge_spend(500_000);
2520 assert_eq!(0, registry.ctx.spend_left(), "the fixture did not spend the allowance");
2521 let mut session = Session::new(fmt!("s1"), fmt!("budget"), fmt!("model"));
2522 // The provider is a port nothing is listening on, so the turn fails at once -- which is
2523 // the point: what is under test happens BEFORE the request goes out.
2524 let _ = a.run_turn(&mut session, fmt!("read the file"), &registry, &mut |_| {}).await;
2525 assert!(!registry.ctx.spend_is_short(),
2526 "this turn began already short, having spent nothing of its own");
2527 assert_eq!(crate::tools::TURN_SPEND_BUDGET, registry.ctx.spend_left(),
2528 "the last turn's spending is still charged to this one");
2529 }
2530
2531 #[tokio::test]
2532 async fn test_a_refused_tool_is_not_a_completed_step_00() {
2533 // Three real calls: one that works, one the fence stops, one that is not registered. The
2534 // audit tallies the tool layer's own verdict -- it does not re-read the sentence, which is
2535 // how a refusal, whose reply opens "Refused" rather than "Error", was drawn as a completed
2536 // step, journalled as a success and reported to the Optimiser as a tool that had worked.
2537 let registry = one_tool();
2538 let end = ran(vec![three_calls(), plain_answer()], &registry, 2).await;
2539
2540 assert_eq!(3, end.calls, "{:?}", end);
2541 assert_eq!(1, end.refused, "the call the fence stopped was booked as work: {:?}", end);
2542 assert_eq!(1, end.failed, "{:?}", end);
2543 assert!(end.unaccounted(),
2544 "a turn holding a refusal and a breakage reported nothing to answer for: {:?}", end);
2545 // And the ending itself is the ordinary one: the turn ANSWERED. What is unaccounted for
2546 // is a separate question from how the turn finished, and collapsing the two would make
2547 // every turn with one refused call look like a turn that fell over.
2548 assert_eq!(TurnEnd::Answered, end.how, "{:?}", end);
2549 }
2550
2551 #[tokio::test]
2552 async fn test_a_file_a_call_said_it_left_and_did_not_is_reported_00() {
2553 // The claim is not in the prose. It is in the ARGUMENTS of the call, which name a path,
2554 // and after the turn a named path that is not there is a fact rather than a reading.
2555 //
2556 // Two calls, identical but for the file: one names a path that is on the store and one
2557 // names a path that is not. The pair is what makes either half mean anything -- an audit
2558 // that reported both, or neither, would satisfy a test of only one.
2559 let registry = one_tool();
2560 let there = registry.ctx.workspace.root().join("there.txt");
2561 if let Err(e) = std::fs::write(&there, "x") {
2562 panic!("the fixture file must exist for the control half to mean anything: {}", e);
2563 }
2564 let mut claims = Claims::default();
2565 claims.record("file_write", r#"{"path":"there.txt","content":"x"}"#, CallOutcome::Done);
2566 claims.record("file_write", r#"{"path":"ghost.txt","content":"x"}"#, CallOutcome::Done);
2567
2568 let a = dead_agent();
2569 let end = a.audit(TurnEnd::Answered, 1, &claims, Some(&registry)).await;
2570 assert_eq!(vec![fmt!("ghost.txt")], end.missing,
2571 "the file that is there, the file that is not, or both: {:?}", end);
2572 assert!(end.unaccounted(), "{:?}", end);
2573 }
2574
2575 #[tokio::test]
2576 async fn test_a_call_that_never_ran_claims_nothing_00() {
2577 // A REFUSED WRITE IS NOT A FILE ANYBODY SHOULD BE LOOKING FOR, and neither is one whose
2578 // call was cut off at the output limit -- that reply opens "Error", so the call is booked
2579 // `Failed`. Both carry a perfectly readable `path`; the outcome is what keeps them out of
2580 // the audit, which is the same rule `AgentEvent::ToolResult` carries.
2581 let registry = one_tool();
2582 let mut claims = Claims::default();
2583 claims.record("file_write", r#"{"path":"stopped.txt","content":"x"}"#,
2584 CallOutcome::Refused);
2585 claims.record("file_write", r#"{"path":"src/big.rs","content":"fn main() {"#,
2586 CallOutcome::Failed);
2587 let a = dead_agent();
2588 let end = a.audit(TurnEnd::Answered, 1, &claims, Some(&registry)).await;
2589 assert!(end.missing.is_empty(),
2590 "a turn was told to go looking for a file no tool ever wrote: {:?}", end);
2591 // Counted, though. They are the other half of the audit and the reason the turn has
2592 // something to answer for at all.
2593 assert_eq!(1, end.refused, "{:?}", end);
2594 assert_eq!(1, end.failed, "{:?}", end);
2595 assert!(end.unaccounted(), "{:?}", end);
2596 }
2597
2598 #[tokio::test]
2599 async fn test_a_file_written_and_then_deleted_is_not_reported_missing_00() {
2600 // The turn's LAST word about a path is the turn's word about it. Without that rule every
2601 // scratch file a turn tidied up after itself would be reported as a write that did not
2602 // happen -- which is a warning on a turn that did exactly what it said, and an app that
2603 // does that is teaching its reader to skip the warning that matters.
2604 let registry = tools_over_scratch(vec![
2605 crate::tools::Tool::FileWrite,
2606 crate::tools::Tool::FileDelete,
2607 ]);
2608 let end = ran(vec![
2609 tool_round(&[
2610 ("file_write", r#"{"path":"scratch.txt","content":"working"}"#),
2611 ("file_delete", r#"{"path":"scratch.txt"}"#),
2612 ]),
2613 plain_answer(),
2614 ], &registry, 2).await;
2615
2616 assert_eq!(2, end.calls, "{:?}", end);
2617 assert_eq!(0, end.refused, "{:?}", end);
2618 assert_eq!(0, end.failed, "{:?}", end);
2619 assert!(end.missing.is_empty(), "{:?}", end);
2620 assert!(!end.unaccounted(), "{:?}", end);
2621 // AND THE FILE REALLY IS GONE. Without this the test would pass just as well against a
2622 // delete that silently did nothing, and would be proving that the audit ignores a file
2623 // that is present rather than that it forgives one that was deliberately removed.
2624 assert!(!registry.ctx.workspace.root().join("scratch.txt").exists(),
2625 "the delete did not happen, so nothing here is being tested");
2626 }
2627
2628 #[tokio::test]
2629 async fn test_a_turn_that_ran_a_shell_command_is_not_told_a_file_is_missing_00() {
2630 // A shell command, a build, a verifier run and a worker each act on paths nobody wrote
2631 // down. The app cannot know which of them removed a file, and the honest answer to a
2632 // question it cannot answer is silence -- an audit that guessed here would raise a finding
2633 // on any turn that wrote a file and then ran `cargo test`.
2634 let registry = one_tool();
2635 let a = dead_agent();
2636 let write = r#"{"path":"ghost.txt","content":"x"}"#;
2637
2638 // The control: with no shell in the turn, the missing file IS reported.
2639 let mut alone = Claims::default();
2640 alone.record("file_write", write, CallOutcome::Done);
2641 let end = a.audit(TurnEnd::Answered, 1, &alone, Some(&registry)).await;
2642 assert_eq!(vec![fmt!("ghost.txt")], end.missing,
2643 "the control half does not report the file, so the half below proves nothing: {:?}",
2644 end);
2645
2646 // The same turn, having also run a command.
2647 let mut with_shell = Claims::default();
2648 with_shell.record("file_write", write, CallOutcome::Done);
2649 with_shell.record("shell", r#"{"command":"rm ghost.txt"}"#, CallOutcome::Done);
2650 let end = a.audit(TurnEnd::Answered, 1, &with_shell, Some(&registry)).await;
2651 assert!(end.missing.is_empty(),
2652 "the app claimed to know what a shell command did not do: {:?}", end);
2653
2654 // A command the fence turned away ran nothing, so it explains nothing and silences
2655 // nothing. Without this the whole check would be switched off by any refused shell call.
2656 let mut refused_shell = Claims::default();
2657 refused_shell.record("file_write", write, CallOutcome::Done);
2658 refused_shell.record("shell", r#"{"command":"rm ghost.txt"}"#, CallOutcome::Refused);
2659 let end = a.audit(TurnEnd::Answered, 1, &refused_shell, Some(&registry)).await;
2660 assert_eq!(vec![fmt!("ghost.txt")], end.missing,
2661 "a refused command switched the audit off: {:?}", end);
2662 }
2663
2664 #[tokio::test]
2665 async fn test_a_turn_that_did_what_it_said_has_nothing_to_answer_for_00() {
2666 // THIS MATTERS AS MUCH AS THE FINDING DOES. An audit that raises something on an ordinary
2667 // turn is an audit nobody reads, and then the one turn that needed reading goes past with
2668 // everything else.
2669 let registry = one_tool();
2670 let end = ran(vec![
2671 tool_round(&[("file_write", r#"{"path":"notes/kept.txt","content":"hi"}"#)]),
2672 plain_answer(),
2673 ], &registry, 2).await;
2674
2675 assert_eq!(TurnEnd::Answered, end.how, "{:?}", end);
2676 assert_eq!(1, end.calls, "{:?}", end);
2677 assert_eq!(0, end.refused, "{:?}", end);
2678 assert_eq!(0, end.failed, "{:?}", end);
2679 assert!(end.missing.is_empty(), "{:?}", end);
2680 assert!(!end.unaccounted(),
2681 "an ordinary turn was given something to answer for: {:?}", end);
2682 // The check RAN and passed, rather than being skipped: the file the call named is there.
2683 assert!(registry.ctx.workspace.root().join("notes/kept.txt").exists(),
2684 "the write did not land, so the audit had nothing to be right about");
2685 }
2686
2687 #[tokio::test]
2688 async fn test_a_turn_that_used_none_of_its_tools_says_so_00() {
2689 // THE CASE THIS WAS BUILT FOR. A model announced "let me rewrite from line 43 through 49
2690 // applying the lessons" and then ended its turn having written nothing. The spinner
2691 // stopped, which was correct, and the user was left saying "I have no visibility on what
2692 // occurred here". Nothing was technically wrong, so nothing was said.
2693 //
2694 // The figures say it without reading a word of the reply: tools on the table, no call
2695 // made. What the model promised is the reader's business; whether it did anything is the
2696 // app's, and this is the app's answer.
2697 let registry = one_tool();
2698 let end = ran(vec![plain_answer()], &registry, 4).await;
2699
2700 assert_eq!(TurnEnd::Answered, end.how, "{:?}", end);
2701 assert_eq!(0, end.calls, "{:?}", end);
2702 assert_eq!(1, end.rounds, "{:?}", end);
2703 assert!(end.offered > 0,
2704 "a turn that held no tools is a different case entirely: {:?}", end);
2705 // AND IT IS NOT A WARNING. Nothing went wrong -- the model may simply have answered a
2706 // question -- so the ending reports the shape of the turn and claims no fault.
2707 assert!(!end.unaccounted(),
2708 "a turn that merely answered was reported as having something wrong: {:?}", end);
2709 }
2710
2711 #[tokio::test]
2712 async fn test_a_turn_that_spent_its_round_budget_says_so_00() {
2713 // The round limit was already announced, in the app's own voice. What it did not do was
2714 // say what the turn had MADE OF that budget, which is the figure that tells a reader
2715 // whether to raise the ceiling or stop the work.
2716 let registry = one_tool();
2717 let end = ran(vec![
2718 tool_round(&[("file_write", r#"{"path":"a.txt","content":"1"}"#)]),
2719 tool_round(&[("file_write", r#"{"path":"b.txt","content":"2"}"#)]),
2720 ], &registry, 2).await;
2721
2722 assert_eq!(TurnEnd::Capped, end.how, "{:?}", end);
2723 assert_eq!(2, end.rounds, "{:?}", end);
2724 assert_eq!(2, end.calls, "{:?}", end);
2725 assert!(end.missing.is_empty(), "{:?}", end);
2726 }
2727
2728 #[tokio::test]
2729 async fn test_a_turn_the_provider_ended_still_says_how_it_ended_00() {
2730 // The ending that was least visible of all: the turn returns an error, and until now
2731 // nothing recorded that a turn had finished at all.
2732 let a = dead_agent();
2733 let registry = no_tools();
2734 let mut session = Session::new(fmt!("s1"), fmt!("dead"), fmt!("model"));
2735 let out = a.run_turn(&mut session, fmt!("hello"), &registry, &mut |_| {}).await;
2736 assert!(out.is_err(), "the stub agent reached a provider, so this proves nothing");
2737
2738 let end = match a.ending() {
2739 Some(e) => e,
2740 None => panic!("a turn ended in an error and said nothing about how it ended"),
2741 };
2742 assert_eq!(TurnEnd::Failed, end.how, "{:?}", end);
2743 assert!(!end.unaccounted(),
2744 "a turn that never reached a tool has no tool to answer for: {:?}", end);
2745 }
2746
2747 #[tokio::test]
2748 async fn test_a_turn_that_said_nothing_says_that_it_said_nothing_00() {
2749 // The one ending the app already explained, kept explained -- and now in the same words
2750 // every other ending uses, so a reader does not have to know which of two mechanisms
2751 // produced the line in front of them.
2752 let quiet = crate::llm::tests::Reply::Sse {
2753 chunks: vec![
2754 "data: {\"choices\":[{\"delta\":{}}]}\n\n".to_string(),
2755 "data: [DONE]\n\n".to_string(),
2756 ],
2757 reset_after: None,
2758 };
2759 let registry = no_tools();
2760 let end = ran(vec![quiet], &registry, 1).await;
2761 assert_eq!(TurnEnd::Silent, end.how, "{:?}", end);
2762 assert_eq!(0, end.offered, "a pure chat holds no tools: {:?}", end);
2763 }
2764
2765 #[test]
2766 fn test_an_ending_travels_as_one_of_five_words_00() {
2767 // Spelled once, for the reason `CallOutcome::wire` is: the browser knows these words and
2768 // no others, so a second speller would not fail loudly -- it would draw an ending nobody
2769 // recognises, which is the silence this whole mechanism replaces.
2770 let all = [
2771 (TurnEnd::Answered, "answered"),
2772 (TurnEnd::Stopped, "stopped"),
2773 (TurnEnd::Capped, "capped"),
2774 (TurnEnd::Silent, "silent"),
2775 (TurnEnd::Failed, "failed"),
2776 ];
2777 for (end, word) in all {
2778 assert_eq!(word, end.wire(), "{:?}", end);
2779 }
2780 // `Stopped` is the browser's alone: the native transport has no cancellation path, so
2781 // `stream_sse` always reports a stream that ran to its end. It is spelled here so the two
2782 // halves of the seam agree about a word only one of them can produce.
2783 assert_eq!("stopped", TurnEnd::Stopped.wire());
2784 }
2785
2786 // ── A reply that ran out of room ────────────────────────────────────
2787
2788 #[test]
2789 fn test_a_call_cut_at_the_limit_is_told_the_truth_rather_than_a_parse_error_00() {
2790 // The failure this replaces is the single most confusing one in a coding
2791 // session: the model asks to write a long file, the reply is cut in the middle
2792 // of the arguments, the dispatcher says the JSON was bad, and the model -- told
2793 // its own JSON was the mistake -- writes exactly the same thing again.
2794 let cut = r#"{"path":"src/big.rs","content":"fn main() {\n let x ="#;
2795 assert!(cut_short(true, cut), "a cut object was not recognised as cut");
2796 let note = truncated_call_note(8_192);
2797 assert!(note.contains("cut off at the output limit"), "{}", note);
2798 assert!(note.contains("8192"), "the figure is what makes it actionable: {}", note);
2799 assert!(note.contains("smaller pieces"), "{}", note);
2800 assert!(note.contains("Nothing was changed"),
2801 "a model that thinks a half-written file exists will read it back: {}", note);
2802 // And it must not read as the model's own mistake, which is what sends it round
2803 // the same loop again.
2804 assert!(!note.to_lowercase().contains("invalid json"), "{}", note);
2805 assert!(!note.to_lowercase().contains("malformed"), "{}", note);
2806 }
2807
2808 #[tokio::test]
2809 async fn test_every_turn_says_how_it_ended_and_says_it_before_done_00() {
2810 // The owner watched a model announce work and then end its turn having done none.
2811 // The spinner stopped, correctly -- the turn HAD ended -- and there was no way to
2812 // tell that from a turn that finished. Three endings were explained before this
2813 // and every other one was silence, so the assertion is that there is no longer any
2814 // such thing as a turn that ends without saying so.
2815 use crate::llm::tests::{start_stub, stub_client, Reply};
2816 let (port, _seen) = start_stub(vec![Reply::Sse {
2817 chunks: vec![
2818 "data: {\"choices\":[{\"delta\":{\"content\":\"I will rewrite it.\"}}]}\n\n"
2819 .to_string(),
2820 "data: {\"choices\":[{\"delta\":{},\"finish_reason\":\"stop\"}]}\n\n"
2821 .to_string(),
2822 "data: [DONE]\n\n".to_string(),
2823 ],
2824 reset_after: None,
2825 }]).await;
2826 let mut llm = stub_client(port);
2827 llm.retry.max_attempts = 1;
2828 let a = Agent::new(llm, "You are Daimond.");
2829 a.set_max_rounds(2);
2830
2831 let registry = one_tool();
2832 let mut session = Session::new(fmt!("s1"), fmt!("ended"), fmt!("model"));
2833 let mut events: Vec<AgentEvent> = Vec::new();
2834 let _ = a.run_turn(&mut session, fmt!("rewrite lines 43 to 49"), &registry,
2835 &mut |ev| events.push(ev)).await;
2836
2837 let end = events.iter().position(|e| matches!(e, AgentEvent::Ended { .. }));
2838 assert!(end.is_some(), "a turn ended and said nothing about how: {:?}",
2839 events.iter().map(|e| fmt!("{:?}", e)).collect::<Vec<_>>());
2840 // ORDER, because a reader draws the closing line under the turn: an ending that
2841 // arrived after `Done` would be drawn under the turn after it.
2842 if let Some(done) = events.iter().position(|e| matches!(e, AgentEvent::Done)) {
2843 assert!(end < Some(done), "the ending arrived after the turn was declared done");
2844 }
2845 // Tools were on the table and none was called, which is exactly the owner's case.
2846 // It is reported as a FACT and not as a fault: what the model promised is the
2847 // reader's business, whether it did anything is the app's.
2848 match events.iter().find(|e| matches!(e, AgentEvent::Ended { .. })) {
2849 Some(AgentEvent::Ended { how, offered, calls, .. }) => {
2850 assert_eq!(how, "answered", "a plain stop is not a failure");
2851 assert!(*offered > 0, "the turn was offered no tools, so it proves nothing here");
2852 assert_eq!(*calls, 0, "the stub called nothing, so the count must say so");
2853 }
2854 _ => panic!("no ending to read"),
2855 }
2856 }
2857
2858 #[tokio::test]
2859 async fn test_a_turn_whose_tool_call_was_cut_tells_the_model_and_the_user_00() {
2860 // End to end, against a real TLS server streaming a real SSE body: the provider
2861 // says `finish_reason: "length"` half-way through a `file_write`'s arguments,
2862 // which is exactly what a model asked for a long file does under an output cap.
2863 use crate::llm::tests::{start_stub, stub_client, Reply};
2864 let (port, _seen) = start_stub(vec![Reply::Sse {
2865 chunks: vec![
2866 "data: {\"choices\":[{\"delta\":{\"tool_calls\":[{\"index\":0,\
2867 \"id\":\"call_1\",\"type\":\"function\",\"function\":{\"name\":\
2868 \"file_write\",\"arguments\":\"\"}}]}}]}\n\n".to_string(),
2869 "data: {\"choices\":[{\"delta\":{\"tool_calls\":[{\"index\":0,\
2870 \"function\":{\"arguments\":\"{\\\"path\\\":\\\"big.rs\\\",\
2871 \\\"content\\\":\\\"fn main() {\"}}]}}]}\n\n".to_string(),
2872 "data: {\"choices\":[{\"delta\":{},\"finish_reason\":\"length\"}]}\n\n"
2873 .to_string(),
2874 "data: [DONE]\n\n".to_string(),
2875 ],
2876 reset_after: None,
2877 },
2878 // The round after it, so the turn ends of its own accord: a turn stopped at
2879 // the round limit emits an error of its own, which would mask the question.
2880 Reply::Sse {
2881 chunks: vec![
2882 "data: {\"choices\":[{\"delta\":{\"content\":\"I will split it.\"}}]}\n\n"
2883 .to_string(),
2884 "data: {\"choices\":[{\"delta\":{},\"finish_reason\":\"stop\"}]}\n\n"
2885 .to_string(),
2886 "data: [DONE]\n\n".to_string(),
2887 ],
2888 reset_after: None,
2889 }]).await;
2890 let mut llm = stub_client(port);
2891 llm.retry.max_attempts = 1;
2892 llm.max_tokens = 8_192;
2893 let a = Agent::new(llm, "You are Daimond.");
2894 a.set_max_rounds(2);
2895
2896 let registry = one_tool();
2897 let mut session = Session::new(fmt!("s1"), fmt!("cut"), fmt!("model"));
2898 let mut events: Vec<AgentEvent> = Vec::new();
2899 let _ = a.run_turn(&mut session, fmt!("write me a long file"), &registry,
2900 &mut |ev| events.push(ev)).await;
2901
2902 // The user is told, in its own event rather than by the browser guessing from
2903 // arguments that will not parse.
2904 assert!(events.iter().any(|e| matches!(e, AgentEvent::Truncated)),
2905 "the reply was cut and nothing said so: {:?}",
2906 events.iter().map(|e| fmt!("{:?}", e)).collect::<Vec<_>>());
2907 // It is not reported as a failure: the request succeeded and a setting was hit.
2908 assert!(!events.iter().any(|e| matches!(e, AgentEvent::Error(_))),
2909 "a reply that hit the output cap was reported as an error");
2910 // And the MODEL is told what actually happened, in the tool reply it will read
2911 // on the next round -- not that its JSON was bad, which sends it round the same
2912 // loop writing the same thing.
2913 let reply = session.messages.iter().rev()
2914 .find_map(|m| match m {
2915 ChatMessage::Tool { content, .. } => Some(content.clone()),
2916 _ => None,
2917 })
2918 .expect("the cut call must still be answered, or the conversation is illegal");
2919 assert!(reply.as_text().contains("cut off at the output limit"), "{}", reply);
2920 assert!(reply.as_text().contains("smaller pieces"), "{}", reply);
2921 // Nothing was written: the arguments never reached the dispatcher.
2922 assert!(!reply.as_text().contains("Wrote"), "{}", reply);
2923 }
2924
2925 // ── A question put to the user ends the turn; a refused one does not ──
2926
2927 /// **The rule the tool loop ends a turn by, in all four of its cases.**
2928 ///
2929 /// Two of them are the defect and two are the controls. A question that was PUT must end the
2930 /// turn, or the model is charged a whole extra request whose only content is a paragraph
2931 /// restating the question under a card that already asks it. A question that was REFUSED
2932 /// must not, because every refusal `ask_step` composes is advice -- put the options back,
2933 /// name the recommendation, ask one thing rather than six -- and advice the model is denied
2934 /// the round to take is advice it never takes.
2935 ///
2936 /// **That second case is a scar and not a hypothesis.** The rule stood here before, read a
2937 /// tool's NAME alone, and named `say`; a worker refused the fold had its turn ended on the
2938 /// refusal telling it to write a report, so the report was never written and the whole errand
2939 /// came back as whatever prose accompanied the call. Work done, paid for and thrown away, on
2940 /// 2026-08-21.
2941 #[test]
2942 fn test_only_a_question_that_reached_the_screen_ends_the_turn_00() {
2943 use crate::tools::CallOutcome;
2944 assert!(ends_turn("ask", CallOutcome::Done),
2945 "a question that is on the user's screen no longer ends the turn, so the model is \
2946 charged a request to say it has asked");
2947 assert!(!ends_turn("ask", CallOutcome::Refused),
2948 "a refused question ended the turn, so the model never got the round in which to put \
2949 it properly -- which is what `say` did to a worker's report");
2950 assert!(!ends_turn("ask", CallOutcome::Failed),
2951 "a question the page could not draw ended the turn, so the user is left with nothing \
2952 on screen and the model with nothing to say");
2953 assert!(!ends_turn("file_write", CallOutcome::Done),
2954 "a tool that is not a question ended the turn");
2955 }
2956
2957 // ── A refused tool call does not end the turn ───────────────────────
2958
2959 /// **A worker whose tool call is refused gets the round it was told to use, and its report
2960 /// survives.**
2961 ///
2962 /// Two rounds from the stub. In the first the worker calls `file_show`, which a dispatched
2963 /// worker may not take -- nobody is reading its transcript. In the second it writes the report
2964 /// the refusal told it to write.
2965 ///
2966 /// **The subject used to be `say`, and it was the tool that ended a turn.** A `say` that
2967 /// ANSWERED ended it; the rule first read the tool's NAME alone, so a `say` that was REFUSED
2968 /// ended it too -- and a worker told to put the detail in its report was denied the round in
2969 /// which to write one, so the errand came back as whatever prose happened to accompany the
2970 /// call, which is usually nothing. Work done, paid for and thrown away. The tool is gone and
2971 /// no tool ends a turn now, which makes this the guard on that: a refusal must still leave the
2972 /// loop running, whichever tool refused.
2973 ///
2974 /// **Asserted on the report's CONTENT and not on the round count**, because a count is
2975 /// satisfied by any second round at all -- including one that says nothing.
2976 #[tokio::test]
2977 async fn test_a_worker_refused_a_tool_still_reports_00() {
2978 use crate::llm::tests::{start_stub, stub_client, Reply};
2979 const REPORT: &str = "THE-CRATE-FAILS-TO-BUILD-ON-LINE-42";
2980 let (port, _seen) = start_stub(vec![
2981 // Round one: the call, which is refused.
2982 Reply::Sse {
2983 chunks: vec![
2984 "data: {\"choices\":[{\"delta\":{\"tool_calls\":[{\"index\":0,\
2985 \"id\":\"call_1\",\"type\":\"function\",\"function\":{\"name\":\
2986 \"file_show\",\"arguments\":\"{\\\"path\\\":\\\"report.md\\\"}\"}}]}}]}\n\n"
2987 .to_string(),
2988 "data: {\"choices\":[{\"delta\":{},\"finish_reason\":\"stop\"}]}\n\n"
2989 .to_string(),
2990 "data: [DONE]\n\n".to_string(),
2991 ],
2992 reset_after: None,
2993 },
2994 // Round two: the report, in the place the refusal told it to put it.
2995 Reply::Sse {
2996 chunks: vec![
2997 fmt!("data: {{\"choices\":[{{\"delta\":{{\"content\":\"{}\"}}}}]}}\n\n",
2998 REPORT),
2999 "data: {\"choices\":[{\"delta\":{},\"finish_reason\":\"stop\"}]}\n\n"
3000 .to_string(),
3001 "data: [DONE]\n\n".to_string(),
3002 ],
3003 reset_after: None,
3004 },
3005 ]).await;
3006 let mut llm = stub_client(port);
3007 llm.retry.max_attempts = 1;
3008 let a = Agent::new(llm, "You are a worker.");
3009 a.set_max_rounds(3);
3010
3011 let mut registry = no_tools();
3012 registry.tools = vec![crate::tools::Tool::FileShow];
3013 registry.ctx.set_unsupervised();
3014
3015 let mut session = Session::new(fmt!("s1"), fmt!("worker"), fmt!("model"));
3016 let mut events: Vec<AgentEvent> = Vec::new();
3017 let _ = a.run_turn(&mut session, fmt!("check the build"), &registry,
3018 &mut |ev| events.push(ev)).await;
3019
3020 // What the agent that dispatched this worker actually receives.
3021 let report = session.messages.iter().rev()
3022 .find_map(|m| match m {
3023 ChatMessage::Assistant { content, tool_calls } if tool_calls.is_empty() =>
3024 Some(content.as_text().into_owned()),
3025 _ => None,
3026 })
3027 .unwrap_or_default();
3028 assert!(report.contains(REPORT),
3029 "the worker's report is not in the transcript: the turn ended on a refused call and \
3030 its findings were discarded. Last assistant turn: {:?}", report);
3031 // And the refusal is still on the record, so the transcript is well formed and the
3032 // model can see why it was asked to write prose.
3033 assert!(session.messages.iter().any(|m| matches!(m, ChatMessage::Tool { .. })),
3034 "the refused call was never answered, which is a malformed conversation");
3035 }
3036
3037 #[test]
3038 fn test_a_whole_call_is_dispatched_even_when_the_reply_was_cut_00() {
3039 // A turn can be cut in its trailing prose with every call already complete.
3040 // Refusing those would break a turn that was working.
3041 for whole in [
3042 r#"{"path":"a.rs"}"#,
3043 r#"{}"#,
3044 r#" {"argv":["cargo","test"]} "#,
3045 // A brace and a quote inside a string are contents, not structure.
3046 r#"{"content":"fn main() { \"hi\" }"}"#,
3047 // And an UNBALANCED brace inside one, which is what the first line of
3048 // almost every source file the model writes looks like.
3049 r#"{"path":"a.rs","content":"fn main() {"}"#,
3050 r#"{"content":"}"}"#,
3051 ] {
3052 assert!(!cut_short(true, whole),
3053 "a complete call was withheld because a later sentence was cut: {}", whole);
3054 assert!(json_object_is_whole(whole), "{}", whole);
3055 }
3056 }
3057
3058 #[test]
3059 fn test_a_broken_call_on_an_uncut_reply_is_still_the_models_own_mistake_00() {
3060 // Only the provider says why it stopped. Without that, an unbalanced object is
3061 // the model writing bad JSON, and telling it otherwise would send it splitting
3062 // work that was never too big.
3063 let broken = r#"{"path":"a.rs""#;
3064 assert!(!cut_short(false, broken));
3065 assert!(!json_object_is_whole(broken));
3066 }
3067
3068 #[test]
3069 fn test_an_unterminated_string_is_the_commonest_cut_of_all_00() {
3070 // A file's contents, stopped mid-word. Every brace is balanced; the quote is not.
3071 let s = r#"{"path":"a.rs","content":"fn main() {}"#;
3072 assert!(!json_object_is_whole(s), "an unterminated string read as whole");
3073 assert!(cut_short(true, s));
3074 }
3075
3076 /// An agent whose provider is Anthropic and whose model takes adaptive thinking, which is the
3077 /// one combination the client silently raises the output cap for.
3078 fn thinking_agent(max_tokens: u32) -> Agent {
3079 let tls = build_test_tls_config();
3080 let llm = LlmClient::new("api.anthropic.com", 443, "/v1/messages", "key",
3081 "claude-opus-5", max_tokens, tls);
3082 Agent::new(llm, "You are Daimond.")
3083 }
3084
3085 /// How a turn recovers from a refusal: fold to the budget, send it, and be refused again
3086 /// whenever the prompt plus the reply the client will ask for still exceeds the window.
3087 ///
3088 /// Returns how many refusals it took to fit, and the prompt budget that finally did -- because
3089 /// the count alone understates the harm. Each refusal costs a whole round trip AND teaches
3090 /// [`Limits::learn_from_refusal`] a smaller window, which is permanent for the session and
3091 /// never revised upward, so a needless refusal leaves the app folding a large model as though
3092 /// it were a small one for the rest of the conversation.
3093 ///
3094 /// The iteration is bounded: what is under test is a loop that need not terminate, and a test
3095 /// that reproduced it faithfully would not return.
3096 ///
3097 /// # Arguments
3098 /// * `a` - The agent, whose limits are ground down in place exactly as a real turn grinds them.
3099 /// * `window` - The provider's real window, which the app does not know and must discover.
3100 /// * `cap` - The output cap the budget is computed against; the bug is passing the wrong one.
3101 fn recovery(a: &Agent, window: u64, cap: u32) -> Option<(usize, u64)> {
3102 for round in 0..40 {
3103 let prompt = a.limits.borrow().budget(cap);
3104 // The whole request: what is sent, plus the room the client asks the provider to
3105 // leave for the answer. That sum is what the provider measures against its window.
3106 if prompt + (a.reply_cap() as u64) <= window {
3107 return Some((round, prompt));
3108 }
3109 // Refused. The app learns from the size of the PROMPT, which is all it sent.
3110 if !a.limits.borrow_mut().learn_from_refusal(prompt) {
3111 return None; // nothing left to learn: the same prompt goes out for ever
3112 }
3113 }
3114 None
3115 }
3116
3117 #[test]
3118 fn test_the_reply_cap_is_the_one_the_client_will_actually_send_00() {
3119 // A streamed Anthropic request to a thinking model is sent 32,000 whatever the client was
3120 // configured with, because thinking is billed as output and counts against the same cap.
3121 assert_eq!(32_000, thinking_agent(4_096).reply_cap());
3122 // A cap already above the floor is its own figure.
3123 assert_eq!(50_000, thinking_agent(50_000).reply_cap());
3124 // And nothing else is raised: not an Anthropic model that takes no adaptive thinking...
3125 let tls = build_test_tls_config();
3126 let old = Agent::new(LlmClient::new("api.anthropic.com", 443, "/v1/messages", "key",
3127 "claude-3-5-sonnet-20241022", 4_096, tls.clone()), "x");
3128 assert_eq!(4_096, old.reply_cap());
3129 // ...nor an OpenAI-dialect endpoint, whose max_tokens bounds the answer alone.
3130 let oa = Agent::new(LlmClient::new("api.test.com", 443, "/v1/chat/completions", "key",
3131 "claude-opus-5", 4_096, tls), "x");
3132 assert_eq!(4_096, oa.reply_cap());
3133 }
3134
3135 #[test]
3136 fn test_a_budget_blind_to_the_real_reply_refuses_its_way_down_the_window_00() {
3137 // `budget` subtracts the reply from the window to decide how big a prompt may be, so it
3138 // has to be given the figure the provider will be SENT. Handed `llm.max_tokens` it
3139 // reserved 5,120 for a reply that may run to 32,000. On the published 131,072 window the
3140 // fraction ceiling is the lower of the two and wins anyway, so nothing shows; on a window
3141 // learned from a refusal -- the mechanism that exists to recover from the first one -- the
3142 // prompt is legal by the app's arithmetic and refused by the provider, over and over.
3143 //
3144 // THE WINDOW IS 80,000 BECAUSE `FOLD_AT` MOVED. At 0.8 the two figures parted on a
3145 // 100,000 window; at 0.65 the fraction ceiling is 65,000 there, which is below the honest
3146 // reserve ceiling of 66,976 -- so both budgets come out at 65,000 and this test could see
3147 // nothing to compare. The reserve binds where the fraction does not, which above 64,000 is
3148 // any window under 94,354; 80,000 sits inside it with room to spare.
3149 let window = 80_000;
3150
3151 // Told the truth, the fold fits first time and the learned window is left alone.
3152 let fixed = thinking_agent(4_096);
3153 fixed.set_context_window(window);
3154 assert_eq!(32_000, fixed.reply_cap());
3155 assert_eq!(Some((0, 46_976)), recovery(&fixed, window, fixed.reply_cap()),
3156 "the budget must leave room for the reply the client asks for");
3157 assert_eq!(window, fixed.limits().window, "and must not have to learn anything");
3158
3159 // Blind, it sends the fold fraction of the window with 32,000 of reply behind it, is
3160 // refused, and pays for the mistake twice: a wasted round trip, and a window permanently
3161 // taught to be 39,000 -- `learn_from_refusal` never revises upward, so every later fold in
3162 // this session is made against a model half the size of the real one.
3163 let blind = thinking_agent(4_096);
3164 blind.set_context_window(window);
3165 assert_eq!(Some((1, 25_350)), recovery(&blind, window, blind.llm.max_tokens),
3166 "the old figure must cost a refusal it did not need to");
3167 assert_eq!(39_000, blind.limits().window, "and mis-teach the window for the rest of the run");
3168 // And the conversation pays: the fold that finally goes out is little more than half the
3169 // one the honest figure would have sent, on the same model, for no reason.
3170 assert!(25_350 < 46_976);
3171 }
3172
3173 #[test]
3174 fn test_a_small_learned_window_recovers_in_one_refusal_rather_than_several_00() {
3175 // The case the fix is really for. At 40,000 the reply is most of the window, so a budget
3176 // that ignores it is wrong by a factor of eight and the grinding-down is visible: the
3177 // blind figure is refused twice and lands on a 4,386-token conversation, the honest one is
3178 // refused once and keeps forty per cent more.
3179 //
3180 // It was THREE refusals and 2,366 tokens while `FOLD_AT` was 0.8. Lowering it to 0.65
3181 // makes the blind figure smaller from the start, so it crosses under the real window one
3182 // round sooner -- the blind path is still worse on both counts, which is the whole claim,
3183 // and the counts themselves are a property of the fraction rather than of the fix.
3184 let honest = thinking_agent(4_096);
3185 honest.set_context_window(40_000);
3186 let (h_rounds, h_prompt) = recovery(&honest, 40_000, honest.reply_cap())
3187 .expect("the honest figure converges");
3188
3189 let blind = thinking_agent(4_096);
3190 blind.set_context_window(40_000);
3191 let (b_rounds, b_prompt) = recovery(&blind, 40_000, blind.llm.max_tokens)
3192 .expect("the blind figure converges too, eventually");
3193
3194 assert_eq!((1, 6_092), (h_rounds, h_prompt));
3195 assert_eq!((2, 4_386), (b_rounds, b_prompt));
3196 assert!(h_rounds < b_rounds, "the fix must cost fewer round trips");
3197 assert!(h_prompt > b_prompt, "and leave more of the conversation standing");
3198 }
3199
3200 #[test]
3201 fn test_a_user_moves_where_their_conversation_folds_00() {
3202 // The setting exists so the person watching the meter can decide, so what is asserted is
3203 // that the BUDGET moves -- not that a field was written. A window of 100,000 makes the
3204 // arithmetic readable: two thirds of it against four fifths of it is a difference of
3205 // 15,000 tokens of conversation, which is several exchanges.
3206 let a = make_test_agent();
3207 a.set_context_window(100_000);
3208 let cap = 0u32; // no reply reserve in the way, so the fraction is the only ceiling
3209 let shipped = a.limits().budget(cap);
3210 assert_eq!(compact::FOLD_AT, a.limits().fold_at,
3211 "a fresh agent must fold at the shipped figure");
3212
3213 a.set_fold_at(0.8);
3214 let later = a.limits().budget(cap);
3215 assert!(later > shipped,
3216 "folding at 0.8 must leave more room than the shipped {}: {} against {}",
3217 compact::FOLD_AT, later, shipped);
3218
3219 // Zero is how a caller says "the user has not chosen". It must leave the figure alone
3220 // rather than fold the conversation to nothing, because that is what the browser passes
3221 // for every chat nobody has set a fraction on.
3222 a.set_fold_at(0.0);
3223 assert_eq!(later, a.limits().budget(cap),
3224 "zero must leave the agent's own figure standing");
3225
3226 // Held at the band, and asserted through the GETTER: a control that draws itself from
3227 // `fold_at` would otherwise show a figure the arithmetic never used.
3228 a.set_fold_at(9.0);
3229 assert_eq!(compact::FOLD_AT_MAX, a.limits().fold_at);
3230 a.set_fold_at(0.001);
3231 assert_eq!(compact::FOLD_AT_MIN, a.limits().fold_at);
3232 }
3233
3234 #[tokio::test]
3235 async fn test_the_fold_does_not_call_a_shortened_answer_a_tool_result_00() {
3236 // A PURE CHAT HAS NO TOOL RESULTS AT ALL, and `compact::elide_bulk` shrinks a long
3237 // assistant turn on the same rule it shrinks a tool reply on. The sentence used to say
3238 // "shortened N tool results" either way, so the one user whose own answers had just been
3239 // clipped to 400 characters was told about tool output they never had.
3240 //
3241 // Built too short to fold -- six messages is `MIN_KEEP_MESSAGES`, so `tail_start` answers
3242 // zero -- and too big to send, which is the only shape where eliding does all the work.
3243 let a = dead_thinking_agent();
3244 a.set_context_window(20_000);
3245 let mut session = Session::new(fmt!("s2"), fmt!("bulky"), fmt!("claude-opus-5"));
3246 for i in 0..2 {
3247 session.messages.push(ChatMessage::user(fmt!("ask {}", i)));
3248 session.messages.push(ChatMessage::Assistant {
3249 content: MessageContent::text("y".repeat(30_000)), tool_calls: Vec::new(),
3250 });
3251 }
3252 // COUNTED AS THE TURN WILL SEE IT. `run_turn` pushes the user's own message first, so a
3253 // fixture measured before that push is one message short -- which is how the first draft
3254 // of this test built five messages, watched `run_turn` make it six, and folded when it
3255 // meant to elide.
3256 let mut as_sent = session.messages.clone();
3257 as_sent.push(ChatMessage::user(fmt!("carry on")));
3258 assert_eq!(0, compact::tail_start(&as_sent,
3259 1_000, compact::MIN_KEEP_MESSAGES, 1_000, &a.llm.open_folds()),
3260 "the fixture must be too short to fold, or this tests the other half");
3261
3262 let registry = no_tools();
3263 let mut events: Vec<AgentEvent> = Vec::new();
3264 let _ = a.run_turn(&mut session, fmt!("carry on"), &registry,
3265 &mut |ev| events.push(ev)).await;
3266 let note = events.iter().find_map(|e| match e {
3267 AgentEvent::Compacted { note, .. } => Some(note.clone()),
3268 _ => None,
3269 }).expect("a conversation many times its window must have been compacted");
3270 assert!(note.starts_with("Shortened "),
3271 "the fixture must have elided rather than folded: {}", note);
3272 assert!(!note.contains("tool results:"),
3273 "the fold called an answer a tool result: {}", note);
3274 assert!(note.contains("tool results and answers"),
3275 "the fold must name both kinds, since it shortens both: {}", note);
3276 // AND WHERE, which is the second half of the 2026-08-28 ruling. A sentence saying the
3277 // answers were shortened, read beside those answers sitting in full on the screen above,
3278 // is a contradiction handed to the reader to resolve.
3279 assert!(note.contains("on the way to the model"),
3280 "the notice does not say the shortening is the request's and not the record's: {}",
3281 note);
3282
3283 // And the claim behind the wording: an ANSWER really was the thing that shrank. Asked of
3284 // the COUNT rather than of the stored text, because since 2026-08-28 the stored text is
3285 // exactly what it was -- see `test_an_answer_shortened_to_fit_stays_whole_in_the_record_00`,
3286 // which reads the request body a stub provider really received and is where the claim
3287 // that something shrank on the wire now lives.
3288 let shrank = note.split("hortened ").nth(1)
3289 .and_then(|s| s.split(' ').next())
3290 .and_then(|s| s.parse::<usize>().ok())
3291 .unwrap_or(0);
3292 assert!(shrank > 0,
3293 "nothing was shortened at all, so the wording is not what this is about: {}", note);
3294 // THE RECORD IS UNTOUCHED, and this is the assertion that was the other way round until
3295 // the ruling. A conversation with no fold and no room has had its request shortened; the
3296 // session it was built from still holds every character of both answers.
3297 for m in session.messages.iter() {
3298 assert!(!m.text().contains("folded away to fit the context window"),
3299 "the engine's own elision note was written into the stored conversation");
3300 }
3301 assert_eq!(2, session.messages.iter()
3302 .filter(|m| matches!(m, ChatMessage::Assistant { .. }) && m.text().len() >= 30_000)
3303 .count(),
3304 "an answer the model was sent a clipped copy of came back short in the record");
3305 }
3306
3307 /// A conversation that is too short to fold and too big to send, with a word past the cap.
3308 ///
3309 /// The marker sits a thousand characters into each answer, which is well past
3310 /// [`compact::TOOL_ELISION_CAP`], so it survives only where the whole answer survives. Both
3311 /// tests below turn on that one word being somewhere and not somewhere else.
3312 fn bulky_session(mark: &str) -> Session {
3313 let mut session = Session::new(fmt!("s"), fmt!("bulky"), fmt!("model"));
3314 for i in 0..2 {
3315 session.messages.push(ChatMessage::user(fmt!("ask {}", i)));
3316 session.messages.push(ChatMessage::Assistant {
3317 content: MessageContent::text(
3318 fmt!("{}{}{}", "y".repeat(1_000), mark, "y".repeat(30_000))),
3319 tool_calls: Vec::new(),
3320 });
3321 }
3322 session
3323 }
3324
3325 #[tokio::test]
3326 async fn test_an_answer_shortened_to_fit_stays_whole_in_the_record_00() {
3327 // THE OWNER'S RULING OF 2026-08-28: the model gets the shortened version, his transcript
3328 // keeps every word. `compact::elide_bulk` edited `session.messages` in place, and that
3329 // list is the one `DaimondApp::export_session` hands the browser to store, to back up and
3330 // to sync -- so a thousand-word answer became four hundred characters in the user's own
3331 // record, permanently, with nothing said and no way back.
3332 //
3333 // Both halves are asked of an artefact rather than of the app. What the model got is read
3334 // out of the request body a stub provider really received; what the user kept is read out
3335 // of the session the turn was run against.
3336 const MARK: &str = "MARKER-PAST-THE-CAP";
3337 let (port, seen) = crate::llm::tests::start_stub(vec![plain_answer()]).await;
3338 let a = Agent::new(crate::llm::tests::stub_client(port), "You are Daimond.");
3339 a.set_context_window(20_000);
3340 let mut session = bulky_session(MARK);
3341 // Six messages is `MIN_KEEP_MESSAGES`, and `run_turn` adds the user's own before any of
3342 // this runs -- so a fixture measured before that push is one short, which is how an
3343 // earlier test in this file folded when it meant to elide.
3344 let mut as_sent = session.messages.clone();
3345 as_sent.push(ChatMessage::user(fmt!("carry on")));
3346 assert_eq!(0, compact::tail_start(&as_sent,
3347 1_000, compact::MIN_KEEP_MESSAGES, 1_000, &a.llm.open_folds()),
3348 "the fixture must be too short to fold, or this tests the other half");
3349
3350 let registry = no_tools();
3351 let mut events: Vec<AgentEvent> = Vec::new();
3352 let _ = a.run_turn(&mut session, fmt!("carry on"), &registry,
3353 &mut |ev| events.push(ev)).await;
3354
3355 let bodies = match seen.lock() {
3356 Ok(g) => g.bodies.clone(),
3357 Err(e) => panic!("the stub's log: {}", e),
3358 };
3359 assert!(!bodies.is_empty(), "the stub was never called, so nothing was sent to read");
3360 let wire = bodies.join("\n");
3361 // The fixture reached the branch: something really was shortened on the way out.
3362 assert!(wire.contains("folded away to fit the context window"),
3363 "nothing was shortened for the model, so this proves nothing about what was kept");
3364 // COUNTED, not looked for. `elide_bulk` stops the moment the conversation fits, so with
3365 // two long answers it may clip one and leave the other -- and an assertion that the
3366 // marker is absent from the wire would then be red for a reason that is the elision
3367 // working. What must be true is that the model saw FEWER whole answers than the record
3368 // holds, and that is what is asked.
3369 let on_wire = wire.matches(MARK).count();
3370 let in_record = session.messages.iter().filter(|m| m.text().contains(MARK)).count();
3371 assert_eq!(2, in_record,
3372 "the stored conversation lost the words past the cap -- the record is still the wire");
3373 assert!(on_wire < in_record,
3374 "the model was sent {} whole answers of the {} the record holds, so nothing was \
3375 shortened on the wire", on_wire, in_record);
3376
3377 // AND THE RECORD IS WHOLE, which is the ruling. Red before it: the two lists were one.
3378 for m in session.messages.iter() {
3379 assert!(!m.text().contains("folded away to fit the context window"),
3380 "the engine's own elision note was written into the stored conversation");
3381 }
3382 }
3383
3384 #[tokio::test]
3385 async fn test_a_later_fold_summarises_the_words_and_not_the_stubs_00() {
3386 // THE SECOND LOSS, which stood behind the first and was quieter. `fold_if_needed` renders
3387 // `session.messages` for the summarising model (`compact::render_for_fold`) and then, in
3388 // the same function, used to clip that same list. So the FIRST fold read the words and
3389 // every fold after it read whatever the elision had left -- four hundred characters and a
3390 // note, per message. A conversation folded twice was summarised from stubs, and the
3391 // summary is the only thing that survives a fold.
3392 //
3393 // Asked of `render_for_fold` itself, over the session a real turn left behind, because
3394 // that is the call the fold makes and the input it makes it on.
3395 const MARK: &str = "MARKER-PAST-THE-CAP";
3396 let (port, _seen) = crate::llm::tests::start_stub(vec![plain_answer()]).await;
3397 let a = Agent::new(crate::llm::tests::stub_client(port), "You are Daimond.");
3398 a.set_context_window(20_000);
3399 let mut session = bulky_session(MARK);
3400 let registry = no_tools();
3401 let mut events: Vec<AgentEvent> = Vec::new();
3402 let _ = a.run_turn(&mut session, fmt!("carry on"), &registry,
3403 &mut |ev| events.push(ev)).await;
3404 // The turn really did shorten something, or there is no "afterwards" to test.
3405 assert!(events.iter().any(|e| matches!(e, AgentEvent::Compacted { .. })),
3406 "the conversation was sent unshortened, so no later fold could read a stub");
3407
3408 let rendered = compact::render_for_fold(&session.messages, compact::FOLD_INPUT_CAP);
3409 assert!(rendered.contains(MARK),
3410 "the next fold would summarise clipped stubs: {} bytes rendered, no marker in them",
3411 rendered.len());
3412 assert!(!rendered.contains("folded away to fit the context window"),
3413 "the next fold would be handed the engine's own elision notes as if they were the \
3414 conversation");
3415 }
3416
3417 /// A dead-endpoint agent whose dialect and model are the ones the client raises the output cap
3418 /// for, so a fold's arithmetic is exercised under the figure that actually goes out.
3419 fn dead_thinking_agent() -> Agent {
3420 let tls = build_test_tls_config();
3421 let mut llm = LlmClient::new("127.0.0.1", 1, "/v1/messages", "key",
3422 "claude-opus-5", 4_096, tls);
3423 llm.retry.max_attempts = 1;
3424 Agent::new(llm, "You are Daimond.")
3425 }
3426
3427 #[tokio::test]
3428 async fn test_whether_to_fold_is_decided_by_the_reply_that_will_be_sent_00() {
3429 // The two arithmetics above are only worth anything if the fold ASKS them, so this pins
3430 // the call site rather than the function. The window is 80,000; the honest budget is
3431 // 46,976 and the blind one 52,000, so a conversation sitting between the two folds under
3432 // the figure that will be sent and does not fold under the figure that was configured --
3433 // the window was 100,000 while `FOLD_AT` was 0.8, and at 0.65 the fraction wins there and
3434 // closes the gap, so the fixture moved down with the constant --
3435 // and not folding is a prompt the provider refuses.
3436 let a = dead_thinking_agent();
3437 a.set_context_window(80_000);
3438 let mut session = Session::new(fmt!("s1"), fmt!("long"), fmt!("claude-opus-5"));
3439 for i in 0..46 {
3440 session.messages.push(ChatMessage::user(fmt!("step {}", i)));
3441 session.messages.push(ChatMessage::Assistant {
3442 content: MessageContent::text("x".repeat(4_000)), tool_calls: Vec::new(),
3443 });
3444 }
3445 // The conversation is deliberately built into the gap, and the gap is asserted rather
3446 // than assumed: a change to the gauge or to `FOLD_AT` that closed it would otherwise make
3447 // this test pass while testing nothing.
3448 let tokens = a.gauge.tokens(
3449 compact::conversation_bytes(&session.messages, &a.llm.open_folds()));
3450 let honest = a.limits().budget(a.reply_cap());
3451 let blind = a.limits().budget(a.llm.max_tokens);
3452 assert!(honest < tokens && tokens <= blind,
3453 "{} tokens is not between the honest budget {} and the blind one {}",
3454 tokens, honest, blind);
3455
3456 let registry = no_tools();
3457 let mut events: Vec<AgentEvent> = Vec::new();
3458 let _ = a.run_turn(&mut session, fmt!("carry on"), &registry,
3459 &mut |ev| events.push(ev)).await;
3460 assert_eq!(1, events.iter().filter(|e| matches!(e, AgentEvent::Compacted { .. })).count(),
3461 "a conversation over the real budget was sent unfolded");
3462 }
3463
3464 #[test]
3465 fn test_the_reserve_is_capped_at_half_the_window_and_that_is_not_this_files_call_00() {
3466 // Recorded rather than asserted away, because the fix here does not finish the job.
3467 // `Limits::budget` never gives the reply more than half the window, so wherever the window
3468 // is below twice the output cap the reserve is short however truthful the figure handed to
3469 // it: on 40,000 the reply may be 32,000 and at most 20,000 is set aside. What this file
3470 // owes is the true figure, and it now passes it; the clamp belongs to `compact.rs`, and a
3471 // build whose learned window falls under twice its cap wants a lower CAP, not a bigger
3472 // prompt.
3473 let a = thinking_agent(4_096);
3474 a.set_context_window(40_000);
3475 let honest = a.limits().budget(a.reply_cap());
3476 let blind = a.limits().budget(a.llm.max_tokens);
3477 // As the FRACTION, not as the figure it happened to come to. That is what the sentence
3478 // beside it claims, and writing it as 32,000 made a test about the reply reserve go red
3479 // when the fold fraction moved.
3480 assert_eq!((40_000.0 * compact::FOLD_AT) as u64, blind,
3481 "blind, the fold fraction was left untouched");
3482 assert_eq!(18_976, honest, "honest, the clamped reserve of 20,000 is taken out");
3483 assert!(honest + 20_000 <= 40_000, "the clamped reserve must at least be honoured");
3484 assert!(honest + (a.reply_cap() as u64) > 40_000,
3485 "and the clamp still leaves a gap, which is compact.rs's to close");
3486 }
3487
3488 #[test]
3489 fn test_the_budget_leaves_room_for_the_reply_00() {
3490 // `max_tokens` on the client is what the model may generate, and it is counted
3491 // against the same window. A budget blind to it is legal arithmetic and an illegal
3492 // request.
3493 let a = make_test_agent();
3494 a.set_context_window(8_192);
3495 let b = a.limits().budget(a.llm.max_tokens);
3496 assert!(b + (a.llm.max_tokens as u64) <= 8_192, "budget {} of 8192", b);
3497 }
3498
3499 #[test]
3500fn test_agent_message_building() {
3501 let mut session = Session::new("s1".to_string(), "Test".to_string(), "model".to_string());
3502 session.messages.push(ChatMessage::user("Hello".to_string()));
3503 assert_eq!(session.messages.len(), 1);
3504 assert_eq!(session.messages[0].role(), "user");
3505 assert_eq!(session.messages[0].text(), "Hello");
3506 }
3507}