Skip to main content

locode_engine/
lib.rs

1//! locode-engine — the sample→dispatch→append loop and the [`Session`] driving API
2//! (ADR-0005, ADR-0004, ADR-0014).
3//!
4//! A [`Session`] drives one run to a terminal [`locode_protocol::Status`] against any
5//! [`locode_provider::Provider`], dispatching tool calls through a
6//! [`locode_tools::Registry`], emitting `stream-json` events to an [`EventSink`], and
7//! returning one [`locode_protocol::Report`]. Proven end-to-end against
8//! `MockProvider` with zero network.
9
10mod 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};
21// The type `Session::cancel_handle` returns (ADR-0018) — re-exported so
22// frontends need no direct tokio-util dependency.
23pub use tokio_util::sync::CancellationToken;
24
25#[cfg(test)]
26mod tests {
27    // Test tools return `&'static str` literals from `description`; the trait ties it
28    // to `&self` so real tools can return a stored field.
29    #![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    // ---- trivial in-test tools ----
47
48    #[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    // ---- harness ----
93
94    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, // no real sleeps in tests
123            // Project-instruction loading off by default in the loop tests: cwd is the
124            // crate dir, so a live loader would inject this repo's own AGENTS.md and
125            // perturb the exact event/message assertions. The injection path has its own
126            // dedicated tests below (`project_instructions_*`).
127            instructions: locode_instructions::InstructionsConfig {
128                enabled: false,
129                ..Default::default()
130            },
131            ..EngineConfig::default()
132        }
133    }
134
135    /// Build a session with a scripted provider + registry, collecting events.
136    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    // ---- terminal-state matrix ----
162
163    #[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        // Init, Message(user), Message(assistant), Result.
174        let evs = dump(&events);
175        assert!(matches!(evs.first(), Some(Event::Init { .. })));
176        assert!(matches!(evs.last(), Some(Event::Result { .. })));
177    }
178
179    // ---- streaming (ADR-0021 slice 1) ----
180
181    #[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        // The deltas concatenate exactly to the finalized assistant text...
199        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        // ...and there was more than one (proving it actually streamed).
208        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        // The whole assistant Message is STILL appended — deltas don't replace it.
214        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        // Deltas precede the finalized assistant Message in the stream.
222        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        // Streaming is display-only: the Report is identical either way.
263        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        // Always asks for a tool → never completes; ceiling of 2.
287        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        // Retryable every time → 1 + resample_retries attempts, then ModelError.
307        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        // Two non-terminal Error retry notes emitted (resample_retries == 2).
318        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        // The boom call still produced a paired (is_error) record.
350        assert_eq!(report.tool_calls.len(), 1);
351        assert!(!report.tool_calls[0].ok);
352    }
353
354    /// An empty completion (no text, no tool calls — e.g. a reasoning-only
355    /// turn truncated by `max_output_tokens`) is resampled, not labeled
356    /// Completed (ADR-0005 amendment 2026-07-19; grok's `is_empty` rule).
357    #[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        // resample_retries = 2 → initial + 2 resamples, all empty → ModelError.
388        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    // ---- transcript hygiene ----
408
409    #[tokio::test]
410    async fn mid_batch_abort_synthesizes_results() {
411        // One assistant turn asks for TWO tools: boom (Fatal) then echo. echo must
412        // not run, yet both tool_use ids must be answered in the transcript.
413        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        // The appended tool-result message pairs BOTH ids.
437        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        // boom recorded (ran, fatal); echo NOT recorded (never executed).
456        assert_eq!(report.tool_calls.len(), 1);
457    }
458
459    // ---- replay + stream fidelity ----
460
461    #[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        // The emitted assistant message preserves the Thinking block + signature.
483        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        // user + assistant(tool_use) + user(tool_result) + assistant(text).
510        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    // ---- the approval seam (ADR-0017) ----
518
519    use std::sync::atomic::{AtomicUsize, Ordering};
520
521    /// A tool that counts its executions — proves a denied call never ran.
522    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    /// Denies tools whose name is in the list; allows everything else. Records
544    /// the `kind` seen on each request so tests can assert it is populated.
545    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        // Soft: the run continued to Completed; the tool never executed.
599        assert_eq!(report.status, Status::Completed);
600        assert_eq!(ran.load(Ordering::SeqCst), 0, "denied tool must not run");
601
602        // The record: ok=false, denial_reason set (and only here), no output.
603        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        // The transcript: a paired is_error result carrying the reason.
613        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        // The trace: a deny Approval event for c1.
631        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        // Both calls answered, in call order, denied first.
669        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        // Records: denied (with reason) then executed (without).
688        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        // Approval trace: deny then allow, in order.
695        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    /// An approver that suspends on a oneshot until an external task resolves
728    /// it — the exact shape of a TUI prompt. Proves the engine awaits the
729    /// decision without deadlocking the run.
730    #[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        // Resolve the prompt from "the UI" after the run has started.
750        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        // The default AllowAll approver still journals every resolution.
764        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        // And denial_reason is absent on ordinary success records.
776        assert_eq!(report.tool_calls[0].denial_reason, None);
777    }
778
779    // ---- cancellation (ADR-0018) ----
780
781    /// A provider whose sample never returns on its own — cancellation is the
782    /// only way out (models a long in-flight request).
783    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    /// A tool that parks on its ctx cancel token and returns cleanly once it
800    /// fires — the cooperative-cancel shape the host implements for real.
801    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(); // idempotent double-cancel
839        });
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        // No assistant message was appended: history is user-prompt only.
848        let roles: Vec<Role> = s.history().iter().map(|m| m.role).collect();
849        assert_eq!(roles, vec![Role::User]);
850        // The stream still terminates in a Result carrying the same report.
851        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        // One turn asks for TWO tools: a cooperative waiter, then echo. The
860        // cancel fires while the waiter runs → its own result is real; echo is
861        // never run (no approval consult, no record) but still paired.
862        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; // provider script has ONE turn: cancel must end the run
883        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        // The waiter executed (cooperatively) and is the only record; no
893        // cancellation synthetic ever carries denial_reason.
894        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        // Both tool_use ids are answered: real result + cancellation synthetic.
900        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        // Only the executed call was consulted for approval.
921        assert_eq!(
922            approvals(&events),
923            vec![("c_wait".into(), "waits".into(), "allow".into())]
924        );
925    }
926
927    /// The token is per-run (ADR-0018 Decision 1): a cancelled run 1 must not
928    /// poison run 2, and run 2 continues the same conversation (with ADR-0016).
929    #[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        // The retired handle stays cancelled, but the session got a fresh
949        // token at run end — run 2 must not see the old cancel.
950        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        // Continuity intact: q1's turns are still in the history.
955        assert!(s.history().len() >= 4, "history: {:?}", s.history().len());
956    }
957
958    // ---- session continuity (ADR-0016) ----
959
960    /// A scripted provider that also records each request's message array, so a
961    /// test can assert what the model actually saw on a follow-up run.
962    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    /// Like `session_with`, but the provider records every request's messages.
982    #[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        // Run 2's request contains run 1's full exchange, then the new prompt.
1028        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        // The public accessor exposes the same transcript (empty test preamble).
1037        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    // ---- project-instruction injection (ADR-0023, Task 30) ----
1045
1046    /// A capturing session over a custom config (for the injection path).
1047    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    /// A config with project instructions **enabled**, rooted at `cwd`, global file off
1061    /// (never read the real `~/.locode` in tests).
1062    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    /// The injected `<system-reminder>` text within a request's messages, if present.
1074    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        // Run 1: injected, with the file content, positioned before the user prompt.
1113        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        // Run 2: NOT re-injected — still exactly the one from run 1 (once per session).
1124        assert_eq!(reminder_count(&reqs[1]), 1, "not re-injected on run 2");
1125    }
1126
1127    /// A session whose transcript is `replayed` — what resume does (the recovered
1128    /// history becomes the preamble).
1129    fn resumed_with_cfg(
1130        script: Vec<Result<Completion, ProviderError>>,
1131        cfg: EngineConfig,
1132        replayed: Vec<Message>,
1133    ) -> (Session, Arc<Mutex<Vec<Vec<Message>>>>) {
1134        let requests = Arc::new(Mutex::new(Vec::new()));
1135        let provider = Arc::new(CapturingProvider {
1136            inner: MockProvider::with_results(script),
1137            requests: Arc::clone(&requests),
1138        });
1139        let session = Session::new(provider, Registry::new(), replayed, cfg, Box::new(NullSink));
1140        (session, requests)
1141    }
1142
1143    /// ADR-0023's Refresh rule says instructions are "never double-injected on
1144    /// fork/resume". A remembered hash cannot deliver that — a resumed session starts
1145    /// with the field empty and re-sends instructions the replayed transcript already
1146    /// carries. The check reads the conversation instead.
1147    #[tokio::test]
1148    async fn resuming_does_not_re_inject_unchanged_instructions() {
1149        let dir = tempfile::tempdir().unwrap();
1150        let root = std::fs::canonicalize(dir.path()).unwrap();
1151        std::fs::create_dir(root.join(".git")).unwrap();
1152        std::fs::write(root.join("AGENTS.md"), "be terse").unwrap();
1153
1154        let (mut first, requests) =
1155            capturing_with_cfg(vec![Ok(text_turn("ok"))], instr_config(root.clone()));
1156        first.run_text("q1").await;
1157        let replayed = requests.lock().unwrap()[0].clone();
1158        assert_eq!(reminder_count(&replayed), 1, "precondition: injected once");
1159
1160        let (mut resumed, resumed_requests) =
1161            resumed_with_cfg(vec![Ok(text_turn("ok2"))], instr_config(root), replayed);
1162        resumed.run_text("q2").await;
1163        let reqs = resumed_requests.lock().unwrap();
1164        assert_eq!(
1165            reminder_count(&reqs[0]),
1166            1,
1167            "still the one from before the resume, not a second copy: {:#?}",
1168            reqs[0]
1169        );
1170    }
1171
1172    /// …but an `AGENTS.md` edited between sessions *is* re-injected, with the replace
1173    /// banner that says the earlier copy no longer applies.
1174    #[tokio::test]
1175    async fn resuming_re_injects_instructions_that_changed_while_away() {
1176        let dir = tempfile::tempdir().unwrap();
1177        let root = std::fs::canonicalize(dir.path()).unwrap();
1178        std::fs::create_dir(root.join(".git")).unwrap();
1179        std::fs::write(root.join("AGENTS.md"), "be terse").unwrap();
1180
1181        let (mut first, requests) =
1182            capturing_with_cfg(vec![Ok(text_turn("ok"))], instr_config(root.clone()));
1183        first.run_text("q1").await;
1184        let replayed = requests.lock().unwrap()[0].clone();
1185
1186        std::fs::write(root.join("AGENTS.md"), "be verbose").unwrap();
1187        let (mut resumed, resumed_requests) =
1188            resumed_with_cfg(vec![Ok(text_turn("ok2"))], instr_config(root), replayed);
1189        resumed.run_text("q2").await;
1190        let reqs = resumed_requests.lock().unwrap();
1191        assert_eq!(reminder_count(&reqs[0]), 2, "the new body joins the old");
1192        let latest = reqs[0]
1193            .iter()
1194            .rev()
1195            .find_map(|m| reminder_text(std::slice::from_ref(m)))
1196            .expect("a reminder");
1197        assert!(latest.contains("be verbose"), "{latest}");
1198        assert!(
1199            latest.contains("replace all previously provided"),
1200            "banner present: {latest}"
1201        );
1202    }
1203
1204    /// The other half of ADR-0023's Refresh rule: instructions dropped from the
1205    /// conversation come back. A remembered hash would keep saying "already sent".
1206    #[tokio::test]
1207    async fn instructions_dropped_from_the_transcript_are_re_injected() {
1208        let dir = tempfile::tempdir().unwrap();
1209        let root = std::fs::canonicalize(dir.path()).unwrap();
1210        std::fs::create_dir(root.join(".git")).unwrap();
1211        std::fs::write(root.join("AGENTS.md"), "be terse").unwrap();
1212
1213        let (mut first, requests) =
1214            capturing_with_cfg(vec![Ok(text_turn("ok"))], instr_config(root.clone()));
1215        first.run_text("q1").await;
1216        // Compaction: keep the conversation, drop the reminder.
1217        let compacted: Vec<Message> = requests.lock().unwrap()[0]
1218            .iter()
1219            .filter(|m| reminder_text(std::slice::from_ref(m)).is_none())
1220            .cloned()
1221            .collect();
1222
1223        let (mut after, after_requests) =
1224            resumed_with_cfg(vec![Ok(text_turn("ok2"))], instr_config(root), compacted);
1225        after.run_text("q2").await;
1226        assert_eq!(
1227            reminder_count(&after_requests.lock().unwrap()[0]),
1228            1,
1229            "re-injected after being compacted away"
1230        );
1231    }
1232
1233    /// A skills config rooted at `cwd` that never reads the real `~/.locode`
1234    /// (`discover` resolves the home root itself, so the temp repo is the only source
1235    /// as long as no skill exists in the developer's home — the project root is what
1236    /// this asserts on).
1237    fn skills_config(cwd: std::path::PathBuf) -> EngineConfig {
1238        EngineConfig {
1239            cwd: cwd.clone(),
1240            skills: locode_skills::SkillsConfig::enabled(),
1241            ..instr_config(cwd)
1242        }
1243    }
1244
1245    fn write_skill(root: &std::path::Path, name: &str, description: &str) {
1246        let dir = root.join(".agents/skills").join(name);
1247        std::fs::create_dir_all(&dir).unwrap();
1248        std::fs::write(
1249            dir.join("SKILL.md"),
1250            format!("---\nname: {name}\ndescription: {description}\n---\n# {name}\n"),
1251        )
1252        .unwrap();
1253    }
1254
1255    /// Injected once, before the prompt; and a second turn with no change sends nothing
1256    /// new — the whole-body comparison is what makes the steady state quiet.
1257    #[tokio::test]
1258    async fn skills_listing_injected_once_then_quiet() {
1259        let dir = tempfile::tempdir().unwrap();
1260        let root = std::fs::canonicalize(dir.path()).unwrap();
1261        std::fs::create_dir(root.join(".git")).unwrap();
1262        write_skill(&root, "commit", "Make a commit");
1263
1264        let (mut s, requests) = capturing_with_cfg(
1265            vec![Ok(text_turn("a")), Ok(text_turn("b"))],
1266            skills_config(root),
1267        );
1268        s.run_text("q1").await;
1269        s.run_text("q2").await;
1270
1271        let reqs = requests.lock().unwrap();
1272        let listing = reqs[0]
1273            .iter()
1274            .filter_map(|m| reminder_text(std::slice::from_ref(m)))
1275            .find(|t| t.contains("skills are available"))
1276            .expect("listing injected");
1277        assert!(listing.contains(r#"<skill name="commit""#), "{listing}");
1278        assert!(listing.contains("Make a commit"), "{listing}");
1279        assert!(
1280            listing.contains("SKILL.md\">"),
1281            "the path attribute: {listing}"
1282        );
1283
1284        let count = |msgs: &[Message]| {
1285            msgs.iter()
1286                .filter(|m| {
1287                    reminder_text(std::slice::from_ref(m))
1288                        .is_some_and(|t| t.contains("skills are available"))
1289                })
1290                .count()
1291        };
1292        assert_eq!(count(&reqs[1]), 1, "unchanged ⇒ not re-sent");
1293    }
1294
1295    /// Adding a skill re-sends the **whole** listing, not just the new entry — the
1296    /// defect the per-skill delta in Claude Code and grok produces (ADR-0025 §3.1).
1297    ///
1298    /// Also pins the timing consequence of §3.2: the scan runs *after* a run finishes,
1299    /// so a skill created while the user is typing is picked up by the **next** run's
1300    /// post-run scan and injected on the turn after that. Writing it during a run — the
1301    /// case the design optimizes for — costs no extra turn.
1302    #[tokio::test]
1303    async fn adding_a_skill_re_sends_the_entire_listing() {
1304        let dir = tempfile::tempdir().unwrap();
1305        let root = std::fs::canonicalize(dir.path()).unwrap();
1306        std::fs::create_dir(root.join(".git")).unwrap();
1307        write_skill(&root, "commit", "Make a commit");
1308
1309        let (mut s, requests) = capturing_with_cfg(
1310            vec![Ok(text_turn("a")), Ok(text_turn("b")), Ok(text_turn("c"))],
1311            skills_config(root.clone()),
1312        );
1313        s.run_text("q1").await;
1314        write_skill(&root, "review", "Review a diff"); // after run 1's post-run scan
1315        s.run_text("q2").await; // still the cached body; run 2's scan picks it up
1316        s.run_text("q3").await;
1317
1318        let reqs = requests.lock().unwrap();
1319        let listing = |msgs: &[Message]| {
1320            msgs.iter()
1321                .filter_map(|m| reminder_text(std::slice::from_ref(m)))
1322                .rfind(|t| t.contains("skills are available"))
1323        };
1324        assert!(
1325            !listing(&reqs[1]).unwrap().contains(r#"name="review""#),
1326            "not yet — the scan that would see it runs at the end of this run"
1327        );
1328        let third = listing(&reqs[2]).expect("re-sent");
1329        assert!(
1330            third.contains(r#"name="commit""#),
1331            "old skill included: {third}"
1332        );
1333        assert!(
1334            third.contains(r#"name="review""#),
1335            "new skill included: {third}"
1336        );
1337    }
1338
1339    /// Removing the last skill says so, rather than going quiet and leaving a stale
1340    /// instruction standing (ADR-0025 §3.1 — codex's behavior, not the other two's).
1341    #[tokio::test]
1342    async fn removing_the_last_skill_announces_it() {
1343        let dir = tempfile::tempdir().unwrap();
1344        let root = std::fs::canonicalize(dir.path()).unwrap();
1345        std::fs::create_dir(root.join(".git")).unwrap();
1346        write_skill(&root, "commit", "Make a commit");
1347
1348        let (mut s, requests) = capturing_with_cfg(
1349            vec![Ok(text_turn("a")), Ok(text_turn("b")), Ok(text_turn("c"))],
1350            skills_config(root.clone()),
1351        );
1352        s.run_text("q1").await;
1353        std::fs::remove_dir_all(root.join(".agents/skills/commit")).unwrap();
1354        s.run_text("q2").await; // run 2's post-run scan observes the removal
1355        s.run_text("q3").await;
1356
1357        let reqs = requests.lock().unwrap();
1358        let last = reqs[2]
1359            .iter()
1360            .filter_map(|m| reminder_text(std::slice::from_ref(m)))
1361            .next_back()
1362            .expect("a reminder");
1363        assert!(last.contains("No skills are currently available"), "{last}");
1364    }
1365
1366    /// A project with no skills at all must not open with a pointless denial.
1367    #[tokio::test]
1368    async fn no_skills_ever_means_no_message_at_all() {
1369        let dir = tempfile::tempdir().unwrap();
1370        let root = std::fs::canonicalize(dir.path()).unwrap();
1371        std::fs::create_dir(root.join(".git")).unwrap();
1372
1373        let (mut s, requests) = capturing_with_cfg(vec![Ok(text_turn("a"))], skills_config(root));
1374        s.run_text("q1").await;
1375
1376        let reqs = requests.lock().unwrap();
1377        assert!(
1378            !reqs[0]
1379                .iter()
1380                .any(|m| reminder_text(std::slice::from_ref(m))
1381                    .is_some_and(|t| t.contains("skills"))),
1382            "silence, not a denial"
1383        );
1384    }
1385
1386    #[tokio::test]
1387    async fn project_instructions_absent_when_disabled() {
1388        let dir = tempfile::tempdir().unwrap();
1389        let root = std::fs::canonicalize(dir.path()).unwrap();
1390        std::fs::create_dir(root.join(".git")).unwrap();
1391        std::fs::write(root.join("AGENTS.md"), "be terse").unwrap();
1392        let mut cfg = instr_config(root);
1393        cfg.instructions.enabled = false;
1394
1395        let (mut s, requests) = capturing_with_cfg(vec![Ok(text_turn("ok"))], cfg);
1396        s.run_text("q").await;
1397        assert!(reminder_text(&requests.lock().unwrap()[0]).is_none());
1398    }
1399
1400    #[tokio::test]
1401    async fn project_instructions_absent_when_no_agents_md() {
1402        let dir = tempfile::tempdir().unwrap();
1403        let root = std::fs::canonicalize(dir.path()).unwrap();
1404        // Empty tempdir — no `.git`, no `AGENTS.md` → nothing discovered.
1405        let (mut s, requests) = capturing_with_cfg(vec![Ok(text_turn("ok"))], instr_config(root));
1406        s.run_text("q").await;
1407        assert!(reminder_text(&requests.lock().unwrap()[0]).is_none());
1408    }
1409
1410    #[tokio::test]
1411    async fn project_instructions_replace_banner_on_edit() {
1412        let dir = tempfile::tempdir().unwrap();
1413        let root = std::fs::canonicalize(dir.path()).unwrap();
1414        std::fs::create_dir(root.join(".git")).unwrap();
1415        let agents = root.join("AGENTS.md");
1416        std::fs::write(&agents, "v1 rules").unwrap();
1417
1418        let (mut s, requests) = capturing_with_cfg(
1419            vec![Ok(text_turn("ok1")), Ok(text_turn("ok2"))],
1420            instr_config(root),
1421        );
1422        s.run_text("q1").await;
1423        std::fs::write(&agents, "v2 rules").unwrap(); // edit between runs
1424        s.run_text("q2").await;
1425
1426        let reqs = requests.lock().unwrap();
1427        let run2 = &reqs[1];
1428        let banner = run2
1429            .iter()
1430            .find_map(|m| match (m.role, m.content.as_slice()) {
1431                (Role::User, [ContentBlock::Text { text }])
1432                    if text.contains("replace all previously provided") =>
1433                {
1434                    Some(text.clone())
1435                }
1436                _ => None,
1437            })
1438            .expect("replace banner on edit");
1439        assert!(banner.contains("v2 rules"), "new content: {banner}");
1440        assert!(!banner.contains("v1 rules"), "not the old content");
1441        // Both the original (v1) and the replacement (v2) reminders are in the transcript.
1442        assert_eq!(reminder_count(run2), 2);
1443    }
1444
1445    #[tokio::test]
1446    async fn project_instructions_removal_banner_on_delete() {
1447        let dir = tempfile::tempdir().unwrap();
1448        let root = std::fs::canonicalize(dir.path()).unwrap();
1449        std::fs::create_dir(root.join(".git")).unwrap();
1450        let agents = root.join("AGENTS.md");
1451        std::fs::write(&agents, "rules").unwrap();
1452
1453        let (mut s, requests) = capturing_with_cfg(
1454            vec![Ok(text_turn("ok1")), Ok(text_turn("ok2"))],
1455            instr_config(root),
1456        );
1457        s.run_text("q1").await;
1458        std::fs::remove_file(&agents).unwrap(); // delete between runs
1459        s.run_text("q2").await;
1460
1461        let reqs = requests.lock().unwrap();
1462        assert!(
1463            reqs[1].iter().any(|m| matches!(
1464                (m.role, m.content.as_slice()),
1465                (Role::User, [ContentBlock::Text { text }]) if text.contains("no longer apply")
1466            )),
1467            "removal notice on delete"
1468        );
1469    }
1470
1471    #[tokio::test]
1472    async fn project_instructions_not_reinjected_when_unchanged() {
1473        let dir = tempfile::tempdir().unwrap();
1474        let root = std::fs::canonicalize(dir.path()).unwrap();
1475        std::fs::create_dir(root.join(".git")).unwrap();
1476        std::fs::write(root.join("AGENTS.md"), "stable").unwrap();
1477
1478        let (mut s, requests) = capturing_with_cfg(
1479            vec![
1480                Ok(text_turn("ok1")),
1481                Ok(text_turn("ok2")),
1482                Ok(text_turn("ok3")),
1483            ],
1484            instr_config(root),
1485        );
1486        s.run_text("q1").await;
1487        s.run_text("q2").await;
1488        s.run_text("q3").await;
1489        // Exactly one injection across three unchanged turns.
1490        assert_eq!(reminder_count(&requests.lock().unwrap()[2]), 1);
1491    }
1492
1493    #[tokio::test]
1494    async fn init_emitted_once_across_runs_with_one_result_each() {
1495        let (mut s, events) = session_with(
1496            vec![Ok(text_turn("one")), Ok(text_turn("two"))],
1497            Registry::new(),
1498            config(),
1499        );
1500        let _ = s.run_text("q1").await;
1501        let _ = s.run_text("q2").await;
1502        let evs = dump(&events);
1503        let inits = evs
1504            .iter()
1505            .filter(|e| matches!(e, Event::Init { .. }))
1506            .count();
1507        let results = evs
1508            .iter()
1509            .filter(|e| matches!(e, Event::Result { .. }))
1510            .count();
1511        assert_eq!(inits, 1, "Init is once per session, not per run");
1512        assert_eq!(results, 2, "one Result per run");
1513        assert!(
1514            matches!(evs.first(), Some(Event::Init { .. })),
1515            "Init still opens the stream"
1516        );
1517    }
1518
1519    #[tokio::test]
1520    async fn report_counts_are_per_run_not_cumulative() {
1521        // Run 1: tool turn + text (2 turns, 1 tool call, 10/5 tokens).
1522        // Run 2: text only (1 turn, 0 tool calls, 20/7 tokens).
1523        let mut t1 = tool_turn("c1", "echo");
1524        t1.usage = Usage {
1525            input_tokens: 10,
1526            output_tokens: 5,
1527            ..Usage::default()
1528        };
1529        let t2 = text_turn("done one");
1530        let mut t3 = text_turn("done two");
1531        t3.usage = Usage {
1532            input_tokens: 20,
1533            output_tokens: 7,
1534            ..Usage::default()
1535        };
1536        let (mut s, _e) = session_with(vec![Ok(t1), Ok(t2), Ok(t3)], echo_registry(), config());
1537        let r1 = s.run_text("q1").await;
1538        let r2 = s.run_text("q2").await;
1539        assert_eq!(r1.turns, 2);
1540        assert_eq!(r1.tool_calls.len(), 1);
1541        assert_eq!(r2.turns, 1, "run 2 counts its own turns only");
1542        assert!(r2.tool_calls.is_empty());
1543        assert_eq!(r2.usage.input_tokens, 20, "usage is per-run");
1544        assert_eq!(r2.usage.output_tokens, 7);
1545    }
1546
1547    /// Golden: a two-run stream (`Init M+ Result M+ Result`) reconstructs the
1548    /// full cross-run conversation (ADR-0014 amendment 2026-07-21).
1549    #[tokio::test]
1550    async fn two_run_stream_reconstructs_the_full_conversation() {
1551        let (mut s, events) = session_with(
1552            vec![
1553                Ok(tool_turn("c1", "echo")),
1554                Ok(text_turn("done one")),
1555                Ok(text_turn("done two")),
1556            ],
1557            echo_registry(),
1558            config(),
1559        );
1560        let _ = s.run_text("q1").await;
1561        let _ = s.run_text("q2").await;
1562        let rebuilt: Conversation = reconstruct_conversation(&dump(&events));
1563        // Run 1: user, assistant(tool_use), user(tool_result), assistant(text);
1564        // run 2: user, assistant(text).
1565        let roles: Vec<Role> = rebuilt.messages.iter().map(|m| m.role).collect();
1566        assert_eq!(
1567            roles,
1568            vec![
1569                Role::User,
1570                Role::Assistant,
1571                Role::User,
1572                Role::Assistant,
1573                Role::User,
1574                Role::Assistant,
1575            ]
1576        );
1577        // And the reconstruction matches the session's own history exactly.
1578        assert_eq!(rebuilt.messages.as_slice(), s.history());
1579    }
1580
1581    /// Continuing after a `ModelError` run is allowed unconditionally
1582    /// (ADR-0016 Resolution): the history simply didn't advance.
1583    #[tokio::test]
1584    async fn continues_after_model_error() {
1585        let (mut s, requests, _e) = capturing_session_with(
1586            vec![Err(ProviderError::ContextOverflow), Ok(text_turn("ok now"))],
1587            Registry::new(),
1588        );
1589        let r1 = s.run_text("q1").await;
1590        let r2 = s.run_text("q2").await;
1591        assert_eq!(r1.status, Status::ModelError);
1592        assert_eq!(r2.status, Status::Completed);
1593        // Run 2's request: q1's user message survived; no phantom assistant turn.
1594        let reqs = requests.lock().unwrap();
1595        let run2 = &reqs[1];
1596        assert_eq!(run2.len(), 2, "user q1 + user q2: {run2:?}");
1597        assert_eq!(user_text(&run2[0]), Some("q1"));
1598        assert_eq!(user_text(&run2[1]), Some("q2"));
1599    }
1600
1601    /// Continuing after a fatal tool `Error` run: the transcript was fully
1602    /// paired before the break, so the next sample sees a valid history.
1603    #[tokio::test]
1604    async fn continues_after_fatal_tool_error_with_valid_pairing() {
1605        let mut reg = Registry::new();
1606        reg.register("boom", Boom);
1607        let (mut s, requests, _e) = capturing_session_with(
1608            vec![Ok(tool_turn("c1", "boom")), Ok(text_turn("recovered"))],
1609            reg,
1610        );
1611        let r1 = s.run_text("q1").await;
1612        let r2 = s.run_text("q2").await;
1613        assert_eq!(r1.status, Status::Error);
1614        assert_eq!(r2.status, Status::Completed);
1615
1616        // Run 2's request replays the failed run intact: the boom tool_use is
1617        // answered by its (is_error) tool_result.
1618        let reqs = requests.lock().unwrap();
1619        let run2 = &reqs[1];
1620        assert_eq!(run2.len(), 4, "q1, assistant, tool_result, q2: {run2:?}");
1621        assert!(
1622            run2[1]
1623                .content
1624                .iter()
1625                .any(|b| matches!(b, ContentBlock::ToolUse { id, .. } if id == "c1"))
1626        );
1627        assert!(run2[2].content.iter().any(|b| matches!(
1628            b,
1629            ContentBlock::ToolResult { tool_use_id, is_error: true, .. } if tool_use_id == "c1"
1630        )));
1631        assert_eq!(user_text(&run2[3]), Some("q2"));
1632    }
1633
1634    #[tokio::test]
1635    async fn usage_is_summed_across_turns() {
1636        let mut first = tool_turn("c1", "echo");
1637        first.usage = Usage {
1638            input_tokens: 10,
1639            output_tokens: 5,
1640            ..Usage::default()
1641        };
1642        let mut second = text_turn("done");
1643        second.usage = Usage {
1644            input_tokens: 20,
1645            output_tokens: 7,
1646            ..Usage::default()
1647        };
1648        let (mut s, _e) = session_with(vec![Ok(first), Ok(second)], echo_registry(), config());
1649        let report = s.run_text("go").await;
1650        assert_eq!(report.usage.input_tokens, 30);
1651        assert_eq!(report.usage.output_tokens, 12);
1652    }
1653
1654    /// Switching the model swaps the provider and **announces** the change instead of
1655    /// rewriting the preamble.
1656    ///
1657    /// A pack's system prompt may name the model, and after a switch that line is
1658    /// stale — but rewriting the `System` message would desync the conversation from
1659    /// the trace, whose `Init` record already captured the original preamble. Appending
1660    /// is the same discipline project instructions and skills follow.
1661    #[tokio::test]
1662    async fn setting_the_model_announces_it_without_touching_the_preamble() {
1663        let preamble = vec![Message {
1664            role: Role::System,
1665            content: vec![ContentBlock::Text {
1666                text: "You are powered by the model old-1.".into(),
1667            }],
1668        }];
1669        let provider = Arc::new(MockProvider::with_results(vec![Ok(text_turn("ok"))]));
1670        let mut s = Session::new(
1671            provider.clone(),
1672            Registry::new(),
1673            preamble.clone(),
1674            config(),
1675            Box::new(NullSink),
1676        );
1677
1678        let notice = s.set_model(provider, "new-2");
1679        s.announce(notice);
1680
1681        assert_eq!(
1682            s.history()[0],
1683            preamble[0],
1684            "the preamble is untouched — the trace already recorded it"
1685        );
1686        let last = s.history().last().expect("announcement appended");
1687        assert_eq!(last.role, Role::User);
1688        let ContentBlock::Text { text } = &last.content[0] else {
1689            panic!("text block")
1690        };
1691        assert!(text.starts_with("<system-reminder>"), "{text}");
1692        assert!(text.contains("is now new-2"), "{text}");
1693        assert!(
1694            text.contains("out of date"),
1695            "corrects the stale line: {text}"
1696        );
1697    }
1698
1699    /// `context_usage` is the **final** turn's, not the sum.
1700    ///
1701    /// Every turn's request re-sends the whole conversation, so summing counts the same
1702    /// history once per turn — a number that only grows and says nothing about how full
1703    /// the context is. The last turn's request *is* the whole conversation.
1704    #[tokio::test]
1705    async fn context_usage_is_the_final_turn_not_the_sum() {
1706        let mut first = tool_turn("c1", "echo");
1707        first.usage = Usage {
1708            input_tokens: 10,
1709            output_tokens: 5,
1710            ..Usage::default()
1711        };
1712        let mut second = text_turn("done");
1713        second.usage = Usage {
1714            input_tokens: 20,
1715            output_tokens: 7,
1716            cache_read_tokens: Some(4),
1717            cache_creation_tokens: Some(3),
1718            ..Usage::default()
1719        };
1720        let (mut s, _e) = session_with(vec![Ok(first), Ok(second)], echo_registry(), config());
1721        let report = s.run_text("go").await;
1722
1723        assert_eq!(report.context_usage.input_tokens, 20, "the last turn only");
1724        assert_eq!(report.context_usage.output_tokens, 7);
1725        assert_eq!(
1726            report.context_usage.context_tokens(),
1727            20 + 4 + 3 + 7,
1728            "both cache counters are prompt tokens"
1729        );
1730        assert_eq!(report.usage.input_tokens, 30, "the sum is still the sum");
1731    }
1732}