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 let Some(value) = error
2019        .find('{')
2020        .and_then(|start| serde_json::from_str::<serde_json::Value>(&error[start..]).ok())
2021    {
2022        let detail = value.get("error").unwrap_or(&value);
2023        if detail
2024            .get("isRetryable")
2025            .and_then(serde_json::Value::as_bool)
2026            == Some(false)
2027            || detail
2028                .get("details")
2029                .and_then(serde_json::Value::as_array)
2030                .is_some_and(|details| {
2031                    details.iter().any(|detail| {
2032                        detail
2033                            .pointer("/debug/details/isRetryable")
2034                            .and_then(serde_json::Value::as_bool)
2035                            == Some(false)
2036                    })
2037                })
2038        {
2039            return false;
2040        }
2041    }
2042    if e.contains("subscription_sharing_usage_limit_exceeded") {
2043        return false;
2044    }
2045    let transient_status = [429, 500, 502, 503, 504, 520].iter().any(|status| {
2046        [
2047            format!("http {status}"),
2048            format!("status {status}"),
2049            format!("status: {status}"),
2050            format!("status code {status}"),
2051        ]
2052        .iter()
2053        .any(|marker| e.contains(marker))
2054    });
2055    transient_status
2056        || [
2057            "overloaded",
2058            "currently experiencing high demand",
2059            "rate limit",
2060            "timeout",
2061            "timed out",
2062            "connection reset",
2063            "connection refused",
2064            "connection closed",
2065            "connection aborted",
2066            "connection error",
2067            "failed to connect",
2068            "network error",
2069            "stream error",
2070            "subscription_sharing_usage_unavailable",
2071            "subscription_sharing_user_unavailable",
2072        ]
2073        .iter()
2074        .any(|needle| e.contains(needle))
2075}
2076
2077#[cfg(test)]
2078mod ephemeral_tests {
2079    use super::*;
2080    use std::collections::BTreeMap;
2081    use std::path::Path;
2082
2083    fn openai_model() -> Model {
2084        Model {
2085            id: "gpt-test".into(),
2086            name: "GPT test".into(),
2087            api: "openai-responses".into(),
2088            provider: "openai".into(),
2089            base_url: "https://api.openai.com/v1".into(),
2090            reasoning: true,
2091            input: vec!["text".into()],
2092            cost: Default::default(),
2093            prompt_cache: None,
2094            context_window: 100_000,
2095            max_tokens: 1_000,
2096            compat: None,
2097            thinking_level_map: BTreeMap::new(),
2098            headers: BTreeMap::new(),
2099            sampling_params: Default::default(),
2100        }
2101    }
2102
2103    #[tokio::test]
2104    async fn context_file_bash_edits_reach_next_request_and_survive_resume() {
2105        let directory = tempfile::tempdir().unwrap();
2106        let mut manager =
2107            SessionManager::create(directory.path(), Some(directory.path().join("sessions")))
2108                .unwrap();
2109        manager
2110            .append_message(AgentMessage::user("obsolete output"))
2111            .unwrap();
2112        manager
2113            .append_compaction(
2114                "portable summary".into(),
2115                100,
2116                vec![AgentMessage::user("obsolete output")],
2117                None,
2118                Some(serde_json::json!({"remoteCompaction": {
2119                    "version": 2,
2120                    "provider": "openai-responses-compaction",
2121                    "modelKey": "openai:openai-responses:gpt-test",
2122                    "replacementHistory": [{"type": "compaction", "encrypted_content": "opaque"}]
2123                }})),
2124            )
2125            .unwrap();
2126        assert!(
2127            manager
2128                .build_openai_compaction_context(&openai_model())
2129                .is_some()
2130        );
2131        let session_path = manager.session_file().unwrap().to_path_buf();
2132        let settings = Settings {
2133            experimental_context_file: true,
2134            compaction: crate::settings::CompactionSettings {
2135                enabled: false,
2136                ..Default::default()
2137            },
2138            retry: crate::settings::RetrySettings {
2139                base_delay_ms: 0,
2140                ..Default::default()
2141            },
2142            ..Default::default()
2143        };
2144        let events = Arc::new(Mutex::new(Vec::new()));
2145        let saved_events = events.clone();
2146        let session = AgentSession::new(
2147            manager,
2148            vec![Arc::new(kiss_agent::tools::bash::BashTool::new(
2149                directory.path().to_path_buf(),
2150            ))],
2151            Registry::load(None),
2152            settings,
2153            "test".into(),
2154            openai_model(),
2155            ThinkingLevel::Off,
2156            None,
2157            Arc::new(move |event| saved_events.lock().unwrap().push(event)),
2158        );
2159        let weak = Arc::downgrade(&session);
2160        let calls = Arc::new(std::sync::atomic::AtomicUsize::new(0));
2161        let observed_calls = calls.clone();
2162        session.set_stream_fn(Some(Arc::new(move |_, context, _| {
2163            let step = observed_calls.fetch_add(1, std::sync::atomic::Ordering::SeqCst);
2164            let prompt = context.system_prompt.as_ref().unwrap();
2165            let path = prompt.lines().find_map(|line| line.strip_prefix("Experimental context file: ")).unwrap();
2166            assert!(Path::new(path).is_file());
2167            assert!(context.openai_responses_input.is_none());
2168            let mut message = kiss_ai::AssistantMessage::empty("openai-responses", "openai", "gpt-test");
2169            let replacement = match step {
2170                0 => {
2171                    weak.upgrade().unwrap().queue_steering(AgentMessage::user("new instruction"));
2172                    Some(r#"[{"role":"user","content":"saved notes","timestamp":0}]"#)
2173                }
2174                1 => {
2175                    let text = serde_json::to_string(&context.messages).unwrap();
2176                    assert!(text.contains("saved notes"));
2177                    assert!(!text.contains("obsolete output"));
2178                    assert!(text.contains("new instruction"));
2179                    assert!(context.messages.iter().any(|message| matches!(message, kiss_ai::Message::ToolResult(result) if result.tool_call_id == "edit_0" && !result.is_error)));
2180                    Some("invalid JSON")
2181                }
2182                2 => {
2183                    let text = serde_json::to_string(&context.messages).unwrap();
2184                    assert!(text.contains("saved notes"));
2185                    assert!(text.contains("Repair the file"));
2186                    assert_eq!(std::fs::read_to_string(path).unwrap(), "invalid JSON");
2187                    Some(r#"[{"role":"user","content":"repaired notes","timestamp":0}]"#)
2188                }
2189                3 => {
2190                    let text = serde_json::to_string(&context.messages).unwrap();
2191                    assert!(text.contains("repaired notes"));
2192                    assert!(context.messages.iter().any(|message| matches!(message, kiss_ai::Message::User(user) if user.content.as_text() == "repaired notes")));
2193                    assert!(context.messages.iter().any(|message| matches!(message, kiss_ai::Message::ToolResult(result) if result.tool_call_id == "edit_2" && !result.is_error)));
2194                    None
2195                }
2196                4 => {
2197                    assert!(!matches!(context.messages.last(), Some(kiss_ai::Message::Assistant(assistant)) if assistant.stop_reason == StopReason::Error));
2198                    assert!(context.messages.iter().any(|message| matches!(message, kiss_ai::Message::User(user) if user.content.as_text() == "repaired notes")));
2199                    None
2200                }
2201                _ => panic!("unexpected model request"),
2202            };
2203            if let Some(replacement) = replacement {
2204                let quoted_path = path.replace('\'', "'\\''");
2205                message.content.push(kiss_ai::ContentBlock::ToolCall(kiss_ai::ToolCall {
2206                    id: format!("edit_{step}"),
2207                    name: "bash".into(),
2208                    arguments: serde_json::json!({"command": format!("printf '%s' '{replacement}' > '{quoted_path}'")}),
2209                    thought_signature: None,
2210                }));
2211                message.stop_reason = StopReason::ToolUse;
2212            } else if step == 3 {
2213                message.stop_reason = StopReason::Error;
2214                message.error_message = Some("HTTP 503".into());
2215            } else {
2216                message.content.push(kiss_ai::ContentBlock::text("done"));
2217                message.stop_reason = StopReason::Stop;
2218            }
2219            let (sink, stream) = kiss_ai::EventStream::channel();
2220            sink.send(kiss_ai::AssistantEvent::Start { partial: message.clone() });
2221            if message.stop_reason == StopReason::Error {
2222                sink.error(message);
2223            } else {
2224                sink.done(message);
2225            }
2226            stream
2227        })));
2228        session
2229            .prompt(vec![AgentMessage::user("manage context")])
2230            .await;
2231        assert_eq!(calls.load(std::sync::atomic::Ordering::SeqCst), 5);
2232        let context = session
2233            .manager
2234            .lock()
2235            .unwrap()
2236            .build_session_context()
2237            .messages;
2238        assert_eq!(
2239            SessionManager::open(&session_path)
2240                .unwrap()
2241                .build_session_context()
2242                .messages,
2243            context
2244        );
2245        assert!(
2246            std::fs::read_to_string(session_path)
2247                .unwrap()
2248                .contains("obsolete output")
2249        );
2250        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();
2251        assert_eq!(errors, 1);
2252    }
2253
2254    #[test]
2255    fn context_file_lifecycle_is_opt_in_and_session_local() {
2256        let session = AgentSession::new(
2257            SessionManager::in_memory(Path::new("/test")),
2258            Vec::new(),
2259            Registry::load(None),
2260            Settings::default(),
2261            "test".into(),
2262            openai_model(),
2263            ThinkingLevel::Off,
2264            None,
2265            Arc::new(|_| {}),
2266        );
2267        assert!(!session.sync_context_file());
2268        assert!(session.context_file.lock().unwrap().is_none());
2269        assert_eq!(session.build_context().system_prompt, "test");
2270        let mut settings = session.settings();
2271        settings.experimental_context_file = true;
2272        session.update_settings(settings);
2273        session.sync_context_file();
2274        let path = session
2275            .context_file
2276            .lock()
2277            .unwrap()
2278            .as_ref()
2279            .unwrap()
2280            .path()
2281            .to_path_buf();
2282        let child = session
2283            .create_subagent_session("task", "/root/task", ForkTurns::None, None, None)
2284            .unwrap();
2285        child.sync_context_file();
2286        let child_path = child
2287            .context_file
2288            .lock()
2289            .unwrap()
2290            .as_ref()
2291            .unwrap()
2292            .path()
2293            .to_path_buf();
2294        assert_ne!(path, child_path);
2295        assert!(child_path.is_file());
2296        session.replace_manager(SessionManager::in_memory(Path::new("/other")));
2297        assert!(!path.exists());
2298        assert!(session.context_file.lock().unwrap().is_none());
2299        session.sync_context_file();
2300        let new_path = session
2301            .context_file
2302            .lock()
2303            .unwrap()
2304            .as_ref()
2305            .unwrap()
2306            .path()
2307            .to_path_buf();
2308        assert_ne!(new_path, child_path);
2309        let mut settings = session.settings();
2310        settings.experimental_context_file = false;
2311        session.update_settings(settings);
2312        assert!(session.sync_context_file());
2313        assert!(!new_path.exists());
2314        assert_eq!(session.build_context().system_prompt, "test");
2315    }
2316
2317    #[test]
2318    fn reasoning_lease_ends_when_model_saved_effort_or_applied_effort_changes() {
2319        let model = openai_model();
2320        let mut run = ReasoningRunState {
2321            lease: crate::jev::ReasoningLease::default(),
2322            provider: model.provider.clone(),
2323            model_id: model.id.clone(),
2324            saved_effort: ThinkingLevel::Medium,
2325        };
2326        let selection = crate::jev::ReasoningSelection {
2327            level: ThinkingLevel::High,
2328            generations: 5,
2329        };
2330        run.lease.install(&selection);
2331        assert_eq!(
2332            run.checkpoint(&model, ThinkingLevel::Medium, ThinkingLevel::High, true),
2333            (false, Some(ThinkingLevel::High))
2334        );
2335
2336        let mut other_model = model.clone();
2337        other_model.provider = "other".into();
2338        assert_eq!(
2339            run.checkpoint(
2340                &other_model,
2341                ThinkingLevel::Medium,
2342                ThinkingLevel::High,
2343                true
2344            ),
2345            (true, None)
2346        );
2347        run.lease.install(&selection);
2348        assert_eq!(
2349            run.checkpoint(&other_model, ThinkingLevel::Low, ThinkingLevel::High, true),
2350            (false, None)
2351        );
2352        run.lease.install(&selection);
2353        assert_eq!(
2354            run.checkpoint(&other_model, ThinkingLevel::Low, ThinkingLevel::Low, true),
2355            (false, None)
2356        );
2357        run.lease.install(&selection);
2358        assert_eq!(
2359            run.checkpoint(&other_model, ThinkingLevel::Low, ThinkingLevel::High, false),
2360            (false, None)
2361        );
2362    }
2363
2364    #[test]
2365    fn late_cache_refreshes_are_skipped_before_the_cache_expires() {
2366        let ttl = std::time::Duration::from_secs(300);
2367        let delay = cache_warming_delay(ttl).unwrap();
2368        assert!(!cache_refresh_deadline_missed(
2369            std::time::Duration::from_secs(284),
2370            ttl,
2371            delay,
2372        ));
2373        assert!(cache_refresh_deadline_missed(
2374            std::time::Duration::from_secs(286),
2375            ttl,
2376            delay,
2377        ));
2378    }
2379
2380    #[test]
2381    fn transient_errors_require_a_status_or_specific_network_failure() {
2382        assert!(is_transient("request failed with HTTP 503"));
2383        assert!(is_transient("Cloudflare returned HTTP 520"));
2384        assert!(is_transient("Azure is currently experiencing high demand"));
2385        assert!(is_transient("connection reset by peer"));
2386        assert!(is_transient("rate limit exceeded"));
2387        assert!(!is_transient("model has a 500 token limit"));
2388        assert!(!is_transient("connection settings are invalid"));
2389        assert!(!is_transient(
2390            "HTTP 429: subscription_sharing_usage_limit_exceeded"
2391        ));
2392        assert!(is_transient("subscription_sharing_usage_unavailable"));
2393        assert!(is_transient("subscription_sharing_user_unavailable"));
2394        assert!(!is_transient(
2395            r#"Cursor request failed: stream ended: {"error":{"code":"not_found","message":"Model name is not valid: auto"}}"#
2396        ));
2397        assert!(!is_transient(
2398            r#"Cursor request failed: stream ended: {"error":{"code":"internal","message":"KISS does not run Cursor-native tools"}}"#
2399        ));
2400        assert!(!is_transient(
2401            r#"Cursor request failed: stream ended: {"error":{"details":[{"debug":{"details":{"isRetryable":false,"detail":"rate limit exceeded"}}}]}}"#
2402        ));
2403    }
2404
2405    #[test]
2406    fn session_tools_install_replace_and_remove_without_changing_base_tools() {
2407        let registry = Registry::load(None);
2408        let session = AgentSession::new(
2409            SessionManager::in_memory(std::path::Path::new("/test")),
2410            Vec::new(),
2411            registry,
2412            Settings::default(),
2413            "test".into(),
2414            openai_model(),
2415            ThinkingLevel::Off,
2416            None,
2417            Arc::new(|_| {}),
2418        );
2419        let tool = || {
2420            Arc::new(crate::tools::grep::GrepTool {
2421                cwd: std::path::PathBuf::from("/test"),
2422            }) as DynTool
2423        };
2424
2425        assert!(!session.available_tool_names().contains(&"grep".into()));
2426        session.install_session_tool(tool());
2427        session.install_session_tool(tool());
2428        assert_eq!(
2429            session
2430                .available_tool_names()
2431                .iter()
2432                .filter(|name| name.as_str() == "grep")
2433                .count(),
2434            1
2435        );
2436        assert!(session.remove_session_tool("grep"));
2437        assert!(!session.remove_session_tool("grep"));
2438        assert!(!session.available_tool_names().contains(&"grep".into()));
2439    }
2440
2441    fn remote_result() -> kiss_ai::api::openai_compaction::RemoteCompactionResult {
2442        kiss_ai::api::openai_compaction::RemoteCompactionResult {
2443            replacement_history: vec![serde_json::json!({
2444                "type": "compaction",
2445                "encrypted_content": "opaque"
2446            })],
2447            usage: Some(Usage {
2448                input: 10,
2449                output: 2,
2450                total_tokens: 12,
2451                ..Default::default()
2452            }),
2453        }
2454    }
2455
2456    fn settings_test_session(settings: Settings, subagents_allowed: bool) -> Arc<AgentSession> {
2457        let registry = Registry::from_builtin();
2458        let model = registry.all().first().expect("built-in model").clone();
2459        AgentSession::new_with_subagents_allowed(
2460            SessionManager::in_memory(std::path::Path::new("/test")),
2461            Vec::new(),
2462            registry,
2463            settings,
2464            "root prompt".into(),
2465            model,
2466            ThinkingLevel::Off,
2467            None,
2468            Arc::new(|_| {}),
2469            subagents_allowed,
2470        )
2471    }
2472
2473    fn benchmark_tools() -> Vec<DynTool> {
2474        let cwd = std::path::PathBuf::from("/synthetic");
2475        vec![
2476            Arc::new(kiss_agent::tools::read::ReadTool { cwd: cwd.clone() }),
2477            Arc::new(kiss_agent::tools::write::WriteTool { cwd: cwd.clone() }),
2478            Arc::new(kiss_agent::tools::edit::EditTool { cwd: cwd.clone() }),
2479            Arc::new(kiss_agent::tools::bash::BashTool::new(cwd)),
2480        ]
2481    }
2482
2483    #[test]
2484    fn subagent_tools_follow_settings_and_command_line_authority() {
2485        let session = settings_test_session(Settings::default(), true);
2486        assert!(session.available_tool_names().is_empty());
2487        assert!(
2488            !session
2489                .build_context()
2490                .system_prompt
2491                .contains("Subagent coordination")
2492        );
2493
2494        let mut enabled = session.settings();
2495        enabled.subagents.enabled = true;
2496        session.update_settings(enabled.clone());
2497        assert_eq!(
2498            session.available_tool_names(),
2499            [
2500                "spawn_agent",
2501                "send_message",
2502                "followup_task",
2503                "wait_agent",
2504                "list_agents",
2505                "interrupt_agent"
2506            ]
2507        );
2508        assert!(
2509            session
2510                .build_context()
2511                .system_prompt
2512                .contains("Subagent coordination")
2513        );
2514
2515        enabled.subagents.enabled = false;
2516        session.update_settings(enabled);
2517        assert!(session.available_tool_names().is_empty());
2518
2519        let mut blocked_settings = Settings::default();
2520        blocked_settings.subagents.enabled = true;
2521        let blocked = settings_test_session(blocked_settings, false);
2522        assert!(blocked.available_tool_names().is_empty());
2523        assert!(
2524            !blocked
2525                .build_context()
2526                .system_prompt
2527                .contains("Subagent coordination")
2528        );
2529    }
2530
2531    #[test]
2532    fn the_workflow_tool_appears_only_in_workflow_prompt_mode() {
2533        let mut settings = Settings::default();
2534        settings.subagents.enabled = true;
2535        let session = settings_test_session(settings, true);
2536
2537        // Subagents on, workflow not armed: an ordinary coding turn pays
2538        // nothing for the feature.
2539        assert!(
2540            !session
2541                .available_tool_names()
2542                .contains(&"run_workflow".into())
2543        );
2544        assert_eq!(
2545            session.prompt_mode_for("run a dynamic workflow for this task"),
2546            PromptMode::Workflow
2547        );
2548        assert_eq!(
2549            session.prompt_mode_for("fix this small function"),
2550            PromptMode::Ordinary
2551        );
2552        assert!(
2553            !session
2554                .build_context()
2555                .system_prompt
2556                .contains("Writing a dynamic workflow")
2557        );
2558
2559        assert!(
2560            session
2561                .available_tool_names_for(PromptMode::Workflow)
2562                .contains(&"run_workflow".into())
2563        );
2564        assert!(
2565            session
2566                .build_context_for(PromptMode::Workflow)
2567                .system_prompt
2568                .contains("Writing a dynamic workflow")
2569        );
2570        assert!(
2571            !session
2572                .available_tool_names()
2573                .contains(&"run_workflow".into())
2574        );
2575    }
2576
2577    #[test]
2578    fn workflow_prompt_mode_does_nothing_while_subagents_are_off() {
2579        // Workflows are built on child agents, so the subagent setting is the
2580        // authority for both.
2581        let session = settings_test_session(Settings::default(), true);
2582        assert!(!session.workflows_enabled());
2583        assert!(
2584            !session
2585                .available_tool_names_for(PromptMode::Workflow)
2586                .contains(&"run_workflow".into())
2587        );
2588
2589        let mut settings = session.settings();
2590        settings.subagents.enabled = true;
2591        settings.workflows.enabled = false;
2592        session.update_settings(settings.clone());
2593        assert!(!session.workflows_enabled());
2594        assert!(
2595            !session
2596                .available_tool_names_for(PromptMode::Workflow)
2597                .contains(&"run_workflow".into())
2598        );
2599
2600        settings.workflows.enabled = true;
2601        session.update_settings(settings);
2602        assert!(session.workflows_enabled());
2603        assert!(
2604            session
2605                .available_tool_names_for(PromptMode::Workflow)
2606                .contains(&"run_workflow".into())
2607        );
2608    }
2609
2610    #[test]
2611    fn one_at_a_time_queues_keep_each_prompts_mode() {
2612        let queue = Arc::new(Mutex::new(VecDeque::from([
2613            QueuedPrompt {
2614                message: AgentMessage::user("ordinary"),
2615                mode: PromptMode::Ordinary,
2616            },
2617            QueuedPrompt {
2618                message: AgentMessage::user("workflow"),
2619                mode: PromptMode::Workflow,
2620            },
2621        ])));
2622
2623        assert_eq!(
2624            queued_mode(&queue, QueueMode::OneAtATime),
2625            Some(PromptMode::Ordinary)
2626        );
2627        assert_eq!(drain_queue(&queue, QueueMode::OneAtATime).len(), 1);
2628        assert_eq!(
2629            queued_mode(&queue, QueueMode::OneAtATime),
2630            Some(PromptMode::Workflow)
2631        );
2632    }
2633
2634    #[test]
2635    fn a_session_without_subagent_authority_has_no_workflow_runtime() {
2636        let mut settings = Settings::default();
2637        settings.subagents.enabled = true;
2638        let child = settings_test_session(settings, false);
2639        assert!(child.workflows().is_none());
2640        assert!(!child.workflows_enabled());
2641    }
2642
2643    #[test]
2644    fn child_session_has_safe_forked_context_without_control_tools() {
2645        let mut settings = Settings::default();
2646        settings.subagents.enabled = true;
2647        let parent = settings_test_session(settings, true);
2648        parent
2649            .manager
2650            .lock()
2651            .unwrap()
2652            .append_message(AgentMessage::user("parent context"))
2653            .unwrap();
2654
2655        let child = parent
2656            .create_subagent_session("inspect", "/root/inspect", ForkTurns::All, None, None)
2657            .unwrap();
2658        assert!(Arc::ptr_eq(&parent.registry, &child.registry));
2659        assert!(child.available_tool_names().is_empty());
2660        let context = child.manager.lock().unwrap().build_session_context();
2661        assert!(matches!(
2662            context.messages.as_slice(),
2663            [AgentMessage::User(user)] if user.content.as_text() == "parent context"
2664        ));
2665    }
2666
2667    #[test]
2668    fn child_without_forked_turns_does_not_copy_parent_context() {
2669        let parent = settings_test_session(Settings::default(), true);
2670        parent
2671            .manager
2672            .lock()
2673            .unwrap()
2674            .append_message(AgentMessage::user("parent context"))
2675            .unwrap();
2676
2677        let child = parent
2678            .create_subagent_session("inspect", "/root/inspect", ForkTurns::None, None, None)
2679            .unwrap();
2680        assert!(
2681            child
2682                .manager
2683                .lock()
2684                .unwrap()
2685                .build_session_context()
2686                .messages
2687                .is_empty()
2688        );
2689    }
2690
2691    #[test]
2692    #[ignore = "release-mode performance benchmark"]
2693    fn benchmark_performance_subagent_overhead() {
2694        let registry = Registry::from_builtin();
2695        let model = registry.all().first().expect("built-in model").clone();
2696        let tools = benchmark_tools();
2697        let make_session = |enabled: bool| {
2698            let mut settings = Settings::default();
2699            settings.subagents.enabled = enabled;
2700            AgentSession::new_with_subagents_allowed(
2701                SessionManager::in_memory(std::path::Path::new("/synthetic")),
2702                tools.clone(),
2703                registry.clone(),
2704                settings,
2705                "benchmark root prompt".into(),
2706                model.clone(),
2707                ThinkingLevel::Off,
2708                None,
2709                Arc::new(|_| {}),
2710                true,
2711            )
2712        };
2713
2714        kiss_bench::measure_pair(
2715            (
2716                "agent_session_create_subagents_off",
2717                "agent_session_create_subagents_on",
2718            ),
2719            21,
2720            500,
2721            (
2722                "new_root_session_4_base_tools_0_control_tools",
2723                "new_root_session_4_base_tools_6_control_tools",
2724            ),
2725            || make_session(false),
2726            || make_session(true),
2727        );
2728
2729        let off = make_session(false);
2730        let on = make_session(true);
2731        kiss_bench::measure_pair(
2732            (
2733                "agent_context_build_subagents_off",
2734                "agent_context_build_subagents_on",
2735            ),
2736            21,
2737            10_000,
2738            (
2739                "empty_session_4_base_tools_0_control_tools",
2740                "empty_session_4_base_tools_6_control_tools",
2741            ),
2742            || off.build_context(),
2743            || on.build_context(),
2744        );
2745    }
2746
2747    #[test]
2748    #[ignore = "release-mode performance benchmark"]
2749    fn benchmark_performance_workflow_tool_exposure() {
2750        // An ordinary coding turn must pay nothing for dynamic workflows. The
2751        // disarmed session is the baseline. The armed one carries the extra
2752        // tool and the authoring instructions.
2753        let registry = Registry::from_builtin();
2754        let model = registry.all().first().expect("built-in model").clone();
2755        let tools = benchmark_tools();
2756        let make_session = || {
2757            let mut settings = Settings::default();
2758            settings.subagents.enabled = true;
2759            AgentSession::new_with_subagents_allowed(
2760                SessionManager::in_memory(std::path::Path::new("/synthetic")),
2761                tools.clone(),
2762                registry.clone(),
2763                settings,
2764                "benchmark root prompt".into(),
2765                model.clone(),
2766                ThinkingLevel::Off,
2767                None,
2768                Arc::new(|_| {}),
2769                true,
2770            )
2771        };
2772
2773        let ordinary = make_session();
2774        let workflow = make_session();
2775        kiss_bench::measure_pair(
2776            (
2777                "agent_context_build_workflow_disarmed",
2778                "agent_context_build_workflow_armed",
2779            ),
2780            21,
2781            10_000,
2782            (
2783                "empty_session_subagents_on_workflow_disarmed",
2784                "empty_session_subagents_on_workflow_armed",
2785            ),
2786            || ordinary.build_context_for(PromptMode::Ordinary),
2787            || workflow.build_context_for(PromptMode::Workflow),
2788        );
2789    }
2790
2791    #[test]
2792    fn transcript_excerpt_keeps_only_recent_user_and_assistant_text() {
2793        let messages = vec![
2794            AgentMessage::user("old"),
2795            AgentMessage::BashExecution(kiss_agent::BashExecutionMessage {
2796                command: "pwd".into(),
2797                output: "ignored".into(),
2798                exit_code: Some(0),
2799                cancelled: false,
2800                truncated: false,
2801                full_output_path: None,
2802                exclude_from_context: false,
2803                timestamp: 1,
2804            }),
2805            AgentMessage::user("new"),
2806        ];
2807        let excerpt = transcript_excerpt(&messages, 1, 100);
2808        assert_eq!(excerpt, "User: new");
2809    }
2810
2811    #[test]
2812    fn transcript_excerpt_enforces_character_budget_from_the_tail() {
2813        let excerpt = transcript_excerpt(&[AgentMessage::user("abcdefghij")], 4, 5);
2814        assert!(excerpt.ends_with("fghij"));
2815        assert!(excerpt.starts_with("[earlier text omitted]"));
2816    }
2817
2818    #[test]
2819    fn hybrid_compaction_keeps_local_summary_and_remote_details() {
2820        let selected = select_compaction_outcome(
2821            &openai_model(),
2822            Ok(compaction::SummaryOutcome {
2823                summary: "portable".into(),
2824                usage: None,
2825            }),
2826            Some(Ok(remote_result())),
2827        )
2828        .unwrap();
2829        assert_eq!(selected.summary, "portable");
2830        assert_eq!(selected.remote_usage.unwrap().input, 10);
2831        assert_eq!(
2832            selected.remote_details.unwrap()["remoteCompaction"]["version"],
2833            2
2834        );
2835    }
2836
2837    #[test]
2838    fn remote_failure_falls_back_to_local_compaction() {
2839        let selected = select_compaction_outcome(
2840            &openai_model(),
2841            Ok(compaction::SummaryOutcome {
2842                summary: "portable".into(),
2843                usage: None,
2844            }),
2845            Some(Err(anyhow::anyhow!("remote unavailable"))),
2846        )
2847        .unwrap();
2848        assert_eq!(selected.summary, "portable");
2849        assert!(selected.remote_details.is_none());
2850    }
2851
2852    #[test]
2853    fn remote_success_survives_local_summary_failure() {
2854        let selected = select_compaction_outcome(
2855            &openai_model(),
2856            Err(anyhow::anyhow!("summary unavailable")),
2857            Some(Ok(remote_result())),
2858        )
2859        .unwrap();
2860        assert!(
2861            selected
2862                .summary
2863                .contains("server-side compaction was applied")
2864        );
2865        assert!(selected.remote_details.is_some());
2866    }
2867
2868    #[test]
2869    fn details_merge_keeps_file_operations_and_remote_artifact() {
2870        let merged = merge_compaction_details(
2871            serde_json::json!({"readFiles": ["a.rs"], "modifiedFiles": []}),
2872            Some(serde_json::json!({"remoteCompaction": {"version": 2}})),
2873        );
2874        assert_eq!(merged["readFiles"][0], "a.rs");
2875        assert_eq!(merged["remoteCompaction"]["version"], 2);
2876    }
2877
2878    #[test]
2879    fn auto_compaction_guard_checks_settings_threshold_and_cancel() {
2880        let mut settings = Settings::default();
2881        settings.compaction.reserve_tokens = 20;
2882        let messages = vec![AgentMessage::user("x".repeat(360))];
2883        let mut model = openai_model();
2884        model.context_window = 100;
2885        assert!(auto_compaction_needed(&settings, &messages, &model, false));
2886        assert!(!auto_compaction_needed(&settings, &messages, &model, true));
2887        settings.compaction.enabled = false;
2888        assert!(!auto_compaction_needed(&settings, &messages, &model, false));
2889        settings.compaction.enabled = true;
2890        let mut assistant = kiss_ai::AssistantMessage::empty("test", "test", "test");
2891        assistant.usage.input = 1000;
2892        assistant
2893            .content
2894            .push(kiss_ai::ContentBlock::text("short note"));
2895        let edited = vec![
2896            AgentMessage::user("task"),
2897            AgentMessage::Assistant(assistant),
2898        ];
2899        assert!(auto_compaction_needed(&settings, &edited, &model, false));
2900        settings.experimental_context_file = true;
2901        assert!(!auto_compaction_needed(&settings, &edited, &model, false));
2902    }
2903
2904    #[test]
2905    fn session_title_normalization_is_safe_and_bounded() {
2906        assert_eq!(
2907            normalize_session_title("  `Fix AUTH-123 login flow!`  \nignored").as_deref(),
2908            Some("Fix AUTH-123 login flow")
2909        );
2910        assert_eq!(normalize_session_title("\n\t"), None);
2911        assert_eq!(
2912            normalize_session_title("🚀".repeat(50).as_str())
2913                .unwrap()
2914                .chars()
2915                .count(),
2916            SESSION_TITLE_MAX_CHARS
2917        );
2918    }
2919
2920    #[test]
2921    fn session_title_prompt_is_utf8_safe_and_bounded() {
2922        let prompt = "🚀".repeat(SESSION_TITLE_PROMPT_MAX_BYTES);
2923        let bounded = bounded_session_title_prompt(&prompt);
2924        assert!(bounded.len() <= SESSION_TITLE_PROMPT_MAX_BYTES);
2925        assert!(std::str::from_utf8(bounded.as_bytes()).is_ok());
2926    }
2927
2928    #[test]
2929    fn cache_warming_uses_ninety_percent_with_ten_second_margin() {
2930        assert_eq!(
2931            cache_warming_delay(std::time::Duration::from_secs(300)),
2932            Some(std::time::Duration::from_secs(270))
2933        );
2934        assert_eq!(
2935            cache_warming_delay(std::time::Duration::from_secs(60)),
2936            Some(std::time::Duration::from_secs(50))
2937        );
2938        assert_eq!(
2939            cache_warming_delay(std::time::Duration::from_secs(10)),
2940            None
2941        );
2942    }
2943}