use std::path::Path;
use std::time::{SystemTime, UNIX_EPOCH};
use sqlx::SqlitePool;
use sqlx::sqlite::{SqliteConnectOptions, SqliteJournalMode, SqlitePoolOptions};
use crate::domain::agent::ReasoningLevel;
use crate::domain::session::{DailyActivity, ReviewRequest, SessionFollowUpTask, SessionStats};
use crate::infra::agent;
#[derive(Debug, thiserror::Error)]
pub enum DbError {
#[error("{0}")]
Query(#[from] sqlx::Error),
#[error("{0}")]
Migration(#[from] sqlx::migrate::MigrateError),
#[error("{0}")]
Io(#[from] std::io::Error),
}
pub const DB_DIR: &str = "db";
pub const DB_FILE: &str = "agentty.db";
pub const DB_POOL_MAX_CONNECTIONS: u32 = 10;
#[derive(Clone)]
pub struct Database {
pool: SqlitePool,
}
pub(crate) struct SessionTurnMetadata<'a> {
pub(crate) follow_up_tasks: &'a [String],
pub(crate) instruction_conversation_id: Option<&'a str>,
pub(crate) model: &'a str,
pub(crate) provider_conversation_id: Option<&'a str>,
pub(crate) questions_json: &'a str,
pub(crate) summary: &'a str,
pub(crate) token_usage_delta: &'a SessionStats,
}
pub struct ProjectRow {
pub created_at: i64,
pub display_name: Option<String>,
pub git_branch: Option<String>,
pub id: i64,
pub is_favorite: bool,
pub last_opened_at: Option<i64>,
pub path: String,
pub updated_at: i64,
}
pub struct ProjectListRow {
pub active_session_count: i64,
pub created_at: i64,
pub display_name: Option<String>,
pub git_branch: Option<String>,
pub id: i64,
pub is_favorite: bool,
pub last_opened_at: Option<i64>,
pub last_session_updated_at: Option<i64>,
pub path: String,
pub session_count: i64,
pub updated_at: i64,
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct SessionReviewRequestRow {
pub display_id: String,
pub forge_kind: String,
pub last_refreshed_at: i64,
pub source_branch: String,
pub state: String,
pub status_summary: Option<String>,
pub target_branch: String,
pub title: String,
pub web_url: String,
}
pub struct SessionRow {
pub added_lines: i64,
pub base_branch: String,
pub created_at: i64,
pub deleted_lines: i64,
pub id: String,
pub in_progress_started_at: Option<i64>,
pub in_progress_total_seconds: i64,
pub input_tokens: i64,
pub is_draft: bool,
pub model: String,
pub output: String,
pub output_tokens: i64,
pub project_id: Option<i64>,
pub prompt: String,
pub reasoning_level_override: Option<String>,
pub published_upstream_ref: Option<String>,
pub questions: Option<String>,
pub review_request: Option<SessionReviewRequestRow>,
pub size: String,
pub status: String,
pub summary: Option<String>,
pub title: Option<String>,
pub updated_at: i64,
}
#[derive(Clone, Debug, Eq, PartialEq, sqlx::FromRow)]
pub struct SessionFollowUpTaskRow {
pub id: i64,
pub launched_session_id: Option<String>,
pub position: i64,
pub session_id: String,
pub text: String,
}
impl SessionFollowUpTaskRow {
pub fn into_session_follow_up_task(self) -> SessionFollowUpTask {
SessionFollowUpTask {
id: self.id,
launched_session_id: self.launched_session_id,
position: usize::try_from(self.position).unwrap_or(usize::MAX),
text: self.text,
}
}
}
pub struct SessionOperationRow {
pub cancel_requested: bool,
pub finished_at: Option<i64>,
pub heartbeat_at: Option<i64>,
pub id: String,
pub kind: String,
pub last_error: Option<String>,
pub queued_at: i64,
pub session_id: String,
pub started_at: Option<i64>,
pub status: String,
}
pub struct SessionUsageRow {
pub created_at: i64,
pub input_tokens: i64,
pub invocation_count: i64,
pub model: String,
pub output_tokens: i64,
pub session_id: Option<String>,
}
struct TimestampValueRow {
pub created_at: i64,
}
struct DailyActivityQueryRow {
day_key: i64,
session_count: i64,
}
impl DailyActivityQueryRow {
fn into_daily_activity(self) -> DailyActivity {
DailyActivity {
day_key: self.day_key,
session_count: u32::try_from(self.session_count).unwrap_or(u32::MAX),
}
}
}
struct SessionTimestampsRow {
created_at: i64,
updated_at: i64,
}
struct OptionalI64ValueRow {
value: Option<i64>,
}
#[derive(sqlx::FromRow)]
struct SessionInstructionStateRow {
app_server_instruction_provider_conversation_id: Option<String>,
}
impl SessionInstructionStateRow {
fn into_instruction_conversation_id(self) -> Option<String> {
agent::normalize_instruction_conversation_id(
self.app_server_instruction_provider_conversation_id
.as_deref(),
)
}
}
struct RequiredStringValueRow {
value: String,
}
struct SessionMetadataRow {
max_updated_at: i64,
session_count: i64,
}
struct RequiredBoolValueRow {
value: bool,
}
struct SessionJoinRow {
added_lines: i64,
base_branch: String,
created_at: i64,
deleted_lines: i64,
id: String,
in_progress_started_at: Option<i64>,
in_progress_total_seconds: i64,
input_tokens: i64,
is_draft: bool,
model: String,
output: String,
output_tokens: i64,
project_id: Option<i64>,
prompt: String,
reasoning_level_override: Option<String>,
published_upstream_ref: Option<String>,
questions: Option<String>,
review_request_display_id: Option<String>,
review_request_forge_kind: Option<String>,
review_request_last_refreshed_at: Option<i64>,
review_request_source_branch: Option<String>,
review_request_state: Option<String>,
review_request_status_summary: Option<String>,
review_request_target_branch: Option<String>,
review_request_title: Option<String>,
review_request_web_url: Option<String>,
size: String,
status: String,
summary: Option<String>,
title: Option<String>,
updated_at: i64,
}
impl SessionJoinRow {
fn into_session_row(self) -> SessionRow {
let Self {
added_lines,
base_branch,
created_at,
deleted_lines,
id,
in_progress_started_at,
in_progress_total_seconds,
input_tokens,
is_draft,
model,
output,
output_tokens,
project_id,
prompt,
reasoning_level_override,
published_upstream_ref,
questions,
review_request_display_id,
review_request_forge_kind,
review_request_last_refreshed_at,
review_request_source_branch,
review_request_state,
review_request_status_summary,
review_request_target_branch,
review_request_title,
review_request_web_url,
size,
status,
summary,
title,
updated_at,
} = self;
let review_request = SessionReviewRequestJoinRow {
display_id: review_request_display_id,
forge_kind: review_request_forge_kind,
last_refreshed_at: review_request_last_refreshed_at,
source_branch: review_request_source_branch,
state: review_request_state,
status_summary: review_request_status_summary,
target_branch: review_request_target_branch,
title: review_request_title,
web_url: review_request_web_url,
}
.into_review_request_row();
SessionRow {
added_lines,
base_branch,
created_at,
deleted_lines,
id,
in_progress_started_at,
in_progress_total_seconds,
input_tokens,
is_draft,
model,
output,
output_tokens,
project_id,
prompt,
reasoning_level_override,
published_upstream_ref,
questions,
review_request,
size,
status,
summary,
title,
updated_at,
}
}
}
struct SessionReviewRequestJoinRow {
display_id: Option<String>,
forge_kind: Option<String>,
last_refreshed_at: Option<i64>,
source_branch: Option<String>,
state: Option<String>,
status_summary: Option<String>,
target_branch: Option<String>,
title: Option<String>,
web_url: Option<String>,
}
impl SessionReviewRequestJoinRow {
fn into_review_request_row(self) -> Option<SessionReviewRequestRow> {
let Self {
display_id,
forge_kind,
last_refreshed_at,
source_branch,
state,
status_summary,
target_branch,
title,
web_url,
} = self;
Some(SessionReviewRequestRow {
display_id: display_id?,
forge_kind: forge_kind?,
last_refreshed_at: last_refreshed_at?,
source_branch: source_branch?,
state: state?,
status_summary,
target_branch: target_branch?,
title: title?,
web_url: web_url?,
})
}
}
impl Database {
pub async fn open(db_path: &Path) -> Result<Self, DbError> {
if let Some(parent) = db_path.parent() {
tokio::fs::create_dir_all(parent).await?;
}
let options = SqliteConnectOptions::new()
.filename(db_path)
.create_if_missing(true)
.journal_mode(SqliteJournalMode::Wal)
.foreign_keys(true);
let pool = SqlitePoolOptions::new()
.max_connections(DB_POOL_MAX_CONNECTIONS)
.connect_with(options)
.await?;
sqlx::migrate!("./migrations").run(&pool).await?;
Ok(Self { pool })
}
pub fn pool(&self) -> &SqlitePool {
&self.pool
}
pub async fn insert_session(
&self,
id: &str,
model: &str,
base_branch: &str,
status: &str,
project_id: i64,
) -> Result<(), DbError> {
self.insert_session_with_draft_mode(id, model, base_branch, status, false, project_id)
.await
}
pub async fn insert_draft_session(
&self,
id: &str,
model: &str,
base_branch: &str,
status: &str,
project_id: i64,
) -> Result<(), DbError> {
self.insert_session_with_draft_mode(id, model, base_branch, status, true, project_id)
.await
}
async fn insert_session_with_draft_mode(
&self,
id: &str,
model: &str,
base_branch: &str,
status: &str,
is_draft: bool,
project_id: i64,
) -> Result<(), DbError> {
sqlx::query(
r"
INSERT INTO session (id, model, base_branch, status, is_draft, project_id, prompt, output)
VALUES (?, ?, ?, ?, ?, ?, ?, ?)
",
)
.bind(id)
.bind(model)
.bind(base_branch)
.bind(status)
.bind(is_draft)
.bind(project_id)
.bind("")
.bind("")
.execute(&self.pool)
.await?;
Ok(())
}
pub async fn insert_session_creation_activity_now(
&self,
session_id: &str,
) -> Result<(), DbError> {
self.insert_session_creation_activity_at(session_id, unix_timestamp_now())
.await
}
pub async fn insert_session_creation_activity_at(
&self,
session_id: &str,
timestamp_seconds: i64,
) -> Result<(), DbError> {
sqlx::query(
r"
INSERT INTO session_activity (session_id, created_at)
VALUES (?, ?)
ON CONFLICT(session_id) DO NOTHING
",
)
.bind(session_id)
.bind(timestamp_seconds)
.execute(&self.pool)
.await?;
Ok(())
}
pub async fn load_sessions_for_project(
&self,
project_id: i64,
) -> Result<Vec<SessionRow>, DbError> {
let rows = sqlx::query_as!(
SessionJoinRow,
r#"
SELECT session.base_branch AS "base_branch!",
session.added_lines AS "added_lines!",
session.created_at AS "created_at!",
session.deleted_lines AS "deleted_lines!",
session.id AS "id!",
session.in_progress_started_at,
session.in_progress_total_seconds AS "in_progress_total_seconds!",
session.input_tokens AS "input_tokens!",
session.is_draft AS "is_draft!: bool",
session.model AS "model!",
session.output AS "output!",
session.output_tokens AS "output_tokens!",
session.project_id,
session.prompt AS "prompt!",
session.reasoning_level AS "reasoning_level_override?",
session.published_upstream_ref,
session.questions,
session_review_request.display_id AS "review_request_display_id?",
session_review_request.forge_kind AS "review_request_forge_kind?",
session_review_request.last_refreshed_at AS "review_request_last_refreshed_at?",
session_review_request.source_branch AS "review_request_source_branch?",
session_review_request.state AS "review_request_state?",
session_review_request.status_summary AS "review_request_status_summary?",
session_review_request.target_branch AS "review_request_target_branch?",
session_review_request.title AS "review_request_title?",
session_review_request.web_url AS "review_request_web_url?",
session.size AS "size!",
session.status AS "status!",
session.summary,
session.title,
session.updated_at AS "updated_at!"
FROM session
LEFT JOIN session_review_request
ON session_review_request.session_id = session.id
WHERE session.project_id = ?
ORDER BY session.updated_at DESC, session.id
"#,
project_id
)
.fetch_all(&self.pool)
.await?;
Ok(rows
.into_iter()
.map(SessionJoinRow::into_session_row)
.collect())
}
pub async fn load_sessions(&self) -> Result<Vec<SessionRow>, DbError> {
let rows = sqlx::query_as!(
SessionJoinRow,
r#"
SELECT session.base_branch AS "base_branch!",
session.added_lines AS "added_lines!",
session.created_at AS "created_at!",
session.deleted_lines AS "deleted_lines!",
session.id AS "id!",
session.in_progress_started_at,
session.in_progress_total_seconds AS "in_progress_total_seconds!",
session.input_tokens AS "input_tokens!",
session.is_draft AS "is_draft!: bool",
session.model AS "model!",
session.output AS "output!",
session.output_tokens AS "output_tokens!",
session.project_id,
session.prompt AS "prompt!",
session.reasoning_level AS "reasoning_level_override?",
session.published_upstream_ref,
session.questions,
session_review_request.display_id AS "review_request_display_id?",
session_review_request.forge_kind AS "review_request_forge_kind?",
session_review_request.last_refreshed_at AS "review_request_last_refreshed_at?",
session_review_request.source_branch AS "review_request_source_branch?",
session_review_request.state AS "review_request_state?",
session_review_request.status_summary AS "review_request_status_summary?",
session_review_request.target_branch AS "review_request_target_branch?",
session_review_request.title AS "review_request_title?",
session_review_request.web_url AS "review_request_web_url?",
session.size AS "size!",
session.status AS "status!",
session.summary,
session.title,
session.updated_at AS "updated_at!"
FROM session
LEFT JOIN session_review_request
ON session_review_request.session_id = session.id
ORDER BY session.updated_at DESC, session.id
"#
)
.fetch_all(&self.pool)
.await?;
Ok(rows
.into_iter()
.map(SessionJoinRow::into_session_row)
.collect())
}
pub async fn load_session_activity_timestamps(&self) -> Result<Vec<i64>, DbError> {
let rows = sqlx::query_as!(
TimestampValueRow,
r"
SELECT created_at
FROM session_activity
ORDER BY id
",
)
.fetch_all(&self.pool)
.await?;
Ok(rows.into_iter().map(|row| row.created_at).collect())
}
pub async fn load_session_activity(&self) -> Result<Vec<DailyActivity>, DbError> {
let rows = sqlx::query_as!(
DailyActivityQueryRow,
r#"
SELECT CAST(
unixepoch(datetime(created_at, 'unixepoch', 'localtime', 'start of day', 'utc')) / 86400
AS INTEGER
) AS "day_key!: _",
COUNT(*) AS "session_count!: _"
FROM session_activity
WHERE created_at IS NOT NULL
GROUP BY 1
ORDER BY 1
"#
)
.fetch_all(&self.pool)
.await?;
Ok(rows
.into_iter()
.map(DailyActivityQueryRow::into_daily_activity)
.collect())
}
pub async fn load_sessions_metadata(&self) -> Result<(i64, i64), DbError> {
let row = sqlx::query_as!(
SessionMetadataRow,
r#"
SELECT (SELECT COUNT(*) FROM session) AS "session_count!: _",
COALESCE(
(
SELECT updated_at
FROM session
ORDER BY updated_at DESC, id
LIMIT 1
),
0
) AS "max_updated_at!: _"
"#
)
.fetch_one(&self.pool)
.await?;
Ok((row.session_count, row.max_updated_at))
}
pub async fn delete_session(&self, id: &str) -> Result<(), DbError> {
sqlx::query(
r"
DELETE FROM session
WHERE id = ?
",
)
.bind(id)
.execute(&self.pool)
.await?;
Ok(())
}
pub async fn update_session_status_with_timing_at(
&self,
id: &str,
status: &str,
timestamp_seconds: i64,
) -> Result<(), DbError> {
sqlx::query(
r"
UPDATE session
SET status = ?,
in_progress_total_seconds = CASE
WHEN ? = 'InProgress' OR in_progress_started_at IS NULL THEN in_progress_total_seconds
ELSE in_progress_total_seconds + MAX(0, ? - in_progress_started_at)
END,
in_progress_started_at = CASE
WHEN ? = 'InProgress' THEN COALESCE(in_progress_started_at, ?)
ELSE NULL
END
WHERE id = ?
",
)
.bind(status)
.bind(status)
.bind(timestamp_seconds)
.bind(status)
.bind(timestamp_seconds)
.bind(id)
.execute(&self.pool)
.await?;
Ok(())
}
pub async fn update_session_updated_at(
&self,
id: &str,
updated_at: i64,
) -> Result<(), DbError> {
sqlx::query(
r"
UPDATE session
SET updated_at = ?
WHERE id = ?
",
)
.bind(updated_at)
.bind(id)
.execute(&self.pool)
.await?;
Ok(())
}
pub async fn update_session_created_at(
&self,
id: &str,
created_at: i64,
) -> Result<(), DbError> {
sqlx::query(
r"
UPDATE session
SET created_at = ?
WHERE id = ?
",
)
.bind(created_at)
.bind(id)
.execute(&self.pool)
.await?;
Ok(())
}
pub async fn clear_session_activity(&self) -> Result<(), DbError> {
sqlx::query(
r"
DELETE FROM session_activity
",
)
.execute(&self.pool)
.await?;
Ok(())
}
pub async fn backfill_session_activity_from_sessions(&self) -> Result<(), DbError> {
sqlx::query(
r"
INSERT INTO session_activity (session_id, created_at)
SELECT id, created_at
FROM session
",
)
.execute(&self.pool)
.await?;
Ok(())
}
pub async fn update_session_diff_stats(
&self,
added_lines: u64,
deleted_lines: u64,
id: &str,
size: &str,
) -> Result<(), DbError> {
sqlx::query(
r"
UPDATE session
SET added_lines = ?,
deleted_lines = ?,
size = ?
WHERE id = ?
AND (
added_lines <> ?
OR deleted_lines <> ?
OR size <> ?
)
",
)
.bind(added_lines.cast_signed())
.bind(deleted_lines.cast_signed())
.bind(size)
.bind(id)
.bind(added_lines.cast_signed())
.bind(deleted_lines.cast_signed())
.bind(size)
.execute(&self.pool)
.await?;
Ok(())
}
pub async fn update_session_questions(&self, id: &str, questions: &str) -> Result<(), DbError> {
sqlx::query(
r"
UPDATE session
SET questions = ?
WHERE id = ?
",
)
.bind(questions)
.bind(id)
.execute(&self.pool)
.await?;
Ok(())
}
pub(crate) async fn persist_session_turn_metadata(
&self,
session_id: &str,
turn_metadata: &SessionTurnMetadata<'_>,
) -> Result<(), DbError> {
let mut transaction = self.pool.begin().await?;
let session_update = sqlx::query(
r"
UPDATE session
SET questions = ?,
summary = ?,
provider_conversation_id = ?,
app_server_instruction_provider_conversation_id = ?
WHERE id = ?
",
)
.bind(turn_metadata.questions_json)
.bind(turn_metadata.summary)
.bind(turn_metadata.provider_conversation_id)
.bind(turn_metadata.instruction_conversation_id)
.bind(session_id)
.execute(&mut *transaction)
.await?;
if session_update.rows_affected() != 1 {
return Err(sqlx::Error::RowNotFound.into());
}
sqlx::query(
r"
DELETE FROM session_follow_up_task
WHERE session_id = ?
",
)
.bind(session_id)
.execute(&mut *transaction)
.await?;
for (position, follow_up_task) in turn_metadata.follow_up_tasks.iter().enumerate() {
sqlx::query(
r"
INSERT INTO session_follow_up_task (session_id, position, text)
VALUES (?, ?, ?)
",
)
.bind(session_id)
.bind(i64::try_from(position).unwrap_or(i64::MAX))
.bind(follow_up_task)
.execute(&mut *transaction)
.await?;
}
if turn_metadata.token_usage_delta.input_tokens != 0
|| turn_metadata.token_usage_delta.output_tokens != 0
{
sqlx::query(
r"
UPDATE session
SET input_tokens = input_tokens + ?,
output_tokens = output_tokens + ?
WHERE id = ?
",
)
.bind(turn_metadata.token_usage_delta.input_tokens.cast_signed())
.bind(turn_metadata.token_usage_delta.output_tokens.cast_signed())
.bind(session_id)
.execute(&mut *transaction)
.await?;
sqlx::query(
r"
INSERT INTO session_usage (session_id, model, input_tokens, output_tokens, invocation_count)
VALUES (?, ?, ?, ?, 1)
ON CONFLICT(session_id, model) DO UPDATE SET
input_tokens = input_tokens + excluded.input_tokens,
output_tokens = output_tokens + excluded.output_tokens,
invocation_count = invocation_count + 1
",
)
.bind(session_id)
.bind(turn_metadata.model)
.bind(turn_metadata.token_usage_delta.input_tokens.cast_signed())
.bind(turn_metadata.token_usage_delta.output_tokens.cast_signed())
.execute(&mut *transaction)
.await?;
}
transaction.commit().await?;
Ok(())
}
pub async fn replace_session_follow_up_tasks(
&self,
session_id: &str,
follow_up_tasks: &[String],
) -> Result<(), DbError> {
let mut transaction = self.pool.begin().await?;
sqlx::query(
r"
DELETE FROM session_follow_up_task
WHERE session_id = ?
",
)
.bind(session_id)
.execute(&mut *transaction)
.await?;
for (position, follow_up_task) in follow_up_tasks.iter().enumerate() {
sqlx::query(
r"
INSERT INTO session_follow_up_task (session_id, position, text)
VALUES (?, ?, ?)
",
)
.bind(session_id)
.bind(i64::try_from(position).unwrap_or(i64::MAX))
.bind(follow_up_task)
.execute(&mut *transaction)
.await?;
}
transaction.commit().await?;
Ok(())
}
pub async fn update_session_follow_up_task_launched_session_id(
&self,
session_id: &str,
position: usize,
launched_session_id: Option<&str>,
) -> Result<(), DbError> {
sqlx::query(
r"
UPDATE session_follow_up_task
SET launched_session_id = ?
WHERE session_id = ?
AND position = ?
",
)
.bind(launched_session_id)
.bind(session_id)
.bind(i64::try_from(position).unwrap_or(i64::MAX))
.execute(&self.pool)
.await?;
Ok(())
}
pub async fn update_session_prompt(&self, id: &str, prompt: &str) -> Result<(), DbError> {
sqlx::query(
r"
UPDATE session
SET prompt = ?
WHERE id = ?
",
)
.bind(prompt)
.bind(id)
.execute(&self.pool)
.await?;
Ok(())
}
pub async fn update_session_title(&self, id: &str, title: &str) -> Result<(), DbError> {
sqlx::query!(
r#"
UPDATE session
SET title = ?
WHERE id = ?
"#,
title,
id,
)
.execute(&self.pool)
.await?;
Ok(())
}
pub async fn update_session_title_for_prompt(
&self,
id: &str,
expected_prompt: &str,
title: &str,
) -> Result<bool, DbError> {
let result = sqlx::query!(
r#"
UPDATE session
SET title = ?
WHERE id = ?
AND prompt = ?
"#,
title,
id,
expected_prompt,
)
.execute(&self.pool)
.await?;
Ok(result.rows_affected() > 0)
}
pub async fn update_session_summary(&self, id: &str, summary: &str) -> Result<(), DbError> {
sqlx::query(
r"
UPDATE session
SET summary = ?
WHERE id = ?
",
)
.bind(summary)
.bind(id)
.execute(&self.pool)
.await?;
Ok(())
}
pub async fn update_session_stats(
&self,
id: &str,
stats: &SessionStats,
) -> Result<(), DbError> {
if stats.input_tokens == 0 && stats.output_tokens == 0 {
return Ok(());
}
sqlx::query(
r"
UPDATE session
SET input_tokens = input_tokens + ?,
output_tokens = output_tokens + ?
WHERE id = ?
",
)
.bind(stats.input_tokens.cast_signed())
.bind(stats.output_tokens.cast_signed())
.bind(id)
.execute(&self.pool)
.await?;
Ok(())
}
pub async fn update_session_model(&self, id: &str, model: &str) -> Result<(), DbError> {
sqlx::query(
r"
UPDATE session
SET model = ?
WHERE id = ?
",
)
.bind(model)
.bind(id)
.execute(&self.pool)
.await?;
Ok(())
}
pub async fn update_session_reasoning_level(
&self,
id: &str,
reasoning_level: Option<&str>,
) -> Result<(), DbError> {
sqlx::query!(
r#"
UPDATE session
SET reasoning_level = ?
WHERE id = ?
"#,
reasoning_level,
id
)
.execute(&self.pool)
.await?;
Ok(())
}
pub async fn update_session_provider_conversation_id(
&self,
id: &str,
provider_conversation_id: Option<&str>,
) -> Result<(), DbError> {
sqlx::query(
r"
UPDATE session
SET provider_conversation_id = ?
WHERE id = ?
",
)
.bind(provider_conversation_id)
.bind(id)
.execute(&self.pool)
.await?;
Ok(())
}
pub(crate) async fn update_session_instruction_conversation_id(
&self,
id: &str,
provider_conversation_id: Option<&str>,
) -> Result<(), DbError> {
sqlx::query(
r"
UPDATE session
SET app_server_instruction_provider_conversation_id = ?
WHERE id = ?
",
)
.bind(provider_conversation_id)
.bind(id)
.execute(&self.pool)
.await?;
Ok(())
}
pub async fn update_session_published_upstream_ref(
&self,
id: &str,
published_upstream_ref: Option<&str>,
) -> Result<(), DbError> {
sqlx::query(
r"
UPDATE session
SET published_upstream_ref = ?
WHERE id = ?
",
)
.bind(published_upstream_ref)
.bind(id)
.execute(&self.pool)
.await?;
Ok(())
}
pub async fn load_session_published_upstream_ref(
&self,
id: &str,
) -> Result<Option<String>, DbError> {
let value = sqlx::query_scalar!(
r"SELECT published_upstream_ref FROM session WHERE id = ?",
id
)
.fetch_optional(&self.pool)
.await?
.flatten();
Ok(value)
}
pub async fn update_session_review_request(
&self,
id: &str,
review_request: Option<&ReviewRequest>,
) -> Result<(), DbError> {
if let Some(review_request) = review_request {
sqlx::query(
r"
INSERT INTO session_review_request (
session_id,
display_id,
forge_kind,
last_refreshed_at,
source_branch,
state,
status_summary,
target_branch,
title,
web_url
)
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
ON CONFLICT(session_id) DO UPDATE
SET display_id = excluded.display_id,
forge_kind = excluded.forge_kind,
last_refreshed_at = excluded.last_refreshed_at,
source_branch = excluded.source_branch,
state = excluded.state,
status_summary = excluded.status_summary,
target_branch = excluded.target_branch,
title = excluded.title,
web_url = excluded.web_url
",
)
.bind(id)
.bind(review_request.summary.display_id.as_str())
.bind(review_request.summary.forge_kind.as_str())
.bind(review_request.last_refreshed_at)
.bind(review_request.summary.source_branch.as_str())
.bind(review_request.summary.state.as_str())
.bind(review_request.summary.status_summary.as_deref())
.bind(review_request.summary.target_branch.as_str())
.bind(review_request.summary.title.as_str())
.bind(review_request.summary.web_url.as_str())
.execute(&self.pool)
.await?;
} else {
sqlx::query(
r"
DELETE FROM session_review_request
WHERE session_id = ?
",
)
.bind(id)
.execute(&self.pool)
.await?;
}
Ok(())
}
pub async fn replace_session_output(&self, id: &str, output: &str) -> Result<(), DbError> {
sqlx::query(
r"
UPDATE session
SET output = ?
WHERE id = ?
",
)
.bind(output)
.bind(id)
.execute(&self.pool)
.await?;
Ok(())
}
pub async fn append_session_output(&self, id: &str, chunk: &str) -> Result<(), DbError> {
sqlx::query(
r"
UPDATE session
SET output = output || ?
WHERE id = ?
",
)
.bind(chunk)
.bind(id)
.execute(&self.pool)
.await?;
Ok(())
}
pub async fn backfill_session_project(&self, project_id: i64) -> Result<(), DbError> {
sqlx::query(
r"
UPDATE session
SET project_id = ?
WHERE project_id IS NULL
",
)
.bind(project_id)
.execute(&self.pool)
.await?;
Ok(())
}
pub async fn get_session_base_branch(&self, id: &str) -> Result<Option<String>, DbError> {
let row = sqlx::query_as!(
RequiredStringValueRow,
r#"
SELECT base_branch AS "value!: _"
FROM session
WHERE id = ?
"#,
id
)
.fetch_optional(&self.pool)
.await?;
Ok(row.map(|row| row.value))
}
pub async fn get_session_provider_conversation_id(
&self,
id: &str,
) -> Result<Option<String>, DbError> {
let value = sqlx::query_scalar!(
r"SELECT provider_conversation_id FROM session WHERE id = ?",
id
)
.fetch_optional(&self.pool)
.await?
.flatten();
Ok(value)
}
pub(crate) async fn get_session_instruction_conversation_id(
&self,
id: &str,
) -> Result<Option<String>, DbError> {
let row = sqlx::query_as::<_, SessionInstructionStateRow>(
r"
SELECT app_server_instruction_provider_conversation_id
FROM session
WHERE id = ?
",
)
.bind(id)
.fetch_optional(&self.pool)
.await?;
Ok(row.and_then(SessionInstructionStateRow::into_instruction_conversation_id))
}
pub async fn insert_session_operation(
&self,
operation_id: &str,
session_id: &str,
kind: &str,
) -> Result<(), DbError> {
let queued_at = unix_timestamp_now();
sqlx::query(
r"
INSERT INTO session_operation (id, session_id, kind, status, queued_at)
VALUES (?, ?, ?, 'queued', ?)
",
)
.bind(operation_id)
.bind(session_id)
.bind(kind)
.bind(queued_at)
.execute(&self.pool)
.await?;
Ok(())
}
pub async fn load_unfinished_session_operations(
&self,
) -> Result<Vec<SessionOperationRow>, DbError> {
let rows = sqlx::query_as!(
SessionOperationRow,
r#"
SELECT id AS "id!", session_id AS "session_id!", kind AS "kind!", status AS "status!",
queued_at, started_at, finished_at,
heartbeat_at, last_error,
cancel_requested AS "cancel_requested: _"
FROM session_operation
WHERE status IN ('queued', 'running')
ORDER BY queued_at ASC, id ASC
"#
)
.fetch_all(&self.pool)
.await?;
Ok(rows)
}
pub async fn is_session_operation_unfinished(
&self,
operation_id: &str,
) -> Result<bool, DbError> {
let row = sqlx::query_as!(
RequiredBoolValueRow,
r#"
SELECT EXISTS(
SELECT 1
FROM session_operation
WHERE id = ?
AND status IN ('queued', 'running')
) AS "value!: _"
"#,
operation_id
)
.fetch_one(&self.pool)
.await?;
Ok(row.value)
}
pub async fn mark_session_operation_running(&self, operation_id: &str) -> Result<(), DbError> {
let now = unix_timestamp_now();
sqlx::query(
r"
UPDATE session_operation
SET status = 'running',
started_at = COALESCE(started_at, ?),
heartbeat_at = ?,
last_error = NULL
WHERE id = ?
",
)
.bind(now)
.bind(now)
.bind(operation_id)
.execute(&self.pool)
.await?;
Ok(())
}
pub async fn mark_session_operation_done(&self, operation_id: &str) -> Result<(), DbError> {
let now = unix_timestamp_now();
sqlx::query(
r"
UPDATE session_operation
SET status = 'done',
finished_at = ?,
heartbeat_at = ?,
last_error = NULL
WHERE id = ?
",
)
.bind(now)
.bind(now)
.bind(operation_id)
.execute(&self.pool)
.await?;
Ok(())
}
pub async fn mark_session_operation_failed(
&self,
operation_id: &str,
error: &str,
) -> Result<(), DbError> {
let now = unix_timestamp_now();
sqlx::query(
r"
UPDATE session_operation
SET status = 'failed',
finished_at = ?,
heartbeat_at = ?,
last_error = ?
WHERE id = ?
",
)
.bind(now)
.bind(now)
.bind(error)
.bind(operation_id)
.execute(&self.pool)
.await?;
Ok(())
}
pub async fn mark_session_operation_canceled(
&self,
operation_id: &str,
reason: &str,
) -> Result<(), DbError> {
let now = unix_timestamp_now();
sqlx::query(
r"
UPDATE session_operation
SET status = 'canceled',
finished_at = ?,
heartbeat_at = ?,
last_error = ?
WHERE id = ?
",
)
.bind(now)
.bind(now)
.bind(reason)
.bind(operation_id)
.execute(&self.pool)
.await?;
Ok(())
}
pub async fn request_cancel_for_session_operations(
&self,
session_id: &str,
) -> Result<(), DbError> {
sqlx::query(
r"
UPDATE session_operation
SET cancel_requested = 1
WHERE session_id = ?
AND status IN ('queued', 'running')
",
)
.bind(session_id)
.execute(&self.pool)
.await?;
Ok(())
}
pub async fn is_cancel_requested_for_operation(
&self,
operation_id: &str,
) -> Result<bool, DbError> {
let row = sqlx::query_as!(
RequiredBoolValueRow,
r#"
SELECT EXISTS(
SELECT 1
FROM session_operation
WHERE id = ?
AND cancel_requested = 1
AND status IN ('queued', 'running')
) AS "value!: _"
"#,
operation_id
)
.fetch_one(&self.pool)
.await?;
Ok(row.value)
}
pub async fn fail_unfinished_session_operations(&self, reason: &str) -> Result<(), DbError> {
let now = unix_timestamp_now();
sqlx::query(
r"
UPDATE session_operation
SET status = 'failed',
finished_at = ?,
heartbeat_at = ?,
last_error = ?,
cancel_requested = 1
WHERE status IN ('queued', 'running')
",
)
.bind(now)
.bind(now)
.bind(reason)
.execute(&self.pool)
.await?;
Ok(())
}
pub async fn load_session_project_id(&self, session_id: &str) -> Result<Option<i64>, DbError> {
let row = sqlx::query_as!(
OptionalI64ValueRow,
r#"
SELECT project_id AS "value: _"
FROM session
WHERE id = ?
"#,
session_id
)
.fetch_optional(&self.pool)
.await?;
Ok(row.and_then(|row| row.value))
}
pub async fn load_session_reasoning_level_override(
&self,
session_id: &str,
) -> Result<Option<ReasoningLevel>, DbError> {
let value = sqlx::query_scalar!(
r"SELECT reasoning_level FROM session WHERE id = ?",
session_id
)
.fetch_optional(&self.pool)
.await?
.flatten();
Ok(value.and_then(|value| value.parse::<ReasoningLevel>().ok()))
}
pub async fn load_session_summary(&self, session_id: &str) -> Result<Option<String>, DbError> {
let row = sqlx::query_scalar::<_, Option<String>>(
r"
SELECT summary
FROM session
WHERE id = ?
",
)
.bind(session_id)
.fetch_optional(&self.pool)
.await?;
Ok(row.flatten())
}
pub async fn load_session_follow_up_tasks(
&self,
) -> Result<Vec<SessionFollowUpTaskRow>, DbError> {
let rows = sqlx::query_as::<_, SessionFollowUpTaskRow>(
r"
SELECT id,
launched_session_id,
position,
session_id,
text
FROM session_follow_up_task
ORDER BY session_id, position, id
",
)
.fetch_all(&self.pool)
.await?;
Ok(rows)
}
pub async fn upsert_session_usage(
&self,
session_id: &str,
model: &str,
stats: &SessionStats,
) -> Result<(), DbError> {
if stats.input_tokens == 0 && stats.output_tokens == 0 {
return Ok(());
}
sqlx::query(
r"
INSERT INTO session_usage (session_id, model, input_tokens, output_tokens, invocation_count)
VALUES (?, ?, ?, ?, 1)
ON CONFLICT(session_id, model) DO UPDATE SET
input_tokens = input_tokens + excluded.input_tokens,
output_tokens = output_tokens + excluded.output_tokens,
invocation_count = invocation_count + 1
",
)
.bind(session_id)
.bind(model)
.bind(stats.input_tokens.cast_signed())
.bind(stats.output_tokens.cast_signed())
.execute(&self.pool)
.await?;
Ok(())
}
pub async fn load_session_usage(
&self,
session_id: &str,
) -> Result<Vec<SessionUsageRow>, DbError> {
let rows = sqlx::query_as!(
SessionUsageRow,
r#"
SELECT session_id, model, created_at, input_tokens, invocation_count, output_tokens
FROM session_usage
WHERE session_id = ?
ORDER BY model
"#,
session_id
)
.fetch_all(&self.pool)
.await?;
Ok(rows)
}
pub async fn load_session_timestamps(
&self,
session_id: &str,
) -> Result<Option<(i64, i64)>, DbError> {
let row = sqlx::query_as!(
SessionTimestampsRow,
r#"
SELECT created_at, updated_at
FROM session
WHERE id = ?
"#,
session_id
)
.fetch_optional(&self.pool)
.await?;
Ok(row.map(|row| (row.created_at, row.updated_at)))
}
}
pub(crate) fn unix_timestamp_now() -> i64 {
SystemTime::now()
.duration_since(UNIX_EPOCH)
.map_or(0, |duration| i64::try_from(duration.as_secs()).unwrap_or(0))
}
impl Database {
pub async fn open_in_memory() -> Result<Self, DbError> {
let options = SqliteConnectOptions::new()
.filename(":memory:")
.journal_mode(SqliteJournalMode::Wal)
.foreign_keys(true);
let pool = SqlitePoolOptions::new()
.max_connections(1)
.connect_with(options)
.await?;
sqlx::migrate!("./migrations").run(&pool).await?;
Ok(Self { pool })
}
}
#[cfg(test)]
mod tests {
use std::env;
use std::process::Command;
use tempfile::tempdir;
use super::*;
use crate::agent::AgentModel;
use crate::domain::agent::ReasoningLevel;
use crate::domain::session::{ForgeKind, ReviewRequestState, ReviewRequestSummary};
use crate::domain::setting::SettingName;
const DST_TEST_SUBPROCESS_ENV: &str = "AGENTTY_DST_TEST_SUBPROCESS";
fn review_request_fixture() -> ReviewRequest {
ReviewRequest {
last_refreshed_at: 456,
summary: ReviewRequestSummary {
display_id: "#42".to_string(),
forge_kind: ForgeKind::GitHub,
source_branch: "feature/forge".to_string(),
state: ReviewRequestState::Open,
status_summary: Some("2 approvals, checks passing".to_string()),
target_branch: "main".to_string(),
title: "Add forge review support".to_string(),
web_url: "https://github.com/agentty-xyz/agentty/pull/42".to_string(),
},
}
}
fn assert_review_request_row(row: &SessionRow) {
assert_eq!(
row.review_request
.as_ref()
.map(|review_request| review_request.display_id.as_str()),
Some("#42")
);
assert_eq!(
row.review_request
.as_ref()
.map(|review_request| review_request.forge_kind.as_str()),
Some("GitHub")
);
assert_eq!(
row.review_request
.as_ref()
.map(|review_request| review_request.last_refreshed_at),
Some(456)
);
assert_eq!(
row.review_request
.as_ref()
.map(|review_request| review_request.source_branch.as_str()),
Some("feature/forge")
);
assert_eq!(
row.review_request
.as_ref()
.map(|review_request| review_request.state.as_str()),
Some("Open")
);
assert_eq!(
row.review_request
.as_ref()
.and_then(|review_request| review_request.status_summary.as_deref()),
Some("2 approvals, checks passing")
);
assert_eq!(
row.review_request
.as_ref()
.map(|review_request| review_request.target_branch.as_str()),
Some("main")
);
assert_eq!(
row.review_request
.as_ref()
.map(|review_request| review_request.title.as_str()),
Some("Add forge review support")
);
assert_eq!(
row.review_request
.as_ref()
.map(|review_request| review_request.web_url.as_str()),
Some("https://github.com/agentty-xyz/agentty/pull/42")
);
}
async fn insert_session_fixture(
database: &Database,
session_id: &str,
base_branch: &str,
status: &str,
project_id: i64,
) {
database
.insert_session(session_id, "gpt-5.4", base_branch, status, project_id)
.await
.expect("failed to insert session fixture");
}
async fn load_session_row(database: &Database, session_id: &str) -> SessionRow {
database
.load_sessions()
.await
.expect("failed to load all sessions")
.into_iter()
.find(|row| row.id == session_id)
.expect("missing session row")
}
async fn load_session_operation_row(
database: &Database,
operation_id: &str,
) -> SessionOperationRow {
sqlx::query_as!(
SessionOperationRow,
r#"
SELECT id AS "id!", session_id AS "session_id!", kind AS "kind!", status AS "status!",
queued_at, started_at, finished_at,
heartbeat_at, last_error, cancel_requested AS "cancel_requested: _"
FROM session_operation
WHERE id = ?
"#,
operation_id
)
.fetch_one(database.pool())
.await
.expect("failed to load session operation row")
}
struct SessionUsageSessionIdRow {
session_id: Option<String>,
}
fn session_join_row_fixture() -> SessionJoinRow {
SessionJoinRow {
added_lines: 14,
base_branch: "main".to_string(),
created_at: 100,
deleted_lines: 6,
id: "session-a".to_string(),
in_progress_started_at: None,
in_progress_total_seconds: 0,
input_tokens: 11,
is_draft: false,
model: "gpt-5.4".to_string(),
output: "Saved output".to_string(),
output_tokens: 29,
project_id: Some(7),
prompt: "Implement feature".to_string(),
reasoning_level_override: None,
published_upstream_ref: Some("origin/session-a".to_string()),
questions: Some("Question text".to_string()),
review_request_display_id: Some("#42".to_string()),
review_request_forge_kind: Some("GitHub".to_string()),
review_request_last_refreshed_at: Some(456),
review_request_source_branch: Some("feature/forge".to_string()),
review_request_state: Some("Open".to_string()),
review_request_status_summary: Some("2 approvals, checks passing".to_string()),
review_request_target_branch: Some("main".to_string()),
review_request_title: Some("Add forge review support".to_string()),
review_request_web_url: Some(
"https://github.com/agentty-xyz/agentty/pull/42".to_string(),
),
size: "M".to_string(),
status: "Review".to_string(),
summary: Some("Summary text".to_string()),
title: Some("Review session".to_string()),
updated_at: 200,
}
}
#[tokio::test]
async fn test_open_creates_missing_parent_directory() {
let temp_dir = tempdir().expect("temp dir should be created");
let db_path = temp_dir.path().join("nested").join("db").join(DB_FILE);
let database = Database::open(&db_path)
.await
.expect("database should open with missing parent directories");
assert!(db_path.parent().is_some_and(std::path::Path::is_dir));
assert!(!database.pool().is_closed());
}
#[tokio::test]
async fn test_load_sessions_maps_joined_session_fields() {
let (database, project_id) = database_with_joined_session_fields().await;
let session_row = load_session_row(&database, "session-a").await;
assert_eq!(session_row.id, "session-a");
assert_eq!(session_row.base_branch, "main");
assert_eq!(session_row.created_at, 100);
assert_eq!(session_row.updated_at, 200);
assert_eq!(session_row.model, "claude-opus-4.1");
assert_eq!(session_row.status, "Review");
assert_eq!(session_row.in_progress_started_at, None);
assert_eq!(session_row.in_progress_total_seconds, 120);
assert_eq!(session_row.project_id, Some(project_id));
assert_eq!(session_row.prompt, "Implement the feature");
assert_eq!(session_row.output, "First line\nSecond line");
assert_eq!(session_row.added_lines, 14);
assert_eq!(session_row.deleted_lines, 6);
assert_eq!(session_row.input_tokens, 11);
assert_eq!(session_row.output_tokens, 29);
assert_eq!(session_row.size, "L");
assert_eq!(
session_row.summary.as_deref(),
Some("Implemented the requested feature")
);
assert_eq!(session_row.questions.as_deref(), Some("[\"Need logs?\"]"));
assert_eq!(session_row.title.as_deref(), Some("Feature work"));
assert_eq!(
session_row.published_upstream_ref.as_deref(),
Some("origin/agentty/session-a")
);
assert_review_request_row(&session_row);
assert_joined_session_follow_up_tasks(&database).await;
}
async fn database_with_joined_session_fields() -> (Database, i64) {
let database = Database::open_in_memory()
.await
.expect("failed to open in-memory db");
let project_id = database
.upsert_project("/tmp/project", Some("main"))
.await
.expect("failed to insert project");
let review_request = review_request_fixture();
insert_session_fixture(&database, "session-a", "main", "Review", project_id).await;
persist_joined_session_metadata(&database, &review_request).await;
persist_joined_session_output(&database).await;
(database, project_id)
}
async fn persist_joined_session_metadata(database: &Database, review_request: &ReviewRequest) {
database
.update_session_created_at("session-a", 100)
.await
.expect("failed to update session created_at");
database
.update_session_updated_at("session-a", 200)
.await
.expect("failed to update session updated_at");
database
.update_session_diff_stats(14, 6, "session-a", "L")
.await
.expect("failed to update session diff stats");
database
.update_session_questions("session-a", "[\"Need logs?\"]")
.await
.expect("failed to update session questions");
database
.replace_session_follow_up_tasks(
"session-a",
&[
"Document the new shortcut.".to_string(),
"Add a session-view regression test.".to_string(),
],
)
.await
.expect("failed to replace session follow-up tasks");
database
.update_session_prompt("session-a", "Implement the feature")
.await
.expect("failed to update session prompt");
database
.update_session_title("session-a", "Feature work")
.await
.expect("failed to update session title");
database
.update_session_summary("session-a", "Implemented the requested feature")
.await
.expect("failed to update session summary");
database
.update_session_stats(
"session-a",
&SessionStats {
added_lines: 0,
deleted_lines: 0,
input_tokens: 11,
output_tokens: 29,
},
)
.await
.expect("failed to update session stats");
database
.update_session_model("session-a", "claude-opus-4.1")
.await
.expect("failed to update session model");
database
.update_session_published_upstream_ref("session-a", Some("origin/agentty/session-a"))
.await
.expect("failed to update published upstream ref");
database
.update_session_review_request("session-a", Some(review_request))
.await
.expect("failed to update review request");
}
async fn persist_joined_session_output(database: &Database) {
database
.update_session_status_with_timing_at("session-a", "InProgress", 50)
.await
.expect("failed to open in-progress timing window");
database
.update_session_status_with_timing_at("session-a", "Review", 170)
.await
.expect("failed to close in-progress timing window");
database
.replace_session_output("session-a", "First line")
.await
.expect("failed to replace session output");
database
.append_session_output("session-a", "\nSecond line")
.await
.expect("failed to append session output");
database
.update_session_updated_at("session-a", 200)
.await
.expect("failed to update session updated_at");
}
async fn assert_joined_session_follow_up_tasks(database: &Database) {
let follow_up_tasks = database
.load_session_follow_up_tasks()
.await
.expect("failed to load session follow-up tasks");
let follow_up_task_text = follow_up_tasks
.into_iter()
.filter(|task| task.session_id == "session-a")
.map(|task| task.text)
.collect::<Vec<_>>();
assert_eq!(
follow_up_task_text,
vec![
"Document the new shortcut.".to_string(),
"Add a session-view regression test.".to_string()
]
);
}
#[tokio::test]
async fn test_update_session_title_for_prompt_requires_matching_prompt() {
let database = Database::open_in_memory()
.await
.expect("failed to open in-memory db");
let project_id = database
.upsert_project("/tmp/project", Some("main"))
.await
.expect("failed to insert project");
insert_session_fixture(&database, "session-a", "main", "New", project_id).await;
database
.update_session_prompt("session-a", "First draft")
.await
.expect("failed to persist first staged prompt");
database
.update_session_title("session-a", "First draft")
.await
.expect("failed to persist fallback title");
let stale_update_applied = database
.update_session_title_for_prompt(
"session-a",
"Second draft",
"Refine draft workflow title",
)
.await
.expect("failed to reject stale title update");
let matching_update_applied = database
.update_session_title_for_prompt(
"session-a",
"First draft",
"Refine draft workflow title",
)
.await
.expect("failed to apply matching title update");
let session_row = load_session_row(&database, "session-a").await;
assert!(!stale_update_applied);
assert!(matching_update_applied);
assert_eq!(
session_row.title.as_deref(),
Some("Refine draft workflow title")
);
}
#[tokio::test]
async fn test_update_session_status_with_timing_at_accumulates_repeated_intervals() {
let database = Database::open_in_memory()
.await
.expect("failed to open in-memory db");
let project_id = database
.upsert_project("/tmp/project", Some("main"))
.await
.expect("failed to insert project");
insert_session_fixture(&database, "session-a", "main", "New", project_id).await;
database
.update_session_status_with_timing_at("session-a", "InProgress", 10)
.await
.expect("failed to enter in-progress the first time");
database
.update_session_status_with_timing_at("session-a", "Review", 70)
.await
.expect("failed to leave in-progress the first time");
database
.update_session_status_with_timing_at("session-a", "InProgress", 100)
.await
.expect("failed to enter in-progress the second time");
database
.update_session_status_with_timing_at("session-a", "Question", 190)
.await
.expect("failed to leave in-progress the second time");
let session_row = load_session_row(&database, "session-a").await;
assert_eq!(session_row.status, "Question");
assert_eq!(session_row.in_progress_started_at, None);
assert_eq!(session_row.in_progress_total_seconds, 150);
}
#[tokio::test]
async fn test_load_sessions_for_project_filters_to_project_rows() {
let database = Database::open_in_memory()
.await
.expect("failed to open in-memory db");
let first_project_id = database
.upsert_project("/tmp/project-a", Some("main"))
.await
.expect("failed to insert first project");
let second_project_id = database
.upsert_project("/tmp/project-b", Some("develop"))
.await
.expect("failed to insert second project");
insert_session_fixture(&database, "session-a", "main", "Review", first_project_id).await;
insert_session_fixture(&database, "session-b", "main", "Done", first_project_id).await;
insert_session_fixture(&database, "session-c", "develop", "Done", second_project_id).await;
database
.update_session_updated_at("session-a", 300)
.await
.expect("failed to update session-a updated_at");
database
.update_session_updated_at("session-b", 200)
.await
.expect("failed to update session-b updated_at");
database
.update_session_updated_at("session-c", 100)
.await
.expect("failed to update session-c updated_at");
let session_rows = database
.load_sessions_for_project(first_project_id)
.await
.expect("failed to load project sessions");
assert_eq!(session_rows.len(), 2);
assert_eq!(session_rows[0].id, "session-a");
assert_eq!(session_rows[1].id, "session-b");
assert!(
session_rows
.iter()
.all(|row| row.project_id == Some(first_project_id))
);
}
#[tokio::test]
async fn test_load_sessions_metadata_returns_count_and_latest_timestamp() {
let database = Database::open_in_memory()
.await
.expect("failed to open in-memory db");
let project_id = database
.upsert_project("/tmp/project", Some("main"))
.await
.expect("failed to insert project");
insert_session_fixture(&database, "session-a", "main", "Review", project_id).await;
insert_session_fixture(&database, "session-b", "main", "Done", project_id).await;
database
.update_session_updated_at("session-a", 200)
.await
.expect("failed to update session-a updated_at");
database
.update_session_updated_at("session-b", 300)
.await
.expect("failed to update session-b updated_at");
let session_metadata = database
.load_sessions_metadata()
.await
.expect("failed to load session metadata");
assert_eq!(session_metadata, (2, 300));
}
#[tokio::test]
async fn test_load_session_timestamps_returns_created_and_updated_values() {
let database = Database::open_in_memory()
.await
.expect("failed to open in-memory db");
let project_id = database
.upsert_project("/tmp/project", Some("main"))
.await
.expect("failed to insert project");
insert_session_fixture(&database, "session-a", "main", "Done", project_id).await;
database
.update_session_created_at("session-a", 111)
.await
.expect("failed to update session created_at");
database
.update_session_updated_at("session-a", 222)
.await
.expect("failed to update session updated_at");
let session_timestamps = database
.load_session_timestamps("session-a")
.await
.expect("failed to load session timestamps");
assert_eq!(session_timestamps, Some((111, 222)));
}
#[tokio::test]
async fn test_get_session_base_branch_returns_persisted_value() {
let database = Database::open_in_memory()
.await
.expect("failed to open in-memory db");
let project_id = database
.upsert_project("/tmp/project", Some("main"))
.await
.expect("failed to insert project");
insert_session_fixture(&database, "session-a", "release", "Done", project_id).await;
let base_branch = database
.get_session_base_branch("session-a")
.await
.expect("failed to load session base branch");
assert_eq!(base_branch.as_deref(), Some("release"));
}
#[tokio::test]
async fn test_delete_session_removes_row_and_nulls_usage_foreign_key() {
let database = Database::open_in_memory()
.await
.expect("failed to open in-memory db");
let project_id = database
.upsert_project("/tmp/project", Some("main"))
.await
.expect("failed to insert project");
insert_session_fixture(&database, "session-a", "main", "Done", project_id).await;
database
.upsert_session_usage(
"session-a",
"claude-opus-4.1",
&SessionStats {
added_lines: 0,
deleted_lines: 0,
input_tokens: 11,
output_tokens: 29,
},
)
.await
.expect("failed to insert usage row");
database
.delete_session("session-a")
.await
.expect("failed to delete session");
let deleted_session = database
.load_session_timestamps("session-a")
.await
.expect("failed to load deleted session timestamps");
let retained_usage_row = sqlx::query_as!(
SessionUsageSessionIdRow,
r#"
SELECT session_id AS "session_id: _"
FROM session_usage
WHERE model = ?
"#,
"claude-opus-4.1"
)
.fetch_one(database.pool())
.await
.expect("failed to load retained usage row");
assert_eq!(deleted_session, None);
assert_eq!(retained_usage_row.session_id, None,);
}
#[tokio::test]
async fn test_load_unfinished_session_operations_returns_only_queued_and_running_rows() {
let database = Database::open_in_memory()
.await
.expect("failed to open in-memory db");
let project_id = database
.upsert_project("/tmp/project", Some("main"))
.await
.expect("failed to insert project");
insert_session_fixture(&database, "session-a", "main", "Review", project_id).await;
database
.insert_session_operation("operation-queued", "session-a", "merge")
.await
.expect("failed to insert queued operation");
database
.insert_session_operation("operation-running", "session-a", "sync")
.await
.expect("failed to insert running operation");
database
.insert_session_operation("operation-done", "session-a", "review")
.await
.expect("failed to insert done operation");
database
.mark_session_operation_running("operation-running")
.await
.expect("failed to mark running operation");
database
.mark_session_operation_running("operation-done")
.await
.expect("failed to mark done operation running");
database
.mark_session_operation_done("operation-done")
.await
.expect("failed to mark done operation");
let unfinished_rows = database
.load_unfinished_session_operations()
.await
.expect("failed to load unfinished operations");
assert_eq!(unfinished_rows.len(), 2);
assert_eq!(unfinished_rows[0].id, "operation-queued");
assert_eq!(unfinished_rows[0].status, "queued");
assert_eq!(unfinished_rows[1].id, "operation-running");
assert_eq!(unfinished_rows[1].status, "running");
}
#[tokio::test]
async fn test_request_cancel_for_session_operations_marks_only_unfinished_rows() {
let database = Database::open_in_memory()
.await
.expect("failed to open in-memory db");
let project_id = database
.upsert_project("/tmp/project", Some("main"))
.await
.expect("failed to insert project");
insert_session_fixture(&database, "session-a", "main", "Review", project_id).await;
database
.insert_session_operation("operation-queued", "session-a", "merge")
.await
.expect("failed to insert queued operation");
database
.insert_session_operation("operation-done", "session-a", "review")
.await
.expect("failed to insert done operation");
database
.mark_session_operation_running("operation-done")
.await
.expect("failed to mark done operation running");
database
.mark_session_operation_done("operation-done")
.await
.expect("failed to mark done operation");
database
.request_cancel_for_session_operations("session-a")
.await
.expect("failed to request cancel");
let queued_row = load_session_operation_row(&database, "operation-queued").await;
let done_row = load_session_operation_row(&database, "operation-done").await;
assert!(queued_row.cancel_requested);
assert!(!done_row.cancel_requested);
}
#[tokio::test]
async fn test_is_session_operation_unfinished_returns_false_for_done_operation() {
let database = Database::open_in_memory()
.await
.expect("failed to open in-memory db");
let project_id = database
.upsert_project("/tmp/project", Some("main"))
.await
.expect("failed to insert project");
insert_session_fixture(&database, "session-a", "main", "Review", project_id).await;
database
.insert_session_operation("operation-a", "session-a", "merge")
.await
.expect("failed to insert operation");
database
.mark_session_operation_running("operation-a")
.await
.expect("failed to mark operation running");
database
.mark_session_operation_done("operation-a")
.await
.expect("failed to mark operation done");
let is_unfinished = database
.is_session_operation_unfinished("operation-a")
.await
.expect("failed to check unfinished operation state");
assert!(!is_unfinished);
}
#[tokio::test]
async fn test_is_cancel_requested_for_operation_scoped_to_single_operation() {
let database = Database::open_in_memory()
.await
.expect("failed to open in-memory db");
let project_id = database
.upsert_project("/tmp/project", Some("main"))
.await
.expect("failed to insert project");
insert_session_fixture(&database, "session-a", "main", "Review", project_id).await;
database
.insert_session_operation("operation-cancelled", "session-a", "reply")
.await
.expect("failed to insert cancelled operation");
database
.insert_session_operation("operation-new", "session-a", "reply")
.await
.expect("failed to insert new operation");
database
.request_cancel_for_session_operations("session-a")
.await
.expect("failed to request cancel");
sqlx::query("UPDATE session_operation SET cancel_requested = 0 WHERE id = 'operation-new'")
.execute(&database.pool)
.await
.expect("failed to reset new operation flag");
let cancelled_flag = database
.is_cancel_requested_for_operation("operation-cancelled")
.await
.expect("failed to check cancelled operation");
let new_flag = database
.is_cancel_requested_for_operation("operation-new")
.await
.expect("failed to check new operation");
assert!(cancelled_flag);
assert!(!new_flag);
}
#[tokio::test]
async fn test_mark_session_operation_running_sets_started_at_and_heartbeat() {
let database = Database::open_in_memory()
.await
.expect("failed to open in-memory db");
let project_id = database
.upsert_project("/tmp/project", Some("main"))
.await
.expect("failed to insert project");
insert_session_fixture(&database, "session-a", "main", "Review", project_id).await;
database
.insert_session_operation("operation-a", "session-a", "merge")
.await
.expect("failed to insert operation");
database
.mark_session_operation_running("operation-a")
.await
.expect("failed to mark operation running");
let running_row = load_session_operation_row(&database, "operation-a").await;
assert_eq!(running_row.status, "running");
assert!(running_row.started_at.is_some());
assert!(running_row.heartbeat_at.is_some());
assert_eq!(running_row.last_error, None);
}
#[tokio::test]
async fn test_mark_session_operation_done_sets_finished_state() {
let database = Database::open_in_memory()
.await
.expect("failed to open in-memory db");
let project_id = database
.upsert_project("/tmp/project", Some("main"))
.await
.expect("failed to insert project");
insert_session_fixture(&database, "session-a", "main", "Review", project_id).await;
database
.insert_session_operation("operation-a", "session-a", "merge")
.await
.expect("failed to insert operation");
database
.mark_session_operation_running("operation-a")
.await
.expect("failed to mark operation running");
database
.mark_session_operation_done("operation-a")
.await
.expect("failed to mark operation done");
let done_row = load_session_operation_row(&database, "operation-a").await;
assert_eq!(done_row.status, "done");
assert!(done_row.finished_at.is_some());
assert!(done_row.heartbeat_at.is_some());
assert_eq!(done_row.last_error, None);
}
#[test]
fn test_session_join_row_ignores_partial_review_request_columns() {
let mut session_join_row = session_join_row_fixture();
session_join_row.review_request_last_refreshed_at = None;
let session_row = session_join_row.into_session_row();
assert_eq!(session_row.id, "session-a");
assert_eq!(session_row.project_id, Some(7));
assert_eq!(session_row.status, "Review");
assert_eq!(session_row.added_lines, 14);
assert_eq!(session_row.deleted_lines, 6);
assert_eq!(session_row.review_request, None);
}
#[test]
fn test_session_join_row_maps_review_request_columns() {
let session_join_row = session_join_row_fixture();
let session_row = session_join_row.into_session_row();
assert_eq!(session_row.id, "session-a");
assert_eq!(session_row.added_lines, 14);
assert_eq!(session_row.deleted_lines, 6);
assert_eq!(session_row.project_id, Some(7));
assert_eq!(
session_row.published_upstream_ref.as_deref(),
Some("origin/session-a")
);
assert_eq!(session_row.questions.as_deref(), Some("Question text"));
assert_eq!(session_row.summary.as_deref(), Some("Summary text"));
assert_eq!(session_row.title.as_deref(), Some("Review session"));
assert_review_request_row(&session_row);
}
#[tokio::test]
async fn test_upsert_session_usage_accumulates_counts_per_model() {
let database = Database::open_in_memory()
.await
.expect("failed to open in-memory db");
let project_id = database
.upsert_project("/tmp/project", Some("main"))
.await
.expect("failed to insert project");
insert_session_fixture(&database, "session-a", "main", "Done", project_id).await;
database
.upsert_session_usage(
"session-a",
"claude-opus-4.1",
&SessionStats {
added_lines: 0,
deleted_lines: 0,
input_tokens: 11,
output_tokens: 29,
},
)
.await
.expect("failed to insert first usage row");
database
.upsert_session_usage(
"session-a",
"claude-opus-4.1",
&SessionStats {
added_lines: 0,
deleted_lines: 0,
input_tokens: 3,
output_tokens: 5,
},
)
.await
.expect("failed to update existing usage row");
database
.upsert_session_usage("session-a", "ignored-model", &SessionStats::default())
.await
.expect("failed to ignore zero-usage update");
let usage_rows = database
.load_session_usage("session-a")
.await
.expect("failed to load session usage");
assert_eq!(usage_rows.len(), 1);
assert_eq!(usage_rows[0].model, "claude-opus-4.1");
assert_eq!(usage_rows[0].input_tokens, 14);
assert_eq!(usage_rows[0].invocation_count, 2);
assert_eq!(usage_rows[0].output_tokens, 34);
assert_eq!(usage_rows[0].session_id.as_deref(), Some("session-a"));
}
#[tokio::test]
async fn test_setting_round_trip_supports_default_smart_fast_and_review_models() {
let database = Database::open_in_memory()
.await
.expect("failed to open in-memory db");
database
.upsert_setting(
SettingName::DefaultSmartModel,
AgentModel::Gemini31ProPreview.as_str(),
)
.await
.expect("failed to persist default smart model");
database
.upsert_setting(SettingName::DefaultFastModel, AgentModel::Gpt54.as_str())
.await
.expect("failed to persist default fast model");
database
.upsert_setting(
SettingName::DefaultReviewModel,
AgentModel::ClaudeOpus46.as_str(),
)
.await
.expect("failed to persist default review model");
let default_smart_model = database
.get_setting(SettingName::DefaultSmartModel)
.await
.expect("failed to load default smart model");
let default_fast_model = database
.get_setting(SettingName::DefaultFastModel)
.await
.expect("failed to load default fast model");
let default_review_model = database
.get_setting(SettingName::DefaultReviewModel)
.await
.expect("failed to load default review model");
assert_eq!(
default_smart_model,
Some(AgentModel::Gemini31ProPreview.as_str().to_string())
);
assert_eq!(
default_fast_model,
Some(AgentModel::Gpt54.as_str().to_string())
);
assert_eq!(
default_review_model,
Some(AgentModel::ClaudeOpus46.as_str().to_string())
);
}
#[tokio::test]
async fn test_project_setting_round_trip_is_isolated_per_project() {
let database = Database::open_in_memory()
.await
.expect("failed to open in-memory db");
let first_project_id = database
.upsert_project("/tmp/project-a", Some("main"))
.await
.expect("failed to insert first project");
let second_project_id = database
.upsert_project("/tmp/project-b", Some("main"))
.await
.expect("failed to insert second project");
database
.upsert_project_setting(first_project_id, SettingName::OpenCommand, "npm run dev")
.await
.expect("failed to persist first project setting");
database
.upsert_project_setting(second_project_id, SettingName::OpenCommand, "cargo test")
.await
.expect("failed to persist second project setting");
let first_project_setting = database
.get_project_setting(first_project_id, SettingName::OpenCommand)
.await
.expect("failed to load first project setting");
let second_project_setting = database
.get_project_setting(second_project_id, SettingName::OpenCommand)
.await
.expect("failed to load second project setting");
assert_eq!(first_project_setting, Some("npm run dev".to_string()));
assert_eq!(second_project_setting, Some("cargo test".to_string()));
}
#[tokio::test]
async fn test_project_reasoning_level_round_trip_uses_typed_setting_helpers() {
let database = Database::open_in_memory()
.await
.expect("failed to open in-memory db");
let project_id = database
.upsert_project("/tmp/project", Some("main"))
.await
.expect("failed to insert project");
database
.set_project_reasoning_level(project_id, ReasoningLevel::Low)
.await
.expect("failed to persist project reasoning level");
let reasoning_level = database
.load_project_reasoning_level(project_id)
.await
.expect("failed to load project reasoning level");
assert_eq!(reasoning_level, ReasoningLevel::Low);
}
#[tokio::test]
async fn test_load_project_reasoning_level_defaults_when_setting_is_missing_or_invalid() {
let database = Database::open_in_memory()
.await
.expect("failed to open in-memory db");
let project_id = database
.upsert_project("/tmp/project", Some("main"))
.await
.expect("failed to insert project");
let missing_setting_level = database
.load_project_reasoning_level(project_id)
.await
.expect("failed to load default project reasoning level");
database
.upsert_project_setting(project_id, SettingName::ReasoningLevel, "unsupported")
.await
.expect("failed to insert unsupported project reasoning level");
let invalid_setting_level = database
.load_project_reasoning_level(project_id)
.await
.expect("failed to load fallback project reasoning level");
assert_eq!(missing_setting_level, ReasoningLevel::High);
assert_eq!(invalid_setting_level, ReasoningLevel::High);
}
#[tokio::test]
async fn test_reasoning_level_round_trip_uses_typed_setting_helpers() {
let database = Database::open_in_memory()
.await
.expect("failed to open in-memory db");
database
.set_reasoning_level(ReasoningLevel::Low)
.await
.expect("failed to persist reasoning level");
let reasoning_level = database
.load_reasoning_level()
.await
.expect("failed to load reasoning level");
assert_eq!(reasoning_level, ReasoningLevel::Low);
}
#[tokio::test]
async fn test_load_reasoning_level_defaults_when_setting_is_missing_or_invalid() {
let database = Database::open_in_memory()
.await
.expect("failed to open in-memory db");
let missing_setting_level = database
.load_reasoning_level()
.await
.expect("failed to load default reasoning level");
database
.upsert_setting(SettingName::ReasoningLevel, "unsupported")
.await
.expect("failed to insert unsupported reasoning level");
let invalid_setting_level = database
.load_reasoning_level()
.await
.expect("failed to load fallback reasoning level");
assert_eq!(missing_setting_level, ReasoningLevel::High);
assert_eq!(invalid_setting_level, ReasoningLevel::High);
}
#[tokio::test]
async fn test_session_provider_conversation_id_round_trip_and_clear() {
let database = Database::open_in_memory()
.await
.expect("failed to open in-memory db");
let project_id = database
.upsert_project("/tmp/project", None)
.await
.expect("failed to upsert project");
database
.insert_session("session-a", "gpt-5.4", "main", "Done", project_id)
.await
.expect("failed to insert session");
database
.update_session_provider_conversation_id("session-a", Some("thread-123"))
.await
.expect("failed to set provider conversation id");
let stored_id = database
.get_session_provider_conversation_id("session-a")
.await
.expect("failed to load provider conversation id");
database
.update_session_provider_conversation_id("session-a", None)
.await
.expect("failed to clear provider conversation id");
let cleared_id = database
.get_session_provider_conversation_id("session-a")
.await
.expect("failed to load cleared provider conversation id");
assert_eq!(stored_id, Some("thread-123".to_string()));
assert_eq!(cleared_id, None);
}
#[tokio::test]
async fn test_session_instruction_conversation_id_round_trip_and_clear() {
let database = Database::open_in_memory()
.await
.expect("failed to open in-memory db");
let project_id = database
.upsert_project("/tmp/project", None)
.await
.expect("failed to upsert project");
database
.insert_session("session-a", "gpt-5.4", "main", "Done", project_id)
.await
.expect("failed to insert session");
let instruction_conversation_id = Some("thread-123");
database
.update_session_instruction_conversation_id("session-a", instruction_conversation_id)
.await
.expect("failed to set instruction conversation id");
let stored_conversation_id = database
.get_session_instruction_conversation_id("session-a")
.await
.expect("failed to load instruction conversation id");
database
.update_session_instruction_conversation_id("session-a", None)
.await
.expect("failed to clear instruction conversation id");
let cleared_conversation_id = database
.get_session_instruction_conversation_id("session-a")
.await
.expect("failed to load cleared instruction conversation id");
assert_eq!(stored_conversation_id, Some("thread-123".to_string()));
assert_eq!(cleared_conversation_id, None);
}
#[tokio::test]
async fn test_session_published_upstream_ref_round_trip_and_clear() {
let database = Database::open_in_memory()
.await
.expect("failed to open in-memory db");
let project_id = database
.upsert_project("/tmp/project", None)
.await
.expect("failed to upsert project");
database
.insert_session("session-a", "gpt-5.4", "main", "Review", project_id)
.await
.expect("failed to insert session");
database
.update_session_published_upstream_ref("session-a", Some("origin/agentty/session-a"))
.await
.expect("failed to persist session published upstream ref");
let persisted_row = database
.load_sessions()
.await
.expect("failed to load sessions")
.into_iter()
.find(|row| row.id == "session-a")
.expect("missing persisted session row");
database
.update_session_published_upstream_ref("session-a", None)
.await
.expect("failed to clear session published upstream ref");
let cleared_row = database
.load_sessions()
.await
.expect("failed to load sessions after clearing")
.into_iter()
.find(|row| row.id == "session-a")
.expect("missing cleared session row");
assert_eq!(
persisted_row.published_upstream_ref.as_deref(),
Some("origin/agentty/session-a")
);
assert_eq!(cleared_row.published_upstream_ref, None);
}
#[tokio::test]
async fn test_load_session_published_upstream_ref_returns_stored_value() {
let database = Database::open_in_memory()
.await
.expect("failed to open in-memory db");
let project_id = database
.upsert_project("/tmp/project", None)
.await
.expect("failed to upsert project");
database
.insert_session("session-load", "gpt-5.4", "main", "Review", project_id)
.await
.expect("failed to insert session");
database
.update_session_published_upstream_ref(
"session-load",
Some("origin/agentty/session-load"),
)
.await
.expect("failed to set published upstream ref");
let loaded_ref = database
.load_session_published_upstream_ref("session-load")
.await
.expect("failed to load published upstream ref");
assert_eq!(loaded_ref.as_deref(), Some("origin/agentty/session-load"));
}
#[tokio::test]
async fn test_load_session_published_upstream_ref_returns_none_when_unset() {
let database = Database::open_in_memory()
.await
.expect("failed to open in-memory db");
let project_id = database
.upsert_project("/tmp/project", None)
.await
.expect("failed to upsert project");
database
.insert_session("session-unset", "gpt-5.4", "main", "Review", project_id)
.await
.expect("failed to insert session");
let loaded_ref = database
.load_session_published_upstream_ref("session-unset")
.await
.expect("failed to load published upstream ref");
assert_eq!(loaded_ref, None);
}
#[tokio::test]
async fn test_load_session_published_upstream_ref_returns_none_for_missing_session() {
let database = Database::open_in_memory()
.await
.expect("failed to open in-memory db");
let loaded_ref = database
.load_session_published_upstream_ref("nonexistent")
.await
.expect("failed to load published upstream ref");
assert_eq!(loaded_ref, None);
}
#[tokio::test]
async fn test_session_review_request_round_trip_and_clear() {
let database = Database::open_in_memory()
.await
.expect("failed to open in-memory db");
let project_id = database
.upsert_project("/tmp/project", None)
.await
.expect("failed to upsert project");
database
.insert_session("session-a", "gpt-5.4", "main", "Review", project_id)
.await
.expect("failed to insert session");
let review_request = review_request_fixture();
database
.update_session_review_request("session-a", Some(&review_request))
.await
.expect("failed to persist session review request");
let persisted_row = database
.load_sessions()
.await
.expect("failed to load sessions")
.into_iter()
.find(|row| row.id == "session-a")
.expect("missing persisted session row");
database
.update_session_review_request("session-a", None)
.await
.expect("failed to clear session review request");
let cleared_row = database
.load_sessions()
.await
.expect("failed to load sessions after clearing")
.into_iter()
.find(|row| row.id == "session-a")
.expect("missing cleared session row");
assert_review_request_row(&persisted_row);
assert_eq!(cleared_row.review_request, None);
}
#[tokio::test]
async fn test_insert_session_creation_activity_at_persists_timestamp() {
let database = Database::open_in_memory()
.await
.expect("failed to open in-memory db");
let project_id = database
.upsert_project("/tmp/project", None)
.await
.expect("failed to upsert project");
database
.insert_session("session-a", "gpt-5.4", "main", "Done", project_id)
.await
.expect("failed to insert session");
database
.insert_session_creation_activity_at("session-a", 123)
.await
.expect("failed to persist activity event");
let activity_timestamps = database
.load_session_activity_timestamps()
.await
.expect("failed to load activity timestamps");
assert_eq!(activity_timestamps, vec![123]);
}
#[tokio::test]
async fn test_insert_session_creation_activity_at_ignores_duplicates_per_session() {
let database = Database::open_in_memory()
.await
.expect("failed to open in-memory db");
let project_id = database
.upsert_project("/tmp/project", None)
.await
.expect("failed to upsert project");
database
.insert_session("session-a", "gpt-5.4", "main", "Done", project_id)
.await
.expect("failed to insert session");
database
.insert_session_creation_activity_at("session-a", 100)
.await
.expect("failed to persist first activity event");
database
.insert_session_creation_activity_at("session-a", 200)
.await
.expect("failed to persist duplicate activity event");
let activity_timestamps = database
.load_session_activity_timestamps()
.await
.expect("failed to load activity timestamps");
assert_eq!(activity_timestamps, vec![100]);
}
#[tokio::test]
async fn test_load_session_activity_timestamps_keeps_deleted_session_history() {
let database = Database::open_in_memory()
.await
.expect("failed to open in-memory db");
let project_id = database
.upsert_project("/tmp/project", None)
.await
.expect("failed to upsert project");
database
.insert_session("session-a", "gpt-5.4", "main", "Done", project_id)
.await
.expect("failed to insert first session");
database
.insert_session_creation_activity_at("session-a", 100)
.await
.expect("failed to persist first activity event");
database
.insert_session("session-b", "gpt-5.4", "main", "Done", project_id)
.await
.expect("failed to insert second session");
database
.insert_session_creation_activity_at("session-b", 200)
.await
.expect("failed to persist second activity event");
database
.delete_session("session-a")
.await
.expect("failed to delete first session");
let activity_timestamps = database
.load_session_activity_timestamps()
.await
.expect("failed to load activity timestamps");
assert_eq!(activity_timestamps, vec![100, 200]);
}
#[tokio::test]
async fn test_load_session_activity_groups_counts_by_local_day() {
let database = Database::open_in_memory()
.await
.expect("failed to open in-memory db");
let project_id = database
.upsert_project("/tmp/project", None)
.await
.expect("failed to upsert project");
database
.insert_session("session-a", "gpt-5.4", "main", "Done", project_id)
.await
.expect("failed to insert first session");
database
.insert_session("session-b", "gpt-5.4", "main", "Done", project_id)
.await
.expect("failed to insert second session");
database
.insert_session("session-c", "gpt-5.4", "main", "Done", project_id)
.await
.expect("failed to insert third session");
let first_day_timestamp = 10 * 86_400 + 10;
let second_timestamp_same_day = 10 * 86_400 + 600;
let second_day_timestamp = 11 * 86_400 + 50;
database
.clear_session_activity()
.await
.expect("failed to clear session activity");
database
.insert_session_creation_activity_at("session-a", first_day_timestamp)
.await
.expect("failed to persist first activity event");
database
.insert_session_creation_activity_at("session-b", second_timestamp_same_day)
.await
.expect("failed to persist second activity event");
database
.insert_session_creation_activity_at("session-c", second_day_timestamp)
.await
.expect("failed to persist third activity event");
let expected_activity = vec![
DailyActivity {
day_key: local_day_key(first_day_timestamp),
session_count: 2,
},
DailyActivity {
day_key: local_day_key(second_day_timestamp),
session_count: 1,
},
];
let activity = database
.load_session_activity()
.await
.expect("failed to load aggregated session activity");
assert_eq!(activity, expected_activity);
}
#[tokio::test]
async fn test_load_projects_with_stats_returns_session_counts_and_last_update() {
let database = Database::open_in_memory()
.await
.expect("failed to open in-memory db");
let project_id = database
.upsert_project("/tmp/project", Some("main"))
.await
.expect("failed to upsert project");
database
.insert_session("session-a", "gpt-5.4", "main", "Done", project_id)
.await
.expect("failed to insert session-a");
database
.insert_session("session-b", "gpt-5.4", "main", "Done", project_id)
.await
.expect("failed to insert session-b");
let projects = database
.load_projects_with_stats()
.await
.expect("failed to load projects with stats");
assert_eq!(projects.len(), 1);
assert_eq!(projects[0].session_count, 2);
assert!(projects[0].last_session_updated_at.is_some());
}
fn local_day_key(timestamp_seconds: i64) -> i64 {
let utc_timestamp = time::OffsetDateTime::from_unix_timestamp(timestamp_seconds)
.expect("timestamp should be valid for test fixture");
let local_offset = time::UtcOffset::local_offset_at(utc_timestamp)
.expect("local offset should resolve for test fixture");
timestamp_seconds
.saturating_add(i64::from(local_offset.whole_seconds()))
.div_euclid(86_400)
}
#[test]
fn test_load_session_activity_matches_rust_grouping_across_dst_transition() {
if !cfg!(unix) {
return;
}
let current_test_binary = env::current_exe().expect("failed to resolve current test bin");
let output = Command::new(current_test_binary)
.env(DST_TEST_SUBPROCESS_ENV, "1")
.env("TZ", "America/Los_Angeles")
.arg(
"test_load_session_activity_matches_rust_grouping_across_dst_transition_subprocess",
)
.arg("--exact")
.arg("--test-threads=1")
.output()
.expect("failed to run DST subprocess test");
assert!(
output.status.success(),
"DST subprocess test failed.\nstdout:\n{}\nstderr:\n{}",
String::from_utf8_lossy(&output.stdout),
String::from_utf8_lossy(&output.stderr)
);
}
#[tokio::test]
async fn test_load_session_activity_matches_rust_grouping_across_dst_transition_subprocess() {
if !cfg!(unix) || env::var_os(DST_TEST_SUBPROCESS_ENV).is_none() {
return;
}
let database = Database::open_in_memory()
.await
.expect("failed to open in-memory db");
let project_id = database
.upsert_project("/tmp/project", None)
.await
.expect("failed to upsert project");
database
.insert_session("session-a", "gpt-5.4", "main", "Done", project_id)
.await
.expect("failed to insert first session");
database
.insert_session("session-b", "gpt-5.4", "main", "Done", project_id)
.await
.expect("failed to insert second session");
database
.insert_session("session-c", "gpt-5.4", "main", "Done", project_id)
.await
.expect("failed to insert third session");
database
.clear_session_activity()
.await
.expect("failed to clear activity history");
let before_dst_jump = 1_710_063_000_i64;
let after_dst_jump = 1_710_066_600_i64;
let next_local_day = 1_710_142_200_i64;
database
.insert_session_creation_activity_at("session-a", before_dst_jump)
.await
.expect("failed to persist pre-DST activity");
database
.insert_session_creation_activity_at("session-b", after_dst_jump)
.await
.expect("failed to persist post-DST activity");
database
.insert_session_creation_activity_at("session-c", next_local_day)
.await
.expect("failed to persist next-day activity");
let first_day_key = local_day_key(before_dst_jump);
let second_day_key = local_day_key(after_dst_jump);
let third_day_key = local_day_key(next_local_day);
let activity = database
.load_session_activity()
.await
.expect("failed to load grouped session activity");
assert_eq!(first_day_key, second_day_key);
assert_ne!(second_day_key, third_day_key);
assert_eq!(
activity,
vec![
DailyActivity {
day_key: first_day_key,
session_count: 2,
},
DailyActivity {
day_key: third_day_key,
session_count: 1,
},
]
);
}
#[tokio::test]
async fn test_set_and_load_active_project_id_round_trip() {
let database = Database::open_in_memory()
.await
.expect("failed to open in-memory db");
let project_id = database
.upsert_project("/tmp/project", Some("main"))
.await
.expect("failed to upsert project");
database
.set_active_project_id(project_id)
.await
.expect("failed to persist active project id");
let active_project_id = database
.load_active_project_id()
.await
.expect("failed to load active project id");
assert_eq!(active_project_id, Some(project_id));
}
#[tokio::test]
async fn test_load_session_project_id_returns_associated_project() {
let database = Database::open_in_memory()
.await
.expect("failed to open in-memory db");
let project_id = database
.upsert_project("/tmp/project", Some("main"))
.await
.expect("failed to insert project");
database
.insert_session("session-a", "gpt-5.4", "main", "Done", project_id)
.await
.expect("failed to insert session");
let loaded_project_id = database
.load_session_project_id("session-a")
.await
.expect("failed to load session project id");
assert_eq!(loaded_project_id, Some(project_id));
}
#[tokio::test]
async fn test_load_session_summary_returns_persisted_summary() {
let database = Database::open_in_memory()
.await
.expect("failed to open in-memory db");
let project_id = database
.upsert_project("/tmp/project", Some("main"))
.await
.expect("failed to insert project");
database
.insert_session("session-a", "gpt-5.4", "main", "Done", project_id)
.await
.expect("failed to insert session");
database
.update_session_summary("session-a", "persisted summary")
.await
.expect("failed to update session summary");
let loaded_summary = database
.load_session_summary("session-a")
.await
.expect("failed to load session summary");
assert_eq!(loaded_summary.as_deref(), Some("persisted summary"));
}
#[tokio::test]
async fn test_replace_session_follow_up_tasks_round_trips_latest_tasks() {
let database = Database::open_in_memory()
.await
.expect("failed to open in-memory db");
let project_id = database
.upsert_project("/tmp/project", Some("main"))
.await
.expect("failed to insert project");
database
.insert_session("session-a", "gpt-5.4", "main", "Done", project_id)
.await
.expect("failed to insert session");
database
.replace_session_follow_up_tasks(
"session-a",
&["Stale task".to_string(), "Remove me".to_string()],
)
.await
.expect("failed to insert initial follow-up tasks");
database
.replace_session_follow_up_tasks(
"session-a",
&[
"Document the release note.".to_string(),
"Add integration coverage.".to_string(),
],
)
.await
.expect("failed to replace follow-up tasks");
let follow_up_tasks = database
.load_session_follow_up_tasks()
.await
.expect("failed to load session follow-up tasks");
let follow_up_task_text = follow_up_tasks
.into_iter()
.filter(|task| task.session_id == "session-a")
.map(|task| task.text)
.collect::<Vec<_>>();
assert_eq!(
follow_up_task_text,
vec![
"Document the release note.".to_string(),
"Add integration coverage.".to_string()
]
);
}
#[tokio::test]
async fn test_persist_session_turn_metadata_rolls_back_on_failure() {
let database = Database::open_in_memory()
.await
.expect("failed to open in-memory db");
let project_id = database
.upsert_project("/tmp/project", Some("main"))
.await
.expect("failed to insert project");
database
.insert_session("session-a", "gpt-5.4", "main", "Review", project_id)
.await
.expect("failed to insert session");
database
.update_session_summary("session-a", "persisted summary")
.await
.expect("failed to seed summary");
sqlx::query("DROP TABLE session_follow_up_task")
.execute(database.pool())
.await
.expect("failed to drop follow-up-task table");
let result = database
.persist_session_turn_metadata(
"session-a",
&SessionTurnMetadata {
follow_up_tasks: &["Document the failure path.".to_string()],
instruction_conversation_id: Some("instruction-thread"),
model: AgentModel::Gpt54.as_str(),
provider_conversation_id: Some("thread-123"),
questions_json: r#"[{"text":"Need tests?"}]"#,
summary: r#"{"turn":"Updated the worker.","session":"Session state changed."}"#,
token_usage_delta: &SessionStats {
added_lines: 0,
deleted_lines: 0,
input_tokens: 3,
output_tokens: 5,
},
},
)
.await;
let session = database
.load_sessions()
.await
.expect("failed to reload sessions")
.into_iter()
.find(|session| session.id == "session-a")
.expect("expected seeded session");
let provider_conversation_id = database
.get_session_provider_conversation_id("session-a")
.await
.expect("failed to load provider conversation id");
assert!(matches!(result, Err(DbError::Query(_))));
assert_eq!(session.summary.as_deref(), Some("persisted summary"));
assert_eq!(session.questions.as_deref(), None);
assert_eq!(session.input_tokens, 0);
assert_eq!(session.output_tokens, 0);
assert_eq!(provider_conversation_id.as_deref(), None);
}
#[tokio::test]
async fn test_update_session_follow_up_task_launched_session_id_round_trips() {
let database = Database::open_in_memory()
.await
.expect("failed to open in-memory db");
let project_id = database
.upsert_project("/tmp/project", Some("main"))
.await
.expect("failed to insert project");
database
.insert_session("session-a", "gpt-5.4", "main", "Done", project_id)
.await
.expect("failed to insert source session");
database
.insert_session("session-b", "gpt-5.4", "main", "New", project_id)
.await
.expect("failed to insert sibling session");
database
.replace_session_follow_up_tasks("session-a", &["Launch the sibling task.".to_string()])
.await
.expect("failed to insert follow-up task");
database
.update_session_follow_up_task_launched_session_id("session-a", 0, Some("session-b"))
.await
.expect("failed to persist launched sibling-session id");
let follow_up_tasks = database
.load_session_follow_up_tasks()
.await
.expect("failed to load follow-up tasks");
let follow_up_task = follow_up_tasks
.into_iter()
.find(|task| task.session_id == "session-a")
.expect("expected persisted follow-up task");
assert_eq!(
follow_up_task.launched_session_id.as_deref(),
Some("session-b")
);
assert_eq!(follow_up_task.position, 0);
}
#[tokio::test]
async fn test_set_project_favorite_updates_project_state() {
let database = Database::open_in_memory()
.await
.expect("failed to open in-memory db");
let project_id = database
.upsert_project("/tmp/project", Some("main"))
.await
.expect("failed to upsert project");
database
.set_project_favorite(project_id, true)
.await
.expect("failed to set project favorite");
let project = database
.get_project(project_id)
.await
.expect("failed to load project")
.expect("expected existing project");
assert!(project.is_favorite);
}
#[tokio::test]
async fn query_on_dropped_table_returns_db_error_query() {
let database = Database::open_in_memory()
.await
.expect("failed to open database");
sqlx::query("DROP TABLE session")
.execute(database.pool())
.await
.expect("failed to drop table");
let result = database.load_sessions_metadata().await;
assert!(
matches!(result, Err(DbError::Query(_))),
"expected DbError::Query variant"
);
}
#[tokio::test]
async fn db_error_display_includes_underlying_message() {
let database = Database::open_in_memory()
.await
.expect("failed to open database");
sqlx::query("DROP TABLE session")
.execute(database.pool())
.await
.expect("failed to drop table");
let result = database.load_sessions_metadata().await;
let error = result.expect_err("expected query on dropped table to fail");
let display_text = error.to_string();
assert!(
!display_text.is_empty(),
"DbError Display should produce a non-empty message"
);
}
#[tokio::test]
async fn open_with_unwritable_parent_returns_db_error_io() {
let temp = tempdir().expect("failed to create temp directory");
let blocking_file = temp.path().join("not_a_dir");
std::fs::write(&blocking_file, b"").expect("failed to create blocking file");
let db_path = blocking_file.join("nested").join("db.sqlite");
let result = Database::open(&db_path).await;
assert!(
matches!(result, Err(DbError::Io(_))),
"expected DbError::Io variant"
);
}
}