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