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 mut end = prompt.len().min(SESSION_TITLE_PROMPT_MAX_BYTES);
123 while !prompt.is_char_boundary(end) {
124 end -= 1;
125 }
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 api_key_override: Option<(String, String)>,
167 sink: SessionEventSink,
168 subagents_allowed: bool,
169 subagents: OnceLock<Arc<SubagentRuntime>>,
170 workflows: OnceLock<Arc<WorkflowRuntime>>,
171 iterative: OnceLock<Arc<IterativeRuntime>>,
172 workflow_approver: Mutex<Option<WorkflowApprover>>,
173 stream_fn: Mutex<Option<StreamFn>>,
176}
177
178struct ReasoningRunState {
179 lease: crate::jev::ReasoningLease,
180 provider: String,
181 model_id: String,
182 saved_effort: ThinkingLevel,
183}
184
185impl ReasoningRunState {
186 fn checkpoint(
187 &mut self,
188 model: &Model,
189 saved_effort: ThinkingLevel,
190 current_effort: ThinkingLevel,
191 dynamic_enabled: bool,
192 ) -> (bool, Option<ThinkingLevel>) {
193 let model_changed = self.provider != model.provider || self.model_id != model.id;
194 if model_changed || self.saved_effort != saved_effort || !dynamic_enabled {
195 self.lease.clear();
196 }
197 self.provider.clone_from(&model.provider);
198 self.model_id.clone_from(&model.id);
199 self.saved_effort = saved_effort;
200 if self
201 .lease
202 .current()
203 .is_some_and(|level| level != current_effort)
204 {
205 self.lease.clear();
206 }
207 (model_changed, self.lease.current())
208 }
209}
210
211impl AgentSession {
212 pub(crate) fn emit_iterative(&self, job: crate::iterative::JobId, version: u64) {
213 (self.sink)(SessionEvent::Iterative { job, version });
214 }
215
216 pub(crate) fn emit_workflow(&self, run: crate::workflows::RunId, version: u64) {
217 (self.sink)(SessionEvent::Workflow { run, version });
218 }
219
220 pub(crate) fn emit_workflow_outcome(
221 &self,
222 run: Option<crate::workflows::RunId>,
223 name: String,
224 status: WorkflowTurnStatus,
225 ) {
226 (self.sink)(SessionEvent::WorkflowOutcome { run, name, status });
227 }
228
229 #[allow(clippy::too_many_arguments)]
230 pub fn new(
231 manager: SessionManager,
232 tools: Vec<DynTool>,
233 registry: impl Into<Arc<Registry>>,
234 settings: Settings,
235 system_prompt: String,
236 model: Model,
237 thinking: ThinkingLevel,
238 api_key_override: Option<(String, String)>,
239 sink: SessionEventSink,
240 ) -> Arc<Self> {
241 Self::new_with_subagents_allowed(
242 manager,
243 tools,
244 registry,
245 settings,
246 system_prompt,
247 model,
248 thinking,
249 api_key_override,
250 sink,
251 true,
252 )
253 }
254
255 #[allow(clippy::too_many_arguments)]
256 pub fn new_with_subagents_allowed(
257 manager: SessionManager,
258 tools: Vec<DynTool>,
259 registry: impl Into<Arc<Registry>>,
260 settings: Settings,
261 system_prompt: String,
262 model: Model,
263 thinking: ThinkingLevel,
264 api_key_override: Option<(String, String)>,
265 sink: SessionEventSink,
266 subagents_allowed: bool,
267 ) -> Arc<Self> {
268 let totals = manager.usage_totals();
269 let session = Arc::new(AgentSession {
270 manager: Mutex::new(manager),
271 registry: registry.into(),
272 base_tools: Mutex::new(tools.clone()),
273 session_tools: Mutex::new(Vec::new()),
274 tools: Mutex::new(tools),
275 settings: Mutex::new(settings),
276 system_prompt: Mutex::new(system_prompt),
277 model: Mutex::new(model),
278 thinking: Mutex::new(thinking),
279 fast_mode: Mutex::new(false),
280 steering: Default::default(),
281 follow_up: Default::default(),
282 cancel: Mutex::new(CancellationToken::new()),
283 cache_warm_cancel: Mutex::new(CancellationToken::new()),
284 running: Mutex::new(false),
285 totals: Mutex::new(totals),
286 context_usage_cache: Default::default(),
287 api_key_override,
288 sink,
289 subagents_allowed,
290 subagents: OnceLock::new(),
291 workflows: OnceLock::new(),
292 iterative: OnceLock::new(),
293 workflow_approver: Mutex::new(None),
294 stream_fn: Mutex::new(None),
295 });
296 if subagents_allowed {
297 let runtime = SubagentRuntime::new(Arc::downgrade(&session));
298 assert!(session.subagents.set(runtime).is_ok());
299 let workflows = WorkflowRuntime::new(Arc::downgrade(&session));
300 assert!(session.workflows.set(workflows).is_ok());
301 let iterative = IterativeRuntime::new(Arc::downgrade(&session));
302 assert!(session.iterative.set(iterative).is_ok());
303 }
304 session.rebuild_tools();
305 session
306 }
307
308 pub fn set_stream_fn(&self, stream_fn: Option<StreamFn>) {
314 *self.stream_fn.lock().unwrap() = stream_fn;
315 }
316
317 pub fn model(&self) -> Model {
318 self.model.lock().unwrap().clone()
319 }
320
321 pub fn thinking_level(&self) -> ThinkingLevel {
322 *self.thinking.lock().unwrap()
323 }
324
325 pub fn fast_mode(&self) -> bool {
326 *self.fast_mode.lock().unwrap()
327 }
328
329 pub fn set_fast_mode(&self, enabled: bool) {
330 *self.fast_mode.lock().unwrap() = enabled;
331 }
332
333 pub fn totals(&self) -> Usage {
334 *self.totals.lock().unwrap()
335 }
336
337 pub(crate) fn record_subagent_usage(&self, usage: Usage) {
338 self.totals.lock().unwrap().add(&usage);
339 }
340
341 pub fn is_running(&self) -> bool {
342 *self.running.lock().unwrap()
343 }
344
345 pub fn settings(&self) -> Settings {
346 self.settings.lock().unwrap().clone()
347 }
348
349 pub fn update_settings(&self, settings: Settings) {
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 let was_enabled = self.subagents_enabled();
371 let workflows_were_enabled = self.workflows_enabled();
372 let cache_warming = settings.cache_warming;
373 *self.settings.lock().unwrap() = settings;
374 *self.system_prompt.lock().unwrap() = system_prompt;
375 *self.base_tools.lock().unwrap() = tools;
376 if cache_warming == CacheWarmingMode::Off
377 || (cache_warming == CacheWarmingMode::Streaming && !self.is_running())
378 {
379 self.cache_warm_cancel.lock().unwrap().cancel();
380 }
381 let is_enabled = self.subagents_enabled();
382 if was_enabled && !is_enabled {
383 self.stop_child_work();
384 } else if workflows_were_enabled && !self.workflows_enabled() {
385 self.stop_workflows();
386 }
387 self.rebuild_tools();
388 }
389
390 pub fn install_session_tool(&self, tool: DynTool) {
392 let mut tools = self.session_tools.lock().unwrap();
393 if let Some(existing) = tools
394 .iter_mut()
395 .find(|existing| existing.name() == tool.name())
396 {
397 *existing = tool;
398 } else {
399 tools.push(tool);
400 }
401 drop(tools);
402 self.rebuild_tools();
403 }
404
405 pub fn remove_session_tool(&self, name: &str) -> bool {
407 let mut tools = self.session_tools.lock().unwrap();
408 let previous_len = tools.len();
409 tools.retain(|tool| tool.name() != name);
410 let removed = tools.len() != previous_len;
411 drop(tools);
412 if removed {
413 self.rebuild_tools();
414 }
415 removed
416 }
417
418 fn stop_child_work(&self) {
420 if let Some(runtime) = self.subagents.get() {
421 runtime.interrupt_all();
422 }
423 if let Some(runtime) = self.workflows.get() {
424 runtime.stop_all();
425 }
426 if let Some(runtime) = self.iterative.get() {
427 runtime.stop_all();
428 }
429 }
430
431 fn stop_workflows(&self) {
432 if let Some(runtime) = self.workflows.get() {
433 runtime.stop_all();
434 }
435 }
436
437 fn subagents_enabled(&self) -> bool {
438 self.subagents_allowed && self.settings.lock().unwrap().subagents.enabled
439 }
440
441 pub fn workflows_enabled(&self) -> bool {
444 self.subagents_enabled() && self.settings.lock().unwrap().workflows.enabled
445 }
446
447 pub fn prompt_mode_for(&self, text: &str) -> PromptMode {
449 if self.workflows_enabled()
450 && self.settings.lock().unwrap().workflows.keyword_trigger
451 && crate::workflows::workflow_trigger(text).is_some()
452 {
453 PromptMode::Workflow
454 } else {
455 PromptMode::Ordinary
456 }
457 }
458
459 pub fn workflows(&self) -> Option<Arc<WorkflowRuntime>> {
460 self.workflows.get().cloned()
461 }
462
463 pub fn iterative_jobs(&self) -> Option<Arc<IterativeRuntime>> {
464 self.iterative.get().cloned()
465 }
466
467 pub fn set_workflow_approver(&self, approver: WorkflowApprover) {
469 *self.workflow_approver.lock().unwrap() = Some(approver);
470 }
471
472 pub(crate) fn workflow_approver(&self) -> Option<WorkflowApprover> {
473 self.workflow_approver.lock().unwrap().clone()
474 }
475
476 fn rebuild_tools(&self) {
477 let mut tools = self.base_tools.lock().unwrap().clone();
478 for session_tool in self.session_tools.lock().unwrap().iter().cloned() {
479 if let Some(existing) = tools
480 .iter_mut()
481 .find(|existing| existing.name() == session_tool.name())
482 {
483 *existing = session_tool;
484 } else {
485 tools.push(session_tool);
486 }
487 }
488 if self.subagents_enabled()
489 && let Some(runtime) = self.subagents.get()
490 {
491 tools.extend(runtime.control_tools());
492 }
493 *self.tools.lock().unwrap() = tools;
494 }
495
496 fn tools_for(&self, mode: PromptMode) -> Vec<DynTool> {
497 let mut tools = self.tools.lock().unwrap().clone();
498 if mode == PromptMode::Workflow
499 && self.workflows_enabled()
500 && let Some(runtime) = self.workflows.get()
501 {
502 tools.push(runtime.tool());
503 }
504 tools
505 }
506
507 pub fn available_tool_names(&self) -> Vec<String> {
508 self.tools
509 .lock()
510 .unwrap()
511 .iter()
512 .map(|tool| tool.name().to_string())
513 .collect()
514 }
515
516 #[cfg(test)]
517 fn available_tool_names_for(&self, mode: PromptMode) -> Vec<String> {
518 self.tools_for(mode)
519 .iter()
520 .map(|tool| tool.name().to_string())
521 .collect()
522 }
523
524 pub fn replace_manager(&self, manager: SessionManager) {
526 self.cache_warm_cancel.lock().unwrap().cancel();
527 if let Some(runtime) = self.subagents.get() {
528 runtime.reset();
529 }
530 if let Some(runtime) = self.workflows.get() {
531 runtime.stop_all();
532 }
533 if let Some(runtime) = self.iterative.get() {
534 runtime.stop_all();
535 }
536 let totals = manager.usage_totals();
537 let context = manager.build_session_context();
538 if let Some((provider, model_id)) = context.model
539 && let Some((model, _)) = self.registry.resolve(&model_id, Some(&provider))
540 {
541 *self.model.lock().unwrap() = model;
542 }
543 if let Some(thinking) = context.thinking_level {
544 *self.thinking.lock().unwrap() = thinking;
545 }
546 *self.manager.lock().unwrap() = manager;
547 *self.context_usage_cache.lock().unwrap() = None;
548 *self.totals.lock().unwrap() = totals;
549 }
550
551 pub fn set_model(&self, model: Model) {
552 self.cache_warm_cancel.lock().unwrap().cancel();
553 {
554 let mut m = self.manager.lock().unwrap();
555 let _ = m.append_model_change(&model.provider, &model.id);
556 }
557 (self.sink)(SessionEvent::ModelChanged {
558 provider: model.provider.clone(),
559 model_id: model.id.clone(),
560 });
561 *self.model.lock().unwrap() = model;
562 }
563
564 pub fn set_thinking_level(&self, level: ThinkingLevel) {
565 let _ = self
566 .manager
567 .lock()
568 .unwrap()
569 .append_thinking_level_change(level);
570 *self.thinking.lock().unwrap() = level;
571 }
572
573 pub async fn navigate_tree(
576 self: &Arc<Self>,
577 target_id: &str,
578 summarize: bool,
579 custom_instructions: Option<String>,
580 cancel: CancellationToken,
581 ) -> anyhow::Result<TreeNavigationOutcome> {
582 let (old_leaf, new_leaf, editor_text, abandoned_messages) = {
583 let manager = self.manager.lock().unwrap();
584 let target = manager
585 .get_entry(target_id)
586 .cloned()
587 .ok_or_else(|| anyhow::anyhow!("unknown tree entry {target_id}"))?;
588 let old_leaf = manager.leaf_id().map(str::to_string);
589 if old_leaf.as_deref() == Some(target_id) {
590 return Ok(TreeNavigationOutcome {
591 editor_text: None,
592 summarized: false,
593 cancelled: false,
594 });
595 }
596
597 let (new_leaf, editor_text) = match &target {
598 crate::session::entry::SessionEntry::Message {
599 message: AgentMessage::User(user),
600 ..
601 } => (
602 target.parent_id().map(str::to_string),
603 Some(user.content.as_text()),
604 ),
605 _ => (Some(target_id.to_string()), None),
606 };
607
608 let target_ancestors: std::collections::HashSet<String> = new_leaf
609 .as_deref()
610 .map(|leaf| {
611 manager
612 .branch_entries(Some(leaf))
613 .into_iter()
614 .map(|entry| entry.id().to_string())
615 .collect()
616 })
617 .unwrap_or_default();
618 let common_ancestor = old_leaf.as_deref().and_then(|leaf| {
619 manager
620 .branch_entries(Some(leaf))
621 .into_iter()
622 .rev()
623 .find(|entry| target_ancestors.contains(entry.id()))
624 .map(|entry| entry.id().to_string())
625 });
626 let abandoned_messages = old_leaf
627 .as_deref()
628 .map(|leaf| manager.branch_messages_after(leaf, common_ancestor.as_deref()))
629 .unwrap_or_default();
630 (old_leaf, new_leaf, editor_text, abandoned_messages)
631 };
632
633 let summary = if summarize && !abandoned_messages.is_empty() {
634 let model = self.model();
635 let credential = self.resolve_credential(&model.provider).await;
636 let serialized = compaction::serialize_agent_messages(&abandoned_messages);
637 let result = compaction::generate_summary(
638 &model,
639 credential,
640 &serialized,
641 None,
642 custom_instructions.as_deref(),
643 false,
644 cancel.clone(),
645 )
646 .await?;
647 if cancel.is_cancelled() {
648 return Ok(TreeNavigationOutcome {
649 editor_text: None,
650 summarized: false,
651 cancelled: true,
652 });
653 }
654 Some(result)
655 } else {
656 None
657 };
658
659 let mut manager = self.manager.lock().unwrap();
660 if manager.leaf_id() != old_leaf.as_deref() {
661 anyhow::bail!("the session tree changed during navigation");
662 }
663 let summarized = if let Some(summary) = summary {
664 let (read, modified) = extract_file_ops(&abandoned_messages);
665 if let Some(usage) = &summary.usage {
666 self.totals.lock().unwrap().add(usage);
667 }
668 manager.branch_with_summary(
669 new_leaf.as_deref(),
670 old_leaf.as_deref().unwrap_or(target_id),
671 summary.summary,
672 summary.usage,
673 Some(file_ops_details(&read, &modified)),
674 )?;
675 true
676 } else {
677 if let Some(new_leaf) = new_leaf.as_deref() {
678 manager.branch(new_leaf)?;
679 } else {
680 manager.reset_leaf();
681 }
682 false
683 };
684
685 Ok(TreeNavigationOutcome {
686 editor_text,
687 summarized,
688 cancelled: false,
689 })
690 }
691
692 pub fn queue_steering(&self, message: AgentMessage) {
693 self.queue_steering_with_mode(message, PromptMode::Ordinary);
694 }
695
696 pub fn queue_steering_with_mode(&self, message: AgentMessage, mode: PromptMode) {
697 self.steering
698 .lock()
699 .unwrap()
700 .push_back(QueuedPrompt { message, mode });
701 self.emit_queues();
702 }
703
704 pub fn queue_follow_up(&self, message: AgentMessage) {
705 self.queue_follow_up_with_mode(message, PromptMode::Ordinary);
706 }
707
708 pub fn queue_follow_up_with_mode(&self, message: AgentMessage, mode: PromptMode) {
709 self.follow_up
710 .lock()
711 .unwrap()
712 .push_back(QueuedPrompt { message, mode });
713 self.emit_queues();
714 }
715
716 pub fn reclaim_queued(&self) -> Vec<AgentMessage> {
718 let mut out: Vec<AgentMessage> = self
719 .steering
720 .lock()
721 .unwrap()
722 .drain(..)
723 .map(|prompt| prompt.message)
724 .collect();
725 out.extend(
726 self.follow_up
727 .lock()
728 .unwrap()
729 .drain(..)
730 .map(|prompt| prompt.message),
731 );
732 self.emit_queues();
733 out
734 }
735
736 pub fn abort(&self) {
737 self.cancel.lock().unwrap().cancel();
738 }
739
740 fn emit_queues(&self) {
741 let preview = |q: &VecDeque<QueuedPrompt>| {
742 q.iter()
743 .map(|prompt| match &prompt.message {
744 AgentMessage::User(u) => u.content.as_text().chars().take(80).collect(),
745 other => other.role().to_string(),
746 })
747 .collect::<Vec<String>>()
748 };
749 (self.sink)(SessionEvent::QueueUpdate {
750 steering: preview(&self.steering.lock().unwrap()),
751 follow_up: preview(&self.follow_up.lock().unwrap()),
752 });
753 }
754
755 async fn resolve_credential(&self, provider: &str) -> Option<kiss_ai::ResolvedCredential> {
756 if let Some((override_provider, key)) = &self.api_key_override
757 && override_provider == provider
758 {
759 return Some(kiss_ai::ResolvedCredential::api_key(key));
760 }
761 let credential_provider = self.registry.credential_provider(provider);
762 kiss_ai::auth::resolve_credential_async(credential_provider, &self.registry.declared_keys)
763 .await
764 .ok()
765 .flatten()
766 }
767
768 fn loop_config(
769 &self,
770 session_arc: &Arc<Self>,
771 active_prompt_mode: Arc<Mutex<PromptMode>>,
772 ) -> AgentLoopConfig {
773 let mut config = AgentLoopConfig::new(self.model());
774 config.thinking_level = self.thinking_level();
775 config.fast_mode = self.fast_mode();
776 config.session_id = Some(self.manager.lock().unwrap().session_id().to_string());
777 let settings = self.settings();
778 config.transport = settings.transport;
779 let registry = self.registry.clone();
780 let declared = self.registry.declared_keys.clone();
781 let api_key_override = self.api_key_override.clone();
782 config.get_credential = Some(Arc::new(move |provider| {
783 let registry = registry.clone();
784 let declared = declared.clone();
785 let api_key_override = api_key_override.clone();
786 Box::pin(async move {
787 if let Some((override_provider, key)) = api_key_override
788 && override_provider == provider
789 {
790 return Some(kiss_ai::ResolvedCredential::api_key(key));
791 }
792 let credential_provider = registry.credential_provider(&provider).to_string();
793 kiss_ai::auth::resolve_credential_async(&credential_provider, &declared)
794 .await
795 .ok()
796 .flatten()
797 })
798 }));
799
800 let steering = self.steering.clone();
801 let steering_for_mode = self.steering.clone();
802 let steering_mode = settings.steering_mode;
803 let session_for_queues = session_arc.clone();
804 config.get_steering_messages = Some(Arc::new(move || {
805 let drained = drain_queue(&steering, steering_mode);
806 session_for_queues.emit_queues();
807 Box::pin(async move { drained })
808 }));
809 let follow_up = self.follow_up.clone();
810 let follow_up_for_mode = self.follow_up.clone();
811 let follow_up_mode = settings.follow_up_mode;
812 let session_for_queues = session_arc.clone();
813 config.get_follow_up_messages = Some(Arc::new(move || {
814 let drained = drain_queue(&follow_up, follow_up_mode);
815 session_for_queues.emit_queues();
816 Box::pin(async move { drained })
817 }));
818 let reasoning_run = Arc::new(Mutex::new(ReasoningRunState {
819 lease: crate::jev::ReasoningLease::default(),
820 provider: config.model.provider.clone(),
821 model_id: config.model.id.clone(),
822 saved_effort: config.thinking_level,
823 }));
824 let session_for_reasoning = session_arc.clone();
825 let reasoning_for_generation = reasoning_run.clone();
826 let prompt_mode_for_reasoning = active_prompt_mode.clone();
827 config.prepare_generation = Some(Arc::new(move |current_effort| {
828 let session = session_for_reasoning.clone();
829 let reasoning_run = reasoning_for_generation.clone();
830 let active_prompt_mode = prompt_mode_for_reasoning.clone();
831 Box::pin(async move {
832 let settings = session.settings();
833 let model = session.model();
834 let saved_effort = session.thinking_level();
835 let dynamic_enabled =
836 settings.reasoning_effort.mode == ReasoningEffortMode::Jev && model.reasoning;
837 let (model_changed, leased) = {
838 let mut run = reasoning_run.lock().unwrap();
839 run.checkpoint(&model, saved_effort, current_effort, dynamic_enabled)
840 };
841 let mut update = TurnUpdate {
842 context: model_changed
843 .then(|| session.build_context_for(*active_prompt_mode.lock().unwrap())),
844 model: model_changed.then(|| model.clone()),
845 thinking_level: Some(saved_effort),
846 fast_mode: Some(session.fast_mode()),
847 ..Default::default()
848 };
849 if !dynamic_enabled {
850 return Some(update);
851 }
852 if let Some(level) = leased {
853 update.thinking_level = Some(level);
854 return Some(update);
855 }
856
857 let selected = match kiss_ai::auth::resolve_api_key_async(
858 "typesafe",
859 &session.registry.declared_keys,
860 )
861 .await
862 {
863 Ok(Some(api_key)) => {
864 let messages = session
865 .manager
866 .lock()
867 .unwrap()
868 .build_session_context()
869 .messages;
870 let cancel = session.cancel.lock().unwrap().clone();
871 crate::jev::select_reasoning(&messages, &[], &model, &api_key, cancel)
872 .await
873 .map_err(|error| error.to_string())
874 }
875 Ok(None) => Err(
876 "TypeSafe credentials are unavailable; run /login typesafe or set TYPESAFE_API_KEY"
877 .into(),
878 ),
879 Err(error) => Err(format!(
880 "TypeSafe credentials could not be read: {error}; run /login typesafe or set TYPESAFE_API_KEY"
881 )),
882 };
883 let latest_model = session.model();
884 let latest_effort = session.thinking_level();
885 if latest_model.provider != model.provider
886 || latest_model.id != model.id
887 || latest_effort != saved_effort
888 || session.settings().reasoning_effort.mode != ReasoningEffortMode::Jev
889 {
890 reasoning_run.lock().unwrap().lease.clear();
891 update.context =
892 Some(session.build_context_for(*active_prompt_mode.lock().unwrap()));
893 update.model = Some(latest_model);
894 update.thinking_level = Some(latest_effort);
895 return Some(update);
896 }
897 match selected {
898 Ok(selection) => {
899 update.thinking_level = Some(selection.level);
900 if selection.level != current_effort {
901 let sink = session.sink.clone();
902 let level = selection.level;
903 let generations = selection.generations;
904 update.on_applied = Some(Box::new(move |_, applied| {
905 if applied == level {
906 sink(SessionEvent::ReasoningEffortChanged {
907 level,
908 generations,
909 });
910 }
911 }));
912 }
913 reasoning_run.lock().unwrap().lease.install(&selection);
914 }
915 Err(reason) => {
916 let sink = session.sink.clone();
917 let reason = reason.chars().take(300).collect();
918 update.on_applied = Some(Box::new(move |_, applied| {
919 sink(SessionEvent::ReasoningEffortFallback {
920 level: applied,
921 reason,
922 });
923 }));
924 }
925 }
926 Some(update)
927 })
928 }));
929 let session_for_compaction = session_arc.clone();
930 config.prepare_next_turn = Some(Arc::new(move |turn| {
931 let has_tool_results = !turn.tool_results.is_empty();
932 let tool_failed = turn.tool_results.iter().any(|result| result.is_error);
933 let will_continue = turn.will_continue;
934 let session = session_for_compaction.clone();
935 let active_prompt_mode = active_prompt_mode.clone();
936 let steering_for_mode = steering_for_mode.clone();
937 let follow_up_for_mode = follow_up_for_mode.clone();
938 let reasoning_run = reasoning_run.clone();
939 Box::pin(async move {
940 let has_queued_user_input = !steering_for_mode.lock().unwrap().is_empty()
941 || (!will_continue && !follow_up_for_mode.lock().unwrap().is_empty());
942 let queued_mode = queued_mode(&steering_for_mode, steering_mode)
943 .or_else(|| queued_mode(&follow_up_for_mode, follow_up_mode));
944 let mode_changed = {
945 let mut active = active_prompt_mode.lock().unwrap();
946 let before = *active;
947 if let Some(queued_mode) = queued_mode {
948 *active = queued_mode;
949 } else if *active == PromptMode::Workflow && !has_tool_results {
950 *active = PromptMode::Ordinary;
954 }
955 before != *active
956 };
957
958 let mut context_changed = mode_changed;
959 if has_tool_results {
960 let settings = session.settings();
961 let cancel = session.cancel.lock().unwrap().clone();
962 let model = session.model();
963 let revision_before = {
964 let manager = session.manager.lock().unwrap();
965 let context = manager.build_session_context();
966 if !auto_compaction_needed(
967 &settings,
968 &context.messages,
969 &model,
970 cancel.is_cancelled(),
971 ) {
972 None
973 } else {
974 Some(manager.context_revision())
975 }
976 };
977
978 if let Some(revision_before) = revision_before {
979 session.compact(None, true).await;
980 let revision_after = session.manager.lock().unwrap().context_revision();
981 context_changed |= revision_after != revision_before;
982 }
983 }
984
985 {
986 let mut run = reasoning_run.lock().unwrap();
987 run.lease.consume_generation();
988 if tool_failed || has_queued_user_input {
989 run.lease.clear();
990 }
991 }
992
993 Some(TurnUpdate {
994 context: context_changed
995 .then(|| session.build_context_for(*active_prompt_mode.lock().unwrap())),
996 fast_mode: Some(session.fast_mode()),
997 ..Default::default()
998 })
999 })
1000 }));
1001 if let Some(stream_fn) = self.stream_fn.lock().unwrap().clone() {
1002 config.stream_fn = stream_fn;
1003 }
1004 config
1005 }
1006
1007 fn build_context(&self) -> AgentContext {
1008 self.build_context_for(PromptMode::Ordinary)
1009 }
1010
1011 fn build_context_for(&self, prompt_mode: PromptMode) -> AgentContext {
1012 let model = self.model();
1013 let manager = self.manager.lock().unwrap();
1014 let (openai_responses_input, messages) =
1015 if kiss_ai::api::openai_compaction::supports_remote_compaction(&model)
1016 && let Some(remote) = manager.build_openai_compaction_context(&model)
1017 {
1018 (Some(remote.replacement_history), remote.messages)
1019 } else {
1020 (None, manager.build_session_context().messages)
1021 };
1022 drop(manager);
1023 let mut system_prompt = self.system_prompt.lock().unwrap().clone();
1024 if self.subagents_enabled() {
1025 system_prompt.push_str("\n\n");
1026 system_prompt.push_str(SUBAGENT_SYSTEM_PROMPT);
1027 }
1028 if prompt_mode == PromptMode::Workflow
1029 && self.workflows_enabled()
1030 && let Some(runtime) = self.workflows.get()
1031 {
1032 let limits = runtime.limits();
1033 let size = self.settings.lock().unwrap().workflows.size;
1034 system_prompt.push_str("\n\n");
1035 system_prompt.push_str(&crate::workflows::authoring_prompt(
1036 size,
1037 limits.max_agents,
1038 limits.max_fanout,
1039 ));
1040 }
1041 AgentContext {
1042 system_prompt,
1043 openai_responses_input,
1044 messages,
1045 tools: self.tools_for(prompt_mode),
1046 }
1047 }
1048
1049 pub(crate) fn create_subagent_session(
1050 self: &Arc<Self>,
1051 task_name: &str,
1052 canonical_path: &str,
1053 fork_turns: ForkTurns,
1054 model_pattern: Option<&str>,
1055 reasoning_effort: Option<&str>,
1056 ) -> anyhow::Result<Arc<Self>> {
1057 let (model, suggested_thinking) = match model_pattern {
1058 Some(pattern) => self
1059 .registry
1060 .resolve(pattern, None)
1061 .with_context(|| format!("no model matches subagent model '{pattern}'"))?,
1062 None => (self.model(), None),
1063 };
1064 let thinking = match reasoning_effort {
1065 Some(level) => ThinkingLevel::parse(level)
1066 .with_context(|| format!("unknown subagent reasoning_effort '{level}'"))?,
1067 None => suggested_thinking.unwrap_or_else(|| self.thinking_level()),
1068 };
1069 let (mut manager, parent_messages, parent_id) = {
1070 let parent = self.manager.lock().unwrap();
1071 (
1072 parent.create_child()?,
1073 if fork_turns == ForkTurns::None {
1074 Vec::new()
1075 } else {
1076 parent.build_session_context().messages
1077 },
1078 parent.session_id().to_string(),
1079 )
1080 };
1081 for message in fork_messages(&parent_messages, fork_turns) {
1082 manager.append_message(message)?;
1083 }
1084 manager.append_custom(
1085 "subagent",
1086 Some(serde_json::json!({
1087 "taskName": task_name,
1088 "canonicalPath": canonical_path,
1089 "parentSessionId": parent_id,
1090 })),
1091 )?;
1092
1093 let mut settings = self.settings();
1094 settings.subagents.enabled = false;
1095 let mut system_prompt = self.system_prompt.lock().unwrap().clone();
1096 system_prompt.push_str(&format!(
1097 "\n\nYou are child agent {canonical_path}. Complete only the assigned task. Return a concise result to the parent agent."
1098 ));
1099
1100 let child = Self::new_with_subagents_allowed(
1101 manager,
1102 self.base_tools.lock().unwrap().clone(),
1103 self.registry.clone(),
1104 settings,
1105 system_prompt,
1106 model,
1107 thinking,
1108 self.api_key_override.clone(),
1109 Arc::new(|_| {}),
1110 false,
1111 );
1112 child.set_stream_fn(self.stream_fn.lock().unwrap().clone());
1113 Ok(child)
1114 }
1115
1116 async fn run_ephemeral(
1117 self: &Arc<Self>,
1118 system_prompt: String,
1119 prompt: String,
1120 tools: Vec<DynTool>,
1121 max_tokens: u64,
1122 cancel: CancellationToken,
1123 ) -> anyhow::Result<EphemeralResponse> {
1124 let mut config = self.loop_config(self, Arc::new(Mutex::new(PromptMode::Ordinary)));
1125 config.thinking_level = ThinkingLevel::Off;
1126 config.max_tokens = Some(max_tokens);
1127 config.session_id = Some(format!("ephemeral-{}", uuid::Uuid::new_v4()));
1128 config.get_steering_messages = None;
1129 config.get_follow_up_messages = None;
1130 config.prepare_next_turn = None;
1131 config.prepare_generation = None;
1132
1133 let context = AgentContext {
1134 system_prompt,
1135 openai_responses_input: None,
1136 messages: Vec::new(),
1137 tools,
1138 };
1139 let sink: EventSink = Arc::new(|_| {});
1140 let messages = kiss_agent::run_agent_loop(
1141 vec![AgentMessage::user(prompt)],
1142 context,
1143 config,
1144 cancel.clone(),
1145 sink,
1146 )
1147 .await;
1148 if cancel.is_cancelled() {
1149 anyhow::bail!("request cancelled");
1150 }
1151
1152 let mut usage = Usage::default();
1153 for message in &messages {
1154 if let AgentMessage::Assistant(assistant) = message {
1155 usage.add(&assistant.usage);
1156 }
1157 }
1158 let assistant = messages.iter().rev().find_map(|message| match message {
1159 AgentMessage::Assistant(assistant) => Some(assistant),
1160 _ => None,
1161 });
1162 let Some(assistant) = assistant else {
1163 anyhow::bail!("the provider returned no answer");
1164 };
1165 if assistant.stop_reason == StopReason::Error {
1166 anyhow::bail!(
1167 "{}",
1168 assistant
1169 .error_message
1170 .as_deref()
1171 .unwrap_or("the provider request failed")
1172 );
1173 }
1174 let text = assistant.text();
1175 if text.trim().is_empty() {
1176 anyhow::bail!("the provider returned an empty answer");
1177 }
1178 self.totals.lock().unwrap().add(&usage);
1179 Ok(EphemeralResponse { text, usage })
1180 }
1181
1182 pub async fn answer_btw(
1184 self: &Arc<Self>,
1185 question: &str,
1186 cancel: CancellationToken,
1187 ) -> anyhow::Result<EphemeralResponse> {
1188 let messages = self
1189 .manager
1190 .lock()
1191 .unwrap()
1192 .build_session_context()
1193 .messages;
1194 let transcript = transcript_excerpt(&messages, 4, 4_000);
1195 let prompt = if transcript.is_empty() {
1196 format!("Side question:\n{question}")
1197 } else {
1198 format!("Recent session context:\n{transcript}\n\nSide question:\n{question}")
1199 };
1200 let read_tools = self
1201 .tools
1202 .lock()
1203 .unwrap()
1204 .iter()
1205 .filter(|tool| tool.name() == "read")
1206 .cloned()
1207 .collect();
1208 self.run_ephemeral(
1209 "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(),
1210 prompt,
1211 read_tools,
1212 500,
1213 cancel,
1214 )
1215 .await
1216 }
1217
1218 pub async fn generate_session_title(
1220 self: &Arc<Self>,
1221 prompt: &str,
1222 cancel: CancellationToken,
1223 ) -> anyhow::Result<String> {
1224 let prompt = bounded_session_title_prompt(prompt);
1225 if prompt.is_empty() {
1226 anyhow::bail!("the session prompt is empty");
1227 }
1228 let response = self
1229 .run_ephemeral(
1230 format!(
1231 "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."
1232 ),
1233 format!("User prompt:\n{prompt}"),
1234 Vec::new(),
1235 64,
1236 cancel,
1237 )
1238 .await?;
1239 normalize_session_title(&response.text)
1240 .context("the provider returned an invalid session title")
1241 }
1242
1243 pub async fn generate_recap(
1245 self: &Arc<Self>,
1246 previous_recap: Option<&str>,
1247 cancel: CancellationToken,
1248 ) -> anyhow::Result<EphemeralResponse> {
1249 let messages = self
1250 .manager
1251 .lock()
1252 .unwrap()
1253 .build_session_context()
1254 .messages;
1255 let transcript = transcript_excerpt(&messages, 12, 12_000);
1256 if transcript.is_empty() {
1257 anyhow::bail!("the session has no conversation to recap");
1258 }
1259 let previous = previous_recap
1260 .filter(|recap| !recap.trim().is_empty())
1261 .map(|recap| format!("\n\nPrevious recap:\n{recap}"))
1262 .unwrap_or_default();
1263 self.run_ephemeral(
1264 "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(),
1265 format!("Session transcript:\n{transcript}{previous}"),
1266 Vec::new(),
1267 160,
1268 cancel,
1269 )
1270 .await
1271 }
1272
1273 pub async fn prompt(self: &Arc<Self>, prompts: Vec<AgentMessage>) {
1275 self.prompt_with_mode(prompts, PromptMode::Ordinary).await;
1276 }
1277
1278 pub async fn prompt_with_mode(
1280 self: &Arc<Self>,
1281 prompts: Vec<AgentMessage>,
1282 prompt_mode: PromptMode,
1283 ) {
1284 self.cache_warm_cancel.lock().unwrap().cancel();
1285 {
1286 let mut running = self.running.lock().unwrap();
1287 if *running {
1288 drop(running);
1290 for p in prompts {
1291 self.queue_steering_with_mode(p, prompt_mode);
1292 }
1293 return;
1294 }
1295 *running = true;
1296 }
1297 let cancel = {
1298 let mut guard = self.cancel.lock().unwrap();
1299 *guard = CancellationToken::new();
1300 guard.clone()
1301 };
1302
1303 {
1305 let mut manager = self.manager.lock().unwrap();
1306 for p in &prompts {
1307 let _ = manager.append_message(p.clone());
1308 }
1309 }
1310
1311 let session = self.clone();
1312 let sink: EventSink = Arc::new(move |event: AgentEvent| {
1313 session.on_agent_event(&event);
1314 (session.sink)(SessionEvent::Agent(Box::new(event)));
1315 });
1316
1317 let active_prompt_mode = Arc::new(Mutex::new(prompt_mode));
1318 let mut config = self.loop_config(self, active_prompt_mode.clone());
1319 let mut context = self.build_context_for(prompt_mode);
1320 let mut attempt: u32 = 0;
1324 loop {
1325 let messages = kiss_agent::run_agent_loop_continue(
1326 context,
1327 config.clone(),
1328 cancel.clone(),
1329 sink.clone(),
1330 )
1331 .await;
1332
1333 let last_error = messages.iter().rev().find_map(|m| match m {
1335 AgentMessage::Assistant(a) if a.stop_reason == StopReason::Error => {
1336 Some(a.error_message.clone().unwrap_or_default())
1337 }
1338 _ => None,
1339 });
1340 let settings = self.settings();
1341 let retry = &settings.retry;
1342 if let Some(error) = last_error
1343 && retry.enabled
1344 && attempt < retry.max_retries
1345 && is_transient(&error)
1346 && !cancel.is_cancelled()
1347 {
1348 attempt += 1;
1349 let delay = retry
1350 .base_delay_ms
1351 .saturating_mul(1u64.checked_shl(attempt - 1).unwrap_or(u64::MAX))
1352 .min(retry.max_agent_delay_ms);
1353 (self.sink)(SessionEvent::Retry {
1354 attempt,
1355 max: retry.max_retries,
1356 delay_ms: delay,
1357 error,
1358 });
1359 tokio::select! {
1360 _ = tokio::time::sleep(std::time::Duration::from_millis(delay)) => {}
1361 _ = cancel.cancelled() => break,
1362 }
1363 context = self.build_context_for(*active_prompt_mode.lock().unwrap());
1364 while matches!(
1366 context.messages.last(),
1367 Some(AgentMessage::Assistant(a)) if a.stop_reason == StopReason::Error
1368 ) {
1369 context.messages.pop();
1370 }
1371 config.model = self.model();
1372 config.thinking_level = self.thinking_level();
1373 config.fast_mode = self.fast_mode();
1374 continue;
1375 }
1376
1377 let ctx = self.manager.lock().unwrap().build_session_context();
1379 if auto_compaction_needed(
1380 &settings,
1381 &ctx.messages,
1382 &self.model(),
1383 cancel.is_cancelled(),
1384 ) {
1385 self.compact(None, true).await;
1386 }
1387 break;
1388 }
1389
1390 *self.running.lock().unwrap() = false;
1391 if self.settings.lock().unwrap().cache_warming == CacheWarmingMode::Streaming {
1392 self.cache_warm_cancel.lock().unwrap().cancel();
1393 }
1394 }
1395
1396 fn on_agent_event(self: &Arc<Self>, event: &AgentEvent) {
1397 match event {
1398 AgentEvent::MessageEnd { message } => {
1399 let persist = match message {
1402 AgentMessage::Assistant(a) => {
1403 self.schedule_cache_warming(a);
1404 let mut totals = self.totals.lock().unwrap();
1405 totals.add(&a.usage);
1406 true
1407 }
1408 AgentMessage::ToolResult(_)
1409 | AgentMessage::User(_)
1410 | AgentMessage::Custom(_) => true,
1411 _ => false,
1412 };
1413 if persist {
1414 let mut manager = self.manager.lock().unwrap();
1417 let duplicate = matches!(
1418 (manager.entries().last(), message),
1419 (Some(crate::session::entry::SessionEntry::Message { message: last, .. }), m) if last == m
1420 );
1421 if !duplicate {
1422 let _ = manager.append_message(message.clone());
1423 }
1424 }
1425 }
1426 AgentEvent::AgentEnd { .. } => {}
1427 _ => {}
1428 }
1429 }
1430
1431 fn schedule_cache_warming(self: &Arc<Self>, assistant: &kiss_ai::AssistantMessage) {
1432 let settings = self.settings();
1433 let model = self.model();
1434 let Some(cache) = model.prompt_cache else {
1435 return;
1436 };
1437 let Some(short_ttl) = cache.short else {
1438 return;
1439 };
1440 let reasoning = self.thinking_level();
1441 let fast_mode = self.fast_mode();
1442 if settings.cache_warming == CacheWarmingMode::Off
1443 || matches!(
1444 assistant.stop_reason,
1445 StopReason::Error | StopReason::Aborted
1446 )
1447 || (assistant.usage.input + assistant.usage.cache_read + assistant.usage.cache_write
1448 == 0)
1449 || (reasoning != ThinkingLevel::Off
1450 && model.api == "anthropic-messages"
1451 && !model
1452 .compat
1453 .as_ref()
1454 .and_then(|compat| compat.force_adaptive_thinking)
1455 .unwrap_or(false))
1456 {
1457 return;
1458 }
1459 let ttl = std::time::Duration::from_secs(short_ttl);
1460 let Some(delay) = cache_warming_delay(ttl) else {
1461 return;
1462 };
1463 let prompt_tokens =
1464 assistant.usage.input + assistant.usage.cache_read + assistant.usage.cache_write;
1465 let agent_context = self.build_context();
1466 let context = kiss_ai::Context {
1467 system_prompt: Some(agent_context.system_prompt),
1468 openai_responses_input: agent_context.openai_responses_input,
1469 messages: kiss_agent::convert_to_llm(&agent_context.messages),
1470 tools: agent_context
1471 .tools
1472 .iter()
1473 .map(|tool| tool.to_def())
1474 .collect(),
1475 };
1476 let cancel = CancellationToken::new();
1477 {
1478 let mut current = self.cache_warm_cancel.lock().unwrap();
1479 current.cancel();
1480 *current = cancel.clone();
1481 }
1482 let session = self.clone();
1483 tokio::spawn(async move {
1484 let started = tokio::time::Instant::now();
1485 loop {
1486 let scheduled = tokio::time::Instant::now();
1487 if tokio::select! {
1488 _ = tokio::time::sleep(delay) => false,
1489 _ = cancel.cancelled() => true,
1490 } {
1491 return;
1492 }
1493 if cache_refresh_deadline_missed(scheduled.elapsed(), ttl, delay) {
1494 return;
1495 }
1496 let idle = !session.is_running();
1497 let max_age = if idle {
1498 std::time::Duration::from_secs(30 * 60)
1499 } else {
1500 std::time::Duration::from_secs(60 * 60)
1501 };
1502 if started.elapsed() > max_age {
1503 return;
1504 }
1505 let priced = |input: u64, output: u64, cache_read: u64, cache_write: u64| {
1506 let mut usage = Usage {
1507 input,
1508 output,
1509 cache_read,
1510 cache_write,
1511 ..Default::default()
1512 };
1513 kiss_ai::api::finalize_cost(&mut usage, &model);
1514 usage.cost.total
1515 };
1516 let hit_cost = priced(0, 0, prompt_tokens, 0);
1517 let miss_cost = if model.cost.cache_write > 0.0 {
1518 priced(0, 0, 0, prompt_tokens)
1519 } else {
1520 priced(prompt_tokens, 0, 0, 0)
1521 };
1522 let warm_cost = priced(0, 1, prompt_tokens, 0);
1523 let probability = if idle { 0.15 } else { 1.0 };
1524 if probability * (miss_cost - hit_cost).max(0.0) - warm_cost < 0.05 {
1525 return;
1526 }
1527 let Some(credential) = session.resolve_credential(&model.provider).await else {
1528 return;
1529 };
1530 let options = kiss_ai::StreamOptions {
1531 credential: Some(credential),
1532 max_tokens: Some(1),
1533 reasoning,
1534 fast_mode,
1535 session_id: Some(session.manager.lock().unwrap().session_id().to_string()),
1536 transport: settings.transport,
1537 cancel: cancel.clone(),
1538 ..Default::default()
1539 };
1540 let stream_fn = session
1541 .stream_fn
1542 .lock()
1543 .unwrap()
1544 .clone()
1545 .unwrap_or_else(|| Arc::new(kiss_ai::stream_simple));
1546 let warmed = stream_fn(&model, &context, &options).result().await;
1547 if matches!(warmed.stop_reason, StopReason::Error | StopReason::Aborted) {
1548 return;
1549 }
1550 session.totals.lock().unwrap().add(&warmed.usage);
1551 let _ = session.manager.lock().unwrap().append_usage(
1552 "cache_warm",
1553 &warmed.provider,
1554 warmed.response_model.as_deref().unwrap_or(&warmed.model),
1555 warmed.usage,
1556 None,
1557 );
1558 }
1559 });
1560 }
1561
1562 pub async fn compact(self: &Arc<Self>, custom_instructions: Option<String>, auto: bool) {
1564 (self.sink)(SessionEvent::CompactionStart { auto });
1565 let ctx = self.manager.lock().unwrap().build_session_context();
1566 let previous_summary = ctx.messages.iter().rev().find_map(|m| match m {
1567 AgentMessage::CompactionSummary(c) => Some(c.summary.clone()),
1568 _ => None,
1569 });
1570 let settings = self.settings();
1571 let model = self.model();
1572 let override_settings = settings
1573 .compaction
1574 .model_overrides
1575 .get(&format!("{}/{}", model.provider, model.id));
1576 let keep_recent_tokens = override_settings
1577 .and_then(|value| value.keep_recent_tokens)
1578 .unwrap_or(settings.compaction.keep_recent_tokens);
1579 let reserve_tokens = override_settings
1580 .and_then(|value| value.reserve_tokens)
1581 .unwrap_or(settings.compaction.reserve_tokens);
1582 let plan = plan_compaction(&ctx.messages, keep_recent_tokens);
1583 if plan.to_summarize.is_empty() && plan.turn_prefix.is_empty() {
1584 (self.sink)(SessionEvent::CompactionEnd {
1585 summary: String::new(),
1586 tokens_before: plan.tokens_before,
1587 error: Some("Nothing to compact".into()),
1588 });
1589 return;
1590 }
1591
1592 if settings.compaction.mode == CompactionMode::Jev
1593 && let Ok(Some(api_key)) =
1594 kiss_ai::auth::resolve_api_key_async("typesafe", &self.registry.declared_keys).await
1595 {
1596 let cancel = self.cancel.lock().unwrap().clone();
1597 let pinned_start = ctx.messages.len().saturating_sub(plan.kept.len());
1598 if let Ok(result) =
1599 crate::jev::compact(&ctx.messages, pinned_start, &api_key, cancel).await
1600 {
1601 let estimated_before: u64 = ctx
1602 .messages
1603 .iter()
1604 .map(compaction::estimate_message_tokens)
1605 .sum();
1606 let estimated_after: u64 = result
1607 .messages
1608 .iter()
1609 .map(compaction::estimate_message_tokens)
1610 .sum();
1611 let removed = estimated_before.saturating_sub(estimated_after);
1612 let tokens_after = plan.tokens_before.saturating_sub(removed);
1613 let useful =
1614 estimated_after.saturating_mul(4) <= estimated_before.saturating_mul(3);
1615 let resolved_auto_threshold =
1616 !auto || !should_compact(tokens_after, model.context_window, reserve_tokens);
1617 if useful && resolved_auto_threshold {
1618 let summary = format!(
1619 "Jev kept {}, truncated {}, and removed {} of {} older tool interactions",
1620 result.stats.kept,
1621 result.stats.truncated,
1622 result.stats.dropped,
1623 result.stats.eligible,
1624 );
1625 let details = serde_json::json!({
1626 "mode": "jev",
1627 "stats": result.stats,
1628 "estimatedTokensBefore": plan.tokens_before,
1629 "estimatedTokensAfter": tokens_after,
1630 });
1631 let mut manager = self.manager.lock().unwrap();
1632 let append = manager.append_compaction(
1633 String::new(),
1634 plan.tokens_before,
1635 result.messages,
1636 None,
1637 Some(details),
1638 );
1639 drop(manager);
1640 let error = append.err().map(|error| format!("{error:#}"));
1641 (self.sink)(SessionEvent::CompactionEnd {
1642 summary: if error.is_none() {
1643 summary
1644 } else {
1645 String::new()
1646 },
1647 tokens_before: plan.tokens_before,
1648 error,
1649 });
1650 return;
1651 }
1652 }
1653 }
1654
1655 let credential = self.resolve_credential(&model.provider).await;
1656
1657 let mut serialized = compaction::serialize_agent_messages(&plan.to_summarize);
1658 if plan.is_split_turn {
1659 serialized.push_str(
1660 "\n\n[The following is the earlier part of the still-active task turn:]\n\n",
1661 );
1662 serialized.push_str(&compaction::serialize_agent_messages(&plan.turn_prefix));
1663 }
1664
1665 let summary_cancel = self.cancel.lock().unwrap().clone();
1666 let remote_request = if kiss_ai::api::openai_compaction::supports_remote_compaction(&model)
1667 {
1668 let context = self.build_context();
1669 Some((
1670 kiss_ai::Context {
1671 system_prompt: Some(context.system_prompt),
1672 openai_responses_input: context.openai_responses_input,
1673 messages: kiss_agent::convert_to_llm(&context.messages),
1674 tools: context.tools.iter().map(|tool| tool.to_def()).collect(),
1675 },
1676 kiss_ai::StreamOptions {
1677 credential: credential.clone(),
1678 reasoning: self.thinking_level(),
1679 fast_mode: self.fast_mode(),
1680 session_id: Some(self.manager.lock().unwrap().session_id().to_string()),
1681 cancel: summary_cancel.clone(),
1682 ..Default::default()
1683 },
1684 ))
1685 } else {
1686 None
1687 };
1688 let local_future = compaction::generate_summary(
1689 &model,
1690 credential.clone(),
1691 &serialized,
1692 previous_summary.as_deref(),
1693 custom_instructions.as_deref(),
1694 plan.is_split_turn,
1695 summary_cancel.clone(),
1696 );
1697 let remote_model = model.clone();
1698 let remote_future = async move {
1699 match remote_request {
1700 Some((context, options)) => Some(
1701 kiss_ai::api::openai_compaction::compact(&remote_model, &context, &options)
1702 .await,
1703 ),
1704 None => None,
1705 }
1706 };
1707 let (local_outcome, remote_outcome) = tokio::join!(local_future, remote_future);
1708
1709 match select_compaction_outcome(&model, local_outcome, remote_outcome) {
1710 Ok(result) => {
1711 let mut summarized_all = plan.to_summarize.clone();
1712 summarized_all.extend(plan.turn_prefix.clone());
1713 let (read, modified) = extract_file_ops(&summarized_all);
1714 {
1715 let mut totals = self.totals.lock().unwrap();
1716 if let Some(u) = &result.local_usage {
1717 totals.add(u);
1718 }
1719 if let Some(u) = &result.remote_usage {
1720 totals.add(u);
1721 }
1722 }
1723 let details = merge_compaction_details(
1724 file_ops_details(&read, &modified),
1725 result.remote_details,
1726 );
1727 let mut manager = self.manager.lock().unwrap();
1728 let _ = manager.append_compaction(
1729 result.summary.clone(),
1730 plan.tokens_before,
1731 plan.kept.clone(),
1732 result.local_usage,
1733 Some(details),
1734 );
1735 (self.sink)(SessionEvent::CompactionEnd {
1736 summary: result.summary,
1737 tokens_before: plan.tokens_before,
1738 error: None,
1739 });
1740 }
1741 Err(error) => {
1742 (self.sink)(SessionEvent::CompactionEnd {
1743 summary: String::new(),
1744 tokens_before: plan.tokens_before,
1745 error: Some(format!("{error:#}")),
1746 });
1747 }
1748 }
1749 }
1750
1751 pub fn context_usage(&self) -> (u64, u64) {
1753 let manager = self.manager.lock().unwrap();
1754 let revision = manager.context_revision();
1755 let used = if let Some((cached_revision, tokens)) =
1756 *self.context_usage_cache.lock().unwrap()
1757 && cached_revision == revision
1758 {
1759 tokens
1760 } else {
1761 let tokens = estimate_context_tokens(&manager.build_session_context().messages);
1762 *self.context_usage_cache.lock().unwrap() = Some((revision, tokens));
1763 tokens
1764 };
1765 drop(manager);
1766 (used, self.model().context_window)
1767 }
1768}
1769
1770struct SelectedCompaction {
1771 summary: String,
1772 local_usage: Option<Usage>,
1773 remote_usage: Option<Usage>,
1774 remote_details: Option<serde_json::Value>,
1775}
1776
1777fn select_compaction_outcome(
1778 model: &Model,
1779 local: anyhow::Result<compaction::SummaryOutcome>,
1780 remote: Option<anyhow::Result<kiss_ai::api::openai_compaction::RemoteCompactionResult>>,
1781) -> anyhow::Result<SelectedCompaction> {
1782 match (local, remote) {
1783 (Ok(local), Some(Ok(remote))) => Ok(SelectedCompaction {
1784 summary: local.summary,
1785 local_usage: local.usage,
1786 remote_usage: remote.usage,
1787 remote_details: Some(
1788 kiss_ai::api::openai_compaction::build_remote_compaction_details(model, &remote),
1789 ),
1790 }),
1791 (Ok(local), Some(Err(_)) | None) => Ok(SelectedCompaction {
1792 summary: local.summary,
1793 local_usage: local.usage,
1794 remote_usage: None,
1795 remote_details: None,
1796 }),
1797 (Err(_), Some(Ok(remote))) => Ok(SelectedCompaction {
1798 summary: format!(
1799 "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.",
1800 model.provider, model.id
1801 ),
1802 local_usage: None,
1803 remote_usage: remote.usage,
1804 remote_details: Some(
1805 kiss_ai::api::openai_compaction::build_remote_compaction_details(model, &remote),
1806 ),
1807 }),
1808 (Err(local), Some(Err(remote))) => anyhow::bail!(
1809 "local compaction failed: {local:#}. OpenAI remote compaction failed: {remote:#}"
1810 ),
1811 (Err(error), None) => Err(error),
1812 }
1813}
1814
1815fn merge_compaction_details(
1816 mut local: serde_json::Value,
1817 remote: Option<serde_json::Value>,
1818) -> serde_json::Value {
1819 let Some(remote) = remote else {
1820 return local;
1821 };
1822 let Some(local_object) = local.as_object_mut() else {
1823 return remote;
1824 };
1825 if let Some(remote_object) = remote.as_object() {
1826 for (key, value) in remote_object {
1827 local_object.insert(key.clone(), value.clone());
1828 }
1829 }
1830 local
1831}
1832
1833fn transcript_excerpt(messages: &[AgentMessage], max_messages: usize, max_chars: usize) -> String {
1834 let mut entries = messages
1835 .iter()
1836 .rev()
1837 .filter_map(|message| match message {
1838 AgentMessage::User(user) => Some(("User", user.content.as_text())),
1839 AgentMessage::Assistant(assistant) => Some(("Assistant", assistant.text())),
1840 _ => None,
1841 })
1842 .filter(|(_, text)| !text.trim().is_empty())
1843 .take(max_messages)
1844 .collect::<Vec<_>>();
1845 entries.reverse();
1846 let transcript = entries
1847 .into_iter()
1848 .map(|(role, text)| format!("{role}: {}", text.trim()))
1849 .collect::<Vec<_>>()
1850 .join("\n\n");
1851 let count = transcript.chars().count();
1852 if count <= max_chars {
1853 return transcript;
1854 }
1855 let omitted = count - max_chars;
1856 let tail = transcript.chars().skip(omitted).collect::<String>();
1857 format!("[earlier text omitted]\n{tail}")
1858}
1859
1860fn drain_queue(queue: &Arc<Mutex<VecDeque<QueuedPrompt>>>, mode: QueueMode) -> Vec<AgentMessage> {
1861 let mut q = queue.lock().unwrap();
1862 match mode {
1863 QueueMode::All => q.drain(..).map(|prompt| prompt.message).collect(),
1864 QueueMode::OneAtATime => q
1865 .pop_front()
1866 .map(|prompt| prompt.message)
1867 .into_iter()
1868 .collect(),
1869 }
1870}
1871
1872fn queued_mode(queue: &Arc<Mutex<VecDeque<QueuedPrompt>>>, mode: QueueMode) -> Option<PromptMode> {
1873 let queue = queue.lock().unwrap();
1874 match mode {
1875 QueueMode::All => queue
1876 .iter()
1877 .any(|prompt| prompt.mode == PromptMode::Workflow)
1878 .then_some(PromptMode::Workflow)
1879 .or_else(|| (!queue.is_empty()).then_some(PromptMode::Ordinary)),
1880 QueueMode::OneAtATime => queue.front().map(|prompt| prompt.mode),
1881 }
1882}
1883
1884fn auto_compaction_needed(
1885 settings: &Settings,
1886 messages: &[AgentMessage],
1887 model: &Model,
1888 cancelled: bool,
1889) -> bool {
1890 let reserve_tokens = settings
1891 .compaction
1892 .model_overrides
1893 .get(&format!("{}/{}", model.provider, model.id))
1894 .and_then(|value| value.reserve_tokens)
1895 .unwrap_or(settings.compaction.reserve_tokens);
1896 settings.compaction.enabled
1897 && !cancelled
1898 && model.context_window > 0
1899 && should_compact(
1900 estimate_context_tokens(messages),
1901 model.context_window,
1902 reserve_tokens,
1903 )
1904}
1905
1906fn cache_warming_delay(ttl: std::time::Duration) -> Option<std::time::Duration> {
1907 (ttl > std::time::Duration::from_secs(10))
1908 .then(|| std::cmp::min(ttl.mul_f64(0.9), ttl - std::time::Duration::from_secs(10)))
1909}
1910
1911fn cache_refresh_deadline_missed(
1912 elapsed: std::time::Duration,
1913 ttl: std::time::Duration,
1914 delay: std::time::Duration,
1915) -> bool {
1916 elapsed > delay + ttl.saturating_sub(delay) / 2
1917}
1918
1919fn is_transient(error: &str) -> bool {
1920 let e = error.to_lowercase();
1921 let transient_status = [429, 500, 502, 503, 504, 520].iter().any(|status| {
1922 [
1923 format!("http {status}"),
1924 format!("status {status}"),
1925 format!("status: {status}"),
1926 format!("status code {status}"),
1927 ]
1928 .iter()
1929 .any(|marker| e.contains(marker))
1930 });
1931 transient_status
1932 || [
1933 "overloaded",
1934 "currently experiencing high demand",
1935 "rate limit",
1936 "timeout",
1937 "timed out",
1938 "connection reset",
1939 "connection refused",
1940 "connection closed",
1941 "connection aborted",
1942 "connection error",
1943 "failed to connect",
1944 "network error",
1945 "stream error",
1946 "request failed",
1947 ]
1948 .iter()
1949 .any(|needle| e.contains(needle))
1950}
1951
1952#[cfg(test)]
1953mod ephemeral_tests {
1954 use super::*;
1955 use std::collections::BTreeMap;
1956
1957 fn openai_model() -> Model {
1958 Model {
1959 id: "gpt-test".into(),
1960 name: "GPT test".into(),
1961 api: "openai-responses".into(),
1962 provider: "openai".into(),
1963 base_url: "https://api.openai.com/v1".into(),
1964 reasoning: true,
1965 input: vec!["text".into()],
1966 cost: Default::default(),
1967 prompt_cache: None,
1968 context_window: 100_000,
1969 max_tokens: 1_000,
1970 compat: None,
1971 thinking_level_map: BTreeMap::new(),
1972 headers: BTreeMap::new(),
1973 }
1974 }
1975
1976 #[test]
1977 fn reasoning_lease_ends_when_model_saved_effort_or_applied_effort_changes() {
1978 let model = openai_model();
1979 let mut run = ReasoningRunState {
1980 lease: crate::jev::ReasoningLease::default(),
1981 provider: model.provider.clone(),
1982 model_id: model.id.clone(),
1983 saved_effort: ThinkingLevel::Medium,
1984 };
1985 let selection = crate::jev::ReasoningSelection {
1986 level: ThinkingLevel::High,
1987 generations: 5,
1988 };
1989 run.lease.install(&selection);
1990 assert_eq!(
1991 run.checkpoint(&model, ThinkingLevel::Medium, ThinkingLevel::High, true),
1992 (false, Some(ThinkingLevel::High))
1993 );
1994
1995 let mut other_model = model.clone();
1996 other_model.provider = "other".into();
1997 assert_eq!(
1998 run.checkpoint(
1999 &other_model,
2000 ThinkingLevel::Medium,
2001 ThinkingLevel::High,
2002 true
2003 ),
2004 (true, None)
2005 );
2006 run.lease.install(&selection);
2007 assert_eq!(
2008 run.checkpoint(&other_model, ThinkingLevel::Low, ThinkingLevel::High, true),
2009 (false, None)
2010 );
2011 run.lease.install(&selection);
2012 assert_eq!(
2013 run.checkpoint(&other_model, ThinkingLevel::Low, ThinkingLevel::Low, true),
2014 (false, None)
2015 );
2016 run.lease.install(&selection);
2017 assert_eq!(
2018 run.checkpoint(&other_model, ThinkingLevel::Low, ThinkingLevel::High, false),
2019 (false, None)
2020 );
2021 }
2022
2023 #[test]
2024 fn late_cache_refreshes_are_skipped_before_the_cache_expires() {
2025 let ttl = std::time::Duration::from_secs(300);
2026 let delay = cache_warming_delay(ttl).unwrap();
2027 assert!(!cache_refresh_deadline_missed(
2028 std::time::Duration::from_secs(284),
2029 ttl,
2030 delay,
2031 ));
2032 assert!(cache_refresh_deadline_missed(
2033 std::time::Duration::from_secs(286),
2034 ttl,
2035 delay,
2036 ));
2037 }
2038
2039 #[test]
2040 fn transient_errors_require_a_status_or_specific_network_failure() {
2041 assert!(is_transient("request failed with HTTP 503"));
2042 assert!(is_transient("Cloudflare returned HTTP 520"));
2043 assert!(is_transient("Azure is currently experiencing high demand"));
2044 assert!(is_transient("connection reset by peer"));
2045 assert!(is_transient("rate limit exceeded"));
2046 assert!(!is_transient("model has a 500 token limit"));
2047 assert!(!is_transient("connection settings are invalid"));
2048 }
2049
2050 #[test]
2051 fn session_tools_install_replace_and_remove_without_changing_base_tools() {
2052 let registry = Registry::load(None);
2053 let session = AgentSession::new(
2054 SessionManager::in_memory(std::path::Path::new("/test")),
2055 Vec::new(),
2056 registry,
2057 Settings::default(),
2058 "test".into(),
2059 openai_model(),
2060 ThinkingLevel::Off,
2061 None,
2062 Arc::new(|_| {}),
2063 );
2064 let tool = || {
2065 Arc::new(crate::tools::grep::GrepTool {
2066 cwd: std::path::PathBuf::from("/test"),
2067 }) as DynTool
2068 };
2069
2070 assert!(!session.available_tool_names().contains(&"grep".into()));
2071 session.install_session_tool(tool());
2072 session.install_session_tool(tool());
2073 assert_eq!(
2074 session
2075 .available_tool_names()
2076 .iter()
2077 .filter(|name| name.as_str() == "grep")
2078 .count(),
2079 1
2080 );
2081 assert!(session.remove_session_tool("grep"));
2082 assert!(!session.remove_session_tool("grep"));
2083 assert!(!session.available_tool_names().contains(&"grep".into()));
2084 }
2085
2086 fn remote_result() -> kiss_ai::api::openai_compaction::RemoteCompactionResult {
2087 kiss_ai::api::openai_compaction::RemoteCompactionResult {
2088 replacement_history: vec![serde_json::json!({
2089 "type": "compaction",
2090 "encrypted_content": "opaque"
2091 })],
2092 usage: Some(Usage {
2093 input: 10,
2094 output: 2,
2095 total_tokens: 12,
2096 ..Default::default()
2097 }),
2098 }
2099 }
2100
2101 fn settings_test_session(settings: Settings, subagents_allowed: bool) -> Arc<AgentSession> {
2102 let registry = Registry::from_builtin();
2103 let model = registry.all().first().expect("built-in model").clone();
2104 AgentSession::new_with_subagents_allowed(
2105 SessionManager::in_memory(std::path::Path::new("/test")),
2106 Vec::new(),
2107 registry,
2108 settings,
2109 "root prompt".into(),
2110 model,
2111 ThinkingLevel::Off,
2112 None,
2113 Arc::new(|_| {}),
2114 subagents_allowed,
2115 )
2116 }
2117
2118 fn benchmark_tools() -> Vec<DynTool> {
2119 let cwd = std::path::PathBuf::from("/synthetic");
2120 vec![
2121 Arc::new(kiss_agent::tools::read::ReadTool { cwd: cwd.clone() }),
2122 Arc::new(kiss_agent::tools::write::WriteTool { cwd: cwd.clone() }),
2123 Arc::new(kiss_agent::tools::edit::EditTool { cwd: cwd.clone() }),
2124 Arc::new(kiss_agent::tools::bash::BashTool::new(cwd)),
2125 ]
2126 }
2127
2128 #[test]
2129 fn subagent_tools_follow_settings_and_command_line_authority() {
2130 let session = settings_test_session(Settings::default(), true);
2131 assert!(session.available_tool_names().is_empty());
2132 assert!(
2133 !session
2134 .build_context()
2135 .system_prompt
2136 .contains("Subagent coordination")
2137 );
2138
2139 let mut enabled = session.settings();
2140 enabled.subagents.enabled = true;
2141 session.update_settings(enabled.clone());
2142 assert_eq!(
2143 session.available_tool_names(),
2144 [
2145 "spawn_agent",
2146 "send_message",
2147 "followup_task",
2148 "wait_agent",
2149 "list_agents",
2150 "interrupt_agent"
2151 ]
2152 );
2153 assert!(
2154 session
2155 .build_context()
2156 .system_prompt
2157 .contains("Subagent coordination")
2158 );
2159
2160 enabled.subagents.enabled = false;
2161 session.update_settings(enabled);
2162 assert!(session.available_tool_names().is_empty());
2163
2164 let mut blocked_settings = Settings::default();
2165 blocked_settings.subagents.enabled = true;
2166 let blocked = settings_test_session(blocked_settings, false);
2167 assert!(blocked.available_tool_names().is_empty());
2168 assert!(
2169 !blocked
2170 .build_context()
2171 .system_prompt
2172 .contains("Subagent coordination")
2173 );
2174 }
2175
2176 #[test]
2177 fn the_workflow_tool_appears_only_in_workflow_prompt_mode() {
2178 let mut settings = Settings::default();
2179 settings.subagents.enabled = true;
2180 let session = settings_test_session(settings, true);
2181
2182 assert!(
2185 !session
2186 .available_tool_names()
2187 .contains(&"run_workflow".into())
2188 );
2189 assert_eq!(
2190 session.prompt_mode_for("run a dynamic workflow for this task"),
2191 PromptMode::Workflow
2192 );
2193 assert_eq!(
2194 session.prompt_mode_for("fix this small function"),
2195 PromptMode::Ordinary
2196 );
2197 assert!(
2198 !session
2199 .build_context()
2200 .system_prompt
2201 .contains("Writing a dynamic workflow")
2202 );
2203
2204 assert!(
2205 session
2206 .available_tool_names_for(PromptMode::Workflow)
2207 .contains(&"run_workflow".into())
2208 );
2209 assert!(
2210 session
2211 .build_context_for(PromptMode::Workflow)
2212 .system_prompt
2213 .contains("Writing a dynamic workflow")
2214 );
2215 assert!(
2216 !session
2217 .available_tool_names()
2218 .contains(&"run_workflow".into())
2219 );
2220 }
2221
2222 #[test]
2223 fn workflow_prompt_mode_does_nothing_while_subagents_are_off() {
2224 let session = settings_test_session(Settings::default(), true);
2227 assert!(!session.workflows_enabled());
2228 assert!(
2229 !session
2230 .available_tool_names_for(PromptMode::Workflow)
2231 .contains(&"run_workflow".into())
2232 );
2233
2234 let mut settings = session.settings();
2235 settings.subagents.enabled = true;
2236 settings.workflows.enabled = false;
2237 session.update_settings(settings.clone());
2238 assert!(!session.workflows_enabled());
2239 assert!(
2240 !session
2241 .available_tool_names_for(PromptMode::Workflow)
2242 .contains(&"run_workflow".into())
2243 );
2244
2245 settings.workflows.enabled = true;
2246 session.update_settings(settings);
2247 assert!(session.workflows_enabled());
2248 assert!(
2249 session
2250 .available_tool_names_for(PromptMode::Workflow)
2251 .contains(&"run_workflow".into())
2252 );
2253 }
2254
2255 #[test]
2256 fn one_at_a_time_queues_keep_each_prompts_mode() {
2257 let queue = Arc::new(Mutex::new(VecDeque::from([
2258 QueuedPrompt {
2259 message: AgentMessage::user("ordinary"),
2260 mode: PromptMode::Ordinary,
2261 },
2262 QueuedPrompt {
2263 message: AgentMessage::user("workflow"),
2264 mode: PromptMode::Workflow,
2265 },
2266 ])));
2267
2268 assert_eq!(
2269 queued_mode(&queue, QueueMode::OneAtATime),
2270 Some(PromptMode::Ordinary)
2271 );
2272 assert_eq!(drain_queue(&queue, QueueMode::OneAtATime).len(), 1);
2273 assert_eq!(
2274 queued_mode(&queue, QueueMode::OneAtATime),
2275 Some(PromptMode::Workflow)
2276 );
2277 }
2278
2279 #[test]
2280 fn a_session_without_subagent_authority_has_no_workflow_runtime() {
2281 let mut settings = Settings::default();
2282 settings.subagents.enabled = true;
2283 let child = settings_test_session(settings, false);
2284 assert!(child.workflows().is_none());
2285 assert!(!child.workflows_enabled());
2286 }
2287
2288 #[test]
2289 fn child_session_has_safe_forked_context_without_control_tools() {
2290 let mut settings = Settings::default();
2291 settings.subagents.enabled = true;
2292 let parent = settings_test_session(settings, true);
2293 parent
2294 .manager
2295 .lock()
2296 .unwrap()
2297 .append_message(AgentMessage::user("parent context"))
2298 .unwrap();
2299
2300 let child = parent
2301 .create_subagent_session("inspect", "/root/inspect", ForkTurns::All, None, None)
2302 .unwrap();
2303 assert!(Arc::ptr_eq(&parent.registry, &child.registry));
2304 assert!(child.available_tool_names().is_empty());
2305 let context = child.manager.lock().unwrap().build_session_context();
2306 assert!(matches!(
2307 context.messages.as_slice(),
2308 [AgentMessage::User(user)] if user.content.as_text() == "parent context"
2309 ));
2310 }
2311
2312 #[test]
2313 fn child_without_forked_turns_does_not_copy_parent_context() {
2314 let parent = settings_test_session(Settings::default(), true);
2315 parent
2316 .manager
2317 .lock()
2318 .unwrap()
2319 .append_message(AgentMessage::user("parent context"))
2320 .unwrap();
2321
2322 let child = parent
2323 .create_subagent_session("inspect", "/root/inspect", ForkTurns::None, None, None)
2324 .unwrap();
2325 assert!(
2326 child
2327 .manager
2328 .lock()
2329 .unwrap()
2330 .build_session_context()
2331 .messages
2332 .is_empty()
2333 );
2334 }
2335
2336 #[test]
2337 #[ignore = "release-mode performance benchmark"]
2338 fn benchmark_performance_subagent_overhead() {
2339 let registry = Registry::from_builtin();
2340 let model = registry.all().first().expect("built-in model").clone();
2341 let tools = benchmark_tools();
2342 let make_session = |enabled: bool| {
2343 let mut settings = Settings::default();
2344 settings.subagents.enabled = enabled;
2345 AgentSession::new_with_subagents_allowed(
2346 SessionManager::in_memory(std::path::Path::new("/synthetic")),
2347 tools.clone(),
2348 registry.clone(),
2349 settings,
2350 "benchmark root prompt".into(),
2351 model.clone(),
2352 ThinkingLevel::Off,
2353 None,
2354 Arc::new(|_| {}),
2355 true,
2356 )
2357 };
2358
2359 kiss_bench::measure_pair(
2360 (
2361 "agent_session_create_subagents_off",
2362 "agent_session_create_subagents_on",
2363 ),
2364 21,
2365 500,
2366 (
2367 "new_root_session_4_base_tools_0_control_tools",
2368 "new_root_session_4_base_tools_6_control_tools",
2369 ),
2370 || make_session(false),
2371 || make_session(true),
2372 );
2373
2374 let off = make_session(false);
2375 let on = make_session(true);
2376 kiss_bench::measure_pair(
2377 (
2378 "agent_context_build_subagents_off",
2379 "agent_context_build_subagents_on",
2380 ),
2381 21,
2382 10_000,
2383 (
2384 "empty_session_4_base_tools_0_control_tools",
2385 "empty_session_4_base_tools_6_control_tools",
2386 ),
2387 || off.build_context(),
2388 || on.build_context(),
2389 );
2390 }
2391
2392 #[test]
2393 #[ignore = "release-mode performance benchmark"]
2394 fn benchmark_performance_workflow_tool_exposure() {
2395 let registry = Registry::from_builtin();
2399 let model = registry.all().first().expect("built-in model").clone();
2400 let tools = benchmark_tools();
2401 let make_session = || {
2402 let mut settings = Settings::default();
2403 settings.subagents.enabled = true;
2404 AgentSession::new_with_subagents_allowed(
2405 SessionManager::in_memory(std::path::Path::new("/synthetic")),
2406 tools.clone(),
2407 registry.clone(),
2408 settings,
2409 "benchmark root prompt".into(),
2410 model.clone(),
2411 ThinkingLevel::Off,
2412 None,
2413 Arc::new(|_| {}),
2414 true,
2415 )
2416 };
2417
2418 let ordinary = make_session();
2419 let workflow = make_session();
2420 kiss_bench::measure_pair(
2421 (
2422 "agent_context_build_workflow_disarmed",
2423 "agent_context_build_workflow_armed",
2424 ),
2425 21,
2426 10_000,
2427 (
2428 "empty_session_subagents_on_workflow_disarmed",
2429 "empty_session_subagents_on_workflow_armed",
2430 ),
2431 || ordinary.build_context_for(PromptMode::Ordinary),
2432 || workflow.build_context_for(PromptMode::Workflow),
2433 );
2434 }
2435
2436 #[test]
2437 fn transcript_excerpt_keeps_only_recent_user_and_assistant_text() {
2438 let messages = vec![
2439 AgentMessage::user("old"),
2440 AgentMessage::BashExecution(kiss_agent::BashExecutionMessage {
2441 command: "pwd".into(),
2442 output: "ignored".into(),
2443 exit_code: Some(0),
2444 cancelled: false,
2445 truncated: false,
2446 full_output_path: None,
2447 exclude_from_context: false,
2448 timestamp: 1,
2449 }),
2450 AgentMessage::user("new"),
2451 ];
2452 let excerpt = transcript_excerpt(&messages, 1, 100);
2453 assert_eq!(excerpt, "User: new");
2454 }
2455
2456 #[test]
2457 fn transcript_excerpt_enforces_character_budget_from_the_tail() {
2458 let excerpt = transcript_excerpt(&[AgentMessage::user("abcdefghij")], 4, 5);
2459 assert!(excerpt.ends_with("fghij"));
2460 assert!(excerpt.starts_with("[earlier text omitted]"));
2461 }
2462
2463 #[test]
2464 fn hybrid_compaction_keeps_local_summary_and_remote_details() {
2465 let selected = select_compaction_outcome(
2466 &openai_model(),
2467 Ok(compaction::SummaryOutcome {
2468 summary: "portable".into(),
2469 usage: None,
2470 }),
2471 Some(Ok(remote_result())),
2472 )
2473 .unwrap();
2474 assert_eq!(selected.summary, "portable");
2475 assert_eq!(selected.remote_usage.unwrap().input, 10);
2476 assert_eq!(
2477 selected.remote_details.unwrap()["remoteCompaction"]["version"],
2478 2
2479 );
2480 }
2481
2482 #[test]
2483 fn remote_failure_falls_back_to_local_compaction() {
2484 let selected = select_compaction_outcome(
2485 &openai_model(),
2486 Ok(compaction::SummaryOutcome {
2487 summary: "portable".into(),
2488 usage: None,
2489 }),
2490 Some(Err(anyhow::anyhow!("remote unavailable"))),
2491 )
2492 .unwrap();
2493 assert_eq!(selected.summary, "portable");
2494 assert!(selected.remote_details.is_none());
2495 }
2496
2497 #[test]
2498 fn remote_success_survives_local_summary_failure() {
2499 let selected = select_compaction_outcome(
2500 &openai_model(),
2501 Err(anyhow::anyhow!("summary unavailable")),
2502 Some(Ok(remote_result())),
2503 )
2504 .unwrap();
2505 assert!(
2506 selected
2507 .summary
2508 .contains("server-side compaction was applied")
2509 );
2510 assert!(selected.remote_details.is_some());
2511 }
2512
2513 #[test]
2514 fn details_merge_keeps_file_operations_and_remote_artifact() {
2515 let merged = merge_compaction_details(
2516 serde_json::json!({"readFiles": ["a.rs"], "modifiedFiles": []}),
2517 Some(serde_json::json!({"remoteCompaction": {"version": 2}})),
2518 );
2519 assert_eq!(merged["readFiles"][0], "a.rs");
2520 assert_eq!(merged["remoteCompaction"]["version"], 2);
2521 }
2522
2523 #[test]
2524 fn auto_compaction_guard_checks_settings_threshold_and_cancel() {
2525 let mut settings = Settings::default();
2526 settings.compaction.reserve_tokens = 20;
2527 let messages = vec![AgentMessage::user("x".repeat(360))];
2528 let mut model = openai_model();
2529 model.context_window = 100;
2530 assert!(auto_compaction_needed(&settings, &messages, &model, false));
2531 assert!(!auto_compaction_needed(&settings, &messages, &model, true));
2532 settings.compaction.enabled = false;
2533 assert!(!auto_compaction_needed(&settings, &messages, &model, false));
2534 }
2535
2536 #[test]
2537 fn session_title_normalization_is_safe_and_bounded() {
2538 assert_eq!(
2539 normalize_session_title(" `Fix AUTH-123 login flow!` \nignored").as_deref(),
2540 Some("Fix AUTH-123 login flow")
2541 );
2542 assert_eq!(normalize_session_title("\n\t"), None);
2543 assert_eq!(
2544 normalize_session_title("🚀".repeat(50).as_str())
2545 .unwrap()
2546 .chars()
2547 .count(),
2548 SESSION_TITLE_MAX_CHARS
2549 );
2550 }
2551
2552 #[test]
2553 fn session_title_prompt_is_utf8_safe_and_bounded() {
2554 let prompt = "🚀".repeat(SESSION_TITLE_PROMPT_MAX_BYTES);
2555 let bounded = bounded_session_title_prompt(&prompt);
2556 assert!(bounded.len() <= SESSION_TITLE_PROMPT_MAX_BYTES);
2557 assert!(std::str::from_utf8(bounded.as_bytes()).is_ok());
2558 }
2559
2560 #[test]
2561 fn cache_warming_uses_ninety_percent_with_ten_second_margin() {
2562 assert_eq!(
2563 cache_warming_delay(std::time::Duration::from_secs(300)),
2564 Some(std::time::Duration::from_secs(270))
2565 );
2566 assert_eq!(
2567 cache_warming_delay(std::time::Duration::from_secs(60)),
2568 Some(std::time::Duration::from_secs(50))
2569 );
2570 assert_eq!(
2571 cache_warming_delay(std::time::Duration::from_secs(10)),
2572 None
2573 );
2574 }
2575}