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