mod compaction_dump;
pub mod dead_session;
pub(crate) mod image_strip;
mod manager;
pub(crate) mod transcript;
pub(crate) use manager::FinalizeOutcome;
pub(crate) use manager::RewriteOutcome;
pub(crate) use manager::Session;
pub(crate) use transcript::{TRANSCRIPT_REGISTRY, TranscriptSnapshot};
use crate::db::{self, IntoParams, Row, TxGuard, Value, params};
use crate::{ChatMessage, ChatRole, Reasoning, ToolCall, ToolResultPayload};
use anyhow::{Context, Result, anyhow};
use chrono::{DateTime, Utc};
use std::borrow::Cow;
pub(crate) const SUMMARIZATION_THRESHOLD: u64 = 200_000;
pub(crate) const SUMMARIZATION_IMAGE_COUNT: usize = 10;
const RETENTION_PER_SIDE: usize = 3;
fn decoded_tool_calls(msg: &ChatMessage) -> Option<Vec<ToolCall>> {
match decode_native_history_message(msg) {
Some(DecodedNativeHistoryMessage::Assistant { tool_calls, .. }) => tool_calls,
_ => None,
}
}
fn is_tool_call_frame(msg: &ChatMessage) -> bool {
decoded_tool_calls(msg).is_some()
}
fn is_dangling_tool_call_frame(role: ChatRole, content: &str) -> bool {
decoded_tool_calls(&ChatMessage {
role,
content: content.to_string(),
})
.is_some_and(|calls| !calls.is_empty())
}
#[must_use]
fn select_retention_window(history: &[ChatMessage]) -> Vec<ChatMessage> {
let mut users: Vec<(usize, &ChatMessage)> = Vec::new();
let mut assistants: Vec<(usize, &ChatMessage)> = Vec::new();
for (idx, msg) in history.iter().enumerate().rev() {
match msg.role {
ChatRole::User if users.len() < RETENTION_PER_SIDE => users.push((idx, msg)),
ChatRole::Assistant
if !is_tool_call_frame(msg) && assistants.len() < RETENTION_PER_SIDE =>
{
assistants.push((idx, msg));
}
_ => {}
}
if users.len() == RETENTION_PER_SIDE && assistants.len() == RETENTION_PER_SIDE {
break;
}
}
let mut selected: Vec<(usize, &ChatMessage)> = users.into_iter().chain(assistants).collect();
selected.sort_by_key(|(idx, _)| *idx);
selected.into_iter().map(|(_, m)| m.clone()).collect()
}
#[must_use]
pub(crate) fn render_timestamp() -> String {
let now = chrono::Local::now();
format!("{} ({})", now.format("%Y-%m-%d %H:%M:%S"), now.format("%Z"))
}
#[must_use]
pub(crate) fn user_msg_with_ts(content: &str, round_ts: Option<&str>) -> ChatMessage {
let ts = round_ts.map_or_else(render_timestamp, str::to_string);
ChatMessage::user(format!("{content}\n\n<timestamp>{ts}</timestamp>"))
}
crate::define_store! {
pub(crate) static SESSIONS: SessionStore,
expect = "SESSIONS not initialized — call init_all_stores() first",
}
crate::columns! {
SESSION_MESSAGE_COLUMNS [SM] {
ROLE => "role",
CONTENT => "content",
}
}
crate::columns! {
SESSION_LIST_COLUMNS [SL] {
AGENT_ID => "sm.agent_id",
LAST_ACTIVITY => "sm.last_activity",
MESSAGE_COUNT => "sm.message_count",
TOKEN_LENGTH => "sm.token_length",
ROLE => "sm.role",
USER_NAME => "sm.user_name",
WORKSPACE_NAME => "sm.workspace_name",
}
}
const TRANSIENT_AGENT_ID_PREFIXES: &[&str] = &[
"ticket_",
"analyze_",
"research_",
"cleanup_",
"maintainer_",
"discovery_",
];
#[derive(Debug, Clone)]
pub(crate) struct SessionMetadata {
pub agent_id: String,
pub last_activity: DateTime<Utc>,
pub message_count: usize,
pub token_length: Option<u64>,
pub role: Option<String>,
pub user_name: Option<String>,
pub workspace_name: Option<String>,
}
#[derive(Debug, Clone)]
pub(crate) struct SessionContext {
pub channel: String,
pub user_name: String,
pub workspace_name: String,
pub role: String,
}
fn session_metadata_from_row(
agent_id: &str,
activity_str: &str,
count: i64,
token_length: Option<i64>,
role: Option<String>,
user_name: Option<String>,
workspace_name: Option<String>,
) -> Result<SessionMetadata> {
let last_activity = db::parse_utc_timestamp(activity_str).with_context(|| {
format!("invalid last_activity {activity_str:?} for session {agent_id}")
})?;
Ok(SessionMetadata {
agent_id: agent_id.to_string(),
last_activity,
message_count: usize::try_from(count).unwrap_or(0),
token_length: token_length.and_then(|t| u64::try_from(t).ok()),
role,
user_name,
workspace_name,
})
}
const INSERT_SESSION_ROW: &str =
"INSERT INTO sessions (agent_id, role, content, created_at) VALUES (?1, ?2, ?3, ?4)";
async fn insert_messages_in_transaction(
tx: &TxGuard<'_>,
agent_id: &str,
messages: &[ChatMessage],
context: Option<(&str, &str, &str, &str)>,
replace: bool,
) -> Result<()> {
let now = db::now();
for msg in messages {
tx.execute(
INSERT_SESSION_ROW,
params![
agent_id,
msg.role.to_string(),
msg.content.clone(),
now.clone()
],
)
.await?;
}
let count = i64::try_from(messages.len()).context("message batch exceeds i64")?;
match context {
Some((channel, user_name, workspace_name, role)) => {
let created_at = now.clone();
let count_clause = message_count_clause(replace);
tx.upsert_row(
&format!(
"UPDATE session_metadata SET last_activity = ?2, channel = ?4, \
user_name = ?5, workspace_name = ?6, role = ?7, {count_clause} \
WHERE agent_id = ?1"
),
|| {
params![
agent_id,
now.as_str(),
count,
channel,
user_name,
workspace_name,
role,
]
},
"INSERT INTO session_metadata (agent_id, last_activity, message_count, \
channel, user_name, workspace_name, role, created_at) \
VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8) \
ON CONFLICT(agent_id) DO NOTHING",
params![
agent_id,
now.as_str(),
count,
channel,
user_name,
workspace_name,
role,
created_at.as_str(),
],
)
.await?;
}
None => upsert_message_count(tx, agent_id, count, &now, replace).await?,
}
Ok(())
}
fn message_count_clause(replace: bool) -> &'static str {
if replace {
"message_count = ?3"
} else {
"message_count = message_count + ?3"
}
}
async fn upsert_message_count(
tx: &TxGuard<'_>,
agent_id: &str,
count: i64,
now: &str,
replace: bool,
) -> Result<()> {
let count_clause = message_count_clause(replace);
tx.upsert_row(
&format!(
"UPDATE session_metadata SET last_activity = ?2, {count_clause} \
WHERE agent_id = ?1"
),
|| params![agent_id, now, count],
"INSERT INTO session_metadata (agent_id, last_activity, message_count, created_at) \
VALUES (?1, ?2, ?3, ?4) \
ON CONFLICT(agent_id) DO NOTHING",
params![agent_id, now, count, now],
)
.await?;
Ok(())
}
async fn query_map_collect<T, E>(
conn: &db::Connection,
sql: &str,
params: impl IntoParams + Send + 'static,
row_parser: impl FnMut(&Row) -> std::result::Result<T, E> + Send + 'static,
warn_context: &str,
agent_id: Option<&str>,
) -> Vec<T>
where
T: Send + 'static,
E: std::fmt::Display + Send + Sync + 'static,
{
let rows = match conn.query_map(sql, params, row_parser).await {
Ok(rows) => rows,
Err(e) => {
tracing::warn!(error = %e, agent_id, "{warn_context}: query failed, returning empty");
return Vec::new();
}
};
rows.into_iter()
.filter_map(|r| match r {
Ok(val) => Some(val),
Err(e) => {
tracing::warn!(error = %e, agent_id, "{warn_context}: row decode failed, skipping");
None
}
})
.collect()
}
fn sessions_list_sql(where_clause: &str) -> String {
format!(
"SELECT {SESSION_LIST_COLUMNS} \
FROM session_metadata sm \
{where_clause} \
ORDER BY sm.last_activity DESC",
)
}
fn session_messages_sql() -> String {
format!("SELECT {SESSION_MESSAGE_COLUMNS} FROM sessions WHERE agent_id = ?1 ORDER BY id ASC")
}
fn session_metadata_row(row: &Row) -> Result<SessionMetadata> {
session_metadata_from_row(
&row.get::<String>(COL_SL_AGENT_ID)?,
&row.get::<String>(COL_SL_LAST_ACTIVITY)?,
row.get::<i64>(COL_SL_MESSAGE_COUNT)?,
row.get::<Option<i64>>(COL_SL_TOKEN_LENGTH)?,
row.get::<Option<String>>(COL_SL_ROLE)?,
row.get::<Option<String>>(COL_SL_USER_NAME)?,
row.get::<Option<String>>(COL_SL_WORKSPACE_NAME)?,
)
}
async fn list_sessions_where(
conn: &db::Connection,
where_clause: &str,
params: impl IntoParams + Send + 'static,
warn_context: &str,
) -> Vec<SessionMetadata> {
query_map_collect(
conn,
&sessions_list_sql(where_clause),
params,
session_metadata_row,
warn_context,
None,
)
.await
}
fn session_message_from_row(row: &Row) -> Result<ChatMessage> {
Ok(ChatMessage {
role: row
.get::<String>(COL_SM_ROLE)?
.parse::<ChatRole>()
.map_err(|e| anyhow!(e))?,
content: row.get(COL_SM_CONTENT)?,
})
}
enum MetadataColumn {
ActiveModels,
TokenLength,
SleepEnded,
}
impl MetadataColumn {
fn as_str(&self) -> &'static str {
match self {
Self::ActiveModels => "active_models",
Self::TokenLength => "token_length",
Self::SleepEnded => "sleep_ended",
}
}
}
#[derive(Debug, Clone)]
pub(crate) struct SettledRow {
pub role: String,
pub content: String,
pub created_at: String,
}
fn build_settle_sequence(
captured: &[SettledRow],
frame_calls: &[crate::ToolCall],
results: &[(String, String)],
follow_up: &[ChatMessage],
now: &str,
) -> Result<Vec<SettledRow>> {
let mut captured_by_call_id: std::collections::HashMap<String, usize> =
std::collections::HashMap::new();
for (i, row) in captured.iter().enumerate() {
if row.role == "tool"
&& let Ok(payload) = serde_json::from_str::<crate::ToolResultPayload>(&row.content)
{
captured_by_call_id.insert(payload.tool_call_id, i);
}
}
let new_results: std::collections::HashMap<&str, &str> = results
.iter()
.map(|(id, content)| (id.as_str(), content.as_str()))
.collect();
let mut sequence = Vec::new();
for call in frame_calls {
if let Some(content) = new_results.get(call.id.as_str()) {
let payload = crate::ToolResultPayload {
tool_call_id: call.id.clone(),
content: (*content).to_string(),
};
sequence.push(SettledRow {
role: "tool".to_string(),
content: serde_json::to_string(&payload)?,
created_at: now.to_string(),
});
} else if let Some(&i) = captured_by_call_id.get(&call.id) {
sequence.push(captured[i].clone());
}
}
for msg in follow_up {
sequence.push(SettledRow {
role: msg.role.to_string(),
content: msg.content.clone(),
created_at: now.to_string(),
});
}
Ok(sequence)
}
impl SessionStore {
pub(crate) async fn load(&self, agent_id: &str) -> Vec<ChatMessage> {
query_map_collect(
&self.conn,
&session_messages_sql(),
params![agent_id],
session_message_from_row,
"load session",
Some(agent_id),
)
.await
}
pub(crate) async fn load_checked(&self, agent_id: &str) -> Result<Vec<ChatMessage>> {
self.conn
.query_map_strict(
&session_messages_sql(),
params![agent_id],
session_message_from_row,
)
.await
}
pub(crate) async fn has_content(&self, agent_id: &str) -> bool {
self.conn
.query_optional(
"SELECT 1 FROM sessions WHERE agent_id = ?1 AND length(content) > 0 LIMIT 1",
params![agent_id],
|_| Ok::<(), anyhow::Error>(()),
)
.await
.ok()
.flatten()
.is_some()
}
async fn append_messages(
&self,
agent_id: &str,
messages: &[ChatMessage],
replace: bool,
context: Option<(&str, &str, &str, &str)>,
) -> Result<()> {
let tx = self.conn.begin_tx().await?;
if replace {
tx.execute(
"DELETE FROM sessions WHERE agent_id = ?1",
params![agent_id],
)
.await?;
}
insert_messages_in_transaction(&tx, agent_id, messages, context, replace).await?;
tx.commit().await?;
Ok(())
}
async fn batch_append(&self, agent_id: &str, messages: &[ChatMessage]) -> Result<()> {
self.append_messages(agent_id, messages, false, None).await
}
async fn batch_append_with_context(
&self,
agent_id: &str,
messages: &[ChatMessage],
channel: &str,
user_name: &str,
workspace_name: &str,
role: &str,
) -> Result<()> {
self.append_messages(
agent_id,
messages,
false,
Some((channel, user_name, workspace_name, role)),
)
.await
}
async fn append_with_context(
&self,
agent_id: &str,
message: &ChatMessage,
channel: &str,
user_name: &str,
workspace_name: &str,
role: &str,
) -> Result<()> {
self.batch_append_with_context(
agent_id,
std::slice::from_ref(message),
channel,
user_name,
workspace_name,
role,
)
.await
}
async fn replace_messages(&self, agent_id: &str, messages: &[ChatMessage]) -> Result<()> {
self.append_messages(agent_id, messages, true, None).await
}
pub(crate) async fn rewrite_last_user_message(
&self,
agent_id: &str,
content: &str,
) -> Result<()> {
let tx = self.conn.begin_tx().await?;
let changed = tx
.execute(
"UPDATE sessions SET content = ?1 WHERE agent_id = ?2 AND id = (
SELECT MAX(id) FROM sessions WHERE agent_id = ?2 AND role = 'user'
)",
params![content, agent_id],
)
.await?;
tx.commit().await?;
if changed == 0 {
anyhow::bail!("no user message row to rewrite for agent {agent_id}");
}
Ok(())
}
#[cfg(test)]
pub(crate) async fn settle_tool_results(
&self,
agent_id: &str,
results: &[(String, String)],
follow_up: &[ChatMessage],
) -> Result<Vec<SettledRow>> {
let tx = self.conn.begin_tx().await?;
let rows = self
.settle_tool_results_tx(&tx, agent_id, results, follow_up)
.await?;
tx.commit().await?;
Ok(rows)
}
pub(crate) async fn settle_tool_results_tx(
&self,
tx: &TxGuard<'_>,
agent_id: &str,
results: &[(String, String)],
follow_up: &[ChatMessage],
) -> Result<Vec<SettledRow>> {
let candidates: Vec<(i64, String)> = tx
.query(
"SELECT id, content FROM sessions \
WHERE agent_id = ?1 AND role = 'assistant' AND content LIKE '%\"tool_calls\"%' \
ORDER BY id DESC",
params![agent_id],
)
.await
.context("find tool-call frame candidates")?
.into_iter()
.map(|r| Result::<_, turso::Error>::Ok((r.get::<i64>(0)?, r.get::<String>(1)?)))
.collect::<Result<Vec<_>, _>>()
.context("decode tool-call frame candidate rows")?;
let mut frame: Option<(i64, String)> = None;
for (id, content) in candidates {
let msg = ChatMessage {
role: ChatRole::Assistant,
content,
};
if is_tool_call_frame(&msg) {
frame = Some((id, msg.content));
break;
}
}
let Some((frame_row_id, frame_content)) = frame else {
tracing::warn!(
agent_id,
"settle_tool_results: no tool-call frame found — no-op"
);
return Ok(Vec::new());
};
let frame_calls = {
let msg = ChatMessage {
role: ChatRole::Assistant,
content: frame_content,
};
let Some(DecodedNativeHistoryMessage::Assistant {
tool_calls: Some(calls),
..
}) = decode_native_history_message(&msg)
else {
unreachable!("is_tool_call_frame verified the frame decodes");
};
calls
};
let captured: Vec<SettledRow> = tx
.query_map_strict(
"SELECT role, content, created_at FROM sessions WHERE agent_id = ?1 AND id > ?2 ORDER BY id",
params![agent_id, frame_row_id],
|row| {
Ok::<_, anyhow::Error>(SettledRow {
role: row.get(0)?,
content: row.get(1)?,
created_at: row.get(2)?,
})
},
)
.await
.context("capture rows after tool-call frame")?;
tx.execute(
"DELETE FROM sessions WHERE agent_id = ?1 AND id > ?2",
params![agent_id, frame_row_id],
)
.await
.context("delete rows after tool-call frame")?;
let now = db::now();
let sequence = build_settle_sequence(&captured, &frame_calls, results, follow_up, &now)?;
for row in &sequence {
tx.execute(
INSERT_SESSION_ROW,
params![
agent_id,
row.role.clone(),
row.content.clone(),
row.created_at.clone()
],
)
.await
.context("insert settled session row")?;
}
let count = i64::try_from(results.len() + follow_up.len())
.context("settle result count exceeds i64")?;
upsert_message_count(tx, agent_id, count, &now, false)
.await
.context("bump settle message count")?;
Ok(sequence)
}
pub(crate) async fn delete(&self, agent_id: &str) -> Result<bool> {
let tx = self.conn.begin_tx().await?;
let deleted = tx
.execute(
"DELETE FROM sessions WHERE agent_id = ?1",
params![agent_id],
)
.await?;
tx.execute(
"DELETE FROM session_metadata WHERE agent_id = ?1",
params![agent_id],
)
.await?;
tx.commit().await?;
Ok(deleted > 0)
}
pub(crate) async fn truncate_from(&self, agent_id: &str, index: usize) -> Result<usize> {
let offset = i64::try_from(index).context("truncate index exceeds i64")?;
let tx = self.conn.begin_tx().await?;
let deleted = tx
.execute(
"DELETE FROM sessions WHERE agent_id = ?1 AND id >= \
(SELECT id FROM sessions WHERE agent_id = ?1 ORDER BY id ASC LIMIT 1 OFFSET ?2)",
params![agent_id, offset],
)
.await?;
if deleted > 0 {
let removed = i64::try_from(deleted).context("removed row count exceeds i64")?;
tx.execute(
"UPDATE session_metadata SET message_count = MAX(message_count - ?2, 0) \
WHERE agent_id = ?1",
params![agent_id, removed],
)
.await?;
}
let removed = usize::try_from(deleted).context("removed row count exceeds usize")?;
tx.commit().await?;
Ok(removed)
}
pub(crate) async fn list_sessions_with_metadata(&self) -> Vec<SessionMetadata> {
list_sessions_where(&self.conn, "", (), "list sessions").await
}
pub(crate) async fn list_sessions_with_metadata_checked(&self) -> Result<Vec<SessionMetadata>> {
self.conn
.query_map_strict(&sessions_list_sql(""), (), session_metadata_row)
.await
}
async fn list_sessions_with_metadata_excluding(
&self,
exclude_prefixes: &[&str],
) -> Vec<SessionMetadata> {
if exclude_prefixes.is_empty() {
return self.list_sessions_with_metadata().await;
}
let where_clause = exclude_prefixes
.iter()
.enumerate()
.map(|(i, _)| format!("sm.agent_id NOT LIKE ?{}", i + 1))
.collect::<Vec<_>>()
.join(" AND ");
let params: Vec<db::Value> = exclude_prefixes
.iter()
.map(|p| db::Value::Text(format!("{p}%")))
.collect();
list_sessions_where(
&self.conn,
&format!("WHERE {where_clause}"),
params,
"list sessions (excluding prefixes)",
)
.await
}
pub(crate) async fn get_last_message_tail(&self, agent_id: &str) -> Option<(ChatRole, String)> {
let rows = self
.conn
.query(
&format!(
"SELECT {SESSION_MESSAGE_COLUMNS} FROM sessions \
WHERE agent_id = ?1 ORDER BY id DESC LIMIT 1"
),
params![agent_id],
)
.await
.ok()?;
rows.first().and_then(|row| {
let role_str: String = row.get(COL_SM_ROLE).ok()?;
let content: String = row.get(COL_SM_CONTENT).ok()?;
Some((role_str.parse::<ChatRole>().ok()?, content))
})
}
async fn get_session_context(&self, agent_id: &str) -> Option<SessionContext> {
let rows = self
.conn
.query(
"SELECT channel, user_name, workspace_name, role FROM session_metadata WHERE agent_id = ?1",
params![agent_id],
)
.await
.ok()?;
rows.first().and_then(|row| {
let channel: Option<String> = row.get(0).ok();
let user_name: Option<String> = row.get(1).ok();
let workspace_name: Option<String> = row.get(2).ok();
let role: Option<String> = row.get(3).ok();
Some(SessionContext {
channel: channel?,
user_name: user_name?,
workspace_name: workspace_name?,
role: role?,
})
})
}
async fn set_metadata_value(
&self,
agent_id: &str,
column: MetadataColumn,
value: Value,
) -> Result<()> {
let col = column.as_str();
let now = db::now();
let created_at = now.clone();
self.conn
.upsert_row(
&format!("UPDATE session_metadata SET {col} = ?1 WHERE agent_id = ?2"),
|| params![value.clone(), agent_id],
&format!(
"INSERT INTO session_metadata (agent_id, last_activity, {col}, created_at) \
VALUES (?1, ?2, ?3, ?4) \
ON CONFLICT(agent_id) DO NOTHING"
),
params![agent_id, now.as_str(), value.clone(), created_at.as_str()],
)
.await?;
Ok(())
}
async fn get_metadata_value<X>(
&self,
agent_id: &str,
column: MetadataColumn,
warn_context: &str,
map: impl FnOnce(&Row) -> std::result::Result<Option<X>, ::turso::Error> + Send + 'static,
) -> Option<X>
where
X: Send + 'static,
{
let col = column.as_str();
let sql = format!("SELECT {col} FROM session_metadata WHERE agent_id = ?1");
match self.conn.query_optional(&sql, params![agent_id], map).await {
Ok(value) => value.flatten(),
Err(e) => {
tracing::warn!(agent_id = %agent_id, error = %e, "{warn_context}");
None
}
}
}
async fn set_active_models(&self, agent_id: &str, snapshot: Option<&str>) -> Result<()> {
self.set_metadata_value(agent_id, MetadataColumn::ActiveModels, snapshot.into())
.await
}
async fn get_active_models(&self, agent_id: &str) -> Option<String> {
self.get_metadata_value(
agent_id,
MetadataColumn::ActiveModels,
"Failed to read active-models snapshot",
|row| row.get::<Option<String>>(0),
)
.await
}
pub(crate) async fn set_token_length(
&self,
agent_id: &str,
token_length: Option<u64>,
) -> Result<()> {
let bound: Option<i64> = token_length.map(i64::try_from).transpose()?;
self.set_metadata_value(agent_id, MetadataColumn::TokenLength, bound.into())
.await
}
async fn get_token_length(&self, agent_id: &str) -> Option<u64> {
self.get_metadata_value(
agent_id,
MetadataColumn::TokenLength,
"Failed to read session token length",
|row| {
row.get::<Option<i64>>(0)
.map(|opt| opt.and_then(|v| u64::try_from(v).ok()))
},
)
.await
}
pub(crate) async fn set_sleep_ended(&self, agent_id: &str, ended: bool) -> Result<()> {
let flag = if ended { "1" } else { "0" };
self.set_metadata_value(
agent_id,
MetadataColumn::SleepEnded,
db::Value::Text(flag.into()),
)
.await
}
pub(crate) async fn get_sleep_ended(&self, agent_id: &str) -> bool {
self.get_metadata_value(
agent_id,
MetadataColumn::SleepEnded,
"Failed to read sleep-ended flag",
|row| row.get::<Option<String>>(0),
)
.await
.is_some_and(|flag| flag == "1")
}
}
pub async fn cleanup_old_transient_sessions(cutoff: &str) -> Result<u64> {
let session_store = store();
let tx = session_store.conn.begin_tx().await?;
let likes = TRANSIENT_AGENT_ID_PREFIXES
.iter()
.map(|_| "agent_id LIKE ?")
.collect::<Vec<_>>()
.join(" OR ");
let prefix_patterns = format!("({likes})");
let mut build_params = vec![Value::Text(cutoff.to_string())];
build_params.extend(
TRANSIENT_AGENT_ID_PREFIXES
.iter()
.map(|prefix| Value::Text(format!("{prefix}%"))),
);
let safety_predicate = format!(
"last_activity < ? AND {prefix_patterns} \
AND agent_id NOT IN (SELECT agent_id FROM agents)"
);
tx.execute(
&format!(
"DELETE FROM sessions WHERE agent_id IN ( \
SELECT agent_id FROM session_metadata WHERE {safety_predicate})"
),
build_params.clone(),
)
.await?;
let deleted = tx
.execute(
&format!("DELETE FROM session_metadata WHERE {safety_predicate}"),
build_params,
)
.await?;
tx.commit().await?;
Ok(deleted)
}
fn reserved_agent_id_prefixes() -> impl Iterator<Item = &'static str> {
std::iter::once("manager_").chain(TRANSIENT_AGENT_ID_PREFIXES.iter().copied())
}
fn starts_with_ignore_ascii_case(s: &str, prefix: &str) -> bool {
s.get(..prefix.len())
.is_some_and(|head| head.eq_ignore_ascii_case(prefix))
}
pub(crate) fn normalize_user_name<'a>(user_name: &'a str, context: &str) -> &'a str {
if user_name.is_empty() {
tracing::warn!(
context = %context,
"empty user_name — falling back to seeded 'admin'",
);
crate::users::ADMIN_USER_NAME
} else {
user_name
}
}
#[must_use]
fn user_segment_collides(user_name: &str) -> bool {
let touches_reserved = reserved_agent_id_prefixes()
.any(|prefix| starts_with_ignore_ascii_case(user_name, prefix.trim_end_matches('_')));
touches_reserved || starts_with_ignore_ascii_case(user_name, "user_")
}
#[must_use]
fn safe_user_segment(user_name: &str) -> Cow<'_, str> {
if user_segment_collides(user_name) {
Cow::Owned(format!("user_{user_name}"))
} else {
Cow::Borrowed(user_name)
}
}
#[must_use]
fn direct_agent_id(user_name: &str, role: &str, ws_name: &str) -> String {
let user_name = normalize_user_name(user_name, "direct_agent_id");
if ws_name.strip_prefix("personal:") == Some(user_name) {
return format!("{}_personal:{role}", safe_user_segment(user_name));
}
format!("{}_{}_{}", safe_user_segment(user_name), ws_name, role)
}
#[must_use]
pub(crate) fn ticket_agent_id(ticket_id: &str, role: &str) -> String {
format!("ticket_{ticket_id}_{role}")
}
#[must_use]
pub(crate) fn manager_agent_id(ws_name: &str) -> String {
format!("manager_{ws_name}")
}
#[must_use]
pub(crate) fn resolve_agent_id(user_name: &str, role: &str, ws_name: &str) -> String {
if role == "manager" {
manager_agent_id(ws_name)
} else {
direct_agent_id(user_name, role, ws_name)
}
}
pub async fn clear_session(user_name: &str, role: &str, ws_name: &str) -> anyhow::Result<String> {
let agent_id = resolve_agent_id(user_name, role, ws_name);
crate::jobs::abandon_session_jobs(&agent_id).await?;
crate::session::store().delete(&agent_id).await?;
Ok("Session cleared. Starting fresh.".to_string())
}
#[must_use]
fn transient_agent_id(prefix: &str, ws_name: &str, label: &str) -> String {
let suffix = crate::generate_suffix();
format!("{prefix}{ws_name}_{suffix}_{label}")
}
#[must_use]
pub(crate) fn maintainer_agent_id(ws_name: &str) -> String {
format!("maintainer_{ws_name}")
}
#[must_use]
pub(crate) fn analyze_agent_id(ws_name: &str, role: &str) -> String {
transient_agent_id("analyze_", ws_name, role)
}
#[must_use]
pub(crate) fn research_agent_id(ws_name: &str, label: &str) -> String {
transient_agent_id("research_", ws_name, label)
}
#[must_use]
pub(crate) fn discovery_agent_id(ws_name: &str, role: &str) -> String {
transient_agent_id("discovery_", ws_name, role)
}
#[cfg(test)]
mod tests {
use super::*;
use std::sync::atomic::{AtomicU32, Ordering};
static TEST_ID: AtomicU32 = AtomicU32::new(0);
fn unique_key() -> String {
format!("s{}", TEST_ID.fetch_add(1, Ordering::Relaxed))
}
async fn session_with_init(agent_id: &str, msg: &str, ws: &crate::Workspace) -> Session {
let mut session = Session::default();
session.init(agent_id).await.unwrap();
session
.append_turn_message(
agent_id,
msg,
ws,
&crate::Role::Assistant,
None,
"gui",
"tester",
None,
false,
)
.await
.unwrap();
session
}
#[test]
fn user_msg_timestamp_suffix_format() {
let msg = user_msg_with_ts("task text", Some("2026-01-01 00:00:00 (UTC)"));
assert_eq!(
msg.content,
"task text\n\n<timestamp>2026-01-01 00:00:00 (UTC)</timestamp>"
);
let fresh = user_msg_with_ts("task text", None);
assert!(
fresh.content.starts_with("task text\n\n<timestamp>")
&& fresh.content.ends_with("</timestamp>"),
"fresh stamp must still be a suffix: {}",
fresh.content
);
}
#[tokio::test]
async fn session_store_create_and_load() {
crate::util::test::init_test_stores().await;
let k = unique_key();
store()
.batch_append(&k, &[ChatMessage::user("hello")])
.await
.unwrap();
let msgs = store().load(&k).await;
assert_eq!(msgs.len(), 1);
assert_eq!(msgs[0].content, "hello");
}
#[tokio::test]
async fn checked_readers_propagate_undecodable_rows() {
crate::util::test::init_test_stores().await;
let k = unique_key();
store()
.batch_append(&k, &[ChatMessage::user("hello")])
.await
.unwrap();
let checked = store().load_checked(&k).await.unwrap();
assert_eq!(checked.len(), 1);
assert_eq!(checked[0].content, "hello");
let listed = store().list_sessions_with_metadata_checked().await.unwrap();
assert!(listed.iter().any(|s| s.agent_id == k));
store()
.conn
.execute(
"INSERT INTO sessions (agent_id, role, content, created_at) \
VALUES (?1, 'bogus', 'x', ?2)",
params![k.as_str(), db::now()],
)
.await
.unwrap();
assert!(
store().load_checked(&k).await.is_err(),
"undecodable row must be an error for the checked reader"
);
assert_eq!(
store().load(&k).await.len(),
1,
"the infallible reader keeps its skip-undecodable-rows contract"
);
}
#[tokio::test]
async fn session_store_replace_messages() {
crate::util::test::init_test_stores().await;
let k = unique_key();
store()
.batch_append(&k, &[ChatMessage::user("old")])
.await
.unwrap();
store()
.replace_messages(&k, &[ChatMessage::user("new")])
.await
.unwrap();
let msgs = store().load(&k).await;
assert_eq!(msgs.len(), 1);
assert_eq!(msgs[0].content, "new");
}
#[tokio::test]
async fn session_store_rewrite_last_user_message() {
crate::util::test::init_test_stores().await;
let k = unique_key();
store()
.batch_append(
&k,
&[
ChatMessage::system("role"),
ChatMessage::user("[IMAGE:/tmp/a.png] first"),
ChatMessage::assistant("answer"),
ChatMessage::user("[IMAGE:/tmp/b.png] second"),
],
)
.await
.unwrap();
store()
.rewrite_last_user_message(&k, "rewritten")
.await
.unwrap();
let msgs = store().load(&k).await;
assert_eq!(msgs[3].role, crate::ChatRole::User);
assert_eq!(msgs[3].content, "rewritten", "last user row rewritten");
assert_eq!(
msgs[1].content, "[IMAGE:/tmp/a.png] first",
"earlier user row untouched"
);
assert_eq!(msgs[2].content, "answer", "assistant row untouched");
}
#[tokio::test]
async fn session_store_rewrite_last_user_message_targets_only_user_rows() {
crate::util::test::init_test_stores().await;
let k = unique_key();
store()
.batch_append(
&k,
&[
ChatMessage::system("role"),
ChatMessage::user("only user"),
ChatMessage::assistant("last row is assistant"),
],
)
.await
.unwrap();
store()
.rewrite_last_user_message(&k, "rewritten")
.await
.unwrap();
let msgs = store().load(&k).await;
assert_eq!(
msgs[1].content, "rewritten",
"last USER row rewritten, not the trailing assistant row"
);
assert_eq!(msgs[2].content, "last row is assistant");
}
#[tokio::test]
async fn session_store_rewrite_last_user_message_no_user_row_errors() {
crate::util::test::init_test_stores().await;
let k = unique_key();
store()
.batch_append(&k, &[ChatMessage::assistant("only assistant")])
.await
.unwrap();
let err = store()
.rewrite_last_user_message(&k, "x")
.await
.unwrap_err();
assert!(err.to_string().contains("no user message row"), "{err}");
}
#[tokio::test]
async fn session_store_delete() {
crate::util::test::init_test_stores().await;
let k = unique_key();
store()
.batch_append(&k, &[ChatMessage::user("a")])
.await
.unwrap();
assert!(store().delete(&k).await.unwrap());
assert!(!store().delete(&k).await.unwrap());
}
#[tokio::test]
async fn session_store_truncate_from_cuts_inclusively() {
crate::util::test::init_test_stores().await;
let k = unique_key();
store()
.batch_append(
&k,
&[
ChatMessage::user("a"),
ChatMessage::assistant("b"),
ChatMessage::user("c"),
ChatMessage::assistant("d"),
],
)
.await
.unwrap();
assert_eq!(store().truncate_from(&k, 2).await.unwrap(), 2);
let kept: Vec<String> = store()
.load(&k)
.await
.into_iter()
.map(|m| m.content)
.collect();
assert_eq!(kept, ["a", "b"]);
let meta = store()
.list_sessions_with_metadata()
.await
.into_iter()
.find(|s| s.agent_id == k)
.expect("session listed");
assert_eq!(meta.message_count, 2, "count follows the deleted rows");
assert_eq!(store().truncate_from(&k, 9).await.unwrap(), 0);
assert_eq!(store().load(&k).await.len(), 2);
}
#[tokio::test]
async fn clear_session_abandons_owned_sync_jobs() {
crate::util::test::init_test_stores().await;
let conn = &crate::session::store().conn;
let agent_id = resolve_agent_id("clear_user", "engineer", "clear_ws");
crate::jobs::spawn_job(
conn,
"clear_sess_sync",
"task",
"clear_ws",
"clear_user",
"telegram",
crate::Role::Engineer,
&[],
&crate::jobs::SpawnChild::Analyze,
Some(&agent_id),
)
.await
.unwrap();
crate::session::clear_session("clear_user", "engineer", "clear_ws")
.await
.unwrap();
let rows = conn
.query("SELECT id FROM jobs WHERE id = 'clear_sess_sync'", ())
.await
.unwrap();
assert!(
rows.is_empty(),
"clear_session abandoned the session's sync job"
);
}
#[tokio::test]
async fn session_context_roundtrip() {
crate::util::test::init_test_stores().await;
let agent_id = unique_key();
assert!(store().get_session_context(&agent_id).await.is_none());
store()
.append_with_context(
&agent_id,
&ChatMessage::user("hello"),
"gui",
"alice",
"work",
"engineer",
)
.await
.unwrap();
let ctx = store()
.get_session_context(&agent_id)
.await
.expect("should have context after set");
assert_eq!(ctx.channel, "gui");
assert_eq!(ctx.user_name, "alice");
assert_eq!(ctx.workspace_name, "work");
assert_eq!(ctx.role, "engineer");
store()
.append_with_context(
&agent_id,
&ChatMessage::user("hello again"),
"telegram",
"bob",
"project-x",
"analyst",
)
.await
.unwrap();
let ctx = store()
.get_session_context(&agent_id)
.await
.expect("should have updated context");
assert_eq!(ctx.channel, "telegram");
assert_eq!(ctx.user_name, "bob");
assert_eq!(ctx.workspace_name, "project-x");
assert_eq!(ctx.role, "analyst");
}
#[tokio::test]
async fn session_get_last_message_tail() {
crate::util::test::init_test_stores().await;
let agent_id = unique_key();
assert!(store().get_last_message_tail(&agent_id).await.is_none());
store()
.batch_append(&agent_id, &[ChatMessage::user("hello")])
.await
.unwrap();
assert_eq!(
store().get_last_message_tail(&agent_id).await,
Some((ChatRole::User, "hello".into()))
);
store()
.batch_append(&agent_id, &[ChatMessage::assistant("world")])
.await
.unwrap();
assert_eq!(
store().get_last_message_tail(&agent_id).await,
Some((ChatRole::Assistant, "world".into()))
);
let frame = crate::providers::reasoning::assistant_replay_payload(
Some(""),
&[crate::ToolCall {
id: "call_tail".into(),
name: "read".into(),
arguments: serde_json::json!({"path": "x"}),
}],
None,
)
.to_string();
store()
.batch_append(
&agent_id,
&[
ChatMessage::user("read the file"),
ChatMessage::assistant(frame.clone()),
],
)
.await
.unwrap();
let (role, content) = store()
.get_last_message_tail(&agent_id)
.await
.expect("tail should be readable");
assert_eq!(role, ChatRole::Assistant);
assert_eq!(content, frame);
assert!(super::is_dangling_tool_call_frame(
ChatRole::Assistant,
&content
));
}
#[tokio::test]
async fn session_token_length_roundtrip() {
crate::util::test::init_test_stores().await;
let k = unique_key();
assert_eq!(store().get_token_length(&k).await, None);
store().set_token_length(&k, Some(12_345)).await.unwrap();
assert_eq!(store().get_token_length(&k).await, Some(12_345));
store().set_token_length(&k, Some(67_890)).await.unwrap();
assert_eq!(store().get_token_length(&k).await, Some(67_890));
store().set_token_length(&k, None).await.unwrap();
assert_eq!(store().get_token_length(&k).await, None);
}
#[tokio::test]
async fn message_count_append_matches_rows_and_token_length_flows() {
crate::util::test::init_test_stores().await;
let k = unique_key();
store()
.append_with_context(
&k,
&ChatMessage::user("u1"),
"test_channel",
"test_user",
"test_ws",
"engineer",
)
.await
.unwrap();
store()
.batch_append(
&k,
&[
ChatMessage::assistant("a1"),
ChatMessage::tool_result("t1", "r1"),
],
)
.await
.unwrap();
let mine = store()
.list_sessions_with_metadata()
.await
.into_iter()
.find(|s| s.agent_id == k)
.expect("session listed");
assert_eq!(
mine.message_count,
store().load(&k).await.len(),
"denormalized count must equal the session row count"
);
assert_eq!(mine.message_count, 3);
assert_eq!(mine.token_length, None);
store().set_token_length(&k, Some(12_300)).await.unwrap();
let mine = store()
.list_sessions_with_metadata()
.await
.into_iter()
.find(|s| s.agent_id == k)
.expect("session listed");
assert_eq!(mine.token_length, Some(12_300));
assert_eq!(mine.message_count, 3);
}
#[tokio::test]
async fn message_count_replace_sets_not_increments() {
crate::util::test::init_test_stores().await;
let k = unique_key();
store()
.batch_append(
&k,
&[
ChatMessage::user("u1"),
ChatMessage::assistant("a1"),
ChatMessage::tool_result("t1", "r1"),
],
)
.await
.unwrap();
let count = async |agent: &str| {
let s = store()
.list_sessions_with_metadata()
.await
.into_iter()
.find(|s| s.agent_id == agent)
.expect("session listed");
s.message_count
};
assert_eq!(count(&k).await, 3);
store()
.replace_messages(
&k,
&[ChatMessage::system("prompt"), ChatMessage::user("u1")],
)
.await
.unwrap();
assert_eq!(count(&k).await, 2);
assert_eq!(store().load(&k).await.len(), 2);
store()
.replace_messages(
&k,
&[
ChatMessage::system("p"),
ChatMessage::user("u"),
ChatMessage::assistant("a"),
ChatMessage::assistant("a2"),
],
)
.await
.unwrap();
assert_eq!(count(&k).await, 4);
assert_eq!(store().load(&k).await.len(), 4);
}
#[tokio::test]
async fn message_count_ttl_cleanup_and_delete_remove_counts_with_sessions() {
crate::util::test::init_test_stores().await;
let transient = format!("ticket_{}", unique_key());
store()
.append_with_context(
&transient,
&ChatMessage::user("u1"),
"test_channel",
"test_user",
"test_ws",
"engineer",
)
.await
.unwrap();
store()
.batch_append(&transient, &[ChatMessage::assistant("a1")])
.await
.unwrap();
assert!(store().delete(&transient).await.unwrap());
assert!(
!store()
.list_sessions_with_metadata()
.await
.iter()
.any(|s| s.agent_id == transient),
"deleted session (and its count) must leave the list"
);
let stale = format!("ticket_{}", unique_key());
store()
.append_with_context(
&stale,
&ChatMessage::user("u1"),
"test_channel",
"test_user",
"test_ws",
"engineer",
)
.await
.unwrap();
store()
.batch_append(
&stale,
&[ChatMessage::assistant("a1"), ChatMessage::assistant("a2")],
)
.await
.unwrap();
let deleted = cleanup_old_transient_sessions(&crate::db::now())
.await
.unwrap();
assert!(deleted >= 1, "TTL cleanup must remove the stale session");
assert!(
!store()
.list_sessions_with_metadata()
.await
.iter()
.any(|s| s.agent_id == stale),
"TTL-cleaned session (and its count) must leave the list"
);
}
#[tokio::test]
async fn session_init_loads_persisted_token_length() {
crate::util::test::init_test_stores().await;
let k = unique_key();
store().set_token_length(&k, Some(123_456)).await.unwrap();
let ws = crate::workspace::test_ws_named("/_test_token_length", "token_length");
let session = session_with_init(&k, "", &ws).await;
assert_eq!(session.token_length(), Some(123_456));
}
#[tokio::test]
async fn session_list_excluding_prefixes() {
let (store, _dir) = crate::open_test_store!(crate::session::SessionStore, "session");
let direct_id = unique_key();
let excluded_prefixes: Vec<&str> = reserved_agent_id_prefixes().collect();
let prefixed_ids: Vec<String> = excluded_prefixes
.iter()
.map(|p| format!("{p}{}", unique_key()))
.collect();
let case_variant_ids: Vec<String> = ["Manager_", "Ticket_", "tIcKeT_"]
.iter()
.map(|p| format!("{p}{}", unique_key()))
.collect();
for id in std::iter::once(&direct_id)
.chain(prefixed_ids.iter())
.chain(case_variant_ids.iter())
{
store
.append_with_context(
id,
&ChatMessage::user("msg"),
"test_channel",
"test_user",
"test_ws",
"engineer",
)
.await
.unwrap();
}
let all = store.list_sessions_with_metadata().await;
let all_ids: Vec<&str> = all.iter().map(|s| s.agent_id.as_str()).collect();
assert!(
all_ids.contains(&direct_id.as_str()),
"direct session should be in full list"
);
for (prefix, id) in excluded_prefixes.iter().zip(&prefixed_ids) {
assert!(
all_ids.contains(&id.as_str()),
"{prefix} session should be in full list"
);
}
for id in &case_variant_ids {
assert!(
all_ids.contains(&id.as_str()),
"case-variant '{id}' should be in full list"
);
}
let excluded = store
.list_sessions_with_metadata_excluding(&excluded_prefixes)
.await;
let excluded_ids: Vec<&str> = excluded.iter().map(|s| s.agent_id.as_str()).collect();
assert!(
excluded_ids.contains(&direct_id.as_str()),
"direct session should survive exclusion"
);
for (prefix, id) in excluded_prefixes.iter().zip(&prefixed_ids) {
assert!(
!excluded_ids.contains(&id.as_str()),
"{prefix} session should be excluded"
);
}
for id in &case_variant_ids {
assert!(
!excluded_ids.contains(&id.as_str()),
"case-variant '{id}' should be excluded"
);
}
}
#[tokio::test]
async fn session_init_empty_message_no_append() {
crate::util::test::init_test_stores().await;
let agent_id = unique_key();
let ws = crate::workspace::test_ws_named("/_test_empty_session_init", "empty_test");
let session = session_with_init(&agent_id, "hello", &ws).await;
let len_after_real = session.history().len();
assert!(
len_after_real >= 2,
"real message should produce system prompt + user message (got {len_after_real})"
);
let session = session_with_init(&agent_id, "", &ws).await;
assert_eq!(
session.history().len(),
len_after_real,
"empty message must not append to session history",
);
}
fn tool_call_frame() -> ChatMessage {
let tc = ToolCall {
id: "t1".into(),
name: "read".into(),
arguments: serde_json::json!({}),
};
ChatMessage::assistant(
crate::providers::reasoning::assistant_replay_payload(
Some("prologue"),
std::slice::from_ref(&tc),
None,
)
.to_string(),
)
}
#[test]
fn retention_window_keeps_latest_per_side_excluding_tools() {
let frame = tool_call_frame();
let messages = vec![
ChatMessage::system("prompt"),
ChatMessage::user("u1"),
ChatMessage::assistant("a1"),
ChatMessage::user("u2"),
frame.clone(),
ChatMessage::tool_result("t1", "file contents"),
ChatMessage::assistant("a2"),
ChatMessage::user("u3"), ];
let window = select_retention_window(&messages);
let contents: Vec<&str> = window.iter().map(|m| m.content.as_str()).collect();
assert_eq!(contents, vec!["u1", "a1", "u2", "a2", "u3"]);
assert_eq!(contents.iter().filter(|c| **c == "u3").count(), 1);
}
#[test]
fn retention_window_drops_oldest_beyond_three() {
let messages: Vec<ChatMessage> = (0..5)
.flat_map(|i| {
vec![
ChatMessage::user(format!("u{i}")),
ChatMessage::assistant(format!("a{i}")),
]
})
.collect();
let window = select_retention_window(&messages);
let contents: Vec<&str> = window.iter().map(|m| m.content.as_str()).collect();
assert_eq!(contents, vec!["u2", "a2", "u3", "a3", "u4", "a4"]);
}
#[test]
fn retention_window_short_history_and_json_answers() {
let window = select_retention_window(&[ChatMessage::user("u1")]);
assert_eq!(window.len(), 1);
assert_eq!(window[0].content, "u1");
let reasoning = Reasoning {
reasoning: Some("thinking".into()),
reasoning_content: None,
reasoning_details: None,
};
let reasoning_msg = ChatMessage::assistant(
crate::providers::reasoning::assistant_replay_payload(Some(""), &[], Some(&reasoning))
.to_string(),
);
let window =
select_retention_window(&[reasoning_msg, ChatMessage::assistant("{\"result\": 42}")]);
assert_eq!(window.len(), 2);
}
#[test]
fn decode_native_history_message_gates_on_native_shape() {
assert!(
decode_native_history_message(&ChatMessage::assistant(r#"{"verdict": "bad"}"#))
.is_none()
);
assert!(
decode_native_history_message(&ChatMessage::assistant(
r#"{"reasoning_content": "think"}"#
))
.is_none()
);
let decoded =
decode_native_history_message(&ChatMessage::assistant(r#"{"content": null}"#));
assert!(matches!(
decoded,
Some(DecodedNativeHistoryMessage::Assistant { content: None, .. })
));
let decoded =
decode_native_history_message(&ChatMessage::assistant(r#"{"content": "hi"}"#));
assert!(matches!(
decoded,
Some(DecodedNativeHistoryMessage::Assistant {
content: Some(_),
..
})
));
let tc = ToolCall {
id: "t1".into(),
name: "read".into(),
arguments: serde_json::json!({}),
};
let frame = crate::providers::reasoning::assistant_replay_payload(
Some(""),
std::slice::from_ref(&tc),
None,
)
.to_string();
let decoded = decode_native_history_message(&ChatMessage::assistant(frame));
assert!(matches!(
decoded,
Some(DecodedNativeHistoryMessage::Assistant {
tool_calls: Some(_),
..
})
));
let decoded = decode_native_history_message(&ChatMessage::tool_result("t1", "r"));
assert!(matches!(
decoded,
Some(DecodedNativeHistoryMessage::ToolResult {
tool_call_id,
content
}) if tool_call_id == "t1" && content == "r"
));
}
#[tokio::test]
async fn finalize_appends_only_new_assistant_answer() {
crate::util::test::init_test_stores().await;
let k = unique_key();
store()
.batch_append(
&k,
&[
ChatMessage::system("prompt"),
ChatMessage::user("u1"),
ChatMessage::assistant("a1"),
],
)
.await
.unwrap();
let ws = crate::workspace::test_ws_named("/_test_finalize_guard", "finalize_test");
let mut session = session_with_init(&k, "", &ws).await;
assert_eq!(
session.finalize(&k).await.unwrap(),
FinalizeOutcome::NoUnpersistedTail
);
assert_eq!(store().load(&k).await.len(), 3);
session.push_assistant("a2".to_string());
assert_eq!(
session.finalize(&k).await.unwrap(),
FinalizeOutcome::Flushed
);
let msgs = store().load(&k).await;
assert_eq!(msgs.len(), 4);
assert_eq!(msgs[3].content, "a2");
}
#[tokio::test]
async fn failed_drain_gap_flushed_by_later_persist() {
crate::util::test::init_test_stores().await;
let k = unique_key();
store()
.batch_append(
&k,
&[
ChatMessage::system("prompt"),
ChatMessage::user("u1"),
ChatMessage::assistant("a1"),
],
)
.await
.unwrap();
let ws = crate::workspace::test_ws_named("/_test_gap_flush", "gap_flush");
let mut session = session_with_init(&k, "", &ws).await;
session.push_messages_unpersisted(&[ChatMessage::user("C1")]);
session
.persist_messages(
&k,
&[
ChatMessage::assistant("A"),
ChatMessage::tool_result("t1", "R1"),
],
)
.await
.unwrap();
session.push_assistant("final".into());
session.finalize(&k).await.unwrap();
let msgs = store().load(&k).await;
let contents: Vec<&str> = msgs.iter().map(|m| m.content.as_str()).collect();
assert_eq!(
contents,
[
"prompt",
"u1",
"a1",
"C1",
"A",
"{\"tool_call_id\":\"t1\",\"content\":\"R1\"}",
"final"
]
);
}
#[tokio::test]
async fn failed_drain_gap_flushed_on_aborted_turn() {
crate::util::test::init_test_stores().await;
let k = unique_key();
store()
.batch_append(
&k,
&[ChatMessage::system("prompt"), ChatMessage::user("u1")],
)
.await
.unwrap();
let ws = crate::workspace::test_ws_named("/_test_gap_abort", "gap_abort");
let mut session = session_with_init(&k, "", &ws).await;
session.push_messages_unpersisted(&[ChatMessage::user("C1")]);
session.finalize(&k).await.unwrap();
let msgs = store().load(&k).await;
let contents: Vec<&str> = msgs.iter().map(|m| m.content.as_str()).collect();
assert_eq!(contents, ["prompt", "u1", "C1"]);
}
}
#[cfg(test)]
mod transient_prefix_tests {
use super::*;
#[test]
fn forward_no_collision_with_user_facing_agent_ids() {
let reserved_prefixes: Vec<&str> = reserved_agent_id_prefixes().collect();
for prefix in &reserved_prefixes {
let bare_word = prefix.trim_end_matches('_');
let capitalized = {
let mut chars = bare_word.chars();
match chars.next() {
Some(first) => first.to_ascii_uppercase().to_string() + chars.as_str(),
None => bare_word.to_string(),
}
};
let user_names: [String; 5] = [
bare_word.to_string(),
format!("{bare_word}_bob"),
format!("{bare_word}x"),
capitalized.clone(),
format!("{capitalized}Bob"),
];
for user_name in &user_names {
let key = direct_agent_id(user_name, "analyst", "test-ws");
assert!(
!starts_with_ignore_ascii_case(&key, bare_word),
"DIRECT AGENT ID COLLISION: bare-word='{bare_word}' \
(case-insensitive) matches id='{key}' (user='{user_name}'). \
Fix: safe_user_segment must escape any user name starting \
(case-insensitively) with '{bare_word}'.",
);
}
let user_word = format!("user_{bare_word}");
let escaped_key = direct_agent_id(&user_word, "analyst", "test-ws");
assert!(
escaped_key.starts_with("user_user_"),
"real user '{user_word}' should be double-escaped to start \
with 'user_user_', got '{escaped_key}'",
);
assert_ne!(
escaped_key,
direct_agent_id(bare_word, "analyst", "test-ws"),
"escape of '{bare_word}' collides with escape of '{user_word}'",
);
}
for prefix in TRANSIENT_AGENT_ID_PREFIXES {
let manager_key = manager_agent_id("test-ws");
assert!(
!manager_key.starts_with(prefix),
"MANAGER AGENT ID COLLISION: prefix='{prefix}' \
matches id='{manager_key}'. \
Fix: remove '{prefix}' from TRANSIENT_AGENT_ID_PREFIXES \
or change the manager_agent_id pattern.",
);
}
assert_eq!(
direct_agent_id("alice", "analyst", "test-ws"),
"alice_test-ws_analyst",
);
for prefix in &reserved_prefixes {
let bare_word = prefix.trim_end_matches('_');
let key = direct_agent_id(bare_word, "analyst", "test-ws");
assert!(
key.starts_with("user_"),
"bare reserved word '{bare_word}' should be escaped to start \
with 'user_', got '{key}'",
);
}
}
fn assert_transient_key(key: &str, expected_prefix: &str, builder_expr: &str) {
assert!(
key.starts_with(expected_prefix),
"{builder_expr} = '{key}' does not start with '{expected_prefix}'.\n\
Fix: update {builder_expr} to produce IDs starting with '{expected_prefix}'.",
);
assert!(
TRANSIENT_AGENT_ID_PREFIXES.contains(&expected_prefix),
"TRANSIENT_AGENT_ID_PREFIXES is missing '{expected_prefix}' — \
{builder_expr} sessions will never be cleaned up.\n\
Fix: add \"{expected_prefix}\" to TRANSIENT_AGENT_ID_PREFIXES.",
);
}
#[test]
fn reverse_transient_builders_use_registered_prefixes() {
assert_transient_key(
&ticket_agent_id("abc123", "analyst"),
"ticket_",
"ticket_agent_id('abc123', 'analyst')",
);
assert_transient_key(
&analyze_agent_id("ws", "coder"),
"analyze_",
"analyze_agent_id('ws', 'coder')",
);
assert_transient_key(
&research_agent_id("ws", "decomposer"),
"research_",
"research_agent_id('ws', 'decomposer')",
);
assert_transient_key(
&crate::research_cleanup::cleanup_agent_id("run_abc123"),
"cleanup_",
"cleanup_agent_id('run_abc123')",
);
assert_transient_key(
&maintainer_agent_id("ws"),
"maintainer_",
"maintainer_agent_id('ws')",
);
assert_transient_key(
&discovery_agent_id("ws", "analyst"),
"discovery_",
"discovery_agent_id('ws', 'analyst')",
);
}
#[test]
fn resolve_agent_id_manager_dispatch() {
let key = resolve_agent_id("alice", "manager", "my-workspace");
assert_eq!(key, "manager_my-workspace");
}
#[test]
fn resolve_agent_id_non_manager_dispatch() {
let key = resolve_agent_id("bob", "engineer", "my-workspace");
assert_eq!(key, "bob_my-workspace_engineer");
}
#[test]
fn resolve_agent_id_lowercase_manager() {
let key = resolve_agent_id("carol", "Manager", "ws");
assert_ne!(key, "manager_ws", "capital-M 'Manager' should NOT match");
assert_eq!(key, "carol_ws_Manager");
}
#[test]
fn direct_agent_id_personal_workspace_dedup() {
assert_eq!(
direct_agent_id("alice", "assistant", "personal:alice"),
"alice_personal:assistant",
);
assert_eq!(
direct_agent_id("bob", "assistant", "personal:admin"),
"bob_personal:admin_assistant",
);
assert_eq!(
direct_agent_id("alice", "analyst", "personal_work"),
"alice_personal_work_analyst",
);
assert_eq!(
direct_agent_id("manager", "engineer", "personal:manager"),
"user_manager_personal:engineer",
);
assert_eq!(
direct_agent_id("user_ticket", "analyst", "personal:user_ticket"),
"user_user_ticket_personal:analyst",
);
}
}
#[derive(Debug)]
pub(crate) enum DecodedNativeHistoryMessage {
Assistant {
content: Option<String>,
tool_calls: Option<Vec<ToolCall>>,
reasoning: Option<Reasoning>,
},
ToolResult {
tool_call_id: String,
content: String,
},
}
pub(crate) fn decode_native_history_message(
message: &ChatMessage,
) -> Option<DecodedNativeHistoryMessage> {
let parsed = serde_json::from_str::<serde_json::Value>(&message.content).ok();
if message.role == ChatRole::Assistant
&& let Some(value) = parsed.as_ref()
&& (value.get("content").is_some() || value.get("tool_calls").is_some())
{
let content = value
.get("content")
.and_then(serde_json::Value::as_str)
.map(ToString::to_string);
let (r, rc, rd) =
crate::providers::reasoning::json_lossless_assistant_reasoning_fields(value);
let reasoning = Reasoning::from_optional_parts(r, rc, rd);
let tool_calls = value
.get("tool_calls")
.and_then(|v| serde_json::from_value::<Vec<ToolCall>>(v.clone()).ok())
.map(|mut parsed_calls| {
for call in &mut parsed_calls {
if let Some(s) = call.arguments.as_str()
&& let Ok(v) = serde_json::from_str::<serde_json::Value>(s)
{
call.arguments = v;
}
}
parsed_calls
});
return Some(DecodedNativeHistoryMessage::Assistant {
content,
tool_calls,
reasoning,
});
}
if message.role == ChatRole::Tool
&& let Ok(payload) = serde_json::from_str::<ToolResultPayload>(&message.content)
{
return Some(DecodedNativeHistoryMessage::ToolResult {
tool_call_id: payload.tool_call_id,
content: payload.content,
});
}
None
}