use std::fmt::Write;
use std::path::Path;
use std::sync::Arc;
use anyhow::Result;
use arc_swap::ArcSwap;
use futures_util::future::join_all;
use ignore::WalkBuilder;
use crate::pipeline::board::BOARD;
use crate::prompt::{build_workspace_context, format_ticket_block, load_prompt, substitute};
use crate::runtime_events::RuntimeEvent;
use crate::workspace::{WORKSPACES, truncate_workspace_notes};
use crate::agent::skills;
use crate::alarms::Alarm;
use crate::tools::active_models::{ModelKind, ModelSnapshot};
use crate::{ChatMessage, ChatRole, Role, Workspace};
use super::TranscriptSnapshot;
#[derive(Default)]
pub(crate) struct Session {
history: Vec<ChatMessage>,
persisted_len: usize,
token_length: Option<u64>,
transcript: Option<Arc<ArcSwap<TranscriptSnapshot>>>,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(crate) enum FinalizeOutcome {
Flushed,
NoUnpersistedTail,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(crate) enum RewriteOutcome {
Rewritten,
UnpersistedTailNoop,
}
pub(crate) struct PendingToolFrame {
pub calls: Vec<crate::ToolCall>,
pub has_text: bool,
pub call_count: usize,
}
impl Session {
pub(crate) async fn init(&mut self, agent_id: &str) -> Result<()> {
self.history = crate::session::store().load(agent_id).await;
self.token_length = crate::session::store().get_token_length(agent_id).await;
self.persisted_len = self.history.len();
self.publish_transcript();
Ok(())
}
#[expect(clippy::too_many_arguments)]
pub(crate) async fn append_turn_message(
&mut self,
agent_id: &str,
msg: &str,
ws: &Workspace,
role: &Role,
ticket: Option<&crate::pipeline::board::Ticket>,
channel: &str,
user_name: &str,
round_ts: Option<&str>,
is_admin: bool,
) -> Result<()> {
if !msg.is_empty() {
let is_new = self.history.is_empty();
if is_new {
let (msgs, snapshot) =
Self::build_turn_messages(msg, ws, role, ticket, round_ts, is_admin, user_name)
.await;
crate::session::store()
.batch_append_with_context(
agent_id,
&msgs,
channel,
user_name,
&ws.name,
role.as_str(),
)
.await?;
self.history.extend(msgs);
if matches!(role, Role::Assistant) {
Self::persist_active_models_snapshot(agent_id, &snapshot).await;
}
} else {
let (content, new_snapshot) = if matches!(role, Role::Assistant) {
Self::prepend_model_change(agent_id, msg, user_name).await
} else {
(msg.to_string(), None)
};
let user_msg = crate::session::user_msg_with_ts(&content, round_ts);
crate::session::store()
.append_with_context(
agent_id,
&user_msg,
channel,
user_name,
&ws.name,
role.as_str(),
)
.await?;
if let Some(snapshot) = new_snapshot {
Self::persist_active_models_snapshot(agent_id, &snapshot).await;
}
self.history.push(user_msg);
}
self.persisted_len = self.history.len();
self.publish_transcript();
}
Ok(())
}
pub(crate) fn attach_transcript(&mut self, holder: Arc<ArcSwap<TranscriptSnapshot>>) {
self.transcript = Some(holder);
}
fn publish_transcript(&self) {
if let Some(holder) = &self.transcript {
let snapshot = TranscriptSnapshot {
history: self.history.clone(),
token_count: self.token_length,
};
holder.store(Arc::new(snapshot));
crate::runtime_events::publish(RuntimeEvent::Registries);
}
}
#[must_use]
pub(crate) fn history(&self) -> &[ChatMessage] {
&self.history
}
#[must_use]
pub(crate) fn token_length(&self) -> Option<u64> {
self.token_length
}
pub(crate) fn set_token_length(&mut self, token_length: Option<u64>) {
self.token_length = token_length;
self.publish_transcript();
}
pub(crate) fn push_assistant(&mut self, content: String) {
self.history.push(ChatMessage::assistant(content));
self.publish_transcript();
}
pub(crate) async fn persist_messages(
&mut self,
agent_id: &str,
messages: &[ChatMessage],
) -> Result<()> {
debug_assert!(self.persisted_len <= self.history.len());
let mut batch =
Vec::with_capacity(self.history.len() - self.persisted_len + messages.len());
batch.extend_from_slice(&self.history[self.persisted_len..]);
batch.extend_from_slice(messages);
crate::session::store()
.batch_append(agent_id, &batch)
.await?;
self.persisted_len += batch.len();
self.history.extend_from_slice(messages);
self.publish_transcript();
Ok(())
}
pub(crate) fn pending_tool_frame(&self) -> Option<PendingToolFrame> {
let frame_idx = self.history.iter().rposition(super::is_tool_call_frame)?;
let Some(super::DecodedNativeHistoryMessage::Assistant {
content,
tool_calls: Some(calls),
..
}) = super::decode_native_history_message(&self.history[frame_idx])
else {
return None;
};
let has_text = content.as_deref().is_some_and(|t| !t.trim().is_empty());
let call_count = calls.len();
let mut completed: std::collections::HashSet<String> = std::collections::HashSet::default();
for msg in &self.history[frame_idx + 1..] {
if msg.role == ChatRole::Tool {
if let Some(super::DecodedNativeHistoryMessage::ToolResult {
tool_call_id, ..
}) = super::decode_native_history_message(msg)
{
completed.insert(tool_call_id);
}
} else {
break;
}
}
let pending: Vec<crate::ToolCall> = calls
.into_iter()
.filter(|c| !completed.contains(&c.id))
.collect();
if pending.is_empty() {
None
} else {
Some(PendingToolFrame {
calls: pending,
has_text,
call_count,
})
}
}
#[cfg(test)]
pub(crate) async fn settle_tool_results(
&mut self,
agent_id: &str,
results: &[(String, String)],
follow_up: &[ChatMessage],
) -> Result<()> {
let settled = crate::session::store()
.settle_tool_results(agent_id, results, follow_up)
.await?;
self.mirror_settled(settled)?;
Ok(())
}
pub(crate) fn mirror_settled(&mut self, settled: Vec<super::SettledRow>) -> Result<()> {
let Some(frame_idx) = self.history.iter().rposition(super::is_tool_call_frame) else {
debug_assert!(
false,
"settle_tool_results called without a tool-call frame in history"
);
return Ok(());
};
self.history.truncate(frame_idx + 1);
for row in settled {
let role = row.role.parse::<ChatRole>().map_err(|e| {
anyhow::anyhow!("settled session role {} is invalid: {e}", row.role)
})?;
self.history.push(ChatMessage {
role,
content: row.content,
});
}
self.persisted_len = self.history.len();
self.publish_transcript();
Ok(())
}
pub(crate) fn push_messages_unpersisted(&mut self, messages: &[ChatMessage]) {
self.history.extend_from_slice(messages);
self.publish_transcript();
}
pub(crate) async fn rewrite_last_user_message(
&mut self,
agent_id: &str,
content: String,
) -> Result<RewriteOutcome> {
debug_assert!(self.persisted_len <= self.history.len());
let Some(idx) = self.history.iter().rposition(|m| m.role == ChatRole::User) else {
anyhow::bail!("no user message in history to rewrite");
};
if idx >= self.persisted_len {
return Ok(RewriteOutcome::UnpersistedTailNoop);
}
crate::session::store()
.rewrite_last_user_message(agent_id, &content)
.await?;
self.history[idx].content = content;
self.publish_transcript();
Ok(RewriteOutcome::Rewritten)
}
#[expect(clippy::too_many_arguments)]
pub(crate) async fn apply_summary(
&mut self,
agent_id: &str,
summary_text: &str,
ws: &Workspace,
role: &Role,
ticket: Option<&crate::pipeline::board::Ticket>,
is_admin: bool,
user_name: &str,
) {
let (mut compacted, snapshot) =
Self::build_context_messages(ws, role, ticket, is_admin, user_name).await;
let prefix = load_prompt("context/summary_prefix.md");
let retained = crate::session::select_retention_window(&self.history);
let dump_path = super::compaction_dump::dump_deleted_messages(
agent_id,
&self.history,
&retained,
role,
user_name,
)
.await;
let mut summary = format!("{prefix}{summary_text}");
if let Some(path) = &dump_path {
summary.push_str("\n\n");
let rendered_path = path.display().to_string();
summary.push_str(&substitute(
&load_prompt("context/summary_spill.md"),
&[("{{path}}", rendered_path.as_str())],
));
}
compacted.push(ChatMessage::system(summary));
compacted.extend(retained);
if let Err(e) = crate::session::store()
.replace_messages(agent_id, &compacted)
.await
{
tracing::error!(
agent_id = %agent_id,
error = %e,
"Failed to persist compacted session after summarization"
);
} else {
self.history = compacted;
self.persisted_len = self.history.len();
self.publish_transcript();
let block_rendered = snapshot.image.is_some() || snapshot.video.is_some();
if matches!(role, Role::Assistant) && block_rendered {
let merged = Self::merge_active_models_snapshot(agent_id, &snapshot).await;
Self::persist_active_models_snapshot(agent_id, &merged).await;
}
}
}
pub(crate) async fn finalize(&mut self, agent_id: &str) -> Result<FinalizeOutcome> {
debug_assert!(self.persisted_len <= self.history.len());
if self.persisted_len == self.history.len() {
return Ok(FinalizeOutcome::NoUnpersistedTail);
}
self.persist_messages(agent_id, &[]).await?;
Ok(FinalizeOutcome::Flushed)
}
async fn build_context_messages(
ws: &Workspace,
role: &Role,
ticket: Option<&crate::pipeline::board::Ticket>,
is_admin: bool,
user_name: &str,
) -> (Vec<ChatMessage>, ModelSnapshot) {
let (stored_context, board_context) = tokio::join!(
lookup_workspace_context(ws, role),
build_board_context(ws, role),
);
let workspace_context = match stored_context.as_deref() {
Some(ctx) => ctx.to_owned(),
None => build_workspace_context(ws.as_path()).await,
};
let workspace_context = {
let mut ctx = String::new();
ctx.push_str(&crate::prompt::wrap_workspace_context(&workspace_context));
if !ws.notes.trim().is_empty() {
let notes = truncate_workspace_notes(&ws.notes);
let _ = write!(ctx, "\n<user-notes>\n{notes}\n</user-notes>\n");
}
ctx
};
let workspace_boilerplate = substitute(
&load_prompt("context/workspace.md"),
&[
(
"{{path_frame}}",
&load_prompt(workspace_frame_key(*role, is_admin)),
),
("{{operating_system}}", std::env::consts::OS),
("{{workspace}}", &ws.as_path().display().to_string()),
("{{workspace_context}}", &workspace_context),
(
"{{system_locale}}",
&sys_locale::get_locale().unwrap_or_else(|| "unknown".to_string()),
),
],
);
let role_description = role.role_description_for(is_admin);
let skills = skills::load_skills(ws).await;
let mut msgs = Vec::with_capacity(7);
msgs.push(ChatMessage::system(&role_description));
if matches!(role, Role::Assistant) && is_admin {
let live = crate::config::CONFIG.onboarding_stage()
!= crate::config::OnboardingState::Finished;
if live || crate::onboarding::take_greeting_guide_pending() {
msgs.push(ChatMessage::system(crate::prompt::load_prompt(
"role/onboarding.md",
)));
}
}
let mut snapshot = ModelSnapshot::default();
if matches!(role, Role::Assistant)
&& let Some((block, rendered)) =
crate::tools::active_models::render_block(user_name).await
{
msgs.push(ChatMessage::system(block));
snapshot = rendered;
}
msgs.push(ChatMessage::system(&workspace_boilerplate));
if !skills.is_empty() {
msgs.push(ChatMessage::system(skills::skills_to_prompt(&skills, ws)));
}
if matches!(role, Role::Assistant) {
let ((alarms, workspaces, personal_files), custom_tools) = tokio::join!(
fetch_assistant_context(user_name, is_admin),
crate::tools::custom::context_block(user_name, is_admin),
);
for block in assistant_context_blocks(
is_admin,
&alarms,
&workspaces,
personal_files.as_deref(),
&custom_tools,
) {
msgs.push(ChatMessage::system(&block));
}
}
if let Some(board_context) = board_context {
msgs.push(ChatMessage::system(&board_context));
}
if let Some(t) = ticket {
msgs.push(ChatMessage::system(format_ticket_block(t)));
}
(msgs, snapshot)
}
async fn build_turn_messages(
msg: &str,
ws: &Workspace,
role: &Role,
ticket: Option<&crate::pipeline::board::Ticket>,
round_ts: Option<&str>,
is_admin: bool,
user_name: &str,
) -> (Vec<ChatMessage>, ModelSnapshot) {
let (mut msgs, snapshot) =
Self::build_context_messages(ws, role, ticket, is_admin, user_name).await;
msgs.push(crate::session::user_msg_with_ts(msg, round_ts));
(msgs, snapshot)
}
async fn persist_active_models_snapshot(agent_id: &str, snapshot: &ModelSnapshot) {
let json =
(snapshot.image.is_some() || snapshot.video.is_some()).then(|| snapshot.to_json());
if let Err(e) = crate::session::store()
.set_active_models(agent_id, json.as_deref())
.await
{
tracing::warn!(agent_id = %agent_id, error = %e, "Failed to persist active-models snapshot");
}
}
async fn merge_active_models_snapshot(
agent_id: &str,
rendered: &ModelSnapshot,
) -> ModelSnapshot {
let previous = crate::session::store()
.get_active_models(agent_id)
.await
.and_then(|json| ModelSnapshot::from_json(&json));
ModelSnapshot {
image: rendered
.image
.clone()
.or_else(|| previous.as_ref().and_then(|p| p.image.clone())),
video: rendered
.video
.clone()
.or_else(|| previous.as_ref().and_then(|p| p.video.clone())),
}
}
async fn prepend_model_change(
agent_id: &str,
msg: &str,
user_name: &str,
) -> (String, Option<ModelSnapshot>) {
let Some(previous) = crate::session::store()
.get_active_models(agent_id)
.await
.and_then(|json| ModelSnapshot::from_json(&json))
else {
return (msg.to_string(), None);
};
let current = ModelSnapshot::from_user(user_name).await;
let mut advanced = previous.clone();
let mut changed = Vec::new();
if let (Some(old), Some(new)) = (previous.image.as_deref(), current.image.as_deref())
&& old != new
{
advanced.image = current.image.clone();
changed.push((ModelKind::Image, old, new));
}
if let (Some(old), Some(new)) = (previous.video.as_deref(), current.video.as_deref())
&& old != new
{
advanced.video = current.video.clone();
changed.push((ModelKind::Video, old, new));
}
if changed.is_empty() {
return (msg.to_string(), None);
}
let blocks = join_all(
changed
.into_iter()
.map(|(kind, old, new)| Self::model_change_block(kind, old, new)),
)
.await;
(format!("{}\n\n{msg}", blocks.join("\n")), Some(advanced))
}
async fn model_change_block(kind: ModelKind, old: &str, new: &str) -> String {
match crate::tools::active_models::render_section(kind, new).await {
Some(section) => format!(
"<model-change>Active {} model changed from {} to {}.\nNew model capabilities:\n{section}</model-change>",
kind.label(),
old,
new
),
None => format!(
"<model-change>Active {} model changed from {} to {}.</model-change>",
kind.label(),
old,
new
),
}
}
}
async fn build_board_context(ws: &Workspace, role: &Role) -> Option<String> {
if !matches!(role, Role::Manager) {
return None;
}
let board = BOARD.get()?;
let tickets = board.list_all_tickets(Some(&ws.name), None).await.ok()?;
let active: Vec<_> = tickets
.into_iter()
.filter(|t| !t.phase.is_unblocking())
.collect();
if active.is_empty() {
return None;
}
let count = active.len();
let mut output = format!(
"<workspace-board>\nTickets in {} ({count} active):\n",
ws.name
);
for t in &active {
let _ = writeln!(output, "{}", t.short_display());
}
output.push_str("</workspace-board>");
Some(output)
}
const MAX_WORKSPACE_SUMMARY_CHARS: usize = 1000;
async fn fetch_assistant_context(
user_name: &str,
is_admin: bool,
) -> (Vec<Alarm>, Vec<(Workspace, Option<String>)>, Option<String>) {
let alarms = async {
if crate::alarms::ALARMS.get().is_none() {
return Vec::new();
}
crate::alarms::list_user_alarms(user_name)
.await
.unwrap_or_default()
};
let workspaces = async {
if !is_admin {
return Vec::new();
}
let Some(workspaces) = WORKSPACES.get() else {
return Vec::new();
};
let Ok(all) = workspaces.list().await else {
return Vec::new();
};
let registered: Vec<Workspace> = all
.into_iter()
.filter(|w| !w.name.starts_with("personal:"))
.collect();
let summaries = join_all(
registered
.iter()
.map(|w| workspaces.get_general_context(&w.name)),
)
.await
.into_iter()
.map(|res| res.ok().flatten())
.collect::<Vec<_>>();
registered.into_iter().zip(summaries).collect()
};
tokio::join!(alarms, workspaces, fetch_personal_files(user_name))
}
fn workspace_frame_key(role: Role, is_admin: bool) -> &'static str {
if matches!(role, Role::Assistant) && is_admin {
"context/workspace_frame_admin.md"
} else {
"context/workspace_frame.md"
}
}
fn assistant_context_blocks(
is_admin: bool,
alarms: &[Alarm],
workspaces: &[(Workspace, Option<String>)],
personal_files: Option<&str>,
custom_tools: &str,
) -> Vec<String> {
let mut blocks = Vec::new();
if let Some(lines) = render_alarm_lines(alarms) {
blocks.push(substitute(
&load_prompt("context/alarms.md"),
&[("{{alarms}}", &lines)],
));
}
if let Some(lines) = personal_files {
blocks.push(substitute(
&load_prompt("context/personal_files.md"),
&[("{{files}}", lines)],
));
}
if is_admin && let Some(lines) = render_workspace_lines(workspaces) {
blocks.push(substitute(
&load_prompt("context/workspaces.md"),
&[("{{workspaces}}", &lines)],
));
}
if !custom_tools.is_empty() {
blocks.push(custom_tools.to_string());
}
blocks
}
fn render_alarm_lines(alarms: &[Alarm]) -> Option<String> {
if alarms.is_empty() {
return None;
}
let mut out = String::new();
for alarm in alarms {
let fire = crate::alarms::format_fire_time(&alarm.next_fire_at)
.unwrap_or_else(|_| alarm.next_fire_at.clone());
let _ = write!(out, "- {}, {}", alarm.id, alarm.kind);
if let Some(interval) = alarm.interval_seconds {
let _ = write!(out, " (every {interval} seconds)");
}
let _ = write!(out, ", {}, next fire: {}", alarm.text, fire);
if let Some(trigger) = &alarm.trigger {
let _ = write!(out, ", trigger: {}", trigger.render());
}
out.push('\n');
}
Some(out.trim_end().to_string())
}
fn render_workspace_lines(workspaces: &[(Workspace, Option<String>)]) -> Option<String> {
let mut out = String::new();
for (ws, summary) in workspaces {
let _ = writeln!(out, "- {} ({}): {}", ws.name, ws.status, ws.path);
if let Some(summary) = summary.as_deref().map(str::trim).filter(|s| !s.is_empty()) {
let _ = writeln!(
out,
" {}",
crate::util::truncate(
&summary.split_whitespace().collect::<Vec<_>>().join(" "),
MAX_WORKSPACE_SUMMARY_CHARS,
)
);
}
}
(!out.is_empty()).then(|| out.trim_end().to_string())
}
const MAX_PERSONAL_FILE_ENTRIES: usize = 200;
const MAX_PERSONAL_FILES_BYTES: usize = 6 * 1024;
async fn fetch_personal_files(user_name: &str) -> Option<String> {
if !crate::users::is_valid_personal_user_name(user_name) {
return None;
}
let root = crate::users::personal_workspace_path(user_name);
let listed = tokio::task::spawn_blocking(move || walk_personal_files(&root))
.await
.ok()?;
render_personal_file_lines(listed)
}
fn walk_personal_files(root: &Path) -> Vec<String> {
let cap = MAX_PERSONAL_FILE_ENTRIES + 1;
let mut out = Vec::new();
for entry in WalkBuilder::new(root)
.filter_entry(|e| {
e.depth() == 0
|| !e.file_type().is_some_and(|t| t.is_dir())
|| !matches!(e.file_name().to_str(), Some("generated" | "uploads"))
})
.build()
.flatten()
{
if entry.depth() == 0 || entry.file_type().is_some_and(|t| t.is_dir()) {
continue; }
let Ok(rel) = entry.path().strip_prefix(root) else {
continue;
};
out.push(rel.display().to_string());
if out.len() >= cap {
break;
}
}
out
}
fn render_personal_file_lines(mut paths: Vec<String>) -> Option<String> {
if paths.is_empty() {
return None;
}
paths.sort();
let mut out = String::new();
let mut shown = 0;
for path in &paths {
if shown == MAX_PERSONAL_FILE_ENTRIES
|| out.len() + path.len() + 1 > MAX_PERSONAL_FILES_BYTES
{
break;
}
out.push_str(path);
out.push('\n');
shown += 1;
}
if shown < paths.len() {
out.push_str("…listing truncated; use the read tool for the full picture");
}
Some(out.trim_end().to_string())
}
async fn lookup_workspace_context(ws: &Workspace, role: &Role) -> Option<String> {
let workspaces = crate::workspace::store();
workspaces.get_context(&ws.name, role.as_str()).await.ok()?
}
#[cfg(test)]
mod tests {
use super::*;
#[tokio::test]
#[serial_test::serial(active_models)]
async fn prepend_model_change_detects_switch_and_refreshes_baseline() {
let user = "test_assistant_change_user";
crate::util::test::init_test_stores().await;
crate::tools::media_catalog::image::seed_cache(Some(std::sync::Arc::new(
crate::tools::media_catalog::image::ImageCatalog::default(),
)));
crate::tools::media_catalog::video::seed_cache(Some(std::sync::Arc::new(
crate::tools::media_catalog::video::VideoCatalog::default(),
)));
let agent_id = "test_assistant_change";
crate::users::store()
.set_image_gen_model(user, "model-a")
.await
.expect("set image model");
crate::users::store()
.set_video_model(user, "video-a")
.await
.expect("set video model");
let baseline = ModelSnapshot::from_user(user).await;
crate::session::store()
.set_active_models(agent_id, Some(&baseline.to_json()))
.await
.expect("baseline persisted");
crate::users::store()
.set_image_gen_model(user, "model-b")
.await
.expect("set image model");
crate::users::store()
.set_video_model(user, "video-b")
.await
.expect("set video model");
let (out, new_snapshot) = Session::prepend_model_change(agent_id, "hello", user).await;
assert!(
out.starts_with(
"<model-change>Active image model changed from model-a to model-b.</model-change>\n\
<model-change>Active video model changed from video-a to video-b.</model-change>"
),
"unexpected change-info: {out:?}"
);
assert!(out.ends_with("\n\nhello"));
let new_snapshot = new_snapshot.expect("snapshot returned for persistence");
crate::session::store()
.set_active_models(agent_id, Some(&new_snapshot.to_json()))
.await
.expect("baseline advanced");
let (out, snapshot) = Session::prepend_model_change(agent_id, "again", user).await;
assert_eq!(out, "again");
assert!(snapshot.is_none());
}
#[tokio::test]
async fn prepend_model_change_without_baseline_is_noop() {
crate::util::test::init_test_stores().await;
let (out, snapshot) =
Session::prepend_model_change("test_assistant_no_baseline", "hello", "no_such_user")
.await;
assert_eq!(out, "hello");
assert!(snapshot.is_none());
}
#[tokio::test]
#[serial_test::serial(active_models)]
async fn prepend_model_change_preserves_absent_sections() {
let user = "test_assistant_partial_user";
crate::util::test::init_test_stores().await;
crate::tools::media_catalog::image::seed_cache(Some(std::sync::Arc::new(
crate::tools::media_catalog::image::ImageCatalog::default(),
)));
crate::tools::media_catalog::video::seed_cache(Some(std::sync::Arc::new(
crate::tools::media_catalog::video::VideoCatalog::default(),
)));
let agent_id = "test_assistant_partial";
crate::users::store()
.set_video_model(user, "video-a")
.await
.expect("set video model");
let baseline = ModelSnapshot {
image: None,
video: Some("video-a".into()),
};
crate::session::store()
.set_active_models(agent_id, Some(&baseline.to_json()))
.await
.expect("baseline persisted");
crate::users::store()
.set_video_model(user, "video-b")
.await
.expect("set video model");
let (out, new_snapshot) = Session::prepend_model_change(agent_id, "hello", user).await;
assert!(
out.starts_with("<model-change>Active video model changed from video-a to video-b.")
);
let new_snapshot = new_snapshot.expect("snapshot returned for persistence");
assert_eq!(new_snapshot.image, None);
assert_eq!(new_snapshot.video.as_deref(), Some("video-b"));
crate::session::store()
.set_active_models(agent_id, Some(&new_snapshot.to_json()))
.await
.expect("baseline advanced");
crate::users::store()
.set_image_gen_model(user, "image-b")
.await
.expect("set image model");
let (out, snapshot) = Session::prepend_model_change(agent_id, "again", user).await;
assert_eq!(out, "again");
assert!(snapshot.is_none());
}
#[tokio::test]
async fn merge_active_models_snapshot_preserves_unrendered_sections() {
crate::util::test::init_test_stores().await;
let previous = ModelSnapshot {
image: Some("model-a".into()),
video: Some("video-a".into()),
};
crate::session::store()
.set_active_models("test_assistant_merge", Some(&previous.to_json()))
.await
.expect("baseline persisted");
let rendered = ModelSnapshot {
image: Some("model-b".into()),
video: None,
};
let merged = Session::merge_active_models_snapshot("test_assistant_merge", &rendered).await;
assert_eq!(
merged,
ModelSnapshot {
image: Some("model-b".into()),
video: Some("video-a".into()),
}
);
let merged =
Session::merge_active_models_snapshot("test_assistant_no_merge_baseline", &rendered)
.await;
assert_eq!(merged, rendered);
}
#[tokio::test]
async fn rewrite_last_user_message_persists_then_swaps() {
crate::util::test::init_test_stores().await;
let agent_id = "test_rewrite_ok";
let seed = vec![
ChatMessage::system("role"),
ChatMessage::user("[IMAGE:/tmp/a.png] first"),
ChatMessage::user("[IMAGE:/tmp/b.png] second"),
];
crate::session::store()
.batch_append(agent_id, &seed)
.await
.unwrap();
let mut session = Session::default();
session.history = seed;
session.persisted_len = session.history.len();
let outcome = session
.rewrite_last_user_message(agent_id, "rewritten".to_string())
.await
.unwrap();
assert_eq!(outcome, RewriteOutcome::Rewritten);
assert_eq!(session.history[2].content, "rewritten");
assert_eq!(
session.history[1].content, "[IMAGE:/tmp/a.png] first",
"earlier user row untouched"
);
let msgs = crate::session::store().load(agent_id).await;
assert_eq!(msgs.len(), 3);
assert_eq!(msgs[2].content, "rewritten");
}
#[tokio::test]
async fn rewrite_last_user_message_unpersisted_tail_is_conservative_noop() {
crate::util::test::init_test_stores().await;
let agent_id = "test_rewrite_noop";
crate::session::store()
.batch_append(
agent_id,
&[ChatMessage::user("[IMAGE:/tmp/a.png] persisted")],
)
.await
.unwrap();
let mut session = Session {
history: vec![
ChatMessage::user("[IMAGE:/tmp/a.png] persisted"),
ChatMessage::user("[IMAGE:/tmp/b.png] unpersisted"),
],
persisted_len: 1,
..Default::default()
};
let outcome = session
.rewrite_last_user_message(agent_id, "rewritten".to_string())
.await
.unwrap();
assert_eq!(outcome, RewriteOutcome::UnpersistedTailNoop);
assert_eq!(
session.history[1].content, "[IMAGE:/tmp/b.png] unpersisted",
"in-memory history untouched"
);
let msgs = crate::session::store().load(agent_id).await;
assert_eq!(msgs[0].content, "[IMAGE:/tmp/a.png] persisted");
}
#[tokio::test]
async fn rewrite_last_user_message_no_user_message_errors() {
crate::util::test::init_test_stores().await;
let mut session = Session {
history: vec![ChatMessage::system("role")],
persisted_len: 1,
..Default::default()
};
let err = session
.rewrite_last_user_message("test_rewrite_err", "x".to_string())
.await
.unwrap_err();
assert!(err.to_string().contains("no user message"), "{err}");
}
#[tokio::test]
async fn rewrite_last_user_message_persist_failure_leaves_history_untouched() {
crate::util::test::init_test_stores().await;
let agent_id = "test_rewrite_persist_fail";
crate::session::store()
.batch_append(
"other_agent",
&[ChatMessage::user("[IMAGE:/tmp/a.png] other")],
)
.await
.unwrap();
let mut session = Session {
history: vec![
ChatMessage::system("role"),
ChatMessage::user("[IMAGE:/tmp/a.png] mine"),
],
persisted_len: 2, ..Default::default()
};
let err = session
.rewrite_last_user_message(agent_id, "rewritten".to_string())
.await
.unwrap_err();
assert!(
err.to_string().contains("no user message row"),
"store error propagates: {err}"
);
assert_eq!(
session.history[1].content, "[IMAGE:/tmp/a.png] mine",
"in-memory history untouched on persist failure (ordering invariant)"
);
}
fn analyze_frame(call_id: &str) -> String {
crate::providers::reasoning::assistant_replay_payload(
Some(""),
&[crate::ToolCall {
id: call_id.to_string(),
name: "analyze".to_string(),
arguments: serde_json::json!({"analyze": "analyze task"}),
}],
None,
)
.to_string()
}
#[test]
fn pending_tool_frame_none_and_dangling_and_partial() {
assert!(Session::default().pending_tool_frame().is_none());
let frame = crate::providers::reasoning::assistant_replay_payload(
Some(""),
&[
crate::ToolCall {
id: "call_a".to_string(),
name: "read".to_string(),
arguments: serde_json::json!({"path": "a"}),
},
crate::ToolCall {
id: "call_b".to_string(),
name: "read".to_string(),
arguments: serde_json::json!({"path": "b"}),
},
],
None,
)
.to_string();
let mut session = Session {
history: vec![ChatMessage::assistant(frame)],
persisted_len: 1,
..Default::default()
};
let pending = session
.pending_tool_frame()
.expect("dangling calls present")
.calls;
assert_eq!(pending.len(), 2);
assert_eq!(pending[0].id, "call_a");
assert_eq!(pending[1].id, "call_b");
session
.history
.push(ChatMessage::tool_result("call_a", "result a"));
let pending = session
.pending_tool_frame()
.expect("still one dangling")
.calls;
assert_eq!(pending.len(), 1);
assert_eq!(pending[0].id, "call_b");
session
.history
.push(ChatMessage::tool_result("call_b", "result b"));
assert!(session.pending_tool_frame().is_none());
}
#[tokio::test]
async fn settle_tool_results_inserts_result_after_frame() {
crate::util::test::init_test_stores().await;
let agent_id = "settle_tool_results_test";
let frame = analyze_frame("call_analyze_1");
let mut session = Session::default();
session
.persist_messages(
agent_id,
&[ChatMessage::user("first"), ChatMessage::assistant(frame)],
)
.await
.unwrap();
assert_eq!(session.persisted_len, 2, "both seeded rows persisted");
let conn = &crate::session::store().conn;
let count_before: i64 = conn
.query_optional(
"SELECT message_count FROM session_metadata WHERE agent_id = ?1",
crate::db::params![agent_id],
|r| r.get::<i64>(0),
)
.await
.unwrap()
.expect("metadata row exists");
session
.settle_tool_results(
agent_id,
&[("call_analyze_1".to_string(), "the result".to_string())],
&[],
)
.await
.unwrap();
let rows = conn
.query(
"SELECT id, role, content FROM sessions WHERE agent_id = ?1 ORDER BY id",
crate::db::params![agent_id],
)
.await
.unwrap();
assert_eq!(rows.len(), 3, "one result row inserted after the frame");
let roles: Vec<String> = rows.iter().map(|r| r.get::<String>(1).unwrap()).collect();
assert_eq!(roles, vec!["user", "assistant", "tool"]);
assert_eq!(
rows[2].get::<String>(2).unwrap(),
serde_json::to_string(&crate::ToolResultPayload {
tool_call_id: "call_analyze_1".to_string(),
content: "the result".to_string(),
})
.unwrap(),
"the tool row carries the settled result with the original call id"
);
let count_after: i64 = conn
.query_optional(
"SELECT message_count FROM session_metadata WHERE agent_id = ?1",
crate::db::params![agent_id],
|r| r.get::<i64>(0),
)
.await
.unwrap()
.unwrap();
assert_eq!(count_after - count_before, 1);
let roles: Vec<crate::ChatRole> = session.history().iter().map(|m| m.role).collect();
assert_eq!(
roles,
vec![ChatRole::User, ChatRole::Assistant, ChatRole::Tool]
);
assert_eq!(
session.history()[2].content,
serde_json::to_string(&crate::ToolResultPayload {
tool_call_id: "call_analyze_1".to_string(),
content: "the result".to_string(),
})
.unwrap()
);
assert_eq!(session.persisted_len, session.history().len());
}
#[tokio::test]
async fn settle_tool_results_preserves_frame_call_order() {
crate::util::test::init_test_stores().await;
let agent_id = "settle_tool_results_order_test";
let frame = crate::providers::reasoning::assistant_replay_payload(
Some(""),
&[
crate::ToolCall {
id: "call_a".to_string(),
name: "read".to_string(),
arguments: serde_json::json!({"path": "a"}),
},
crate::ToolCall {
id: "call_b".to_string(),
name: "read".to_string(),
arguments: serde_json::json!({"path": "b"}),
},
],
None,
)
.to_string();
let mut session = Session::default();
session
.persist_messages(
agent_id,
&[
ChatMessage::user("first"),
ChatMessage::assistant(frame),
ChatMessage::tool_result("call_b", "result b"),
],
)
.await
.unwrap();
session
.settle_tool_results(
agent_id,
&[("call_a".to_string(), "result a".to_string())],
&[],
)
.await
.unwrap();
let conn = &crate::session::store().conn;
let rows = conn
.query(
"SELECT id, role, content FROM sessions WHERE agent_id = ?1 ORDER BY id",
crate::db::params![agent_id],
)
.await
.unwrap();
assert_eq!(
rows.len(),
4,
"frame + captured sibling + one recovered result"
);
let roles: Vec<String> = rows.iter().map(|r| r.get::<String>(1).unwrap()).collect();
assert_eq!(roles, vec!["user", "assistant", "tool", "tool"]);
let mut ids: Vec<String> = Vec::new();
for row in &rows[2..] {
let payload: crate::ToolResultPayload =
serde_json::from_str(&row.get::<String>(2).unwrap()).unwrap();
ids.push(payload.tool_call_id);
}
assert_eq!(ids, vec!["call_a", "call_b"], "frame call order preserved");
let b_payload: crate::ToolResultPayload =
serde_json::from_str(&rows[3].get::<String>(2).unwrap()).unwrap();
assert_eq!(b_payload.content, "result b", "captured sibling verbatim");
}
#[test]
fn the_admin_workspace_frame_states_what_is_closed() {
let shared = crate::prompt::load_prompt("context/workspace.md");
assert!(
shared.contains("{{path_frame}}"),
"the shared boilerplate takes the role's frame: {shared}"
);
let frame = crate::prompt::load_prompt("context/workspace_frame.md");
assert!(
frame.contains("outside the workspace"),
"the shared frame keeps the workspace rule: {frame}"
);
let admin = crate::prompt::load_prompt("context/workspace_frame_admin.md");
assert!(
!admin.contains("outside the workspace"),
"the admin's frame must not read as confinement: {admin}"
);
for named in ["registered project workspaces", "live databases"] {
assert!(
admin.contains(named),
"the admin's frame must name what stays closed ({named}): {admin}"
);
}
for asset in [&frame, &admin] {
assert!(
!asset.ends_with('\n'),
"the frame assets are fragments and must not end with a newline"
);
}
assert_eq!(
workspace_frame_key(Role::Assistant, true),
"context/workspace_frame_admin.md"
);
for pick in [
workspace_frame_key(Role::Assistant, false),
workspace_frame_key(Role::Manager, true),
workspace_frame_key(Role::Engineer, false),
] {
assert_eq!(pick, "context/workspace_frame.md");
}
}
#[test]
fn assistant_blocks_gating() {
let alarm = Alarm {
id: "a1".to_string(),
session_id: "s1".to_string(),
user_name: "u".to_string(),
kind: "one-shot".to_string(),
text: "hello".to_string(),
trigger: None,
interval_seconds: None,
next_fire_at: "2026-01-01T00:00:00+00:00".to_string(),
};
let workspaces = vec![(
crate::workspace::test_ws_named("/tmp/proj", "proj"),
Some("summary".to_string()),
)];
let files = Some("MEMORY.md\nnotes/projects.md");
let custom = "<custom-tools>\n- weather\n</custom-tools>";
let blocks = assistant_context_blocks(
false,
std::slice::from_ref(&alarm),
&workspaces,
files,
custom,
);
assert_eq!(blocks.len(), 3);
assert!(blocks[0].contains("<user-alarms>"));
assert!(blocks[1].contains("<personal-files>"));
assert!(!blocks[1].contains("<registered-workspaces>"));
assert_eq!(blocks[2], custom);
let blocks = assistant_context_blocks(
true,
std::slice::from_ref(&alarm),
&workspaces,
files,
custom,
);
assert_eq!(blocks.len(), 4);
assert!(blocks[0].contains("<user-alarms>"));
assert!(blocks[1].contains("<personal-files>"));
assert!(blocks[2].contains("<registered-workspaces>"));
assert_eq!(blocks[3], custom);
assert_eq!(
assistant_context_blocks(true, &[], &[], None, custom),
vec![custom.to_string()]
);
}
#[test]
fn personal_files_walk_prunes_excluded_dirs() {
let root = tempfile::tempdir().expect("tempdir");
std::fs::write(root.path().join("MEMORY.md"), "m").expect("write");
std::fs::create_dir_all(root.path().join("notes")).expect("mkdir");
std::fs::write(root.path().join("notes/topic.md"), "t").expect("write");
std::fs::create_dir_all(root.path().join("generated")).expect("mkdir");
std::fs::write(root.path().join("generated/big.bin"), "b").expect("write");
std::fs::create_dir_all(root.path().join(".git")).expect("mkdir");
std::fs::write(root.path().join(".git/config"), "c").expect("write");
let listed = walk_personal_files(root.path());
assert!(listed.contains(&"MEMORY.md".to_string()));
assert!(listed.contains(&Path::new("notes").join("topic.md").display().to_string()));
assert!(!listed.iter().any(|p| p.contains("generated")));
assert!(!listed.iter().any(|p| p.contains(".git")));
assert!(!listed.iter().any(|p| p == "notes"));
assert!(walk_personal_files(&root.path().join("absent")).is_empty());
}
#[tokio::test]
async fn personal_files_fetch_rejects_synthetic_names() {
for name in ["", " ", ".", "..", "a/b", "a\\b"] {
assert!(fetch_personal_files(name).await.is_none(), "name: {name}");
}
}
#[test]
fn personal_files_render_caps_budgets_and_sorts() {
let paths: Vec<String> = (0..MAX_PERSONAL_FILE_ENTRIES + 50)
.map(|i| format!("f{i:03}.txt"))
.collect();
let rendered = render_personal_file_lines(paths).expect("render");
let lines: Vec<&str> = rendered.lines().collect();
assert_eq!(lines.len(), MAX_PERSONAL_FILE_ENTRIES + 1); assert!(lines[0] < lines[1]);
assert!(rendered.ends_with("use the read tool for the full picture"));
let long = "x".repeat(MAX_PERSONAL_FILES_BYTES / 2);
let rendered =
render_personal_file_lines(vec![long.clone(), long.clone(), long]).expect("render");
assert_eq!(rendered.lines().count(), 2);
assert!(render_personal_file_lines(Vec::new()).is_none());
}
#[test]
fn alarm_lines_show_trigger_and_interval() {
let one_shot = Alarm {
id: "a1".to_string(),
session_id: "s1".to_string(),
user_name: "u".to_string(),
kind: "one-shot".to_string(),
text: "hello".to_string(),
trigger: None,
interval_seconds: None,
next_fire_at: "2026-01-01T00:00:00+00:00".to_string(),
};
let periodic = Alarm {
id: "a2".to_string(),
session_id: "s1".to_string(),
user_name: "u".to_string(),
kind: "periodic".to_string(),
text: "check".to_string(),
trigger: Some(crate::alarms::StoredTrigger::Tool(crate::alarms::Trigger {
tool: "weather".to_string(),
args: serde_json::Map::from_iter([
("city".to_string(), serde_json::json!("Minsk")),
(
"api_key".to_string(),
serde_json::json!("supersecretvalue123"),
),
]),
})),
interval_seconds: Some(3600),
next_fire_at: "2026-01-01T00:00:00+00:00".to_string(),
};
let lines = render_alarm_lines(&[one_shot, periodic]).expect("alarms render");
assert!(lines.contains("every 3600 seconds"));
assert!(lines.contains("trigger: weather "));
assert!(lines.contains("\"city\":\"Minsk\""));
assert!(lines.contains("\"api_key\":\"supersecretvalue123\""));
assert!(lines.contains("next fire:"));
assert!(render_alarm_lines(&[]).is_none());
}
#[test]
fn workspace_lines_summary_collapse() {
let workspaces = vec![
(
crate::workspace::test_ws_named("/data/proj", "proj"),
Some("Multi line\nsummary text".to_string()),
),
(
crate::workspace::test_ws_named("/data/other", "other"),
None,
),
];
let lines = render_workspace_lines(&workspaces).expect("workspaces render");
assert!(lines.contains("- proj (pending): /data/proj"));
assert!(lines.contains("\n Multi line summary text"));
assert!(lines.contains("- other (pending): /data/other"));
assert!(render_workspace_lines(&[]).is_none());
}
}