use std::path::{Path, PathBuf};
use async_openai::types::chat::{
ChatCompletionRequestAssistantMessage, ChatCompletionRequestMessage,
ChatCompletionRequestToolMessage, ChatCompletionRequestUserMessage,
ChatCompletionRequestUserMessageContent,
};
use rusqlite::{params, types::ToSqlOutput, Connection, Result as SqliteResult};
use serde::{Deserialize, Serialize};
use crate::datetime::current_timestamp;
use crate::error::Result;
const ROBIT_DIR: &str = ".robit";
const MEMORY_DIR: &str = "memory";
const DB_FILE: &str = "robit.db";
pub fn resolve_db_path(working_dir: &Path, global_storage: bool) -> Result<PathBuf> {
if global_storage {
let home = dirs::home_dir().ok_or_else(|| {
crate::error::AgentError::InternalError("Cannot determine home directory".to_string())
})?;
Ok(home.join(ROBIT_DIR).join(MEMORY_DIR).join(DB_FILE))
} else {
Ok(working_dir.join(ROBIT_DIR).join(MEMORY_DIR).join(DB_FILE))
}
}
#[derive(Debug, Clone, Serialize)]
pub struct SessionInfo {
pub id: String,
pub chat_id: Option<String>,
pub title: String,
pub model: String,
pub source: String,
pub status: String, pub created_at: String,
pub updated_at: String,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct MessageData {
pub id: i64,
pub role: String,
pub content: String,
pub tool_name: Option<String>,
pub tool_call_id: Option<String>,
pub tool_info: Option<serde_json::Value>,
pub created_at: String,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct ToolCallInfoData {
pub tool_call_id: String,
pub name: String,
pub arguments: String,
pub status: String,
pub output: Option<String>,
pub requires_confirm: bool,
}
const CURRENT_SCHEMA_VERSION: i32 = 4;
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub enum MemoryType {
Fact,
Preference,
Note,
Task,
Custom(String),
}
impl MemoryType {
pub fn as_str(&self) -> &str {
match self {
MemoryType::Fact => "fact",
MemoryType::Preference => "preference",
MemoryType::Note => "note",
MemoryType::Task => "task",
MemoryType::Custom(s) => s,
}
}
pub fn from_str(s: &str) -> Self {
match s.to_lowercase().as_str() {
"fact" => MemoryType::Fact,
"preference" => MemoryType::Preference,
"note" => MemoryType::Note,
"task" => MemoryType::Task,
_ => MemoryType::Custom(s.to_string()),
}
}
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct Memory {
pub id: String,
pub session_id: Option<String>,
pub chat_id: Option<String>,
pub memory_type: MemoryType,
pub title: String,
pub content: String,
pub tags: Vec<String>,
pub is_active: bool,
pub created_at: String,
pub updated_at: String,
}
#[derive(Debug, Clone, Default)]
pub struct MemoryFilter {
pub memory_type: Option<MemoryType>,
pub tags: Option<Vec<String>>,
pub session_id: Option<String>,
pub chat_id: Option<String>,
pub since: Option<String>,
pub only_active: bool,
}
impl Memory {
pub fn new(
title: String,
content: String,
memory_type: MemoryType,
tags: Vec<String>,
) -> Self {
let now = current_timestamp();
Memory {
id: uuid::Uuid::new_v4().to_string(),
session_id: None,
chat_id: None,
memory_type,
title,
content,
tags,
is_active: true,
created_at: now.clone(),
updated_at: now,
}
}
pub fn with_session_id(mut self, session_id: String) -> Self {
self.session_id = Some(session_id);
self
}
pub fn with_chat_id(mut self, chat_id: String) -> Self {
self.chat_id = Some(chat_id);
self
}
}
pub fn init_db(conn: &Connection) -> SqliteResult<()> {
ensure_meta_table(conn)?;
let version = read_schema_version(conn)?;
if version == 0 {
create_all_tables(conn)?;
write_schema_version(conn, CURRENT_SCHEMA_VERSION)?;
tracing::info!(
"Database initialized at schema v{}",
CURRENT_SCHEMA_VERSION
);
return Ok(());
}
migrate(conn, version, CURRENT_SCHEMA_VERSION)?;
Ok(())
}
fn read_schema_version(conn: &Connection) -> SqliteResult<i32> {
match conn.query_row(
"SELECT value FROM _schema_meta WHERE key = 'version'",
[],
|row| row.get::<_, String>(0),
) {
Ok(v) => v.parse().map_err(|_| {
rusqlite::Error::InvalidParameterName(format!("Invalid schema version: {}", v))
}),
Err(rusqlite::Error::QueryReturnedNoRows) => {
if sessions_table_exists(conn)? {
Ok(1) } else {
Ok(0) }
}
Err(e) => Err(e),
}
}
fn sessions_table_exists(conn: &Connection) -> SqliteResult<bool> {
let count: i64 = conn.query_row(
"SELECT COUNT(*) FROM sqlite_master WHERE type = 'table' AND name = 'sessions'",
[],
|row| row.get(0),
)?;
Ok(count > 0)
}
fn create_all_tables(conn: &Connection) -> SqliteResult<()> {
conn.execute_batch(
"CREATE TABLE IF NOT EXISTS sessions (
id TEXT PRIMARY KEY,
chat_id TEXT,
title TEXT NOT NULL,
model TEXT NOT NULL,
source TEXT NOT NULL DEFAULT 'gui',
created_at TEXT NOT NULL,
updated_at TEXT NOT NULL,
is_active INTEGER DEFAULT 1
);
CREATE TABLE IF NOT EXISTS messages (
id INTEGER PRIMARY KEY AUTOINCREMENT,
session_id TEXT NOT NULL REFERENCES sessions(id),
role TEXT NOT NULL,
content TEXT NOT NULL,
tool_name TEXT,
tool_call_id TEXT,
tool_info TEXT,
tokens INTEGER,
created_at TEXT NOT NULL
);
CREATE INDEX IF NOT EXISTS idx_messages_session
ON messages(session_id);
CREATE INDEX IF NOT EXISTS idx_messages_created
ON messages(session_id, created_at);
CREATE UNIQUE INDEX IF NOT EXISTS idx_sessions_chat_id
ON sessions(chat_id) WHERE chat_id IS NOT NULL;
CREATE TABLE IF NOT EXISTS memories (
id TEXT PRIMARY KEY,
session_id TEXT,
chat_id TEXT,
memory_type TEXT NOT NULL,
title TEXT NOT NULL,
content TEXT NOT NULL,
tags TEXT,
is_active INTEGER DEFAULT 1,
created_at TEXT NOT NULL,
updated_at TEXT NOT NULL,
FOREIGN KEY (session_id) REFERENCES sessions(id) ON DELETE SET NULL
);
CREATE INDEX IF NOT EXISTS idx_memories_type ON memories(memory_type);
CREATE INDEX IF NOT EXISTS idx_memories_created ON memories(created_at DESC);
CREATE INDEX IF NOT EXISTS idx_memories_session ON memories(session_id);
CREATE INDEX IF NOT EXISTS idx_memories_chat ON memories(chat_id);
CREATE INDEX IF NOT EXISTS idx_memories_active ON memories(is_active) WHERE is_active = 1;
-- FTS5 virtual table for full-text search on messages
CREATE VIRTUAL TABLE IF NOT EXISTS messages_fts USING fts5(
content,
content='messages',
content_rowid='id',
tokenize='unicode61'
);
-- Triggers to keep FTS index in sync with messages table
CREATE TRIGGER IF NOT EXISTS messages_ai AFTER INSERT ON messages BEGIN
INSERT INTO messages_fts(rowid, content) VALUES (new.id, new.content);
END;
CREATE TRIGGER IF NOT EXISTS messages_ad AFTER DELETE ON messages BEGIN
INSERT INTO messages_fts(messages_fts, rowid, content) VALUES ('delete', old.id, old.content);
END;
CREATE TRIGGER IF NOT EXISTS messages_au AFTER UPDATE OF content ON messages BEGIN
INSERT INTO messages_fts(messages_fts, rowid, content) VALUES ('delete', old.id, old.content);
INSERT INTO messages_fts(rowid, content) VALUES (new.id, new.content);
END;",
)?;
Ok(())
}
fn migrate(conn: &Connection, from: i32, to: i32) -> SqliteResult<()> {
let mut current = from;
while current < to {
tracing::info!("Migrating database: v{} → v{}", current, current + 1);
match current {
1 => migrate_v1_to_v2(conn)?,
2 => migrate_v2_to_v3(conn)?,
3 => migrate_v3_to_v4(conn)?,
other => {
return Err(rusqlite::Error::InvalidParameterName(format!(
"Unknown schema version: {}",
other
)))
}
}
current += 1;
write_schema_version(conn, current)?;
tracing::info!("Database migrated to v{}", current);
}
Ok(())
}
fn migrate_v1_to_v2(conn: &Connection) -> SqliteResult<()> {
let _ = conn.execute("ALTER TABLE sessions ADD COLUMN chat_id TEXT", []);
let _ = conn.execute(
"ALTER TABLE sessions ADD COLUMN source TEXT NOT NULL DEFAULT 'gui'",
[],
);
let _ = conn.execute("ALTER TABLE messages ADD COLUMN tool_info TEXT", []);
conn.execute_batch(
"CREATE UNIQUE INDEX IF NOT EXISTS idx_sessions_chat_id
ON sessions(chat_id) WHERE chat_id IS NOT NULL;",
)?;
Ok(())
}
fn migrate_v2_to_v3(conn: &Connection) -> SqliteResult<()> {
conn.execute_batch(
"CREATE TABLE IF NOT EXISTS memories (
id TEXT PRIMARY KEY,
session_id TEXT,
chat_id TEXT,
memory_type TEXT NOT NULL,
title TEXT NOT NULL,
content TEXT NOT NULL,
tags TEXT,
is_active INTEGER DEFAULT 1,
created_at TEXT NOT NULL,
updated_at TEXT NOT NULL,
FOREIGN KEY (session_id) REFERENCES sessions(id) ON DELETE SET NULL
);
CREATE INDEX IF NOT EXISTS idx_memories_type ON memories(memory_type);
CREATE INDEX IF NOT EXISTS idx_memories_created ON memories(created_at DESC);
CREATE INDEX IF NOT EXISTS idx_memories_session ON memories(session_id);
CREATE INDEX IF NOT EXISTS idx_memories_chat ON memories(chat_id);
CREATE INDEX IF NOT EXISTS idx_memories_active ON memories(is_active) WHERE is_active = 1;",
)?;
Ok(())
}
fn migrate_v3_to_v4(conn: &Connection) -> SqliteResult<()> {
conn.execute_batch(
"CREATE VIRTUAL TABLE IF NOT EXISTS messages_fts USING fts5(
content,
content='messages',
content_rowid='id',
tokenize='unicode61'
);
CREATE TRIGGER IF NOT EXISTS messages_ai AFTER INSERT ON messages BEGIN
INSERT INTO messages_fts(rowid, content) VALUES (new.id, new.content);
END;
CREATE TRIGGER IF NOT EXISTS messages_ad AFTER DELETE ON messages BEGIN
INSERT INTO messages_fts(messages_fts, rowid, content) VALUES ('delete', old.id, old.content);
END;
CREATE TRIGGER IF NOT EXISTS messages_au AFTER UPDATE OF content ON messages BEGIN
INSERT INTO messages_fts(messages_fts, rowid, content) VALUES ('delete', old.id, old.content);
INSERT INTO messages_fts(rowid, content) VALUES (new.id, new.content);
END;
-- Backfill existing messages into the FTS index
INSERT INTO messages_fts(rowid, content)
SELECT id, content FROM messages;",
)?;
Ok(())
}
fn ensure_meta_table(conn: &Connection) -> SqliteResult<()> {
conn.execute_batch(
"CREATE TABLE IF NOT EXISTS _schema_meta (
key TEXT PRIMARY KEY,
value TEXT NOT NULL
)",
)
}
fn write_schema_version(conn: &Connection, version: i32) -> SqliteResult<()> {
conn.execute(
"INSERT OR REPLACE INTO _schema_meta (key, value) VALUES ('version', ?1)",
rusqlite::params![version.to_string()],
)?;
Ok(())
}
pub fn insert_session(
conn: &Connection,
id: &str,
chat_id: Option<&str>,
title: &str,
model: &str,
source: &str,
) -> SqliteResult<()> {
let now = current_timestamp();
conn.execute(
"INSERT INTO sessions (id, chat_id, title, model, source, created_at, updated_at) \
VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7)",
params![id, chat_id, title, model, source, now, now],
)?;
Ok(())
}
pub fn list_sessions(
conn: &Connection,
source_filter: Option<&str>,
) -> SqliteResult<Vec<SessionInfo>> {
let sql = if source_filter.is_some() {
"SELECT id, chat_id, title, model, source, created_at, updated_at \
FROM sessions WHERE is_active = 1 AND source = ?1 ORDER BY updated_at DESC"
} else {
"SELECT id, chat_id, title, model, source, created_at, updated_at \
FROM sessions WHERE is_active = 1 ORDER BY updated_at DESC"
};
let mut stmt = conn.prepare(sql)?;
let rows = if let Some(source) = source_filter {
stmt.query_map(params![source], map_session_row)?
} else {
stmt.query_map([], map_session_row)?
};
rows.collect()
}
pub fn find_session_by_chat_id(
conn: &Connection,
chat_id: &str,
) -> SqliteResult<Option<SessionInfo>> {
let mut stmt = conn.prepare(
"SELECT id, chat_id, title, model, source, created_at, updated_at \
FROM sessions WHERE chat_id = ?1 AND is_active = 1",
)?;
let mut rows = stmt.query_map(params![chat_id], map_session_row)?;
match rows.next() {
Some(Ok(session)) => Ok(Some(session)),
_ => Ok(None),
}
}
pub fn list_all_sessions_by_chat_id(
conn: &Connection,
chat_id: &str,
) -> SqliteResult<Vec<SessionInfo>> {
let mut stmt = conn.prepare(
"SELECT id, chat_id, title, model, source, created_at, updated_at \
FROM sessions WHERE chat_id = ?1 ORDER BY updated_at DESC",
)?;
let rows = stmt.query_map(params![chat_id], map_session_row)?;
rows.collect()
}
pub fn activate_session(
conn: &Connection,
session_id: &str,
chat_id: &str,
) -> SqliteResult<()> {
let now = current_timestamp();
conn.execute(
"UPDATE sessions SET is_active = 0 WHERE chat_id = ?1",
params![chat_id],
)?;
conn.execute(
"UPDATE sessions SET is_active = 1, updated_at = ?1 WHERE id = ?2",
params![now, session_id],
)?;
Ok(())
}
pub fn get_session(conn: &Connection, id: &str) -> SqliteResult<Option<SessionInfo>> {
let mut stmt = conn.prepare(
"SELECT id, chat_id, title, model, source, created_at, updated_at \
FROM sessions WHERE id = ?1 AND is_active = 1",
)?;
let mut rows = stmt.query_map(params![id], map_session_row)?;
match rows.next() {
Some(Ok(session)) => Ok(Some(session)),
_ => Ok(None),
}
}
fn map_session_row(row: &rusqlite::Row<'_>) -> SqliteResult<SessionInfo> {
Ok(SessionInfo {
id: row.get(0)?,
chat_id: row.get(1)?,
title: row.get(2)?,
model: row.get(3)?,
source: row.get(4)?,
status: "idle".to_string(),
created_at: row.get(5)?,
updated_at: row.get(6)?,
})
}
pub fn update_session_title(conn: &Connection, id: &str, title: &str) -> SqliteResult<()> {
let now = current_timestamp();
conn.execute(
"UPDATE sessions SET title = ?1, updated_at = ?2 WHERE id = ?3",
params![title, now, id],
)?;
Ok(())
}
pub fn touch_session(conn: &Connection, id: &str) -> SqliteResult<()> {
let now = current_timestamp();
conn.execute(
"UPDATE sessions SET updated_at = ?1 WHERE id = ?2",
params![now, id],
)?;
Ok(())
}
pub fn delete_session(conn: &Connection, id: &str) -> SqliteResult<()> {
conn.execute(
"UPDATE sessions SET is_active = 0 WHERE id = ?1",
params![id],
)?;
Ok(())
}
pub fn insert_message(
conn: &Connection,
session_id: &str,
role: &str,
content: &str,
tool_name: Option<&str>,
tool_call_id: Option<&str>,
tool_info: Option<&str>,
) -> SqliteResult<i64> {
let now = current_timestamp();
conn.execute(
"INSERT INTO messages (session_id, role, content, tool_name, tool_call_id, tool_info, created_at) VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7)",
params![session_id, role, content, tool_name, tool_call_id, tool_info, now],
)?;
Ok(conn.last_insert_rowid())
}
pub fn get_messages(conn: &Connection, session_id: &str) -> SqliteResult<Vec<MessageData>> {
let mut stmt = conn.prepare(
"SELECT id, role, content, tool_name, tool_call_id, tool_info, created_at FROM messages WHERE session_id = ?1 ORDER BY id ASC"
)?;
let rows = stmt.query_map(params![session_id], |row| {
let tool_info_str: Option<String> = row.get(5)?;
let tool_info = tool_info_str.and_then(|s| serde_json::from_str(&s).ok());
Ok(MessageData {
id: row.get(0)?,
role: row.get(1)?,
content: row.get(2)?,
tool_name: row.get(3)?,
tool_call_id: row.get(4)?,
tool_info,
created_at: row.get(6)?,
})
})?;
rows.collect()
}
pub fn update_tool_message(
conn: &Connection,
session_id: &str,
tool_call_id: &str,
tool_info: &str,
) -> SqliteResult<()> {
conn.execute(
"UPDATE messages SET tool_info = ?1 WHERE session_id = ?2 AND tool_call_id = ?3",
params![tool_info, session_id, tool_call_id],
)?;
Ok(())
}
pub fn message_to_chat_message(data: &MessageData) -> Result<ChatCompletionRequestMessage> {
match data.role.as_str() {
"user" => Ok(ChatCompletionRequestMessage::User(
ChatCompletionRequestUserMessage {
content: ChatCompletionRequestUserMessageContent::Text(data.content.clone()),
name: None,
}
.into(),
)),
"assistant" => {
let tool_calls = if let Some(tool_info) = &data.tool_info {
if let serde_json::Value::Object(obj) = tool_info {
if let Some(serde_json::Value::Array(arr)) = obj.get("tool_calls") {
use async_openai::types::chat::{
ChatCompletionMessageToolCall, ChatCompletionMessageToolCalls,
};
let mut calls = Vec::new();
for call_val in arr {
if let Ok(call) =
serde_json::from_value::<ChatCompletionMessageToolCall>(
call_val.clone(),
)
{
calls.push(ChatCompletionMessageToolCalls::Function(call));
}
}
if !calls.is_empty() {
Some(calls)
} else {
None
}
} else {
None
}
} else {
None
}
} else {
None
};
let content = if data.content.is_empty() {
None
} else {
Some(data.content.clone().into())
};
if content.is_none() && tool_calls.is_none() {
tracing::warn!("Skipping invalid assistant message: both content and tool_calls are None (message id: {:?})", data.id);
return Err(crate::error::AgentError::InternalError("Invalid assistant message".to_string()));
}
Ok(ChatCompletionRequestMessage::Assistant(
ChatCompletionRequestAssistantMessage {
content,
name: None,
tool_calls,
refusal: None,
audio: None,
#[allow(deprecated)]
function_call: None,
}
.into(),
))
}
"tool" => {
let tool_call_id = data.tool_call_id.clone().unwrap_or_default();
Ok(ChatCompletionRequestMessage::Tool(
ChatCompletionRequestToolMessage {
content: data.content.clone().into(),
tool_call_id,
}
.into(),
))
}
_ => Err(crate::error::AgentError::InternalError(format!(
"Unknown role: {}",
data.role
))),
}
}
#[derive(Debug, Clone, Serialize)]
pub struct MessageSearchResult {
pub message_id: i64,
pub session_id: String,
pub session_title: String,
pub role: String,
pub content_snippet: String,
pub created_at: String,
}
#[derive(Debug, Clone, Default)]
pub struct MessageSearchFilter<'a> {
pub session_id: Option<&'a str>,
pub role: Option<&'a str>,
pub since: Option<&'a str>,
pub until: Option<&'a str>,
}
pub fn search_messages(
conn: &Connection,
query: &str,
filter: &MessageSearchFilter,
limit: usize,
) -> SqliteResult<Vec<MessageSearchResult>> {
if query.trim().is_empty() {
return Ok(Vec::new());
}
let mut conditions: Vec<String> = Vec::new();
let mut params: Vec<ToSqlOutput> = Vec::new();
params.push(ToSqlOutput::from(query));
if let Some(session_id) = filter.session_id {
conditions.push("m.session_id = ?".to_string());
params.push(ToSqlOutput::from(session_id));
}
if let Some(role) = filter.role {
conditions.push("m.role = ?".to_string());
params.push(ToSqlOutput::from(role));
}
if let Some(since) = filter.since {
conditions.push("m.created_at >= ?".to_string());
params.push(ToSqlOutput::from(since));
}
if let Some(until) = filter.until {
conditions.push("m.created_at <= ?".to_string());
params.push(ToSqlOutput::from(until));
}
let where_extra = if conditions.is_empty() {
String::new()
} else {
format!(" AND {}", conditions.join(" AND "))
};
let sql = format!(
"SELECT
m.id,
m.session_id,
s.title,
m.role,
snippet(messages_fts, 0, '<b>', '</b>', '...', 16),
m.created_at
FROM messages_fts
JOIN messages m ON m.id = messages_fts.rowid
JOIN sessions s ON s.id = m.session_id
WHERE messages_fts MATCH ?1{}
ORDER BY bm25(messages_fts)
LIMIT {}",
where_extra,
limit
);
let mut stmt = conn.prepare(&sql)?;
let param_refs: Vec<&dyn rusqlite::ToSql> = params.iter().map(|p| p as &dyn rusqlite::ToSql).collect();
let rows = stmt.query_map(param_refs.as_slice(), |row| {
Ok(MessageSearchResult {
message_id: row.get(0)?,
session_id: row.get(1)?,
session_title: row.get(2)?,
role: row.get(3)?,
content_snippet: row.get(4)?,
created_at: row.get(5)?,
})
})?;
rows.collect()
}
pub fn load_chat_messages(
conn: &Connection,
session_id: &str,
) -> Result<Vec<ChatCompletionRequestMessage>> {
let messages = get_messages(conn, session_id)?;
tracing::debug!(
"load_chat_messages: session_id={}, loaded {} messages from DB",
session_id,
messages.len()
);
let mut result = Vec::with_capacity(messages.len());
for (idx, msg) in messages.iter().enumerate() {
match message_to_chat_message(&msg) {
Ok(chat_msg) => {
let is_valid = match &chat_msg {
ChatCompletionRequestMessage::Assistant(assistant_msg) => {
assistant_msg.content.is_some() || assistant_msg.tool_calls.is_some()
}
_ => true,
};
if is_valid {
tracing::debug!(
" Message {}: role={}, content_len={}",
idx,
msg.role,
msg.content.len()
);
result.push(chat_msg);
} else {
tracing::warn!(
"Skipping invalid assistant message {} (has neither content nor tool_calls)",
idx
);
}
}
Err(e) => {
tracing::warn!("Skipping invalid message {}: {}", idx, e);
}
}
}
tracing::debug!("load_chat_messages: successfully converted {} messages", result.len());
Ok(result)
}
pub fn insert_memory(conn: &Connection, memory: &Memory) -> SqliteResult<()> {
let tags_str = if memory.tags.is_empty() {
None
} else {
Some(memory.tags.join(","))
};
conn.execute(
"INSERT INTO memories (
id, session_id, chat_id, memory_type, title, content, tags, is_active, created_at, updated_at
) VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10)",
params![
memory.id,
memory.session_id,
memory.chat_id,
memory.memory_type.as_str(),
memory.title,
memory.content,
tags_str,
memory.is_active,
memory.created_at,
memory.updated_at,
],
)?;
Ok(())
}
pub fn update_memory(conn: &Connection, memory: &Memory) -> SqliteResult<()> {
let tags_str = if memory.tags.is_empty() {
None
} else {
Some(memory.tags.join(","))
};
conn.execute(
"UPDATE memories SET
title = ?1,
content = ?2,
tags = ?3,
memory_type = ?4,
updated_at = ?5
WHERE id = ?6",
params![
memory.title,
memory.content,
tags_str,
memory.memory_type.as_str(),
current_timestamp(),
memory.id,
],
)?;
Ok(())
}
pub fn deactivate_memory(conn: &Connection, memory_id: &str) -> SqliteResult<()> {
conn.execute(
"UPDATE memories SET is_active = 0, updated_at = ?1 WHERE id = ?2",
params![current_timestamp(), memory_id],
)?;
Ok(())
}
pub fn delete_memory_permanently(conn: &Connection, memory_id: &str) -> SqliteResult<()> {
conn.execute("DELETE FROM memories WHERE id = ?1", params![memory_id])?;
Ok(())
}
pub fn get_memory(conn: &Connection, memory_id: &str) -> SqliteResult<Option<Memory>> {
let mut stmt = conn.prepare(
"SELECT id, session_id, chat_id, memory_type, title, content, tags, is_active, created_at, updated_at
FROM memories WHERE id = ?1",
)?;
let mut rows = stmt.query_map(params![memory_id], map_memory_row)?;
rows.next().transpose()
}
pub fn find_memories_by_title(
conn: &Connection,
title_part: &str,
filter: &MemoryFilter,
limit: Option<usize>,
) -> SqliteResult<Vec<Memory>> {
let (sql, params) = build_memory_query(Some(title_part), filter, limit);
let mut stmt = conn.prepare(&sql)?;
let rows = stmt.query_map(rusqlite::params_from_iter(params), map_memory_row)?;
rows.collect()
}
pub fn list_memories(
conn: &Connection,
filter: &MemoryFilter,
limit: Option<usize>,
) -> SqliteResult<Vec<Memory>> {
let (sql, params) = build_memory_query(None, filter, limit);
let mut stmt = conn.prepare(&sql)?;
let rows = stmt.query_map(rusqlite::params_from_iter(params), map_memory_row)?;
rows.collect()
}
pub fn recall_memories(
conn: &Connection,
query: &str,
filter: &MemoryFilter,
limit: usize,
) -> SqliteResult<Vec<Memory>> {
let mut results = find_memories_by_title(conn, query, filter, Some(limit))?;
if results.len() >= limit {
results.truncate(limit);
return Ok(results);
}
let remaining = limit - results.len();
let (sql, params) = build_recall_query(query, filter, remaining);
let mut stmt = conn.prepare(&sql)?;
let rows = stmt.query_map(rusqlite::params_from_iter(params), map_memory_row)?;
for row in rows {
let memory = row?;
if !results.iter().any(|m| m.id == memory.id) {
results.push(memory);
}
}
results.truncate(limit);
Ok(results)
}
fn map_memory_row(row: &rusqlite::Row) -> SqliteResult<Memory> {
let tags_str: Option<String> = row.get(6)?;
let tags = tags_str
.map(|s| {
s.split(',')
.map(|t| t.trim().to_string())
.filter(|t| !t.is_empty())
.collect()
})
.unwrap_or_default();
Ok(Memory {
id: row.get(0)?,
session_id: row.get(1)?,
chat_id: row.get(2)?,
memory_type: MemoryType::from_str(&row.get::<_, String>(3)?),
title: row.get(4)?,
content: row.get(5)?,
tags,
is_active: row.get::<_, i32>(7)? != 0,
created_at: row.get(8)?,
updated_at: row.get(9)?,
})
}
fn build_memory_query<'a>(
title_search: Option<&'a str>,
filter: &'a MemoryFilter,
limit: Option<usize>,
) -> (String, Vec<rusqlite::types::ToSqlOutput<'a>>) {
let mut conditions = Vec::new();
let mut params: Vec<rusqlite::types::ToSqlOutput> = Vec::new();
if filter.only_active {
conditions.push("is_active = 1".to_string());
}
if let Some(memory_type) = &filter.memory_type {
conditions.push("memory_type = ?".to_string());
params.push(memory_type.as_str().into());
}
if let Some(session_id) = &filter.session_id {
conditions.push("session_id = ?".to_string());
params.push(session_id.as_str().into());
}
if let Some(chat_id) = &filter.chat_id {
conditions.push("chat_id = ?".to_string());
params.push(chat_id.as_str().into());
}
if let Some(since) = &filter.since {
conditions.push("created_at >= ?".to_string());
params.push(since.as_str().into());
}
if let Some(title) = title_search {
conditions.push("title LIKE ?".to_string());
params.push(format!("%{}%", title).into());
}
if let Some(tags) = &filter.tags {
if !tags.is_empty() {
let tag_conditions: Vec<_> = tags.iter().map(|_| "tags LIKE ?").collect();
conditions.push(format!("({})", tag_conditions.join(" OR ")));
for tag in tags {
params.push(format!("%{}%", tag).into());
}
}
}
let where_clause = if conditions.is_empty() {
String::new()
} else {
format!("WHERE {}", conditions.join(" AND "))
};
let limit_clause = limit.map(|l| format!("LIMIT {}", l)).unwrap_or_default();
let sql = format!(
"SELECT id, session_id, chat_id, memory_type, title, content, tags, is_active, created_at, updated_at
FROM memories
{}
ORDER BY created_at DESC
{}",
where_clause, limit_clause
);
(sql, params)
}
fn build_recall_query<'a>(
query: &'a str,
filter: &'a MemoryFilter,
limit: usize,
) -> (String, Vec<rusqlite::types::ToSqlOutput<'a>>) {
let mut conditions = Vec::new();
let mut params: Vec<rusqlite::types::ToSqlOutput> = Vec::new();
if filter.only_active {
conditions.push("is_active = 1".to_string());
}
if let Some(memory_type) = &filter.memory_type {
conditions.push("memory_type = ?".to_string());
params.push(memory_type.as_str().into());
}
if let Some(session_id) = &filter.session_id {
conditions.push("session_id = ?".to_string());
params.push(session_id.as_str().into());
}
if let Some(chat_id) = &filter.chat_id {
conditions.push("chat_id = ?".to_string());
params.push(chat_id.as_str().into());
}
conditions.push("(title LIKE ? OR content LIKE ? OR tags LIKE ?)".to_string());
let pattern = format!("%{}%", query);
params.push(pattern.clone().into());
params.push(pattern.clone().into());
params.push(pattern.into());
let where_clause = format!("WHERE {}", conditions.join(" AND "));
let sql = format!(
"SELECT id, session_id, chat_id, memory_type, title, content, tags, is_active, created_at, updated_at
FROM memories
{}
ORDER BY created_at DESC
LIMIT {}",
where_clause, limit
);
(sql, params)
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn resolves_local_db_path() {
let working_dir = PathBuf::from("project");
let path = resolve_db_path(&working_dir, false).unwrap();
assert_eq!(
path,
working_dir.join(ROBIT_DIR).join(MEMORY_DIR).join(DB_FILE)
);
}
#[test]
fn session_crud() {
let conn = Connection::open_in_memory().unwrap();
init_db(&conn).unwrap();
insert_session(
&conn,
"test-123",
None,
"Test Session",
"deepseek/deepseek-chat",
"gui",
)
.unwrap();
let sessions = list_sessions(&conn, None).unwrap();
assert_eq!(sessions.len(), 1);
assert_eq!(sessions[0].id, "test-123");
assert_eq!(sessions[0].title, "Test Session");
assert_eq!(sessions[0].source, "gui");
assert_eq!(sessions[0].chat_id, None);
assert_eq!(sessions[0].status, "idle");
let session = get_session(&conn, "test-123").unwrap().unwrap();
assert_eq!(session.title, "Test Session");
assert_eq!(session.source, "gui");
update_session_title(&conn, "test-123", "Updated Title").unwrap();
let updated = get_session(&conn, "test-123").unwrap().unwrap();
assert_eq!(updated.title, "Updated Title");
delete_session(&conn, "test-123").unwrap();
assert!(get_session(&conn, "test-123").unwrap().is_none());
assert!(list_sessions(&conn, None).unwrap().is_empty());
}
#[test]
fn message_operations() {
let conn = Connection::open_in_memory().unwrap();
init_db(&conn).unwrap();
insert_session(&conn, "session-msg", None, "Chat Session", "model", "gui").unwrap();
let user_id = insert_message(
&conn,
"session-msg",
"user",
"Hello Robit",
None,
None,
None,
)
.unwrap();
let assistant_id = insert_message(
&conn,
"session-msg",
"assistant",
"Hello! How can I help?",
None,
None,
None,
)
.unwrap();
let messages = get_messages(&conn, "session-msg").unwrap();
assert_eq!(messages.len(), 2);
assert_eq!(messages[0].id, user_id);
assert_eq!(messages[0].role, "user");
assert_eq!(messages[0].content, "Hello Robit");
assert_eq!(messages[1].id, assistant_id);
assert_eq!(messages[1].role, "assistant");
assert_eq!(messages[1].content, "Hello! How can I help?");
}
#[test]
fn empty_sessions() {
let conn = Connection::open_in_memory().unwrap();
init_db(&conn).unwrap();
let sessions = list_sessions(&conn, None).unwrap();
assert_eq!(sessions.len(), 0);
}
#[test]
fn get_nonexistent_session() {
let conn = Connection::open_in_memory().unwrap();
init_db(&conn).unwrap();
let session = get_session(&conn, "nonexistent").unwrap();
assert!(session.is_none());
}
#[test]
fn tool_message_update() {
let conn = Connection::open_in_memory().unwrap();
init_db(&conn).unwrap();
insert_session(&conn, "session-tool", None, "Tool Session", "model", "gui").unwrap();
let initial = serde_json::json!({
"tool_call_id": "tool-1",
"name": "bash",
"arguments": "{}",
"status": "pending",
"requires_confirm": true
})
.to_string();
insert_message(
&conn,
"session-tool",
"tool",
"{}",
Some("bash"),
Some("tool-1"),
Some(&initial),
)
.unwrap();
let updated = serde_json::json!({
"tool_call_id": "tool-1",
"status": "success",
"output": "done"
})
.to_string();
update_tool_message(&conn, "session-tool", "tool-1", &updated).unwrap();
let messages = get_messages(&conn, "session-tool").unwrap();
assert_eq!(messages.len(), 1);
assert_eq!(messages[0].tool_name.as_deref(), Some("bash"));
assert_eq!(messages[0].tool_call_id.as_deref(), Some("tool-1"));
assert_eq!(messages[0].tool_info.as_ref().unwrap()["status"], "success");
assert_eq!(messages[0].tool_info.as_ref().unwrap()["output"], "done");
}
#[test]
fn chat_id_lookup_and_source_filter() {
let conn = Connection::open_in_memory().unwrap();
init_db(&conn).unwrap();
insert_session(&conn, "gui-1", None, "GUI Session", "model", "gui").unwrap();
insert_session(
&conn,
"qq-1",
Some("group:abc"),
"技术讨论群",
"model",
"qq",
)
.unwrap();
insert_session(
&conn,
"qq-2",
Some("private:xyz"),
"私聊",
"model",
"qq",
)
.unwrap();
let found = find_session_by_chat_id(&conn, "group:abc").unwrap().unwrap();
assert_eq!(found.id, "qq-1");
assert_eq!(found.source, "qq");
assert_eq!(found.chat_id.as_deref(), Some("group:abc"));
assert!(find_session_by_chat_id(&conn, "does-not-exist")
.unwrap()
.is_none());
let qq_sessions = list_sessions(&conn, Some("qq")).unwrap();
assert_eq!(qq_sessions.len(), 2);
assert!(qq_sessions.iter().all(|s| s.source == "qq"));
let gui_sessions = list_sessions(&conn, Some("gui")).unwrap();
assert_eq!(gui_sessions.len(), 1);
assert_eq!(gui_sessions[0].id, "gui-1");
assert_eq!(list_sessions(&conn, None).unwrap().len(), 3);
}
#[test]
fn chat_id_unique_per_chat() {
let conn = Connection::open_in_memory().unwrap();
init_db(&conn).unwrap();
insert_session(
&conn,
"qq-1",
Some("group:abc"),
"First",
"model",
"qq",
)
.unwrap();
let err = insert_session(&conn, "qq-2", Some("group:abc"), "Second", "model", "qq");
assert!(err.is_err());
}
#[test]
fn migrates_legacy_v1_database() {
let conn = Connection::open_in_memory().unwrap();
conn.execute_batch(
"CREATE TABLE sessions (
id TEXT PRIMARY KEY,
title TEXT NOT NULL,
model TEXT NOT NULL,
created_at TEXT NOT NULL,
updated_at TEXT NOT NULL,
is_active INTEGER DEFAULT 1
);
CREATE TABLE messages (
id INTEGER PRIMARY KEY AUTOINCREMENT,
session_id TEXT NOT NULL REFERENCES sessions(id),
role TEXT NOT NULL,
content TEXT NOT NULL,
tool_name TEXT,
tool_call_id TEXT,
tokens INTEGER,
created_at TEXT NOT NULL
);",
)
.unwrap();
conn.execute(
"INSERT INTO sessions (id, title, model, created_at, updated_at) \
VALUES ('legacy-1', 'Legacy', 'model', '2020-01-01', '2020-01-01')",
[],
)
.unwrap();
init_db(&conn).unwrap();
let v: i32 = read_schema_version(&conn).unwrap();
assert_eq!(v, CURRENT_SCHEMA_VERSION);
let session = get_session(&conn, "legacy-1").unwrap().unwrap();
assert_eq!(session.title, "Legacy");
assert_eq!(session.source, "gui");
assert_eq!(session.chat_id, None);
}
#[test]
fn init_db_is_idempotent() {
let conn = Connection::open_in_memory().unwrap();
init_db(&conn).unwrap();
init_db(&conn).unwrap();
assert_eq!(read_schema_version(&conn).unwrap(), CURRENT_SCHEMA_VERSION);
}
fn setup_search_test(conn: &Connection) {
init_db(conn).unwrap();
insert_session(conn, "sess-1", None, "Session One", "model", "gui").unwrap();
insert_session(conn, "sess-2", None, "Session Two", "model", "gui").unwrap();
insert_message(conn, "sess-1", "user", "Hello, how do I write Rust code?", None, None, None).unwrap();
insert_message(conn, "sess-1", "assistant", "To write Rust code, start with cargo new.", None, None, None).unwrap();
insert_message(conn, "sess-1", "user", "What about Python?", None, None, None).unwrap();
insert_message(conn, "sess-1", "assistant", "Python is also a great language.", None, None, None).unwrap();
insert_message(conn, "sess-2", "user", "How to deploy a Rust application?", None, None, None).unwrap();
insert_message(conn, "sess-2", "assistant", "You can deploy Rust apps with Docker.", None, None, None).unwrap();
}
#[test]
fn search_messages_basic() {
let conn = Connection::open_in_memory().unwrap();
setup_search_test(&conn);
let filter = MessageSearchFilter {
session_id: None,
role: None,
since: None,
until: None,
};
let results = search_messages(&conn, "Rust", &filter, 10).unwrap();
assert!(results.len() >= 2, "Expected at least 2 results for 'Rust', got {}", results.len());
assert!(results[0].content_snippet.contains("<b>"));
}
#[test]
fn search_messages_session_filter() {
let conn = Connection::open_in_memory().unwrap();
setup_search_test(&conn);
let filter = MessageSearchFilter {
session_id: Some("sess-1"),
role: None,
since: None,
until: None,
};
let results = search_messages(&conn, "Rust", &filter, 10).unwrap();
assert_eq!(results.len(), 2);
let filter2 = MessageSearchFilter {
session_id: Some("sess-2"),
role: None,
since: None,
until: None,
};
let results2 = search_messages(&conn, "Rust", &filter2, 10).unwrap();
assert_eq!(results2.len(), 2);
}
#[test]
fn search_messages_role_filter() {
let conn = Connection::open_in_memory().unwrap();
setup_search_test(&conn);
let filter = MessageSearchFilter {
session_id: Some("sess-1"),
role: Some("user"),
since: None,
until: None,
};
let results = search_messages(&conn, "Rust", &filter, 10).unwrap();
assert_eq!(results.len(), 1);
assert_eq!(results[0].role, "user");
}
#[test]
fn search_messages_empty_query() {
let conn = Connection::open_in_memory().unwrap();
setup_search_test(&conn);
let filter = MessageSearchFilter::default();
let results = search_messages(&conn, "", &filter, 10).unwrap();
assert!(results.is_empty());
}
#[test]
fn search_messages_no_results() {
let conn = Connection::open_in_memory().unwrap();
setup_search_test(&conn);
let filter = MessageSearchFilter::default();
let results = search_messages(&conn, "nonexistent_keyword_xyz", &filter, 10).unwrap();
assert!(results.is_empty());
}
#[test]
fn search_messages_limit() {
let conn = Connection::open_in_memory().unwrap();
setup_search_test(&conn);
let filter = MessageSearchFilter {
session_id: None,
role: None,
since: None,
until: None,
};
let results = search_messages(&conn, "Rust", &filter, 2).unwrap();
assert_eq!(results.len(), 2);
}
#[test]
fn search_messages_cross_session_has_session_title() {
let conn = Connection::open_in_memory().unwrap();
setup_search_test(&conn);
let filter = MessageSearchFilter {
session_id: None,
role: None,
since: None,
until: None,
};
let results = search_messages(&conn, "Rust", &filter, 10).unwrap();
for r in &results {
assert!(!r.session_title.is_empty());
}
}
#[test]
fn migration_v3_to_v4_backfill() {
let conn = Connection::open_in_memory().unwrap();
conn.execute_batch(
"CREATE TABLE sessions (
id TEXT PRIMARY KEY,
chat_id TEXT,
title TEXT NOT NULL,
model TEXT NOT NULL,
source TEXT NOT NULL DEFAULT 'gui',
created_at TEXT NOT NULL,
updated_at TEXT NOT NULL,
is_active INTEGER DEFAULT 1
);
CREATE TABLE messages (
id INTEGER PRIMARY KEY AUTOINCREMENT,
session_id TEXT NOT NULL REFERENCES sessions(id),
role TEXT NOT NULL,
content TEXT NOT NULL,
tool_name TEXT,
tool_call_id TEXT,
tool_info TEXT,
tokens INTEGER,
created_at TEXT NOT NULL
);
CREATE TABLE memories (
id TEXT PRIMARY KEY,
session_id TEXT,
chat_id TEXT,
memory_type TEXT NOT NULL,
title TEXT NOT NULL,
content TEXT NOT NULL,
tags TEXT,
is_active INTEGER DEFAULT 1,
created_at TEXT NOT NULL,
updated_at TEXT NOT NULL
);
CREATE TABLE _schema_meta (
key TEXT PRIMARY KEY,
value TEXT NOT NULL
);
INSERT INTO _schema_meta (key, value) VALUES ('version', '3');
INSERT INTO sessions (id, title, model, source, created_at, updated_at)
VALUES ('old-sess', 'Old Session', 'model', 'gui', '2025-01-01', '2025-01-01');
INSERT INTO messages (session_id, role, content, created_at)
VALUES ('old-sess', 'user', 'This is a legacy message about testing', '2025-01-01');",
)
.unwrap();
init_db(&conn).unwrap();
assert_eq!(read_schema_version(&conn).unwrap(), 4);
let filter = MessageSearchFilter {
session_id: Some("old-sess"),
role: None,
since: None,
until: None,
};
let results = search_messages(&conn, "legacy", &filter, 10).unwrap();
assert_eq!(results.len(), 1);
assert_eq!(results[0].session_title, "Old Session");
}
}