1use 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#[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 Workflow {
60 run: crate::workflows::RunId,
61 version: u64,
62 },
63 WorkflowOutcome {
66 run: Option<crate::workflows::RunId>,
67 name: String,
68 status: WorkflowTurnStatus,
69 },
70 Iterative {
72 job: crate::iterative::JobId,
73 version: u64,
74 },
75 AccountSwitched(kiss_ai::auth::accounts::AccountSwitch),
78}
79
80pub type SessionEventSink = Arc<dyn Fn(SessionEvent) + Send + Sync>;
81
82#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
87pub enum PromptMode {
88 #[default]
89 Ordinary,
90 Workflow,
91}
92
93#[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 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 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 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 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 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 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 pub fn workflows_enabled(&self) -> bool {
448 self.subagents_enabled() && self.settings.lock().unwrap().workflows.enabled
449 }
450
451 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 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 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 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 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 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 *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 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 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 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 pub async fn prompt(self: &Arc<Self>, prompts: Vec<AgentMessage>) {
1364 self.prompt_with_mode(prompts, PromptMode::Ordinary).await;
1365 }
1366
1367 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 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 {
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 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 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 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 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 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 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 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 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 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 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 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 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}