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 e.contains("subscription_sharing_usage_limit_exceeded") {
2019 return false;
2020 }
2021 let transient_status = [429, 500, 502, 503, 504, 520].iter().any(|status| {
2022 [
2023 format!("http {status}"),
2024 format!("status {status}"),
2025 format!("status: {status}"),
2026 format!("status code {status}"),
2027 ]
2028 .iter()
2029 .any(|marker| e.contains(marker))
2030 });
2031 transient_status
2032 || [
2033 "overloaded",
2034 "currently experiencing high demand",
2035 "rate limit",
2036 "timeout",
2037 "timed out",
2038 "connection reset",
2039 "connection refused",
2040 "connection closed",
2041 "connection aborted",
2042 "connection error",
2043 "failed to connect",
2044 "network error",
2045 "stream error",
2046 "request failed",
2047 "subscription_sharing_usage_unavailable",
2048 "subscription_sharing_user_unavailable",
2049 ]
2050 .iter()
2051 .any(|needle| e.contains(needle))
2052}
2053
2054#[cfg(test)]
2055mod ephemeral_tests {
2056 use super::*;
2057 use std::collections::BTreeMap;
2058 use std::path::Path;
2059
2060 fn openai_model() -> Model {
2061 Model {
2062 id: "gpt-test".into(),
2063 name: "GPT test".into(),
2064 api: "openai-responses".into(),
2065 provider: "openai".into(),
2066 base_url: "https://api.openai.com/v1".into(),
2067 reasoning: true,
2068 input: vec!["text".into()],
2069 cost: Default::default(),
2070 prompt_cache: None,
2071 context_window: 100_000,
2072 max_tokens: 1_000,
2073 compat: None,
2074 thinking_level_map: BTreeMap::new(),
2075 headers: BTreeMap::new(),
2076 sampling_params: Default::default(),
2077 }
2078 }
2079
2080 #[tokio::test]
2081 async fn context_file_bash_edits_reach_next_request_and_survive_resume() {
2082 let directory = tempfile::tempdir().unwrap();
2083 let mut manager =
2084 SessionManager::create(directory.path(), Some(directory.path().join("sessions")))
2085 .unwrap();
2086 manager
2087 .append_message(AgentMessage::user("obsolete output"))
2088 .unwrap();
2089 manager
2090 .append_compaction(
2091 "portable summary".into(),
2092 100,
2093 vec![AgentMessage::user("obsolete output")],
2094 None,
2095 Some(serde_json::json!({"remoteCompaction": {
2096 "version": 2,
2097 "provider": "openai-responses-compaction",
2098 "modelKey": "openai:openai-responses:gpt-test",
2099 "replacementHistory": [{"type": "compaction", "encrypted_content": "opaque"}]
2100 }})),
2101 )
2102 .unwrap();
2103 assert!(
2104 manager
2105 .build_openai_compaction_context(&openai_model())
2106 .is_some()
2107 );
2108 let session_path = manager.session_file().unwrap().to_path_buf();
2109 let settings = Settings {
2110 experimental_context_file: true,
2111 compaction: crate::settings::CompactionSettings {
2112 enabled: false,
2113 ..Default::default()
2114 },
2115 retry: crate::settings::RetrySettings {
2116 base_delay_ms: 0,
2117 ..Default::default()
2118 },
2119 ..Default::default()
2120 };
2121 let events = Arc::new(Mutex::new(Vec::new()));
2122 let saved_events = events.clone();
2123 let session = AgentSession::new(
2124 manager,
2125 vec![Arc::new(kiss_agent::tools::bash::BashTool::new(
2126 directory.path().to_path_buf(),
2127 ))],
2128 Registry::load(None),
2129 settings,
2130 "test".into(),
2131 openai_model(),
2132 ThinkingLevel::Off,
2133 None,
2134 Arc::new(move |event| saved_events.lock().unwrap().push(event)),
2135 );
2136 let weak = Arc::downgrade(&session);
2137 let calls = Arc::new(std::sync::atomic::AtomicUsize::new(0));
2138 let observed_calls = calls.clone();
2139 session.set_stream_fn(Some(Arc::new(move |_, context, _| {
2140 let step = observed_calls.fetch_add(1, std::sync::atomic::Ordering::SeqCst);
2141 let prompt = context.system_prompt.as_ref().unwrap();
2142 let path = prompt.lines().find_map(|line| line.strip_prefix("Experimental context file: ")).unwrap();
2143 assert!(Path::new(path).is_file());
2144 assert!(context.openai_responses_input.is_none());
2145 let mut message = kiss_ai::AssistantMessage::empty("openai-responses", "openai", "gpt-test");
2146 let replacement = match step {
2147 0 => {
2148 weak.upgrade().unwrap().queue_steering(AgentMessage::user("new instruction"));
2149 Some(r#"[{"role":"user","content":"saved notes","timestamp":0}]"#)
2150 }
2151 1 => {
2152 let text = serde_json::to_string(&context.messages).unwrap();
2153 assert!(text.contains("saved notes"));
2154 assert!(!text.contains("obsolete output"));
2155 assert!(text.contains("new instruction"));
2156 assert!(context.messages.iter().any(|message| matches!(message, kiss_ai::Message::ToolResult(result) if result.tool_call_id == "edit_0" && !result.is_error)));
2157 Some("invalid JSON")
2158 }
2159 2 => {
2160 let text = serde_json::to_string(&context.messages).unwrap();
2161 assert!(text.contains("saved notes"));
2162 assert!(text.contains("Repair the file"));
2163 assert_eq!(std::fs::read_to_string(path).unwrap(), "invalid JSON");
2164 Some(r#"[{"role":"user","content":"repaired notes","timestamp":0}]"#)
2165 }
2166 3 => {
2167 let text = serde_json::to_string(&context.messages).unwrap();
2168 assert!(text.contains("repaired notes"));
2169 assert!(context.messages.iter().any(|message| matches!(message, kiss_ai::Message::User(user) if user.content.as_text() == "repaired notes")));
2170 assert!(context.messages.iter().any(|message| matches!(message, kiss_ai::Message::ToolResult(result) if result.tool_call_id == "edit_2" && !result.is_error)));
2171 None
2172 }
2173 4 => {
2174 assert!(!matches!(context.messages.last(), Some(kiss_ai::Message::Assistant(assistant)) if assistant.stop_reason == StopReason::Error));
2175 assert!(context.messages.iter().any(|message| matches!(message, kiss_ai::Message::User(user) if user.content.as_text() == "repaired notes")));
2176 None
2177 }
2178 _ => panic!("unexpected model request"),
2179 };
2180 if let Some(replacement) = replacement {
2181 let quoted_path = path.replace('\'', "'\\''");
2182 message.content.push(kiss_ai::ContentBlock::ToolCall(kiss_ai::ToolCall {
2183 id: format!("edit_{step}"),
2184 name: "bash".into(),
2185 arguments: serde_json::json!({"command": format!("printf '%s' '{replacement}' > '{quoted_path}'")}),
2186 thought_signature: None,
2187 }));
2188 message.stop_reason = StopReason::ToolUse;
2189 } else if step == 3 {
2190 message.stop_reason = StopReason::Error;
2191 message.error_message = Some("HTTP 503".into());
2192 } else {
2193 message.content.push(kiss_ai::ContentBlock::text("done"));
2194 message.stop_reason = StopReason::Stop;
2195 }
2196 let (sink, stream) = kiss_ai::EventStream::channel();
2197 sink.send(kiss_ai::AssistantEvent::Start { partial: message.clone() });
2198 if message.stop_reason == StopReason::Error {
2199 sink.error(message);
2200 } else {
2201 sink.done(message);
2202 }
2203 stream
2204 })));
2205 session
2206 .prompt(vec![AgentMessage::user("manage context")])
2207 .await;
2208 assert_eq!(calls.load(std::sync::atomic::Ordering::SeqCst), 5);
2209 let context = session
2210 .manager
2211 .lock()
2212 .unwrap()
2213 .build_session_context()
2214 .messages;
2215 assert_eq!(
2216 SessionManager::open(&session_path)
2217 .unwrap()
2218 .build_session_context()
2219 .messages,
2220 context
2221 );
2222 assert!(
2223 std::fs::read_to_string(session_path)
2224 .unwrap()
2225 .contains("obsolete output")
2226 );
2227 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();
2228 assert_eq!(errors, 1);
2229 }
2230
2231 #[test]
2232 fn context_file_lifecycle_is_opt_in_and_session_local() {
2233 let session = AgentSession::new(
2234 SessionManager::in_memory(Path::new("/test")),
2235 Vec::new(),
2236 Registry::load(None),
2237 Settings::default(),
2238 "test".into(),
2239 openai_model(),
2240 ThinkingLevel::Off,
2241 None,
2242 Arc::new(|_| {}),
2243 );
2244 assert!(!session.sync_context_file());
2245 assert!(session.context_file.lock().unwrap().is_none());
2246 assert_eq!(session.build_context().system_prompt, "test");
2247 let mut settings = session.settings();
2248 settings.experimental_context_file = true;
2249 session.update_settings(settings);
2250 session.sync_context_file();
2251 let path = session
2252 .context_file
2253 .lock()
2254 .unwrap()
2255 .as_ref()
2256 .unwrap()
2257 .path()
2258 .to_path_buf();
2259 let child = session
2260 .create_subagent_session("task", "/root/task", ForkTurns::None, None, None)
2261 .unwrap();
2262 child.sync_context_file();
2263 let child_path = child
2264 .context_file
2265 .lock()
2266 .unwrap()
2267 .as_ref()
2268 .unwrap()
2269 .path()
2270 .to_path_buf();
2271 assert_ne!(path, child_path);
2272 assert!(child_path.is_file());
2273 session.replace_manager(SessionManager::in_memory(Path::new("/other")));
2274 assert!(!path.exists());
2275 assert!(session.context_file.lock().unwrap().is_none());
2276 session.sync_context_file();
2277 let new_path = session
2278 .context_file
2279 .lock()
2280 .unwrap()
2281 .as_ref()
2282 .unwrap()
2283 .path()
2284 .to_path_buf();
2285 assert_ne!(new_path, child_path);
2286 let mut settings = session.settings();
2287 settings.experimental_context_file = false;
2288 session.update_settings(settings);
2289 assert!(session.sync_context_file());
2290 assert!(!new_path.exists());
2291 assert_eq!(session.build_context().system_prompt, "test");
2292 }
2293
2294 #[test]
2295 fn reasoning_lease_ends_when_model_saved_effort_or_applied_effort_changes() {
2296 let model = openai_model();
2297 let mut run = ReasoningRunState {
2298 lease: crate::jev::ReasoningLease::default(),
2299 provider: model.provider.clone(),
2300 model_id: model.id.clone(),
2301 saved_effort: ThinkingLevel::Medium,
2302 };
2303 let selection = crate::jev::ReasoningSelection {
2304 level: ThinkingLevel::High,
2305 generations: 5,
2306 };
2307 run.lease.install(&selection);
2308 assert_eq!(
2309 run.checkpoint(&model, ThinkingLevel::Medium, ThinkingLevel::High, true),
2310 (false, Some(ThinkingLevel::High))
2311 );
2312
2313 let mut other_model = model.clone();
2314 other_model.provider = "other".into();
2315 assert_eq!(
2316 run.checkpoint(
2317 &other_model,
2318 ThinkingLevel::Medium,
2319 ThinkingLevel::High,
2320 true
2321 ),
2322 (true, None)
2323 );
2324 run.lease.install(&selection);
2325 assert_eq!(
2326 run.checkpoint(&other_model, ThinkingLevel::Low, ThinkingLevel::High, true),
2327 (false, None)
2328 );
2329 run.lease.install(&selection);
2330 assert_eq!(
2331 run.checkpoint(&other_model, ThinkingLevel::Low, ThinkingLevel::Low, true),
2332 (false, None)
2333 );
2334 run.lease.install(&selection);
2335 assert_eq!(
2336 run.checkpoint(&other_model, ThinkingLevel::Low, ThinkingLevel::High, false),
2337 (false, None)
2338 );
2339 }
2340
2341 #[test]
2342 fn late_cache_refreshes_are_skipped_before_the_cache_expires() {
2343 let ttl = std::time::Duration::from_secs(300);
2344 let delay = cache_warming_delay(ttl).unwrap();
2345 assert!(!cache_refresh_deadline_missed(
2346 std::time::Duration::from_secs(284),
2347 ttl,
2348 delay,
2349 ));
2350 assert!(cache_refresh_deadline_missed(
2351 std::time::Duration::from_secs(286),
2352 ttl,
2353 delay,
2354 ));
2355 }
2356
2357 #[test]
2358 fn transient_errors_require_a_status_or_specific_network_failure() {
2359 assert!(is_transient("request failed with HTTP 503"));
2360 assert!(is_transient("Cloudflare returned HTTP 520"));
2361 assert!(is_transient("Azure is currently experiencing high demand"));
2362 assert!(is_transient("connection reset by peer"));
2363 assert!(is_transient("rate limit exceeded"));
2364 assert!(!is_transient("model has a 500 token limit"));
2365 assert!(!is_transient("connection settings are invalid"));
2366 assert!(!is_transient(
2367 "HTTP 429: subscription_sharing_usage_limit_exceeded"
2368 ));
2369 assert!(is_transient("subscription_sharing_usage_unavailable"));
2370 assert!(is_transient("subscription_sharing_user_unavailable"));
2371 }
2372
2373 #[test]
2374 fn session_tools_install_replace_and_remove_without_changing_base_tools() {
2375 let registry = Registry::load(None);
2376 let session = AgentSession::new(
2377 SessionManager::in_memory(std::path::Path::new("/test")),
2378 Vec::new(),
2379 registry,
2380 Settings::default(),
2381 "test".into(),
2382 openai_model(),
2383 ThinkingLevel::Off,
2384 None,
2385 Arc::new(|_| {}),
2386 );
2387 let tool = || {
2388 Arc::new(crate::tools::grep::GrepTool {
2389 cwd: std::path::PathBuf::from("/test"),
2390 }) as DynTool
2391 };
2392
2393 assert!(!session.available_tool_names().contains(&"grep".into()));
2394 session.install_session_tool(tool());
2395 session.install_session_tool(tool());
2396 assert_eq!(
2397 session
2398 .available_tool_names()
2399 .iter()
2400 .filter(|name| name.as_str() == "grep")
2401 .count(),
2402 1
2403 );
2404 assert!(session.remove_session_tool("grep"));
2405 assert!(!session.remove_session_tool("grep"));
2406 assert!(!session.available_tool_names().contains(&"grep".into()));
2407 }
2408
2409 fn remote_result() -> kiss_ai::api::openai_compaction::RemoteCompactionResult {
2410 kiss_ai::api::openai_compaction::RemoteCompactionResult {
2411 replacement_history: vec![serde_json::json!({
2412 "type": "compaction",
2413 "encrypted_content": "opaque"
2414 })],
2415 usage: Some(Usage {
2416 input: 10,
2417 output: 2,
2418 total_tokens: 12,
2419 ..Default::default()
2420 }),
2421 }
2422 }
2423
2424 fn settings_test_session(settings: Settings, subagents_allowed: bool) -> Arc<AgentSession> {
2425 let registry = Registry::from_builtin();
2426 let model = registry.all().first().expect("built-in model").clone();
2427 AgentSession::new_with_subagents_allowed(
2428 SessionManager::in_memory(std::path::Path::new("/test")),
2429 Vec::new(),
2430 registry,
2431 settings,
2432 "root prompt".into(),
2433 model,
2434 ThinkingLevel::Off,
2435 None,
2436 Arc::new(|_| {}),
2437 subagents_allowed,
2438 )
2439 }
2440
2441 fn benchmark_tools() -> Vec<DynTool> {
2442 let cwd = std::path::PathBuf::from("/synthetic");
2443 vec![
2444 Arc::new(kiss_agent::tools::read::ReadTool { cwd: cwd.clone() }),
2445 Arc::new(kiss_agent::tools::write::WriteTool { cwd: cwd.clone() }),
2446 Arc::new(kiss_agent::tools::edit::EditTool { cwd: cwd.clone() }),
2447 Arc::new(kiss_agent::tools::bash::BashTool::new(cwd)),
2448 ]
2449 }
2450
2451 #[test]
2452 fn subagent_tools_follow_settings_and_command_line_authority() {
2453 let session = settings_test_session(Settings::default(), true);
2454 assert!(session.available_tool_names().is_empty());
2455 assert!(
2456 !session
2457 .build_context()
2458 .system_prompt
2459 .contains("Subagent coordination")
2460 );
2461
2462 let mut enabled = session.settings();
2463 enabled.subagents.enabled = true;
2464 session.update_settings(enabled.clone());
2465 assert_eq!(
2466 session.available_tool_names(),
2467 [
2468 "spawn_agent",
2469 "send_message",
2470 "followup_task",
2471 "wait_agent",
2472 "list_agents",
2473 "interrupt_agent"
2474 ]
2475 );
2476 assert!(
2477 session
2478 .build_context()
2479 .system_prompt
2480 .contains("Subagent coordination")
2481 );
2482
2483 enabled.subagents.enabled = false;
2484 session.update_settings(enabled);
2485 assert!(session.available_tool_names().is_empty());
2486
2487 let mut blocked_settings = Settings::default();
2488 blocked_settings.subagents.enabled = true;
2489 let blocked = settings_test_session(blocked_settings, false);
2490 assert!(blocked.available_tool_names().is_empty());
2491 assert!(
2492 !blocked
2493 .build_context()
2494 .system_prompt
2495 .contains("Subagent coordination")
2496 );
2497 }
2498
2499 #[test]
2500 fn the_workflow_tool_appears_only_in_workflow_prompt_mode() {
2501 let mut settings = Settings::default();
2502 settings.subagents.enabled = true;
2503 let session = settings_test_session(settings, true);
2504
2505 assert!(
2508 !session
2509 .available_tool_names()
2510 .contains(&"run_workflow".into())
2511 );
2512 assert_eq!(
2513 session.prompt_mode_for("run a dynamic workflow for this task"),
2514 PromptMode::Workflow
2515 );
2516 assert_eq!(
2517 session.prompt_mode_for("fix this small function"),
2518 PromptMode::Ordinary
2519 );
2520 assert!(
2521 !session
2522 .build_context()
2523 .system_prompt
2524 .contains("Writing a dynamic workflow")
2525 );
2526
2527 assert!(
2528 session
2529 .available_tool_names_for(PromptMode::Workflow)
2530 .contains(&"run_workflow".into())
2531 );
2532 assert!(
2533 session
2534 .build_context_for(PromptMode::Workflow)
2535 .system_prompt
2536 .contains("Writing a dynamic workflow")
2537 );
2538 assert!(
2539 !session
2540 .available_tool_names()
2541 .contains(&"run_workflow".into())
2542 );
2543 }
2544
2545 #[test]
2546 fn workflow_prompt_mode_does_nothing_while_subagents_are_off() {
2547 let session = settings_test_session(Settings::default(), true);
2550 assert!(!session.workflows_enabled());
2551 assert!(
2552 !session
2553 .available_tool_names_for(PromptMode::Workflow)
2554 .contains(&"run_workflow".into())
2555 );
2556
2557 let mut settings = session.settings();
2558 settings.subagents.enabled = true;
2559 settings.workflows.enabled = false;
2560 session.update_settings(settings.clone());
2561 assert!(!session.workflows_enabled());
2562 assert!(
2563 !session
2564 .available_tool_names_for(PromptMode::Workflow)
2565 .contains(&"run_workflow".into())
2566 );
2567
2568 settings.workflows.enabled = true;
2569 session.update_settings(settings);
2570 assert!(session.workflows_enabled());
2571 assert!(
2572 session
2573 .available_tool_names_for(PromptMode::Workflow)
2574 .contains(&"run_workflow".into())
2575 );
2576 }
2577
2578 #[test]
2579 fn one_at_a_time_queues_keep_each_prompts_mode() {
2580 let queue = Arc::new(Mutex::new(VecDeque::from([
2581 QueuedPrompt {
2582 message: AgentMessage::user("ordinary"),
2583 mode: PromptMode::Ordinary,
2584 },
2585 QueuedPrompt {
2586 message: AgentMessage::user("workflow"),
2587 mode: PromptMode::Workflow,
2588 },
2589 ])));
2590
2591 assert_eq!(
2592 queued_mode(&queue, QueueMode::OneAtATime),
2593 Some(PromptMode::Ordinary)
2594 );
2595 assert_eq!(drain_queue(&queue, QueueMode::OneAtATime).len(), 1);
2596 assert_eq!(
2597 queued_mode(&queue, QueueMode::OneAtATime),
2598 Some(PromptMode::Workflow)
2599 );
2600 }
2601
2602 #[test]
2603 fn a_session_without_subagent_authority_has_no_workflow_runtime() {
2604 let mut settings = Settings::default();
2605 settings.subagents.enabled = true;
2606 let child = settings_test_session(settings, false);
2607 assert!(child.workflows().is_none());
2608 assert!(!child.workflows_enabled());
2609 }
2610
2611 #[test]
2612 fn child_session_has_safe_forked_context_without_control_tools() {
2613 let mut settings = Settings::default();
2614 settings.subagents.enabled = true;
2615 let parent = settings_test_session(settings, true);
2616 parent
2617 .manager
2618 .lock()
2619 .unwrap()
2620 .append_message(AgentMessage::user("parent context"))
2621 .unwrap();
2622
2623 let child = parent
2624 .create_subagent_session("inspect", "/root/inspect", ForkTurns::All, None, None)
2625 .unwrap();
2626 assert!(Arc::ptr_eq(&parent.registry, &child.registry));
2627 assert!(child.available_tool_names().is_empty());
2628 let context = child.manager.lock().unwrap().build_session_context();
2629 assert!(matches!(
2630 context.messages.as_slice(),
2631 [AgentMessage::User(user)] if user.content.as_text() == "parent context"
2632 ));
2633 }
2634
2635 #[test]
2636 fn child_without_forked_turns_does_not_copy_parent_context() {
2637 let parent = settings_test_session(Settings::default(), true);
2638 parent
2639 .manager
2640 .lock()
2641 .unwrap()
2642 .append_message(AgentMessage::user("parent context"))
2643 .unwrap();
2644
2645 let child = parent
2646 .create_subagent_session("inspect", "/root/inspect", ForkTurns::None, None, None)
2647 .unwrap();
2648 assert!(
2649 child
2650 .manager
2651 .lock()
2652 .unwrap()
2653 .build_session_context()
2654 .messages
2655 .is_empty()
2656 );
2657 }
2658
2659 #[test]
2660 #[ignore = "release-mode performance benchmark"]
2661 fn benchmark_performance_subagent_overhead() {
2662 let registry = Registry::from_builtin();
2663 let model = registry.all().first().expect("built-in model").clone();
2664 let tools = benchmark_tools();
2665 let make_session = |enabled: bool| {
2666 let mut settings = Settings::default();
2667 settings.subagents.enabled = enabled;
2668 AgentSession::new_with_subagents_allowed(
2669 SessionManager::in_memory(std::path::Path::new("/synthetic")),
2670 tools.clone(),
2671 registry.clone(),
2672 settings,
2673 "benchmark root prompt".into(),
2674 model.clone(),
2675 ThinkingLevel::Off,
2676 None,
2677 Arc::new(|_| {}),
2678 true,
2679 )
2680 };
2681
2682 kiss_bench::measure_pair(
2683 (
2684 "agent_session_create_subagents_off",
2685 "agent_session_create_subagents_on",
2686 ),
2687 21,
2688 500,
2689 (
2690 "new_root_session_4_base_tools_0_control_tools",
2691 "new_root_session_4_base_tools_6_control_tools",
2692 ),
2693 || make_session(false),
2694 || make_session(true),
2695 );
2696
2697 let off = make_session(false);
2698 let on = make_session(true);
2699 kiss_bench::measure_pair(
2700 (
2701 "agent_context_build_subagents_off",
2702 "agent_context_build_subagents_on",
2703 ),
2704 21,
2705 10_000,
2706 (
2707 "empty_session_4_base_tools_0_control_tools",
2708 "empty_session_4_base_tools_6_control_tools",
2709 ),
2710 || off.build_context(),
2711 || on.build_context(),
2712 );
2713 }
2714
2715 #[test]
2716 #[ignore = "release-mode performance benchmark"]
2717 fn benchmark_performance_workflow_tool_exposure() {
2718 let registry = Registry::from_builtin();
2722 let model = registry.all().first().expect("built-in model").clone();
2723 let tools = benchmark_tools();
2724 let make_session = || {
2725 let mut settings = Settings::default();
2726 settings.subagents.enabled = true;
2727 AgentSession::new_with_subagents_allowed(
2728 SessionManager::in_memory(std::path::Path::new("/synthetic")),
2729 tools.clone(),
2730 registry.clone(),
2731 settings,
2732 "benchmark root prompt".into(),
2733 model.clone(),
2734 ThinkingLevel::Off,
2735 None,
2736 Arc::new(|_| {}),
2737 true,
2738 )
2739 };
2740
2741 let ordinary = make_session();
2742 let workflow = make_session();
2743 kiss_bench::measure_pair(
2744 (
2745 "agent_context_build_workflow_disarmed",
2746 "agent_context_build_workflow_armed",
2747 ),
2748 21,
2749 10_000,
2750 (
2751 "empty_session_subagents_on_workflow_disarmed",
2752 "empty_session_subagents_on_workflow_armed",
2753 ),
2754 || ordinary.build_context_for(PromptMode::Ordinary),
2755 || workflow.build_context_for(PromptMode::Workflow),
2756 );
2757 }
2758
2759 #[test]
2760 fn transcript_excerpt_keeps_only_recent_user_and_assistant_text() {
2761 let messages = vec![
2762 AgentMessage::user("old"),
2763 AgentMessage::BashExecution(kiss_agent::BashExecutionMessage {
2764 command: "pwd".into(),
2765 output: "ignored".into(),
2766 exit_code: Some(0),
2767 cancelled: false,
2768 truncated: false,
2769 full_output_path: None,
2770 exclude_from_context: false,
2771 timestamp: 1,
2772 }),
2773 AgentMessage::user("new"),
2774 ];
2775 let excerpt = transcript_excerpt(&messages, 1, 100);
2776 assert_eq!(excerpt, "User: new");
2777 }
2778
2779 #[test]
2780 fn transcript_excerpt_enforces_character_budget_from_the_tail() {
2781 let excerpt = transcript_excerpt(&[AgentMessage::user("abcdefghij")], 4, 5);
2782 assert!(excerpt.ends_with("fghij"));
2783 assert!(excerpt.starts_with("[earlier text omitted]"));
2784 }
2785
2786 #[test]
2787 fn hybrid_compaction_keeps_local_summary_and_remote_details() {
2788 let selected = select_compaction_outcome(
2789 &openai_model(),
2790 Ok(compaction::SummaryOutcome {
2791 summary: "portable".into(),
2792 usage: None,
2793 }),
2794 Some(Ok(remote_result())),
2795 )
2796 .unwrap();
2797 assert_eq!(selected.summary, "portable");
2798 assert_eq!(selected.remote_usage.unwrap().input, 10);
2799 assert_eq!(
2800 selected.remote_details.unwrap()["remoteCompaction"]["version"],
2801 2
2802 );
2803 }
2804
2805 #[test]
2806 fn remote_failure_falls_back_to_local_compaction() {
2807 let selected = select_compaction_outcome(
2808 &openai_model(),
2809 Ok(compaction::SummaryOutcome {
2810 summary: "portable".into(),
2811 usage: None,
2812 }),
2813 Some(Err(anyhow::anyhow!("remote unavailable"))),
2814 )
2815 .unwrap();
2816 assert_eq!(selected.summary, "portable");
2817 assert!(selected.remote_details.is_none());
2818 }
2819
2820 #[test]
2821 fn remote_success_survives_local_summary_failure() {
2822 let selected = select_compaction_outcome(
2823 &openai_model(),
2824 Err(anyhow::anyhow!("summary unavailable")),
2825 Some(Ok(remote_result())),
2826 )
2827 .unwrap();
2828 assert!(
2829 selected
2830 .summary
2831 .contains("server-side compaction was applied")
2832 );
2833 assert!(selected.remote_details.is_some());
2834 }
2835
2836 #[test]
2837 fn details_merge_keeps_file_operations_and_remote_artifact() {
2838 let merged = merge_compaction_details(
2839 serde_json::json!({"readFiles": ["a.rs"], "modifiedFiles": []}),
2840 Some(serde_json::json!({"remoteCompaction": {"version": 2}})),
2841 );
2842 assert_eq!(merged["readFiles"][0], "a.rs");
2843 assert_eq!(merged["remoteCompaction"]["version"], 2);
2844 }
2845
2846 #[test]
2847 fn auto_compaction_guard_checks_settings_threshold_and_cancel() {
2848 let mut settings = Settings::default();
2849 settings.compaction.reserve_tokens = 20;
2850 let messages = vec![AgentMessage::user("x".repeat(360))];
2851 let mut model = openai_model();
2852 model.context_window = 100;
2853 assert!(auto_compaction_needed(&settings, &messages, &model, false));
2854 assert!(!auto_compaction_needed(&settings, &messages, &model, true));
2855 settings.compaction.enabled = false;
2856 assert!(!auto_compaction_needed(&settings, &messages, &model, false));
2857 settings.compaction.enabled = true;
2858 let mut assistant = kiss_ai::AssistantMessage::empty("test", "test", "test");
2859 assistant.usage.input = 1000;
2860 assistant
2861 .content
2862 .push(kiss_ai::ContentBlock::text("short note"));
2863 let edited = vec![
2864 AgentMessage::user("task"),
2865 AgentMessage::Assistant(assistant),
2866 ];
2867 assert!(auto_compaction_needed(&settings, &edited, &model, false));
2868 settings.experimental_context_file = true;
2869 assert!(!auto_compaction_needed(&settings, &edited, &model, false));
2870 }
2871
2872 #[test]
2873 fn session_title_normalization_is_safe_and_bounded() {
2874 assert_eq!(
2875 normalize_session_title(" `Fix AUTH-123 login flow!` \nignored").as_deref(),
2876 Some("Fix AUTH-123 login flow")
2877 );
2878 assert_eq!(normalize_session_title("\n\t"), None);
2879 assert_eq!(
2880 normalize_session_title("🚀".repeat(50).as_str())
2881 .unwrap()
2882 .chars()
2883 .count(),
2884 SESSION_TITLE_MAX_CHARS
2885 );
2886 }
2887
2888 #[test]
2889 fn session_title_prompt_is_utf8_safe_and_bounded() {
2890 let prompt = "🚀".repeat(SESSION_TITLE_PROMPT_MAX_BYTES);
2891 let bounded = bounded_session_title_prompt(&prompt);
2892 assert!(bounded.len() <= SESSION_TITLE_PROMPT_MAX_BYTES);
2893 assert!(std::str::from_utf8(bounded.as_bytes()).is_ok());
2894 }
2895
2896 #[test]
2897 fn cache_warming_uses_ninety_percent_with_ten_second_margin() {
2898 assert_eq!(
2899 cache_warming_delay(std::time::Duration::from_secs(300)),
2900 Some(std::time::Duration::from_secs(270))
2901 );
2902 assert_eq!(
2903 cache_warming_delay(std::time::Duration::from_secs(60)),
2904 Some(std::time::Duration::from_secs(50))
2905 );
2906 assert_eq!(
2907 cache_warming_delay(std::time::Duration::from_secs(10)),
2908 None
2909 );
2910 }
2911}