1mod approve;
11mod config;
12mod run;
13mod session;
14mod sink;
15mod terminal;
16
17pub use approve::{AllowAll, ApprovalRequest, Approver, Decision};
18pub use config::EngineConfig;
19pub use session::Session;
20pub use sink::{EventSink, FnSink, NullSink};
21pub use tokio_util::sync::CancellationToken;
24
25#[cfg(test)]
26mod tests {
27 #![allow(clippy::unnecessary_literal_bound)]
30
31 use super::*;
32 use async_trait::async_trait;
33 use locode_protocol::{
34 ContentBlock, Conversation, Event, Message, ReasoningFormat, Role, Status, Usage,
35 reconstruct_conversation,
36 };
37 use locode_provider::{
38 Completion, ConversationRequest, MockProvider, Provider, ProviderError, StopReason,
39 };
40 use locode_tools::{Registry, Tool, ToolCtx, ToolError, ToolKind, ToolOutput};
41 use serde::Serialize;
42 use serde_json::{Value, json};
43 use std::sync::{Arc, Mutex};
44 use std::time::Duration;
45
46 #[derive(Serialize)]
49 struct EchoOut {
50 echoed: String,
51 }
52 impl ToolOutput for EchoOut {
53 fn to_prompt_text(&self) -> String {
54 self.echoed.clone()
55 }
56 }
57
58 struct Echo;
59 #[async_trait]
60 impl Tool for Echo {
61 type Args = Value;
62 type Output = EchoOut;
63 fn kind(&self) -> ToolKind {
64 ToolKind::Shell
65 }
66 fn description(&self) -> &str {
67 "echo"
68 }
69 async fn run(&self, _ctx: &ToolCtx, args: Value) -> Result<EchoOut, ToolError> {
70 Ok(EchoOut {
71 echoed: args.to_string(),
72 })
73 }
74 }
75
76 struct Boom;
77 #[async_trait]
78 impl Tool for Boom {
79 type Args = Value;
80 type Output = EchoOut;
81 fn kind(&self) -> ToolKind {
82 ToolKind::Shell
83 }
84 fn description(&self) -> &str {
85 "boom"
86 }
87 async fn run(&self, _ctx: &ToolCtx, _args: Value) -> Result<EchoOut, ToolError> {
88 Err(ToolError::Fatal("boom aborted the turn".into()))
89 }
90 }
91
92 fn text_turn(text: &str) -> Completion {
95 Completion {
96 content: vec![ContentBlock::Text { text: text.into() }],
97 usage: Usage::default(),
98 stop: StopReason::EndTurn,
99 }
100 }
101
102 fn tool_turn(id: &str, name: &str) -> Completion {
103 Completion {
104 content: vec![ContentBlock::ToolUse {
105 id: id.into(),
106 name: name.into(),
107 input: json!({}),
108 }],
109 usage: Usage::default(),
110 stop: StopReason::ToolUse,
111 }
112 }
113
114 fn config() -> EngineConfig {
115 EngineConfig {
116 session_id: "sess-1".into(),
117 harness: "grok".into(),
118 api_schema: "mock".into(),
119 model: "mock-1".into(),
120 max_turns: None,
121 resample_retries: 2,
122 resample_backoff: Duration::ZERO, instructions: locode_instructions::InstructionsConfig {
128 enabled: false,
129 ..Default::default()
130 },
131 ..EngineConfig::default()
132 }
133 }
134
135 fn session_with(
137 script: Vec<Result<Completion, ProviderError>>,
138 registry: Registry,
139 cfg: EngineConfig,
140 ) -> (Session, Arc<Mutex<Vec<Event>>>) {
141 let events = Arc::new(Mutex::new(Vec::new()));
142 let sink_events = Arc::clone(&events);
143 let sink = Box::new(FnSink(move |event| {
144 sink_events.lock().unwrap().push(event);
145 }));
146 let provider = Arc::new(MockProvider::with_results(script));
147 let session = Session::new(provider, registry, vec![], cfg, sink);
148 (session, events)
149 }
150
151 fn echo_registry() -> Registry {
152 let mut reg = Registry::new();
153 reg.register("echo", Echo);
154 reg
155 }
156
157 fn dump(events: &Arc<Mutex<Vec<Event>>>) -> Vec<Event> {
158 events.lock().unwrap().clone()
159 }
160
161 #[tokio::test]
164 async fn completed_with_no_tools() {
165 let (mut s, events) =
166 session_with(vec![Ok(text_turn("all done"))], Registry::new(), config());
167 let report = s.run_text("hi").await;
168 assert_eq!(report.status, Status::Completed);
169 assert_eq!(report.final_message.as_deref(), Some("all done"));
170 assert_eq!(report.turns, 1);
171 assert!(report.tool_calls.is_empty());
172 assert_eq!(report.api_schema, "mock");
173 let evs = dump(&events);
175 assert!(matches!(evs.first(), Some(Event::Init { .. })));
176 assert!(matches!(evs.last(), Some(Event::Result { .. })));
177 }
178
179 #[tokio::test]
182 async fn streaming_run_emits_text_deltas_and_the_whole_message() {
183 let mut cfg = config();
184 cfg.streaming = true;
185 let (mut s, events) = session_with(
186 vec![Ok(text_turn("hello streamed world"))],
187 Registry::new(),
188 cfg,
189 );
190 let report = s.run_text("hi").await;
191 assert_eq!(report.status, Status::Completed);
192 assert_eq!(
193 report.final_message.as_deref(),
194 Some("hello streamed world")
195 );
196
197 let evs = dump(&events);
198 let delta_text: String = evs
200 .iter()
201 .filter_map(|e| match e {
202 Event::MessageDelta { text } => Some(text.as_str()),
203 _ => None,
204 })
205 .collect();
206 assert_eq!(delta_text, "hello streamed world", "{evs:?}");
207 let n_deltas = evs
209 .iter()
210 .filter(|e| matches!(e, Event::MessageDelta { .. }))
211 .count();
212 assert!(n_deltas > 1, "expected multiple deltas, got {n_deltas}");
213 assert!(
215 evs.iter().any(|e| matches!(
216 e,
217 Event::Message { message } if message.role == Role::Assistant
218 )),
219 "whole assistant Message still emitted: {evs:?}"
220 );
221 let first_delta = evs
223 .iter()
224 .position(|e| matches!(e, Event::MessageDelta { .. }))
225 .expect("a delta");
226 let asst_msg = evs
227 .iter()
228 .position(
229 |e| matches!(e, Event::Message { message } if message.role == Role::Assistant),
230 )
231 .expect("assistant message");
232 assert!(
233 first_delta < asst_msg,
234 "deltas come before the whole message"
235 );
236 }
237
238 #[tokio::test]
239 async fn non_streaming_run_emits_no_deltas() {
240 let (mut s, events) =
241 session_with(vec![Ok(text_turn("no stream"))], Registry::new(), config());
242 let _ = s.run_text("hi").await;
243 let evs = dump(&events);
244 assert!(
245 !evs.iter().any(|e| matches!(e, Event::MessageDelta { .. })),
246 "default (non-streaming) run must not emit deltas: {evs:?}"
247 );
248 }
249
250 #[tokio::test]
251 async fn streaming_and_non_streaming_reports_match() {
252 let (mut a, _ea) = session_with(
253 vec![Ok(text_turn("same result"))],
254 Registry::new(),
255 config(),
256 );
257 let mut cfg = config();
258 cfg.streaming = true;
259 let (mut b, _eb) = session_with(vec![Ok(text_turn("same result"))], Registry::new(), cfg);
260 let ra = a.run_text("go").await;
261 let rb = b.run_text("go").await;
262 assert_eq!(ra.status, rb.status);
264 assert_eq!(ra.final_message, rb.final_message);
265 assert_eq!(ra.turns, rb.turns);
266 assert_eq!(ra.tool_calls.len(), rb.tool_calls.len());
267 }
268
269 #[tokio::test]
270 async fn tool_call_then_complete() {
271 let (mut s, _e) = session_with(
272 vec![Ok(tool_turn("c1", "echo")), Ok(text_turn("done"))],
273 echo_registry(),
274 config(),
275 );
276 let report = s.run_text("go").await;
277 assert_eq!(report.status, Status::Completed);
278 assert_eq!(report.turns, 2);
279 assert_eq!(report.tool_calls.len(), 1);
280 assert!(report.tool_calls[0].ok);
281 assert_eq!(report.tool_calls[0].name, "echo");
282 }
283
284 #[tokio::test]
285 async fn hits_max_turns_after_dispatch() {
286 let mut cfg = config();
288 cfg.max_turns = Some(2);
289 let (mut s, _e) = session_with(
290 vec![
291 Ok(tool_turn("c1", "echo")),
292 Ok(tool_turn("c2", "echo")),
293 Ok(tool_turn("c3", "echo")),
294 ],
295 echo_registry(),
296 cfg,
297 );
298 let report = s.run_text("go").await;
299 assert_eq!(report.status, Status::MaxTurns);
300 assert_eq!(report.turns, 2);
301 assert_eq!(report.tool_calls.len(), 2);
302 }
303
304 #[tokio::test]
305 async fn model_error_after_bounded_retry() {
306 let script = vec![
308 Err(ProviderError::Transport("reset".into())),
309 Err(ProviderError::Transport("reset".into())),
310 Err(ProviderError::Transport("reset".into())),
311 ];
312 let (mut s, events) = session_with(script, Registry::new(), config());
313 let report = s.run_text("go").await;
314 assert_eq!(report.status, Status::ModelError);
315 assert!(report.error.is_some());
316 assert_eq!(report.turns, 0);
317 let retries = dump(&events)
319 .iter()
320 .filter(|e| matches!(e, Event::Error { .. }))
321 .count();
322 assert_eq!(retries, 2);
323 }
324
325 #[tokio::test]
326 async fn model_error_non_retryable_is_immediate() {
327 let (mut s, events) = session_with(
328 vec![Err(ProviderError::ContextOverflow)],
329 Registry::new(),
330 config(),
331 );
332 let report = s.run_text("go").await;
333 assert_eq!(report.status, Status::ModelError);
334 let retries = dump(&events)
335 .iter()
336 .filter(|e| matches!(e, Event::Error { .. }))
337 .count();
338 assert_eq!(retries, 0, "a non-retryable error must not resample");
339 }
340
341 #[tokio::test]
342 async fn fatal_tool_error_ends_the_run() {
343 let mut reg = Registry::new();
344 reg.register("boom", Boom);
345 let (mut s, _e) = session_with(vec![Ok(tool_turn("c1", "boom"))], reg, config());
346 let report = s.run_text("go").await;
347 assert_eq!(report.status, Status::Error);
348 assert!(report.error.is_some());
349 assert_eq!(report.tool_calls.len(), 1);
351 assert!(!report.tool_calls[0].ok);
352 }
353
354 #[tokio::test]
358 async fn empty_completion_resamples_then_succeeds() {
359 let empty = Completion {
360 content: vec![ContentBlock::Reasoning {
361 format: ReasoningFormat::Anthropic,
362 text: "thinking only".into(),
363 signature: Some("sig".into()),
364 payload: None,
365 }],
366 usage: Usage::default(),
367 stop: StopReason::MaxTokens,
368 };
369 let (mut session, _events) = session_with(
370 vec![Ok(empty), Ok(text_turn("recovered"))],
371 echo_registry(),
372 config(),
373 );
374 let report = session.run_text("go").await;
375 assert_eq!(report.status, Status::Completed);
376 assert_eq!(report.final_message.as_deref(), Some("recovered"));
377 assert_eq!(report.stop_reason.as_deref(), Some("end_turn"));
378 }
379
380 #[tokio::test]
381 async fn persistent_empty_completions_are_model_error() {
382 let empty = || Completion {
383 content: vec![],
384 usage: Usage::default(),
385 stop: StopReason::MaxTokens,
386 };
387 let (mut session, _events) = session_with(
389 vec![Ok(empty()), Ok(empty()), Ok(empty())],
390 echo_registry(),
391 config(),
392 );
393 let report = session.run_text("go").await;
394 assert_eq!(report.status, Status::ModelError);
395 assert!(
396 report
397 .error
398 .as_deref()
399 .unwrap_or("")
400 .contains("empty completion"),
401 "error names the cause: {:?}",
402 report.error
403 );
404 assert_eq!(report.stop_reason, None, "no completion was accepted");
405 }
406
407 #[tokio::test]
410 async fn mid_batch_abort_synthesizes_results() {
411 let mut reg = Registry::new();
414 reg.register("boom", Boom);
415 reg.register("echo", Echo);
416 let completion = Completion {
417 content: vec![
418 ContentBlock::ToolUse {
419 id: "c_boom".into(),
420 name: "boom".into(),
421 input: json!({}),
422 },
423 ContentBlock::ToolUse {
424 id: "c_echo".into(),
425 name: "echo".into(),
426 input: json!({}),
427 },
428 ],
429 usage: Usage::default(),
430 stop: StopReason::ToolUse,
431 };
432 let (mut s, events) = session_with(vec![Ok(completion)], reg, config());
433 let report = s.run_text("go").await;
434 assert_eq!(report.status, Status::Error);
435
436 let evs = dump(&events);
438 let answered: Vec<String> = evs
439 .iter()
440 .filter_map(|e| match e {
441 Event::Message { message } if message.role == Role::User => Some(&message.content),
442 _ => None,
443 })
444 .flatten()
445 .filter_map(|b| match b {
446 ContentBlock::ToolResult { tool_use_id, .. } => Some(tool_use_id.clone()),
447 _ => None,
448 })
449 .collect();
450 assert!(answered.iter().any(|id| id == "c_boom"));
451 assert!(
452 answered.iter().any(|id| id == "c_echo"),
453 "the un-run echo must be paired"
454 );
455 assert_eq!(report.tool_calls.len(), 1);
457 }
458
459 #[tokio::test]
462 async fn thinking_block_is_appended_verbatim() {
463 let completion = Completion {
464 content: vec![
465 ContentBlock::Reasoning {
466 format: ReasoningFormat::Anthropic,
467 text: "reasoning".into(),
468 signature: Some("sig-xyz".into()),
469 payload: None,
470 },
471 ContentBlock::Text {
472 text: "answer".into(),
473 },
474 ],
475 usage: Usage::default(),
476 stop: StopReason::EndTurn,
477 };
478 let (mut s, events) = session_with(vec![Ok(completion)], Registry::new(), config());
479 let report = s.run_text("think").await;
480 assert_eq!(report.status, Status::Completed);
481 assert_eq!(report.final_message.as_deref(), Some("answer"));
482 let has_thinking = dump(&events).iter().any(|e| match e {
484 Event::Message { message } if message.role == Role::Assistant => {
485 message.content.iter().any(|b| {
486 matches!(
487 b,
488 ContentBlock::Reasoning { signature: Some(sig), .. } if sig == "sig-xyz"
489 )
490 })
491 }
492 _ => false,
493 });
494 assert!(
495 has_thinking,
496 "thinking + signature must survive into history"
497 );
498 }
499
500 #[tokio::test]
501 async fn events_reconstruct_the_history() {
502 let (mut s, events) = session_with(
503 vec![Ok(tool_turn("c1", "echo")), Ok(text_turn("done"))],
504 echo_registry(),
505 config(),
506 );
507 let _ = s.run_text("go").await;
508 let rebuilt: Conversation = reconstruct_conversation(&dump(&events));
509 let roles: Vec<Role> = rebuilt.messages.iter().map(|m| m.role).collect();
511 assert_eq!(
512 roles,
513 vec![Role::User, Role::Assistant, Role::User, Role::Assistant]
514 );
515 }
516
517 use std::sync::atomic::{AtomicUsize, Ordering};
520
521 struct Counting(Arc<AtomicUsize>);
523 #[async_trait]
524 impl Tool for Counting {
525 type Args = Value;
526 type Output = EchoOut;
527 fn kind(&self) -> ToolKind {
528 ToolKind::Shell
529 }
530 fn description(&self) -> &str {
531 "counting"
532 }
533 async fn run(&self, _ctx: &ToolCtx, _args: Value) -> Result<EchoOut, ToolError> {
534 self.0.fetch_add(1, Ordering::SeqCst);
535 Ok(EchoOut {
536 echoed: "ran".into(),
537 })
538 }
539 }
540
541 type SeenKinds = Arc<Mutex<Vec<(String, Option<ToolKind>)>>>;
542
543 struct DenyNamed {
546 deny: Vec<&'static str>,
547 seen_kinds: SeenKinds,
548 }
549 #[async_trait]
550 impl Approver for DenyNamed {
551 async fn decide(&self, request: &ApprovalRequest<'_>) -> Decision {
552 self.seen_kinds
553 .lock()
554 .unwrap()
555 .push((request.tool_name.to_owned(), request.kind));
556 if self.deny.contains(&request.tool_name) {
557 Decision::Deny {
558 reason: format!("{} is not allowed here", request.tool_name),
559 }
560 } else {
561 Decision::Allow
562 }
563 }
564 }
565
566 fn approvals(events: &Arc<Mutex<Vec<Event>>>) -> Vec<(String, String, String)> {
567 dump(events)
568 .iter()
569 .filter_map(|e| match e {
570 Event::Approval {
571 tool_use_id,
572 tool_name,
573 decision,
574 ..
575 } => Some((tool_use_id.clone(), tool_name.clone(), decision.clone())),
576 _ => None,
577 })
578 .collect()
579 }
580
581 #[tokio::test]
582 async fn deny_is_a_soft_paired_error_and_the_run_continues() {
583 let ran = Arc::new(AtomicUsize::new(0));
584 let mut reg = Registry::new();
585 reg.register("counting", Counting(Arc::clone(&ran)));
586 let (s, events) = session_with(
587 vec![Ok(tool_turn("c1", "counting")), Ok(text_turn("done"))],
588 reg,
589 config(),
590 );
591 let seen = Arc::new(Mutex::new(Vec::new()));
592 let mut s = s.with_approver(Arc::new(DenyNamed {
593 deny: vec!["counting"],
594 seen_kinds: Arc::clone(&seen),
595 }));
596 let report = s.run_text("go").await;
597
598 assert_eq!(report.status, Status::Completed);
600 assert_eq!(ran.load(Ordering::SeqCst), 0, "denied tool must not run");
601
602 assert_eq!(report.tool_calls.len(), 1);
604 let record = &report.tool_calls[0];
605 assert!(!record.ok);
606 assert_eq!(
607 record.denial_reason.as_deref(),
608 Some("counting is not allowed here")
609 );
610 assert_eq!(record.kind, "shell", "kind still recorded on denial");
611
612 let denied_result = dump(&events).iter().any(|e| match e {
614 Event::Message { message } => message.content.iter().any(|b| {
615 matches!(
616 b,
617 ContentBlock::ToolResult { tool_use_id, is_error: true, content, .. }
618 if tool_use_id == "c1"
619 && content.iter().any(|c| matches!(
620 c,
621 locode_protocol::ResultChunk::Text { text }
622 if text == "tool call denied: counting is not allowed here"
623 ))
624 )
625 }),
626 _ => false,
627 });
628 assert!(denied_result, "the model sees the denial reason, paired");
629
630 assert_eq!(
632 approvals(&events),
633 vec![("c1".into(), "counting".into(), "deny".into())]
634 );
635 }
636
637 #[tokio::test]
638 async fn deny_then_allow_within_one_batch_keeps_order_and_pairing() {
639 let ran = Arc::new(AtomicUsize::new(0));
640 let mut reg = Registry::new();
641 reg.register("blocked", Counting(Arc::clone(&ran)));
642 reg.register("echo", Echo);
643 let batch = Completion {
644 content: vec![
645 ContentBlock::ToolUse {
646 id: "c1".into(),
647 name: "blocked".into(),
648 input: json!({}),
649 },
650 ContentBlock::ToolUse {
651 id: "c2".into(),
652 name: "echo".into(),
653 input: json!({}),
654 },
655 ],
656 usage: Usage::default(),
657 stop: StopReason::ToolUse,
658 };
659 let (s, events) = session_with(vec![Ok(batch), Ok(text_turn("done"))], reg, config());
660 let mut s = s.with_approver(Arc::new(DenyNamed {
661 deny: vec!["blocked"],
662 seen_kinds: Arc::new(Mutex::new(Vec::new())),
663 }));
664 let report = s.run_text("go").await;
665 assert_eq!(report.status, Status::Completed);
666 assert_eq!(ran.load(Ordering::SeqCst), 0);
667
668 let pairs: Vec<(String, bool)> = dump(&events)
670 .iter()
671 .filter_map(|e| match e {
672 Event::Message { message } if message.role == Role::User => Some(&message.content),
673 _ => None,
674 })
675 .flatten()
676 .filter_map(|b| match b {
677 ContentBlock::ToolResult {
678 tool_use_id,
679 is_error,
680 ..
681 } => Some((tool_use_id.clone(), *is_error)),
682 _ => None,
683 })
684 .collect();
685 assert_eq!(pairs, vec![("c1".into(), true), ("c2".into(), false)]);
686
687 assert_eq!(report.tool_calls.len(), 2);
689 assert!(report.tool_calls[0].denial_reason.is_some());
690 assert_eq!(report.tool_calls[0].kind, "shell");
691 assert!(report.tool_calls[1].ok);
692 assert_eq!(report.tool_calls[1].denial_reason, None);
693
694 assert_eq!(
696 approvals(&events),
697 vec![
698 ("c1".into(), "blocked".into(), "deny".into()),
699 ("c2".into(), "echo".into(), "allow".into()),
700 ]
701 );
702 }
703
704 #[tokio::test]
705 async fn approval_request_carries_the_registry_kind() {
706 let seen = Arc::new(Mutex::new(Vec::new()));
707 let (s, _e) = session_with(
708 vec![Ok(tool_turn("c1", "echo")), Ok(text_turn("done"))],
709 echo_registry(),
710 config(),
711 );
712 let mut s = s.with_approver(Arc::new(DenyNamed {
713 deny: vec![],
714 seen_kinds: Arc::clone(&seen),
715 }));
716 let _ = s.run_text("go").await;
717 let seen = seen.lock().unwrap();
718 assert_eq!(seen.len(), 1);
719 assert_eq!(seen[0].0, "echo");
720 assert_eq!(
721 seen[0].1,
722 Some(ToolKind::Shell),
723 "kind resolves from the registry pre-dispatch"
724 );
725 }
726
727 #[tokio::test]
731 async fn async_approver_suspends_the_call_until_resolved() {
732 struct OneshotApprover(Mutex<Option<tokio::sync::oneshot::Receiver<Decision>>>);
733 #[async_trait]
734 impl Approver for OneshotApprover {
735 async fn decide(&self, _request: &ApprovalRequest<'_>) -> Decision {
736 let rx = self.0.lock().unwrap().take().expect("one decision");
737 rx.await.expect("decider dropped")
738 }
739 }
740
741 let (tx, rx) = tokio::sync::oneshot::channel();
742 let (s, _e) = session_with(
743 vec![Ok(tool_turn("c1", "echo")), Ok(text_turn("done"))],
744 echo_registry(),
745 config(),
746 );
747 let mut s = s.with_approver(Arc::new(OneshotApprover(Mutex::new(Some(rx)))));
748
749 let ui = tokio::spawn(async move {
751 tokio::task::yield_now().await;
752 let _ = tx.send(Decision::Allow);
753 });
754 let report = s.run_text("go").await;
755 ui.await.expect("ui task");
756 assert_eq!(report.status, Status::Completed);
757 assert_eq!(report.tool_calls.len(), 1);
758 assert!(report.tool_calls[0].ok);
759 }
760
761 #[tokio::test]
762 async fn allowed_calls_emit_approval_events_by_default() {
763 let (mut s, events) = session_with(
765 vec![Ok(tool_turn("c1", "echo")), Ok(text_turn("done"))],
766 echo_registry(),
767 config(),
768 );
769 let report = s.run_text("go").await;
770 assert_eq!(report.status, Status::Completed);
771 assert_eq!(
772 approvals(&events),
773 vec![("c1".into(), "echo".into(), "allow".into())]
774 );
775 assert_eq!(report.tool_calls[0].denial_reason, None);
777 }
778
779 struct HangingProvider;
784 #[async_trait]
785 impl Provider for HangingProvider {
786 #[allow(clippy::unnecessary_literal_bound)]
787 fn api_schema(&self) -> &str {
788 "mock"
789 }
790 async fn complete(
791 &self,
792 _request: &ConversationRequest,
793 ) -> Result<Completion, ProviderError> {
794 tokio::time::sleep(Duration::from_hours(1)).await;
795 Err(ProviderError::Transport("unreachable".into()))
796 }
797 }
798
799 struct WaitsForCancel;
802 #[async_trait]
803 impl Tool for WaitsForCancel {
804 type Args = Value;
805 type Output = EchoOut;
806 fn kind(&self) -> ToolKind {
807 ToolKind::Shell
808 }
809 fn description(&self) -> &str {
810 "waits"
811 }
812 async fn run(&self, ctx: &ToolCtx, _args: Value) -> Result<EchoOut, ToolError> {
813 ctx.cancel.cancelled().await;
814 Ok(EchoOut {
815 echoed: "stopped cooperatively".into(),
816 })
817 }
818 }
819
820 #[tokio::test]
821 async fn cancel_mid_sample_yields_cancelled_report() {
822 let events = Arc::new(Mutex::new(Vec::new()));
823 let sink_events = Arc::clone(&events);
824 let sink = Box::new(FnSink(move |event| {
825 sink_events.lock().unwrap().push(event);
826 }));
827 let mut s = Session::new(
828 Arc::new(HangingProvider),
829 Registry::new(),
830 vec![],
831 config(),
832 sink,
833 );
834 let handle = s.cancel_handle();
835 let canceller = tokio::spawn(async move {
836 tokio::time::sleep(Duration::from_millis(20)).await;
837 handle.cancel();
838 handle.cancel(); });
840 let report = s.run_text("go").await;
841 canceller.await.expect("canceller");
842
843 assert_eq!(report.status, Status::Cancelled);
844 assert_eq!(report.error, None, "cancelled is a stop, not a fault");
845 assert_eq!(report.final_message, None, "no assistant text this run");
846 assert_eq!(report.turns, 0, "no completion was accepted");
847 let roles: Vec<Role> = s.history().iter().map(|m| m.role).collect();
849 assert_eq!(roles, vec![Role::User]);
850 let evs = dump(&events);
852 assert!(
853 matches!(evs.last(), Some(Event::Result { report }) if report.status == Status::Cancelled)
854 );
855 }
856
857 #[tokio::test]
858 async fn cancel_mid_batch_pairs_the_rest_synthetically() {
859 let mut reg = Registry::new();
863 reg.register("waits", WaitsForCancel);
864 reg.register("echo", Echo);
865 let batch = Completion {
866 content: vec![
867 ContentBlock::ToolUse {
868 id: "c_wait".into(),
869 name: "waits".into(),
870 input: json!({}),
871 },
872 ContentBlock::ToolUse {
873 id: "c_echo".into(),
874 name: "echo".into(),
875 input: json!({}),
876 },
877 ],
878 usage: Usage::default(),
879 stop: StopReason::ToolUse,
880 };
881 let (s, events) = session_with(vec![Ok(batch)], reg, config());
882 let mut s = s; let handle = s.cancel_handle();
884 let canceller = tokio::spawn(async move {
885 tokio::time::sleep(Duration::from_millis(20)).await;
886 handle.cancel();
887 });
888 let report = s.run_text("go").await;
889 canceller.await.expect("canceller");
890
891 assert_eq!(report.status, Status::Cancelled);
892 assert_eq!(report.tool_calls.len(), 1);
895 assert_eq!(report.tool_calls[0].id, "c_wait");
896 assert!(report.tool_calls[0].ok);
897 assert_eq!(report.tool_calls[0].denial_reason, None);
898
899 let pairs: Vec<(String, bool)> = dump(&events)
901 .iter()
902 .filter_map(|e| match e {
903 Event::Message { message } if message.role == Role::User => Some(&message.content),
904 _ => None,
905 })
906 .flatten()
907 .filter_map(|b| match b {
908 ContentBlock::ToolResult {
909 tool_use_id,
910 is_error,
911 ..
912 } => Some((tool_use_id.clone(), *is_error)),
913 _ => None,
914 })
915 .collect();
916 assert_eq!(
917 pairs,
918 vec![("c_wait".into(), false), ("c_echo".into(), true)]
919 );
920 assert_eq!(
922 approvals(&events),
923 vec![("c_wait".into(), "waits".into(), "allow".into())]
924 );
925 }
926
927 #[tokio::test]
930 async fn cancelled_session_continues_on_the_next_run_with_a_fresh_token() {
931 let mut reg = Registry::new();
932 reg.register("waits", WaitsForCancel);
933 let (s, _e) = session_with(
934 vec![Ok(tool_turn("c1", "waits")), Ok(text_turn("second run"))],
935 reg,
936 config(),
937 );
938 let mut s = s;
939 let handle1 = s.cancel_handle();
940 let canceller = tokio::spawn(async move {
941 tokio::time::sleep(Duration::from_millis(20)).await;
942 handle1.cancel();
943 });
944 let r1 = s.run_text("q1").await;
945 canceller.await.expect("canceller");
946 assert_eq!(r1.status, Status::Cancelled);
947
948 assert!(!s.cancel_handle().is_cancelled());
951 let r2 = s.run_text("q2").await;
952 assert_eq!(r2.status, Status::Completed);
953 assert_eq!(r2.final_message.as_deref(), Some("second run"));
954 assert!(s.history().len() >= 4, "history: {:?}", s.history().len());
956 }
957
958 struct CapturingProvider {
963 inner: MockProvider,
964 requests: Arc<Mutex<Vec<Vec<Message>>>>,
965 }
966 #[async_trait]
967 impl Provider for CapturingProvider {
968 #[allow(clippy::unnecessary_literal_bound)]
969 fn api_schema(&self) -> &str {
970 "mock"
971 }
972 async fn complete(
973 &self,
974 request: &ConversationRequest,
975 ) -> Result<Completion, ProviderError> {
976 self.requests.lock().unwrap().push(request.messages.clone());
977 self.inner.complete(request).await
978 }
979 }
980
981 #[allow(clippy::type_complexity)]
983 fn capturing_session_with(
984 script: Vec<Result<Completion, ProviderError>>,
985 registry: Registry,
986 ) -> (
987 Session,
988 Arc<Mutex<Vec<Vec<Message>>>>,
989 Arc<Mutex<Vec<Event>>>,
990 ) {
991 let requests = Arc::new(Mutex::new(Vec::new()));
992 let events = Arc::new(Mutex::new(Vec::new()));
993 let sink_events = Arc::clone(&events);
994 let sink = Box::new(FnSink(move |event| {
995 sink_events.lock().unwrap().push(event);
996 }));
997 let provider = Arc::new(CapturingProvider {
998 inner: MockProvider::with_results(script),
999 requests: Arc::clone(&requests),
1000 });
1001 let session = Session::new(provider, registry, vec![], config(), sink);
1002 (session, requests, events)
1003 }
1004
1005 fn user_text(message: &Message) -> Option<&str> {
1006 match (message.role, message.content.as_slice()) {
1007 (Role::User, [ContentBlock::Text { text }]) => Some(text.as_str()),
1008 _ => None,
1009 }
1010 }
1011
1012 #[tokio::test]
1013 async fn second_run_continues_the_conversation() {
1014 let (mut s, requests, _e) = capturing_session_with(
1015 vec![
1016 Ok(text_turn("first answer")),
1017 Ok(text_turn("second answer")),
1018 ],
1019 Registry::new(),
1020 );
1021 let r1 = s.run_text("q1").await;
1022 let r2 = s.run_text("q2").await;
1023 assert_eq!(r1.status, Status::Completed);
1024 assert_eq!(r2.status, Status::Completed);
1025 assert_eq!(r2.final_message.as_deref(), Some("second answer"));
1026
1027 let reqs = requests.lock().unwrap();
1029 assert_eq!(reqs.len(), 2);
1030 let run2 = &reqs[1];
1031 assert_eq!(run2.len(), 3, "user q1, assistant, user q2: {run2:?}");
1032 assert_eq!(user_text(&run2[0]), Some("q1"));
1033 assert_eq!(run2[1].role, Role::Assistant);
1034 assert_eq!(user_text(&run2[2]), Some("q2"));
1035
1036 let roles: Vec<Role> = s.history().iter().map(|m| m.role).collect();
1038 assert_eq!(
1039 roles,
1040 vec![Role::User, Role::Assistant, Role::User, Role::Assistant]
1041 );
1042 }
1043
1044 fn capturing_with_cfg(
1048 script: Vec<Result<Completion, ProviderError>>,
1049 cfg: EngineConfig,
1050 ) -> (Session, Arc<Mutex<Vec<Vec<Message>>>>) {
1051 let requests = Arc::new(Mutex::new(Vec::new()));
1052 let provider = Arc::new(CapturingProvider {
1053 inner: MockProvider::with_results(script),
1054 requests: Arc::clone(&requests),
1055 });
1056 let session = Session::new(provider, Registry::new(), vec![], cfg, Box::new(NullSink));
1057 (session, requests)
1058 }
1059
1060 fn instr_config(cwd: std::path::PathBuf) -> EngineConfig {
1063 EngineConfig {
1064 cwd,
1065 instructions: locode_instructions::InstructionsConfig {
1066 global_file: false,
1067 ..Default::default()
1068 },
1069 ..config()
1070 }
1071 }
1072
1073 fn reminder_text(msgs: &[Message]) -> Option<String> {
1075 msgs.iter()
1076 .find_map(|m| match (m.role, m.content.as_slice()) {
1077 (Role::User, [ContentBlock::Text { text }])
1078 if text.starts_with("<system-reminder>") =>
1079 {
1080 Some(text.clone())
1081 }
1082 _ => None,
1083 })
1084 }
1085
1086 fn reminder_count(msgs: &[Message]) -> usize {
1087 msgs.iter()
1088 .filter(|m| {
1089 matches!(
1090 (m.role, m.content.as_slice()),
1091 (Role::User, [ContentBlock::Text { text }]) if text.starts_with("<system-reminder>")
1092 )
1093 })
1094 .count()
1095 }
1096
1097 #[tokio::test]
1098 async fn project_instructions_injected_once_before_prompt() {
1099 let dir = tempfile::tempdir().unwrap();
1100 let root = std::fs::canonicalize(dir.path()).unwrap();
1101 std::fs::create_dir(root.join(".git")).unwrap();
1102 std::fs::write(root.join("AGENTS.md"), "be terse").unwrap();
1103
1104 let (mut s, requests) = capturing_with_cfg(
1105 vec![Ok(text_turn("ok1")), Ok(text_turn("ok2"))],
1106 instr_config(root),
1107 );
1108 s.run_text("q1").await;
1109 s.run_text("q2").await;
1110
1111 let reqs = requests.lock().unwrap();
1112 let run1 = &reqs[0];
1114 let rem = reminder_text(run1).expect("instructions injected on run 1");
1115 assert!(rem.contains("## From:"), "labeled: {rem}");
1116 assert!(rem.contains("be terse"), "content present: {rem}");
1117 let rem_idx = run1
1118 .iter()
1119 .position(|m| reminder_text(std::slice::from_ref(m)).is_some());
1120 let q1_idx = run1.iter().position(|m| user_text(m) == Some("q1"));
1121 assert!(rem_idx < q1_idx, "reminder comes before the prompt");
1122
1123 assert_eq!(reminder_count(&reqs[1]), 1, "not re-injected on run 2");
1125 }
1126
1127 fn skills_config(cwd: std::path::PathBuf) -> EngineConfig {
1132 EngineConfig {
1133 cwd: cwd.clone(),
1134 skills: locode_skills::SkillsConfig::enabled(),
1135 ..instr_config(cwd)
1136 }
1137 }
1138
1139 fn write_skill(root: &std::path::Path, name: &str, description: &str) {
1140 let dir = root.join(".agents/skills").join(name);
1141 std::fs::create_dir_all(&dir).unwrap();
1142 std::fs::write(
1143 dir.join("SKILL.md"),
1144 format!("---\nname: {name}\ndescription: {description}\n---\n# {name}\n"),
1145 )
1146 .unwrap();
1147 }
1148
1149 #[tokio::test]
1152 async fn skills_listing_injected_once_then_quiet() {
1153 let dir = tempfile::tempdir().unwrap();
1154 let root = std::fs::canonicalize(dir.path()).unwrap();
1155 std::fs::create_dir(root.join(".git")).unwrap();
1156 write_skill(&root, "commit", "Make a commit");
1157
1158 let (mut s, requests) = capturing_with_cfg(
1159 vec![Ok(text_turn("a")), Ok(text_turn("b"))],
1160 skills_config(root),
1161 );
1162 s.run_text("q1").await;
1163 s.run_text("q2").await;
1164
1165 let reqs = requests.lock().unwrap();
1166 let listing = reqs[0]
1167 .iter()
1168 .filter_map(|m| reminder_text(std::slice::from_ref(m)))
1169 .find(|t| t.contains("skills are available"))
1170 .expect("listing injected");
1171 assert!(listing.contains("- commit: Make a commit"), "{listing}");
1172 assert!(listing.contains("Absolute path:"), "{listing}");
1173
1174 let count = |msgs: &[Message]| {
1175 msgs.iter()
1176 .filter(|m| {
1177 reminder_text(std::slice::from_ref(m))
1178 .is_some_and(|t| t.contains("skills are available"))
1179 })
1180 .count()
1181 };
1182 assert_eq!(count(&reqs[1]), 1, "unchanged ⇒ not re-sent");
1183 }
1184
1185 #[tokio::test]
1193 async fn adding_a_skill_re_sends_the_entire_listing() {
1194 let dir = tempfile::tempdir().unwrap();
1195 let root = std::fs::canonicalize(dir.path()).unwrap();
1196 std::fs::create_dir(root.join(".git")).unwrap();
1197 write_skill(&root, "commit", "Make a commit");
1198
1199 let (mut s, requests) = capturing_with_cfg(
1200 vec![Ok(text_turn("a")), Ok(text_turn("b")), Ok(text_turn("c"))],
1201 skills_config(root.clone()),
1202 );
1203 s.run_text("q1").await;
1204 write_skill(&root, "review", "Review a diff"); s.run_text("q2").await; s.run_text("q3").await;
1207
1208 let reqs = requests.lock().unwrap();
1209 let listing = |msgs: &[Message]| {
1210 msgs.iter()
1211 .filter_map(|m| reminder_text(std::slice::from_ref(m)))
1212 .rfind(|t| t.contains("skills are available"))
1213 };
1214 assert!(
1215 !listing(&reqs[1]).unwrap().contains("- review:"),
1216 "not yet — the scan that would see it runs at the end of this run"
1217 );
1218 let third = listing(&reqs[2]).expect("re-sent");
1219 assert!(third.contains("- commit:"), "old skill included: {third}");
1220 assert!(third.contains("- review:"), "new skill included: {third}");
1221 }
1222
1223 #[tokio::test]
1226 async fn removing_the_last_skill_announces_it() {
1227 let dir = tempfile::tempdir().unwrap();
1228 let root = std::fs::canonicalize(dir.path()).unwrap();
1229 std::fs::create_dir(root.join(".git")).unwrap();
1230 write_skill(&root, "commit", "Make a commit");
1231
1232 let (mut s, requests) = capturing_with_cfg(
1233 vec![Ok(text_turn("a")), Ok(text_turn("b")), Ok(text_turn("c"))],
1234 skills_config(root.clone()),
1235 );
1236 s.run_text("q1").await;
1237 std::fs::remove_dir_all(root.join(".agents/skills/commit")).unwrap();
1238 s.run_text("q2").await; s.run_text("q3").await;
1240
1241 let reqs = requests.lock().unwrap();
1242 let last = reqs[2]
1243 .iter()
1244 .filter_map(|m| reminder_text(std::slice::from_ref(m)))
1245 .next_back()
1246 .expect("a reminder");
1247 assert!(last.contains("No skills are currently available"), "{last}");
1248 }
1249
1250 #[tokio::test]
1252 async fn no_skills_ever_means_no_message_at_all() {
1253 let dir = tempfile::tempdir().unwrap();
1254 let root = std::fs::canonicalize(dir.path()).unwrap();
1255 std::fs::create_dir(root.join(".git")).unwrap();
1256
1257 let (mut s, requests) = capturing_with_cfg(vec![Ok(text_turn("a"))], skills_config(root));
1258 s.run_text("q1").await;
1259
1260 let reqs = requests.lock().unwrap();
1261 assert!(
1262 !reqs[0]
1263 .iter()
1264 .any(|m| reminder_text(std::slice::from_ref(m))
1265 .is_some_and(|t| t.contains("skills"))),
1266 "silence, not a denial"
1267 );
1268 }
1269
1270 #[tokio::test]
1271 async fn project_instructions_absent_when_disabled() {
1272 let dir = tempfile::tempdir().unwrap();
1273 let root = std::fs::canonicalize(dir.path()).unwrap();
1274 std::fs::create_dir(root.join(".git")).unwrap();
1275 std::fs::write(root.join("AGENTS.md"), "be terse").unwrap();
1276 let mut cfg = instr_config(root);
1277 cfg.instructions.enabled = false;
1278
1279 let (mut s, requests) = capturing_with_cfg(vec![Ok(text_turn("ok"))], cfg);
1280 s.run_text("q").await;
1281 assert!(reminder_text(&requests.lock().unwrap()[0]).is_none());
1282 }
1283
1284 #[tokio::test]
1285 async fn project_instructions_absent_when_no_agents_md() {
1286 let dir = tempfile::tempdir().unwrap();
1287 let root = std::fs::canonicalize(dir.path()).unwrap();
1288 let (mut s, requests) = capturing_with_cfg(vec![Ok(text_turn("ok"))], instr_config(root));
1290 s.run_text("q").await;
1291 assert!(reminder_text(&requests.lock().unwrap()[0]).is_none());
1292 }
1293
1294 #[tokio::test]
1295 async fn project_instructions_replace_banner_on_edit() {
1296 let dir = tempfile::tempdir().unwrap();
1297 let root = std::fs::canonicalize(dir.path()).unwrap();
1298 std::fs::create_dir(root.join(".git")).unwrap();
1299 let agents = root.join("AGENTS.md");
1300 std::fs::write(&agents, "v1 rules").unwrap();
1301
1302 let (mut s, requests) = capturing_with_cfg(
1303 vec![Ok(text_turn("ok1")), Ok(text_turn("ok2"))],
1304 instr_config(root),
1305 );
1306 s.run_text("q1").await;
1307 std::fs::write(&agents, "v2 rules").unwrap(); s.run_text("q2").await;
1309
1310 let reqs = requests.lock().unwrap();
1311 let run2 = &reqs[1];
1312 let banner = run2
1313 .iter()
1314 .find_map(|m| match (m.role, m.content.as_slice()) {
1315 (Role::User, [ContentBlock::Text { text }])
1316 if text.contains("replace all previously provided") =>
1317 {
1318 Some(text.clone())
1319 }
1320 _ => None,
1321 })
1322 .expect("replace banner on edit");
1323 assert!(banner.contains("v2 rules"), "new content: {banner}");
1324 assert!(!banner.contains("v1 rules"), "not the old content");
1325 assert_eq!(reminder_count(run2), 2);
1327 }
1328
1329 #[tokio::test]
1330 async fn project_instructions_removal_banner_on_delete() {
1331 let dir = tempfile::tempdir().unwrap();
1332 let root = std::fs::canonicalize(dir.path()).unwrap();
1333 std::fs::create_dir(root.join(".git")).unwrap();
1334 let agents = root.join("AGENTS.md");
1335 std::fs::write(&agents, "rules").unwrap();
1336
1337 let (mut s, requests) = capturing_with_cfg(
1338 vec![Ok(text_turn("ok1")), Ok(text_turn("ok2"))],
1339 instr_config(root),
1340 );
1341 s.run_text("q1").await;
1342 std::fs::remove_file(&agents).unwrap(); s.run_text("q2").await;
1344
1345 let reqs = requests.lock().unwrap();
1346 assert!(
1347 reqs[1].iter().any(|m| matches!(
1348 (m.role, m.content.as_slice()),
1349 (Role::User, [ContentBlock::Text { text }]) if text.contains("no longer apply")
1350 )),
1351 "removal notice on delete"
1352 );
1353 }
1354
1355 #[tokio::test]
1356 async fn project_instructions_not_reinjected_when_unchanged() {
1357 let dir = tempfile::tempdir().unwrap();
1358 let root = std::fs::canonicalize(dir.path()).unwrap();
1359 std::fs::create_dir(root.join(".git")).unwrap();
1360 std::fs::write(root.join("AGENTS.md"), "stable").unwrap();
1361
1362 let (mut s, requests) = capturing_with_cfg(
1363 vec![
1364 Ok(text_turn("ok1")),
1365 Ok(text_turn("ok2")),
1366 Ok(text_turn("ok3")),
1367 ],
1368 instr_config(root),
1369 );
1370 s.run_text("q1").await;
1371 s.run_text("q2").await;
1372 s.run_text("q3").await;
1373 assert_eq!(reminder_count(&requests.lock().unwrap()[2]), 1);
1375 }
1376
1377 #[tokio::test]
1378 async fn init_emitted_once_across_runs_with_one_result_each() {
1379 let (mut s, events) = session_with(
1380 vec![Ok(text_turn("one")), Ok(text_turn("two"))],
1381 Registry::new(),
1382 config(),
1383 );
1384 let _ = s.run_text("q1").await;
1385 let _ = s.run_text("q2").await;
1386 let evs = dump(&events);
1387 let inits = evs
1388 .iter()
1389 .filter(|e| matches!(e, Event::Init { .. }))
1390 .count();
1391 let results = evs
1392 .iter()
1393 .filter(|e| matches!(e, Event::Result { .. }))
1394 .count();
1395 assert_eq!(inits, 1, "Init is once per session, not per run");
1396 assert_eq!(results, 2, "one Result per run");
1397 assert!(
1398 matches!(evs.first(), Some(Event::Init { .. })),
1399 "Init still opens the stream"
1400 );
1401 }
1402
1403 #[tokio::test]
1404 async fn report_counts_are_per_run_not_cumulative() {
1405 let mut t1 = tool_turn("c1", "echo");
1408 t1.usage = Usage {
1409 input_tokens: 10,
1410 output_tokens: 5,
1411 ..Usage::default()
1412 };
1413 let t2 = text_turn("done one");
1414 let mut t3 = text_turn("done two");
1415 t3.usage = Usage {
1416 input_tokens: 20,
1417 output_tokens: 7,
1418 ..Usage::default()
1419 };
1420 let (mut s, _e) = session_with(vec![Ok(t1), Ok(t2), Ok(t3)], echo_registry(), config());
1421 let r1 = s.run_text("q1").await;
1422 let r2 = s.run_text("q2").await;
1423 assert_eq!(r1.turns, 2);
1424 assert_eq!(r1.tool_calls.len(), 1);
1425 assert_eq!(r2.turns, 1, "run 2 counts its own turns only");
1426 assert!(r2.tool_calls.is_empty());
1427 assert_eq!(r2.usage.input_tokens, 20, "usage is per-run");
1428 assert_eq!(r2.usage.output_tokens, 7);
1429 }
1430
1431 #[tokio::test]
1434 async fn two_run_stream_reconstructs_the_full_conversation() {
1435 let (mut s, events) = session_with(
1436 vec![
1437 Ok(tool_turn("c1", "echo")),
1438 Ok(text_turn("done one")),
1439 Ok(text_turn("done two")),
1440 ],
1441 echo_registry(),
1442 config(),
1443 );
1444 let _ = s.run_text("q1").await;
1445 let _ = s.run_text("q2").await;
1446 let rebuilt: Conversation = reconstruct_conversation(&dump(&events));
1447 let roles: Vec<Role> = rebuilt.messages.iter().map(|m| m.role).collect();
1450 assert_eq!(
1451 roles,
1452 vec![
1453 Role::User,
1454 Role::Assistant,
1455 Role::User,
1456 Role::Assistant,
1457 Role::User,
1458 Role::Assistant,
1459 ]
1460 );
1461 assert_eq!(rebuilt.messages.as_slice(), s.history());
1463 }
1464
1465 #[tokio::test]
1468 async fn continues_after_model_error() {
1469 let (mut s, requests, _e) = capturing_session_with(
1470 vec![Err(ProviderError::ContextOverflow), Ok(text_turn("ok now"))],
1471 Registry::new(),
1472 );
1473 let r1 = s.run_text("q1").await;
1474 let r2 = s.run_text("q2").await;
1475 assert_eq!(r1.status, Status::ModelError);
1476 assert_eq!(r2.status, Status::Completed);
1477 let reqs = requests.lock().unwrap();
1479 let run2 = &reqs[1];
1480 assert_eq!(run2.len(), 2, "user q1 + user q2: {run2:?}");
1481 assert_eq!(user_text(&run2[0]), Some("q1"));
1482 assert_eq!(user_text(&run2[1]), Some("q2"));
1483 }
1484
1485 #[tokio::test]
1488 async fn continues_after_fatal_tool_error_with_valid_pairing() {
1489 let mut reg = Registry::new();
1490 reg.register("boom", Boom);
1491 let (mut s, requests, _e) = capturing_session_with(
1492 vec![Ok(tool_turn("c1", "boom")), Ok(text_turn("recovered"))],
1493 reg,
1494 );
1495 let r1 = s.run_text("q1").await;
1496 let r2 = s.run_text("q2").await;
1497 assert_eq!(r1.status, Status::Error);
1498 assert_eq!(r2.status, Status::Completed);
1499
1500 let reqs = requests.lock().unwrap();
1503 let run2 = &reqs[1];
1504 assert_eq!(run2.len(), 4, "q1, assistant, tool_result, q2: {run2:?}");
1505 assert!(
1506 run2[1]
1507 .content
1508 .iter()
1509 .any(|b| matches!(b, ContentBlock::ToolUse { id, .. } if id == "c1"))
1510 );
1511 assert!(run2[2].content.iter().any(|b| matches!(
1512 b,
1513 ContentBlock::ToolResult { tool_use_id, is_error: true, .. } if tool_use_id == "c1"
1514 )));
1515 assert_eq!(user_text(&run2[3]), Some("q2"));
1516 }
1517
1518 #[tokio::test]
1519 async fn usage_is_summed_across_turns() {
1520 let mut first = tool_turn("c1", "echo");
1521 first.usage = Usage {
1522 input_tokens: 10,
1523 output_tokens: 5,
1524 ..Usage::default()
1525 };
1526 let mut second = text_turn("done");
1527 second.usage = Usage {
1528 input_tokens: 20,
1529 output_tokens: 7,
1530 ..Usage::default()
1531 };
1532 let (mut s, _e) = session_with(vec![Ok(first), Ok(second)], echo_registry(), config());
1533 let report = s.run_text("go").await;
1534 assert_eq!(report.usage.input_tokens, 30);
1535 assert_eq!(report.usage.output_tokens, 12);
1536 }
1537}