Skip to main content

kiss_coding/
session_runner.rs

1//! AgentSession: the high-level facade every mode drives. Owns the session
2//! manager, tool set, queues, retry, auto-compaction, and cost accounting.
3
4use crate::compaction::{
5    self, estimate_context_tokens, extract_file_ops, file_ops_details, plan_compaction,
6    should_compact,
7};
8use crate::iterative::IterativeRuntime;
9use crate::session::manager::SessionManager;
10use crate::settings::{CacheWarmingMode, CompactionMode, QueueMode, ReasoningEffortMode, Settings};
11use crate::subagents::{ForkTurns, SUBAGENT_SYSTEM_PROMPT, SubagentRuntime, fork_messages};
12use crate::workflows::{WorkflowApprover, WorkflowRuntime};
13use anyhow::Context as _;
14use kiss_agent::{
15    AgentContext, AgentEvent, AgentLoopConfig, AgentMessage, DynTool, EventSink, StreamFn,
16    TurnUpdate,
17};
18use kiss_ai::{Model, Registry, StopReason, ThinkingLevel, Usage};
19use std::collections::VecDeque;
20use std::sync::{Arc, Mutex, OnceLock};
21use tokio_util::sync::CancellationToken;
22
23/// Harness-level events layered over the loop's AgentEvents.
24#[derive(Debug, Clone)]
25pub enum SessionEvent {
26    Agent(Box<AgentEvent>),
27    QueueUpdate {
28        steering: Vec<String>,
29        follow_up: Vec<String>,
30    },
31    CompactionStart {
32        auto: bool,
33    },
34    CompactionEnd {
35        summary: String,
36        tokens_before: u64,
37        error: Option<String>,
38    },
39    Retry {
40        attempt: u32,
41        max: u32,
42        delay_ms: u64,
43        error: String,
44    },
45    ModelChanged {
46        provider: String,
47        model_id: String,
48    },
49    ReasoningEffortChanged {
50        level: ThinkingLevel,
51        generations: u8,
52    },
53    ReasoningEffortFallback {
54        level: ThinkingLevel,
55        reason: String,
56    },
57    /// A dynamic workflow changed. The terminal redraws from the shared
58    /// snapshot. The event carries only the cheap version marker.
59    Workflow {
60        run: crate::workflows::RunId,
61        version: u64,
62    },
63    /// The verified result of a workflow tool call. Frontends show this after
64    /// the assistant answer so model text cannot replace the real run state.
65    WorkflowOutcome {
66        run: Option<crate::workflows::RunId>,
67        name: String,
68        status: WorkflowTurnStatus,
69    },
70    /// A loop or autoresearch job changed. The terminal reads its snapshot.
71    Iterative {
72        job: crate::iterative::JobId,
73        version: u64,
74    },
75}
76
77pub type SessionEventSink = Arc<dyn Fn(SessionEvent) + Send + Sync>;
78
79/// Tools and instructions that belong to one submitted user turn.
80///
81/// Workflow mode is explicit and local to the prompt invocation. This keeps
82/// one prompt from enabling or disabling workflow tools for another prompt.
83#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
84pub enum PromptMode {
85    #[default]
86    Ordinary,
87    Workflow,
88}
89
90/// The result that the workflow runtime, rather than the model, observed.
91#[derive(Debug, Clone, Copy, PartialEq, Eq)]
92pub enum WorkflowTurnStatus {
93    Cancelled,
94    Completed,
95    Failed,
96    Stopped,
97}
98
99#[derive(Clone)]
100struct QueuedPrompt {
101    message: AgentMessage,
102    mode: PromptMode,
103}
104
105pub struct TreeNavigationOutcome {
106    pub editor_text: Option<String>,
107    pub summarized: bool,
108    pub cancelled: bool,
109}
110
111#[derive(Debug, Clone)]
112pub struct EphemeralResponse {
113    pub text: String,
114    pub usage: Usage,
115}
116
117const SESSION_TITLE_MAX_CHARS: usize = 36;
118const SESSION_TITLE_PROMPT_MAX_BYTES: usize = 960;
119
120fn bounded_session_title_prompt(prompt: &str) -> &str {
121    let prompt = prompt.trim();
122    let end = prompt.floor_char_boundary(SESSION_TITLE_PROMPT_MAX_BYTES);
123    &prompt[..end]
124}
125
126fn normalize_session_title(text: &str) -> Option<String> {
127    let line = text.lines().find(|line| !line.trim().is_empty())?;
128    let cleaned: String = line
129        .chars()
130        .filter(|character| !character.is_control())
131        .collect();
132    let normalized = cleaned
133        .trim()
134        .trim_matches(|character| matches!(character, '"' | '\'' | '`' | '“' | '”' | '‘' | '’'))
135        .split_whitespace()
136        .collect::<Vec<_>>()
137        .join(" ");
138    let title: String = normalized.chars().take(SESSION_TITLE_MAX_CHARS).collect();
139    let title = title
140        .trim_end_matches(['.', '?', '!', ',', ':', ';'])
141        .trim();
142    (!title.is_empty()).then(|| title.to_string())
143}
144
145pub struct AgentSession {
146    pub manager: Mutex<SessionManager>,
147    pub registry: Arc<Registry>,
148    base_tools: Mutex<Vec<DynTool>>,
149    session_tools: Mutex<Vec<DynTool>>,
150    tools: Mutex<Vec<DynTool>>,
151    settings: Mutex<Settings>,
152    system_prompt: Mutex<String>,
153    model: Mutex<Model>,
154    thinking: Mutex<ThinkingLevel>,
155    fast_mode: Mutex<bool>,
156    steering: Arc<Mutex<VecDeque<QueuedPrompt>>>,
157    follow_up: Arc<Mutex<VecDeque<QueuedPrompt>>>,
158    cancel: Mutex<CancellationToken>,
159    cache_warm_cancel: Mutex<CancellationToken>,
160    running: Mutex<bool>,
161    totals: Mutex<Usage>,
162    context_usage_cache: Mutex<Option<(u64, u64)>>,
163    context_file: Mutex<Option<crate::context_file::ContextFile>>,
164    api_key_override: Option<(String, String)>,
165    sink: SessionEventSink,
166    subagents_allowed: bool,
167    subagents: OnceLock<Arc<SubagentRuntime>>,
168    workflows: OnceLock<Arc<WorkflowRuntime>>,
169    iterative: OnceLock<Arc<IterativeRuntime>>,
170    workflow_approver: Mutex<Option<WorkflowApprover>>,
171    /// Optional replacement for the provider streaming function. Embedders and
172    /// tests install one to run the whole loop against a scripted fake model.
173    stream_fn: Mutex<Option<StreamFn>>,
174}
175
176struct ReasoningRunState {
177    lease: crate::jev::ReasoningLease,
178    provider: String,
179    model_id: String,
180    saved_effort: ThinkingLevel,
181}
182
183impl ReasoningRunState {
184    fn checkpoint(
185        &mut self,
186        model: &Model,
187        saved_effort: ThinkingLevel,
188        current_effort: ThinkingLevel,
189        dynamic_enabled: bool,
190    ) -> (bool, Option<ThinkingLevel>) {
191        let model_changed = self.provider != model.provider || self.model_id != model.id;
192        if model_changed || self.saved_effort != saved_effort || !dynamic_enabled {
193            self.lease.clear();
194        }
195        self.provider.clone_from(&model.provider);
196        self.model_id.clone_from(&model.id);
197        self.saved_effort = saved_effort;
198        if self
199            .lease
200            .current()
201            .is_some_and(|level| level != current_effort)
202        {
203            self.lease.clear();
204        }
205        (model_changed, self.lease.current())
206    }
207}
208
209impl AgentSession {
210    pub(crate) fn emit_iterative(&self, job: crate::iterative::JobId, version: u64) {
211        (self.sink)(SessionEvent::Iterative { job, version });
212    }
213
214    pub(crate) fn emit_workflow(&self, run: crate::workflows::RunId, version: u64) {
215        (self.sink)(SessionEvent::Workflow { run, version });
216    }
217
218    pub(crate) fn emit_workflow_outcome(
219        &self,
220        run: Option<crate::workflows::RunId>,
221        name: String,
222        status: WorkflowTurnStatus,
223    ) {
224        (self.sink)(SessionEvent::WorkflowOutcome { run, name, status });
225    }
226
227    #[allow(clippy::too_many_arguments)]
228    pub fn new(
229        manager: SessionManager,
230        tools: Vec<DynTool>,
231        registry: impl Into<Arc<Registry>>,
232        settings: Settings,
233        system_prompt: String,
234        model: Model,
235        thinking: ThinkingLevel,
236        api_key_override: Option<(String, String)>,
237        sink: SessionEventSink,
238    ) -> Arc<Self> {
239        Self::new_with_subagents_allowed(
240            manager,
241            tools,
242            registry,
243            settings,
244            system_prompt,
245            model,
246            thinking,
247            api_key_override,
248            sink,
249            true,
250        )
251    }
252
253    #[allow(clippy::too_many_arguments)]
254    pub fn new_with_subagents_allowed(
255        manager: SessionManager,
256        tools: Vec<DynTool>,
257        registry: impl Into<Arc<Registry>>,
258        settings: Settings,
259        system_prompt: String,
260        model: Model,
261        thinking: ThinkingLevel,
262        api_key_override: Option<(String, String)>,
263        sink: SessionEventSink,
264        subagents_allowed: bool,
265    ) -> Arc<Self> {
266        let totals = manager.usage_totals();
267        let session = Arc::new(AgentSession {
268            manager: Mutex::new(manager),
269            registry: registry.into(),
270            base_tools: Mutex::new(tools.clone()),
271            session_tools: Mutex::new(Vec::new()),
272            tools: Mutex::new(tools),
273            settings: Mutex::new(settings),
274            system_prompt: Mutex::new(system_prompt),
275            model: Mutex::new(model),
276            thinking: Mutex::new(thinking),
277            fast_mode: Mutex::new(false),
278            steering: Default::default(),
279            follow_up: Default::default(),
280            cancel: Mutex::new(CancellationToken::new()),
281            cache_warm_cancel: Mutex::new(CancellationToken::new()),
282            running: Mutex::new(false),
283            totals: Mutex::new(totals),
284            context_usage_cache: Default::default(),
285            context_file: Mutex::new(None),
286            api_key_override,
287            sink,
288            subagents_allowed,
289            subagents: OnceLock::new(),
290            workflows: OnceLock::new(),
291            iterative: OnceLock::new(),
292            workflow_approver: Mutex::new(None),
293            stream_fn: Mutex::new(None),
294        });
295        if subagents_allowed {
296            let runtime = SubagentRuntime::new(Arc::downgrade(&session));
297            assert!(session.subagents.set(runtime).is_ok());
298            let workflows = WorkflowRuntime::new(Arc::downgrade(&session));
299            assert!(session.workflows.set(workflows).is_ok());
300            let iterative = IterativeRuntime::new(Arc::downgrade(&session));
301            assert!(session.iterative.set(iterative).is_ok());
302        }
303        session.rebuild_tools();
304        session
305    }
306
307    /// Replace the provider streaming function used by every later request.
308    ///
309    /// Pass `None` to restore the real provider. Child sessions created after
310    /// this call inherit the override so a scripted fake model also serves
311    /// subagents and compaction.
312    pub fn set_stream_fn(&self, stream_fn: Option<StreamFn>) {
313        *self.stream_fn.lock().unwrap() = stream_fn;
314    }
315
316    pub fn model(&self) -> Model {
317        self.model.lock().unwrap().clone()
318    }
319
320    pub fn thinking_level(&self) -> ThinkingLevel {
321        *self.thinking.lock().unwrap()
322    }
323
324    pub fn fast_mode(&self) -> bool {
325        *self.fast_mode.lock().unwrap()
326    }
327
328    pub fn set_fast_mode(&self, enabled: bool) {
329        *self.fast_mode.lock().unwrap() = enabled;
330    }
331
332    pub fn totals(&self) -> Usage {
333        *self.totals.lock().unwrap()
334    }
335
336    pub(crate) fn record_subagent_usage(&self, usage: Usage) {
337        self.totals.lock().unwrap().add(&usage);
338    }
339
340    pub fn is_running(&self) -> bool {
341        *self.running.lock().unwrap()
342    }
343
344    pub fn settings(&self) -> Settings {
345        self.settings.lock().unwrap().clone()
346    }
347
348    pub fn update_settings(&self, settings: Settings) {
349        *self.context_usage_cache.lock().unwrap() = None;
350        let was_enabled = self.subagents_enabled();
351        let workflows_were_enabled = self.workflows_enabled();
352        let cache_warming = settings.cache_warming;
353        *self.settings.lock().unwrap() = settings;
354        if cache_warming == CacheWarmingMode::Off
355            || (cache_warming == CacheWarmingMode::Streaming && !self.is_running())
356        {
357            self.cache_warm_cancel.lock().unwrap().cancel();
358        }
359        let is_enabled = self.subagents_enabled();
360        if was_enabled && !is_enabled {
361            self.stop_child_work();
362        } else if workflows_were_enabled && !self.workflows_enabled() {
363            self.stop_workflows();
364        }
365        self.rebuild_tools();
366    }
367
368    /// Replace resources used by the next model request.
369    pub fn reload_runtime(&self, settings: Settings, system_prompt: String, tools: Vec<DynTool>) {
370        *self.context_usage_cache.lock().unwrap() = None;
371        let was_enabled = self.subagents_enabled();
372        let workflows_were_enabled = self.workflows_enabled();
373        let cache_warming = settings.cache_warming;
374        *self.settings.lock().unwrap() = settings;
375        *self.system_prompt.lock().unwrap() = system_prompt;
376        *self.base_tools.lock().unwrap() = tools;
377        if cache_warming == CacheWarmingMode::Off
378            || (cache_warming == CacheWarmingMode::Streaming && !self.is_running())
379        {
380            self.cache_warm_cancel.lock().unwrap().cancel();
381        }
382        let is_enabled = self.subagents_enabled();
383        if was_enabled && !is_enabled {
384            self.stop_child_work();
385        } else if workflows_were_enabled && !self.workflows_enabled() {
386            self.stop_workflows();
387        }
388        self.rebuild_tools();
389    }
390
391    /// Add or replace one tool for this process without changing saved tool configuration.
392    pub fn install_session_tool(&self, tool: DynTool) {
393        let mut tools = self.session_tools.lock().unwrap();
394        if let Some(existing) = tools
395            .iter_mut()
396            .find(|existing| existing.name() == tool.name())
397        {
398            *existing = tool;
399        } else {
400            tools.push(tool);
401        }
402        drop(tools);
403        self.rebuild_tools();
404    }
405
406    /// Remove one process-only tool. Return true when a tool was removed.
407    pub fn remove_session_tool(&self, name: &str) -> bool {
408        let mut tools = self.session_tools.lock().unwrap();
409        let previous_len = tools.len();
410        tools.retain(|tool| tool.name() != name);
411        let removed = tools.len() != previous_len;
412        drop(tools);
413        if removed {
414            self.rebuild_tools();
415        }
416        removed
417    }
418
419    /// Stop every child agent and workflow run this session started.
420    fn stop_child_work(&self) {
421        if let Some(runtime) = self.subagents.get() {
422            runtime.interrupt_all();
423        }
424        if let Some(runtime) = self.workflows.get() {
425            runtime.stop_all();
426        }
427        if let Some(runtime) = self.iterative.get() {
428            runtime.stop_all();
429        }
430    }
431
432    fn stop_workflows(&self) {
433        if let Some(runtime) = self.workflows.get() {
434            runtime.stop_all();
435        }
436    }
437
438    fn subagents_enabled(&self) -> bool {
439        self.subagents_allowed && self.settings.lock().unwrap().subagents.enabled
440    }
441
442    /// Dynamic workflows are built on child agents, so they need subagents on
443    /// as well as their own setting.
444    pub fn workflows_enabled(&self) -> bool {
445        self.subagents_enabled() && self.settings.lock().unwrap().workflows.enabled
446    }
447
448    /// Select workflow mode from one submitted user prompt.
449    pub fn prompt_mode_for(&self, text: &str) -> PromptMode {
450        if self.workflows_enabled()
451            && self.settings.lock().unwrap().workflows.keyword_trigger
452            && crate::workflows::workflow_trigger(text).is_some()
453        {
454            PromptMode::Workflow
455        } else {
456            PromptMode::Ordinary
457        }
458    }
459
460    pub fn workflows(&self) -> Option<Arc<WorkflowRuntime>> {
461        self.workflows.get().cloned()
462    }
463
464    pub fn iterative_jobs(&self) -> Option<Arc<IterativeRuntime>> {
465        self.iterative.get().cloned()
466    }
467
468    /// Install the callback that asks the user to approve a run.
469    pub fn set_workflow_approver(&self, approver: WorkflowApprover) {
470        *self.workflow_approver.lock().unwrap() = Some(approver);
471    }
472
473    pub(crate) fn workflow_approver(&self) -> Option<WorkflowApprover> {
474        self.workflow_approver.lock().unwrap().clone()
475    }
476
477    fn rebuild_tools(&self) {
478        let mut tools = self.base_tools.lock().unwrap().clone();
479        for session_tool in self.session_tools.lock().unwrap().iter().cloned() {
480            if let Some(existing) = tools
481                .iter_mut()
482                .find(|existing| existing.name() == session_tool.name())
483            {
484                *existing = session_tool;
485            } else {
486                tools.push(session_tool);
487            }
488        }
489        if self.subagents_enabled()
490            && let Some(runtime) = self.subagents.get()
491        {
492            tools.extend(runtime.control_tools());
493        }
494        *self.tools.lock().unwrap() = tools;
495    }
496
497    fn tools_for(&self, mode: PromptMode) -> Vec<DynTool> {
498        let mut tools = self.tools.lock().unwrap().clone();
499        if mode == PromptMode::Workflow
500            && self.workflows_enabled()
501            && let Some(runtime) = self.workflows.get()
502        {
503            tools.push(runtime.tool());
504        }
505        tools
506    }
507
508    pub fn available_tool_names(&self) -> Vec<String> {
509        self.tools
510            .lock()
511            .unwrap()
512            .iter()
513            .map(|tool| tool.name().to_string())
514            .collect()
515    }
516
517    #[cfg(test)]
518    fn available_tool_names_for(&self, mode: PromptMode) -> Vec<String> {
519        self.tools_for(mode)
520            .iter()
521            .map(|tool| tool.name().to_string())
522            .collect()
523    }
524
525    /// Switch the active session without appending synthetic history.
526    pub fn replace_manager(&self, manager: SessionManager) {
527        self.cache_warm_cancel.lock().unwrap().cancel();
528        *self.context_file.lock().unwrap() = None;
529        if let Some(runtime) = self.subagents.get() {
530            runtime.reset();
531        }
532        if let Some(runtime) = self.workflows.get() {
533            runtime.stop_all();
534        }
535        if let Some(runtime) = self.iterative.get() {
536            runtime.stop_all();
537        }
538        let totals = manager.usage_totals();
539        let context = manager.build_session_context();
540        if let Some((provider, model_id)) = context.model
541            && let Some((model, _)) = self.registry.resolve(&model_id, Some(&provider))
542        {
543            *self.model.lock().unwrap() = model;
544        }
545        if let Some(thinking) = context.thinking_level {
546            *self.thinking.lock().unwrap() = thinking;
547        }
548        *self.manager.lock().unwrap() = manager;
549        *self.context_usage_cache.lock().unwrap() = None;
550        *self.totals.lock().unwrap() = totals;
551    }
552
553    pub fn set_model(&self, model: Model) {
554        self.cache_warm_cancel.lock().unwrap().cancel();
555        {
556            let mut m = self.manager.lock().unwrap();
557            let _ = m.append_model_change(&model.provider, &model.id);
558        }
559        (self.sink)(SessionEvent::ModelChanged {
560            provider: model.provider.clone(),
561            model_id: model.id.clone(),
562        });
563        *self.model.lock().unwrap() = model;
564    }
565
566    pub fn set_thinking_level(&self, level: ThinkingLevel) {
567        let _ = self
568            .manager
569            .lock()
570            .unwrap()
571            .append_thinking_level_change(level);
572        *self.thinking.lock().unwrap() = level;
573    }
574
575    /// Move the active leaf in the current session tree, with an optional
576    /// summary of the branch that is no longer active.
577    pub async fn navigate_tree(
578        self: &Arc<Self>,
579        target_id: &str,
580        summarize: bool,
581        custom_instructions: Option<String>,
582        cancel: CancellationToken,
583    ) -> anyhow::Result<TreeNavigationOutcome> {
584        let (old_leaf, new_leaf, editor_text, abandoned_messages) = {
585            let manager = self.manager.lock().unwrap();
586            let target = manager
587                .get_entry(target_id)
588                .cloned()
589                .ok_or_else(|| anyhow::anyhow!("unknown tree entry {target_id}"))?;
590            let old_leaf = manager.leaf_id().map(str::to_string);
591            if old_leaf.as_deref() == Some(target_id) {
592                return Ok(TreeNavigationOutcome {
593                    editor_text: None,
594                    summarized: false,
595                    cancelled: false,
596                });
597            }
598
599            let (new_leaf, editor_text) = match &target {
600                crate::session::entry::SessionEntry::Message {
601                    message: AgentMessage::User(user),
602                    ..
603                } => (
604                    target.parent_id().map(str::to_string),
605                    Some(user.content.as_text()),
606                ),
607                _ => (Some(target_id.to_string()), None),
608            };
609
610            let target_ancestors: std::collections::HashSet<String> = new_leaf
611                .as_deref()
612                .map(|leaf| {
613                    manager
614                        .branch_entries(Some(leaf))
615                        .into_iter()
616                        .map(|entry| entry.id().to_string())
617                        .collect()
618                })
619                .unwrap_or_default();
620            let common_ancestor = old_leaf.as_deref().and_then(|leaf| {
621                manager
622                    .branch_entries(Some(leaf))
623                    .into_iter()
624                    .rev()
625                    .find(|entry| target_ancestors.contains(entry.id()))
626                    .map(|entry| entry.id().to_string())
627            });
628            let abandoned_messages = old_leaf
629                .as_deref()
630                .map(|leaf| manager.branch_messages_after(leaf, common_ancestor.as_deref()))
631                .unwrap_or_default();
632            (old_leaf, new_leaf, editor_text, abandoned_messages)
633        };
634
635        let summary = if summarize && !abandoned_messages.is_empty() {
636            let model = self.model();
637            let credential = self.resolve_credential(&model.provider).await;
638            let serialized = compaction::serialize_agent_messages(&abandoned_messages);
639            let result = compaction::generate_summary(
640                &model,
641                credential,
642                &serialized,
643                None,
644                custom_instructions.as_deref(),
645                false,
646                cancel.clone(),
647            )
648            .await?;
649            if cancel.is_cancelled() {
650                return Ok(TreeNavigationOutcome {
651                    editor_text: None,
652                    summarized: false,
653                    cancelled: true,
654                });
655            }
656            Some(result)
657        } else {
658            None
659        };
660
661        let mut manager = self.manager.lock().unwrap();
662        if manager.leaf_id() != old_leaf.as_deref() {
663            anyhow::bail!("the session tree changed during navigation");
664        }
665        let summarized = if let Some(summary) = summary {
666            let (read, modified) = extract_file_ops(&abandoned_messages);
667            if let Some(usage) = &summary.usage {
668                self.totals.lock().unwrap().add(usage);
669            }
670            manager.branch_with_summary(
671                new_leaf.as_deref(),
672                old_leaf.as_deref().unwrap_or(target_id),
673                summary.summary,
674                summary.usage,
675                Some(file_ops_details(&read, &modified)),
676            )?;
677            true
678        } else {
679            if let Some(new_leaf) = new_leaf.as_deref() {
680                manager.branch(new_leaf)?;
681            } else {
682                manager.reset_leaf();
683            }
684            false
685        };
686
687        Ok(TreeNavigationOutcome {
688            editor_text,
689            summarized,
690            cancelled: false,
691        })
692    }
693
694    pub fn queue_steering(&self, message: AgentMessage) {
695        self.queue_steering_with_mode(message, PromptMode::Ordinary);
696    }
697
698    pub fn queue_steering_with_mode(&self, message: AgentMessage, mode: PromptMode) {
699        self.steering
700            .lock()
701            .unwrap()
702            .push_back(QueuedPrompt { message, mode });
703        self.emit_queues();
704    }
705
706    pub fn queue_follow_up(&self, message: AgentMessage) {
707        self.queue_follow_up_with_mode(message, PromptMode::Ordinary);
708    }
709
710    pub fn queue_follow_up_with_mode(&self, message: AgentMessage, mode: PromptMode) {
711        self.follow_up
712            .lock()
713            .unwrap()
714            .push_back(QueuedPrompt { message, mode });
715        self.emit_queues();
716    }
717
718    /// Drain both queues back to the caller (Escape restores to editor).
719    pub fn reclaim_queued(&self) -> Vec<AgentMessage> {
720        let mut out: Vec<AgentMessage> = self
721            .steering
722            .lock()
723            .unwrap()
724            .drain(..)
725            .map(|prompt| prompt.message)
726            .collect();
727        out.extend(
728            self.follow_up
729                .lock()
730                .unwrap()
731                .drain(..)
732                .map(|prompt| prompt.message),
733        );
734        self.emit_queues();
735        out
736    }
737
738    pub fn abort(&self) {
739        self.cancel.lock().unwrap().cancel();
740    }
741
742    fn emit_queues(&self) {
743        let preview = |q: &VecDeque<QueuedPrompt>| {
744            q.iter()
745                .map(|prompt| match &prompt.message {
746                    AgentMessage::User(u) => u.content.as_text().chars().take(80).collect(),
747                    other => other.role().to_string(),
748                })
749                .collect::<Vec<String>>()
750        };
751        (self.sink)(SessionEvent::QueueUpdate {
752            steering: preview(&self.steering.lock().unwrap()),
753            follow_up: preview(&self.follow_up.lock().unwrap()),
754        });
755    }
756
757    async fn resolve_credential(&self, provider: &str) -> Option<kiss_ai::ResolvedCredential> {
758        if let Some((override_provider, key)) = &self.api_key_override
759            && override_provider == provider
760        {
761            return Some(kiss_ai::ResolvedCredential::api_key(key));
762        }
763        let credential_provider = self.registry.credential_provider(provider);
764        kiss_ai::auth::resolve_credential_async(credential_provider, &self.registry.declared_keys)
765            .await
766            .ok()
767            .flatten()
768    }
769
770    fn loop_config(
771        &self,
772        session_arc: &Arc<Self>,
773        active_prompt_mode: Arc<Mutex<PromptMode>>,
774    ) -> AgentLoopConfig {
775        let mut config = AgentLoopConfig::new(self.model());
776        config.thinking_level = self.thinking_level();
777        config.fast_mode = self.fast_mode();
778        config.session_id = Some(self.manager.lock().unwrap().session_id().to_string());
779        let settings = self.settings();
780        config.transport = settings.transport;
781        let registry = self.registry.clone();
782        let declared = self.registry.declared_keys.clone();
783        let api_key_override = self.api_key_override.clone();
784        config.get_credential = Some(Arc::new(move |provider| {
785            let registry = registry.clone();
786            let declared = declared.clone();
787            let api_key_override = api_key_override.clone();
788            Box::pin(async move {
789                if let Some((override_provider, key)) = api_key_override
790                    && override_provider == provider
791                {
792                    return Some(kiss_ai::ResolvedCredential::api_key(key));
793                }
794                let credential_provider = registry.credential_provider(&provider).to_string();
795                kiss_ai::auth::resolve_credential_async(&credential_provider, &declared)
796                    .await
797                    .ok()
798                    .flatten()
799            })
800        }));
801
802        let steering = self.steering.clone();
803        let steering_for_mode = self.steering.clone();
804        let steering_mode = settings.steering_mode;
805        let session_for_queues = session_arc.clone();
806        config.get_steering_messages = Some(Arc::new(move || {
807            let drained = drain_queue(&steering, steering_mode);
808            session_for_queues.emit_queues();
809            Box::pin(async move { drained })
810        }));
811        let follow_up = self.follow_up.clone();
812        let follow_up_for_mode = self.follow_up.clone();
813        let follow_up_mode = settings.follow_up_mode;
814        let session_for_queues = session_arc.clone();
815        config.get_follow_up_messages = Some(Arc::new(move || {
816            let drained = drain_queue(&follow_up, follow_up_mode);
817            session_for_queues.emit_queues();
818            Box::pin(async move { drained })
819        }));
820        let reasoning_run = Arc::new(Mutex::new(ReasoningRunState {
821            lease: crate::jev::ReasoningLease::default(),
822            provider: config.model.provider.clone(),
823            model_id: config.model.id.clone(),
824            saved_effort: config.thinking_level,
825        }));
826        let session_for_reasoning = session_arc.clone();
827        let reasoning_for_generation = reasoning_run.clone();
828        let prompt_mode_for_reasoning = active_prompt_mode.clone();
829        config.prepare_generation = Some(Arc::new(move |current_effort| {
830            let session = session_for_reasoning.clone();
831            let reasoning_run = reasoning_for_generation.clone();
832            let active_prompt_mode = prompt_mode_for_reasoning.clone();
833            Box::pin(async move {
834                let settings = session.settings();
835                let context_file_changed = session.sync_context_file();
836                let model = session.model();
837                let saved_effort = session.thinking_level();
838                let dynamic_enabled =
839                    settings.reasoning_effort.mode == ReasoningEffortMode::Jev && model.reasoning;
840                let (model_changed, leased) = {
841                    let mut run = reasoning_run.lock().unwrap();
842                    run.checkpoint(&model, saved_effort, current_effort, dynamic_enabled)
843                };
844                let mut update = TurnUpdate {
845                    context: (model_changed
846                        || context_file_changed
847                        || settings.experimental_context_file)
848                        .then(|| session.build_context_for(*active_prompt_mode.lock().unwrap())),
849                    model: model_changed.then(|| model.clone()),
850                    thinking_level: Some(saved_effort),
851                    fast_mode: Some(session.fast_mode()),
852                    ..Default::default()
853                };
854                if settings.experimental_context_file
855                    && let Some(context) = &mut update.context
856                {
857                    // A retry must not put the failed answer back in the
858                    // provider request when the file refreshes the context.
859                    while matches!(context.messages.last(), Some(AgentMessage::Assistant(assistant)) if assistant.stop_reason == StopReason::Error)
860                    {
861                        context.messages.pop();
862                    }
863                }
864                if !dynamic_enabled {
865                    return Some(update);
866                }
867                if let Some(level) = leased {
868                    update.thinking_level = Some(level);
869                    return Some(update);
870                }
871
872                let selected = match kiss_ai::auth::resolve_api_key_async(
873                    "typesafe",
874                    &session.registry.declared_keys,
875                )
876                .await
877                {
878                    Ok(Some(api_key)) => {
879                        let messages = session
880                            .manager
881                            .lock()
882                            .unwrap()
883                            .build_session_context()
884                            .messages;
885                        let cancel = session.cancel.lock().unwrap().clone();
886                        crate::jev::select_reasoning(&messages, &[], &model, &api_key, cancel)
887                            .await
888                            .map_err(|error| error.to_string())
889                    }
890                    Ok(None) => Err(
891                        "TypeSafe credentials are unavailable; run /login typesafe or set TYPESAFE_API_KEY"
892                            .into(),
893                    ),
894                    Err(error) => Err(format!(
895                        "TypeSafe credentials could not be read: {error}; run /login typesafe or set TYPESAFE_API_KEY"
896                    )),
897                };
898                let latest_model = session.model();
899                let latest_effort = session.thinking_level();
900                if latest_model.provider != model.provider
901                    || latest_model.id != model.id
902                    || latest_effort != saved_effort
903                    || session.settings().reasoning_effort.mode != ReasoningEffortMode::Jev
904                {
905                    reasoning_run.lock().unwrap().lease.clear();
906                    update.context =
907                        Some(session.build_context_for(*active_prompt_mode.lock().unwrap()));
908                    update.model = Some(latest_model);
909                    update.thinking_level = Some(latest_effort);
910                    return Some(update);
911                }
912                match selected {
913                    Ok(selection) => {
914                        update.thinking_level = Some(selection.level);
915                        if selection.level != current_effort {
916                            let sink = session.sink.clone();
917                            let level = selection.level;
918                            let generations = selection.generations;
919                            update.on_applied = Some(Box::new(move |_, applied| {
920                                if applied == level {
921                                    sink(SessionEvent::ReasoningEffortChanged {
922                                        level,
923                                        generations,
924                                    });
925                                }
926                            }));
927                        }
928                        reasoning_run.lock().unwrap().lease.install(&selection);
929                    }
930                    Err(reason) => {
931                        let sink = session.sink.clone();
932                        let reason = reason.chars().take(300).collect();
933                        update.on_applied = Some(Box::new(move |_, applied| {
934                            sink(SessionEvent::ReasoningEffortFallback {
935                                level: applied,
936                                reason,
937                            });
938                        }));
939                    }
940                }
941                Some(update)
942            })
943        }));
944        let session_for_compaction = session_arc.clone();
945        config.prepare_next_turn = Some(Arc::new(move |turn| {
946            let has_tool_results = !turn.tool_results.is_empty();
947            let tool_failed = turn.tool_results.iter().any(|result| result.is_error);
948            let will_continue = turn.will_continue;
949            let session = session_for_compaction.clone();
950            let active_prompt_mode = active_prompt_mode.clone();
951            let steering_for_mode = steering_for_mode.clone();
952            let follow_up_for_mode = follow_up_for_mode.clone();
953            let reasoning_run = reasoning_run.clone();
954            Box::pin(async move {
955                let has_queued_user_input = !steering_for_mode.lock().unwrap().is_empty()
956                    || (!will_continue && !follow_up_for_mode.lock().unwrap().is_empty());
957                let queued_mode = queued_mode(&steering_for_mode, steering_mode)
958                    .or_else(|| queued_mode(&follow_up_for_mode, follow_up_mode));
959                let mode_changed = {
960                    let mut active = active_prompt_mode.lock().unwrap();
961                    let before = *active;
962                    if let Some(queued_mode) = queued_mode {
963                        *active = queued_mode;
964                    } else if *active == PromptMode::Workflow && !has_tool_results {
965                        // A workflow turn remains active across its tool call.
966                        // Its final assistant response has no tool results, so
967                        // the next queued or follow-up turn is ordinary again.
968                        *active = PromptMode::Ordinary;
969                    }
970                    before != *active
971                };
972
973                let context_file_changed = session.sync_context_file();
974                let mut context_changed = mode_changed || context_file_changed;
975                if has_tool_results {
976                    let settings = session.settings();
977                    let cancel = session.cancel.lock().unwrap().clone();
978                    let model = session.model();
979                    let revision_before = {
980                        let manager = session.manager.lock().unwrap();
981                        let context = manager.build_session_context();
982                        if !auto_compaction_needed(
983                            &settings,
984                            &context.messages,
985                            &model,
986                            cancel.is_cancelled(),
987                        ) {
988                            None
989                        } else {
990                            Some(manager.context_revision())
991                        }
992                    };
993
994                    if let Some(revision_before) = revision_before {
995                        session.compact(None, true).await;
996                        let revision_after = session.manager.lock().unwrap().context_revision();
997                        context_changed |= revision_after != revision_before;
998                    }
999                }
1000
1001                {
1002                    let mut run = reasoning_run.lock().unwrap();
1003                    run.lease.consume_generation();
1004                    if tool_failed || has_queued_user_input {
1005                        run.lease.clear();
1006                    }
1007                }
1008
1009                Some(TurnUpdate {
1010                    context: context_changed
1011                        .then(|| session.build_context_for(*active_prompt_mode.lock().unwrap())),
1012                    fast_mode: Some(session.fast_mode()),
1013                    ..Default::default()
1014                })
1015            })
1016        }));
1017        if let Some(stream_fn) = self.stream_fn.lock().unwrap().clone() {
1018            config.stream_fn = stream_fn;
1019        }
1020        config
1021    }
1022
1023    fn sync_context_file(&self) -> bool {
1024        let mut file = self.context_file.lock().unwrap();
1025        if !self.settings().experimental_context_file {
1026            return file.take().is_some();
1027        }
1028        let mut manager = self.manager.lock().unwrap();
1029        let result = (|| {
1030            if file.is_none() {
1031                *file = Some(crate::context_file::ContextFile::new().context(
1032                    "Cannot create a context file. A writable temporary directory is required. Set TMPDIR to a writable directory and restart KISS",
1033                )?);
1034            }
1035            file.as_mut().unwrap().sync(&mut manager)
1036        })();
1037        match result {
1038            Ok(changed) => {
1039                if changed {
1040                    *self.context_usage_cache.lock().unwrap() = None;
1041                    self.cache_warm_cancel.lock().unwrap().cancel();
1042                }
1043                changed
1044            }
1045            Err(error) => {
1046                let message = AgentMessage::Custom(kiss_agent::CustomMessage {
1047                    custom_type: "context_file_error".into(),
1048                    content: kiss_ai::UserContent::Text(format!("{error:#}")),
1049                    display: true,
1050                    details: None,
1051                    timestamp: kiss_ai::now_ms(),
1052                });
1053                if manager.append_message(message.clone()).is_err() {
1054                    return false;
1055                }
1056                drop(manager);
1057                drop(file);
1058                (self.sink)(SessionEvent::Agent(Box::new(AgentEvent::MessageEnd {
1059                    message,
1060                })));
1061                *self.context_usage_cache.lock().unwrap() = None;
1062                true
1063            }
1064        }
1065    }
1066
1067    fn build_context(&self) -> AgentContext {
1068        self.build_context_for(PromptMode::Ordinary)
1069    }
1070
1071    fn build_context_for(&self, prompt_mode: PromptMode) -> AgentContext {
1072        let model = self.model();
1073        let manager = self.manager.lock().unwrap();
1074        let (openai_responses_input, messages) = if !self.settings().experimental_context_file
1075            && kiss_ai::api::openai_compaction::supports_remote_compaction(&model)
1076            && let Some(remote) = manager.build_openai_compaction_context(&model)
1077        {
1078            (Some(remote.replacement_history), remote.messages)
1079        } else {
1080            (None, manager.build_session_context().messages)
1081        };
1082        drop(manager);
1083        let mut system_prompt = self.system_prompt.lock().unwrap().clone();
1084        if self.settings().experimental_context_file
1085            && let Some(file) = self.context_file.lock().unwrap().as_ref()
1086        {
1087            system_prompt.push_str(&format!(
1088                "\n\nExperimental context file: {}\n\
1089                 This JSON array contains your live KISS conversation, without system instructions. \
1090                 Use your normal file tools to edit it. Changes apply after the tool batch and before the next model request. \
1091                 New user messages, your current answer, and tool results are added automatically. \
1092                 Keep important task instructions, exact facts, progress, and next steps. Replace stale output with useful notes. \
1093                 Keep complete assistant tool-call/result groups, or replace the whole group with a user note. \
1094                 A user note has the form {{\"role\":\"user\",\"content\":\"Notes here\",\"timestamp\":0}}. \
1095                 Preserve fields of messages you keep, including image and reasoning data. \
1096                 The file must be a valid JSON array at most 16 MiB. Invalid edits leave the live conversation unchanged; repair the file after an error. \
1097                 Read only the parts you need, because the conversation is already in your context. \
1098                 Your model context window is {} tokens. Manage the file before it fills; automatic compaction remains an emergency fallback. \
1099                 Batch edits: changing early text can require the provider to process all later text again. \
1100                 The full session record remains separate from this editable file.",
1101                file.path().display(), model.context_window
1102            ));
1103        }
1104        if self.subagents_enabled() {
1105            system_prompt.push_str("\n\n");
1106            system_prompt.push_str(SUBAGENT_SYSTEM_PROMPT);
1107        }
1108        if prompt_mode == PromptMode::Workflow
1109            && self.workflows_enabled()
1110            && let Some(runtime) = self.workflows.get()
1111        {
1112            let limits = runtime.limits();
1113            let size = self.settings.lock().unwrap().workflows.size;
1114            system_prompt.push_str("\n\n");
1115            system_prompt.push_str(&crate::workflows::authoring_prompt(
1116                size,
1117                limits.max_agents,
1118                limits.max_fanout,
1119            ));
1120        }
1121        AgentContext {
1122            system_prompt,
1123            openai_responses_input,
1124            messages,
1125            tools: self.tools_for(prompt_mode),
1126        }
1127    }
1128
1129    pub(crate) fn create_subagent_session(
1130        self: &Arc<Self>,
1131        task_name: &str,
1132        canonical_path: &str,
1133        fork_turns: ForkTurns,
1134        model_pattern: Option<&str>,
1135        reasoning_effort: Option<&str>,
1136    ) -> anyhow::Result<Arc<Self>> {
1137        let (model, suggested_thinking) = match model_pattern {
1138            Some(pattern) => self
1139                .registry
1140                .resolve(pattern, None)
1141                .with_context(|| format!("no model matches subagent model '{pattern}'"))?,
1142            None => (self.model(), None),
1143        };
1144        let thinking = match reasoning_effort {
1145            Some(level) => ThinkingLevel::parse(level)
1146                .with_context(|| format!("unknown subagent reasoning_effort '{level}'"))?,
1147            None => suggested_thinking.unwrap_or_else(|| self.thinking_level()),
1148        };
1149        let (mut manager, parent_messages, parent_id) = {
1150            let parent = self.manager.lock().unwrap();
1151            (
1152                parent.create_child()?,
1153                if fork_turns == ForkTurns::None {
1154                    Vec::new()
1155                } else {
1156                    parent.build_session_context().messages
1157                },
1158                parent.session_id().to_string(),
1159            )
1160        };
1161        for message in fork_messages(&parent_messages, fork_turns) {
1162            manager.append_message(message)?;
1163        }
1164        manager.append_custom(
1165            "subagent",
1166            Some(serde_json::json!({
1167                "taskName": task_name,
1168                "canonicalPath": canonical_path,
1169                "parentSessionId": parent_id,
1170            })),
1171        )?;
1172
1173        let mut settings = self.settings();
1174        settings.subagents.enabled = false;
1175        let mut system_prompt = self.system_prompt.lock().unwrap().clone();
1176        system_prompt.push_str(&format!(
1177            "\n\nYou are child agent {canonical_path}. Complete only the assigned task. Return a concise result to the parent agent."
1178        ));
1179
1180        let child = Self::new_with_subagents_allowed(
1181            manager,
1182            self.base_tools.lock().unwrap().clone(),
1183            self.registry.clone(),
1184            settings,
1185            system_prompt,
1186            model,
1187            thinking,
1188            self.api_key_override.clone(),
1189            Arc::new(|_| {}),
1190            false,
1191        );
1192        child.set_stream_fn(self.stream_fn.lock().unwrap().clone());
1193        Ok(child)
1194    }
1195
1196    async fn run_ephemeral(
1197        self: &Arc<Self>,
1198        system_prompt: String,
1199        prompt: String,
1200        tools: Vec<DynTool>,
1201        max_tokens: u64,
1202        cancel: CancellationToken,
1203    ) -> anyhow::Result<EphemeralResponse> {
1204        let mut config = self.loop_config(self, Arc::new(Mutex::new(PromptMode::Ordinary)));
1205        config.thinking_level = ThinkingLevel::Off;
1206        config.max_tokens = Some(max_tokens);
1207        config.session_id = Some(format!("ephemeral-{}", uuid::Uuid::new_v4()));
1208        config.get_steering_messages = None;
1209        config.get_follow_up_messages = None;
1210        config.prepare_next_turn = None;
1211        config.prepare_generation = None;
1212
1213        let context = AgentContext {
1214            system_prompt,
1215            openai_responses_input: None,
1216            messages: Vec::new(),
1217            tools,
1218        };
1219        let sink: EventSink = Arc::new(|_| {});
1220        let messages = kiss_agent::run_agent_loop(
1221            vec![AgentMessage::user(prompt)],
1222            context,
1223            config,
1224            cancel.clone(),
1225            sink,
1226        )
1227        .await;
1228        if cancel.is_cancelled() {
1229            anyhow::bail!("request cancelled");
1230        }
1231
1232        let mut usage = Usage::default();
1233        for message in &messages {
1234            if let AgentMessage::Assistant(assistant) = message {
1235                usage.add(&assistant.usage);
1236            }
1237        }
1238        let assistant = messages.iter().rev().find_map(|message| match message {
1239            AgentMessage::Assistant(assistant) => Some(assistant),
1240            _ => None,
1241        });
1242        let Some(assistant) = assistant else {
1243            anyhow::bail!("the provider returned no answer");
1244        };
1245        if assistant.stop_reason == StopReason::Error {
1246            anyhow::bail!(
1247                "{}",
1248                assistant
1249                    .error_message
1250                    .as_deref()
1251                    .unwrap_or("the provider request failed")
1252            );
1253        }
1254        let text = assistant.text();
1255        if text.trim().is_empty() {
1256            anyhow::bail!("the provider returned an empty answer");
1257        }
1258        self.totals.lock().unwrap().add(&usage);
1259        Ok(EphemeralResponse { text, usage })
1260    }
1261
1262    /// Answer a short side question without changing the active session.
1263    pub async fn answer_btw(
1264        self: &Arc<Self>,
1265        question: &str,
1266        cancel: CancellationToken,
1267    ) -> anyhow::Result<EphemeralResponse> {
1268        let messages = self
1269            .manager
1270            .lock()
1271            .unwrap()
1272            .build_session_context()
1273            .messages;
1274        let transcript = transcript_excerpt(&messages, 4, 4_000);
1275        let prompt = if transcript.is_empty() {
1276            format!("Side question:\n{question}")
1277        } else {
1278            format!("Recent session context:\n{transcript}\n\nSide question:\n{question}")
1279        };
1280        let read_tools = self
1281            .tools
1282            .lock()
1283            .unwrap()
1284            .iter()
1285            .filter(|tool| tool.name() == "read")
1286            .cloned()
1287            .collect();
1288        self.run_ephemeral(
1289            "Answer the side question from the supplied session context. This is a read-only request. Use the read tool only when a file is needed. Do not propose or perform edits. Give no more than 150 words or 600 characters. Use no more than five bullets. Return only the answer.".into(),
1290            prompt,
1291            read_tools,
1292            500,
1293            cancel,
1294        )
1295        .await
1296    }
1297
1298    /// Generate a short session title without changing conversation history.
1299    pub async fn generate_session_title(
1300        self: &Arc<Self>,
1301        prompt: &str,
1302        cancel: CancellationToken,
1303    ) -> anyhow::Result<String> {
1304        let prompt = bounded_session_title_prompt(prompt);
1305        if prompt.is_empty() {
1306            anyhow::bail!("the session prompt is empty");
1307        }
1308        let response = self
1309            .run_ephemeral(
1310                format!(
1311                    "Write a one-line title for this task, no more than {SESSION_TITLE_MAX_CHARS} characters. Aim for fewer than five words and start with an imperative verb. Keep ticket IDs and code terms unchanged. Match the user's language. Return the title in sentence case without quotes, Markdown, or ending punctuation."
1312                ),
1313                format!("User prompt:\n{prompt}"),
1314                Vec::new(),
1315                64,
1316                cancel,
1317            )
1318            .await?;
1319        normalize_session_title(&response.text)
1320            .context("the provider returned an invalid session title")
1321    }
1322
1323    /// Create a one-line recap without changing the active session.
1324    pub async fn generate_recap(
1325        self: &Arc<Self>,
1326        previous_recap: Option<&str>,
1327        cancel: CancellationToken,
1328    ) -> anyhow::Result<EphemeralResponse> {
1329        let messages = self
1330            .manager
1331            .lock()
1332            .unwrap()
1333            .build_session_context()
1334            .messages;
1335        let transcript = transcript_excerpt(&messages, 12, 12_000);
1336        if transcript.is_empty() {
1337            anyhow::bail!("the session has no conversation to recap");
1338        }
1339        let previous = previous_recap
1340            .filter(|recap| !recap.trim().is_empty())
1341            .map(|recap| format!("\n\nPrevious recap:\n{recap}"))
1342            .unwrap_or_default();
1343        self.run_ephemeral(
1344            "Summarize the supplied coding session in one plain-text line of at most 120 characters. State what was done and the next action when one is clear. Do not use a prefix, Markdown, or a newline. Return only the recap.".into(),
1345            format!("Session transcript:\n{transcript}{previous}"),
1346            Vec::new(),
1347            160,
1348            cancel,
1349        )
1350        .await
1351    }
1352
1353    /// Run one prompt to completion, including retry and auto-compaction.
1354    pub async fn prompt(self: &Arc<Self>, prompts: Vec<AgentMessage>) {
1355        self.prompt_with_mode(prompts, PromptMode::Ordinary).await;
1356    }
1357
1358    /// Run one prompt with tools and instructions selected for this turn.
1359    pub async fn prompt_with_mode(
1360        self: &Arc<Self>,
1361        prompts: Vec<AgentMessage>,
1362        prompt_mode: PromptMode,
1363    ) {
1364        self.cache_warm_cancel.lock().unwrap().cancel();
1365        {
1366            let mut running = self.running.lock().unwrap();
1367            if *running {
1368                // Already running: enqueue as steering instead.
1369                drop(running);
1370                for p in prompts {
1371                    self.queue_steering_with_mode(p, prompt_mode);
1372                }
1373                return;
1374            }
1375            *running = true;
1376        }
1377        let cancel = {
1378            let mut guard = self.cancel.lock().unwrap();
1379            *guard = CancellationToken::new();
1380            guard.clone()
1381        };
1382
1383        // Persist prompts and run.
1384        {
1385            let mut manager = self.manager.lock().unwrap();
1386            for p in &prompts {
1387                let _ = manager.append_message(p.clone());
1388            }
1389        }
1390
1391        let session = self.clone();
1392        let sink: EventSink = Arc::new(move |event: AgentEvent| {
1393            session.on_agent_event(&event);
1394            (session.sink)(SessionEvent::Agent(Box::new(event)));
1395        });
1396
1397        let active_prompt_mode = Arc::new(Mutex::new(prompt_mode));
1398        let mut config = self.loop_config(self, active_prompt_mode.clone());
1399        let mut context = self.build_context_for(prompt_mode);
1400        // The prompts were already persisted. Context includes them, so run
1401        // as a continuation without a second prompt list.
1402
1403        let mut attempt: u32 = 0;
1404        loop {
1405            let messages = kiss_agent::run_agent_loop_continue(
1406                context,
1407                config.clone(),
1408                cancel.clone(),
1409                sink.clone(),
1410            )
1411            .await;
1412
1413            // Retry on transient error stops.
1414            let last_error = messages.iter().rev().find_map(|m| match m {
1415                AgentMessage::Assistant(a) if a.stop_reason == StopReason::Error => {
1416                    Some(a.error_message.clone().unwrap_or_default())
1417                }
1418                _ => None,
1419            });
1420            let settings = self.settings();
1421            let retry = &settings.retry;
1422            if let Some(error) = last_error
1423                && retry.enabled
1424                && attempt < retry.max_retries
1425                && is_transient(&error)
1426                && !cancel.is_cancelled()
1427            {
1428                attempt += 1;
1429                let delay = retry
1430                    .base_delay_ms
1431                    .saturating_mul(1u64.checked_shl(attempt - 1).unwrap_or(u64::MAX))
1432                    .min(retry.max_agent_delay_ms);
1433                (self.sink)(SessionEvent::Retry {
1434                    attempt,
1435                    max: retry.max_retries,
1436                    delay_ms: delay,
1437                    error,
1438                });
1439                tokio::select! {
1440                    _ = tokio::time::sleep(std::time::Duration::from_millis(delay)) => {}
1441                    _ = cancel.cancelled() => break,
1442                }
1443                context = self.build_context_for(*active_prompt_mode.lock().unwrap());
1444                // Drop the trailing error assistant message from context.
1445                while matches!(
1446                    context.messages.last(),
1447                    Some(AgentMessage::Assistant(a)) if a.stop_reason == StopReason::Error
1448                ) {
1449                    context.messages.pop();
1450                }
1451                config.model = self.model();
1452                config.thinking_level = self.thinking_level();
1453                config.fast_mode = self.fast_mode();
1454                continue;
1455            }
1456
1457            // Auto-compaction check after a completed run.
1458            let ctx = self.manager.lock().unwrap().build_session_context();
1459            if auto_compaction_needed(
1460                &settings,
1461                &ctx.messages,
1462                &self.model(),
1463                cancel.is_cancelled(),
1464            ) {
1465                self.compact(None, true).await;
1466            }
1467            break;
1468        }
1469
1470        *self.running.lock().unwrap() = false;
1471        if self.settings.lock().unwrap().cache_warming == CacheWarmingMode::Streaming {
1472            self.cache_warm_cancel.lock().unwrap().cancel();
1473        }
1474    }
1475
1476    fn on_agent_event(self: &Arc<Self>, event: &AgentEvent) {
1477        match event {
1478            AgentEvent::MessageEnd { message } => {
1479                // Persist assistant + tool results (prompts persisted earlier;
1480                // steering/follow-up user messages arrive here too).
1481                let persist = match message {
1482                    AgentMessage::Assistant(a) => {
1483                        self.schedule_cache_warming(a);
1484                        let mut totals = self.totals.lock().unwrap();
1485                        totals.add(&a.usage);
1486                        true
1487                    }
1488                    AgentMessage::ToolResult(_)
1489                    | AgentMessage::User(_)
1490                    | AgentMessage::Custom(_) => true,
1491                    _ => false,
1492                };
1493                if persist {
1494                    // User prompts were persisted in prompt(). Avoid double
1495                    // writes by checking the current leaf message identity.
1496                    let mut manager = self.manager.lock().unwrap();
1497                    let duplicate = matches!(
1498                        (manager.entries().last(), message),
1499                        (Some(crate::session::entry::SessionEntry::Message { message: last, .. }), m) if last == m
1500                    );
1501                    if !duplicate {
1502                        let _ = manager.append_message(message.clone());
1503                    }
1504                }
1505            }
1506            AgentEvent::AgentEnd { .. } => {}
1507            _ => {}
1508        }
1509    }
1510
1511    fn schedule_cache_warming(self: &Arc<Self>, assistant: &kiss_ai::AssistantMessage) {
1512        let settings = self.settings();
1513        let model = self.model();
1514        let Some(cache) = model.prompt_cache else {
1515            return;
1516        };
1517        let Some(short_ttl) = cache.short else {
1518            return;
1519        };
1520        let reasoning = self.thinking_level();
1521        let fast_mode = self.fast_mode();
1522        if settings.cache_warming == CacheWarmingMode::Off
1523            || matches!(
1524                assistant.stop_reason,
1525                StopReason::Error | StopReason::Aborted
1526            )
1527            || (assistant.usage.input + assistant.usage.cache_read + assistant.usage.cache_write
1528                == 0)
1529            || (reasoning != ThinkingLevel::Off
1530                && model.api == "anthropic-messages"
1531                && !model
1532                    .compat
1533                    .as_ref()
1534                    .and_then(|compat| compat.force_adaptive_thinking)
1535                    .unwrap_or(false))
1536        {
1537            return;
1538        }
1539        let ttl = std::time::Duration::from_secs(short_ttl);
1540        let Some(delay) = cache_warming_delay(ttl) else {
1541            return;
1542        };
1543        let prompt_tokens =
1544            assistant.usage.input + assistant.usage.cache_read + assistant.usage.cache_write;
1545        let agent_context = self.build_context();
1546        let context = kiss_ai::Context {
1547            system_prompt: Some(agent_context.system_prompt),
1548            openai_responses_input: agent_context.openai_responses_input,
1549            messages: kiss_agent::convert_to_llm(&agent_context.messages),
1550            tools: agent_context
1551                .tools
1552                .iter()
1553                .map(|tool| tool.to_def())
1554                .collect(),
1555        };
1556        let cancel = CancellationToken::new();
1557        {
1558            let mut current = self.cache_warm_cancel.lock().unwrap();
1559            current.cancel();
1560            *current = cancel.clone();
1561        }
1562        let session = self.clone();
1563        tokio::spawn(async move {
1564            let started = tokio::time::Instant::now();
1565            loop {
1566                let scheduled = tokio::time::Instant::now();
1567                if tokio::select! {
1568                    _ = tokio::time::sleep(delay) => false,
1569                    _ = cancel.cancelled() => true,
1570                } {
1571                    return;
1572                }
1573                if cache_refresh_deadline_missed(scheduled.elapsed(), ttl, delay) {
1574                    return;
1575                }
1576                let idle = !session.is_running();
1577                let max_age = if idle {
1578                    std::time::Duration::from_secs(30 * 60)
1579                } else {
1580                    std::time::Duration::from_secs(60 * 60)
1581                };
1582                if started.elapsed() > max_age {
1583                    return;
1584                }
1585                let priced = |input: u64, output: u64, cache_read: u64, cache_write: u64| {
1586                    let mut usage = Usage {
1587                        input,
1588                        output,
1589                        cache_read,
1590                        cache_write,
1591                        ..Default::default()
1592                    };
1593                    kiss_ai::api::finalize_cost(&mut usage, &model);
1594                    usage.cost.total
1595                };
1596                let hit_cost = priced(0, 0, prompt_tokens, 0);
1597                let miss_cost = if model.cost.cache_write > 0.0 {
1598                    priced(0, 0, 0, prompt_tokens)
1599                } else {
1600                    priced(prompt_tokens, 0, 0, 0)
1601                };
1602                let warm_cost = priced(0, 1, prompt_tokens, 0);
1603                let probability = if idle { 0.15 } else { 1.0 };
1604                if probability * (miss_cost - hit_cost).max(0.0) - warm_cost < 0.05 {
1605                    return;
1606                }
1607                let Some(credential) = session.resolve_credential(&model.provider).await else {
1608                    return;
1609                };
1610                let options = kiss_ai::StreamOptions {
1611                    credential: Some(credential),
1612                    max_tokens: Some(1),
1613                    reasoning,
1614                    fast_mode,
1615                    session_id: Some(session.manager.lock().unwrap().session_id().to_string()),
1616                    transport: settings.transport,
1617                    cancel: cancel.clone(),
1618                    ..Default::default()
1619                };
1620                let stream_fn = session
1621                    .stream_fn
1622                    .lock()
1623                    .unwrap()
1624                    .clone()
1625                    .unwrap_or_else(|| Arc::new(kiss_ai::stream_simple));
1626                let warmed = stream_fn(&model, &context, &options).result().await;
1627                if matches!(warmed.stop_reason, StopReason::Error | StopReason::Aborted) {
1628                    return;
1629                }
1630                session.totals.lock().unwrap().add(&warmed.usage);
1631                let _ = session.manager.lock().unwrap().append_usage(
1632                    "cache_warm",
1633                    &warmed.provider,
1634                    warmed.response_model.as_deref().unwrap_or(&warmed.model),
1635                    warmed.usage,
1636                    None,
1637                );
1638            }
1639        });
1640    }
1641
1642    /// Manual or automatic compaction.
1643    pub async fn compact(self: &Arc<Self>, custom_instructions: Option<String>, auto: bool) {
1644        (self.sink)(SessionEvent::CompactionStart { auto });
1645        let ctx = self.manager.lock().unwrap().build_session_context();
1646        let previous_summary = ctx.messages.iter().rev().find_map(|m| match m {
1647            AgentMessage::CompactionSummary(c) => Some(c.summary.clone()),
1648            _ => None,
1649        });
1650        let settings = self.settings();
1651        let model = self.model();
1652        let override_settings = settings
1653            .compaction
1654            .model_overrides
1655            .get(&format!("{}/{}", model.provider, model.id));
1656        let keep_recent_tokens = override_settings
1657            .and_then(|value| value.keep_recent_tokens)
1658            .unwrap_or(settings.compaction.keep_recent_tokens);
1659        let reserve_tokens = override_settings
1660            .and_then(|value| value.reserve_tokens)
1661            .unwrap_or(settings.compaction.reserve_tokens);
1662        let plan = plan_compaction(&ctx.messages, keep_recent_tokens);
1663        if plan.to_summarize.is_empty() && plan.turn_prefix.is_empty() {
1664            (self.sink)(SessionEvent::CompactionEnd {
1665                summary: String::new(),
1666                tokens_before: plan.tokens_before,
1667                error: Some("Nothing to compact".into()),
1668            });
1669            return;
1670        }
1671
1672        if settings.compaction.mode == CompactionMode::Jev
1673            && let Ok(Some(api_key)) =
1674                kiss_ai::auth::resolve_api_key_async("typesafe", &self.registry.declared_keys).await
1675        {
1676            let cancel = self.cancel.lock().unwrap().clone();
1677            let pinned_start = ctx.messages.len().saturating_sub(plan.kept.len());
1678            if let Ok(result) =
1679                crate::jev::compact(&ctx.messages, pinned_start, &api_key, cancel).await
1680            {
1681                let estimated_before: u64 = ctx
1682                    .messages
1683                    .iter()
1684                    .map(compaction::estimate_message_tokens)
1685                    .sum();
1686                let estimated_after: u64 = result
1687                    .messages
1688                    .iter()
1689                    .map(compaction::estimate_message_tokens)
1690                    .sum();
1691                let removed = estimated_before.saturating_sub(estimated_after);
1692                let tokens_after = plan.tokens_before.saturating_sub(removed);
1693                let useful =
1694                    estimated_after.saturating_mul(4) <= estimated_before.saturating_mul(3);
1695                let resolved_auto_threshold =
1696                    !auto || !should_compact(tokens_after, model.context_window, reserve_tokens);
1697                if useful && resolved_auto_threshold {
1698                    let summary = format!(
1699                        "Jev kept {}, truncated {}, and removed {} of {} older tool interactions",
1700                        result.stats.kept,
1701                        result.stats.truncated,
1702                        result.stats.dropped,
1703                        result.stats.eligible,
1704                    );
1705                    let details = serde_json::json!({
1706                        "mode": "jev",
1707                        "stats": result.stats,
1708                        "estimatedTokensBefore": plan.tokens_before,
1709                        "estimatedTokensAfter": tokens_after,
1710                    });
1711                    let mut manager = self.manager.lock().unwrap();
1712                    let append = manager.append_compaction(
1713                        String::new(),
1714                        plan.tokens_before,
1715                        result.messages,
1716                        None,
1717                        Some(details),
1718                    );
1719                    drop(manager);
1720                    let error = append.err().map(|error| format!("{error:#}"));
1721                    (self.sink)(SessionEvent::CompactionEnd {
1722                        summary: if error.is_none() {
1723                            summary
1724                        } else {
1725                            String::new()
1726                        },
1727                        tokens_before: plan.tokens_before,
1728                        error,
1729                    });
1730                    return;
1731                }
1732            }
1733        }
1734
1735        let credential = self.resolve_credential(&model.provider).await;
1736
1737        let mut serialized = compaction::serialize_agent_messages(&plan.to_summarize);
1738        if plan.is_split_turn {
1739            serialized.push_str(
1740                "\n\n[The following is the earlier part of the still-active task turn:]\n\n",
1741            );
1742            serialized.push_str(&compaction::serialize_agent_messages(&plan.turn_prefix));
1743        }
1744
1745        let summary_cancel = self.cancel.lock().unwrap().clone();
1746        let remote_request = if kiss_ai::api::openai_compaction::supports_remote_compaction(&model)
1747        {
1748            let context = self.build_context();
1749            Some((
1750                kiss_ai::Context {
1751                    system_prompt: Some(context.system_prompt),
1752                    openai_responses_input: context.openai_responses_input,
1753                    messages: kiss_agent::convert_to_llm(&context.messages),
1754                    tools: context.tools.iter().map(|tool| tool.to_def()).collect(),
1755                },
1756                kiss_ai::StreamOptions {
1757                    credential: credential.clone(),
1758                    reasoning: self.thinking_level(),
1759                    fast_mode: self.fast_mode(),
1760                    session_id: Some(self.manager.lock().unwrap().session_id().to_string()),
1761                    cancel: summary_cancel.clone(),
1762                    ..Default::default()
1763                },
1764            ))
1765        } else {
1766            None
1767        };
1768        let local_future = compaction::generate_summary(
1769            &model,
1770            credential.clone(),
1771            &serialized,
1772            previous_summary.as_deref(),
1773            custom_instructions.as_deref(),
1774            plan.is_split_turn,
1775            summary_cancel.clone(),
1776        );
1777        let remote_model = model.clone();
1778        let remote_future = async move {
1779            match remote_request {
1780                Some((context, options)) => Some(
1781                    kiss_ai::api::openai_compaction::compact(&remote_model, &context, &options)
1782                        .await,
1783                ),
1784                None => None,
1785            }
1786        };
1787        let (local_outcome, remote_outcome) = tokio::join!(local_future, remote_future);
1788
1789        match select_compaction_outcome(&model, local_outcome, remote_outcome) {
1790            Ok(result) => {
1791                let mut summarized_all = plan.to_summarize.clone();
1792                summarized_all.extend(plan.turn_prefix.clone());
1793                let (read, modified) = extract_file_ops(&summarized_all);
1794                {
1795                    let mut totals = self.totals.lock().unwrap();
1796                    if let Some(u) = &result.local_usage {
1797                        totals.add(u);
1798                    }
1799                    if let Some(u) = &result.remote_usage {
1800                        totals.add(u);
1801                    }
1802                }
1803                let details = merge_compaction_details(
1804                    file_ops_details(&read, &modified),
1805                    result.remote_details,
1806                );
1807                let mut manager = self.manager.lock().unwrap();
1808                let _ = manager.append_compaction(
1809                    result.summary.clone(),
1810                    plan.tokens_before,
1811                    plan.kept.clone(),
1812                    result.local_usage,
1813                    Some(details),
1814                );
1815                (self.sink)(SessionEvent::CompactionEnd {
1816                    summary: result.summary,
1817                    tokens_before: plan.tokens_before,
1818                    error: None,
1819                });
1820            }
1821            Err(error) => {
1822                (self.sink)(SessionEvent::CompactionEnd {
1823                    summary: String::new(),
1824                    tokens_before: plan.tokens_before,
1825                    error: Some(format!("{error:#}")),
1826                });
1827            }
1828        }
1829    }
1830
1831    /// Current estimated context tokens and window fraction.
1832    pub fn context_usage(&self) -> (u64, u64) {
1833        let manager = self.manager.lock().unwrap();
1834        let revision = manager.context_revision();
1835        let used = if let Some((cached_revision, tokens)) =
1836            *self.context_usage_cache.lock().unwrap()
1837            && cached_revision == revision
1838        {
1839            tokens
1840        } else {
1841            let messages = manager.build_session_context().messages;
1842            let tokens = if self.settings().experimental_context_file {
1843                messages
1844                    .iter()
1845                    .map(compaction::estimate_message_tokens)
1846                    .sum()
1847            } else {
1848                estimate_context_tokens(&messages)
1849            };
1850            *self.context_usage_cache.lock().unwrap() = Some((revision, tokens));
1851            tokens
1852        };
1853        drop(manager);
1854        (used, self.model().context_window)
1855    }
1856}
1857
1858struct SelectedCompaction {
1859    summary: String,
1860    local_usage: Option<Usage>,
1861    remote_usage: Option<Usage>,
1862    remote_details: Option<serde_json::Value>,
1863}
1864
1865fn select_compaction_outcome(
1866    model: &Model,
1867    local: anyhow::Result<compaction::SummaryOutcome>,
1868    remote: Option<anyhow::Result<kiss_ai::api::openai_compaction::RemoteCompactionResult>>,
1869) -> anyhow::Result<SelectedCompaction> {
1870    match (local, remote) {
1871        (Ok(local), Some(Ok(remote))) => Ok(SelectedCompaction {
1872            summary: local.summary,
1873            local_usage: local.usage,
1874            remote_usage: remote.usage,
1875            remote_details: Some(
1876                kiss_ai::api::openai_compaction::build_remote_compaction_details(model, &remote),
1877            ),
1878        }),
1879        (Ok(local), Some(Err(_)) | None) => Ok(SelectedCompaction {
1880            summary: local.summary,
1881            local_usage: local.usage,
1882            remote_usage: None,
1883            remote_details: None,
1884        }),
1885        (Err(_), Some(Ok(remote))) => Ok(SelectedCompaction {
1886            summary: format!(
1887                "OpenAI server-side compaction was applied for {}/{}. The provider-native context is stored in this session, and this notice keeps the compaction boundary readable for other providers.",
1888                model.provider, model.id
1889            ),
1890            local_usage: None,
1891            remote_usage: remote.usage,
1892            remote_details: Some(
1893                kiss_ai::api::openai_compaction::build_remote_compaction_details(model, &remote),
1894            ),
1895        }),
1896        (Err(local), Some(Err(remote))) => anyhow::bail!(
1897            "local compaction failed: {local:#}. OpenAI remote compaction failed: {remote:#}"
1898        ),
1899        (Err(error), None) => Err(error),
1900    }
1901}
1902
1903fn merge_compaction_details(
1904    mut local: serde_json::Value,
1905    remote: Option<serde_json::Value>,
1906) -> serde_json::Value {
1907    let Some(remote) = remote else {
1908        return local;
1909    };
1910    let Some(local_object) = local.as_object_mut() else {
1911        return remote;
1912    };
1913    if let Some(remote_object) = remote.as_object() {
1914        for (key, value) in remote_object {
1915            local_object.insert(key.clone(), value.clone());
1916        }
1917    }
1918    local
1919}
1920
1921fn transcript_excerpt(messages: &[AgentMessage], max_messages: usize, max_chars: usize) -> String {
1922    let mut entries = messages
1923        .iter()
1924        .rev()
1925        .filter_map(|message| match message {
1926            AgentMessage::User(user) => Some(("User", user.content.as_text())),
1927            AgentMessage::Assistant(assistant) => Some(("Assistant", assistant.text())),
1928            _ => None,
1929        })
1930        .filter(|(_, text)| !text.trim().is_empty())
1931        .take(max_messages)
1932        .collect::<Vec<_>>();
1933    entries.reverse();
1934    let transcript = entries
1935        .into_iter()
1936        .map(|(role, text)| format!("{role}: {}", text.trim()))
1937        .collect::<Vec<_>>()
1938        .join("\n\n");
1939    let count = transcript.chars().count();
1940    if count <= max_chars {
1941        return transcript;
1942    }
1943    let omitted = count - max_chars;
1944    let tail = transcript.chars().skip(omitted).collect::<String>();
1945    format!("[earlier text omitted]\n{tail}")
1946}
1947
1948fn drain_queue(queue: &Arc<Mutex<VecDeque<QueuedPrompt>>>, mode: QueueMode) -> Vec<AgentMessage> {
1949    let mut q = queue.lock().unwrap();
1950    match mode {
1951        QueueMode::All => q.drain(..).map(|prompt| prompt.message).collect(),
1952        QueueMode::OneAtATime => q
1953            .pop_front()
1954            .map(|prompt| prompt.message)
1955            .into_iter()
1956            .collect(),
1957    }
1958}
1959
1960fn queued_mode(queue: &Arc<Mutex<VecDeque<QueuedPrompt>>>, mode: QueueMode) -> Option<PromptMode> {
1961    let queue = queue.lock().unwrap();
1962    match mode {
1963        QueueMode::All => queue
1964            .iter()
1965            .any(|prompt| prompt.mode == PromptMode::Workflow)
1966            .then_some(PromptMode::Workflow)
1967            .or_else(|| (!queue.is_empty()).then_some(PromptMode::Ordinary)),
1968        QueueMode::OneAtATime => queue.front().map(|prompt| prompt.mode),
1969    }
1970}
1971
1972fn auto_compaction_needed(
1973    settings: &Settings,
1974    messages: &[AgentMessage],
1975    model: &Model,
1976    cancelled: bool,
1977) -> bool {
1978    let reserve_tokens = settings
1979        .compaction
1980        .model_overrides
1981        .get(&format!("{}/{}", model.provider, model.id))
1982        .and_then(|value| value.reserve_tokens)
1983        .unwrap_or(settings.compaction.reserve_tokens);
1984    settings.compaction.enabled
1985        && !cancelled
1986        && model.context_window > 0
1987        && should_compact(
1988            if settings.experimental_context_file {
1989                // Usage in retained assistant messages describes the old
1990                // request, not the conversation after a model-owned edit.
1991                messages
1992                    .iter()
1993                    .map(compaction::estimate_message_tokens)
1994                    .sum()
1995            } else {
1996                estimate_context_tokens(messages)
1997            },
1998            model.context_window,
1999            reserve_tokens,
2000        )
2001}
2002
2003fn cache_warming_delay(ttl: std::time::Duration) -> Option<std::time::Duration> {
2004    (ttl > std::time::Duration::from_secs(10))
2005        .then(|| std::cmp::min(ttl.mul_f64(0.9), ttl - std::time::Duration::from_secs(10)))
2006}
2007
2008fn cache_refresh_deadline_missed(
2009    elapsed: std::time::Duration,
2010    ttl: std::time::Duration,
2011    delay: std::time::Duration,
2012) -> bool {
2013    elapsed > delay + ttl.saturating_sub(delay) / 2
2014}
2015
2016fn is_transient(error: &str) -> bool {
2017    let e = error.to_lowercase();
2018    if e.contains("subscription_sharing_usage_limit_exceeded") {
2019        return false;
2020    }
2021    let transient_status = [429, 500, 502, 503, 504, 520].iter().any(|status| {
2022        [
2023            format!("http {status}"),
2024            format!("status {status}"),
2025            format!("status: {status}"),
2026            format!("status code {status}"),
2027        ]
2028        .iter()
2029        .any(|marker| e.contains(marker))
2030    });
2031    transient_status
2032        || [
2033            "overloaded",
2034            "currently experiencing high demand",
2035            "rate limit",
2036            "timeout",
2037            "timed out",
2038            "connection reset",
2039            "connection refused",
2040            "connection closed",
2041            "connection aborted",
2042            "connection error",
2043            "failed to connect",
2044            "network error",
2045            "stream error",
2046            "request failed",
2047            "subscription_sharing_usage_unavailable",
2048            "subscription_sharing_user_unavailable",
2049        ]
2050        .iter()
2051        .any(|needle| e.contains(needle))
2052}
2053
2054#[cfg(test)]
2055mod ephemeral_tests {
2056    use super::*;
2057    use std::collections::BTreeMap;
2058    use std::path::Path;
2059
2060    fn openai_model() -> Model {
2061        Model {
2062            id: "gpt-test".into(),
2063            name: "GPT test".into(),
2064            api: "openai-responses".into(),
2065            provider: "openai".into(),
2066            base_url: "https://api.openai.com/v1".into(),
2067            reasoning: true,
2068            input: vec!["text".into()],
2069            cost: Default::default(),
2070            prompt_cache: None,
2071            context_window: 100_000,
2072            max_tokens: 1_000,
2073            compat: None,
2074            thinking_level_map: BTreeMap::new(),
2075            headers: BTreeMap::new(),
2076            sampling_params: Default::default(),
2077        }
2078    }
2079
2080    #[tokio::test]
2081    async fn context_file_bash_edits_reach_next_request_and_survive_resume() {
2082        let directory = tempfile::tempdir().unwrap();
2083        let mut manager =
2084            SessionManager::create(directory.path(), Some(directory.path().join("sessions")))
2085                .unwrap();
2086        manager
2087            .append_message(AgentMessage::user("obsolete output"))
2088            .unwrap();
2089        manager
2090            .append_compaction(
2091                "portable summary".into(),
2092                100,
2093                vec![AgentMessage::user("obsolete output")],
2094                None,
2095                Some(serde_json::json!({"remoteCompaction": {
2096                    "version": 2,
2097                    "provider": "openai-responses-compaction",
2098                    "modelKey": "openai:openai-responses:gpt-test",
2099                    "replacementHistory": [{"type": "compaction", "encrypted_content": "opaque"}]
2100                }})),
2101            )
2102            .unwrap();
2103        assert!(
2104            manager
2105                .build_openai_compaction_context(&openai_model())
2106                .is_some()
2107        );
2108        let session_path = manager.session_file().unwrap().to_path_buf();
2109        let settings = Settings {
2110            experimental_context_file: true,
2111            compaction: crate::settings::CompactionSettings {
2112                enabled: false,
2113                ..Default::default()
2114            },
2115            retry: crate::settings::RetrySettings {
2116                base_delay_ms: 0,
2117                ..Default::default()
2118            },
2119            ..Default::default()
2120        };
2121        let events = Arc::new(Mutex::new(Vec::new()));
2122        let saved_events = events.clone();
2123        let session = AgentSession::new(
2124            manager,
2125            vec![Arc::new(kiss_agent::tools::bash::BashTool::new(
2126                directory.path().to_path_buf(),
2127            ))],
2128            Registry::load(None),
2129            settings,
2130            "test".into(),
2131            openai_model(),
2132            ThinkingLevel::Off,
2133            None,
2134            Arc::new(move |event| saved_events.lock().unwrap().push(event)),
2135        );
2136        let weak = Arc::downgrade(&session);
2137        let calls = Arc::new(std::sync::atomic::AtomicUsize::new(0));
2138        let observed_calls = calls.clone();
2139        session.set_stream_fn(Some(Arc::new(move |_, context, _| {
2140            let step = observed_calls.fetch_add(1, std::sync::atomic::Ordering::SeqCst);
2141            let prompt = context.system_prompt.as_ref().unwrap();
2142            let path = prompt.lines().find_map(|line| line.strip_prefix("Experimental context file: ")).unwrap();
2143            assert!(Path::new(path).is_file());
2144            assert!(context.openai_responses_input.is_none());
2145            let mut message = kiss_ai::AssistantMessage::empty("openai-responses", "openai", "gpt-test");
2146            let replacement = match step {
2147                0 => {
2148                    weak.upgrade().unwrap().queue_steering(AgentMessage::user("new instruction"));
2149                    Some(r#"[{"role":"user","content":"saved notes","timestamp":0}]"#)
2150                }
2151                1 => {
2152                    let text = serde_json::to_string(&context.messages).unwrap();
2153                    assert!(text.contains("saved notes"));
2154                    assert!(!text.contains("obsolete output"));
2155                    assert!(text.contains("new instruction"));
2156                    assert!(context.messages.iter().any(|message| matches!(message, kiss_ai::Message::ToolResult(result) if result.tool_call_id == "edit_0" && !result.is_error)));
2157                    Some("invalid JSON")
2158                }
2159                2 => {
2160                    let text = serde_json::to_string(&context.messages).unwrap();
2161                    assert!(text.contains("saved notes"));
2162                    assert!(text.contains("Repair the file"));
2163                    assert_eq!(std::fs::read_to_string(path).unwrap(), "invalid JSON");
2164                    Some(r#"[{"role":"user","content":"repaired notes","timestamp":0}]"#)
2165                }
2166                3 => {
2167                    let text = serde_json::to_string(&context.messages).unwrap();
2168                    assert!(text.contains("repaired notes"));
2169                    assert!(context.messages.iter().any(|message| matches!(message, kiss_ai::Message::User(user) if user.content.as_text() == "repaired notes")));
2170                    assert!(context.messages.iter().any(|message| matches!(message, kiss_ai::Message::ToolResult(result) if result.tool_call_id == "edit_2" && !result.is_error)));
2171                    None
2172                }
2173                4 => {
2174                    assert!(!matches!(context.messages.last(), Some(kiss_ai::Message::Assistant(assistant)) if assistant.stop_reason == StopReason::Error));
2175                    assert!(context.messages.iter().any(|message| matches!(message, kiss_ai::Message::User(user) if user.content.as_text() == "repaired notes")));
2176                    None
2177                }
2178                _ => panic!("unexpected model request"),
2179            };
2180            if let Some(replacement) = replacement {
2181                let quoted_path = path.replace('\'', "'\\''");
2182                message.content.push(kiss_ai::ContentBlock::ToolCall(kiss_ai::ToolCall {
2183                    id: format!("edit_{step}"),
2184                    name: "bash".into(),
2185                    arguments: serde_json::json!({"command": format!("printf '%s' '{replacement}' > '{quoted_path}'")}),
2186                    thought_signature: None,
2187                }));
2188                message.stop_reason = StopReason::ToolUse;
2189            } else if step == 3 {
2190                message.stop_reason = StopReason::Error;
2191                message.error_message = Some("HTTP 503".into());
2192            } else {
2193                message.content.push(kiss_ai::ContentBlock::text("done"));
2194                message.stop_reason = StopReason::Stop;
2195            }
2196            let (sink, stream) = kiss_ai::EventStream::channel();
2197            sink.send(kiss_ai::AssistantEvent::Start { partial: message.clone() });
2198            if message.stop_reason == StopReason::Error {
2199                sink.error(message);
2200            } else {
2201                sink.done(message);
2202            }
2203            stream
2204        })));
2205        session
2206            .prompt(vec![AgentMessage::user("manage context")])
2207            .await;
2208        assert_eq!(calls.load(std::sync::atomic::Ordering::SeqCst), 5);
2209        let context = session
2210            .manager
2211            .lock()
2212            .unwrap()
2213            .build_session_context()
2214            .messages;
2215        assert_eq!(
2216            SessionManager::open(&session_path)
2217                .unwrap()
2218                .build_session_context()
2219                .messages,
2220            context
2221        );
2222        assert!(
2223            std::fs::read_to_string(session_path)
2224                .unwrap()
2225                .contains("obsolete output")
2226        );
2227        let errors = events.lock().unwrap().iter().filter(|event| matches!(event, SessionEvent::Agent(event) if matches!(event.as_ref(), AgentEvent::MessageEnd { message: AgentMessage::Custom(custom) } if custom.custom_type == "context_file_error"))).count();
2228        assert_eq!(errors, 1);
2229    }
2230
2231    #[test]
2232    fn context_file_lifecycle_is_opt_in_and_session_local() {
2233        let session = AgentSession::new(
2234            SessionManager::in_memory(Path::new("/test")),
2235            Vec::new(),
2236            Registry::load(None),
2237            Settings::default(),
2238            "test".into(),
2239            openai_model(),
2240            ThinkingLevel::Off,
2241            None,
2242            Arc::new(|_| {}),
2243        );
2244        assert!(!session.sync_context_file());
2245        assert!(session.context_file.lock().unwrap().is_none());
2246        assert_eq!(session.build_context().system_prompt, "test");
2247        let mut settings = session.settings();
2248        settings.experimental_context_file = true;
2249        session.update_settings(settings);
2250        session.sync_context_file();
2251        let path = session
2252            .context_file
2253            .lock()
2254            .unwrap()
2255            .as_ref()
2256            .unwrap()
2257            .path()
2258            .to_path_buf();
2259        let child = session
2260            .create_subagent_session("task", "/root/task", ForkTurns::None, None, None)
2261            .unwrap();
2262        child.sync_context_file();
2263        let child_path = child
2264            .context_file
2265            .lock()
2266            .unwrap()
2267            .as_ref()
2268            .unwrap()
2269            .path()
2270            .to_path_buf();
2271        assert_ne!(path, child_path);
2272        assert!(child_path.is_file());
2273        session.replace_manager(SessionManager::in_memory(Path::new("/other")));
2274        assert!(!path.exists());
2275        assert!(session.context_file.lock().unwrap().is_none());
2276        session.sync_context_file();
2277        let new_path = session
2278            .context_file
2279            .lock()
2280            .unwrap()
2281            .as_ref()
2282            .unwrap()
2283            .path()
2284            .to_path_buf();
2285        assert_ne!(new_path, child_path);
2286        let mut settings = session.settings();
2287        settings.experimental_context_file = false;
2288        session.update_settings(settings);
2289        assert!(session.sync_context_file());
2290        assert!(!new_path.exists());
2291        assert_eq!(session.build_context().system_prompt, "test");
2292    }
2293
2294    #[test]
2295    fn reasoning_lease_ends_when_model_saved_effort_or_applied_effort_changes() {
2296        let model = openai_model();
2297        let mut run = ReasoningRunState {
2298            lease: crate::jev::ReasoningLease::default(),
2299            provider: model.provider.clone(),
2300            model_id: model.id.clone(),
2301            saved_effort: ThinkingLevel::Medium,
2302        };
2303        let selection = crate::jev::ReasoningSelection {
2304            level: ThinkingLevel::High,
2305            generations: 5,
2306        };
2307        run.lease.install(&selection);
2308        assert_eq!(
2309            run.checkpoint(&model, ThinkingLevel::Medium, ThinkingLevel::High, true),
2310            (false, Some(ThinkingLevel::High))
2311        );
2312
2313        let mut other_model = model.clone();
2314        other_model.provider = "other".into();
2315        assert_eq!(
2316            run.checkpoint(
2317                &other_model,
2318                ThinkingLevel::Medium,
2319                ThinkingLevel::High,
2320                true
2321            ),
2322            (true, None)
2323        );
2324        run.lease.install(&selection);
2325        assert_eq!(
2326            run.checkpoint(&other_model, ThinkingLevel::Low, ThinkingLevel::High, true),
2327            (false, None)
2328        );
2329        run.lease.install(&selection);
2330        assert_eq!(
2331            run.checkpoint(&other_model, ThinkingLevel::Low, ThinkingLevel::Low, true),
2332            (false, None)
2333        );
2334        run.lease.install(&selection);
2335        assert_eq!(
2336            run.checkpoint(&other_model, ThinkingLevel::Low, ThinkingLevel::High, false),
2337            (false, None)
2338        );
2339    }
2340
2341    #[test]
2342    fn late_cache_refreshes_are_skipped_before_the_cache_expires() {
2343        let ttl = std::time::Duration::from_secs(300);
2344        let delay = cache_warming_delay(ttl).unwrap();
2345        assert!(!cache_refresh_deadline_missed(
2346            std::time::Duration::from_secs(284),
2347            ttl,
2348            delay,
2349        ));
2350        assert!(cache_refresh_deadline_missed(
2351            std::time::Duration::from_secs(286),
2352            ttl,
2353            delay,
2354        ));
2355    }
2356
2357    #[test]
2358    fn transient_errors_require_a_status_or_specific_network_failure() {
2359        assert!(is_transient("request failed with HTTP 503"));
2360        assert!(is_transient("Cloudflare returned HTTP 520"));
2361        assert!(is_transient("Azure is currently experiencing high demand"));
2362        assert!(is_transient("connection reset by peer"));
2363        assert!(is_transient("rate limit exceeded"));
2364        assert!(!is_transient("model has a 500 token limit"));
2365        assert!(!is_transient("connection settings are invalid"));
2366        assert!(!is_transient(
2367            "HTTP 429: subscription_sharing_usage_limit_exceeded"
2368        ));
2369        assert!(is_transient("subscription_sharing_usage_unavailable"));
2370        assert!(is_transient("subscription_sharing_user_unavailable"));
2371    }
2372
2373    #[test]
2374    fn session_tools_install_replace_and_remove_without_changing_base_tools() {
2375        let registry = Registry::load(None);
2376        let session = AgentSession::new(
2377            SessionManager::in_memory(std::path::Path::new("/test")),
2378            Vec::new(),
2379            registry,
2380            Settings::default(),
2381            "test".into(),
2382            openai_model(),
2383            ThinkingLevel::Off,
2384            None,
2385            Arc::new(|_| {}),
2386        );
2387        let tool = || {
2388            Arc::new(crate::tools::grep::GrepTool {
2389                cwd: std::path::PathBuf::from("/test"),
2390            }) as DynTool
2391        };
2392
2393        assert!(!session.available_tool_names().contains(&"grep".into()));
2394        session.install_session_tool(tool());
2395        session.install_session_tool(tool());
2396        assert_eq!(
2397            session
2398                .available_tool_names()
2399                .iter()
2400                .filter(|name| name.as_str() == "grep")
2401                .count(),
2402            1
2403        );
2404        assert!(session.remove_session_tool("grep"));
2405        assert!(!session.remove_session_tool("grep"));
2406        assert!(!session.available_tool_names().contains(&"grep".into()));
2407    }
2408
2409    fn remote_result() -> kiss_ai::api::openai_compaction::RemoteCompactionResult {
2410        kiss_ai::api::openai_compaction::RemoteCompactionResult {
2411            replacement_history: vec![serde_json::json!({
2412                "type": "compaction",
2413                "encrypted_content": "opaque"
2414            })],
2415            usage: Some(Usage {
2416                input: 10,
2417                output: 2,
2418                total_tokens: 12,
2419                ..Default::default()
2420            }),
2421        }
2422    }
2423
2424    fn settings_test_session(settings: Settings, subagents_allowed: bool) -> Arc<AgentSession> {
2425        let registry = Registry::from_builtin();
2426        let model = registry.all().first().expect("built-in model").clone();
2427        AgentSession::new_with_subagents_allowed(
2428            SessionManager::in_memory(std::path::Path::new("/test")),
2429            Vec::new(),
2430            registry,
2431            settings,
2432            "root prompt".into(),
2433            model,
2434            ThinkingLevel::Off,
2435            None,
2436            Arc::new(|_| {}),
2437            subagents_allowed,
2438        )
2439    }
2440
2441    fn benchmark_tools() -> Vec<DynTool> {
2442        let cwd = std::path::PathBuf::from("/synthetic");
2443        vec![
2444            Arc::new(kiss_agent::tools::read::ReadTool { cwd: cwd.clone() }),
2445            Arc::new(kiss_agent::tools::write::WriteTool { cwd: cwd.clone() }),
2446            Arc::new(kiss_agent::tools::edit::EditTool { cwd: cwd.clone() }),
2447            Arc::new(kiss_agent::tools::bash::BashTool::new(cwd)),
2448        ]
2449    }
2450
2451    #[test]
2452    fn subagent_tools_follow_settings_and_command_line_authority() {
2453        let session = settings_test_session(Settings::default(), true);
2454        assert!(session.available_tool_names().is_empty());
2455        assert!(
2456            !session
2457                .build_context()
2458                .system_prompt
2459                .contains("Subagent coordination")
2460        );
2461
2462        let mut enabled = session.settings();
2463        enabled.subagents.enabled = true;
2464        session.update_settings(enabled.clone());
2465        assert_eq!(
2466            session.available_tool_names(),
2467            [
2468                "spawn_agent",
2469                "send_message",
2470                "followup_task",
2471                "wait_agent",
2472                "list_agents",
2473                "interrupt_agent"
2474            ]
2475        );
2476        assert!(
2477            session
2478                .build_context()
2479                .system_prompt
2480                .contains("Subagent coordination")
2481        );
2482
2483        enabled.subagents.enabled = false;
2484        session.update_settings(enabled);
2485        assert!(session.available_tool_names().is_empty());
2486
2487        let mut blocked_settings = Settings::default();
2488        blocked_settings.subagents.enabled = true;
2489        let blocked = settings_test_session(blocked_settings, false);
2490        assert!(blocked.available_tool_names().is_empty());
2491        assert!(
2492            !blocked
2493                .build_context()
2494                .system_prompt
2495                .contains("Subagent coordination")
2496        );
2497    }
2498
2499    #[test]
2500    fn the_workflow_tool_appears_only_in_workflow_prompt_mode() {
2501        let mut settings = Settings::default();
2502        settings.subagents.enabled = true;
2503        let session = settings_test_session(settings, true);
2504
2505        // Subagents on, workflow not armed: an ordinary coding turn pays
2506        // nothing for the feature.
2507        assert!(
2508            !session
2509                .available_tool_names()
2510                .contains(&"run_workflow".into())
2511        );
2512        assert_eq!(
2513            session.prompt_mode_for("run a dynamic workflow for this task"),
2514            PromptMode::Workflow
2515        );
2516        assert_eq!(
2517            session.prompt_mode_for("fix this small function"),
2518            PromptMode::Ordinary
2519        );
2520        assert!(
2521            !session
2522                .build_context()
2523                .system_prompt
2524                .contains("Writing a dynamic workflow")
2525        );
2526
2527        assert!(
2528            session
2529                .available_tool_names_for(PromptMode::Workflow)
2530                .contains(&"run_workflow".into())
2531        );
2532        assert!(
2533            session
2534                .build_context_for(PromptMode::Workflow)
2535                .system_prompt
2536                .contains("Writing a dynamic workflow")
2537        );
2538        assert!(
2539            !session
2540                .available_tool_names()
2541                .contains(&"run_workflow".into())
2542        );
2543    }
2544
2545    #[test]
2546    fn workflow_prompt_mode_does_nothing_while_subagents_are_off() {
2547        // Workflows are built on child agents, so the subagent setting is the
2548        // authority for both.
2549        let session = settings_test_session(Settings::default(), true);
2550        assert!(!session.workflows_enabled());
2551        assert!(
2552            !session
2553                .available_tool_names_for(PromptMode::Workflow)
2554                .contains(&"run_workflow".into())
2555        );
2556
2557        let mut settings = session.settings();
2558        settings.subagents.enabled = true;
2559        settings.workflows.enabled = false;
2560        session.update_settings(settings.clone());
2561        assert!(!session.workflows_enabled());
2562        assert!(
2563            !session
2564                .available_tool_names_for(PromptMode::Workflow)
2565                .contains(&"run_workflow".into())
2566        );
2567
2568        settings.workflows.enabled = true;
2569        session.update_settings(settings);
2570        assert!(session.workflows_enabled());
2571        assert!(
2572            session
2573                .available_tool_names_for(PromptMode::Workflow)
2574                .contains(&"run_workflow".into())
2575        );
2576    }
2577
2578    #[test]
2579    fn one_at_a_time_queues_keep_each_prompts_mode() {
2580        let queue = Arc::new(Mutex::new(VecDeque::from([
2581            QueuedPrompt {
2582                message: AgentMessage::user("ordinary"),
2583                mode: PromptMode::Ordinary,
2584            },
2585            QueuedPrompt {
2586                message: AgentMessage::user("workflow"),
2587                mode: PromptMode::Workflow,
2588            },
2589        ])));
2590
2591        assert_eq!(
2592            queued_mode(&queue, QueueMode::OneAtATime),
2593            Some(PromptMode::Ordinary)
2594        );
2595        assert_eq!(drain_queue(&queue, QueueMode::OneAtATime).len(), 1);
2596        assert_eq!(
2597            queued_mode(&queue, QueueMode::OneAtATime),
2598            Some(PromptMode::Workflow)
2599        );
2600    }
2601
2602    #[test]
2603    fn a_session_without_subagent_authority_has_no_workflow_runtime() {
2604        let mut settings = Settings::default();
2605        settings.subagents.enabled = true;
2606        let child = settings_test_session(settings, false);
2607        assert!(child.workflows().is_none());
2608        assert!(!child.workflows_enabled());
2609    }
2610
2611    #[test]
2612    fn child_session_has_safe_forked_context_without_control_tools() {
2613        let mut settings = Settings::default();
2614        settings.subagents.enabled = true;
2615        let parent = settings_test_session(settings, true);
2616        parent
2617            .manager
2618            .lock()
2619            .unwrap()
2620            .append_message(AgentMessage::user("parent context"))
2621            .unwrap();
2622
2623        let child = parent
2624            .create_subagent_session("inspect", "/root/inspect", ForkTurns::All, None, None)
2625            .unwrap();
2626        assert!(Arc::ptr_eq(&parent.registry, &child.registry));
2627        assert!(child.available_tool_names().is_empty());
2628        let context = child.manager.lock().unwrap().build_session_context();
2629        assert!(matches!(
2630            context.messages.as_slice(),
2631            [AgentMessage::User(user)] if user.content.as_text() == "parent context"
2632        ));
2633    }
2634
2635    #[test]
2636    fn child_without_forked_turns_does_not_copy_parent_context() {
2637        let parent = settings_test_session(Settings::default(), true);
2638        parent
2639            .manager
2640            .lock()
2641            .unwrap()
2642            .append_message(AgentMessage::user("parent context"))
2643            .unwrap();
2644
2645        let child = parent
2646            .create_subagent_session("inspect", "/root/inspect", ForkTurns::None, None, None)
2647            .unwrap();
2648        assert!(
2649            child
2650                .manager
2651                .lock()
2652                .unwrap()
2653                .build_session_context()
2654                .messages
2655                .is_empty()
2656        );
2657    }
2658
2659    #[test]
2660    #[ignore = "release-mode performance benchmark"]
2661    fn benchmark_performance_subagent_overhead() {
2662        let registry = Registry::from_builtin();
2663        let model = registry.all().first().expect("built-in model").clone();
2664        let tools = benchmark_tools();
2665        let make_session = |enabled: bool| {
2666            let mut settings = Settings::default();
2667            settings.subagents.enabled = enabled;
2668            AgentSession::new_with_subagents_allowed(
2669                SessionManager::in_memory(std::path::Path::new("/synthetic")),
2670                tools.clone(),
2671                registry.clone(),
2672                settings,
2673                "benchmark root prompt".into(),
2674                model.clone(),
2675                ThinkingLevel::Off,
2676                None,
2677                Arc::new(|_| {}),
2678                true,
2679            )
2680        };
2681
2682        kiss_bench::measure_pair(
2683            (
2684                "agent_session_create_subagents_off",
2685                "agent_session_create_subagents_on",
2686            ),
2687            21,
2688            500,
2689            (
2690                "new_root_session_4_base_tools_0_control_tools",
2691                "new_root_session_4_base_tools_6_control_tools",
2692            ),
2693            || make_session(false),
2694            || make_session(true),
2695        );
2696
2697        let off = make_session(false);
2698        let on = make_session(true);
2699        kiss_bench::measure_pair(
2700            (
2701                "agent_context_build_subagents_off",
2702                "agent_context_build_subagents_on",
2703            ),
2704            21,
2705            10_000,
2706            (
2707                "empty_session_4_base_tools_0_control_tools",
2708                "empty_session_4_base_tools_6_control_tools",
2709            ),
2710            || off.build_context(),
2711            || on.build_context(),
2712        );
2713    }
2714
2715    #[test]
2716    #[ignore = "release-mode performance benchmark"]
2717    fn benchmark_performance_workflow_tool_exposure() {
2718        // An ordinary coding turn must pay nothing for dynamic workflows. The
2719        // disarmed session is the baseline. The armed one carries the extra
2720        // tool and the authoring instructions.
2721        let registry = Registry::from_builtin();
2722        let model = registry.all().first().expect("built-in model").clone();
2723        let tools = benchmark_tools();
2724        let make_session = || {
2725            let mut settings = Settings::default();
2726            settings.subagents.enabled = true;
2727            AgentSession::new_with_subagents_allowed(
2728                SessionManager::in_memory(std::path::Path::new("/synthetic")),
2729                tools.clone(),
2730                registry.clone(),
2731                settings,
2732                "benchmark root prompt".into(),
2733                model.clone(),
2734                ThinkingLevel::Off,
2735                None,
2736                Arc::new(|_| {}),
2737                true,
2738            )
2739        };
2740
2741        let ordinary = make_session();
2742        let workflow = make_session();
2743        kiss_bench::measure_pair(
2744            (
2745                "agent_context_build_workflow_disarmed",
2746                "agent_context_build_workflow_armed",
2747            ),
2748            21,
2749            10_000,
2750            (
2751                "empty_session_subagents_on_workflow_disarmed",
2752                "empty_session_subagents_on_workflow_armed",
2753            ),
2754            || ordinary.build_context_for(PromptMode::Ordinary),
2755            || workflow.build_context_for(PromptMode::Workflow),
2756        );
2757    }
2758
2759    #[test]
2760    fn transcript_excerpt_keeps_only_recent_user_and_assistant_text() {
2761        let messages = vec![
2762            AgentMessage::user("old"),
2763            AgentMessage::BashExecution(kiss_agent::BashExecutionMessage {
2764                command: "pwd".into(),
2765                output: "ignored".into(),
2766                exit_code: Some(0),
2767                cancelled: false,
2768                truncated: false,
2769                full_output_path: None,
2770                exclude_from_context: false,
2771                timestamp: 1,
2772            }),
2773            AgentMessage::user("new"),
2774        ];
2775        let excerpt = transcript_excerpt(&messages, 1, 100);
2776        assert_eq!(excerpt, "User: new");
2777    }
2778
2779    #[test]
2780    fn transcript_excerpt_enforces_character_budget_from_the_tail() {
2781        let excerpt = transcript_excerpt(&[AgentMessage::user("abcdefghij")], 4, 5);
2782        assert!(excerpt.ends_with("fghij"));
2783        assert!(excerpt.starts_with("[earlier text omitted]"));
2784    }
2785
2786    #[test]
2787    fn hybrid_compaction_keeps_local_summary_and_remote_details() {
2788        let selected = select_compaction_outcome(
2789            &openai_model(),
2790            Ok(compaction::SummaryOutcome {
2791                summary: "portable".into(),
2792                usage: None,
2793            }),
2794            Some(Ok(remote_result())),
2795        )
2796        .unwrap();
2797        assert_eq!(selected.summary, "portable");
2798        assert_eq!(selected.remote_usage.unwrap().input, 10);
2799        assert_eq!(
2800            selected.remote_details.unwrap()["remoteCompaction"]["version"],
2801            2
2802        );
2803    }
2804
2805    #[test]
2806    fn remote_failure_falls_back_to_local_compaction() {
2807        let selected = select_compaction_outcome(
2808            &openai_model(),
2809            Ok(compaction::SummaryOutcome {
2810                summary: "portable".into(),
2811                usage: None,
2812            }),
2813            Some(Err(anyhow::anyhow!("remote unavailable"))),
2814        )
2815        .unwrap();
2816        assert_eq!(selected.summary, "portable");
2817        assert!(selected.remote_details.is_none());
2818    }
2819
2820    #[test]
2821    fn remote_success_survives_local_summary_failure() {
2822        let selected = select_compaction_outcome(
2823            &openai_model(),
2824            Err(anyhow::anyhow!("summary unavailable")),
2825            Some(Ok(remote_result())),
2826        )
2827        .unwrap();
2828        assert!(
2829            selected
2830                .summary
2831                .contains("server-side compaction was applied")
2832        );
2833        assert!(selected.remote_details.is_some());
2834    }
2835
2836    #[test]
2837    fn details_merge_keeps_file_operations_and_remote_artifact() {
2838        let merged = merge_compaction_details(
2839            serde_json::json!({"readFiles": ["a.rs"], "modifiedFiles": []}),
2840            Some(serde_json::json!({"remoteCompaction": {"version": 2}})),
2841        );
2842        assert_eq!(merged["readFiles"][0], "a.rs");
2843        assert_eq!(merged["remoteCompaction"]["version"], 2);
2844    }
2845
2846    #[test]
2847    fn auto_compaction_guard_checks_settings_threshold_and_cancel() {
2848        let mut settings = Settings::default();
2849        settings.compaction.reserve_tokens = 20;
2850        let messages = vec![AgentMessage::user("x".repeat(360))];
2851        let mut model = openai_model();
2852        model.context_window = 100;
2853        assert!(auto_compaction_needed(&settings, &messages, &model, false));
2854        assert!(!auto_compaction_needed(&settings, &messages, &model, true));
2855        settings.compaction.enabled = false;
2856        assert!(!auto_compaction_needed(&settings, &messages, &model, false));
2857        settings.compaction.enabled = true;
2858        let mut assistant = kiss_ai::AssistantMessage::empty("test", "test", "test");
2859        assistant.usage.input = 1000;
2860        assistant
2861            .content
2862            .push(kiss_ai::ContentBlock::text("short note"));
2863        let edited = vec![
2864            AgentMessage::user("task"),
2865            AgentMessage::Assistant(assistant),
2866        ];
2867        assert!(auto_compaction_needed(&settings, &edited, &model, false));
2868        settings.experimental_context_file = true;
2869        assert!(!auto_compaction_needed(&settings, &edited, &model, false));
2870    }
2871
2872    #[test]
2873    fn session_title_normalization_is_safe_and_bounded() {
2874        assert_eq!(
2875            normalize_session_title("  `Fix AUTH-123 login flow!`  \nignored").as_deref(),
2876            Some("Fix AUTH-123 login flow")
2877        );
2878        assert_eq!(normalize_session_title("\n\t"), None);
2879        assert_eq!(
2880            normalize_session_title("🚀".repeat(50).as_str())
2881                .unwrap()
2882                .chars()
2883                .count(),
2884            SESSION_TITLE_MAX_CHARS
2885        );
2886    }
2887
2888    #[test]
2889    fn session_title_prompt_is_utf8_safe_and_bounded() {
2890        let prompt = "🚀".repeat(SESSION_TITLE_PROMPT_MAX_BYTES);
2891        let bounded = bounded_session_title_prompt(&prompt);
2892        assert!(bounded.len() <= SESSION_TITLE_PROMPT_MAX_BYTES);
2893        assert!(std::str::from_utf8(bounded.as_bytes()).is_ok());
2894    }
2895
2896    #[test]
2897    fn cache_warming_uses_ninety_percent_with_ten_second_margin() {
2898        assert_eq!(
2899            cache_warming_delay(std::time::Duration::from_secs(300)),
2900            Some(std::time::Duration::from_secs(270))
2901        );
2902        assert_eq!(
2903            cache_warming_delay(std::time::Duration::from_secs(60)),
2904            Some(std::time::Duration::from_secs(50))
2905        );
2906        assert_eq!(
2907            cache_warming_delay(std::time::Duration::from_secs(10)),
2908            None
2909        );
2910    }
2911}