use std::collections::{HashMap, HashSet};
use std::path::{Path, PathBuf};
use std::sync::Arc;
use std::time::{Duration, Instant, SystemTime, UNIX_EPOCH};
use ratatui::widgets::TableState;
use tokio::task::{JoinHandle, JoinSet};
use super::workflow::merge::SessionMergeService;
pub(crate) use super::workflow::merge::{SyncMainOutcome, SyncSessionStartError};
pub(crate) use super::workflow::task::{RunAgentAssistTaskInput, SessionTaskService};
use super::workflow::worker::SessionWorkerService;
use crate::app::session_state::SessionGitStatus;
use crate::app::{AppServices, SessionState, setting};
use crate::domain::agent::{AgentModel, ReasoningLevel};
use crate::domain::file_entry::FileEntry;
use crate::domain::question::QuestionItem;
use crate::domain::session::{
DailyActivity, FollowUpTaskAction, PublishedBranchSyncStatus, ReviewRequest, Session,
SessionFollowUpTask, SessionId, SessionStats,
};
use crate::infra::git;
type SessionRenderParts<'a> = (&'a [Session], &'a [DailyActivity], &'a mut TableState);
pub(crate) const SESSION_REFRESH_INTERVAL: Duration = Duration::from_secs(5);
pub(crate) const AT_MENTION_INDEX_TTL: Duration = Duration::from_secs(30);
#[derive(Clone, Copy)]
pub(crate) struct SessionDefaults {
pub(crate) model: AgentModel,
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub(crate) struct TurnAppliedState {
pub(crate) follow_up_tasks: Vec<SessionFollowUpTask>,
pub(crate) questions: Vec<QuestionItem>,
pub(crate) summary: Option<String>,
pub(crate) token_usage_delta: SessionStats,
}
impl TurnAppliedState {
pub(crate) fn merge_newer(&mut self, newer_turn_applied_state: Self) {
self.follow_up_tasks = newer_turn_applied_state.follow_up_tasks;
self.questions = newer_turn_applied_state.questions;
self.summary = newer_turn_applied_state.summary;
self.token_usage_delta.input_tokens = self
.token_usage_delta
.input_tokens
.saturating_add(newer_turn_applied_state.token_usage_delta.input_tokens);
self.token_usage_delta.output_tokens = self
.token_usage_delta
.output_tokens
.saturating_add(newer_turn_applied_state.token_usage_delta.output_tokens);
}
}
pub(crate) use crate::infra::clock::Clock;
pub struct SessionManager {
pub(super) active_prompt_outputs: HashMap<SessionId, String>,
at_mention_indexes: HashMap<PathBuf, AtMentionIndex>,
pub(super) default_session_model: AgentModel,
pub(super) git_client: Arc<dyn git::GitClient>,
pub(super) merge_service: SessionMergeService,
pub(super) pending_history_replay: HashSet<SessionId>,
pub(super) published_branch_sync_operations: HashMap<SessionId, String>,
pub(super) state: SessionState,
pub(super) stats_activity: Vec<DailyActivity>,
pub(super) title_generation_tasks: HashMap<SessionId, TitleGenerationTask>,
pub(super) worker_service: SessionWorkerService,
}
pub(crate) struct TitleGenerationTask {
generation: u64,
join_handle: JoinHandle<()>,
}
#[derive(Debug)]
struct AtMentionIndex {
created_at: Instant,
entries: Vec<FileEntry>,
}
impl SessionManager {
pub(crate) fn new(
defaults: SessionDefaults,
git_client: Arc<dyn git::GitClient>,
state: SessionState,
stats_activity: Vec<DailyActivity>,
) -> Self {
let pending_history_replay = Self::startup_history_replay_set(&state.sessions);
Self {
active_prompt_outputs: HashMap::new(),
at_mention_indexes: HashMap::new(),
default_session_model: defaults.model,
git_client,
merge_service: SessionMergeService,
pending_history_replay,
published_branch_sync_operations: HashMap::new(),
state,
stats_activity,
title_generation_tasks: HashMap::new(),
worker_service: SessionWorkerService::new(),
}
}
pub(crate) fn replace_title_generation_task(
&mut self,
session_id: &str,
generation: u64,
title_generation_task: JoinHandle<()>,
) {
if let Some(existing_task) = self.title_generation_tasks.remove(session_id)
&& !existing_task.join_handle.is_finished()
{
existing_task.join_handle.abort();
}
self.title_generation_tasks.insert(
SessionId::from(session_id),
TitleGenerationTask {
generation,
join_handle: title_generation_task,
},
);
}
pub(crate) fn abort_title_generation_task(&mut self, session_id: &str) {
if let Some(existing_task) = self.title_generation_tasks.remove(session_id)
&& !existing_task.join_handle.is_finished()
{
existing_task.join_handle.abort();
}
}
pub(crate) fn next_title_generation_task_generation(&self, session_id: &str) -> u64 {
self.title_generation_tasks
.get(session_id)
.map_or(1, |tracked_task| tracked_task.generation.saturating_add(1))
}
pub(crate) fn clear_title_generation_task_if_matches(
&mut self,
session_id: &str,
generation: u64,
) {
let should_clear = self
.title_generation_tasks
.get(session_id)
.is_some_and(|tracked_task| tracked_task.generation == generation);
if should_clear {
self.title_generation_tasks.remove(session_id);
}
}
pub(crate) fn merge_service(&self) -> &SessionMergeService {
&self.merge_service
}
pub(crate) fn git_client(&self) -> Arc<dyn git::GitClient> {
Arc::clone(&self.git_client)
}
pub(crate) fn worker_service_mut(&mut self) -> &mut SessionWorkerService {
&mut self.worker_service
}
pub(crate) fn default_session_model(&self) -> AgentModel {
self.default_session_model
}
pub(crate) fn set_default_session_model(&mut self, session_model: AgentModel) {
self.default_session_model = session_model;
}
pub(crate) async fn load_default_session_model(
services: &AppServices,
project_id: Option<i64>,
fallback_model: AgentModel,
) -> AgentModel {
setting::load_default_smart_model_setting(services, project_id, fallback_model).await
}
pub(crate) fn render_parts(&mut self) -> SessionRenderParts<'_> {
(
&self.state.sessions,
&self.stats_activity,
&mut self.state.table_state,
)
}
pub(crate) fn sessions(&self) -> &[Session] {
&self.state.sessions
}
pub(crate) fn sessions_mut(&mut self) -> &mut [Session] {
&mut self.state.sessions
}
#[cfg(test)]
pub(crate) fn push_session(&mut self, session: Session) {
self.state.push_session(session);
}
pub(crate) fn remove_session_at(&mut self, session_index: usize) -> Option<Session> {
self.state.remove_session_at(session_index)
}
pub(crate) fn selected_session_index(&self) -> Option<usize> {
self.state.table_state.selected()
}
pub(crate) fn select_session_index(&mut self, session_index: Option<usize>) {
self.state.table_state.select(session_index);
}
pub(crate) fn session_at_mut(&mut self, session_index: usize) -> Option<&mut Session> {
self.state.sessions.get_mut(session_index)
}
pub(crate) fn session_for_id(&self, session_id: &str) -> Option<&Session> {
self.state.session_for_id(session_id)
}
#[cfg(test)]
pub(crate) fn sync_from_handles(&mut self) {
self.state.sync_from_handles();
}
pub(crate) fn sync_session_from_handle(&mut self, session_id: &str) {
self.state.sync_session_from_handle(session_id);
}
pub(crate) fn apply_session_size_updated(
&mut self,
session_id: &str,
added_lines: u64,
deleted_lines: u64,
session_size: crate::domain::session::SessionSize,
) {
self.state
.apply_session_size_updated(session_id, added_lines, deleted_lines, session_size);
}
pub(crate) fn session_handles(
&self,
) -> &HashMap<SessionId, crate::domain::session::SessionHandles> {
&self.state.handles
}
#[cfg(test)]
pub(crate) fn session_handles_mut(
&mut self,
) -> &mut HashMap<SessionId, crate::domain::session::SessionHandles> {
&mut self.state.handles
}
pub(crate) fn active_prompt_outputs(&self) -> &HashMap<SessionId, String> {
&self.active_prompt_outputs
}
pub(crate) fn state(&self) -> &SessionState {
&self.state
}
pub(crate) fn state_mut(&mut self) -> &mut SessionState {
&mut self.state
}
pub(crate) fn apply_session_model_updated(
&mut self,
session_id: &str,
session_model: AgentModel,
) {
if let Some(session) = self
.state
.sessions
.iter_mut()
.find(|session| session.id == session_id)
{
session.model = session_model;
}
}
pub(crate) fn apply_session_reasoning_level_updated(
&mut self,
session_id: &str,
reasoning_level_override: Option<ReasoningLevel>,
) {
if let Some(session) = self
.state
.sessions
.iter_mut()
.find(|session| session.id == session_id)
{
session.reasoning_level_override = reasoning_level_override;
}
}
pub(crate) fn apply_published_upstream_ref(
&mut self,
session_id: &str,
published_upstream_ref: String,
) {
if let Some(session) = self
.state
.sessions
.iter_mut()
.find(|session| session.id == session_id)
{
session.published_upstream_ref = Some(published_upstream_ref);
}
}
pub(crate) fn apply_review_request(&mut self, session_id: &str, review_request: ReviewRequest) {
if let Some(session) = self
.state
.sessions
.iter_mut()
.find(|session| session.id == session_id)
{
session.review_request = Some(review_request);
}
}
pub(crate) fn start_published_branch_sync(
&mut self,
session_id: &str,
sync_operation_id: String,
) {
self.published_branch_sync_operations
.insert(SessionId::from(session_id), sync_operation_id);
if let Some(session) = self
.state
.sessions
.iter_mut()
.find(|session| session.id == session_id)
{
session.published_branch_sync_status = PublishedBranchSyncStatus::InProgress;
}
}
pub(crate) fn finish_published_branch_sync(
&mut self,
session_id: &str,
sync_operation_id: &str,
sync_status: PublishedBranchSyncStatus,
) {
let Some(current_operation_id) = self.published_branch_sync_operations.get(session_id)
else {
return;
};
if current_operation_id != sync_operation_id {
return;
}
self.published_branch_sync_operations.remove(session_id);
if let Some(session) = self
.state
.sessions
.iter_mut()
.find(|session| session.id == session_id)
{
session.published_branch_sync_status = sync_status;
}
}
pub(crate) fn append_workflow_notice(&mut self, session_id: &str, notice: String) {
if let Some(session) = self
.state
.sessions
.iter_mut()
.find(|session| session.id == session_id)
{
if let Some(workflow_notice) = &mut session.workflow_notice {
workflow_notice.push_str("\n\n");
workflow_notice.push_str(¬ice);
} else {
session.workflow_notice = Some(notice);
}
}
}
pub(crate) fn apply_turn_applied_state(
&mut self,
session_id: &str,
turn_applied_state: &TurnAppliedState,
) {
let Some(session) = self
.state
.sessions
.iter_mut()
.find(|session| session.id == session_id)
else {
return;
};
session
.follow_up_tasks
.clone_from(&turn_applied_state.follow_up_tasks);
session.questions.clone_from(&turn_applied_state.questions);
session.summary.clone_from(&turn_applied_state.summary);
session.stats.input_tokens = session
.stats
.input_tokens
.saturating_add(turn_applied_state.token_usage_delta.input_tokens);
session.stats.output_tokens = session
.stats
.output_tokens
.saturating_add(turn_applied_state.token_usage_delta.output_tokens);
self.active_prompt_outputs.remove(session_id);
}
pub(crate) fn set_active_prompt_output(&mut self, session_id: &str, prompt_output: String) {
self.active_prompt_outputs
.insert(SessionId::from(session_id), prompt_output);
}
pub(crate) fn at_mention_index_for_root(
&mut self,
lookup_root: &Path,
) -> Option<Vec<FileEntry>> {
let cached_index = self.at_mention_indexes.get(lookup_root)?;
if self
.state
.clock
.now_instant()
.saturating_duration_since(cached_index.created_at)
> AT_MENTION_INDEX_TTL
{
self.at_mention_indexes.remove(lookup_root);
return None;
}
Some(cached_index.entries.clone())
}
pub(crate) fn set_at_mention_index_for_root(
&mut self,
lookup_root: PathBuf,
entries: Vec<FileEntry>,
) {
self.at_mention_indexes.insert(
lookup_root,
AtMentionIndex {
created_at: self.state.clock.now_instant(),
entries,
},
);
}
pub(crate) fn remove_at_mention_index_for_root(&mut self, lookup_root: &Path) {
self.at_mention_indexes.remove(lookup_root);
}
pub(crate) fn retain_active_prompt_outputs(&mut self) {
self.active_prompt_outputs.retain(|session_id, _| {
self.state
.sessions
.iter()
.find(|session| session.id == *session_id)
.is_some_and(|session| {
matches!(
session.status,
crate::domain::session::Status::InProgress
| crate::domain::session::Status::Queued
| crate::domain::session::Status::Rebasing
| crate::domain::session::Status::Merging
)
})
});
self.prune_expired_at_mention_indexes();
}
pub(crate) fn replace_session_git_statuses(
&mut self,
session_git_statuses: HashMap<SessionId, SessionGitStatus>,
) {
self.state
.replace_session_git_statuses(session_git_statuses);
}
pub(crate) fn session_git_statuses(&self) -> &HashMap<SessionId, SessionGitStatus> {
&self.state.session_git_statuses
}
fn prune_expired_at_mention_indexes(&mut self) {
let now = self.state.clock.now_instant();
self.at_mention_indexes.retain(|_, cached_index| {
now.saturating_duration_since(cached_index.created_at) <= AT_MENTION_INDEX_TTL
});
}
pub(crate) fn replace_session_worktree_availability(
&mut self,
session_worktree_availability: HashMap<SessionId, bool>,
) {
self.state
.replace_session_worktree_availability(session_worktree_availability);
}
pub(crate) fn session_worktree_availability(&self) -> &HashMap<SessionId, bool> {
&self.state.session_worktree_availability
}
pub(crate) fn replace_session_branch_names(
&mut self,
session_branch_names: HashMap<SessionId, String>,
) {
self.state
.replace_session_branch_names(session_branch_names);
}
pub(crate) fn session_branch_names(&self) -> &HashMap<SessionId, String> {
&self.state.session_branch_names
}
pub(crate) fn session_branch_name(&self, session_id: &str) -> Option<&str> {
self.state
.session_branch_names
.get(session_id)
.map(String::as_str)
}
pub(crate) fn set_session_worktree_available(&mut self, session_id: &str, is_available: bool) {
self.state
.set_session_worktree_available(session_id, is_available);
}
pub(crate) fn remove_session_worktree_availability(&mut self, session_id: &str) {
self.state.remove_session_worktree_availability(session_id);
}
pub(crate) async fn refresh_session_branch_names(&mut self) {
let session_inputs = self
.state
.sessions
.iter()
.map(|session| {
let session_id = session.id.clone();
let default_branch_name = session_branch(&session_id);
(session_id, session.folder.clone(), default_branch_name)
})
.collect::<Vec<_>>();
let mut session_branch_names = session_inputs
.iter()
.map(|(session_id, _, default_branch_name)| {
(session_id.clone(), default_branch_name.clone())
})
.collect::<HashMap<_, _>>();
let mut branch_detection_tasks = JoinSet::new();
for (session_id, session_folder, default_branch_name) in session_inputs {
let git_client = Arc::clone(&self.git_client);
branch_detection_tasks.spawn(async move {
let branch_name = git_client
.detect_git_info(session_folder)
.await
.unwrap_or(default_branch_name);
(session_id, branch_name)
});
}
while let Some(branch_detection_result) = branch_detection_tasks.join_next().await {
let Ok((session_id, branch_name)) = branch_detection_result else {
continue;
};
session_branch_names.insert(session_id, branch_name);
}
self.replace_session_branch_names(session_branch_names);
}
pub(crate) fn selected_follow_up_task_position(&self, session_id: &str) -> Option<usize> {
self.state.selected_follow_up_task_position(session_id)
}
pub(crate) fn selected_follow_up_task_action(
&self,
session_id: &str,
) -> Option<FollowUpTaskAction> {
let position = self.selected_follow_up_task_position(session_id)?;
let session = self
.state
.sessions
.iter()
.find(|session| session.id == session_id)?;
session
.follow_up_task(position)
.map(crate::domain::session::SessionFollowUpTask::action)
}
pub(crate) fn has_multiple_follow_up_tasks(&self, session_id: &str) -> bool {
self.state
.sessions
.iter()
.find(|session| session.id == session_id)
.is_some_and(|session| session.follow_up_tasks.len() > 1)
}
pub(crate) fn select_next_follow_up_task(&mut self, session_id: &str) {
self.state.select_next_follow_up_task(session_id);
}
pub(crate) fn select_previous_follow_up_task(&mut self, session_id: &str) {
self.state.select_previous_follow_up_task(session_id);
}
pub(crate) fn set_follow_up_task_launched_session_id(
&mut self,
session_id: &str,
position: usize,
launched_session_id: Option<SessionId>,
) {
self.state.set_follow_up_task_launched_session_id(
session_id,
position,
launched_session_id,
);
}
}
const SESSION_BRANCH_PREFIX: &str = "wt/";
pub(crate) fn session_folder(base: &Path, session_id: &str) -> PathBuf {
let len = session_id.len().min(8);
base.join(&session_id[..len])
}
pub(crate) fn session_branch(session_id: &str) -> String {
let len = session_id.len().min(8);
format!("{SESSION_BRANCH_PREFIX}{}", &session_id[..len])
}
pub(crate) fn remote_branch_name_from_upstream_ref(upstream_ref: &str) -> String {
upstream_ref.split_once('/').map_or_else(
|| upstream_ref.to_string(),
|(_, branch_name)| branch_name.to_string(),
)
}
pub(crate) fn unix_timestamp_from_system_time(system_time: SystemTime) -> i64 {
system_time
.duration_since(UNIX_EPOCH)
.map_or(0, |duration| i64::try_from(duration.as_secs()).unwrap_or(0))
}
#[cfg(test)]
mod tests {
use std::collections::{BTreeMap, HashMap};
use std::path::{Path, PathBuf};
use std::process::{Command, Stdio};
use std::sync::{Arc, Mutex};
use std::time::{Duration, Instant, SystemTime};
use tempfile::tempdir;
use tokio::sync::Barrier;
use super::*;
use crate::app::{App, AppEvent, SyncSessionStartError, Tab};
use crate::domain::agent::{AgentKind, AgentModel, ReasoningLevel};
use crate::domain::session::{
DailyActivity, SESSION_DATA_DIR, Session, SessionHandles, SessionSize, SessionStats, Status,
};
use crate::domain::setting::SettingName;
use crate::infra::agent::AgentResponse;
use crate::infra::agent::tests::MockAgentBackend;
use crate::infra::app_server::{AppServerTurnResponse, MockAppServerClient};
use crate::infra::channel::{
AgentRequestKind, MockAgentChannel, TurnPrompt, TurnPromptAttachment, TurnPromptTextSource,
TurnResult,
};
use crate::infra::clock::RealClock;
use crate::infra::db::AppRepositories;
use crate::infra::fs::{self as fs, FsClient};
use crate::infra::{app_server, git};
use crate::ui::activity_heatmap;
use crate::ui::state::app_mode::AppMode;
fn mock_app_server() -> Arc<dyn app_server::AppServerClient> {
Arc::new(MockAppServerClient::new())
}
fn create_passthrough_mock_fs_client() -> fs::MockFsClient {
let mut mock_fs_client = fs::MockFsClient::new();
mock_fs_client
.expect_create_dir_all()
.times(0..)
.returning(|path| {
Box::pin(async move {
tokio::fs::create_dir_all(path)
.await
.map_err(fs::FsError::from)
})
});
mock_fs_client
.expect_remove_dir_all()
.times(0..)
.returning(|path| {
Box::pin(async move {
tokio::fs::remove_dir_all(path)
.await
.map_err(fs::FsError::from)
})
});
mock_fs_client
.expect_read_file()
.times(0..)
.returning(|path| {
Box::pin(async move { tokio::fs::read(path).await.map_err(fs::FsError::from) })
});
mock_fs_client
.expect_remove_file()
.times(0..)
.returning(|path| {
Box::pin(async move {
match tokio::fs::remove_file(path).await {
Ok(()) => Ok(()),
Err(error) if error.kind() == std::io::ErrorKind::NotFound => Ok(()),
Err(error) => Err(fs::FsError::from(error)),
}
})
});
mock_fs_client
.expect_is_dir()
.times(0..)
.returning(|path| path.is_dir());
mock_fs_client
}
struct TestClock {
instant: Mutex<Instant>,
system_time: Mutex<SystemTime>,
}
impl TestClock {
fn new(instant: Instant, system_time: SystemTime) -> Self {
Self {
instant: Mutex::new(instant),
system_time: Mutex::new(system_time),
}
}
fn advance(&self, duration: Duration) {
if let Ok(mut instant) = self.instant.lock() {
*instant += duration;
}
if let Ok(mut system_time) = self.system_time.lock() {
*system_time += duration;
}
}
}
impl Clock for TestClock {
fn now_instant(&self) -> Instant {
*self
.instant
.lock()
.expect("test clock instant lock should not be poisoned")
}
fn now_system_time(&self) -> SystemTime {
*self
.system_time
.lock()
.expect("test clock system-time lock should not be poisoned")
}
}
fn create_mock_backend() -> MockAgentBackend {
let mut mock = MockAgentBackend::new();
mock.expect_build_command().returning(|request| {
let mut cmd = Command::new("sh");
cmd.arg("-c")
.arg("printf '{\"answer\":\"mock-start\",\"questions\":[],\"summary\":null}'")
.current_dir(request.folder)
.stdout(Stdio::piped())
.stderr(Stdio::null());
Ok(cmd)
});
mock
}
fn allow_detect_git_info(mock: &mut git::MockGitClient) {
mock.expect_detect_git_info().times(0..).returning(|path| {
let branch_name = path
.file_name()
.and_then(|file_name| file_name.to_str())
.filter(|folder_name| folder_name.len() == 8)
.map_or_else(
|| "main".to_string(),
|folder_name| format!("wt/{folder_name}"),
);
Box::pin(async move { Some(branch_name) })
});
mock.expect_main_repo_root().times(0..).returning(|path| {
let repo_root = path
.parent()
.map(std::path::Path::to_path_buf)
.unwrap_or(path);
Box::pin(async move { Ok(repo_root) })
});
mock.expect_worktree_status()
.times(0..)
.returning(|_| Box::pin(async { Ok(String::new()) }));
mock.expect_tracked_worktree_status()
.times(0..)
.returning(|_| Box::pin(async { Ok(String::new()) }));
mock.expect_fetch_remote()
.times(0..)
.returning(|_| Box::pin(async { Ok(()) }));
mock.expect_branch_tracking_statuses()
.times(0..)
.returning(|_| Box::pin(async { Ok(HashMap::new()) }));
mock.expect_get_ref_ahead_behind()
.times(0..)
.returning(|_, _, _| Box::pin(async { Ok((0, 0)) }));
mock.expect_get_ahead_behind()
.times(0..)
.returning(|_| Box::pin(async { Ok((0, 0)) }));
}
fn create_mock_git_client_for_successful_noop_merges(
expected_merge_count: usize,
repo_root: PathBuf,
) -> git::MockGitClient {
let mut mock = git::MockGitClient::new();
allow_detect_git_info(&mut mock);
mock.expect_find_git_repo_root()
.times(expected_merge_count)
.returning(move |_| {
let repo_root = repo_root.clone();
Box::pin(async move { Some(repo_root) })
});
mock.expect_is_worktree_clean()
.times(expected_merge_count * 2)
.returning(|_| Box::pin(async { Ok(true) }));
mock.expect_is_rebase_in_progress()
.times(expected_merge_count)
.returning(|_| Box::pin(async { Ok(false) }));
mock.expect_rebase_start()
.times(expected_merge_count)
.returning(|_, _| Box::pin(async { Ok(git::RebaseStepResult::Completed) }));
mock.expect_squash_merge_diff()
.times(expected_merge_count)
.returning(|_, _, _| Box::pin(async { Ok(String::new()) }));
mock.expect_remove_worktree()
.times(expected_merge_count)
.returning(|worktree_path| {
Box::pin(async move {
let fs_client = create_passthrough_mock_fs_client();
let _ = fs_client.remove_dir_all(worktree_path).await;
Ok(())
})
});
mock.expect_delete_branch()
.times(expected_merge_count)
.returning(|_, _| Box::pin(async { Ok(()) }));
mock
}
fn create_default_mock_git_client(repo_root: PathBuf) -> git::MockGitClient {
let mut mock = git::MockGitClient::new();
setup_mock_worktree_expectations(&mut mock, repo_root);
setup_mock_merge_and_rebase_expectations(&mut mock);
setup_mock_commit_and_branch_expectations(&mut mock);
mock
}
fn setup_mock_worktree_expectations(mock: &mut git::MockGitClient, repo_root: PathBuf) {
let find_repo_root = repo_root.clone();
mock.expect_detect_git_info().times(0..).returning({
let repo_root = repo_root.clone();
move |path| {
let branch_name = if path == repo_root {
"main".to_string()
} else {
path.file_name()
.and_then(|file_name| file_name.to_str())
.map_or_else(
|| "main".to_string(),
|folder_name| format!("wt/{folder_name}"),
)
};
Box::pin(async move { Some(branch_name) })
}
});
mock.expect_current_upstream_reference()
.times(0..)
.returning(|_| Box::pin(async { Ok("origin/main".to_string()) }));
mock.expect_find_git_repo_root()
.times(0..)
.returning(move |_| {
let repo_root = find_repo_root.clone();
Box::pin(async move { Some(repo_root) })
});
mock.expect_create_worktree()
.times(0..)
.returning(|_, worktree_path, _, _| {
Box::pin(async move {
let fs_client = create_passthrough_mock_fs_client();
fs_client
.create_dir_all(worktree_path.clone())
.await
.map_err(|error| {
git::GitError::OutputParse(format!(
"Failed to create mock worktree directory: {error}"
))
})?;
fs_client
.create_dir_all(worktree_path.join(SESSION_DATA_DIR))
.await
.map_err(|error| {
git::GitError::OutputParse(format!(
"Failed to create mock session data directory: {error}"
))
})?;
Ok(())
})
});
mock.expect_remove_worktree()
.times(0..)
.returning(|worktree_path| {
Box::pin(async move {
let fs_client = create_passthrough_mock_fs_client();
let _ = fs_client.remove_dir_all(worktree_path).await;
Ok(())
})
});
mock.expect_pull_rebase().times(0..).returning(|_| {
Box::pin(async {
Err(git::GitError::OutputParse(
"No upstream branch configured for pull".to_string(),
))
})
});
mock.expect_push_current_branch()
.times(0..)
.returning(|_| Box::pin(async { Ok("origin/main".to_string()) }));
mock.expect_fetch_remote()
.times(0..)
.returning(|_| Box::pin(async { Ok(()) }));
mock.expect_branch_tracking_statuses()
.times(0..)
.returning(|_| Box::pin(async { Ok(HashMap::new()) }));
mock.expect_get_ahead_behind()
.times(0..)
.returning(|_| Box::pin(async { Ok((0, 0)) }));
mock.expect_get_ref_ahead_behind()
.times(0..)
.returning(|_, _, _| Box::pin(async { Ok((0, 0)) }));
mock.expect_list_upstream_commit_titles()
.times(0..)
.returning(|_| Box::pin(async { Ok(Vec::new()) }));
mock.expect_list_local_commit_titles()
.times(0..)
.returning(|_| Box::pin(async { Ok(Vec::new()) }));
mock.expect_repo_url()
.times(0..)
.returning(|_| Box::pin(async { Ok("https://example.invalid/repo.git".to_string()) }));
mock.expect_main_repo_root().times(0..).returning(move |_| {
let repo_root = repo_root.clone();
Box::pin(async move { Ok(repo_root) })
});
}
fn setup_mock_merge_and_rebase_expectations(mock: &mut git::MockGitClient) {
mock.expect_squash_merge_diff()
.times(0..)
.returning(|_, _, _| Box::pin(async { Ok(String::new()) }));
mock.expect_squash_merge()
.times(0..)
.returning(|_, _, _, _| Box::pin(async { Ok(git::SquashMergeOutcome::Committed) }));
mock.expect_rebase()
.times(0..)
.returning(|_, _| Box::pin(async { Ok(()) }));
mock.expect_rebase_start()
.times(0..)
.returning(|_, _| Box::pin(async { Ok(git::RebaseStepResult::Completed) }));
mock.expect_rebase_continue()
.times(0..)
.returning(|_| Box::pin(async { Ok(git::RebaseStepResult::Completed) }));
mock.expect_abort_rebase()
.times(0..)
.returning(|_| Box::pin(async { Ok(()) }));
mock.expect_is_rebase_in_progress()
.times(0..)
.returning(|_| Box::pin(async { Ok(false) }));
mock.expect_has_unmerged_paths()
.times(0..)
.returning(|_| Box::pin(async { Ok(false) }));
mock.expect_list_staged_conflict_marker_files()
.times(0..)
.returning(|_, _| Box::pin(async { Ok(Vec::new()) }));
mock.expect_list_conflicted_files()
.times(0..)
.returning(|_| Box::pin(async { Ok(Vec::new()) }));
}
async fn synthetic_diff_from_session_folder(folder: &Path) -> String {
async fn count_lines_recursively(fs_client: &dyn fs::FsClient, root: &Path) -> usize {
let mut pending_entries = vec![root.to_path_buf()];
let mut line_count = 0;
while let Some(entry) = pending_entries.pop() {
if !entry.is_dir() {
if let Ok(content) = fs_client.read_file(entry).await {
line_count += String::from_utf8_lossy(&content).lines().count();
}
continue;
}
if entry
.file_name()
.is_some_and(|name| name == SESSION_DATA_DIR)
{
continue;
}
if let Ok(entries) = std::fs::read_dir(entry) {
pending_entries.extend(
entries
.filter_map(Result::ok)
.map(|dir_entry| dir_entry.path()),
);
}
}
line_count
}
let fs_client = create_passthrough_mock_fs_client();
let line_count = count_lines_recursively(&fs_client, folder).await;
if line_count == 0 {
String::new()
} else {
"+\n".repeat(line_count)
}
}
fn setup_mock_commit_and_branch_expectations(mock: &mut git::MockGitClient) {
mock.expect_commit_all()
.times(0..)
.returning(|_, _, _| Box::pin(async { Ok(()) }));
mock.expect_commit_all_preserving_single_commit()
.times(0..)
.returning(|_, _, _, _, _| {
Box::pin(async {
Err(git::GitError::OutputParse(
"Nothing to commit: no changes detected".to_string(),
))
})
});
mock.expect_stage_all()
.times(0..)
.returning(|_| Box::pin(async { Ok(()) }));
mock.expect_head_short_hash()
.times(0..)
.returning(|_| Box::pin(async { Ok("abc1234".to_string()) }));
mock.expect_delete_branch()
.times(0..)
.returning(|_, _| Box::pin(async { Ok(()) }));
mock.expect_diff().times(0..).returning(|folder, _| {
Box::pin(async move { Ok(synthetic_diff_from_session_folder(&folder).await) })
});
mock.expect_is_worktree_clean()
.times(0..)
.returning(|_| Box::pin(async { Ok(true) }));
mock.expect_worktree_status()
.times(0..)
.returning(|_| Box::pin(async { Ok(String::new()) }));
mock.expect_tracked_worktree_status()
.times(0..)
.returning(|_| Box::pin(async { Ok(String::new()) }));
mock.expect_has_commits_since()
.times(0..)
.returning(|_, _| Box::pin(async { Ok(false) }));
mock.expect_head_commit_message()
.times(0..)
.returning(|_| Box::pin(async { Ok(None) }));
}
fn install_mock_git_client(app: &mut App, mock_git_client: git::MockGitClient) {
let mock_git_client: Arc<dyn git::GitClient> = Arc::new(mock_git_client);
let base_path = app.services.base_path().to_path_buf();
let db = app.services.db().clone();
let event_sender = app.services.event_sender();
let available_agent_kinds = app.services.available_agent_kinds();
let app_server_client_override = app.services.app_server_client_override();
let fs_client = app.services.fs_client();
let review_request_client = app.services.review_request_client();
app.services = AppServices::new(
base_path,
app.services.clock(),
event_sender,
crate::app::service::AppServiceDeps {
app_server_client_override,
available_agent_kinds,
fs_client,
git_client: Arc::clone(&mock_git_client),
repositories: db,
review_request_client,
},
);
app.sessions.git_client = mock_git_client;
}
fn test_app_clients() -> crate::app::AppClients {
crate::app::AppClients::new().with_agent_availability_probe(Arc::new(
crate::infra::agent::StaticAgentAvailabilityProbe {
available_agent_kinds: AgentKind::ALL.to_vec(),
},
))
}
async fn new_test_app_with_db_and_app_server(
path: PathBuf,
working_dir: PathBuf,
git_branch: Option<String>,
db: AppRepositories,
app_server_client: Arc<dyn app_server::AppServerClient>,
) -> App {
let clients = test_app_clients().with_app_server_client_override(app_server_client);
let mut app = App::new_with_clients(path, working_dir.clone(), git_branch, db, clients)
.await
.expect("failed to build app");
let mock_git_client = create_default_mock_git_client(working_dir);
install_mock_git_client(&mut app, mock_git_client);
app
}
async fn new_test_app_with_db(
path: PathBuf,
working_dir: PathBuf,
git_branch: Option<String>,
db: AppRepositories,
) -> App {
new_test_app_with_db_and_app_server(path, working_dir, git_branch, db, mock_app_server())
.await
}
async fn new_test_app(path: PathBuf) -> App {
let working_dir = PathBuf::from("/tmp/test");
let db = AppRepositories::in_memory().await;
new_test_app_with_db(path, working_dir, None, db).await
}
async fn new_test_app_with_git(path: &Path) -> App {
let db = AppRepositories::in_memory().await;
new_test_app_with_git_and_db(path, db).await
}
async fn new_test_app_with_git_and_db(path: &Path, db: AppRepositories) -> App {
new_test_app_with_db(
path.to_path_buf(),
path.to_path_buf(),
Some("main".to_string()),
db,
)
.await
}
fn add_manual_session(app: &mut App, base_path: &Path, id: &str, prompt: &str) {
add_manual_session_with_status(app, base_path, id, prompt, Status::Review);
}
fn add_manual_session_with_status(
app: &mut App,
base_path: &Path,
id: &str,
prompt: &str,
status: Status,
) {
let folder = session_folder(base_path, id);
let data_dir = folder.join(SESSION_DATA_DIR);
std::fs::create_dir_all(&data_dir).expect("failed to create data dir");
app.sessions.session_handles_mut().insert(
id.to_string().into(),
SessionHandles::new(String::new(), status),
);
app.sessions.push_session(Session {
base_branch: "main".to_string(),
created_at: 0,
draft_attachments: Vec::new(),
folder,
follow_up_tasks: Vec::new(),
id: id.into(),
in_progress_started_at: None,
in_progress_total_seconds: 0,
is_draft: false,
model: AgentModel::Gemini3FlashPreview,
output: String::new(),
parent_session_id: None,
project_name: String::new(),
prompt: prompt.to_string(),
queued_messages: Vec::new(),
reasoning_level_override: None,
published_upstream_ref: None,
published_branch_sync_status: crate::domain::session::PublishedBranchSyncStatus::Idle,
questions: Vec::new(),
review_request: None,
size: SessionSize::Xs,
stats: SessionStats::default(),
status,
summary: None,
title: Some(prompt.to_string()),
updated_at: 0,
workflow_notice: None,
});
if app.sessions.selected_session_index().is_none() {
app.sessions.select_session_index(Some(0));
}
}
fn test_session_manager(
session_id: &str,
reasoning_level_override: Option<ReasoningLevel>,
) -> SessionManager {
test_session_manager_with_clock(session_id, reasoning_level_override, Arc::new(RealClock))
}
fn test_session_manager_with_clock(
session_id: &str,
reasoning_level_override: Option<ReasoningLevel>,
clock: Arc<dyn Clock>,
) -> SessionManager {
let mut handles = HashMap::new();
handles.insert(
session_id.to_string().into(),
SessionHandles::new(String::new(), Status::Review),
);
let state = SessionState::new(
handles,
vec![Session {
base_branch: "main".to_string(),
created_at: 0,
draft_attachments: Vec::new(),
folder: PathBuf::from(format!("/tmp/{session_id}")),
follow_up_tasks: Vec::new(),
id: session_id.into(),
in_progress_started_at: None,
in_progress_total_seconds: 0,
is_draft: false,
model: AgentModel::Gpt54,
output: String::new(),
parent_session_id: None,
project_name: "project".to_string(),
prompt: String::new(),
queued_messages: Vec::new(),
reasoning_level_override,
published_upstream_ref: None,
published_branch_sync_status:
crate::domain::session::PublishedBranchSyncStatus::Idle,
questions: Vec::new(),
review_request: None,
size: SessionSize::Xs,
stats: SessionStats::default(),
status: Status::Review,
summary: None,
title: Some("Title".to_string()),
updated_at: 0,
workflow_notice: None,
}],
ratatui::widgets::TableState::default(),
clock,
1,
0,
);
SessionManager::new(
SessionDefaults {
model: AgentModel::Gpt54,
},
Arc::new(git::MockGitClient::new()),
state,
Vec::new(),
)
}
#[test]
fn test_set_and_get_at_mention_index_for_root_cache() {
let mut session_manager = test_session_manager("session-id", None);
let temp_dir = tempfile::tempdir().expect("create temp dir");
let lookup_root = temp_dir.path().to_path_buf();
let entries = vec![FileEntry {
is_dir: false,
path: "src/main.rs".to_string(),
}];
session_manager.set_at_mention_index_for_root(lookup_root.clone(), entries.clone());
assert_eq!(
session_manager
.at_mention_index_for_root(&lookup_root)
.expect("expected cached entries"),
entries
);
}
#[test]
fn test_at_mention_index_for_root_invalidates_after_ttl_expires() {
let initial_instant = Instant::now();
let initial_system_time = SystemTime::UNIX_EPOCH + Duration::from_mins(1);
let clock = Arc::new(TestClock::new(initial_instant, initial_system_time));
let mut session_manager =
test_session_manager_with_clock("session-id", None, clock.clone());
let temp_dir = tempfile::tempdir().expect("create temp dir");
let lookup_root = temp_dir.path().to_path_buf();
let entries = vec![FileEntry {
is_dir: true,
path: "src".to_string(),
}];
session_manager.set_at_mention_index_for_root(lookup_root.clone(), entries);
clock.advance(AT_MENTION_INDEX_TTL + Duration::from_secs(1));
let cached_entries = session_manager.at_mention_index_for_root(&lookup_root);
assert!(cached_entries.is_none());
}
#[test]
fn test_remove_at_mention_index_for_root_drops_cached_entries() {
let mut session_manager = test_session_manager("session-id", None);
let temp_dir = tempfile::tempdir().expect("create temp dir");
let lookup_root = temp_dir.path().to_path_buf();
let entries = vec![FileEntry {
is_dir: true,
path: "src".to_string(),
}];
session_manager.set_at_mention_index_for_root(lookup_root.clone(), entries);
session_manager.remove_at_mention_index_for_root(&lookup_root);
assert!(
session_manager
.at_mention_index_for_root(&lookup_root)
.is_none()
);
}
#[test]
fn test_apply_session_reasoning_level_updated_updates_only_matching_session() {
let mut session_manager = test_session_manager("session-id", Some(ReasoningLevel::Low));
session_manager
.apply_session_reasoning_level_updated("other-session", Some(ReasoningLevel::High));
let reasoning_level_after_non_matching_update =
session_manager.state.sessions[0].reasoning_level_override;
session_manager.apply_session_reasoning_level_updated("session-id", None);
assert_eq!(
reasoning_level_after_non_matching_update,
Some(ReasoningLevel::Low)
);
assert_eq!(
session_manager.state.sessions[0].reasoning_level_override,
None
);
}
#[tokio::test]
async fn test_apply_turn_applied_state_clears_active_prompt_output() {
let temp_dir = tempfile::tempdir().expect("failed to create temp dir");
let mut app = new_test_app(temp_dir.path().to_path_buf()).await;
add_manual_session_with_status(
&mut app,
temp_dir.path(),
"session-id",
"Prompt",
Status::InProgress,
);
app.sessions
.set_active_prompt_output("session-id", " › Prompt\n\n".to_string());
app.sessions.apply_turn_applied_state(
"session-id",
&TurnAppliedState {
follow_up_tasks: Vec::new(),
questions: Vec::new(),
summary: None,
token_usage_delta: SessionStats::default(),
},
);
assert!(
!app.sessions
.active_prompt_outputs()
.contains_key("session-id")
);
}
#[tokio::test]
async fn test_retain_active_prompt_outputs_keeps_only_active_sessions() {
let temp_dir = tempfile::tempdir().expect("failed to create temp dir");
let mut app = new_test_app(temp_dir.path().to_path_buf()).await;
add_manual_session_with_status(
&mut app,
temp_dir.path(),
"active-session",
"Prompt",
Status::InProgress,
);
add_manual_session_with_status(
&mut app,
temp_dir.path(),
"review-session",
"Prompt",
Status::Review,
);
app.sessions
.set_active_prompt_output("active-session", " › Active\n\n".to_string());
app.sessions
.set_active_prompt_output("review-session", " › Review\n\n".to_string());
app.sessions.retain_active_prompt_outputs();
assert!(
app.sessions
.active_prompt_outputs()
.contains_key("active-session")
);
assert!(
!app.sessions
.active_prompt_outputs()
.contains_key("review-session")
);
}
#[test]
fn test_retain_active_prompt_outputs_prunes_expired_at_mention_indexes() {
let initial_instant = Instant::now();
let initial_system_time = SystemTime::UNIX_EPOCH + Duration::from_mins(1);
let clock = Arc::new(TestClock::new(initial_instant, initial_system_time));
let mut session_manager =
test_session_manager_with_clock("session-id", None, clock.clone());
let temp_dir = tempfile::tempdir().expect("create temp dir");
let lookup_root = temp_dir.path().to_path_buf();
session_manager.set_at_mention_index_for_root(
lookup_root.clone(),
vec![FileEntry {
is_dir: false,
path: "src/main.rs".to_string(),
}],
);
clock.advance(AT_MENTION_INDEX_TTL + Duration::from_secs(1));
session_manager.retain_active_prompt_outputs();
assert!(
session_manager
.at_mention_index_for_root(&lookup_root)
.is_none()
);
}
async fn create_and_start_session(app: &mut App, prompt: &str) {
let session_id = app
.create_session()
.await
.expect("failed to create session");
let start_backend = create_mock_backend();
app.sessions
.reply_with_backend(
&app.services,
&session_id,
prompt,
Arc::new(start_backend),
AgentModel::ClaudeOpus48,
)
.await;
}
async fn wait_for_status(app: &mut App, session_id: &str, expected: Status) {
wait_for_status_with_retries(app, session_id, expected, 2000).await;
}
async fn wait_for_status_with_retries(
app: &mut App,
session_id: &str,
expected: Status,
retries: usize,
) {
for _ in 0..retries {
app.process_pending_app_events().await;
app.sessions.sync_from_handles();
let Some(session) = app
.sessions
.sessions()
.iter()
.find(|session| session.id == session_id)
else {
break;
};
if session.status == expected {
return;
}
tokio::time::sleep(std::time::Duration::from_millis(50)).await;
}
app.process_pending_app_events().await;
app.sessions.sync_from_handles();
let session = app
.sessions
.sessions()
.iter()
.find(|session| session.id == session_id)
.expect("missing session while waiting for status");
assert_eq!(
session.status, expected,
"session output while waiting for status: {}",
session.output
);
}
fn set_session_status_for_test(app: &mut App, session_id: &str, status: Status) {
if let Some(session) = app
.sessions
.sessions_mut()
.iter_mut()
.find(|session| session.id == session_id)
{
session.status = status;
}
if let Some(handles) = app.sessions.session_handles().get(session_id)
&& let Ok(mut current_status) = handles.status.lock()
{
*current_status = status;
}
}
async fn wait_for_path_absent(path: &Path) {
for _ in 0..500 {
if !path.exists() {
return;
}
tokio::time::sleep(std::time::Duration::from_millis(20)).await;
}
assert!(!path.exists(), "timed out waiting for path cleanup");
}
async fn wait_for_output_contains(
app: &mut App,
session_id: &str,
expected_output: &str,
retries: usize,
) {
for _ in 0..retries {
app.sessions.sync_from_handles();
let Some(session) = app
.sessions
.sessions()
.iter()
.find(|session| session.id == session_id)
else {
break;
};
if session.output.contains(expected_output) {
return;
}
tokio::time::sleep(std::time::Duration::from_millis(50)).await;
}
app.sessions.sync_from_handles();
let session = app
.sessions
.sessions()
.iter()
.find(|session| session.id == session_id)
.expect("missing session while waiting for output");
assert!(
session.output.contains(expected_output),
"expected output to contain: {expected_output}, actual output: {}",
session.output
);
}
fn session_status_or_done(app: &App, session_id: &str) -> Status {
app.sessions
.sessions()
.iter()
.find(|session| session.id == session_id)
.map_or(Status::Done, |session| session.status)
}
fn is_session_done(app: &App, session_id: &str) -> bool {
app.sessions
.sessions()
.iter()
.find(|session| session.id == session_id)
.is_some_and(|session| session.status == Status::Done)
}
async fn wait_for_first_merge_to_complete_before_second_starts(
app: &mut App,
first_session_id: &str,
second_session_id: &str,
) {
let mut first_merge_completed = false;
let mut first_merge_pending_observed = false;
let mut second_merge_was_queued = false;
for _ in 0..5000 {
app.process_pending_app_events().await;
app.sessions.sync_from_handles();
let first_status = session_status_or_done(app, first_session_id);
let second_status = session_status_or_done(app, second_session_id);
if second_status == Status::Queued {
second_merge_was_queued = true;
}
if first_status == Status::Done {
first_merge_completed = true;
break;
}
first_merge_pending_observed = true;
assert_ne!(
second_status,
Status::Merging,
"second merge started before first completed"
);
tokio::time::sleep(std::time::Duration::from_millis(50)).await;
}
assert!(
first_merge_completed,
"first merge did not complete within timeout"
);
if first_merge_pending_observed {
assert!(
second_merge_was_queued,
"second merge never entered queued status before first completed"
);
}
}
async fn wait_for_second_merge_to_start(app: &mut App, second_session_id: &str) {
let mut second_merge_started = false;
for _ in 0..5000 {
app.process_pending_app_events().await;
app.sessions.sync_from_handles();
let second_status = session_status_or_done(app, second_session_id);
if matches!(second_status, Status::Merging | Status::Done) {
second_merge_started = true;
break;
}
tokio::time::sleep(std::time::Duration::from_millis(50)).await;
}
assert!(
second_merge_started,
"second merge did not start after first completed"
);
}
async fn wait_for_all_sessions_done(
app: &mut App,
first_session_id: &str,
second_session_id: &str,
) {
for _ in 0..5000 {
app.process_pending_app_events().await;
app.sessions.sync_from_handles();
if is_session_done(app, first_session_id) && is_session_done(app, second_session_id) {
return;
}
tokio::time::sleep(std::time::Duration::from_millis(50)).await;
}
}
#[tokio::test]
async fn test_new_app_empty() {
let dir = tempdir().expect("failed to create temp dir");
let app = new_test_app(dir.path().to_path_buf()).await;
assert!(app.sessions.sessions().is_empty());
assert_eq!(app.sessions.selected_session_index(), None);
}
#[tokio::test]
async fn test_working_dir_getter() {
let dir = tempdir().expect("failed to create temp dir");
let app = new_test_app(dir.path().to_path_buf()).await;
let working_dir = app.working_dir();
assert_eq!(working_dir, Path::new("/tmp/test"));
}
#[tokio::test]
async fn test_git_branch_getter_with_branch() {
let dir = tempdir().expect("failed to create temp dir");
let working_dir = PathBuf::from("/tmp/test");
let db = AppRepositories::in_memory().await;
let app = new_test_app_with_db(
dir.path().to_path_buf(),
working_dir,
Some("main".to_string()),
db,
)
.await;
let branch = app.git_branch();
assert_eq!(branch, Some("main"));
}
#[tokio::test]
async fn test_git_branch_getter_without_branch() {
let dir = tempdir().expect("failed to create temp dir");
let app = new_test_app(dir.path().to_path_buf()).await;
let branch = app.git_branch();
assert_eq!(branch, None);
}
#[tokio::test]
async fn test_navigation() {
let dir = tempdir().expect("failed to create temp dir");
let mut app = new_test_app_with_git(dir.path()).await;
create_and_start_session(&mut app, "A").await;
create_and_start_session(&mut app, "B").await;
app.sessions.select_session_index(Some(0));
app.next();
assert_eq!(app.sessions.selected_session_index(), Some(1));
app.next();
assert_eq!(app.sessions.selected_session_index(), Some(0));
app.previous();
assert_eq!(app.sessions.selected_session_index(), Some(1)); app.previous();
assert_eq!(app.sessions.selected_session_index(), Some(0));
}
#[tokio::test]
async fn test_navigation_follows_grouped_order_and_skips_group_headers() {
let dir = tempdir().expect("failed to create temp dir");
let mut app = new_test_app(dir.path().to_path_buf()).await;
add_manual_session_with_status(
&mut app,
dir.path(),
"archive-1",
"Archive 1",
Status::Done,
);
add_manual_session_with_status(
&mut app,
dir.path(),
"active-1",
"Active 1",
Status::Review,
);
add_manual_session_with_status(
&mut app,
dir.path(),
"queued-1",
"Queued 1",
Status::Queued,
);
add_manual_session_with_status(
&mut app,
dir.path(),
"archive-2",
"Archive 2",
Status::Canceled,
);
add_manual_session_with_status(&mut app, dir.path(), "merge-1", "Merge 1", Status::Merging);
add_manual_session_with_status(&mut app, dir.path(), "active-2", "Active 2", Status::Draft);
app.sessions.select_session_index(Some(3));
app.next();
assert_eq!(
app.selected_session().map(|session| session.id.as_str()),
Some("queued-1")
);
app.next();
assert_eq!(
app.selected_session().map(|session| session.id.as_str()),
Some("merge-1")
);
app.next();
assert_eq!(
app.selected_session().map(|session| session.id.as_str()),
Some("active-1")
);
app.next();
assert_eq!(
app.selected_session().map(|session| session.id.as_str()),
Some("active-2")
);
app.next();
assert_eq!(
app.selected_session().map(|session| session.id.as_str()),
Some("archive-1")
);
app.next();
assert_eq!(
app.selected_session().map(|session| session.id.as_str()),
Some("archive-2")
);
app.previous();
assert_eq!(
app.selected_session().map(|session| session.id.as_str()),
Some("archive-1")
);
}
#[tokio::test]
async fn test_navigation_empty() {
let dir = tempdir().expect("failed to create temp dir");
let mut app = new_test_app(dir.path().to_path_buf()).await;
app.next();
assert_eq!(app.sessions.selected_session_index(), None);
app.previous();
assert_eq!(app.sessions.selected_session_index(), None);
}
#[tokio::test]
async fn test_navigation_recovery() {
let dir = tempdir().expect("failed to create temp dir");
let mut app = new_test_app_with_git(dir.path()).await;
create_and_start_session(&mut app, "A").await;
app.sessions.select_session_index(None);
app.next();
assert_eq!(app.sessions.selected_session_index(), Some(0));
app.sessions.select_session_index(None);
app.previous();
assert_eq!(app.sessions.selected_session_index(), Some(0));
}
#[tokio::test]
async fn test_create_session() {
let dir = tempdir().expect("failed to create temp dir");
let mut app = new_test_app_with_git(dir.path()).await;
let session_id = app
.create_session()
.await
.expect("failed to create session");
assert_eq!(app.sessions.sessions().len(), 1);
assert_eq!(app.sessions.sessions()[0].id, session_id);
assert!(app.sessions.sessions()[0].prompt.is_empty());
assert_eq!(app.sessions.sessions()[0].title, None);
assert_eq!(app.sessions.sessions()[0].display_title(), "No title");
assert!(!app.sessions.sessions()[0].is_draft_session());
assert_eq!(app.sessions.sessions()[0].status, Status::Draft);
assert_eq!(app.sessions.selected_session_index(), Some(0));
assert_eq!(
app.sessions.sessions()[0].model,
AgentKind::Gemini.default_model()
);
let session_dir = &app.sessions.sessions()[0].folder;
let data_dir = session_dir.join(SESSION_DATA_DIR);
assert!(session_dir.exists());
assert!(data_dir.is_dir());
let db_sessions = app
.services
.db()
.sessions()
.load_sessions()
.await
.expect("failed to load");
let activity_timestamps = app
.services
.db()
.activity()
.load_session_activity_timestamps()
.await
.expect("failed to load session activity timestamps");
assert_eq!(db_sessions.len(), 1);
assert_eq!(db_sessions[0].base_branch, "main");
assert_eq!(
db_sessions[0].model,
AgentKind::Gemini.default_model().as_str()
);
assert!(!db_sessions[0].is_draft);
assert_eq!(db_sessions[0].status, "Draft");
assert_eq!(activity_timestamps.len(), 1);
}
#[tokio::test]
async fn test_create_draft_session() {
let dir = tempdir().expect("failed to create temp dir");
let mut app = new_test_app_with_git(dir.path()).await;
let session_id = app
.create_draft_session()
.await
.expect("failed to create draft session");
assert_eq!(app.sessions.sessions().len(), 1);
assert_eq!(app.sessions.sessions()[0].id, session_id);
assert!(app.sessions.sessions()[0].is_draft_session());
assert_eq!(app.sessions.sessions()[0].status, Status::Draft);
assert!(!app.sessions.sessions()[0].folder.exists());
let db_sessions = app
.services
.db()
.sessions()
.load_sessions()
.await
.expect("failed to load");
assert!(db_sessions[0].is_draft);
}
#[tokio::test]
async fn test_create_stacked_draft_session_persists_parent_and_base_branch() {
let dir = tempdir().expect("failed to create temp dir");
let mut app = new_test_app_with_git(dir.path()).await;
let parent_session_id = app.create_session().await.expect("failed to create parent");
let expected_base_branch = session_branch(&parent_session_id);
let child_session_id = app
.create_stacked_draft_session(&parent_session_id)
.await
.expect("failed to create stacked draft session");
let child_session = app
.sessions
.sessions()
.iter()
.find(|session| session.id == child_session_id)
.expect("missing child session");
assert!(child_session.is_draft_session());
assert_eq!(child_session.status, Status::Draft);
assert_eq!(
child_session.parent_session_id.as_deref(),
Some(parent_session_id.as_str())
);
assert_eq!(child_session.base_branch, expected_base_branch);
assert!(!child_session.folder.exists());
let db_sessions = app
.services
.db()
.sessions()
.load_sessions()
.await
.expect("failed to load");
let db_child_session = db_sessions
.iter()
.find(|session| session.id == child_session_id)
.expect("missing persisted child session");
assert_eq!(
db_child_session.parent_session_id.as_deref(),
Some(parent_session_id.as_str())
);
assert_eq!(db_child_session.base_branch, expected_base_branch);
}
#[tokio::test]
async fn test_create_stacked_draft_session_rejects_nested_child() {
let dir = tempdir().expect("failed to create temp dir");
let mut app = new_test_app_with_git(dir.path()).await;
let parent_session_id = app.create_session().await.expect("failed to create parent");
let child_session_id = app
.create_stacked_draft_session(&parent_session_id)
.await
.expect("failed to create stacked draft session");
let result = app.create_stacked_draft_session(&child_session_id).await;
let error = result.expect_err("nested stacked draft creation should fail");
assert!(
error
.to_string()
.contains("Stacked sessions can only be created from root sessions")
);
}
#[tokio::test]
async fn test_create_session_keeps_default_smart_model_setting_when_session_model_changes() {
let dir = tempdir().expect("failed to create temp dir");
let mut app = new_test_app_with_git(dir.path()).await;
let first_session_id = app
.create_session()
.await
.expect("failed to create first session");
app.set_session_model(&first_session_id, AgentModel::Gpt54)
.await
.expect("failed to set session model");
let active_project_id = app.active_project_id();
let default_smart_model_setting = app
.services
.db()
.settings()
.get_project_setting(active_project_id, SettingName::DefaultSmartModel)
.await
.expect("failed to load setting");
let second_session_id = app
.create_session()
.await
.expect("failed to create second session");
let second_session = app
.sessions
.sessions()
.iter()
.find(|session| session.id == second_session_id)
.expect("missing second session");
assert_eq!(second_session.model, AgentKind::Gemini.default_model());
assert_eq!(default_smart_model_setting, None);
let db_sessions = app
.services
.db()
.sessions()
.load_sessions()
.await
.expect("failed to load");
let db_second_session = db_sessions
.iter()
.find(|session| session.id == second_session_id)
.expect("missing second session in db");
assert_eq!(
db_second_session.model,
AgentKind::Gemini.default_model().as_str()
);
}
#[tokio::test]
async fn test_create_session_persists_default_smart_model_setting_when_last_used_model_is_enabled()
{
let dir = tempdir().expect("failed to create temp dir");
let db = AppRepositories::in_memory().await;
let mut app = new_test_app_with_git_and_db(dir.path(), db.clone()).await;
let active_project_id = app.active_project_id();
app.services
.db()
.settings()
.upsert_project_setting(
active_project_id,
SettingName::LastUsedModelAsDefault,
"true",
)
.await
.expect("failed to upsert last-used-model setting");
let first_session_id = app
.create_session()
.await
.expect("failed to create first session");
app.set_session_model(&first_session_id, AgentModel::Gpt54)
.await
.expect("failed to set session model");
let default_smart_model_setting = app
.services
.db()
.settings()
.get_project_setting(active_project_id, SettingName::DefaultSmartModel)
.await
.expect("failed to load setting");
drop(app);
let mut restarted_app = new_test_app_with_git_and_db(dir.path(), db).await;
let second_session_id = restarted_app
.create_session()
.await
.expect("failed to create second session");
assert_eq!(
default_smart_model_setting,
Some(AgentModel::Gpt54.as_str().to_string())
);
let second_session = restarted_app
.sessions
.sessions()
.iter()
.find(|session| session.id == second_session_id)
.expect("missing second session");
assert_eq!(second_session.model, AgentModel::Gpt54);
}
#[tokio::test]
async fn test_create_session_reads_default_smart_model_from_db_setting() {
let dir = tempdir().expect("failed to create temp dir");
let mut app = new_test_app_with_git(dir.path()).await;
let active_project_id = app.active_project_id();
app.services
.db()
.settings()
.upsert_project_setting(
active_project_id,
SettingName::DefaultSmartModel,
AgentModel::ClaudeHaiku4520251001.as_str(),
)
.await
.expect("failed to upsert default smart model setting");
let session_id = app
.create_session()
.await
.expect("failed to create session");
let created_session = app
.sessions
.sessions()
.iter()
.find(|session| session.id == session_id)
.expect("missing created session");
assert_eq!(created_session.model, AgentModel::ClaudeHaiku4520251001);
}
#[tokio::test]
async fn test_start_session() {
let dir = tempdir().expect("failed to create temp dir");
let mut app = new_test_app_with_git(dir.path()).await;
let session_id = app
.create_session()
.await
.expect("failed to create session");
app.start_session(&session_id, "Hello".to_string())
.await
.expect("failed to start session");
assert_eq!(app.sessions.sessions()[0].prompt, "Hello");
assert_eq!(app.sessions.sessions()[0].title, Some("Hello".to_string()));
app.sessions.sync_from_handles();
let output = app.sessions.sessions()[0].output.clone();
assert!(output.contains("Hello"));
let db_sessions = app
.services
.db()
.sessions()
.load_sessions()
.await
.expect("failed to load");
let activity_timestamps = app
.services
.db()
.activity()
.load_session_activity_timestamps()
.await
.expect("failed to load session activity timestamps");
assert_eq!(db_sessions[0].prompt, "Hello");
assert_eq!(db_sessions[0].output, " › Hello\n\n");
assert_eq!(activity_timestamps.len(), 1);
}
#[tokio::test]
async fn test_stage_draft_message_persists_bundle_without_starting_session() {
let dir = tempdir().expect("failed to create temp dir");
let mut app = new_test_app_with_git(dir.path()).await;
let session_id = app
.create_draft_session()
.await
.expect("failed to create draft session");
let session_update_versions = app.services.session_update_versions();
SessionTaskService::remove_session_update_version(
&session_update_versions,
session_id.as_str(),
);
app.stage_draft_message(&session_id, "First draft")
.await
.expect("failed to stage first draft");
app.stage_draft_message(&session_id, "Second draft")
.await
.expect("failed to stage second draft");
let session_update_version = {
let session_update_versions = session_update_versions
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
*session_update_versions
.get(session_id.as_str())
.unwrap_or(&0)
};
assert_eq!(app.sessions.sessions()[0].status, Status::Draft);
assert_eq!(
app.sessions.sessions()[0].prompt,
"First draft\n\nSecond draft"
);
assert_eq!(
app.sessions.sessions()[0].title,
Some("First draft".to_string())
);
assert!(app.sessions.sessions()[0].draft_attachments.is_empty());
let db_sessions = app
.services
.db()
.sessions()
.load_sessions()
.await
.expect("failed to load");
assert_eq!(db_sessions[0].prompt, "First draft\n\nSecond draft");
assert_eq!(db_sessions[0].title, Some("First draft".to_string()));
assert_eq!(session_update_version, 2);
}
#[tokio::test]
async fn test_stage_draft_message_keeps_generated_title_until_replacement_finishes() {
let dir = tempdir().expect("failed to create temp dir");
let mut app = new_test_app_with_git(dir.path()).await;
let session_id = app
.create_draft_session()
.await
.expect("failed to create draft session");
app.stage_draft_message(&session_id, "First draft")
.await
.expect("failed to stage first draft");
app.services
.db()
.sessions()
.update_session_title(&session_id, "Generated draft title")
.await
.expect("failed to persist generated draft title");
app.sessions.sessions_mut()[0].title = Some("Generated draft title".to_string());
app.stage_draft_message(&session_id, "Second draft")
.await
.expect("failed to stage second draft");
assert_eq!(
app.sessions.sessions()[0].prompt,
"First draft\n\nSecond draft"
);
assert_eq!(
app.sessions.sessions()[0].title,
Some("Generated draft title".to_string())
);
let db_sessions = app
.services
.db()
.sessions()
.load_sessions()
.await
.expect("failed to load");
assert_eq!(
db_sessions[0].title,
Some("Generated draft title".to_string())
);
}
#[tokio::test]
async fn test_replace_title_generation_task_aborts_superseded_task() {
let state = SessionState::new(
HashMap::new(),
Vec::new(),
TableState::default(),
Arc::new(RealClock),
1,
0,
);
let mut session_manager = SessionManager::new(
SessionDefaults {
model: AgentModel::Gpt54,
},
Arc::new(git::MockGitClient::new()),
state,
Vec::new(),
);
let session_id = "session-id".to_string();
let first_task_aborted = Arc::new(std::sync::atomic::AtomicBool::new(false));
let first_task_flag = Arc::clone(&first_task_aborted);
let first_task = tokio::spawn(async move {
struct AbortFlagGuard(Arc<std::sync::atomic::AtomicBool>);
impl Drop for AbortFlagGuard {
fn drop(&mut self) {
self.0.store(true, std::sync::atomic::Ordering::SeqCst);
}
}
let _abort_flag_guard = AbortFlagGuard(first_task_flag);
std::future::pending::<()>().await;
});
let second_task = tokio::spawn(async {});
session_manager.replace_title_generation_task(&session_id, 1, first_task);
tokio::task::yield_now().await;
session_manager.replace_title_generation_task(&session_id, 2, second_task);
tokio::time::timeout(Duration::from_secs(1), async {
while !first_task_aborted.load(std::sync::atomic::Ordering::SeqCst) {
tokio::task::yield_now().await;
}
})
.await
.expect("superseded task should abort promptly");
assert_eq!(session_manager.title_generation_tasks.len(), 1);
assert!(
session_manager
.title_generation_tasks
.contains_key(session_id.as_str())
);
}
#[tokio::test]
async fn test_clear_title_generation_task_if_matches_ignores_stale_generation() {
let state = SessionState::new(
HashMap::new(),
Vec::new(),
TableState::default(),
Arc::new(RealClock),
1,
0,
);
let mut session_manager = SessionManager::new(
SessionDefaults {
model: AgentModel::Gpt54,
},
Arc::new(git::MockGitClient::new()),
state,
Vec::new(),
);
let session_id = "session-id".to_string();
let task = tokio::spawn(async {});
session_manager.replace_title_generation_task(&session_id, 2, task);
session_manager.clear_title_generation_task_if_matches(&session_id, 1);
assert_eq!(session_manager.title_generation_tasks.len(), 1);
assert!(
session_manager
.title_generation_tasks
.contains_key(session_id.as_str())
);
}
#[tokio::test]
async fn test_clear_title_generation_task_if_matches_removes_matching_generation() {
let state = SessionState::new(
HashMap::new(),
Vec::new(),
TableState::default(),
Arc::new(RealClock),
1,
0,
);
let mut session_manager = SessionManager::new(
SessionDefaults {
model: AgentModel::Gpt54,
},
Arc::new(git::MockGitClient::new()),
state,
Vec::new(),
);
let session_id = "session-id".to_string();
let task = tokio::spawn(async {});
session_manager.replace_title_generation_task(&session_id, 2, task);
session_manager.clear_title_generation_task_if_matches(&session_id, 2);
assert!(session_manager.title_generation_tasks.is_empty());
}
#[tokio::test]
async fn test_start_staged_session_launches_bundle_and_clears_staged_drafts() {
let dir = tempdir().expect("failed to create temp dir");
let mut app = new_test_app_with_git(dir.path()).await;
let session_id = app
.create_draft_session()
.await
.expect("failed to create draft session");
app.stage_draft_message(&session_id, "First draft")
.await
.expect("failed to stage first draft");
app.stage_draft_message(&session_id, "Second draft")
.await
.expect("failed to stage second draft");
app.start_staged_session(&session_id)
.await
.expect("failed to start staged session");
assert_eq!(
app.sessions.sessions()[0].prompt,
"First draft\n\nSecond draft"
);
assert!(app.sessions.sessions()[0].draft_attachments.is_empty());
assert!(app.sessions.sessions()[0].folder.exists());
assert!(
app.sessions.sessions()[0]
.folder
.join(SESSION_DATA_DIR)
.is_dir()
);
let db_sessions = app
.services
.db()
.sessions()
.load_sessions()
.await
.expect("failed to load");
assert_eq!(db_sessions[0].prompt, "First draft\n\nSecond draft");
}
#[tokio::test]
async fn test_start_staged_session_rejects_stacked_draft_child() {
let dir = tempdir().expect("failed to create temp dir");
let mut app = new_test_app_with_git(dir.path()).await;
let parent_session_id = app.create_session().await.expect("failed to create parent");
let child_session_id = app
.create_stacked_draft_session(&parent_session_id)
.await
.expect("failed to create stacked draft session");
app.stage_draft_message(&child_session_id, "Stacked draft")
.await
.expect("failed to stage stacked draft message");
let result = app.start_staged_session(&child_session_id).await;
let error = result.expect_err("stacked draft should not start before restack");
assert!(
error.to_string().contains(
"Stacked draft sessions can only be started after their parent is merged"
)
);
let child_session = app
.sessions
.sessions()
.iter()
.find(|session| session.id == child_session_id)
.expect("missing child session");
assert_eq!(child_session.status, Status::Draft);
assert_eq!(
child_session.parent_session_id.as_deref(),
Some(parent_session_id.as_str())
);
}
#[tokio::test]
async fn test_delete_selected_draft_session_removes_staged_draft_metadata() {
let dir = tempdir().expect("failed to create temp dir");
let mut app = new_test_app_with_git(dir.path()).await;
let session_id = app
.create_draft_session()
.await
.expect("failed to create draft session");
app.stage_draft_message(
&session_id,
TurnPrompt {
attachments: vec![TurnPromptAttachment {
placeholder: "[Image #1]".to_string(),
local_image_path: dir.path().join("draft-image.png"),
}],
text: "First draft".to_string(),
text_source: TurnPromptTextSource::UserPrompt,
},
)
.await
.expect("failed to stage first draft");
let staged_draft_root = app.services.base_path().join(&session_id);
assert!(staged_draft_root.exists());
app.delete_selected_session().await;
assert!(!staged_draft_root.exists());
}
#[tokio::test]
async fn test_start_session_uses_full_prompt_text_as_title() {
let dir = tempdir().expect("failed to create temp dir");
let mut app = new_test_app_with_git(dir.path()).await;
let session_id = app
.create_session()
.await
.expect("failed to create session");
let prompt = "First line\nSecond line is intentionally long to avoid truncation.";
app.start_session(&session_id, prompt.to_string())
.await
.expect("failed to start session");
assert_eq!(app.sessions.sessions()[0].title, Some(prompt.to_string()));
let db_sessions = app
.services
.db()
.sessions()
.load_sessions()
.await
.expect("failed to load");
assert_eq!(db_sessions[0].title, Some(prompt.to_string()));
}
#[tokio::test]
async fn test_reply_first_message_uses_full_prompt_text_as_title() {
let dir = tempdir().expect("failed to create temp dir");
let mut app = new_test_app_with_git(dir.path()).await;
let session_id = app
.create_session()
.await
.expect("failed to create session");
let prompt = "Line one\nLine two with more words for title text";
let backend = create_mock_backend();
app.sessions
.reply_with_backend(
&app.services,
&session_id,
prompt,
Arc::new(backend),
AgentModel::Gemini3FlashPreview,
)
.await;
assert_eq!(app.sessions.sessions()[0].title, Some(prompt.to_string()));
let db_sessions = app
.services
.db()
.sessions()
.load_sessions()
.await
.expect("failed to load");
assert_eq!(db_sessions[0].title, Some(prompt.to_string()));
}
#[tokio::test]
async fn test_esc_deletes_blank_session() {
let dir = tempdir().expect("failed to create temp dir");
let mut app = new_test_app_with_git(dir.path()).await;
let session_id = app
.create_session()
.await
.expect("failed to create session");
let session_index = app
.session_index_for_id(&session_id)
.expect("missing session index");
let session_folder = app.sessions.sessions()[session_index].folder.clone();
assert!(session_folder.exists());
app.delete_selected_session().await;
assert!(app.sessions.sessions().is_empty());
assert!(!session_folder.exists());
let db_sessions = app
.services
.db()
.sessions()
.load_sessions()
.await
.expect("failed to load");
assert!(db_sessions.is_empty());
}
#[tokio::test]
async fn test_reply() {
let dir = tempdir().expect("failed to create temp dir");
let mut app = new_test_app_with_git(dir.path()).await;
create_and_start_session(&mut app, "Initial").await;
let session_id = app.sessions.sessions()[0].id.clone();
wait_for_status(&mut app, &session_id, Status::Review).await;
app.reply(&session_id, "Reply").await;
app.sessions.sync_from_handles();
let session = &app.sessions.sessions()[0];
let output = &session.output;
let activity_timestamps = app
.services
.db()
.activity()
.load_session_activity_timestamps()
.await
.expect("failed to load session activity timestamps");
assert!(output.contains("Reply"));
assert_eq!(activity_timestamps.len(), 1);
}
#[tokio::test]
async fn test_enqueue_message_pushes_prompt_onto_in_memory_queue() {
let dir = tempdir().expect("failed to create temp dir");
let mut app = new_test_app_with_git(dir.path()).await;
create_and_start_session(&mut app, "Initial").await;
let session_id = app.sessions.sessions()[0].id.clone();
wait_for_status(&mut app, &session_id, Status::Review).await;
set_session_status_for_test(&mut app, &session_id, Status::InProgress);
app.enqueue_message(&session_id, "queued reply")
.expect("enqueue_message should succeed for InProgress session");
let session = app
.sessions
.sessions()
.iter()
.find(|session| session.id == session_id)
.expect("session present");
assert_eq!(session.queued_messages, vec!["queued reply".to_string()]);
let handles = app
.sessions
.session_handles()
.get(session_id.as_str())
.expect("handles present");
let queued_len = handles.queued_messages.lock().expect("queue lock").len();
assert_eq!(queued_len, 1);
}
#[tokio::test]
async fn test_enqueue_message_survives_refresh_sessions_reducer_pass() {
let dir = tempdir().expect("failed to create temp dir");
let mut app = new_test_app_with_git(dir.path()).await;
create_and_start_session(&mut app, "Initial").await;
let session_id = app.sessions.sessions()[0].id.clone();
wait_for_status(&mut app, &session_id, Status::Review).await;
set_session_status_for_test(&mut app, &session_id, Status::InProgress);
app.enqueue_message(&session_id, "queued reply")
.expect("enqueue_message should succeed for InProgress session");
app.apply_app_events(AppEvent::RefreshSessions).await;
let session = app
.sessions
.sessions()
.iter()
.find(|session| session.id == session_id)
.expect("session present");
assert_eq!(
session.queued_messages,
vec!["queued reply".to_string()],
"queued_messages snapshot must be re-projected from handles after a RefreshSessions \
reducer pass instead of being wiped to an empty vec"
);
}
#[tokio::test]
async fn test_enqueue_message_rejects_empty_payload() {
let dir = tempdir().expect("failed to create temp dir");
let mut app = new_test_app_with_git(dir.path()).await;
create_and_start_session(&mut app, "Initial").await;
let session_id = app.sessions.sessions()[0].id.clone();
wait_for_status(&mut app, &session_id, Status::Review).await;
set_session_status_for_test(&mut app, &session_id, Status::InProgress);
let outcome = app.enqueue_message(&session_id, "");
let error = outcome.expect_err("empty payload should error");
assert!(matches!(
error,
crate::app::session::SessionError::Workflow(_)
));
let session = app
.sessions
.sessions()
.iter()
.find(|session| session.id == session_id)
.expect("session present");
assert!(session.queued_messages.is_empty());
}
#[tokio::test]
async fn test_selected_session() {
let dir = tempdir().expect("failed to create temp dir");
let mut app = new_test_app_with_git(dir.path()).await;
create_and_start_session(&mut app, "Test").await;
assert!(app.selected_session().is_some());
app.sessions.select_session_index(None);
assert!(app.selected_session().is_none());
}
#[tokio::test]
async fn test_delete_session() {
let dir = tempdir().expect("failed to create temp dir");
let mut app = new_test_app_with_git(dir.path()).await;
create_and_start_session(&mut app, "A").await;
let session_id = app.sessions.sessions()[0].id.clone();
let session_folder = app.sessions.sessions()[0].folder.clone();
app.sessions.set_at_mention_index_for_root(
session_folder.clone(),
vec![FileEntry {
is_dir: false,
path: "src/main.rs".to_string(),
}],
);
let session_update_versions = app.services.session_update_versions();
SessionTaskService::remove_session_update_version(
&session_update_versions,
session_id.as_str(),
);
let initial_version = SessionTaskService::next_session_update_version(
&session_update_versions,
session_id.as_str(),
);
assert_eq!(initial_version, 1);
app.delete_selected_session().await;
let reset_version = SessionTaskService::next_session_update_version(
&session_update_versions,
session_id.as_str(),
);
assert!(app.sessions.sessions().is_empty());
assert_eq!(app.sessions.selected_session_index(), None);
assert!(!session_folder.exists());
let db_sessions = app
.services
.db()
.sessions()
.load_sessions()
.await
.expect("failed to load");
assert!(db_sessions.is_empty());
assert!(
app.sessions
.at_mention_index_for_root(&session_folder)
.is_none()
);
assert_eq!(reset_version, 1);
SessionTaskService::remove_session_update_version(
&session_update_versions,
session_id.as_str(),
);
}
#[tokio::test]
async fn test_delete_selected_session_edge_cases() {
let dir = tempdir().expect("failed to create temp dir");
let mut app = new_test_app_with_git(dir.path()).await;
create_and_start_session(&mut app, "1").await;
create_and_start_session(&mut app, "2").await;
app.sessions.select_session_index(Some(99));
app.delete_selected_session().await;
assert_eq!(app.sessions.sessions().len(), 2);
app.sessions.select_session_index(None);
app.delete_selected_session().await;
assert_eq!(app.sessions.sessions().len(), 2);
}
#[tokio::test]
async fn test_delete_last_session_update_selection() {
let dir = tempdir().expect("failed to create temp dir");
let mut app = new_test_app_with_git(dir.path()).await;
create_and_start_session(&mut app, "1").await;
create_and_start_session(&mut app, "2").await;
app.sessions.select_session_index(Some(1));
app.delete_selected_session().await;
assert_eq!(app.sessions.sessions().len(), 1);
assert_eq!(app.sessions.selected_session_index(), Some(0));
app.delete_selected_session().await;
assert!(app.sessions.sessions().is_empty());
assert_eq!(app.sessions.selected_session_index(), None);
}
#[tokio::test]
async fn test_load_existing_sessions() {
let dir = tempdir().expect("failed to create temp dir");
let db = AppRepositories::in_memory().await;
let project_id = db
.projects()
.upsert_project("/tmp/test", None)
.await
.expect("failed to upsert project");
db.sessions()
.insert_session("12345678", "claude-opus-4-6", "main", "Done", project_id)
.await
.expect("failed to insert");
let session_dir = dir.path().join("12345678");
let data_dir = session_dir.join(SESSION_DATA_DIR);
std::fs::create_dir(&session_dir).expect("failed to create session dir");
std::fs::create_dir(&data_dir).expect("failed to create data dir");
db.sessions()
.update_session_prompt("12345678", "Existing")
.await
.expect("failed to update prompt");
db.sessions()
.append_session_output("12345678", "Output")
.await
.expect("failed to update output");
let app = new_test_app_with_db(
dir.path().to_path_buf(),
PathBuf::from("/tmp/test"),
None,
db,
)
.await;
assert_eq!(app.sessions.sessions().len(), 1);
assert_eq!(app.sessions.sessions()[0].id, "12345678");
assert_eq!(app.sessions.sessions()[0].model, AgentModel::ClaudeOpus48);
assert_eq!(app.sessions.sessions()[0].prompt, "Existing");
let output = app.sessions.sessions()[0].output.clone();
assert_eq!(output, "Output");
assert_eq!(app.sessions.selected_session_index(), Some(0));
}
#[tokio::test]
async fn test_create_session_uses_default_smart_model_setting_and_most_recent_permission_mode()
{
let dir = tempdir().expect("failed to create temp dir");
let db = AppRepositories::in_memory().await;
let project_id = db
.projects()
.upsert_project(&dir.path().to_string_lossy(), Some("main".to_string()))
.await
.expect("failed to upsert project");
db.sessions()
.insert_session(
"alpha0001",
"gemini-3-flash-preview",
"main",
"Done",
project_id,
)
.await
.expect("failed to insert alpha0001");
db.sessions()
.insert_session(
"beta00002",
AgentModel::ClaudeHaiku4520251001.as_str(),
"main",
"Done",
project_id,
)
.await
.expect("failed to insert beta00002");
db.settings()
.upsert_project_setting(
project_id,
SettingName::DefaultSmartModel,
AgentModel::ClaudeHaiku4520251001.as_str(),
)
.await
.expect("failed to upsert default smart model setting");
db.sessions()
.update_session_updated_at("alpha0001", 1_i64)
.await
.expect("failed to update alpha0001 timestamp");
db.sessions()
.update_session_updated_at("beta00002", 2_i64)
.await
.expect("failed to update beta00002 timestamp");
for session_id in ["alpha0001", "beta00002"] {
let session_dir = session_folder(dir.path(), session_id);
let data_dir = session_dir.join(SESSION_DATA_DIR);
std::fs::create_dir_all(&data_dir).expect("failed to create session data dir");
}
let mut app = new_test_app_with_git_and_db(dir.path(), db).await;
let created_session_id = app
.create_session()
.await
.expect("failed to create session");
let created_session = app
.sessions
.sessions()
.iter()
.find(|session| session.id == created_session_id)
.expect("missing created session");
assert_eq!(created_session.model, AgentModel::ClaudeHaiku4520251001);
}
#[tokio::test]
async fn test_load_existing_sessions_ordered_by_updated_at_desc() {
let dir = tempdir().expect("failed to create temp dir");
let db = AppRepositories::in_memory().await;
let project_id = db
.projects()
.upsert_project("/tmp/test", None)
.await
.expect("failed to upsert project");
db.sessions()
.insert_session("alpha000", "claude-opus-4-8", "main", "Done", project_id)
.await
.expect("failed to insert alpha000");
db.sessions()
.insert_session(
"beta0000",
"gemini-3-flash-preview",
"main",
"Done",
project_id,
)
.await
.expect("failed to insert beta0000");
db.sessions()
.update_session_updated_at("alpha000", 1_i64)
.await
.expect("failed to update alpha000 timestamp");
db.sessions()
.update_session_updated_at("beta0000", 2_i64)
.await
.expect("failed to update beta0000 timestamp");
for session_id in ["alpha000", "beta0000"] {
let session_dir = session_folder(dir.path(), session_id);
let data_dir = session_dir.join(SESSION_DATA_DIR);
std::fs::create_dir_all(&data_dir).expect("failed to create data dir");
}
let app = new_test_app_with_db(
dir.path().to_path_buf(),
PathBuf::from("/tmp/test"),
None,
db,
)
.await;
let session_names: Vec<&str> = app
.sessions
.sessions()
.iter()
.map(|session| session.id.as_str())
.collect();
assert_eq!(session_names, vec!["beta0000", "alpha000"]);
}
#[tokio::test]
async fn test_load_sessions_aggregates_daily_activity() {
let dir = tempdir().expect("failed to create temp dir");
let db = AppRepositories::in_memory().await;
let project_id = db
.projects()
.upsert_project("/tmp/test", None)
.await
.expect("failed to upsert project");
db.sessions()
.insert_session("alpha000", "claude-opus-4-8", "main", "Done", project_id)
.await
.expect("failed to insert alpha000");
db.sessions()
.insert_session("beta0000", "claude-opus-4-8", "main", "Done", project_id)
.await
.expect("failed to insert beta0000");
db.sessions()
.insert_session("gamma000", "claude-opus-4-8", "main", "Done", project_id)
.await
.expect("failed to insert gamma000");
let seconds_per_day = 86_400_i64;
let day_key_one = 10_i64;
let day_key_two = 11_i64;
db.sessions()
.update_session_created_at("alpha000", day_key_one * seconds_per_day + 10)
.await
.expect("failed to update alpha000 created_at");
db.sessions()
.update_session_created_at("beta0000", day_key_one * seconds_per_day + 600)
.await
.expect("failed to update beta0000 created_at");
db.sessions()
.update_session_created_at("gamma000", day_key_two * seconds_per_day + 50)
.await
.expect("failed to update gamma000 created_at");
db.activity()
.clear_session_activity()
.await
.expect("failed to clear session activity");
db.activity()
.backfill_session_activity_from_sessions()
.await
.expect("failed to backfill session activity from session rows");
let working_dir = PathBuf::from("/tmp/test");
let mut handles: HashMap<SessionId, SessionHandles> = HashMap::new();
let mut expected_activity_by_day: BTreeMap<i64, u32> = BTreeMap::new();
for timestamp_seconds in [
day_key_one * seconds_per_day + 10,
day_key_one * seconds_per_day + 600,
day_key_two * seconds_per_day + 50,
] {
let day_key = activity_heatmap::activity_day_key_local(timestamp_seconds);
let day_count = expected_activity_by_day.entry(day_key).or_insert(0);
*day_count = day_count.saturating_add(1);
}
let expected_activity: Vec<DailyActivity> = expected_activity_by_day
.into_iter()
.map(|(day_key, session_count)| DailyActivity {
day_key,
session_count,
})
.collect();
let fs_client = fs::RealFsClient;
let (sessions, stats_activity, _) = SessionManager::load_sessions_with_fs_client(
dir.path(),
&db,
project_id,
&working_dir,
&mut handles,
&fs_client,
None,
)
.await;
assert_eq!(sessions.len(), 3);
assert_eq!(stats_activity, expected_activity);
}
#[tokio::test]
async fn test_load_sessions_keeps_daily_activity_after_session_deletion() {
let dir = tempdir().expect("failed to create temp dir");
let db = AppRepositories::in_memory().await;
let project_id = db
.projects()
.upsert_project("/tmp/test", None)
.await
.expect("failed to upsert project");
db.sessions()
.insert_session("alpha000", "claude-opus-4-8", "main", "Done", project_id)
.await
.expect("failed to insert alpha000");
db.sessions()
.insert_session("beta0000", "claude-opus-4-8", "main", "Done", project_id)
.await
.expect("failed to insert beta0000");
db.activity()
.insert_session_creation_activity_at("alpha000", 10)
.await
.expect("failed to persist first activity event");
db.activity()
.insert_session_creation_activity_at("beta0000", 20)
.await
.expect("failed to persist second activity event");
db.sessions()
.delete_session("alpha000")
.await
.expect("failed to delete alpha000");
let working_dir = PathBuf::from("/tmp/test");
let mut handles: HashMap<SessionId, SessionHandles> = HashMap::new();
let fs_client = fs::RealFsClient;
let (sessions, stats_activity, _) = SessionManager::load_sessions_with_fs_client(
dir.path(),
&db,
project_id,
&working_dir,
&mut handles,
&fs_client,
None,
)
.await;
assert_eq!(sessions.len(), 1);
let total_activity_count: u32 = stats_activity
.iter()
.map(|daily_activity| daily_activity.session_count)
.sum();
assert_eq!(total_activity_count, 2);
}
#[tokio::test]
async fn test_refresh_sessions_if_needed_reloads_and_preserves_selection() {
let dir = tempdir().expect("failed to create temp dir");
let db = AppRepositories::in_memory().await;
let project_id = db
.projects()
.upsert_project("/tmp/test", None)
.await
.expect("failed to upsert project");
db.sessions()
.insert_session(
"alpha000",
"gemini-3-flash-preview",
"main",
"InProgress",
project_id,
)
.await
.expect("failed to insert alpha000");
db.sessions()
.insert_session("beta0000", "claude-opus-4-8", "main", "Done", project_id)
.await
.expect("failed to insert beta0000");
db.sessions()
.update_session_updated_at("alpha000", 1)
.await
.expect("failed to set alpha000 timestamp");
db.sessions()
.update_session_updated_at("beta0000", 2)
.await
.expect("failed to set beta0000 timestamp");
for session_id in ["alpha000", "beta0000"] {
let session_dir = session_folder(dir.path(), session_id);
let data_dir = session_dir.join(SESSION_DATA_DIR);
std::fs::create_dir_all(&data_dir).expect("failed to create data dir");
}
let mut app = new_test_app_with_db(
dir.path().to_path_buf(),
PathBuf::from("/tmp/test"),
None,
db,
)
.await;
app.sessions.select_session_index(Some(1));
app.services
.db()
.sessions()
.update_session_status_with_timing_at("alpha000", "Done", 0)
.await
.expect("failed to update session status");
app.refresh_sessions_now().await;
assert_eq!(app.sessions.sessions()[0].id, "alpha000");
let selected_index = app
.sessions
.selected_session_index()
.expect("missing selection");
assert_eq!(app.sessions.sessions()[selected_index].id, "alpha000");
}
#[tokio::test]
async fn test_refresh_sessions_if_needed_remaps_view_mode_index() {
let dir = tempdir().expect("failed to create temp dir");
let db = AppRepositories::in_memory().await;
let project_id = db
.projects()
.upsert_project("/tmp/test", None)
.await
.expect("failed to upsert project");
db.sessions()
.insert_session(
"alpha000",
"gemini-3-flash-preview",
"main",
"InProgress",
project_id,
)
.await
.expect("failed to insert alpha000");
db.sessions()
.insert_session("beta0000", "claude-opus-4-8", "main", "Done", project_id)
.await
.expect("failed to insert beta0000");
db.sessions()
.update_session_updated_at("alpha000", 1)
.await
.expect("failed to set alpha000 timestamp");
db.sessions()
.update_session_updated_at("beta0000", 2)
.await
.expect("failed to set beta0000 timestamp");
for session_id in ["alpha000", "beta0000"] {
let session_dir = session_folder(dir.path(), session_id);
let data_dir = session_dir.join(SESSION_DATA_DIR);
std::fs::create_dir_all(&data_dir).expect("failed to create data dir");
}
let mut app = new_test_app_with_db(
dir.path().to_path_buf(),
PathBuf::from("/tmp/test"),
None,
db,
)
.await;
let selected_session_id = app.sessions.sessions()[1].id.clone();
app.mode = AppMode::View {
review_status_message: None,
review_text: None,
session_id: selected_session_id.clone(),
scroll_offset: None,
};
app.services
.db()
.sessions()
.update_session_status_with_timing_at("alpha000", "Done", 0)
.await
.expect("failed to update session status");
app.refresh_sessions_now().await;
assert_eq!(app.sessions.sessions()[0].id, "alpha000");
assert!(matches!(app.mode, AppMode::View { .. }));
if let AppMode::View { session_id, .. } = app.mode {
assert_eq!(session_id, selected_session_id);
}
}
#[tokio::test]
async fn test_load_sessions_invalid_path() {
let path = PathBuf::from("/invalid/path/that/does/not/exist");
let app = new_test_app(path).await;
assert!(app.sessions.sessions().is_empty());
}
#[tokio::test]
async fn test_load_done_session_without_folder_kept() {
let dir = tempdir().expect("failed to create temp dir");
let db = AppRepositories::in_memory().await;
let project_id = db
.projects()
.upsert_project("/tmp/test", None)
.await
.expect("failed to upsert project");
db.sessions()
.insert_session(
"missing01",
"gemini-3-flash-preview",
"main",
"Done",
project_id,
)
.await
.expect("failed to insert");
let app = new_test_app_with_db(
dir.path().to_path_buf(),
PathBuf::from("/tmp/test"),
None,
db,
)
.await;
assert_eq!(app.sessions.sessions().len(), 1);
assert_eq!(app.sessions.sessions()[0].id, "missing01");
assert_eq!(app.sessions.sessions()[0].status, Status::Done);
}
#[tokio::test]
async fn test_load_in_progress_session_without_folder_skipped() {
let dir = tempdir().expect("failed to create temp dir");
let db = AppRepositories::in_memory().await;
let project_id = db
.projects()
.upsert_project("/tmp/test", None)
.await
.expect("failed to upsert project");
db.sessions()
.insert_session(
"missing02",
"gemini-3-flash-preview",
"main",
"InProgress",
project_id,
)
.await
.expect("failed to insert");
let app = new_test_app_with_db(
dir.path().to_path_buf(),
PathBuf::from("/tmp/test"),
None,
db,
)
.await;
assert!(app.sessions.sessions().is_empty());
}
#[tokio::test]
async fn test_load_sessions_uses_persisted_size_for_non_terminal_status() {
let dir = tempdir().expect("failed to create temp dir");
let mut app = new_test_app_with_git(dir.path()).await;
let session_id = app
.create_session()
.await
.expect("failed to create session");
app.services
.db()
.sessions()
.update_session_diff_stats(8, 3, &session_id, "S")
.await
.expect("failed to update size");
let session_index = app
.session_index_for_id(&session_id)
.expect("missing created session");
let session_folder = app.sessions.sessions()[session_index].folder.clone();
let changed_lines = "line\n".repeat(700);
std::fs::write(session_folder.join("size-test.txt"), changed_lines)
.expect("failed to write test file");
let fs_client = fs::RealFsClient;
let (reloaded_sessions, _, _) = SessionManager::load_sessions_with_fs_client(
app.services.base_path(),
app.services.db(),
app.projects.active_project_id(),
app.projects.working_dir(),
app.sessions.session_handles_mut(),
&fs_client,
None,
)
.await;
let db_sessions = app
.services
.db()
.sessions()
.load_sessions()
.await
.expect("failed to load");
let reloaded_session = reloaded_sessions
.iter()
.find(|session| session.id == session_id)
.expect("missing reloaded session");
let db_session = db_sessions
.iter()
.find(|session| session.id == session_id)
.expect("missing persisted session");
assert_eq!(reloaded_session.size, SessionSize::S);
assert_eq!(reloaded_session.stats.added_lines, 8);
assert_eq!(reloaded_session.stats.deleted_lines, 3);
assert_eq!(db_session.added_lines, 8);
assert_eq!(db_session.deleted_lines, 3);
assert_eq!(db_session.size, "S");
}
#[tokio::test]
async fn test_reply_turn_completion_persists_session_size() {
let dir = tempdir().expect("failed to create temp dir");
let mut app = new_test_app_with_git(dir.path()).await;
let session_id = app
.create_session()
.await
.expect("failed to create session");
let mut backend = MockAgentBackend::new();
backend.expect_build_command().returning(|request| {
let mut command = Command::new("sh");
command
.args([
"-lc",
"yes line | head -n 20 > turn-size-test.txt; echo turn-complete",
])
.current_dir(request.folder)
.stdout(Stdio::piped())
.stderr(Stdio::null());
Ok(command)
});
app.sessions
.reply_with_backend(
&app.services,
&session_id,
"compute size after turn",
Arc::new(backend),
AgentModel::ClaudeOpus48,
)
.await;
wait_for_status(&mut app, &session_id, Status::AgentReview).await;
app.process_pending_app_events().await;
let db_sessions = app
.services
.db()
.sessions()
.load_sessions()
.await
.expect("failed to load");
let session = app
.sessions
.sessions()
.iter()
.find(|session| session.id == session_id)
.expect("missing in-memory session");
let db_session = db_sessions
.iter()
.find(|db_session| db_session.id == session_id)
.expect("missing persisted session");
assert_eq!(session.size, SessionSize::S);
assert_eq!(session.stats.added_lines, 20);
assert_eq!(session.stats.deleted_lines, 0);
assert_eq!(db_session.added_lines, 20);
assert_eq!(db_session.deleted_lines, 0);
assert_eq!(db_session.size, "S");
}
#[tokio::test]
async fn test_load_sessions_uses_persisted_size_for_done_status() {
let dir = tempdir().expect("failed to create temp dir");
let mut app = new_test_app_with_git(dir.path()).await;
let session_id = app
.create_session()
.await
.expect("failed to create session");
app.services
.db()
.sessions()
.update_session_diff_stats(21, 9, &session_id, "L")
.await
.expect("failed to update size");
app.services
.db()
.sessions()
.update_session_status_with_timing_at(&session_id, "Done", 0)
.await
.expect("failed to update status");
let session_index = app
.session_index_for_id(&session_id)
.expect("missing created session");
let session_folder = app.sessions.sessions()[session_index].folder.clone();
let changed_lines = "line\n".repeat(700);
std::fs::write(session_folder.join("done-size-test.txt"), changed_lines)
.expect("failed to write test file");
let fs_client = fs::RealFsClient;
let (reloaded_sessions, _, _) = SessionManager::load_sessions_with_fs_client(
app.services.base_path(),
app.services.db(),
app.projects.active_project_id(),
app.projects.working_dir(),
app.sessions.session_handles_mut(),
&fs_client,
None,
)
.await;
let db_sessions = app
.services
.db()
.sessions()
.load_sessions()
.await
.expect("failed to load");
let reloaded_session = reloaded_sessions
.iter()
.find(|session| session.id == session_id)
.expect("missing reloaded session");
let db_session = db_sessions
.iter()
.find(|session| session.id == session_id)
.expect("missing persisted session");
assert_eq!(reloaded_session.status, Status::Done);
assert_eq!(reloaded_session.size, SessionSize::L);
assert_eq!(reloaded_session.stats.added_lines, 21);
assert_eq!(reloaded_session.stats.deleted_lines, 9);
assert_eq!(db_session.added_lines, 21);
assert_eq!(db_session.deleted_lines, 9);
assert_eq!(db_session.size, "L");
}
#[tokio::test]
async fn test_spawn_integration() {
let dir = tempdir().expect("failed to create temp dir");
let db = AppRepositories::in_memory().await;
let mut app = new_test_app_with_git_and_db(dir.path(), db).await;
let turn_count = Arc::new(Mutex::new(0usize));
let (done_tx, mut done_rx) = tokio::sync::mpsc::unbounded_channel::<()>();
let mut mock_channel = MockAgentChannel::new();
let turn_count_capture = Arc::clone(&turn_count);
let done_capture = done_tx.clone();
mock_channel
.expect_run_turn()
.returning(move |_, req, _event_tx| {
let turn_index = {
let mut count = turn_count_capture.lock().expect("lock poisoned");
let current = *count;
*count += 1;
current
};
let delta_text = if turn_index == 0 {
assert!(
matches!(req.request_kind, AgentRequestKind::SessionStart),
"expected AgentRequestKind::SessionStart on first turn"
);
format!("--prompt {}\n", req.prompt)
} else {
assert!(
matches!(req.request_kind, AgentRequestKind::SessionResume { .. }),
"expected AgentRequestKind::SessionResume on second turn"
);
format!("--prompt {} --resume latest\n", req.prompt)
};
let done = done_capture.clone();
Box::pin(async move {
let _ = done.send(());
Ok(TurnResult {
assistant_message: AgentResponse::plain(&delta_text),
context_reset: false,
input_tokens: 0,
output_tokens: 0,
provider_conversation_id: None,
})
})
});
mock_channel
.expect_shutdown_session()
.returning(|_| Box::pin(async { Ok(()) }));
let session_id = app
.create_session()
.await
.expect("failed to create session");
app.sessions
.worker_service
.test_agent_channels
.insert(session_id.clone().into(), Arc::new(mock_channel));
app.sessions
.reply(&app.services, &session_id, "SpawnInit")
.await;
done_rx.recv().await.expect("first turn completion signal");
wait_for_status(&mut app, &session_id, Status::Review).await;
wait_for_output_contains(&mut app, &session_id, "SpawnInit", 200).await;
{
app.sessions.sync_from_handles();
let session = &app.sessions.sessions()[0];
let output = session.output.clone();
assert!(output.contains("--prompt"));
assert!(output.contains("SpawnInit"));
assert!(!output.contains("--resume"));
assert_eq!(session.status, Status::Review);
}
let session_id = app.sessions.sessions()[0].id.clone();
app.sessions
.reply(&app.services, &session_id, "SpawnReply")
.await;
done_rx.recv().await.expect("second turn completion signal");
wait_for_output_contains(&mut app, &session_id, "--resume", 200).await;
wait_for_status(&mut app, &session_id, Status::Review).await;
{
app.sessions.sync_from_handles();
let session = &app.sessions.sessions()[0];
let output = session.output.clone();
assert!(output.contains("SpawnReply"));
assert!(output.contains("--resume"));
assert!(output.contains("latest"));
assert_eq!(session.status, Status::Review);
}
}
#[tokio::test]
async fn test_reply_with_backend_replays_history_once_after_model_switch() {
let dir = tempdir().expect("failed to create temp dir");
let db = AppRepositories::in_memory().await;
let mut app = new_test_app_with_git_and_db(dir.path(), db).await;
let session_id = app
.create_session()
.await
.expect("failed to create session");
let initial_output = " › Initial prompt\n\nmock-start\n".to_string();
if let Some(session) = app
.sessions
.sessions_mut()
.iter_mut()
.find(|session| session.id == session_id)
{
session.output.clone_from(&initial_output);
session.prompt = "Initial prompt".to_string();
session.status = Status::Review;
}
if let Some(handles) = app.sessions.session_handles().get(session_id.as_str()) {
if let Ok(mut output) = handles.output.lock() {
output.clone_from(&initial_output);
}
if let Ok(mut status) = handles.status.lock() {
*status = Status::Review;
}
}
app.services
.db()
.sessions()
.update_session_prompt(&session_id, "Initial prompt")
.await
.expect("failed to persist initial prompt");
app.set_session_model(&session_id, AgentModel::ClaudeSonnet46)
.await
.expect("failed to switch model");
let captured_session_outputs: Arc<Mutex<Vec<Option<String>>>> =
Arc::new(Mutex::new(Vec::new()));
let (done_tx, mut done_rx) = tokio::sync::mpsc::unbounded_channel::<()>();
let mut mock_channel = MockAgentChannel::new();
let captured = Arc::clone(&captured_session_outputs);
let done_capture = done_tx.clone();
mock_channel.expect_run_turn().returning(move |_, req, _| {
if let AgentRequestKind::SessionResume { session_output } = req.request_kind {
captured.lock().expect("lock poisoned").push(session_output);
}
let done = done_capture.clone();
Box::pin(async move {
let _ = done.send(());
Ok(TurnResult {
assistant_message: AgentResponse::plain(""),
context_reset: false,
input_tokens: 0,
output_tokens: 0,
provider_conversation_id: None,
})
})
});
mock_channel
.expect_shutdown_session()
.returning(|_| Box::pin(async { Ok(()) }));
app.sessions
.worker_service
.test_agent_channels
.insert(session_id.clone().into(), Arc::new(mock_channel));
app.sessions
.reply(&app.services, &session_id, "Switch reply")
.await;
done_rx.recv().await.expect("first turn completion signal");
wait_for_status(&mut app, &session_id, Status::Review).await;
app.sessions
.reply(&app.services, &session_id, "Second reply")
.await;
done_rx.recv().await.expect("second turn completion signal");
wait_for_status(&mut app, &session_id, Status::Review).await;
let outputs = captured_session_outputs
.lock()
.expect("lock poisoned")
.clone();
assert_eq!(outputs.len(), 2, "expected exactly two Resume turns");
let first_session_output = outputs[0]
.as_deref()
.expect("first reply should include session_output");
assert!(
first_session_output.contains("Initial prompt"),
"first reply should replay history containing 'Initial prompt'"
);
assert!(
outputs[1].is_none(),
"second reply should not replay history"
);
}
#[tokio::test]
async fn test_reply_with_backend_replays_history_after_app_restart_for_review_session() {
let dir = tempdir().expect("failed to create temp dir");
let db = AppRepositories::in_memory().await;
let mut first_app = new_test_app_with_git_and_db(dir.path(), db.clone()).await;
let session_id = first_app
.create_session()
.await
.expect("failed to create session");
let start_backend = create_mock_backend();
first_app
.sessions
.reply_with_backend(
&first_app.services,
&session_id,
"Initial prompt",
Arc::new(start_backend),
AgentModel::ClaudeSonnet46,
)
.await;
wait_for_status(&mut first_app, &session_id, Status::Review).await;
first_app.sessions.sync_from_handles();
let first_run_output = first_app
.sessions
.sessions()
.iter()
.find(|session| session.id == session_id)
.map(|session| session.output.clone())
.expect("missing persisted session");
assert!(first_run_output.contains("Initial prompt"));
assert!(first_run_output.contains("mock-start"));
drop(first_app);
let mut resumed_app = new_test_app_with_git_and_db(dir.path(), db).await;
let resumed_session = resumed_app
.sessions
.sessions()
.iter()
.find(|session| session.id == session_id)
.expect("missing resumed session");
assert_eq!(resumed_session.status, Status::Review);
let mut resume_backend = MockAgentBackend::new();
resume_backend.expect_build_command().returning(|request| {
assert!(request.request_kind.is_resume());
let session_output = request
.request_kind
.session_output()
.expect("expected replayed session output after restart");
assert!(session_output.contains("Initial prompt"));
assert!(session_output.contains("mock-start"));
let mut cmd = Command::new("sh");
cmd.arg("-c")
.arg(
"printf '{\"answer\":\"replayed-after-restart\",\"questions\":[],\"summary\":\
null}'",
)
.current_dir(request.folder)
.stdout(Stdio::piped())
.stderr(Stdio::null());
Ok(cmd)
});
resumed_app
.sessions
.reply_with_backend(
&resumed_app.services,
&session_id,
"Restart reply",
Arc::new(resume_backend),
AgentModel::ClaudeSonnet46,
)
.await;
wait_for_output_contains(
&mut resumed_app,
&session_id,
"replayed-after-restart",
2000,
)
.await;
wait_for_status(&mut resumed_app, &session_id, Status::Review).await;
}
#[tokio::test]
async fn test_spawn_session_task_auto_commits_changes() {
let dir = tempdir().expect("failed to create temp dir");
let db = AppRepositories::in_memory().await;
let mut app = new_test_app_with_git_and_db(dir.path(), db).await;
let repo_root = dir.path().to_path_buf();
let mut mock_git_client = git::MockGitClient::new();
allow_detect_git_info(&mut mock_git_client);
mock_git_client
.expect_find_git_repo_root()
.times(0..)
.returning(move |_| {
let repo_root = repo_root.clone();
Box::pin(async move { Some(repo_root) })
});
mock_git_client
.expect_create_worktree()
.times(1)
.returning(|_, _, _, _| Box::pin(async { Ok(()) }));
mock_git_client
.expect_is_worktree_clean()
.times(1)
.returning(|_| Box::pin(async { Ok(false) }));
mock_git_client
.expect_has_commits_since()
.times(1)
.returning(|_, _| Box::pin(async { Ok(true) }));
mock_git_client
.expect_head_commit_message()
.times(1)
.returning(|_| Box::pin(async { Ok(Some("Existing session commit".to_string())) }));
mock_git_client
.expect_commit_all_preserving_single_commit()
.times(1)
.returning(|_, _, _, _, _| Box::pin(async { Ok(()) }));
mock_git_client
.expect_head_short_hash()
.times(1)
.returning(|_| Box::pin(async { Ok("abc1234".to_string()) }));
mock_git_client
.expect_diff()
.times(2)
.returning(|_, _| Box::pin(async { Ok(String::new()) }));
mock_git_client
.expect_fetch_remote()
.times(0..)
.returning(|_| Box::pin(async { Ok(()) }));
mock_git_client
.expect_branch_tracking_statuses()
.times(0..)
.returning(|_| Box::pin(async { Ok(HashMap::new()) }));
mock_git_client
.expect_get_ref_ahead_behind()
.times(0..)
.returning(|_, _, _| Box::pin(async { Ok((0, 0)) }));
install_mock_git_client(&mut app, mock_git_client);
let mut mock = MockAgentBackend::new();
mock.expect_build_command().returning(|request| {
let mut cmd = Command::new("bash");
cmd.arg("-c")
.arg(
"echo auto-content > auto-committed.txt; printf '{\"answer\":\"Auto commit \
done\",\"questions\":[],\"summary\":null}'",
)
.current_dir(request.folder)
.stdout(Stdio::piped())
.stderr(Stdio::null());
Ok(cmd)
});
let session_id = app
.create_session()
.await
.expect("failed to create session");
app.sessions
.reply_with_backend(
&app.services,
&session_id,
"AutoCommit",
Arc::new(mock),
AgentModel::ClaudeSonnet46,
)
.await;
wait_for_status(&mut app, &session_id, Status::Review).await;
let session = &app.sessions.sessions()[0];
let output = session.output.clone();
let workflow_notice = session.workflow_notice.as_deref();
assert!(
!output.contains("[Commit] committed with hash"),
"commit completion should not be persisted, got: {output}"
);
assert_eq!(
workflow_notice,
Some("[Commit] committed with hash `abc1234`")
);
}
#[tokio::test]
async fn test_commit_changes_reuses_existing_session_commit_message_in_tests() {
let session_folder = PathBuf::from("/tmp/session-worktree");
let mut mock_git_client = git::MockGitClient::new();
let mut sequence = mockall::Sequence::new();
mock_git_client
.expect_is_worktree_clean()
.times(1)
.in_sequence(&mut sequence)
.returning(|_| Box::pin(async { Ok(false) }));
mock_git_client
.expect_has_commits_since()
.times(1)
.in_sequence(&mut sequence)
.returning(|_, _| Box::pin(async { Ok(true) }));
mock_git_client
.expect_head_commit_message()
.times(1)
.in_sequence(&mut sequence)
.returning(|_| Box::pin(async { Ok(Some("Refine session work".to_string())) }));
mock_git_client
.expect_commit_all_preserving_single_commit()
.times(1)
.withf(|_, base_branch, commit_message, strategy, no_verify| {
base_branch == "main"
&& commit_message == "Refine session work"
&& *strategy == git::SingleCommitMessageStrategy::Replace
&& !*no_verify
})
.in_sequence(&mut sequence)
.returning(|_, _, _, _, _| Box::pin(async { Ok(()) }));
mock_git_client
.expect_head_short_hash()
.times(1)
.in_sequence(&mut sequence)
.returning(|_| Box::pin(async { Ok("def5678".to_string()) }));
let outcome = SessionTaskService::commit_session_changes(
&mock_git_client,
&session_folder,
"main",
AgentModel::ClaudeSonnet46,
false,
false,
)
.await
.expect("failed to commit existing session message");
assert_eq!(outcome.commit_hash, "def5678");
assert_eq!(outcome.commit_message, "Refine session work");
}
#[tokio::test]
async fn test_spawn_session_task_skips_commit_when_nothing_to_commit() {
let dir = tempdir().expect("failed to create temp dir");
let mut app = new_test_app_with_git(dir.path()).await;
let repo_root = dir.path().to_path_buf();
let mut mock_git_client = git::MockGitClient::new();
allow_detect_git_info(&mut mock_git_client);
mock_git_client
.expect_find_git_repo_root()
.times(0..)
.returning(move |_| {
let repo_root = repo_root.clone();
Box::pin(async move { Some(repo_root) })
});
mock_git_client
.expect_create_worktree()
.times(1)
.returning(|_, worktree_path, _, _| {
Box::pin(async move {
let fs_client = create_passthrough_mock_fs_client();
fs_client
.create_dir_all(worktree_path.clone())
.await
.map_err(|error| {
git::GitError::OutputParse(format!(
"Failed to create mock worktree: {error}"
))
})?;
fs_client
.create_dir_all(worktree_path.join(SESSION_DATA_DIR))
.await
.map_err(|error| {
git::GitError::OutputParse(format!(
"Failed to create mock worktree data dir: {error}"
))
})?;
Ok(())
})
});
mock_git_client
.expect_is_worktree_clean()
.times(1)
.returning(|_| Box::pin(async { Ok(true) }));
mock_git_client
.expect_diff()
.times(0..)
.returning(|_, _| Box::pin(async { Ok(String::new()) }));
mock_git_client
.expect_fetch_remote()
.times(0..)
.returning(|_| Box::pin(async { Ok(()) }));
mock_git_client
.expect_branch_tracking_statuses()
.times(0..)
.returning(|_| Box::pin(async { Ok(HashMap::new()) }));
mock_git_client
.expect_get_ref_ahead_behind()
.times(0..)
.returning(|_, _, _| Box::pin(async { Ok((0, 0)) }));
install_mock_git_client(&mut app, mock_git_client);
let mut mock = MockAgentBackend::new();
mock.expect_build_command().returning(|request| {
let mut cmd = Command::new("sh");
cmd.arg("-c")
.arg("printf '{\"answer\":\"no-changes\",\"questions\":[],\"summary\":null}'")
.current_dir(request.folder)
.stdout(Stdio::piped())
.stderr(Stdio::null());
Ok(cmd)
});
let session_id = app
.create_session()
.await
.expect("failed to create session");
app.sessions
.reply_with_backend(
&app.services,
&session_id,
"NoChanges",
Arc::new(mock),
AgentModel::ClaudeOpus48,
)
.await;
wait_for_status(&mut app, &session_id, Status::Review).await;
let session = &app.sessions.sessions()[0];
let output = session.output.clone();
let workflow_notice = session.workflow_notice.as_deref();
assert!(
!output.contains("[Commit] No changes to commit."),
"no-op commit output should not be persisted when nothing to commit"
);
assert_eq!(workflow_notice, Some("[Commit] No changes to commit."));
assert!(
!output.contains("[Commit Error]"),
"should not contain commit error when nothing to commit"
);
}
#[tokio::test]
async fn test_next_tab() {
let dir = tempdir().expect("failed to create temp dir");
let mut app = new_test_app(dir.path().to_path_buf()).await;
assert_eq!(app.tabs.current(), Tab::Projects);
app.next_tab();
assert_eq!(app.tabs.current(), Tab::Sessions);
app.next_tab();
assert_eq!(app.tabs.current(), Tab::Review);
app.next_tab();
assert_eq!(app.tabs.current(), Tab::Settings);
app.next_tab();
assert_eq!(app.tabs.current(), Tab::Projects);
}
#[tokio::test]
async fn test_next_tab_includes_tasks_when_active_project_has_roadmap() {
let dir = tempdir().expect("failed to create temp dir");
let database = AppRepositories::in_memory().await;
let mut app = new_test_app_with_db(
dir.path().to_path_buf(),
dir.path().to_path_buf(),
None,
database,
)
.await;
assert_eq!(app.tabs.current(), Tab::Projects);
app.next_tab();
assert_eq!(app.tabs.current(), Tab::Sessions);
app.next_tab();
assert_eq!(app.tabs.current(), Tab::Review);
app.next_tab();
assert_eq!(app.tabs.current(), Tab::Settings);
app.next_tab();
assert_eq!(app.tabs.current(), Tab::Projects);
}
#[tokio::test]
async fn test_create_session_without_git() {
let dir = tempdir().expect("failed to create temp dir");
let mut app = new_test_app(dir.path().to_path_buf()).await;
let result = app.create_session().await;
assert!(result.is_err());
assert!(
result
.expect_err("should be error")
.to_string()
.contains("Git branch is required")
);
assert!(app.sessions.sessions().is_empty());
}
#[tokio::test]
async fn test_create_session_with_git_no_actual_repo() {
let dir = tempdir().expect("failed to create temp dir");
let db = AppRepositories::in_memory().await;
let mut app = new_test_app_with_db(
dir.path().to_path_buf(),
PathBuf::from("/tmp/test"),
Some("main".to_string()),
db,
)
.await;
let mut mock_git_client = git::MockGitClient::new();
mock_git_client
.expect_find_git_repo_root()
.times(1)
.returning(|_| Box::pin(async { None }));
install_mock_git_client(&mut app, mock_git_client);
let result = app.create_session().await;
assert!(result.is_err());
assert!(
result
.expect_err("should be error")
.to_string()
.contains("git repository root")
);
}
#[tokio::test]
async fn test_create_session_cleans_up_on_error() {
let dir = tempdir().expect("failed to create temp dir");
let db = AppRepositories::in_memory().await;
let mut app = new_test_app_with_db(
dir.path().to_path_buf(),
PathBuf::from("/tmp/test"),
Some("main".to_string()),
db,
)
.await;
let repo_root = dir.path().to_path_buf();
let mut mock_git_client = git::MockGitClient::new();
allow_detect_git_info(&mut mock_git_client);
mock_git_client
.expect_find_git_repo_root()
.times(1)
.returning(move |_| {
let repo_root = repo_root.clone();
Box::pin(async move { Some(repo_root) })
});
mock_git_client
.expect_create_worktree()
.times(1)
.returning(|_, _, _, _| {
Box::pin(async {
Err(git::GitError::OutputParse(
"mock create_worktree failed".to_string(),
))
})
});
install_mock_git_client(&mut app, mock_git_client);
let result = app.create_session().await;
assert!(result.is_err());
assert_eq!(app.sessions.sessions().len(), 0);
let entries = std::fs::read_dir(dir.path())
.expect("failed to read dir")
.count();
assert_eq!(entries, 0, "Session folder should be cleaned up on error");
}
#[tokio::test]
async fn test_delete_session_without_git() {
let dir = tempdir().expect("failed to create temp dir");
let mut app = new_test_app(dir.path().to_path_buf()).await;
add_manual_session(&mut app, dir.path(), "manual01", "Test");
app.delete_selected_session().await;
assert_eq!(app.sessions.sessions().len(), 0);
}
#[tokio::test]
async fn test_merge_session_no_git() {
let dir = tempdir().expect("failed to create temp dir");
let mut app = new_test_app(dir.path().to_path_buf()).await;
add_manual_session(&mut app, dir.path(), "manual01", "Test");
let result = app.merge_session("manual01").await;
assert!(result.is_err());
assert!(
result
.expect_err("should be error")
.to_string()
.contains("No git worktree")
);
}
#[tokio::test]
async fn test_merge_session_invalid_id() {
let dir = tempdir().expect("failed to create temp dir");
let mut app = new_test_app(dir.path().to_path_buf()).await;
let result = app.merge_session("missing").await;
assert!(result.is_err());
assert!(
result
.expect_err("should be error")
.to_string()
.contains("Session not found")
);
}
#[tokio::test]
async fn test_merge_session_removes_worktree_and_branch_after_success() {
let dir = tempdir().expect("failed to create temp dir");
let mut app = new_test_app_with_git(dir.path()).await;
let session_id = app
.create_session()
.await
.expect("failed to create merge session");
set_session_status_for_test(&mut app, &session_id, Status::Review);
let session_folder = app
.sessions
.sessions()
.iter()
.find(|session| session.id == session_id)
.expect("missing created session")
.folder
.clone();
let mock_git =
create_mock_git_client_for_successful_noop_merges(1, dir.path().to_path_buf());
app.sessions.git_client = Arc::new(mock_git);
let result = app.merge_session(&session_id).await;
assert!(result.is_ok(), "merge should enqueue successfully");
wait_for_status_with_retries(&mut app, &session_id, Status::Done, 200).await;
app.sessions.sync_from_handles();
let merged_session = app
.sessions
.sessions()
.iter()
.find(|session| session.id == session_id)
.expect("missing merged session");
assert!(!merged_session.output.contains("[Merge Error]"));
assert!(!session_folder.exists(), "worktree should be removed");
}
#[tokio::test]
async fn test_merge_session_restacks_stacked_child_after_success() {
let dir = tempdir().expect("failed to create temp dir");
let mut app = new_test_app_with_git(dir.path()).await;
let parent_session_id = app
.create_session()
.await
.expect("failed to create merge session");
let child_session_id = app
.create_stacked_draft_session(&parent_session_id)
.await
.expect("failed to create stacked draft session");
app.stage_draft_message(&child_session_id, "Ready after parent merge")
.await
.expect("failed to stage child draft message");
set_session_status_for_test(&mut app, &parent_session_id, Status::Review);
let mock_git =
create_mock_git_client_for_successful_noop_merges(1, dir.path().to_path_buf());
app.sessions.git_client = Arc::new(mock_git);
let result = app.merge_session(&parent_session_id).await;
assert!(result.is_ok(), "merge should enqueue successfully");
wait_for_status_with_retries(&mut app, &parent_session_id, Status::Done, 200).await;
let db_sessions = app
.services
.db()
.sessions()
.load_sessions()
.await
.expect("failed to load sessions");
let db_child_session = db_sessions
.iter()
.find(|session| session.id == child_session_id)
.expect("missing persisted child session");
assert_eq!(db_child_session.parent_session_id, None);
assert_eq!(db_child_session.base_branch, "main");
}
#[tokio::test]
async fn test_merge_session_marks_done_when_changes_are_already_in_base() {
let dir = tempdir().expect("failed to create temp dir");
let mut app = new_test_app_with_git(dir.path()).await;
let session_id = app
.create_session()
.await
.expect("failed to create merge session");
set_session_status_for_test(&mut app, &session_id, Status::Review);
let mock_git =
create_mock_git_client_for_successful_noop_merges(1, dir.path().to_path_buf());
app.sessions.git_client = Arc::new(mock_git);
let result = app.merge_session(&session_id).await;
assert!(result.is_ok(), "merge should enqueue successfully");
wait_for_status_with_retries(&mut app, &session_id, Status::Done, 200).await;
app.sessions.sync_from_handles();
let session = app
.sessions
.sessions()
.iter()
.find(|session| session.id == session_id)
.expect("missing session after merge");
assert!(!session.output.contains("[Merge Error]"));
}
#[tokio::test]
async fn test_merge_session_queue_processes_sessions_in_fifo_order() {
let dir = tempdir().expect("failed to create temp dir");
let mut app = new_test_app_with_git(dir.path()).await;
let first_session_id = app
.create_session()
.await
.expect("failed to create first queue session");
let second_session_id = app
.create_session()
.await
.expect("failed to create second queue session");
set_session_status_for_test(&mut app, &first_session_id, Status::Review);
set_session_status_for_test(&mut app, &second_session_id, Status::Review);
let mock_git =
create_mock_git_client_for_successful_noop_merges(2, dir.path().to_path_buf());
app.sessions.git_client = Arc::new(mock_git);
let first_merge_result = app.merge_session(&first_session_id).await;
let second_merge_result = app.merge_session(&second_session_id).await;
assert!(
first_merge_result.is_ok(),
"first merge request should succeed: {:?}",
first_merge_result.err()
);
assert!(
second_merge_result.is_ok(),
"second merge request should enqueue: {:?}",
second_merge_result.err()
);
wait_for_first_merge_to_complete_before_second_starts(
&mut app,
&first_session_id,
&second_session_id,
)
.await;
wait_for_second_merge_to_start(&mut app, &second_session_id).await;
assert!(
session_status_or_done(&app, &first_session_id) == Status::Done,
"first merge should be complete before second starts"
);
wait_for_all_sessions_done(&mut app, &first_session_id, &second_session_id).await;
app.sessions.sync_from_handles();
let first_status = session_status_or_done(&app, &first_session_id);
let second_status = session_status_or_done(&app, &second_session_id);
assert_eq!(first_status, Status::Done);
assert_eq!(second_status, Status::Done);
}
#[tokio::test]
async fn test_rebase_session_no_git() {
let dir = tempdir().expect("failed to create temp dir");
let mut app = new_test_app(dir.path().to_path_buf()).await;
add_manual_session(&mut app, dir.path(), "manual01", "Test");
let result = app.rebase_session("manual01").await;
assert!(result.is_err());
assert!(
result
.expect_err("should be error")
.to_string()
.contains("No git worktree")
);
}
#[tokio::test]
async fn test_rebase_session_requires_review_status() {
let dir = tempdir().expect("failed to create temp dir");
let mut app = new_test_app_with_git(dir.path()).await;
let session_id = app
.create_session()
.await
.expect("failed to create session");
let result = app.rebase_session(&session_id).await;
assert!(result.is_err());
assert!(
result
.expect_err("should be error")
.to_string()
.contains("must be in review")
);
}
#[tokio::test]
async fn test_rebase_session_invalid_id() {
let dir = tempdir().expect("failed to create temp dir");
let app = new_test_app_with_git(dir.path()).await;
let result = app.rebase_session("missing").await;
assert!(result.is_err());
assert!(
result
.expect_err("should be error")
.to_string()
.contains("Session not found")
);
}
#[tokio::test]
async fn test_rebase_session_updates_session_worktree_to_base_head() {
let dir = tempdir().expect("failed to create temp dir");
let mut app = new_test_app_with_git(dir.path()).await;
let repo_root = dir.path().to_path_buf();
let mut mock_git_client = git::MockGitClient::new();
allow_detect_git_info(&mut mock_git_client);
mock_git_client
.expect_find_git_repo_root()
.times(0..)
.returning(move |_| {
let repo_root = repo_root.clone();
Box::pin(async move { Some(repo_root) })
});
mock_git_client
.expect_create_worktree()
.times(1)
.returning(|_, worktree_path, _, _| {
Box::pin(async move {
let fs_client = create_passthrough_mock_fs_client();
fs_client
.create_dir_all(worktree_path.clone())
.await
.map_err(|error| {
git::GitError::OutputParse(format!(
"Failed to create mock worktree: {error}"
))
})?;
fs_client
.create_dir_all(worktree_path.join(SESSION_DATA_DIR))
.await
.map_err(|error| {
git::GitError::OutputParse(format!(
"Failed to create mock worktree data dir: {error}"
))
})?;
Ok(())
})
});
mock_git_client
.expect_is_worktree_clean()
.times(1)
.returning(|_| Box::pin(async { Ok(true) }));
mock_git_client
.expect_is_rebase_in_progress()
.times(1)
.returning(|_| Box::pin(async { Ok(false) }));
mock_git_client
.expect_rebase_start()
.times(1)
.returning(|_, _| Box::pin(async { Ok(git::RebaseStepResult::Completed) }));
install_mock_git_client(&mut app, mock_git_client);
let session_id = app
.create_session()
.await
.expect("failed to create session");
app.sessions.sessions_mut()[0].status = Status::Review;
if let Some(handles) = app.sessions.session_handles().get(session_id.as_str())
&& let Ok(mut session_status) = handles.status.lock()
{
*session_status = Status::Review;
}
let result = app.rebase_session(&session_id).await;
assert!(result.is_ok(), "rebase should succeed: {:?}", result.err());
wait_for_output_contains(&mut app, &session_id, "[Rebase] Successfully rebased", 200).await;
}
#[tokio::test]
async fn test_rebase_session_auto_commits_uncommitted_changes() {
let dir = tempdir().expect("failed to create temp dir");
let mut app = new_test_app_with_git(dir.path()).await;
let repo_root = dir.path().to_path_buf();
let mut mock_git_client = git::MockGitClient::new();
allow_detect_git_info(&mut mock_git_client);
mock_git_client
.expect_find_git_repo_root()
.times(1)
.returning(move |_| {
let repo_root = repo_root.clone();
Box::pin(async move { Some(repo_root) })
});
mock_git_client
.expect_create_worktree()
.times(1)
.returning(|_, worktree_path, _, _| {
Box::pin(async move {
let fs_client = create_passthrough_mock_fs_client();
fs_client
.create_dir_all(worktree_path.clone())
.await
.map_err(|error| {
git::GitError::OutputParse(format!(
"Failed to create mock worktree: {error}"
))
})?;
fs_client
.create_dir_all(worktree_path.join(SESSION_DATA_DIR))
.await
.map_err(|error| {
git::GitError::OutputParse(format!(
"Failed to create mock worktree data dir: {error}"
))
})?;
Ok(())
})
});
mock_git_client
.expect_is_worktree_clean()
.times(1)
.returning(|_| Box::pin(async { Ok(false) }));
mock_git_client
.expect_has_commits_since()
.times(1)
.returning(|_, _| Box::pin(async { Ok(true) }));
mock_git_client
.expect_head_commit_message()
.times(1)
.returning(|_| Box::pin(async { Ok(Some("Existing session commit".to_string())) }));
mock_git_client
.expect_commit_all_preserving_single_commit()
.times(1)
.withf(|_, _, _, _, no_verify| *no_verify)
.returning(|_, _, _, _, _| Box::pin(async { Ok(()) }));
mock_git_client
.expect_head_short_hash()
.times(1)
.returning(|_| Box::pin(async { Ok("cafe123".to_string()) }));
mock_git_client
.expect_is_rebase_in_progress()
.times(1)
.returning(|_| Box::pin(async { Ok(false) }));
mock_git_client
.expect_rebase_start()
.times(1)
.returning(|_, _| Box::pin(async { Ok(git::RebaseStepResult::Completed) }));
install_mock_git_client(&mut app, mock_git_client);
let session_id = app
.create_session()
.await
.expect("failed to create session");
let session_folder = app.sessions.sessions()[0].folder.clone();
app.sessions.sessions_mut()[0].status = Status::Review;
if let Some(handles) = app.sessions.session_handles().get(session_id.as_str())
&& let Ok(mut session_status) = handles.status.lock()
{
*session_status = Status::Review;
}
std::fs::write(session_folder.join("dirty.txt"), "uncommitted content")
.expect("failed to write dirty file");
let result = app.rebase_session(&session_id).await;
assert!(result.is_ok(), "rebase should succeed: {:?}", result.err());
wait_for_output_contains(&mut app, &session_id, "[Rebase] Successfully rebased", 200).await;
app.refresh_sessions_now().await;
}
#[tokio::test]
async fn test_sync_main_uses_active_project_branch_from_context() {
let dir = tempdir().expect("failed to create temp dir");
let mut app = new_test_app_with_git(dir.path()).await;
app.projects.update_active_project_context(
app.active_project_id(),
app.projects.project_name().to_string(),
Some("develop".to_string()),
None,
dir.path().to_path_buf(),
);
let repo_root = dir.path().to_path_buf();
let mut mock_git_client = git::MockGitClient::new();
mock_git_client
.expect_find_git_repo_root()
.times(1)
.returning(move |_| {
let repo_root = repo_root.clone();
Box::pin(async move { Some(repo_root) })
});
mock_git_client
.expect_is_worktree_clean()
.times(1)
.returning(|_| Box::pin(async { Ok(false) }));
let result = SessionManager::sync_main_for_project(
app.projects.git_branch().map(str::to_string),
app.projects.working_dir().to_path_buf(),
Arc::new(mock_git_client),
AgentModel::Gemini3FlashPreview,
)
.await;
assert_eq!(
result,
Err(SyncSessionStartError::MainHasUncommittedChanges {
default_branch: "develop".to_string(),
})
);
}
#[tokio::test]
async fn test_sync_main_requires_clean_selected_project_branch() {
let dir = tempdir().expect("failed to create temp dir");
let app = new_test_app_with_git(dir.path()).await;
let repo_root = dir.path().to_path_buf();
let mut mock_git_client = git::MockGitClient::new();
mock_git_client
.expect_find_git_repo_root()
.times(1)
.returning(move |_| {
let repo_root = repo_root.clone();
Box::pin(async move { Some(repo_root) })
});
mock_git_client
.expect_is_worktree_clean()
.times(1)
.returning(|_| Box::pin(async { Ok(false) }));
let result = SessionManager::sync_main_for_project(
app.projects.git_branch().map(str::to_string),
app.projects.working_dir().to_path_buf(),
Arc::new(mock_git_client),
AgentModel::Gemini3FlashPreview,
)
.await;
assert_eq!(
result,
Err(SyncSessionStartError::MainHasUncommittedChanges {
default_branch: "main".to_string(),
})
);
}
#[tokio::test]
async fn test_sync_main_returns_error_without_upstream_remote() {
let dir = tempdir().expect("failed to create temp dir");
let app = new_test_app_with_git(dir.path()).await;
let result = SessionManager::sync_main_for_project(
app.projects.git_branch().map(str::to_string),
app.projects.working_dir().to_path_buf(),
app.services.git_client(),
AgentModel::Gemini3FlashPreview,
)
.await;
assert!(matches!(result, Err(SyncSessionStartError::Other(_))));
}
#[tokio::test]
async fn test_sync_main_pushes_local_commits_to_remote() {
let dir = tempdir().expect("failed to create temp dir");
let repo_root = dir.path().to_path_buf();
let mut mock_git_client = git::MockGitClient::new();
mock_git_client
.expect_find_git_repo_root()
.times(1)
.returning(move |_| {
let repo_root = repo_root.clone();
Box::pin(async move { Some(repo_root) })
});
mock_git_client
.expect_is_worktree_clean()
.times(1)
.returning(|_| Box::pin(async { Ok(true) }));
let mut ahead_behind_calls = 0_u8;
mock_git_client
.expect_get_ahead_behind()
.times(2)
.returning(move |_| {
ahead_behind_calls = ahead_behind_calls.saturating_add(1);
let value = if ahead_behind_calls == 1 {
(1, 2)
} else {
(0, 0)
};
Box::pin(async move { Ok(value) })
});
mock_git_client
.expect_list_upstream_commit_titles()
.times(1)
.returning(|_| Box::pin(async { Ok(vec!["remote fix".to_string()]) }));
mock_git_client
.expect_pull_rebase()
.times(1)
.returning(|_| Box::pin(async { Ok(git::PullRebaseResult::Completed) }));
mock_git_client
.expect_list_local_commit_titles()
.times(1)
.returning(|_| Box::pin(async { Ok(vec!["local work".to_string()]) }));
mock_git_client
.expect_push_current_branch()
.times(1)
.returning(|_| Box::pin(async { Ok("origin/main".to_string()) }));
let result = SessionManager::sync_main_for_project(
Some("main".to_string()),
dir.path().to_path_buf(),
Arc::new(mock_git_client),
AgentModel::Gemini3FlashPreview,
)
.await;
let outcome = result.expect("sync should succeed");
assert_eq!(outcome.pulled_commits, Some(2));
assert_eq!(outcome.pushed_commits, Some(0));
assert_eq!(outcome.pulled_commit_titles, vec!["remote fix".to_string()]);
assert_eq!(outcome.pushed_commit_titles, vec!["local work".to_string()]);
assert!(outcome.resolved_conflict_files.is_empty());
}
#[tokio::test]
async fn test_cancel_session() {
let dir = tempdir().expect("failed to create temp dir");
let mut app = new_test_app_with_git(dir.path()).await;
let session_id = app
.create_session()
.await
.expect("failed to create session");
let session_folder = app
.sessions
.sessions()
.iter()
.find(|session| session.id == session_id)
.expect("missing session")
.folder
.clone();
set_session_status_for_test(&mut app, &session_id, Status::Review);
app.sessions
.cancel_session(&app.services, &session_id)
.await
.expect("failed to cancel session");
app.sessions.sync_from_handles();
let session = app
.sessions
.sessions()
.iter()
.find(|session| session.id == session_id)
.expect("missing session");
assert_eq!(session.status, Status::Canceled);
let db_sessions = app
.services
.db()
.sessions()
.load_sessions()
.await
.expect("failed to load");
let db_session = db_sessions
.iter()
.find(|session| session.id == session_id)
.expect("missing persisted session");
assert_eq!(db_session.status, "Canceled");
wait_for_path_absent(&session_folder).await;
}
#[tokio::test]
async fn test_cancel_draft_session() {
let dir = tempdir().expect("failed to create temp dir");
let mut app = new_test_app_with_git(dir.path()).await;
let session_id = app
.create_draft_session()
.await
.expect("failed to create draft session");
let session_folder = app
.sessions
.sessions()
.iter()
.find(|session| session.id == session_id)
.expect("missing session")
.folder
.clone();
app.sessions
.cancel_session(&app.services, &session_id)
.await
.expect("failed to cancel draft session");
app.sessions.sync_from_handles();
let session = app
.sessions
.sessions()
.iter()
.find(|session| session.id == session_id)
.expect("missing session");
assert_eq!(session.status, Status::Canceled);
let db_sessions = app
.services
.db()
.sessions()
.load_sessions()
.await
.expect("failed to load");
let db_session = db_sessions
.iter()
.find(|session| session.id == session_id)
.expect("missing persisted session");
assert_eq!(db_session.status, "Canceled");
wait_for_path_absent(&session_folder).await;
}
#[tokio::test]
async fn test_cancel_session_cascades_to_stacked_child() {
let dir = tempdir().expect("failed to create temp dir");
let mut app = new_test_app_with_git(dir.path()).await;
let parent_session_id = app.create_session().await.expect("failed to create parent");
let child_session_id = app
.create_stacked_draft_session(&parent_session_id)
.await
.expect("failed to create stacked draft session");
set_session_status_for_test(&mut app, &parent_session_id, Status::Review);
app.sessions
.cancel_session(&app.services, &parent_session_id)
.await
.expect("failed to cancel parent session");
app.sessions.sync_from_handles();
let parent_session = app
.sessions
.sessions()
.iter()
.find(|session| session.id == parent_session_id)
.expect("missing parent session");
let child_session = app
.sessions
.sessions()
.iter()
.find(|session| session.id == child_session_id)
.expect("missing child session");
assert_eq!(parent_session.status, Status::Canceled);
assert_eq!(child_session.status, Status::Canceled);
let db_sessions = app
.services
.db()
.sessions()
.load_sessions()
.await
.expect("failed to load");
let db_parent_session = db_sessions
.iter()
.find(|session| session.id == parent_session_id)
.expect("missing persisted parent session");
let db_child_session = db_sessions
.iter()
.find(|session| session.id == child_session_id)
.expect("missing persisted child session");
assert_eq!(db_parent_session.status, "Canceled");
assert_eq!(db_child_session.status, "Canceled");
}
#[tokio::test]
async fn test_cancel_running_session_stops_turn_and_cancels_session() {
let dir = tempdir().expect("failed to create temp dir");
let mut app = new_test_app_with_git(dir.path()).await;
let session_id = app
.create_session()
.await
.expect("failed to create session");
set_session_status_for_test(&mut app, &session_id, Status::InProgress);
app.services
.db()
.operations()
.insert_session_operation("operation-id", &session_id, "reply")
.await
.expect("failed to insert operation");
app.sessions
.enqueue_message(&app.services, &session_id, "queued reply")
.expect("failed to enqueue message");
app.sessions
.cancel_session(&app.services, &session_id)
.await
.expect("failed to cancel session");
app.sessions.sync_from_handles();
let session = app
.sessions
.sessions()
.iter()
.find(|session| session.id == session_id)
.expect("missing session");
assert_eq!(session.status, Status::Canceled);
let handles = app
.sessions
.session_handles_or_err(&session_id)
.expect("missing session handles");
assert!(
handles
.cancel_token
.lock()
.expect("cancel token lock")
.is_cancelled()
);
assert!(
handles
.queued_messages
.lock()
.expect("queue lock")
.is_empty()
);
assert!(
app.services
.db()
.operations()
.is_cancel_requested_for_operation("operation-id")
.await
.expect("failed to load operation cancel status")
);
let db_sessions = app
.services
.db()
.sessions()
.load_sessions()
.await
.expect("failed to load");
let db_session = db_sessions
.iter()
.find(|session| session.id == session_id)
.expect("missing persisted session");
assert_eq!(db_session.status, "Canceled");
}
#[tokio::test]
async fn test_cancel_session_triggers_app_server_shutdown() {
let dir = tempdir().expect("failed to create temp dir");
let db = AppRepositories::in_memory().await;
let (shutdown_tx, mut shutdown_rx) = tokio::sync::mpsc::unbounded_channel::<String>();
let mut mock_app_server = MockAppServerClient::new();
mock_app_server
.expect_run_turn()
.times(1)
.returning(|_, _| {
Box::pin(async {
Ok(AppServerTurnResponse {
assistant_message: r#"{"answer":"ready","questions":[],"summary":null}"#
.to_string(),
context_reset: false,
input_tokens: 0,
output_tokens: 0,
pid: None,
provider_conversation_id: None,
})
})
});
mock_app_server
.expect_shutdown_session()
.times(1)
.returning(move |session_id| {
let shutdown_tx = shutdown_tx.clone();
Box::pin(async move {
let _ = shutdown_tx.send(session_id);
})
});
let app_server_client: Arc<dyn app_server::AppServerClient> = Arc::new(mock_app_server);
let mut app = new_test_app_with_db_and_app_server(
dir.path().to_path_buf(),
dir.path().to_path_buf(),
Some("main".to_string()),
db,
app_server_client,
)
.await;
let session_id = app
.create_session()
.await
.expect("failed to create session");
app.sessions
.reply(&app.services, &session_id, "Start")
.await;
wait_for_status(&mut app, &session_id, Status::Review).await;
app.cancel_session(&session_id)
.await
.expect("failed to cancel session");
app.process_pending_app_events().await;
wait_for_status(&mut app, &session_id, Status::Canceled).await;
let shutdown_session_id =
tokio::time::timeout(std::time::Duration::from_secs(1), shutdown_rx.recv())
.await
.expect("timed out waiting for app-server shutdown")
.expect("missing shutdown session id");
assert_eq!(shutdown_session_id, session_id);
}
#[tokio::test]
async fn test_done_status_triggers_app_server_shutdown() {
let dir = tempdir().expect("failed to create temp dir");
let db = AppRepositories::in_memory().await;
let (shutdown_tx, mut shutdown_rx) = tokio::sync::mpsc::unbounded_channel::<String>();
let mut mock_app_server = MockAppServerClient::new();
mock_app_server
.expect_run_turn()
.times(1)
.returning(|_, _| {
Box::pin(async {
Ok(AppServerTurnResponse {
assistant_message: r#"{"answer":"ready","questions":[],"summary":null}"#
.to_string(),
context_reset: false,
input_tokens: 0,
output_tokens: 0,
pid: None,
provider_conversation_id: None,
})
})
});
mock_app_server
.expect_shutdown_session()
.times(1)
.returning(move |session_id| {
let shutdown_tx = shutdown_tx.clone();
Box::pin(async move {
let _ = shutdown_tx.send(session_id);
})
});
let app_server_client: Arc<dyn app_server::AppServerClient> = Arc::new(mock_app_server);
let mut app = new_test_app_with_db_and_app_server(
dir.path().to_path_buf(),
dir.path().to_path_buf(),
Some("main".to_string()),
db,
app_server_client,
)
.await;
let session_id = app
.create_session()
.await
.expect("failed to create session");
app.sessions
.reply(&app.services, &session_id, "Start")
.await;
wait_for_status(&mut app, &session_id, Status::Review).await;
let handles = app
.sessions
.session_handles_or_err(&session_id)
.expect("missing session handles");
let session_status = Arc::clone(&handles.status);
let app_event_tx = app.services.event_sender();
let transitioned_to_merging = SessionTaskService::update_status(
&session_status,
app.services.clock().as_ref(),
app.services.db(),
&app_event_tx,
&app.services.session_update_versions(),
&session_id,
Status::Merging,
)
.await;
assert!(
transitioned_to_merging,
"status transition to Merging should succeed"
);
let transitioned_to_done = SessionTaskService::update_status(
&session_status,
app.services.clock().as_ref(),
app.services.db(),
&app_event_tx,
&app.services.session_update_versions(),
&session_id,
Status::Done,
)
.await;
assert!(
transitioned_to_done,
"status transition to Done should succeed"
);
app.process_pending_app_events().await;
wait_for_status(&mut app, &session_id, Status::Done).await;
let shutdown_session_id =
tokio::time::timeout(std::time::Duration::from_secs(1), shutdown_rx.recv())
.await
.expect("timed out waiting for app-server shutdown")
.expect("missing shutdown session id");
assert_eq!(shutdown_session_id, session_id);
}
#[tokio::test]
async fn test_cancel_session_requires_cancelable_status() {
let dir = tempdir().expect("failed to create temp dir");
let mut app = new_test_app_with_git(dir.path()).await;
let session_id = app
.create_session()
.await
.expect("failed to create session");
let result = app
.sessions
.cancel_session(&app.services, &session_id)
.await;
assert!(result.is_err());
assert!(
result
.expect_err("should be error")
.to_string()
.contains("must be running, in review")
);
}
#[tokio::test]
async fn test_cancel_session_invalid_id() {
let dir = tempdir().expect("failed to create temp dir");
let app = new_test_app(dir.path().to_path_buf()).await;
let result = app.sessions.cancel_session(&app.services, "missing").await;
assert!(result.is_err());
assert!(
result
.expect_err("should be error")
.to_string()
.contains("Session not found")
);
}
#[tokio::test]
async fn test_cleanup_merged_session_worktree_without_repo_hint() {
let dir = tempdir().expect("failed to create temp dir");
let worktree_folder = dir.path().join("merged-worktree");
let branch_name = "wt/cleanup123";
std::fs::create_dir_all(&worktree_folder).expect("failed to create worktree folder");
assert!(
worktree_folder.exists(),
"worktree should exist before cleanup"
);
let repo_root = dir.path().to_path_buf();
let mut mock_git_client = git::MockGitClient::new();
mock_git_client
.expect_main_repo_root()
.times(1)
.returning(move |_| {
let repo_root = repo_root.clone();
Box::pin(async move { Ok(repo_root) })
});
mock_git_client
.expect_remove_worktree()
.times(1)
.returning(|worktree_path| {
Box::pin(async move {
let fs_client = create_passthrough_mock_fs_client();
let _ = fs_client.remove_dir_all(worktree_path).await;
Ok(())
})
});
mock_git_client
.expect_delete_branch()
.times(1)
.withf(|_, branch| branch == "wt/cleanup123")
.returning(|_, _| Box::pin(async { Ok(()) }));
let result = SessionManager::cleanup_merged_session_worktree(
worktree_folder.clone(),
Arc::new(create_passthrough_mock_fs_client()),
Arc::new(mock_git_client),
branch_name.to_string(),
None,
)
.await;
assert!(result.is_ok(), "cleanup should succeed: {:?}", result.err());
assert!(
!worktree_folder.exists(),
"worktree should be removed after cleanup"
);
}
#[tokio::test]
async fn test_active_project_id_getter() {
let dir = tempdir().expect("failed to create temp dir");
let app = new_test_app(dir.path().to_path_buf()).await;
assert!(app.active_project_id() > 0);
}
#[tokio::test]
async fn test_create_session_scoped_to_project() {
let dir = tempdir().expect("failed to create temp dir");
let mut app = new_test_app_with_git(dir.path()).await;
let project_id = app.active_project_id();
app.create_session()
.await
.expect("failed to create session");
let sessions = app
.services
.db()
.sessions()
.load_sessions()
.await
.expect("failed to load");
assert_eq!(sessions.len(), 1);
assert_eq!(sessions[0].project_id, Some(project_id));
}
#[test]
fn test_parse_merge_commit_message_response_with_protocol_message() {
let content = r#"{"answer":"Title\n\n- Detail","questions":[],"summary":null}"#;
let parsed = crate::infra::agent::protocol::parse_agent_response_strict(content)
.ok()
.map(|response| response.to_answer_display_text())
.filter(|answer_text| !answer_text.trim().is_empty());
assert!(parsed.is_some());
assert_eq!(parsed.as_deref(), Some("Title\n\n- Detail"));
}
#[test]
fn test_parse_merge_commit_message_response_rejects_non_protocol_json() {
let content = r#"{"title":"Title","description":"- Detail"}"#;
let parsed = crate::infra::agent::protocol::parse_agent_response_strict(content)
.ok()
.map(|response| response.to_answer_display_text())
.filter(|answer_text| !answer_text.trim().is_empty());
assert!(parsed.is_none());
}
#[test]
fn test_session_folder_uses_first_8_chars() {
let base = Path::new("/home/user/.agentty/wt");
let session_id = "a1b2c3d4-e5f6-7890-abcd-ef1234567890";
let folder = session_folder(base, session_id);
assert_eq!(folder, PathBuf::from("/home/user/.agentty/wt/a1b2c3d4"));
}
#[test]
fn test_session_branch_uses_first_8_chars() {
let session_id = "a1b2c3d4-e5f6-7890-abcd-ef1234567890";
let branch = session_branch(session_id);
assert_eq!(branch, "wt/a1b2c3d4");
}
#[tokio::test]
async fn test_refresh_session_branch_names_runs_detection_concurrently() {
let dir = tempdir().expect("failed to create temp dir");
let mut app = new_test_app(dir.path().to_path_buf()).await;
add_manual_session(&mut app, dir.path(), "alpha001", "First");
add_manual_session(&mut app, dir.path(), "bravo002", "Second");
let barrier = Arc::new(Barrier::new(2));
let mut mock_git_client = git::MockGitClient::new();
mock_git_client
.expect_detect_git_info()
.times(2)
.returning({
let barrier = Arc::clone(&barrier);
move |_| {
let barrier = Arc::clone(&barrier);
Box::pin(async move {
barrier.wait().await;
None
})
}
});
install_mock_git_client(&mut app, mock_git_client);
let refresh_result = tokio::time::timeout(
Duration::from_millis(100),
app.sessions.refresh_session_branch_names(),
)
.await;
assert!(
refresh_result.is_ok(),
"branch refresh should complete without serially blocking"
);
assert_eq!(
app.sessions.session_branch_name("alpha001"),
Some("wt/alpha001")
);
assert_eq!(
app.sessions.session_branch_name("bravo002"),
Some("wt/bravo002")
);
}
#[test]
fn test_remote_branch_name_strips_remote_prefix() {
let branch = remote_branch_name_from_upstream_ref("origin/wt/abc12345");
assert_eq!(branch, "wt/abc12345");
}
#[test]
fn test_remote_branch_name_returns_input_for_bare_ref() {
let branch = remote_branch_name_from_upstream_ref("no-slash");
assert_eq!(branch, "no-slash");
}
}