Skip to main content

kiss_coding/
session_runner.rs

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