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