Skip to main content

kiss_coding/
session_runner.rs

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