pub mod dead_session;
pub mod manager;
pub(crate) use manager::FinalizeOutcome;
pub(crate) use manager::RewriteOutcome;
pub use manager::Session;
use crate::turso::{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 RETENTION_PER_SIDE: usize = 3;
fn is_tool_call_frame(msg: &ChatMessage) -> bool {
matches!(
decode_native_history_message(msg),
Some(DecodedNativeHistoryMessage::Assistant {
tool_calls: Some(_),
..
})
)
}
#[must_use]
pub(crate) 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,
db_name = "sessions",
schema = SCHEMA,
post_open = run_migrations,
expect = "SESSIONS not initialized",
}
const SCHEMA: &str = "CREATE TABLE IF NOT EXISTS sessions (
id INTEGER PRIMARY KEY AUTOINCREMENT,
agent_id TEXT NOT NULL,
role TEXT NOT NULL,
content TEXT NOT NULL,
created_at TEXT NOT NULL
);
CREATE INDEX IF NOT EXISTS idx_sessions_agent_id ON sessions(agent_id, id);
CREATE TABLE IF NOT EXISTS session_metadata (
agent_id TEXT PRIMARY KEY,
last_activity TEXT NOT NULL,
channel TEXT,
user_name TEXT,
workspace_name TEXT,
role TEXT,
active_models TEXT,
token_length INTEGER,
message_count INTEGER NOT NULL DEFAULT 0
);
-- ── Durability/resume substrate (see src/jobs.rs) ─────────────────────
-- jobs must be declared BEFORE agents (lazy FK resolution).
CREATE TABLE IF NOT EXISTS jobs (
id TEXT PRIMARY KEY,
kind TEXT NOT NULL,
status TEXT NOT NULL DEFAULT 'launched',
task TEXT NOT NULL DEFAULT '',
workspace_name TEXT NOT NULL,
user_name TEXT NOT NULL DEFAULT '',
channel TEXT NOT NULL DEFAULT '',
role TEXT NOT NULL,
retry_count INTEGER NOT NULL DEFAULT 0,
created_at TEXT NOT NULL,
updated_at TEXT NOT NULL
);
CREATE INDEX IF NOT EXISTS idx_jobs_kind_status ON jobs(kind, status);
CREATE INDEX IF NOT EXISTS idx_jobs_updated_at ON jobs(updated_at);
CREATE TABLE IF NOT EXISTS agents (
job_id TEXT REFERENCES jobs(id) ON DELETE CASCADE,
agent_id TEXT NOT NULL,
kind TEXT NOT NULL,
idx INTEGER,
status TEXT NOT NULL DEFAULT 'launched',
outcome TEXT,
task TEXT NOT NULL,
PRIMARY KEY (job_id, agent_id)
);
CREATE UNIQUE INDEX IF NOT EXISTS idx_agents_anchor ON agents(agent_id) WHERE job_id IS NULL;
CREATE TABLE IF NOT EXISTS pending_jobs (
id TEXT PRIMARY KEY,
target_agent_id TEXT NOT NULL,
envelope TEXT NOT NULL,
created_at TEXT NOT NULL
);
CREATE INDEX IF NOT EXISTS idx_pending_jobs_agent_created ON pending_jobs(target_agent_id, created_at);
CREATE TABLE IF NOT EXISTS ticket_stage_jobs (
id TEXT PRIMARY KEY REFERENCES jobs(id) ON DELETE CASCADE,
ticket_id TEXT NOT NULL,
stage TEXT NOT NULL,
phase TEXT NOT NULL,
round INTEGER NOT NULL
);
CREATE TABLE IF NOT EXISTS research_jobs (
id TEXT PRIMARY KEY REFERENCES jobs(id) ON DELETE CASCADE,
state TEXT NOT NULL
);";
impl SessionStore {
async fn run_migrations(&self) -> anyhow::Result<()> {
crate::turso::run_pending_migrations(
&self.conn,
"sessions",
crate::migrations::SESSION_MIGRATIONS,
)
.await
}
}
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",
}
}
pub(crate) 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>,
}
#[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>,
) -> Result<SessionMetadata> {
let last_activity = turso::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()),
})
}
async fn insert_messages_in_transaction(
tx: &TxGuard<'_>,
agent_id: &str,
messages: &[ChatMessage],
context: Option<(&str, &str, &str, &str)>,
replace: bool,
) -> Result<()> {
let now = turso::now();
for msg in messages {
tx.execute(
"INSERT INTO sessions (agent_id, role, content, created_at) VALUES (?1, ?2, ?3, ?4)",
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")?;
let count_clause = if replace {
"message_count = excluded.message_count"
} else {
"message_count = message_count + excluded.message_count"
};
match context {
Some((channel, user_name, workspace_name, role)) => {
tx.execute(
&format!(
"INSERT INTO session_metadata (agent_id, last_activity, message_count, \
channel, user_name, workspace_name, role) \
VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7) \
ON CONFLICT(agent_id) DO UPDATE SET \
last_activity = excluded.last_activity, \
channel = excluded.channel, \
user_name = excluded.user_name, \
workspace_name = excluded.workspace_name, \
role = excluded.role, \
{count_clause}"
),
params![
agent_id,
now,
count,
channel,
user_name,
workspace_name,
role,
],
)
.await?;
}
None => {
tx.execute(
&format!(
"INSERT INTO session_metadata (agent_id, last_activity, message_count) \
VALUES (?1, ?2, ?3) \
ON CONFLICT(agent_id) DO UPDATE SET \
last_activity = excluded.last_activity, \
{count_clause}"
),
params![agent_id, now, count],
)
.await?;
}
}
Ok(())
}
async fn query_map_collect<T, E>(
conn: &turso::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()
}
async fn list_sessions_where(
conn: &turso::Connection,
where_clause: &str,
params: impl IntoParams + Send + 'static,
warn_context: &str,
) -> Vec<SessionMetadata> {
query_map_collect(
conn,
&format!(
"SELECT {SESSION_LIST_COLUMNS} \
FROM session_metadata sm \
{where_clause} \
ORDER BY sm.last_activity DESC",
),
params,
|row| {
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)?,
)
},
warn_context,
None,
)
.await
}
impl SessionStore {
pub(crate) async fn load(&self, agent_id: &str) -> Vec<ChatMessage> {
query_map_collect(
&self.conn,
&format!(
"SELECT {SESSION_MESSAGE_COLUMNS} FROM sessions WHERE agent_id = ?1 ORDER BY id ASC"
),
params![agent_id],
|row| {
Ok::<_, anyhow::Error>(ChatMessage {
role: row
.get::<String>(COL_SM_ROLE)?
.parse::<ChatRole>()
.map_err(|e| anyhow!(e))?,
content: row.get(COL_SM_CONTENT)?,
})
},
"load session",
Some(agent_id),
)
.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(())
}
pub(crate) async fn batch_append(
&self,
agent_id: &str,
messages: &[ChatMessage],
) -> Result<()> {
self.append_messages(agent_id, messages, false, None).await
}
pub(crate) 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
}
pub(crate) 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
}
pub(crate) 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(())
}
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 list_sessions_with_metadata(&self) -> Vec<SessionMetadata> {
list_sessions_where(&self.conn, "", (), "list sessions").await
}
pub(crate) 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<turso::Value> = exclude_prefixes
.iter()
.map(|p| turso::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_role(&self, agent_id: &str) -> Option<ChatRole> {
let rows = self
.conn
.query(
"SELECT role 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(0).ok()?;
role_str.parse::<ChatRole>().ok()
})
}
pub(crate) 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?,
})
})
}
pub(crate) async fn set_active_models(
&self,
agent_id: &str,
snapshot: Option<&str>,
) -> Result<()> {
let now = turso::now();
self.conn
.execute(
"INSERT INTO session_metadata (agent_id, last_activity, active_models) \
VALUES (?1, ?2, ?3) \
ON CONFLICT(agent_id) DO UPDATE SET active_models = excluded.active_models",
params![agent_id, now, snapshot],
)
.await?;
Ok(())
}
pub(crate) async fn get_active_models(&self, agent_id: &str) -> Option<String> {
match self
.conn
.query_optional(
"SELECT active_models FROM session_metadata WHERE agent_id = ?1",
params![agent_id],
|row| row.get::<Option<String>>(0),
)
.await
{
Ok(Some(snapshot)) => snapshot,
Ok(None) => None,
Err(e) => {
tracing::warn!(agent_id = %agent_id, error = %e, "Failed to read active-models snapshot");
None
}
}
}
pub(crate) async fn set_token_length(
&self,
agent_id: &str,
token_length: Option<u64>,
) -> Result<()> {
let now = turso::now();
let bound: Option<i64> = token_length.map(i64::try_from).transpose()?;
self.conn
.execute(
"INSERT INTO session_metadata (agent_id, last_activity, token_length) \
VALUES (?1, ?2, ?3) \
ON CONFLICT(agent_id) DO UPDATE SET token_length = excluded.token_length",
params![agent_id, now, bound],
)
.await?;
Ok(())
}
pub(crate) async fn get_token_length(&self, agent_id: &str) -> Option<u64> {
match self
.conn
.query_optional(
"SELECT token_length FROM session_metadata WHERE agent_id = ?1",
params![agent_id],
|row| row.get::<Option<i64>>(0),
)
.await
{
Ok(Some(Some(tokens))) => u64::try_from(tokens).ok(),
Ok(_) => None,
Err(e) => {
tracing::warn!(agent_id = %agent_id, error = %e, "Failed to read session token length");
None
}
}
}
}
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 build_params = {
let mut p = vec![Value::Text(cutoff.to_string())];
p.extend(
TRANSIENT_AGENT_ID_PREFIXES
.iter()
.map(|prefix| Value::Text(format!("{prefix}%"))),
);
p
};
tx.execute(
&format!(
"DELETE FROM sessions WHERE agent_id IN ( \
SELECT agent_id FROM session_metadata \
WHERE last_activity < ? AND {prefix_patterns} \
AND agent_id NOT IN (SELECT agent_id FROM agents))"
),
build_params.clone(),
)
.await?;
let deleted = tx
.execute(
&format!(
"DELETE FROM session_metadata WHERE last_activity < ? AND {prefix_patterns} \
AND agent_id NOT IN (SELECT agent_id FROM agents)"
),
build_params.clone(),
)
.await?;
tx.commit().await?;
Ok(deleted)
}
pub(crate) fn reserved_agent_id_prefixes() -> impl Iterator<Item = &'static str> {
std::iter::once("manager_").chain(TRANSIENT_AGENT_ID_PREFIXES.iter().copied())
}
pub(crate) 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'",
);
"admin"
} 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]
pub fn direct_agent_id(user_name: &str, role: &str, ws_name: &str) -> String {
let user_name = normalize_user_name(user_name, "direct_agent_id");
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 fn manager_agent_id(ws_name: &str) -> String {
format!("manager_{ws_name}")
}
#[must_use]
pub 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) -> String {
Session::delete(&resolve_agent_id(user_name, role, ws_name)).await
}
#[must_use]
pub(crate) fn maintainer_agent_id(ws_name: &str) -> String {
format!("maintainer_{}_{}", ws_name, crate::generate_suffix())
}
#[must_use]
pub(crate) fn analyze_agent_id(ws_name: &str, role: &str) -> String {
format!("analyze_{}_{}_{}", ws_name, crate::generate_suffix(), role)
}
#[must_use]
pub(crate) fn research_agent_id(ws_name: &str, label: &str) -> String {
format!(
"research_{}_{}_{}",
ws_name,
crate::generate_suffix(),
label
)
}
#[must_use]
pub(crate) fn discovery_agent_id(ws_name: &str, role: &str) -> String {
format!(
"discovery_{}_{}_{}",
ws_name,
crate::generate_suffix(),
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))
}
#[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 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_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_role() {
crate::util::test::init_test_stores().await;
let agent_id = unique_key();
assert!(store().get_last_message_role(&agent_id).await.is_none());
store()
.batch_append(&agent_id, &[ChatMessage::user("hello")])
.await
.unwrap();
assert_eq!(
store().get_last_message_role(&agent_id).await,
Some(ChatRole::User)
);
store()
.batch_append(&agent_id, &[ChatMessage::assistant("world")])
.await
.unwrap();
assert_eq!(
store().get_last_message_role(&agent_id).await,
Some(ChatRole::Assistant)
);
}
#[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::turso::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 mut session = Session::default();
session
.init(
&k,
"",
&ws,
&crate::Role::Assistant,
None,
"gui",
"tester",
None,
)
.await
.unwrap();
assert_eq!(session.token_length(), Some(123_456));
}
#[tokio::test]
async fn session_list_excluding_prefixes() {
crate::util::test::init_test_stores().await;
let direct_id = unique_key();
let excluded_prefixes: Vec<&str> = std::iter::once("manager_")
.chain(TRANSIENT_AGENT_ID_PREFIXES.iter().copied())
.collect();
let prefixed_ids: Vec<String> = excluded_prefixes
.iter()
.map(|p| format!("{p}{}", unique_key()))
.collect();
for id in std::iter::once(&direct_id).chain(prefixed_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"
);
}
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"
);
}
}
#[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 role = crate::Role::Assistant;
let mut session = Session::default();
session
.init(&agent_id, "hello", &ws, &role, None, "gui", "tester", None)
.await
.unwrap();
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 mut session = Session::default();
session
.init(&agent_id, "", &ws, &role, None, "gui", "tester", None)
.await
.unwrap();
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_roundtrip::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_roundtrip::assistant_replay_payload(
Some(""),
&[],
Some(&reasoning),
)
.to_string(),
);
let window =
select_retention_window(&[reasoning_msg, ChatMessage::assistant("{\"result\": 42}")]);
assert_eq!(window.len(), 2);
}
#[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::default();
session
.init(
&k,
"",
&ws,
&crate::Role::Assistant,
None,
"gui",
"tester",
None,
)
.await
.unwrap();
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::default();
session
.init(
&k,
"",
&ws,
&crate::Role::Assistant,
None,
"gui",
"tester",
None,
)
.await
.unwrap();
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::default();
session
.init(
&k,
"",
&ws,
&crate::Role::Assistant,
None,
"gui",
"tester",
None,
)
.await
.unwrap();
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 = std::iter::once("manager_")
.chain(TRANSIENT_AGENT_ID_PREFIXES.iter().copied())
.collect::<Vec<_>>();
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");
}
}
#[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()
{
let content = value
.get("content")
.and_then(serde_json::Value::as_str)
.map(ToString::to_string);
let (r, rc, rd) =
crate::providers::reasoning_roundtrip::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
}