use crate::compaction::{
self, estimate_context_tokens, extract_file_ops, file_ops_details, plan_compaction,
should_compact,
};
use crate::iterative::IterativeRuntime;
use crate::session::manager::SessionManager;
use crate::settings::{CacheWarmingMode, CompactionMode, QueueMode, ReasoningEffortMode, Settings};
use crate::subagents::{ForkTurns, SUBAGENT_SYSTEM_PROMPT, SubagentRuntime, fork_messages};
use crate::workflows::{WorkflowApprover, WorkflowRuntime};
use anyhow::Context as _;
use kiss_agent::{
AgentContext, AgentEvent, AgentLoopConfig, AgentMessage, DynTool, EventSink, StreamFn,
TurnUpdate,
};
use kiss_ai::{Model, Registry, StopReason, ThinkingLevel, Usage};
use std::collections::VecDeque;
use std::sync::{Arc, Mutex, OnceLock};
use tokio_util::sync::CancellationToken;
#[derive(Debug, Clone)]
pub enum SessionEvent {
Agent(Box<AgentEvent>),
QueueUpdate {
steering: Vec<String>,
follow_up: Vec<String>,
},
CompactionStart {
auto: bool,
},
CompactionEnd {
summary: String,
tokens_before: u64,
error: Option<String>,
},
Retry {
attempt: u32,
max: u32,
delay_ms: u64,
error: String,
},
ModelChanged {
provider: String,
model_id: String,
},
ReasoningEffortChanged {
level: ThinkingLevel,
generations: u8,
},
ReasoningEffortFallback {
level: ThinkingLevel,
reason: String,
},
Workflow {
run: crate::workflows::RunId,
version: u64,
},
WorkflowOutcome {
run: Option<crate::workflows::RunId>,
name: String,
status: WorkflowTurnStatus,
},
Iterative {
job: crate::iterative::JobId,
version: u64,
},
}
pub type SessionEventSink = Arc<dyn Fn(SessionEvent) + Send + Sync>;
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
pub enum PromptMode {
#[default]
Ordinary,
Workflow,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum WorkflowTurnStatus {
Cancelled,
Completed,
Failed,
Stopped,
}
#[derive(Clone)]
struct QueuedPrompt {
message: AgentMessage,
mode: PromptMode,
}
pub struct TreeNavigationOutcome {
pub editor_text: Option<String>,
pub summarized: bool,
pub cancelled: bool,
}
#[derive(Debug, Clone)]
pub struct EphemeralResponse {
pub text: String,
pub usage: Usage,
}
const SESSION_TITLE_MAX_CHARS: usize = 36;
const SESSION_TITLE_PROMPT_MAX_BYTES: usize = 960;
fn bounded_session_title_prompt(prompt: &str) -> &str {
let prompt = prompt.trim();
let mut end = prompt.len().min(SESSION_TITLE_PROMPT_MAX_BYTES);
while !prompt.is_char_boundary(end) {
end -= 1;
}
&prompt[..end]
}
fn normalize_session_title(text: &str) -> Option<String> {
let line = text.lines().find(|line| !line.trim().is_empty())?;
let cleaned: String = line
.chars()
.filter(|character| !character.is_control())
.collect();
let normalized = cleaned
.trim()
.trim_matches(|character| matches!(character, '"' | '\'' | '`' | '“' | '”' | '‘' | '’'))
.split_whitespace()
.collect::<Vec<_>>()
.join(" ");
let title: String = normalized.chars().take(SESSION_TITLE_MAX_CHARS).collect();
let title = title
.trim_end_matches(['.', '?', '!', ',', ':', ';'])
.trim();
(!title.is_empty()).then(|| title.to_string())
}
pub struct AgentSession {
pub manager: Mutex<SessionManager>,
pub registry: Arc<Registry>,
base_tools: Mutex<Vec<DynTool>>,
session_tools: Mutex<Vec<DynTool>>,
tools: Mutex<Vec<DynTool>>,
settings: Mutex<Settings>,
system_prompt: Mutex<String>,
model: Mutex<Model>,
thinking: Mutex<ThinkingLevel>,
fast_mode: Mutex<bool>,
steering: Arc<Mutex<VecDeque<QueuedPrompt>>>,
follow_up: Arc<Mutex<VecDeque<QueuedPrompt>>>,
cancel: Mutex<CancellationToken>,
cache_warm_cancel: Mutex<CancellationToken>,
running: Mutex<bool>,
totals: Mutex<Usage>,
context_usage_cache: Mutex<Option<(u64, u64)>>,
api_key_override: Option<(String, String)>,
sink: SessionEventSink,
subagents_allowed: bool,
subagents: OnceLock<Arc<SubagentRuntime>>,
workflows: OnceLock<Arc<WorkflowRuntime>>,
iterative: OnceLock<Arc<IterativeRuntime>>,
workflow_approver: Mutex<Option<WorkflowApprover>>,
stream_fn: Mutex<Option<StreamFn>>,
}
struct ReasoningRunState {
lease: crate::jev::ReasoningLease,
provider: String,
model_id: String,
saved_effort: ThinkingLevel,
}
impl ReasoningRunState {
fn checkpoint(
&mut self,
model: &Model,
saved_effort: ThinkingLevel,
current_effort: ThinkingLevel,
dynamic_enabled: bool,
) -> (bool, Option<ThinkingLevel>) {
let model_changed = self.provider != model.provider || self.model_id != model.id;
if model_changed || self.saved_effort != saved_effort || !dynamic_enabled {
self.lease.clear();
}
self.provider.clone_from(&model.provider);
self.model_id.clone_from(&model.id);
self.saved_effort = saved_effort;
if self
.lease
.current()
.is_some_and(|level| level != current_effort)
{
self.lease.clear();
}
(model_changed, self.lease.current())
}
}
impl AgentSession {
pub(crate) fn emit_iterative(&self, job: crate::iterative::JobId, version: u64) {
(self.sink)(SessionEvent::Iterative { job, version });
}
pub(crate) fn emit_workflow(&self, run: crate::workflows::RunId, version: u64) {
(self.sink)(SessionEvent::Workflow { run, version });
}
pub(crate) fn emit_workflow_outcome(
&self,
run: Option<crate::workflows::RunId>,
name: String,
status: WorkflowTurnStatus,
) {
(self.sink)(SessionEvent::WorkflowOutcome { run, name, status });
}
#[allow(clippy::too_many_arguments)]
pub fn new(
manager: SessionManager,
tools: Vec<DynTool>,
registry: impl Into<Arc<Registry>>,
settings: Settings,
system_prompt: String,
model: Model,
thinking: ThinkingLevel,
api_key_override: Option<(String, String)>,
sink: SessionEventSink,
) -> Arc<Self> {
Self::new_with_subagents_allowed(
manager,
tools,
registry,
settings,
system_prompt,
model,
thinking,
api_key_override,
sink,
true,
)
}
#[allow(clippy::too_many_arguments)]
pub fn new_with_subagents_allowed(
manager: SessionManager,
tools: Vec<DynTool>,
registry: impl Into<Arc<Registry>>,
settings: Settings,
system_prompt: String,
model: Model,
thinking: ThinkingLevel,
api_key_override: Option<(String, String)>,
sink: SessionEventSink,
subagents_allowed: bool,
) -> Arc<Self> {
let totals = manager.usage_totals();
let session = Arc::new(AgentSession {
manager: Mutex::new(manager),
registry: registry.into(),
base_tools: Mutex::new(tools.clone()),
session_tools: Mutex::new(Vec::new()),
tools: Mutex::new(tools),
settings: Mutex::new(settings),
system_prompt: Mutex::new(system_prompt),
model: Mutex::new(model),
thinking: Mutex::new(thinking),
fast_mode: Mutex::new(false),
steering: Default::default(),
follow_up: Default::default(),
cancel: Mutex::new(CancellationToken::new()),
cache_warm_cancel: Mutex::new(CancellationToken::new()),
running: Mutex::new(false),
totals: Mutex::new(totals),
context_usage_cache: Default::default(),
api_key_override,
sink,
subagents_allowed,
subagents: OnceLock::new(),
workflows: OnceLock::new(),
iterative: OnceLock::new(),
workflow_approver: Mutex::new(None),
stream_fn: Mutex::new(None),
});
if subagents_allowed {
let runtime = SubagentRuntime::new(Arc::downgrade(&session));
assert!(session.subagents.set(runtime).is_ok());
let workflows = WorkflowRuntime::new(Arc::downgrade(&session));
assert!(session.workflows.set(workflows).is_ok());
let iterative = IterativeRuntime::new(Arc::downgrade(&session));
assert!(session.iterative.set(iterative).is_ok());
}
session.rebuild_tools();
session
}
pub fn set_stream_fn(&self, stream_fn: Option<StreamFn>) {
*self.stream_fn.lock().unwrap() = stream_fn;
}
pub fn model(&self) -> Model {
self.model.lock().unwrap().clone()
}
pub fn thinking_level(&self) -> ThinkingLevel {
*self.thinking.lock().unwrap()
}
pub fn fast_mode(&self) -> bool {
*self.fast_mode.lock().unwrap()
}
pub fn set_fast_mode(&self, enabled: bool) {
*self.fast_mode.lock().unwrap() = enabled;
}
pub fn totals(&self) -> Usage {
*self.totals.lock().unwrap()
}
pub(crate) fn record_subagent_usage(&self, usage: Usage) {
self.totals.lock().unwrap().add(&usage);
}
pub fn is_running(&self) -> bool {
*self.running.lock().unwrap()
}
pub fn settings(&self) -> Settings {
self.settings.lock().unwrap().clone()
}
pub fn update_settings(&self, settings: Settings) {
let was_enabled = self.subagents_enabled();
let workflows_were_enabled = self.workflows_enabled();
let cache_warming = settings.cache_warming;
*self.settings.lock().unwrap() = settings;
if cache_warming == CacheWarmingMode::Off
|| (cache_warming == CacheWarmingMode::Streaming && !self.is_running())
{
self.cache_warm_cancel.lock().unwrap().cancel();
}
let is_enabled = self.subagents_enabled();
if was_enabled && !is_enabled {
self.stop_child_work();
} else if workflows_were_enabled && !self.workflows_enabled() {
self.stop_workflows();
}
self.rebuild_tools();
}
pub fn reload_runtime(&self, settings: Settings, system_prompt: String, tools: Vec<DynTool>) {
let was_enabled = self.subagents_enabled();
let workflows_were_enabled = self.workflows_enabled();
let cache_warming = settings.cache_warming;
*self.settings.lock().unwrap() = settings;
*self.system_prompt.lock().unwrap() = system_prompt;
*self.base_tools.lock().unwrap() = tools;
if cache_warming == CacheWarmingMode::Off
|| (cache_warming == CacheWarmingMode::Streaming && !self.is_running())
{
self.cache_warm_cancel.lock().unwrap().cancel();
}
let is_enabled = self.subagents_enabled();
if was_enabled && !is_enabled {
self.stop_child_work();
} else if workflows_were_enabled && !self.workflows_enabled() {
self.stop_workflows();
}
self.rebuild_tools();
}
pub fn install_session_tool(&self, tool: DynTool) {
let mut tools = self.session_tools.lock().unwrap();
if let Some(existing) = tools
.iter_mut()
.find(|existing| existing.name() == tool.name())
{
*existing = tool;
} else {
tools.push(tool);
}
drop(tools);
self.rebuild_tools();
}
pub fn remove_session_tool(&self, name: &str) -> bool {
let mut tools = self.session_tools.lock().unwrap();
let previous_len = tools.len();
tools.retain(|tool| tool.name() != name);
let removed = tools.len() != previous_len;
drop(tools);
if removed {
self.rebuild_tools();
}
removed
}
fn stop_child_work(&self) {
if let Some(runtime) = self.subagents.get() {
runtime.interrupt_all();
}
if let Some(runtime) = self.workflows.get() {
runtime.stop_all();
}
if let Some(runtime) = self.iterative.get() {
runtime.stop_all();
}
}
fn stop_workflows(&self) {
if let Some(runtime) = self.workflows.get() {
runtime.stop_all();
}
}
fn subagents_enabled(&self) -> bool {
self.subagents_allowed && self.settings.lock().unwrap().subagents.enabled
}
pub fn workflows_enabled(&self) -> bool {
self.subagents_enabled() && self.settings.lock().unwrap().workflows.enabled
}
pub fn prompt_mode_for(&self, text: &str) -> PromptMode {
if self.workflows_enabled()
&& self.settings.lock().unwrap().workflows.keyword_trigger
&& crate::workflows::workflow_trigger(text).is_some()
{
PromptMode::Workflow
} else {
PromptMode::Ordinary
}
}
pub fn workflows(&self) -> Option<Arc<WorkflowRuntime>> {
self.workflows.get().cloned()
}
pub fn iterative_jobs(&self) -> Option<Arc<IterativeRuntime>> {
self.iterative.get().cloned()
}
pub fn set_workflow_approver(&self, approver: WorkflowApprover) {
*self.workflow_approver.lock().unwrap() = Some(approver);
}
pub(crate) fn workflow_approver(&self) -> Option<WorkflowApprover> {
self.workflow_approver.lock().unwrap().clone()
}
fn rebuild_tools(&self) {
let mut tools = self.base_tools.lock().unwrap().clone();
for session_tool in self.session_tools.lock().unwrap().iter().cloned() {
if let Some(existing) = tools
.iter_mut()
.find(|existing| existing.name() == session_tool.name())
{
*existing = session_tool;
} else {
tools.push(session_tool);
}
}
if self.subagents_enabled()
&& let Some(runtime) = self.subagents.get()
{
tools.extend(runtime.control_tools());
}
*self.tools.lock().unwrap() = tools;
}
fn tools_for(&self, mode: PromptMode) -> Vec<DynTool> {
let mut tools = self.tools.lock().unwrap().clone();
if mode == PromptMode::Workflow
&& self.workflows_enabled()
&& let Some(runtime) = self.workflows.get()
{
tools.push(runtime.tool());
}
tools
}
pub fn available_tool_names(&self) -> Vec<String> {
self.tools
.lock()
.unwrap()
.iter()
.map(|tool| tool.name().to_string())
.collect()
}
#[cfg(test)]
fn available_tool_names_for(&self, mode: PromptMode) -> Vec<String> {
self.tools_for(mode)
.iter()
.map(|tool| tool.name().to_string())
.collect()
}
pub fn replace_manager(&self, manager: SessionManager) {
self.cache_warm_cancel.lock().unwrap().cancel();
if let Some(runtime) = self.subagents.get() {
runtime.reset();
}
if let Some(runtime) = self.workflows.get() {
runtime.stop_all();
}
if let Some(runtime) = self.iterative.get() {
runtime.stop_all();
}
let totals = manager.usage_totals();
let context = manager.build_session_context();
if let Some((provider, model_id)) = context.model
&& let Some((model, _)) = self.registry.resolve(&model_id, Some(&provider))
{
*self.model.lock().unwrap() = model;
}
if let Some(thinking) = context.thinking_level {
*self.thinking.lock().unwrap() = thinking;
}
*self.manager.lock().unwrap() = manager;
*self.context_usage_cache.lock().unwrap() = None;
*self.totals.lock().unwrap() = totals;
}
pub fn set_model(&self, model: Model) {
self.cache_warm_cancel.lock().unwrap().cancel();
{
let mut m = self.manager.lock().unwrap();
let _ = m.append_model_change(&model.provider, &model.id);
}
(self.sink)(SessionEvent::ModelChanged {
provider: model.provider.clone(),
model_id: model.id.clone(),
});
*self.model.lock().unwrap() = model;
}
pub fn set_thinking_level(&self, level: ThinkingLevel) {
let _ = self
.manager
.lock()
.unwrap()
.append_thinking_level_change(level);
*self.thinking.lock().unwrap() = level;
}
pub async fn navigate_tree(
self: &Arc<Self>,
target_id: &str,
summarize: bool,
custom_instructions: Option<String>,
cancel: CancellationToken,
) -> anyhow::Result<TreeNavigationOutcome> {
let (old_leaf, new_leaf, editor_text, abandoned_messages) = {
let manager = self.manager.lock().unwrap();
let target = manager
.get_entry(target_id)
.cloned()
.ok_or_else(|| anyhow::anyhow!("unknown tree entry {target_id}"))?;
let old_leaf = manager.leaf_id().map(str::to_string);
if old_leaf.as_deref() == Some(target_id) {
return Ok(TreeNavigationOutcome {
editor_text: None,
summarized: false,
cancelled: false,
});
}
let (new_leaf, editor_text) = match &target {
crate::session::entry::SessionEntry::Message {
message: AgentMessage::User(user),
..
} => (
target.parent_id().map(str::to_string),
Some(user.content.as_text()),
),
_ => (Some(target_id.to_string()), None),
};
let target_ancestors: std::collections::HashSet<String> = new_leaf
.as_deref()
.map(|leaf| {
manager
.branch_entries(Some(leaf))
.into_iter()
.map(|entry| entry.id().to_string())
.collect()
})
.unwrap_or_default();
let common_ancestor = old_leaf.as_deref().and_then(|leaf| {
manager
.branch_entries(Some(leaf))
.into_iter()
.rev()
.find(|entry| target_ancestors.contains(entry.id()))
.map(|entry| entry.id().to_string())
});
let abandoned_messages = old_leaf
.as_deref()
.map(|leaf| manager.branch_messages_after(leaf, common_ancestor.as_deref()))
.unwrap_or_default();
(old_leaf, new_leaf, editor_text, abandoned_messages)
};
let summary = if summarize && !abandoned_messages.is_empty() {
let model = self.model();
let credential = self.resolve_credential(&model.provider).await;
let serialized = compaction::serialize_agent_messages(&abandoned_messages);
let result = compaction::generate_summary(
&model,
credential,
&serialized,
None,
custom_instructions.as_deref(),
false,
cancel.clone(),
)
.await?;
if cancel.is_cancelled() {
return Ok(TreeNavigationOutcome {
editor_text: None,
summarized: false,
cancelled: true,
});
}
Some(result)
} else {
None
};
let mut manager = self.manager.lock().unwrap();
if manager.leaf_id() != old_leaf.as_deref() {
anyhow::bail!("the session tree changed during navigation");
}
let summarized = if let Some(summary) = summary {
let (read, modified) = extract_file_ops(&abandoned_messages);
if let Some(usage) = &summary.usage {
self.totals.lock().unwrap().add(usage);
}
manager.branch_with_summary(
new_leaf.as_deref(),
old_leaf.as_deref().unwrap_or(target_id),
summary.summary,
summary.usage,
Some(file_ops_details(&read, &modified)),
)?;
true
} else {
if let Some(new_leaf) = new_leaf.as_deref() {
manager.branch(new_leaf)?;
} else {
manager.reset_leaf();
}
false
};
Ok(TreeNavigationOutcome {
editor_text,
summarized,
cancelled: false,
})
}
pub fn queue_steering(&self, message: AgentMessage) {
self.queue_steering_with_mode(message, PromptMode::Ordinary);
}
pub fn queue_steering_with_mode(&self, message: AgentMessage, mode: PromptMode) {
self.steering
.lock()
.unwrap()
.push_back(QueuedPrompt { message, mode });
self.emit_queues();
}
pub fn queue_follow_up(&self, message: AgentMessage) {
self.queue_follow_up_with_mode(message, PromptMode::Ordinary);
}
pub fn queue_follow_up_with_mode(&self, message: AgentMessage, mode: PromptMode) {
self.follow_up
.lock()
.unwrap()
.push_back(QueuedPrompt { message, mode });
self.emit_queues();
}
pub fn reclaim_queued(&self) -> Vec<AgentMessage> {
let mut out: Vec<AgentMessage> = self
.steering
.lock()
.unwrap()
.drain(..)
.map(|prompt| prompt.message)
.collect();
out.extend(
self.follow_up
.lock()
.unwrap()
.drain(..)
.map(|prompt| prompt.message),
);
self.emit_queues();
out
}
pub fn abort(&self) {
self.cancel.lock().unwrap().cancel();
}
fn emit_queues(&self) {
let preview = |q: &VecDeque<QueuedPrompt>| {
q.iter()
.map(|prompt| match &prompt.message {
AgentMessage::User(u) => u.content.as_text().chars().take(80).collect(),
other => other.role().to_string(),
})
.collect::<Vec<String>>()
};
(self.sink)(SessionEvent::QueueUpdate {
steering: preview(&self.steering.lock().unwrap()),
follow_up: preview(&self.follow_up.lock().unwrap()),
});
}
async fn resolve_credential(&self, provider: &str) -> Option<kiss_ai::ResolvedCredential> {
if let Some((override_provider, key)) = &self.api_key_override
&& override_provider == provider
{
return Some(kiss_ai::ResolvedCredential::api_key(key));
}
let credential_provider = self.registry.credential_provider(provider);
kiss_ai::auth::resolve_credential_async(credential_provider, &self.registry.declared_keys)
.await
.ok()
.flatten()
}
fn loop_config(
&self,
session_arc: &Arc<Self>,
active_prompt_mode: Arc<Mutex<PromptMode>>,
) -> AgentLoopConfig {
let mut config = AgentLoopConfig::new(self.model());
config.thinking_level = self.thinking_level();
config.fast_mode = self.fast_mode();
config.session_id = Some(self.manager.lock().unwrap().session_id().to_string());
let settings = self.settings();
config.transport = settings.transport;
let registry = self.registry.clone();
let declared = self.registry.declared_keys.clone();
let api_key_override = self.api_key_override.clone();
config.get_credential = Some(Arc::new(move |provider| {
let registry = registry.clone();
let declared = declared.clone();
let api_key_override = api_key_override.clone();
Box::pin(async move {
if let Some((override_provider, key)) = api_key_override
&& override_provider == provider
{
return Some(kiss_ai::ResolvedCredential::api_key(key));
}
let credential_provider = registry.credential_provider(&provider).to_string();
kiss_ai::auth::resolve_credential_async(&credential_provider, &declared)
.await
.ok()
.flatten()
})
}));
let steering = self.steering.clone();
let steering_for_mode = self.steering.clone();
let steering_mode = settings.steering_mode;
let session_for_queues = session_arc.clone();
config.get_steering_messages = Some(Arc::new(move || {
let drained = drain_queue(&steering, steering_mode);
session_for_queues.emit_queues();
Box::pin(async move { drained })
}));
let follow_up = self.follow_up.clone();
let follow_up_for_mode = self.follow_up.clone();
let follow_up_mode = settings.follow_up_mode;
let session_for_queues = session_arc.clone();
config.get_follow_up_messages = Some(Arc::new(move || {
let drained = drain_queue(&follow_up, follow_up_mode);
session_for_queues.emit_queues();
Box::pin(async move { drained })
}));
let reasoning_run = Arc::new(Mutex::new(ReasoningRunState {
lease: crate::jev::ReasoningLease::default(),
provider: config.model.provider.clone(),
model_id: config.model.id.clone(),
saved_effort: config.thinking_level,
}));
let session_for_reasoning = session_arc.clone();
let reasoning_for_generation = reasoning_run.clone();
let prompt_mode_for_reasoning = active_prompt_mode.clone();
config.prepare_generation = Some(Arc::new(move |current_effort| {
let session = session_for_reasoning.clone();
let reasoning_run = reasoning_for_generation.clone();
let active_prompt_mode = prompt_mode_for_reasoning.clone();
Box::pin(async move {
let settings = session.settings();
let model = session.model();
let saved_effort = session.thinking_level();
let dynamic_enabled =
settings.reasoning_effort.mode == ReasoningEffortMode::Jev && model.reasoning;
let (model_changed, leased) = {
let mut run = reasoning_run.lock().unwrap();
run.checkpoint(&model, saved_effort, current_effort, dynamic_enabled)
};
let mut update = TurnUpdate {
context: model_changed
.then(|| session.build_context_for(*active_prompt_mode.lock().unwrap())),
model: model_changed.then(|| model.clone()),
thinking_level: Some(saved_effort),
fast_mode: Some(session.fast_mode()),
..Default::default()
};
if !dynamic_enabled {
return Some(update);
}
if let Some(level) = leased {
update.thinking_level = Some(level);
return Some(update);
}
let selected = match kiss_ai::auth::resolve_api_key_async(
"typesafe",
&session.registry.declared_keys,
)
.await
{
Ok(Some(api_key)) => {
let messages = session
.manager
.lock()
.unwrap()
.build_session_context()
.messages;
let cancel = session.cancel.lock().unwrap().clone();
crate::jev::select_reasoning(&messages, &[], &model, &api_key, cancel)
.await
.map_err(|error| error.to_string())
}
Ok(None) => Err(
"TypeSafe credentials are unavailable; run /login typesafe or set TYPESAFE_API_KEY"
.into(),
),
Err(error) => Err(format!(
"TypeSafe credentials could not be read: {error}; run /login typesafe or set TYPESAFE_API_KEY"
)),
};
let latest_model = session.model();
let latest_effort = session.thinking_level();
if latest_model.provider != model.provider
|| latest_model.id != model.id
|| latest_effort != saved_effort
|| session.settings().reasoning_effort.mode != ReasoningEffortMode::Jev
{
reasoning_run.lock().unwrap().lease.clear();
update.context =
Some(session.build_context_for(*active_prompt_mode.lock().unwrap()));
update.model = Some(latest_model);
update.thinking_level = Some(latest_effort);
return Some(update);
}
match selected {
Ok(selection) => {
update.thinking_level = Some(selection.level);
if selection.level != current_effort {
let sink = session.sink.clone();
let level = selection.level;
let generations = selection.generations;
update.on_applied = Some(Box::new(move |_, applied| {
if applied == level {
sink(SessionEvent::ReasoningEffortChanged {
level,
generations,
});
}
}));
}
reasoning_run.lock().unwrap().lease.install(&selection);
}
Err(reason) => {
let sink = session.sink.clone();
let reason = reason.chars().take(300).collect();
update.on_applied = Some(Box::new(move |_, applied| {
sink(SessionEvent::ReasoningEffortFallback {
level: applied,
reason,
});
}));
}
}
Some(update)
})
}));
let session_for_compaction = session_arc.clone();
config.prepare_next_turn = Some(Arc::new(move |turn| {
let has_tool_results = !turn.tool_results.is_empty();
let tool_failed = turn.tool_results.iter().any(|result| result.is_error);
let will_continue = turn.will_continue;
let session = session_for_compaction.clone();
let active_prompt_mode = active_prompt_mode.clone();
let steering_for_mode = steering_for_mode.clone();
let follow_up_for_mode = follow_up_for_mode.clone();
let reasoning_run = reasoning_run.clone();
Box::pin(async move {
let has_queued_user_input = !steering_for_mode.lock().unwrap().is_empty()
|| (!will_continue && !follow_up_for_mode.lock().unwrap().is_empty());
let queued_mode = queued_mode(&steering_for_mode, steering_mode)
.or_else(|| queued_mode(&follow_up_for_mode, follow_up_mode));
let mode_changed = {
let mut active = active_prompt_mode.lock().unwrap();
let before = *active;
if let Some(queued_mode) = queued_mode {
*active = queued_mode;
} else if *active == PromptMode::Workflow && !has_tool_results {
*active = PromptMode::Ordinary;
}
before != *active
};
let mut context_changed = mode_changed;
if has_tool_results {
let settings = session.settings();
let cancel = session.cancel.lock().unwrap().clone();
let model = session.model();
let revision_before = {
let manager = session.manager.lock().unwrap();
let context = manager.build_session_context();
if !auto_compaction_needed(
&settings,
&context.messages,
&model,
cancel.is_cancelled(),
) {
None
} else {
Some(manager.context_revision())
}
};
if let Some(revision_before) = revision_before {
session.compact(None, true).await;
let revision_after = session.manager.lock().unwrap().context_revision();
context_changed |= revision_after != revision_before;
}
}
{
let mut run = reasoning_run.lock().unwrap();
run.lease.consume_generation();
if tool_failed || has_queued_user_input {
run.lease.clear();
}
}
Some(TurnUpdate {
context: context_changed
.then(|| session.build_context_for(*active_prompt_mode.lock().unwrap())),
fast_mode: Some(session.fast_mode()),
..Default::default()
})
})
}));
if let Some(stream_fn) = self.stream_fn.lock().unwrap().clone() {
config.stream_fn = stream_fn;
}
config
}
fn build_context(&self) -> AgentContext {
self.build_context_for(PromptMode::Ordinary)
}
fn build_context_for(&self, prompt_mode: PromptMode) -> AgentContext {
let model = self.model();
let manager = self.manager.lock().unwrap();
let (openai_responses_input, messages) =
if kiss_ai::api::openai_compaction::supports_remote_compaction(&model)
&& let Some(remote) = manager.build_openai_compaction_context(&model)
{
(Some(remote.replacement_history), remote.messages)
} else {
(None, manager.build_session_context().messages)
};
drop(manager);
let mut system_prompt = self.system_prompt.lock().unwrap().clone();
if self.subagents_enabled() {
system_prompt.push_str("\n\n");
system_prompt.push_str(SUBAGENT_SYSTEM_PROMPT);
}
if prompt_mode == PromptMode::Workflow
&& self.workflows_enabled()
&& let Some(runtime) = self.workflows.get()
{
let limits = runtime.limits();
let size = self.settings.lock().unwrap().workflows.size;
system_prompt.push_str("\n\n");
system_prompt.push_str(&crate::workflows::authoring_prompt(
size,
limits.max_agents,
limits.max_fanout,
));
}
AgentContext {
system_prompt,
openai_responses_input,
messages,
tools: self.tools_for(prompt_mode),
}
}
pub(crate) fn create_subagent_session(
self: &Arc<Self>,
task_name: &str,
canonical_path: &str,
fork_turns: ForkTurns,
model_pattern: Option<&str>,
reasoning_effort: Option<&str>,
) -> anyhow::Result<Arc<Self>> {
let (model, suggested_thinking) = match model_pattern {
Some(pattern) => self
.registry
.resolve(pattern, None)
.with_context(|| format!("no model matches subagent model '{pattern}'"))?,
None => (self.model(), None),
};
let thinking = match reasoning_effort {
Some(level) => ThinkingLevel::parse(level)
.with_context(|| format!("unknown subagent reasoning_effort '{level}'"))?,
None => suggested_thinking.unwrap_or_else(|| self.thinking_level()),
};
let (mut manager, parent_messages, parent_id) = {
let parent = self.manager.lock().unwrap();
(
parent.create_child()?,
if fork_turns == ForkTurns::None {
Vec::new()
} else {
parent.build_session_context().messages
},
parent.session_id().to_string(),
)
};
for message in fork_messages(&parent_messages, fork_turns) {
manager.append_message(message)?;
}
manager.append_custom(
"subagent",
Some(serde_json::json!({
"taskName": task_name,
"canonicalPath": canonical_path,
"parentSessionId": parent_id,
})),
)?;
let mut settings = self.settings();
settings.subagents.enabled = false;
let mut system_prompt = self.system_prompt.lock().unwrap().clone();
system_prompt.push_str(&format!(
"\n\nYou are child agent {canonical_path}. Complete only the assigned task. Return a concise result to the parent agent."
));
let child = Self::new_with_subagents_allowed(
manager,
self.base_tools.lock().unwrap().clone(),
self.registry.clone(),
settings,
system_prompt,
model,
thinking,
self.api_key_override.clone(),
Arc::new(|_| {}),
false,
);
child.set_stream_fn(self.stream_fn.lock().unwrap().clone());
Ok(child)
}
async fn run_ephemeral(
self: &Arc<Self>,
system_prompt: String,
prompt: String,
tools: Vec<DynTool>,
max_tokens: u64,
cancel: CancellationToken,
) -> anyhow::Result<EphemeralResponse> {
let mut config = self.loop_config(self, Arc::new(Mutex::new(PromptMode::Ordinary)));
config.thinking_level = ThinkingLevel::Off;
config.max_tokens = Some(max_tokens);
config.session_id = Some(format!("ephemeral-{}", uuid::Uuid::new_v4()));
config.get_steering_messages = None;
config.get_follow_up_messages = None;
config.prepare_next_turn = None;
config.prepare_generation = None;
let context = AgentContext {
system_prompt,
openai_responses_input: None,
messages: Vec::new(),
tools,
};
let sink: EventSink = Arc::new(|_| {});
let messages = kiss_agent::run_agent_loop(
vec![AgentMessage::user(prompt)],
context,
config,
cancel.clone(),
sink,
)
.await;
if cancel.is_cancelled() {
anyhow::bail!("request cancelled");
}
let mut usage = Usage::default();
for message in &messages {
if let AgentMessage::Assistant(assistant) = message {
usage.add(&assistant.usage);
}
}
let assistant = messages.iter().rev().find_map(|message| match message {
AgentMessage::Assistant(assistant) => Some(assistant),
_ => None,
});
let Some(assistant) = assistant else {
anyhow::bail!("the provider returned no answer");
};
if assistant.stop_reason == StopReason::Error {
anyhow::bail!(
"{}",
assistant
.error_message
.as_deref()
.unwrap_or("the provider request failed")
);
}
let text = assistant.text();
if text.trim().is_empty() {
anyhow::bail!("the provider returned an empty answer");
}
self.totals.lock().unwrap().add(&usage);
Ok(EphemeralResponse { text, usage })
}
pub async fn answer_btw(
self: &Arc<Self>,
question: &str,
cancel: CancellationToken,
) -> anyhow::Result<EphemeralResponse> {
let messages = self
.manager
.lock()
.unwrap()
.build_session_context()
.messages;
let transcript = transcript_excerpt(&messages, 4, 4_000);
let prompt = if transcript.is_empty() {
format!("Side question:\n{question}")
} else {
format!("Recent session context:\n{transcript}\n\nSide question:\n{question}")
};
let read_tools = self
.tools
.lock()
.unwrap()
.iter()
.filter(|tool| tool.name() == "read")
.cloned()
.collect();
self.run_ephemeral(
"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(),
prompt,
read_tools,
500,
cancel,
)
.await
}
pub async fn generate_session_title(
self: &Arc<Self>,
prompt: &str,
cancel: CancellationToken,
) -> anyhow::Result<String> {
let prompt = bounded_session_title_prompt(prompt);
if prompt.is_empty() {
anyhow::bail!("the session prompt is empty");
}
let response = self
.run_ephemeral(
format!(
"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."
),
format!("User prompt:\n{prompt}"),
Vec::new(),
64,
cancel,
)
.await?;
normalize_session_title(&response.text)
.context("the provider returned an invalid session title")
}
pub async fn generate_recap(
self: &Arc<Self>,
previous_recap: Option<&str>,
cancel: CancellationToken,
) -> anyhow::Result<EphemeralResponse> {
let messages = self
.manager
.lock()
.unwrap()
.build_session_context()
.messages;
let transcript = transcript_excerpt(&messages, 12, 12_000);
if transcript.is_empty() {
anyhow::bail!("the session has no conversation to recap");
}
let previous = previous_recap
.filter(|recap| !recap.trim().is_empty())
.map(|recap| format!("\n\nPrevious recap:\n{recap}"))
.unwrap_or_default();
self.run_ephemeral(
"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(),
format!("Session transcript:\n{transcript}{previous}"),
Vec::new(),
160,
cancel,
)
.await
}
pub async fn prompt(self: &Arc<Self>, prompts: Vec<AgentMessage>) {
self.prompt_with_mode(prompts, PromptMode::Ordinary).await;
}
pub async fn prompt_with_mode(
self: &Arc<Self>,
prompts: Vec<AgentMessage>,
prompt_mode: PromptMode,
) {
self.cache_warm_cancel.lock().unwrap().cancel();
{
let mut running = self.running.lock().unwrap();
if *running {
drop(running);
for p in prompts {
self.queue_steering_with_mode(p, prompt_mode);
}
return;
}
*running = true;
}
let cancel = {
let mut guard = self.cancel.lock().unwrap();
*guard = CancellationToken::new();
guard.clone()
};
{
let mut manager = self.manager.lock().unwrap();
for p in &prompts {
let _ = manager.append_message(p.clone());
}
}
let session = self.clone();
let sink: EventSink = Arc::new(move |event: AgentEvent| {
session.on_agent_event(&event);
(session.sink)(SessionEvent::Agent(Box::new(event)));
});
let active_prompt_mode = Arc::new(Mutex::new(prompt_mode));
let mut config = self.loop_config(self, active_prompt_mode.clone());
let mut context = self.build_context_for(prompt_mode);
let mut attempt: u32 = 0;
loop {
let messages = kiss_agent::run_agent_loop_continue(
context,
config.clone(),
cancel.clone(),
sink.clone(),
)
.await;
let last_error = messages.iter().rev().find_map(|m| match m {
AgentMessage::Assistant(a) if a.stop_reason == StopReason::Error => {
Some(a.error_message.clone().unwrap_or_default())
}
_ => None,
});
let settings = self.settings();
let retry = &settings.retry;
if let Some(error) = last_error
&& retry.enabled
&& attempt < retry.max_retries
&& is_transient(&error)
&& !cancel.is_cancelled()
{
attempt += 1;
let delay = retry
.base_delay_ms
.saturating_mul(1u64.checked_shl(attempt - 1).unwrap_or(u64::MAX))
.min(retry.max_agent_delay_ms);
(self.sink)(SessionEvent::Retry {
attempt,
max: retry.max_retries,
delay_ms: delay,
error,
});
tokio::select! {
_ = tokio::time::sleep(std::time::Duration::from_millis(delay)) => {}
_ = cancel.cancelled() => break,
}
context = self.build_context_for(*active_prompt_mode.lock().unwrap());
while matches!(
context.messages.last(),
Some(AgentMessage::Assistant(a)) if a.stop_reason == StopReason::Error
) {
context.messages.pop();
}
config.model = self.model();
config.thinking_level = self.thinking_level();
config.fast_mode = self.fast_mode();
continue;
}
let ctx = self.manager.lock().unwrap().build_session_context();
if auto_compaction_needed(
&settings,
&ctx.messages,
&self.model(),
cancel.is_cancelled(),
) {
self.compact(None, true).await;
}
break;
}
*self.running.lock().unwrap() = false;
if self.settings.lock().unwrap().cache_warming == CacheWarmingMode::Streaming {
self.cache_warm_cancel.lock().unwrap().cancel();
}
}
fn on_agent_event(self: &Arc<Self>, event: &AgentEvent) {
match event {
AgentEvent::MessageEnd { message } => {
let persist = match message {
AgentMessage::Assistant(a) => {
self.schedule_cache_warming(a);
let mut totals = self.totals.lock().unwrap();
totals.add(&a.usage);
true
}
AgentMessage::ToolResult(_)
| AgentMessage::User(_)
| AgentMessage::Custom(_) => true,
_ => false,
};
if persist {
let mut manager = self.manager.lock().unwrap();
let duplicate = matches!(
(manager.entries().last(), message),
(Some(crate::session::entry::SessionEntry::Message { message: last, .. }), m) if last == m
);
if !duplicate {
let _ = manager.append_message(message.clone());
}
}
}
AgentEvent::AgentEnd { .. } => {}
_ => {}
}
}
fn schedule_cache_warming(self: &Arc<Self>, assistant: &kiss_ai::AssistantMessage) {
let settings = self.settings();
let model = self.model();
let Some(cache) = model.prompt_cache else {
return;
};
let Some(short_ttl) = cache.short else {
return;
};
let reasoning = self.thinking_level();
let fast_mode = self.fast_mode();
if settings.cache_warming == CacheWarmingMode::Off
|| matches!(
assistant.stop_reason,
StopReason::Error | StopReason::Aborted
)
|| (assistant.usage.input + assistant.usage.cache_read + assistant.usage.cache_write
== 0)
|| (reasoning != ThinkingLevel::Off
&& model.api == "anthropic-messages"
&& !model
.compat
.as_ref()
.and_then(|compat| compat.force_adaptive_thinking)
.unwrap_or(false))
{
return;
}
let ttl = std::time::Duration::from_secs(short_ttl);
let Some(delay) = cache_warming_delay(ttl) else {
return;
};
let prompt_tokens =
assistant.usage.input + assistant.usage.cache_read + assistant.usage.cache_write;
let agent_context = self.build_context();
let context = kiss_ai::Context {
system_prompt: Some(agent_context.system_prompt),
openai_responses_input: agent_context.openai_responses_input,
messages: kiss_agent::convert_to_llm(&agent_context.messages),
tools: agent_context
.tools
.iter()
.map(|tool| tool.to_def())
.collect(),
};
let cancel = CancellationToken::new();
{
let mut current = self.cache_warm_cancel.lock().unwrap();
current.cancel();
*current = cancel.clone();
}
let session = self.clone();
tokio::spawn(async move {
let started = tokio::time::Instant::now();
loop {
let scheduled = tokio::time::Instant::now();
if tokio::select! {
_ = tokio::time::sleep(delay) => false,
_ = cancel.cancelled() => true,
} {
return;
}
if cache_refresh_deadline_missed(scheduled.elapsed(), ttl, delay) {
return;
}
let idle = !session.is_running();
let max_age = if idle {
std::time::Duration::from_secs(30 * 60)
} else {
std::time::Duration::from_secs(60 * 60)
};
if started.elapsed() > max_age {
return;
}
let priced = |input: u64, output: u64, cache_read: u64, cache_write: u64| {
let mut usage = Usage {
input,
output,
cache_read,
cache_write,
..Default::default()
};
kiss_ai::api::finalize_cost(&mut usage, &model);
usage.cost.total
};
let hit_cost = priced(0, 0, prompt_tokens, 0);
let miss_cost = if model.cost.cache_write > 0.0 {
priced(0, 0, 0, prompt_tokens)
} else {
priced(prompt_tokens, 0, 0, 0)
};
let warm_cost = priced(0, 1, prompt_tokens, 0);
let probability = if idle { 0.15 } else { 1.0 };
if probability * (miss_cost - hit_cost).max(0.0) - warm_cost < 0.05 {
return;
}
let Some(credential) = session.resolve_credential(&model.provider).await else {
return;
};
let options = kiss_ai::StreamOptions {
credential: Some(credential),
max_tokens: Some(1),
reasoning,
fast_mode,
session_id: Some(session.manager.lock().unwrap().session_id().to_string()),
transport: settings.transport,
cancel: cancel.clone(),
..Default::default()
};
let stream_fn = session
.stream_fn
.lock()
.unwrap()
.clone()
.unwrap_or_else(|| Arc::new(kiss_ai::stream_simple));
let warmed = stream_fn(&model, &context, &options).result().await;
if matches!(warmed.stop_reason, StopReason::Error | StopReason::Aborted) {
return;
}
session.totals.lock().unwrap().add(&warmed.usage);
let _ = session.manager.lock().unwrap().append_usage(
"cache_warm",
&warmed.provider,
warmed.response_model.as_deref().unwrap_or(&warmed.model),
warmed.usage,
None,
);
}
});
}
pub async fn compact(self: &Arc<Self>, custom_instructions: Option<String>, auto: bool) {
(self.sink)(SessionEvent::CompactionStart { auto });
let ctx = self.manager.lock().unwrap().build_session_context();
let previous_summary = ctx.messages.iter().rev().find_map(|m| match m {
AgentMessage::CompactionSummary(c) => Some(c.summary.clone()),
_ => None,
});
let settings = self.settings();
let model = self.model();
let override_settings = settings
.compaction
.model_overrides
.get(&format!("{}/{}", model.provider, model.id));
let keep_recent_tokens = override_settings
.and_then(|value| value.keep_recent_tokens)
.unwrap_or(settings.compaction.keep_recent_tokens);
let reserve_tokens = override_settings
.and_then(|value| value.reserve_tokens)
.unwrap_or(settings.compaction.reserve_tokens);
let plan = plan_compaction(&ctx.messages, keep_recent_tokens);
if plan.to_summarize.is_empty() && plan.turn_prefix.is_empty() {
(self.sink)(SessionEvent::CompactionEnd {
summary: String::new(),
tokens_before: plan.tokens_before,
error: Some("Nothing to compact".into()),
});
return;
}
if settings.compaction.mode == CompactionMode::Jev
&& let Ok(Some(api_key)) =
kiss_ai::auth::resolve_api_key_async("typesafe", &self.registry.declared_keys).await
{
let cancel = self.cancel.lock().unwrap().clone();
let pinned_start = ctx.messages.len().saturating_sub(plan.kept.len());
if let Ok(result) =
crate::jev::compact(&ctx.messages, pinned_start, &api_key, cancel).await
{
let estimated_before: u64 = ctx
.messages
.iter()
.map(compaction::estimate_message_tokens)
.sum();
let estimated_after: u64 = result
.messages
.iter()
.map(compaction::estimate_message_tokens)
.sum();
let removed = estimated_before.saturating_sub(estimated_after);
let tokens_after = plan.tokens_before.saturating_sub(removed);
let useful =
estimated_after.saturating_mul(4) <= estimated_before.saturating_mul(3);
let resolved_auto_threshold =
!auto || !should_compact(tokens_after, model.context_window, reserve_tokens);
if useful && resolved_auto_threshold {
let summary = format!(
"Jev kept {}, truncated {}, and removed {} of {} older tool interactions",
result.stats.kept,
result.stats.truncated,
result.stats.dropped,
result.stats.eligible,
);
let details = serde_json::json!({
"mode": "jev",
"stats": result.stats,
"estimatedTokensBefore": plan.tokens_before,
"estimatedTokensAfter": tokens_after,
});
let mut manager = self.manager.lock().unwrap();
let append = manager.append_compaction(
String::new(),
plan.tokens_before,
result.messages,
None,
Some(details),
);
drop(manager);
let error = append.err().map(|error| format!("{error:#}"));
(self.sink)(SessionEvent::CompactionEnd {
summary: if error.is_none() {
summary
} else {
String::new()
},
tokens_before: plan.tokens_before,
error,
});
return;
}
}
}
let credential = self.resolve_credential(&model.provider).await;
let mut serialized = compaction::serialize_agent_messages(&plan.to_summarize);
if plan.is_split_turn {
serialized.push_str(
"\n\n[The following is the earlier part of the still-active task turn:]\n\n",
);
serialized.push_str(&compaction::serialize_agent_messages(&plan.turn_prefix));
}
let summary_cancel = self.cancel.lock().unwrap().clone();
let remote_request = if kiss_ai::api::openai_compaction::supports_remote_compaction(&model)
{
let context = self.build_context();
Some((
kiss_ai::Context {
system_prompt: Some(context.system_prompt),
openai_responses_input: context.openai_responses_input,
messages: kiss_agent::convert_to_llm(&context.messages),
tools: context.tools.iter().map(|tool| tool.to_def()).collect(),
},
kiss_ai::StreamOptions {
credential: credential.clone(),
reasoning: self.thinking_level(),
fast_mode: self.fast_mode(),
session_id: Some(self.manager.lock().unwrap().session_id().to_string()),
cancel: summary_cancel.clone(),
..Default::default()
},
))
} else {
None
};
let local_future = compaction::generate_summary(
&model,
credential.clone(),
&serialized,
previous_summary.as_deref(),
custom_instructions.as_deref(),
plan.is_split_turn,
summary_cancel.clone(),
);
let remote_model = model.clone();
let remote_future = async move {
match remote_request {
Some((context, options)) => Some(
kiss_ai::api::openai_compaction::compact(&remote_model, &context, &options)
.await,
),
None => None,
}
};
let (local_outcome, remote_outcome) = tokio::join!(local_future, remote_future);
match select_compaction_outcome(&model, local_outcome, remote_outcome) {
Ok(result) => {
let mut summarized_all = plan.to_summarize.clone();
summarized_all.extend(plan.turn_prefix.clone());
let (read, modified) = extract_file_ops(&summarized_all);
{
let mut totals = self.totals.lock().unwrap();
if let Some(u) = &result.local_usage {
totals.add(u);
}
if let Some(u) = &result.remote_usage {
totals.add(u);
}
}
let details = merge_compaction_details(
file_ops_details(&read, &modified),
result.remote_details,
);
let mut manager = self.manager.lock().unwrap();
let _ = manager.append_compaction(
result.summary.clone(),
plan.tokens_before,
plan.kept.clone(),
result.local_usage,
Some(details),
);
(self.sink)(SessionEvent::CompactionEnd {
summary: result.summary,
tokens_before: plan.tokens_before,
error: None,
});
}
Err(error) => {
(self.sink)(SessionEvent::CompactionEnd {
summary: String::new(),
tokens_before: plan.tokens_before,
error: Some(format!("{error:#}")),
});
}
}
}
pub fn context_usage(&self) -> (u64, u64) {
let manager = self.manager.lock().unwrap();
let revision = manager.context_revision();
let used = if let Some((cached_revision, tokens)) =
*self.context_usage_cache.lock().unwrap()
&& cached_revision == revision
{
tokens
} else {
let tokens = estimate_context_tokens(&manager.build_session_context().messages);
*self.context_usage_cache.lock().unwrap() = Some((revision, tokens));
tokens
};
drop(manager);
(used, self.model().context_window)
}
}
struct SelectedCompaction {
summary: String,
local_usage: Option<Usage>,
remote_usage: Option<Usage>,
remote_details: Option<serde_json::Value>,
}
fn select_compaction_outcome(
model: &Model,
local: anyhow::Result<compaction::SummaryOutcome>,
remote: Option<anyhow::Result<kiss_ai::api::openai_compaction::RemoteCompactionResult>>,
) -> anyhow::Result<SelectedCompaction> {
match (local, remote) {
(Ok(local), Some(Ok(remote))) => Ok(SelectedCompaction {
summary: local.summary,
local_usage: local.usage,
remote_usage: remote.usage,
remote_details: Some(
kiss_ai::api::openai_compaction::build_remote_compaction_details(model, &remote),
),
}),
(Ok(local), Some(Err(_)) | None) => Ok(SelectedCompaction {
summary: local.summary,
local_usage: local.usage,
remote_usage: None,
remote_details: None,
}),
(Err(_), Some(Ok(remote))) => Ok(SelectedCompaction {
summary: format!(
"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.",
model.provider, model.id
),
local_usage: None,
remote_usage: remote.usage,
remote_details: Some(
kiss_ai::api::openai_compaction::build_remote_compaction_details(model, &remote),
),
}),
(Err(local), Some(Err(remote))) => anyhow::bail!(
"local compaction failed: {local:#}. OpenAI remote compaction failed: {remote:#}"
),
(Err(error), None) => Err(error),
}
}
fn merge_compaction_details(
mut local: serde_json::Value,
remote: Option<serde_json::Value>,
) -> serde_json::Value {
let Some(remote) = remote else {
return local;
};
let Some(local_object) = local.as_object_mut() else {
return remote;
};
if let Some(remote_object) = remote.as_object() {
for (key, value) in remote_object {
local_object.insert(key.clone(), value.clone());
}
}
local
}
fn transcript_excerpt(messages: &[AgentMessage], max_messages: usize, max_chars: usize) -> String {
let mut entries = messages
.iter()
.rev()
.filter_map(|message| match message {
AgentMessage::User(user) => Some(("User", user.content.as_text())),
AgentMessage::Assistant(assistant) => Some(("Assistant", assistant.text())),
_ => None,
})
.filter(|(_, text)| !text.trim().is_empty())
.take(max_messages)
.collect::<Vec<_>>();
entries.reverse();
let transcript = entries
.into_iter()
.map(|(role, text)| format!("{role}: {}", text.trim()))
.collect::<Vec<_>>()
.join("\n\n");
let count = transcript.chars().count();
if count <= max_chars {
return transcript;
}
let omitted = count - max_chars;
let tail = transcript.chars().skip(omitted).collect::<String>();
format!("[earlier text omitted]\n{tail}")
}
fn drain_queue(queue: &Arc<Mutex<VecDeque<QueuedPrompt>>>, mode: QueueMode) -> Vec<AgentMessage> {
let mut q = queue.lock().unwrap();
match mode {
QueueMode::All => q.drain(..).map(|prompt| prompt.message).collect(),
QueueMode::OneAtATime => q
.pop_front()
.map(|prompt| prompt.message)
.into_iter()
.collect(),
}
}
fn queued_mode(queue: &Arc<Mutex<VecDeque<QueuedPrompt>>>, mode: QueueMode) -> Option<PromptMode> {
let queue = queue.lock().unwrap();
match mode {
QueueMode::All => queue
.iter()
.any(|prompt| prompt.mode == PromptMode::Workflow)
.then_some(PromptMode::Workflow)
.or_else(|| (!queue.is_empty()).then_some(PromptMode::Ordinary)),
QueueMode::OneAtATime => queue.front().map(|prompt| prompt.mode),
}
}
fn auto_compaction_needed(
settings: &Settings,
messages: &[AgentMessage],
model: &Model,
cancelled: bool,
) -> bool {
let reserve_tokens = settings
.compaction
.model_overrides
.get(&format!("{}/{}", model.provider, model.id))
.and_then(|value| value.reserve_tokens)
.unwrap_or(settings.compaction.reserve_tokens);
settings.compaction.enabled
&& !cancelled
&& model.context_window > 0
&& should_compact(
estimate_context_tokens(messages),
model.context_window,
reserve_tokens,
)
}
fn cache_warming_delay(ttl: std::time::Duration) -> Option<std::time::Duration> {
(ttl > std::time::Duration::from_secs(10))
.then(|| std::cmp::min(ttl.mul_f64(0.9), ttl - std::time::Duration::from_secs(10)))
}
fn cache_refresh_deadline_missed(
elapsed: std::time::Duration,
ttl: std::time::Duration,
delay: std::time::Duration,
) -> bool {
elapsed > delay + ttl.saturating_sub(delay) / 2
}
fn is_transient(error: &str) -> bool {
let e = error.to_lowercase();
let transient_status = [429, 500, 502, 503, 504, 520].iter().any(|status| {
[
format!("http {status}"),
format!("status {status}"),
format!("status: {status}"),
format!("status code {status}"),
]
.iter()
.any(|marker| e.contains(marker))
});
transient_status
|| [
"overloaded",
"currently experiencing high demand",
"rate limit",
"timeout",
"timed out",
"connection reset",
"connection refused",
"connection closed",
"connection aborted",
"connection error",
"failed to connect",
"network error",
"stream error",
"request failed",
]
.iter()
.any(|needle| e.contains(needle))
}
#[cfg(test)]
mod ephemeral_tests {
use super::*;
use std::collections::BTreeMap;
fn openai_model() -> Model {
Model {
id: "gpt-test".into(),
name: "GPT test".into(),
api: "openai-responses".into(),
provider: "openai".into(),
base_url: "https://api.openai.com/v1".into(),
reasoning: true,
input: vec!["text".into()],
cost: Default::default(),
prompt_cache: None,
context_window: 100_000,
max_tokens: 1_000,
compat: None,
thinking_level_map: BTreeMap::new(),
headers: BTreeMap::new(),
}
}
#[test]
fn reasoning_lease_ends_when_model_saved_effort_or_applied_effort_changes() {
let model = openai_model();
let mut run = ReasoningRunState {
lease: crate::jev::ReasoningLease::default(),
provider: model.provider.clone(),
model_id: model.id.clone(),
saved_effort: ThinkingLevel::Medium,
};
let selection = crate::jev::ReasoningSelection {
level: ThinkingLevel::High,
generations: 5,
};
run.lease.install(&selection);
assert_eq!(
run.checkpoint(&model, ThinkingLevel::Medium, ThinkingLevel::High, true),
(false, Some(ThinkingLevel::High))
);
let mut other_model = model.clone();
other_model.provider = "other".into();
assert_eq!(
run.checkpoint(
&other_model,
ThinkingLevel::Medium,
ThinkingLevel::High,
true
),
(true, None)
);
run.lease.install(&selection);
assert_eq!(
run.checkpoint(&other_model, ThinkingLevel::Low, ThinkingLevel::High, true),
(false, None)
);
run.lease.install(&selection);
assert_eq!(
run.checkpoint(&other_model, ThinkingLevel::Low, ThinkingLevel::Low, true),
(false, None)
);
run.lease.install(&selection);
assert_eq!(
run.checkpoint(&other_model, ThinkingLevel::Low, ThinkingLevel::High, false),
(false, None)
);
}
#[test]
fn late_cache_refreshes_are_skipped_before_the_cache_expires() {
let ttl = std::time::Duration::from_secs(300);
let delay = cache_warming_delay(ttl).unwrap();
assert!(!cache_refresh_deadline_missed(
std::time::Duration::from_secs(284),
ttl,
delay,
));
assert!(cache_refresh_deadline_missed(
std::time::Duration::from_secs(286),
ttl,
delay,
));
}
#[test]
fn transient_errors_require_a_status_or_specific_network_failure() {
assert!(is_transient("request failed with HTTP 503"));
assert!(is_transient("Cloudflare returned HTTP 520"));
assert!(is_transient("Azure is currently experiencing high demand"));
assert!(is_transient("connection reset by peer"));
assert!(is_transient("rate limit exceeded"));
assert!(!is_transient("model has a 500 token limit"));
assert!(!is_transient("connection settings are invalid"));
}
#[test]
fn session_tools_install_replace_and_remove_without_changing_base_tools() {
let registry = Registry::load(None);
let session = AgentSession::new(
SessionManager::in_memory(std::path::Path::new("/test")),
Vec::new(),
registry,
Settings::default(),
"test".into(),
openai_model(),
ThinkingLevel::Off,
None,
Arc::new(|_| {}),
);
let tool = || {
Arc::new(crate::tools::grep::GrepTool {
cwd: std::path::PathBuf::from("/test"),
}) as DynTool
};
assert!(!session.available_tool_names().contains(&"grep".into()));
session.install_session_tool(tool());
session.install_session_tool(tool());
assert_eq!(
session
.available_tool_names()
.iter()
.filter(|name| name.as_str() == "grep")
.count(),
1
);
assert!(session.remove_session_tool("grep"));
assert!(!session.remove_session_tool("grep"));
assert!(!session.available_tool_names().contains(&"grep".into()));
}
fn remote_result() -> kiss_ai::api::openai_compaction::RemoteCompactionResult {
kiss_ai::api::openai_compaction::RemoteCompactionResult {
replacement_history: vec![serde_json::json!({
"type": "compaction",
"encrypted_content": "opaque"
})],
usage: Some(Usage {
input: 10,
output: 2,
total_tokens: 12,
..Default::default()
}),
}
}
fn settings_test_session(settings: Settings, subagents_allowed: bool) -> Arc<AgentSession> {
let registry = Registry::from_builtin();
let model = registry.all().first().expect("built-in model").clone();
AgentSession::new_with_subagents_allowed(
SessionManager::in_memory(std::path::Path::new("/test")),
Vec::new(),
registry,
settings,
"root prompt".into(),
model,
ThinkingLevel::Off,
None,
Arc::new(|_| {}),
subagents_allowed,
)
}
fn benchmark_tools() -> Vec<DynTool> {
let cwd = std::path::PathBuf::from("/synthetic");
vec![
Arc::new(kiss_agent::tools::read::ReadTool { cwd: cwd.clone() }),
Arc::new(kiss_agent::tools::write::WriteTool { cwd: cwd.clone() }),
Arc::new(kiss_agent::tools::edit::EditTool { cwd: cwd.clone() }),
Arc::new(kiss_agent::tools::bash::BashTool::new(cwd)),
]
}
#[test]
fn subagent_tools_follow_settings_and_command_line_authority() {
let session = settings_test_session(Settings::default(), true);
assert!(session.available_tool_names().is_empty());
assert!(
!session
.build_context()
.system_prompt
.contains("Subagent coordination")
);
let mut enabled = session.settings();
enabled.subagents.enabled = true;
session.update_settings(enabled.clone());
assert_eq!(
session.available_tool_names(),
[
"spawn_agent",
"send_message",
"followup_task",
"wait_agent",
"list_agents",
"interrupt_agent"
]
);
assert!(
session
.build_context()
.system_prompt
.contains("Subagent coordination")
);
enabled.subagents.enabled = false;
session.update_settings(enabled);
assert!(session.available_tool_names().is_empty());
let mut blocked_settings = Settings::default();
blocked_settings.subagents.enabled = true;
let blocked = settings_test_session(blocked_settings, false);
assert!(blocked.available_tool_names().is_empty());
assert!(
!blocked
.build_context()
.system_prompt
.contains("Subagent coordination")
);
}
#[test]
fn the_workflow_tool_appears_only_in_workflow_prompt_mode() {
let mut settings = Settings::default();
settings.subagents.enabled = true;
let session = settings_test_session(settings, true);
assert!(
!session
.available_tool_names()
.contains(&"run_workflow".into())
);
assert_eq!(
session.prompt_mode_for("run a dynamic workflow for this task"),
PromptMode::Workflow
);
assert_eq!(
session.prompt_mode_for("fix this small function"),
PromptMode::Ordinary
);
assert!(
!session
.build_context()
.system_prompt
.contains("Writing a dynamic workflow")
);
assert!(
session
.available_tool_names_for(PromptMode::Workflow)
.contains(&"run_workflow".into())
);
assert!(
session
.build_context_for(PromptMode::Workflow)
.system_prompt
.contains("Writing a dynamic workflow")
);
assert!(
!session
.available_tool_names()
.contains(&"run_workflow".into())
);
}
#[test]
fn workflow_prompt_mode_does_nothing_while_subagents_are_off() {
let session = settings_test_session(Settings::default(), true);
assert!(!session.workflows_enabled());
assert!(
!session
.available_tool_names_for(PromptMode::Workflow)
.contains(&"run_workflow".into())
);
let mut settings = session.settings();
settings.subagents.enabled = true;
settings.workflows.enabled = false;
session.update_settings(settings.clone());
assert!(!session.workflows_enabled());
assert!(
!session
.available_tool_names_for(PromptMode::Workflow)
.contains(&"run_workflow".into())
);
settings.workflows.enabled = true;
session.update_settings(settings);
assert!(session.workflows_enabled());
assert!(
session
.available_tool_names_for(PromptMode::Workflow)
.contains(&"run_workflow".into())
);
}
#[test]
fn one_at_a_time_queues_keep_each_prompts_mode() {
let queue = Arc::new(Mutex::new(VecDeque::from([
QueuedPrompt {
message: AgentMessage::user("ordinary"),
mode: PromptMode::Ordinary,
},
QueuedPrompt {
message: AgentMessage::user("workflow"),
mode: PromptMode::Workflow,
},
])));
assert_eq!(
queued_mode(&queue, QueueMode::OneAtATime),
Some(PromptMode::Ordinary)
);
assert_eq!(drain_queue(&queue, QueueMode::OneAtATime).len(), 1);
assert_eq!(
queued_mode(&queue, QueueMode::OneAtATime),
Some(PromptMode::Workflow)
);
}
#[test]
fn a_session_without_subagent_authority_has_no_workflow_runtime() {
let mut settings = Settings::default();
settings.subagents.enabled = true;
let child = settings_test_session(settings, false);
assert!(child.workflows().is_none());
assert!(!child.workflows_enabled());
}
#[test]
fn child_session_has_safe_forked_context_without_control_tools() {
let mut settings = Settings::default();
settings.subagents.enabled = true;
let parent = settings_test_session(settings, true);
parent
.manager
.lock()
.unwrap()
.append_message(AgentMessage::user("parent context"))
.unwrap();
let child = parent
.create_subagent_session("inspect", "/root/inspect", ForkTurns::All, None, None)
.unwrap();
assert!(Arc::ptr_eq(&parent.registry, &child.registry));
assert!(child.available_tool_names().is_empty());
let context = child.manager.lock().unwrap().build_session_context();
assert!(matches!(
context.messages.as_slice(),
[AgentMessage::User(user)] if user.content.as_text() == "parent context"
));
}
#[test]
fn child_without_forked_turns_does_not_copy_parent_context() {
let parent = settings_test_session(Settings::default(), true);
parent
.manager
.lock()
.unwrap()
.append_message(AgentMessage::user("parent context"))
.unwrap();
let child = parent
.create_subagent_session("inspect", "/root/inspect", ForkTurns::None, None, None)
.unwrap();
assert!(
child
.manager
.lock()
.unwrap()
.build_session_context()
.messages
.is_empty()
);
}
#[test]
#[ignore = "release-mode performance benchmark"]
fn benchmark_performance_subagent_overhead() {
let registry = Registry::from_builtin();
let model = registry.all().first().expect("built-in model").clone();
let tools = benchmark_tools();
let make_session = |enabled: bool| {
let mut settings = Settings::default();
settings.subagents.enabled = enabled;
AgentSession::new_with_subagents_allowed(
SessionManager::in_memory(std::path::Path::new("/synthetic")),
tools.clone(),
registry.clone(),
settings,
"benchmark root prompt".into(),
model.clone(),
ThinkingLevel::Off,
None,
Arc::new(|_| {}),
true,
)
};
kiss_bench::measure_pair(
(
"agent_session_create_subagents_off",
"agent_session_create_subagents_on",
),
21,
500,
(
"new_root_session_4_base_tools_0_control_tools",
"new_root_session_4_base_tools_6_control_tools",
),
|| make_session(false),
|| make_session(true),
);
let off = make_session(false);
let on = make_session(true);
kiss_bench::measure_pair(
(
"agent_context_build_subagents_off",
"agent_context_build_subagents_on",
),
21,
10_000,
(
"empty_session_4_base_tools_0_control_tools",
"empty_session_4_base_tools_6_control_tools",
),
|| off.build_context(),
|| on.build_context(),
);
}
#[test]
#[ignore = "release-mode performance benchmark"]
fn benchmark_performance_workflow_tool_exposure() {
let registry = Registry::from_builtin();
let model = registry.all().first().expect("built-in model").clone();
let tools = benchmark_tools();
let make_session = || {
let mut settings = Settings::default();
settings.subagents.enabled = true;
AgentSession::new_with_subagents_allowed(
SessionManager::in_memory(std::path::Path::new("/synthetic")),
tools.clone(),
registry.clone(),
settings,
"benchmark root prompt".into(),
model.clone(),
ThinkingLevel::Off,
None,
Arc::new(|_| {}),
true,
)
};
let ordinary = make_session();
let workflow = make_session();
kiss_bench::measure_pair(
(
"agent_context_build_workflow_disarmed",
"agent_context_build_workflow_armed",
),
21,
10_000,
(
"empty_session_subagents_on_workflow_disarmed",
"empty_session_subagents_on_workflow_armed",
),
|| ordinary.build_context_for(PromptMode::Ordinary),
|| workflow.build_context_for(PromptMode::Workflow),
);
}
#[test]
fn transcript_excerpt_keeps_only_recent_user_and_assistant_text() {
let messages = vec![
AgentMessage::user("old"),
AgentMessage::BashExecution(kiss_agent::BashExecutionMessage {
command: "pwd".into(),
output: "ignored".into(),
exit_code: Some(0),
cancelled: false,
truncated: false,
full_output_path: None,
exclude_from_context: false,
timestamp: 1,
}),
AgentMessage::user("new"),
];
let excerpt = transcript_excerpt(&messages, 1, 100);
assert_eq!(excerpt, "User: new");
}
#[test]
fn transcript_excerpt_enforces_character_budget_from_the_tail() {
let excerpt = transcript_excerpt(&[AgentMessage::user("abcdefghij")], 4, 5);
assert!(excerpt.ends_with("fghij"));
assert!(excerpt.starts_with("[earlier text omitted]"));
}
#[test]
fn hybrid_compaction_keeps_local_summary_and_remote_details() {
let selected = select_compaction_outcome(
&openai_model(),
Ok(compaction::SummaryOutcome {
summary: "portable".into(),
usage: None,
}),
Some(Ok(remote_result())),
)
.unwrap();
assert_eq!(selected.summary, "portable");
assert_eq!(selected.remote_usage.unwrap().input, 10);
assert_eq!(
selected.remote_details.unwrap()["remoteCompaction"]["version"],
2
);
}
#[test]
fn remote_failure_falls_back_to_local_compaction() {
let selected = select_compaction_outcome(
&openai_model(),
Ok(compaction::SummaryOutcome {
summary: "portable".into(),
usage: None,
}),
Some(Err(anyhow::anyhow!("remote unavailable"))),
)
.unwrap();
assert_eq!(selected.summary, "portable");
assert!(selected.remote_details.is_none());
}
#[test]
fn remote_success_survives_local_summary_failure() {
let selected = select_compaction_outcome(
&openai_model(),
Err(anyhow::anyhow!("summary unavailable")),
Some(Ok(remote_result())),
)
.unwrap();
assert!(
selected
.summary
.contains("server-side compaction was applied")
);
assert!(selected.remote_details.is_some());
}
#[test]
fn details_merge_keeps_file_operations_and_remote_artifact() {
let merged = merge_compaction_details(
serde_json::json!({"readFiles": ["a.rs"], "modifiedFiles": []}),
Some(serde_json::json!({"remoteCompaction": {"version": 2}})),
);
assert_eq!(merged["readFiles"][0], "a.rs");
assert_eq!(merged["remoteCompaction"]["version"], 2);
}
#[test]
fn auto_compaction_guard_checks_settings_threshold_and_cancel() {
let mut settings = Settings::default();
settings.compaction.reserve_tokens = 20;
let messages = vec![AgentMessage::user("x".repeat(360))];
let mut model = openai_model();
model.context_window = 100;
assert!(auto_compaction_needed(&settings, &messages, &model, false));
assert!(!auto_compaction_needed(&settings, &messages, &model, true));
settings.compaction.enabled = false;
assert!(!auto_compaction_needed(&settings, &messages, &model, false));
}
#[test]
fn session_title_normalization_is_safe_and_bounded() {
assert_eq!(
normalize_session_title(" `Fix AUTH-123 login flow!` \nignored").as_deref(),
Some("Fix AUTH-123 login flow")
);
assert_eq!(normalize_session_title("\n\t"), None);
assert_eq!(
normalize_session_title("🚀".repeat(50).as_str())
.unwrap()
.chars()
.count(),
SESSION_TITLE_MAX_CHARS
);
}
#[test]
fn session_title_prompt_is_utf8_safe_and_bounded() {
let prompt = "🚀".repeat(SESSION_TITLE_PROMPT_MAX_BYTES);
let bounded = bounded_session_title_prompt(&prompt);
assert!(bounded.len() <= SESSION_TITLE_PROMPT_MAX_BYTES);
assert!(std::str::from_utf8(bounded.as_bytes()).is_ok());
}
#[test]
fn cache_warming_uses_ninety_percent_with_ten_second_margin() {
assert_eq!(
cache_warming_delay(std::time::Duration::from_secs(300)),
Some(std::time::Duration::from_secs(270))
);
assert_eq!(
cache_warming_delay(std::time::Duration::from_secs(60)),
Some(std::time::Duration::from_secs(50))
);
assert_eq!(
cache_warming_delay(std::time::Duration::from_secs(10)),
None
);
}
}