use std::path::{Path, PathBuf};
use std::sync::{Arc, Mutex};
use ag_forge as forge;
use askama::Template;
use tokio::sync::mpsc;
use tracing::warn;
use uuid::Uuid;
use super::worker::{SessionCommand, TurnMetadata};
use super::{
SessionTaskService, draft, session_branch, session_folder, unix_timestamp_from_system_time,
};
use crate::app::session::SessionError;
use crate::app::{
AppEvent, AppServices, ProjectManager, SessionManager, agentty_home, review_request, setting,
};
use crate::domain::agent::{AgentModel, ReasoningLevel};
use crate::domain::session::{ReviewRequest, SESSION_DATA_DIR, Session, SessionId, Status};
use crate::domain::setting::SettingName;
use crate::infra::channel::{
AgentRequestKind, TurnPrompt, TurnPromptAttachment, TurnPromptTextSource,
};
use crate::infra::fs::FsClient;
use crate::infra::{agent, db, git};
use crate::ui::page::session_list::grouped_session_indexes;
const GENERATED_SESSION_TITLE_MAX_CHARACTERS: usize = 72;
const USER_PROMPT_PREFIX: &str = " › ";
const USER_PROMPT_CONTINUATION_PREFIX: &str = " ";
struct BuildSessionCommandInput {
is_first_message: bool,
published_upstream_ref: Option<String>,
prompt: TurnPrompt,
session_model: AgentModel,
session_output: Option<String>,
}
type ReplyContext = (Option<String>, bool, SessionId, Option<String>);
struct DeletedSessionCleanup {
branch_name: String,
folder: PathBuf,
has_git_branch: bool,
session_id: SessionId,
staged_draft_root: PathBuf,
working_dir: PathBuf,
}
#[derive(Template)]
#[template(path = "session_title_generation_prompt.md", escape = "none")]
struct SessionTitleGenerationPromptTemplate<'a> {
prompt: &'a str,
}
struct TitleGenerationTaskCompletion {
generation: u64,
session_id: SessionId,
}
impl SessionManager {
pub fn next(&mut self) {
let grouped_indexes = grouped_session_indexes(&self.state.sessions);
if grouped_indexes.is_empty() {
return;
}
let index = match self
.state
.table_state
.selected()
.and_then(|selected_index| {
grouped_indexes
.iter()
.position(|session_index| *session_index == selected_index)
}) {
Some(position) => {
if position >= grouped_indexes.len() - 1 {
0
} else {
position + 1
}
}
None => 0,
};
self.state.table_state.select(Some(grouped_indexes[index]));
}
pub fn previous(&mut self) {
let grouped_indexes = grouped_session_indexes(&self.state.sessions);
if grouped_indexes.is_empty() {
return;
}
let index = match self
.state
.table_state
.selected()
.and_then(|selected_index| {
grouped_indexes
.iter()
.position(|session_index| *session_index == selected_index)
}) {
Some(position) => {
if position == 0 {
grouped_indexes.len() - 1
} else {
position - 1
}
}
None => 0,
};
self.state.table_state.select(Some(grouped_indexes[index]));
}
pub async fn create_session(
&mut self,
projects: &ProjectManager,
services: &AppServices,
) -> Result<String, SessionError> {
self.create_live_session(projects, services).await
}
pub async fn create_draft_session(
&mut self,
projects: &ProjectManager,
services: &AppServices,
) -> Result<String, SessionError> {
let base_branch = projects.git_branch().ok_or_else(|| {
SessionError::Workflow("Git branch is required to create a session".to_string())
})?;
self.create_draft_session_for_project(services, projects.active_project_id(), base_branch)
.await
}
pub async fn create_draft_session_for_project(
&mut self,
services: &AppServices,
project_id: i64,
base_branch: &str,
) -> Result<String, SessionError> {
let session_model = self
.resolve_default_session_model(services, project_id)
.await;
self.default_session_model = session_model;
let session_id = Uuid::new_v4().to_string();
let folder = session_folder(services.base_path(), &session_id);
if services.fs_client().exists(folder.clone()) {
return Err(SessionError::Workflow(format!(
"Session folder {session_id} already exists"
)));
}
services
.db()
.insert_draft_session(
&session_id,
session_model.as_str(),
base_branch,
&Status::New.to_string(),
project_id,
)
.await
.map_err(|error| {
SessionError::Workflow(format!("Failed to save session metadata: {error}"))
})?;
Self::record_session_creation_activity(services, &session_id).await;
services.emit_app_event(AppEvent::RefreshSessions);
Ok(session_id)
}
async fn create_live_session(
&mut self,
projects: &ProjectManager,
services: &AppServices,
) -> Result<String, SessionError> {
let base_branch = projects.git_branch().ok_or_else(|| {
SessionError::Workflow("Git branch is required to create a session".to_string())
})?;
let session_model = self
.resolve_default_session_model(services, projects.active_project_id())
.await;
self.default_session_model = session_model;
let session_id = Uuid::new_v4().to_string();
let folder = session_folder(services.base_path(), &session_id);
let fs_client = services.fs_client();
if fs_client.exists(folder.clone()) {
return Err(SessionError::Workflow(format!(
"Session folder {session_id} already exists"
)));
}
let worktree_branch = session_branch(&session_id);
let working_dir = projects.working_dir().to_path_buf();
let git_client = services.git_client();
let repo_root = git_client
.find_git_repo_root(working_dir)
.await
.ok_or_else(|| {
SessionError::Workflow("Failed to find git repository root".to_string())
})?;
self.create_session_worktree(
services,
&session_id,
&folder,
&repo_root,
&worktree_branch,
base_branch,
)
.await?;
if let Err(error) = services
.db()
.insert_session(
&session_id,
session_model.as_str(),
base_branch,
&Status::New.to_string(),
projects.active_project_id(),
)
.await
{
self.rollback_failed_session_creation(
services,
&folder,
&repo_root,
&session_id,
&worktree_branch,
false,
)
.await;
return Err(SessionError::Workflow(format!(
"Failed to save session metadata: {error}"
)));
}
Self::record_session_creation_activity(services, &session_id).await;
if let Err(error) = agent::create_backend(session_model.kind()).setup(&folder) {
self.rollback_failed_session_creation(
services,
&folder,
&repo_root,
&session_id,
&worktree_branch,
true,
)
.await;
return Err(SessionError::Workflow(format!(
"Failed to setup session backend: {error}"
)));
}
services.emit_app_event(AppEvent::RefreshSessions);
Ok(session_id)
}
async fn create_session_worktree(
&self,
services: &AppServices,
session_id: &str,
folder: &Path,
repo_root: &Path,
worktree_branch: &str,
base_branch: &str,
) -> Result<(), SessionError> {
let git_client = services.git_client();
git_client
.create_worktree(
repo_root.to_path_buf(),
folder.to_path_buf(),
worktree_branch.to_string(),
base_branch.to_string(),
)
.await
.map_err(|error| {
SessionError::Workflow(format!("Failed to create git worktree: {error}"))
})?;
let data_dir = folder.join(SESSION_DATA_DIR);
if let Err(error) = services.fs_client().create_dir_all(data_dir).await {
self.rollback_failed_session_creation(
services,
folder,
repo_root,
session_id,
worktree_branch,
false,
)
.await;
return Err(SessionError::Workflow(format!(
"Failed to create session metadata directory: {error}"
)));
}
Ok(())
}
async fn ensure_session_worktree_ready(
&mut self,
services: &AppServices,
session_id: &str,
) -> Result<(), SessionError> {
let (base_branch, folder, persisted_session_id, session_model) = {
let session = self.session_or_err(session_id)?;
if !session.is_draft_session() {
return Ok(());
}
(
session.base_branch.clone(),
session.folder.clone(),
session.id.clone(),
session.model,
)
};
if services.fs_client().is_dir(folder.clone()) {
agent::create_backend(session_model.kind())
.setup(&folder)
.map_err(|error| {
SessionError::Workflow(format!("Failed to setup session backend: {error}"))
})?;
self.set_session_worktree_available(session_id, true);
return Ok(());
}
let repo_root = self.load_session_repo_root(services, session_id).await?;
let worktree_branch = session_branch(&persisted_session_id);
self.create_session_worktree(
services,
&persisted_session_id,
&folder,
&repo_root,
&worktree_branch,
&base_branch,
)
.await?;
if let Err(error) = agent::create_backend(session_model.kind()).setup(&folder) {
let cleanup_errors = Self::cleanup_session_worktree_resources(
services.fs_client().clone(),
services.git_client(),
folder,
worktree_branch,
Some(repo_root),
true,
)
.await;
if !cleanup_errors.is_empty() {
return Err(SessionError::Workflow(format!(
"Failed to setup session backend: {error}. Cleanup also failed: {}",
cleanup_errors.join("; ")
)));
}
return Err(SessionError::Workflow(format!(
"Failed to setup session backend: {error}"
)));
}
self.set_session_worktree_available(session_id, true);
Ok(())
}
async fn load_session_repo_root(
&self,
services: &AppServices,
session_id: &str,
) -> Result<PathBuf, SessionError> {
let project_id = services
.db()
.load_session_project_id(session_id)
.await?
.ok_or_else(|| {
SessionError::Workflow(
"Session project is required to create a worktree".to_string(),
)
})?;
let project_path = self.load_project_path(services, project_id).await?;
services
.git_client()
.find_git_repo_root(project_path)
.await
.ok_or_else(|| SessionError::Workflow("Failed to find git repository root".to_string()))
}
async fn load_project_path(
&self,
services: &AppServices,
project_id: i64,
) -> Result<PathBuf, SessionError> {
let project_row = services
.db()
.get_project(project_id)
.await?
.ok_or_else(|| {
SessionError::Workflow(format!("Project with id `{project_id}` was not found"))
})?;
Ok(PathBuf::from(project_row.path))
}
async fn persist_staged_draft(
services: &AppServices,
session_id: &str,
staged_attachments: &[TurnPromptAttachment],
staged_prompt: &str,
title_to_save: Option<&str>,
) -> Result<(), SessionError> {
draft::store_staged_draft_attachments(
services.fs_client().as_ref(),
services.base_path(),
session_id,
staged_attachments,
)
.await?;
services
.db()
.update_session_prompt(session_id, staged_prompt)
.await?;
if let Some(title) = title_to_save {
services.db().update_session_title(session_id, title).await?;
}
Ok(())
}
pub async fn stage_draft_message(
&mut self,
services: &AppServices,
session_id: &str,
prompt: impl Into<TurnPrompt>,
) -> Result<(), SessionError> {
let prompt = prompt.into();
let session_index = self.session_index_or_err(session_id)?;
let (
folder,
persisted_session_id,
session_model,
staged_attachments,
staged_prompt,
title_to_save,
) = {
let session = self
.sessions
.get(session_index)
.ok_or(SessionError::NotFound)?;
if !session.is_draft_session() {
return Err(SessionError::Workflow(
"Only draft sessions can stage drafts".to_string(),
));
}
if session.status != Status::New {
return Err(SessionError::Workflow(
"Only `New` sessions can stage drafts".to_string(),
));
}
let next_attachment_number = session.draft_attachments.len().saturating_add(1);
let staged_prompt =
Self::append_staged_prompt(&session.prompt, &prompt, next_attachment_number);
let mut staged_attachments = session.draft_attachments.clone();
staged_attachments.extend(Self::renumbered_attachments(
&prompt,
next_attachment_number,
));
let title_to_save = session.title.is_none().then(|| prompt.transcript_text());
(
session.folder.clone(),
session.id.clone(),
session.model,
staged_attachments,
staged_prompt,
title_to_save,
)
};
let project_id = services
.db()
.load_session_project_id(&persisted_session_id)
.await?
.ok_or_else(|| {
SessionError::Workflow(
"Session project is required to stage draft prompts".to_string(),
)
})?;
let project_working_dir = self.load_project_path(services, project_id).await?;
let title_generation_model =
setting::load_default_fast_model_setting(services, Some(project_id), session_model)
.await;
let title_generation_folder = if services.fs_client().is_dir(folder.clone()) {
folder.clone()
} else {
project_working_dir
};
Self::persist_staged_draft(
services,
&persisted_session_id,
&staged_attachments,
&staged_prompt,
title_to_save.as_deref(),
)
.await?;
let title_generation_prompt = staged_prompt.clone();
if let Some(session) = self.sessions.get_mut(session_index) {
session.prompt = staged_prompt;
session.draft_attachments = staged_attachments;
if let Some(title_to_save) = title_to_save {
session.title = Some(title_to_save);
}
}
let title_generation_task_generation =
self.next_title_generation_task_generation(&persisted_session_id);
let title_generation_task = Self::spawn_session_title_generation_task(
services.event_sender(),
services.db().clone(),
&persisted_session_id,
&title_generation_folder,
&title_generation_prompt,
title_generation_model,
Some(title_generation_task_generation),
);
self.replace_title_generation_task(
&persisted_session_id,
title_generation_task_generation,
title_generation_task,
);
SessionTaskService::emit_session_updated(
&services.event_sender(),
&services.session_update_versions(),
persisted_session_id.as_str(),
);
Ok(())
}
pub async fn start_staged_session(
&mut self,
services: &AppServices,
session_id: &str,
) -> Result<(), SessionError> {
let prompt = {
let session = self.session_or_err(session_id)?;
if !session.is_draft_session() {
return Err(SessionError::Workflow(
"Only draft sessions can be started from staged drafts".to_string(),
));
}
if session.status != Status::New {
return Err(SessionError::Workflow(
"Only `New` sessions can be started from staged drafts".to_string(),
));
}
if session.prompt.is_empty() {
return Err(SessionError::Workflow(
"Stage at least one draft before starting the session".to_string(),
));
}
TurnPrompt {
attachments: session.draft_attachments.clone(),
text: session.prompt.clone(),
text_source: TurnPromptTextSource::UserPrompt,
}
};
self.start_session(services, session_id, prompt).await?;
if let Ok(session_index) = self.session_index_or_err(session_id)
&& let Some(session) = self.sessions.get_mut(session_index)
{
session.draft_attachments.clear();
}
if let Err(error) = draft::store_staged_draft_attachments(
services.fs_client().as_ref(),
services.base_path(),
session_id,
&[],
)
.await
{
warn!(
session_id = session_id,
error = %error,
"failed to clear staged draft attachments after session start"
);
}
Ok(())
}
pub async fn start_session(
&mut self,
services: &AppServices,
session_id: &str,
prompt: impl Into<TurnPrompt>,
) -> Result<(), SessionError> {
let prompt = prompt.into();
self.ensure_session_worktree_ready(services, session_id)
.await?;
let session_index = self.session_index_or_err(session_id)?;
let (persisted_session_id, session_model, title) = {
let session = self
.sessions
.get_mut(session_index)
.ok_or(SessionError::NotFound)?;
session.prompt.clone_from(&prompt.text);
let title = prompt.text.clone();
session.title = Some(title.clone());
let session_model = session.model;
(session.id.clone(), session_model, title)
};
let handles = self.session_handles_or_err(&persisted_session_id)?;
let output = Arc::clone(&handles.output);
let status = Arc::clone(&handles.status);
let app_event_tx = services.event_sender();
self.persist_first_message_metadata(services, &persisted_session_id, &prompt.text, &title)
.await;
let initial_output = Self::formatted_prompt_output(&prompt, false);
SessionTaskService::append_session_output(
&output,
services.db(),
&app_event_tx,
&services.session_update_versions(),
&persisted_session_id,
&initial_output,
)
.await;
self.set_active_prompt_output(&persisted_session_id, initial_output);
if !SessionTaskService::update_status(
&status,
services.clock().as_ref(),
services.db(),
&app_event_tx,
&services.session_update_versions(),
&persisted_session_id,
Status::InProgress,
)
.await
{
warn!(
session_id = %persisted_session_id,
"skipped session start status update because the in-memory status did not transition to in-progress"
);
}
let operation_id = Uuid::new_v4().to_string();
let command = SessionCommand::Run {
operation_id,
request_kind: AgentRequestKind::SessionStart,
prompt: prompt.clone(),
turn_metadata: TurnMetadata {
published_upstream_ref: None,
session_model,
},
};
if let Err(error) = self
.enqueue_session_command(services, &persisted_session_id, command)
.await
{
self.cleanup_prompt_attachment_files(services, &prompt)
.await;
return Err(error);
}
Ok(())
}
pub async fn reply(
&mut self,
services: &AppServices,
session_id: &str,
prompt: impl Into<TurnPrompt>,
) {
let prompt = prompt.into();
let Ok(session) = self.session_or_err(session_id) else {
return;
};
let session_model = session.model;
self.reply_impl(services, session_id, prompt, session_model)
.await;
}
pub async fn set_session_model(
&mut self,
services: &AppServices,
session_id: &str,
session_model: AgentModel,
) -> Result<(), SessionError> {
let session_index = self.session_index_or_err(session_id)?;
let model_changed = self
.sessions
.get(session_index)
.is_some_and(|session| session.model != session_model);
services
.db()
.update_session_model(session_id, session_model.as_str())
.await?;
if model_changed {
services
.db()
.update_session_provider_conversation_id(session_id, None)
.await?;
services
.db()
.update_session_instruction_conversation_id(session_id, None)
.await?;
self.clear_session_worker(session_id);
}
let session_project_id = services.db().load_session_project_id(session_id).await?;
if Self::should_persist_last_used_model_as_default(services, session_project_id).await?
&& let Some(project_id) = session_project_id
{
services
.db()
.upsert_project_setting(
project_id,
SettingName::DefaultSmartModel,
session_model.as_str(),
)
.await?;
}
services.emit_app_event(AppEvent::SessionModelUpdated {
session_id: SessionId::from(session_id),
session_model,
});
if model_changed {
self.mark_history_replay_pending(session_id);
}
Ok(())
}
pub async fn set_session_reasoning_level(
&mut self,
services: &AppServices,
session_id: &str,
reasoning_level_override: Option<ReasoningLevel>,
) -> Result<(), SessionError> {
self.session_index_or_err(session_id)?;
services
.db()
.update_session_reasoning_level(
session_id,
reasoning_level_override.map(ReasoningLevel::as_str),
)
.await?;
services.emit_app_event(AppEvent::SessionReasoningLevelUpdated {
reasoning_level_override,
session_id: SessionId::from(session_id),
});
Ok(())
}
async fn should_persist_last_used_model_as_default(
services: &AppServices,
project_id: Option<i64>,
) -> Result<bool, SessionError> {
let Some(project_id) = project_id else {
return Ok(false);
};
let should_persist = services
.db()
.get_project_setting(project_id, SettingName::LastUsedModelAsDefault)
.await?
.and_then(|setting_value| setting_value.parse::<bool>().ok())
.unwrap_or(false);
Ok(should_persist)
}
pub fn selected_session(&self) -> Option<&Session> {
self.state
.table_state
.selected()
.and_then(|index| self.state.sessions.get(index))
}
pub fn session_at(&self, session_index: usize) -> Option<&Session> {
self.state.sessions.get(session_index)
}
pub fn session_id_for_index(&self, session_index: usize) -> Option<SessionId> {
self.state
.sessions
.get(session_index)
.map(|session| session.id.clone())
}
pub fn session_index_for_id(&self, session_id: &str) -> Option<usize> {
self.state.session_index_for_id(session_id)
}
pub async fn publish_review_request(
&mut self,
services: &AppServices,
session_id: &str,
) -> Result<ReviewRequest, SessionError> {
let session_index = self.session_index_or_err(session_id)?;
let Some(session) = self.state.sessions.get(session_index) else {
return Err(SessionError::NotFound);
};
if !session.status.allows_review_actions() {
return Err(SessionError::Workflow(
"Session must be in review to create a review request".to_string(),
));
}
let folder = session.folder.clone();
let source_branch = session_branch(session_id);
let linked_review_request = session.review_request.clone();
let git_client = services.git_client();
let published_upstream_ref = git_client
.push_current_branch(folder.clone())
.await
.map_err(|error| {
SessionError::Workflow(format!("Failed to publish session branch: {error}"))
})?;
self.store_published_upstream_ref(services, session_id, published_upstream_ref)
.await?;
let session = self
.state
.sessions
.get(session_index)
.ok_or(SessionError::NotFound)?;
let review_request_client = services.review_request_client();
let remote = self
.review_request_remote(services, session, linked_review_request.as_ref())
.await?;
let review_request_summary = if let Some(review_request) = linked_review_request {
review_request_client
.refresh_review_request(remote, review_request.summary.display_id)
.await
.map_err(|error| SessionError::Workflow(error.detail_message()))?
} else {
match review_request_client
.find_by_source_branch(remote.clone(), source_branch.clone())
.await
.map_err(|error| SessionError::Workflow(error.detail_message()))?
{
Some(existing_review_request) => review_request_client
.refresh_review_request(remote, existing_review_request.display_id)
.await
.map_err(|error| SessionError::Workflow(error.detail_message()))?,
None => review_request_client
.create_review_request(
remote,
Self::load_review_request_create_input(
git_client.as_ref(),
session,
source_branch.clone(),
)
.await?,
)
.await
.map_err(|error| SessionError::Workflow(error.detail_message()))?,
}
};
self.store_review_request_summary(services, session_id, review_request_summary)
.await
}
pub fn review_request_web_url(
&self,
services: &AppServices,
session_id: &str,
) -> Result<String, SessionError> {
let session = self.session_or_err(session_id)?;
let review_request = session.review_request.as_ref().ok_or_else(|| {
SessionError::Workflow("Session has no linked review request".to_string())
})?;
services
.review_request_client()
.review_request_web_url(&review_request.summary)
.map_err(|error| SessionError::Workflow(error.detail_message()))
}
pub async fn delete_selected_session(
&mut self,
projects: &ProjectManager,
services: &AppServices,
) {
let Some(cleanup) = self
.remove_selected_session_from_state_and_db(projects, services)
.await
else {
return;
};
Self::cleanup_deleted_session_resources(
services.fs_client(),
services.git_client(),
cleanup,
)
.await;
}
pub async fn delete_selected_session_deferred_cleanup(
&mut self,
projects: &ProjectManager,
services: &AppServices,
) {
let Some(cleanup) = self
.remove_selected_session_from_state_and_db(projects, services)
.await
else {
return;
};
let fs_client = services.fs_client();
let git_client = services.git_client();
tokio::spawn(async move {
SessionManager::cleanup_deleted_session_resources(fs_client, git_client, cleanup).await;
});
}
async fn remove_selected_session_from_state_and_db(
&mut self,
projects: &ProjectManager,
services: &AppServices,
) -> Option<DeletedSessionCleanup> {
let selected_index = self.state.table_state.selected()?;
if selected_index >= self.state.sessions.len() {
return None;
}
let session = self.remove_session_at(selected_index)?;
self.state.handles.remove(&session.id);
self.remove_session_worktree_availability(&session.id);
self.remove_at_mention_index_for_root(&session.folder);
self.abort_title_generation_task(&session.id);
self.clear_history_replay_pending(&session.id);
SessionTaskService::remove_session_update_version(
&services.session_update_versions(),
&session.id,
);
if let Err(error) = services
.db()
.request_cancel_for_session_operations(&session.id)
.await
{
warn!(
session_id = %session.id,
error = %error,
"failed to cancel pending session operations during deletion"
);
}
self.clear_session_worker(&session.id);
services.review_comment_cache().forget(&session.id);
if let Err(error) = services.db().delete_session(&session.id).await {
warn!(
session_id = %session.id,
error = %error,
"failed to delete session record during session deletion"
);
}
services.emit_app_event(AppEvent::RefreshSessions);
let staged_draft_root = services.base_path().join(&session.id);
Some(DeletedSessionCleanup {
branch_name: session_branch(&session.id),
folder: session.folder,
has_git_branch: projects.has_git_branch(),
session_id: session.id,
staged_draft_root,
working_dir: projects.working_dir().to_path_buf(),
})
}
async fn cleanup_deleted_session_resources(
fs_client: Arc<dyn FsClient>,
git_client: Arc<dyn git::GitClient>,
cleanup: DeletedSessionCleanup,
) {
let repo_root = if cleanup.has_git_branch {
git_client.find_git_repo_root(cleanup.working_dir).await
} else {
None
};
let cleanup_errors = Self::cleanup_session_worktree_resources(
fs_client.clone(),
git_client,
cleanup.folder,
cleanup.branch_name,
repo_root,
cleanup.has_git_branch,
)
.await;
Self::warn_cleanup_errors(&cleanup.session_id, &cleanup_errors);
if fs_client.is_dir(cleanup.staged_draft_root.clone())
&& let Err(error) = fs_client.remove_dir_all(cleanup.staged_draft_root).await
{
warn!(
session_id = %cleanup.session_id,
error = %error,
"failed to remove staged draft directory during session deletion"
);
}
Self::cleanup_session_temp_directory(fs_client, &cleanup.session_id).await;
}
async fn load_review_request_create_input(
git_client: &dyn git::GitClient,
session: &Session,
source_branch: String,
) -> Result<forge::CreateReviewRequestInput, SessionError> {
let commit_message = git_client
.head_commit_message(session.folder.clone())
.await
.map_err(|error| {
SessionError::Workflow(format!(
"Failed to load session branch commit message: {error}"
))
})?
.ok_or_else(|| {
SessionError::Workflow(
"Session branch has no commit message for review-request publishing."
.to_string(),
)
})?;
let review_request_commit_message = review_request::parse_review_request_commit_message(
&commit_message,
)
.ok_or_else(|| {
SessionError::Workflow(
"Session branch commit message must have a non-empty title for review-request \
publishing."
.to_string(),
)
})?;
Ok(forge::CreateReviewRequestInput {
body: review_request_commit_message.body,
source_branch,
target_branch: session.base_branch.clone(),
title: review_request_commit_message.title,
})
}
pub(super) fn build_review_request(
&self,
summary: forge::ReviewRequestSummary,
) -> ReviewRequest {
ReviewRequest {
last_refreshed_at: unix_timestamp_from_system_time(self.state.clock.now_system_time()),
summary,
}
}
pub(crate) async fn store_review_request_summary(
&mut self,
services: &AppServices,
session_id: &str,
summary: forge::ReviewRequestSummary,
) -> Result<ReviewRequest, SessionError> {
let session_index = self.session_index_or_err(session_id)?;
let review_request = self.build_review_request(summary);
self.store_review_request(services, session_index, review_request)
.await
}
pub(super) async fn store_review_request(
&mut self,
services: &AppServices,
session_index: usize,
review_request: ReviewRequest,
) -> Result<ReviewRequest, SessionError> {
let session_id = self
.state
.sessions
.get(session_index)
.map(|session| session.id.clone())
.ok_or(SessionError::NotFound)?;
services
.db()
.update_session_review_request(&session_id, Some(&review_request))
.await?;
let Some(session) = self.state.sessions.get_mut(session_index) else {
return Err(SessionError::NotFound);
};
session.review_request = Some(review_request.clone());
Ok(review_request)
}
pub(super) async fn store_published_upstream_ref(
&mut self,
services: &AppServices,
session_id: &str,
published_upstream_ref: String,
) -> Result<(), SessionError> {
services
.db()
.update_session_published_upstream_ref(session_id, Some(&published_upstream_ref))
.await?;
let session_index = self.session_index_or_err(session_id)?;
let Some(session) = self.state.sessions.get_mut(session_index) else {
return Err(SessionError::NotFound);
};
session.published_upstream_ref = Some(published_upstream_ref);
Ok(())
}
async fn reply_impl(
&mut self,
services: &AppServices,
session_id: &str,
prompt: TurnPrompt,
session_model: AgentModel,
) {
let Ok(session_index) = self.session_index_or_err(session_id) else {
return;
};
let should_replay_history = self.should_replay_history(session_id);
let (session_output, is_first_message, persisted_session_id, title_to_save) = match self
.prepare_reply_context(session_index, &prompt, session_model, should_replay_history)
{
Ok(Some(reply_context)) => reply_context,
Ok(None) => return,
Err(error) => {
self.append_reply_status_error(services, session_id, &error)
.await;
return;
}
};
if should_replay_history {
self.clear_history_replay_pending(&persisted_session_id);
}
let app_event_tx = services.event_sender();
let Ok(handles) = self.session_handles_or_err(&persisted_session_id) else {
return;
};
let output = Arc::clone(&handles.output);
let status = Arc::clone(&handles.status);
let effective_prompt = prompt;
if let Some(title) = title_to_save {
self.persist_first_message_metadata(
services,
&persisted_session_id,
&effective_prompt.text,
&title,
)
.await;
if !SessionTaskService::update_status(
&status,
services.clock().as_ref(),
services.db(),
&app_event_tx,
&services.session_update_versions(),
&persisted_session_id,
Status::InProgress,
)
.await
{
warn!(
session_id = %persisted_session_id,
"skipped reply status update because the in-memory status did not transition to in-progress"
);
}
}
self.append_reply_prompt_line(
services,
&output,
&app_event_tx,
&persisted_session_id,
&effective_prompt,
)
.await;
let published_upstream_ref = self
.session_or_err(&persisted_session_id)
.ok()
.and_then(|session| session.published_upstream_ref.clone());
let command = Self::build_session_command(BuildSessionCommandInput {
is_first_message,
published_upstream_ref,
prompt: effective_prompt.clone(),
session_model,
session_output,
});
self.enqueue_reply_command(
services,
&output,
&app_event_tx,
&persisted_session_id,
&effective_prompt,
command,
)
.await;
}
fn prepare_reply_context(
&mut self,
session_index: usize,
prompt: &TurnPrompt,
session_model: AgentModel,
should_replay_history: bool,
) -> Result<Option<ReplyContext>, SessionError> {
let Some(session) = self.state.sessions.get_mut(session_index) else {
return Ok(None);
};
let is_first_message = session.prompt.is_empty();
let allowed = session.status.allows_review_actions()
|| session.status == Status::Question
|| (is_first_message && session.status == Status::New);
if !allowed {
return Err(SessionError::Workflow(
"Session must be in review status".to_string(),
));
}
let mut title_to_save = None;
if is_first_message {
session.prompt.clone_from(&prompt.text);
let title = prompt.text.clone();
session.title = Some(title.clone());
title_to_save = Some(title);
}
let session_output = if !is_first_message
&& (should_replay_history
|| agent::transport_mode(session_model.kind()).uses_app_server())
{
Some(session.output.clone())
} else {
None
};
Ok(Some((
session_output,
is_first_message,
session.id.clone(),
title_to_save,
)))
}
async fn persist_first_message_metadata(
&self,
services: &AppServices,
session_id: &str,
prompt: &str,
title: &str,
) {
if let Err(error) = services.db().update_session_title(session_id, title).await {
warn!(
session_id = session_id,
error = %error,
"failed to persist first-message session title"
);
}
if let Err(error) = services
.db()
.update_session_prompt(session_id, prompt)
.await
{
warn!(
session_id = session_id,
error = %error,
"failed to persist first-message session prompt"
);
}
}
async fn append_reply_prompt_line(
&mut self,
services: &AppServices,
output: &Arc<Mutex<String>>,
app_event_tx: &mpsc::UnboundedSender<AppEvent>,
session_id: &str,
prompt: &TurnPrompt,
) {
let reply_line = Self::formatted_prompt_output(prompt, true);
SessionTaskService::append_session_output(
output,
services.db(),
app_event_tx,
&services.session_update_versions(),
session_id,
&reply_line,
)
.await;
self.set_active_prompt_output(session_id, reply_line);
}
fn formatted_prompt_output(prompt: &TurnPrompt, prepend_newline: bool) -> String {
let prompt_text = prompt.transcript_text();
let prompt_lines = prompt_text.split('\n').collect::<Vec<_>>();
let mut formatted_lines = Vec::with_capacity(prompt_lines.len());
for (index, prompt_line) in prompt_lines.into_iter().enumerate() {
let prefix = if index == 0 {
USER_PROMPT_PREFIX
} else {
USER_PROMPT_CONTINUATION_PREFIX
};
formatted_lines.push(format!("{prefix}{prompt_line}"));
}
let prompt_block = formatted_lines.join("\n");
if prepend_newline {
return format!("\n{prompt_block}\n\n");
}
format!("{prompt_block}\n\n")
}
fn append_staged_prompt(
existing_prompt: &str,
prompt: &TurnPrompt,
next_attachment_number: usize,
) -> String {
let staged_prompt = Self::renumbered_prompt_text(prompt, next_attachment_number);
if existing_prompt.is_empty() {
return staged_prompt;
}
format!("{existing_prompt}\n\n{staged_prompt}")
}
fn renumbered_prompt_text(prompt: &TurnPrompt, next_attachment_number: usize) -> String {
let mut prompt_text = prompt.text.clone();
for (offset, attachment) in prompt.attachments.iter().enumerate() {
let placeholder = format!("[Image #{}]", next_attachment_number.saturating_add(offset));
prompt_text = replace_first(&prompt_text, &attachment.placeholder, &placeholder);
}
prompt_text
}
fn renumbered_attachments(
prompt: &TurnPrompt,
next_attachment_number: usize,
) -> Vec<TurnPromptAttachment> {
prompt
.attachments
.iter()
.enumerate()
.map(|(offset, attachment)| TurnPromptAttachment {
placeholder: format!("[Image #{}]", next_attachment_number.saturating_add(offset)),
local_image_path: attachment.local_image_path.clone(),
})
.collect()
}
fn build_session_command(input: BuildSessionCommandInput) -> SessionCommand {
let BuildSessionCommandInput {
is_first_message,
published_upstream_ref,
prompt,
session_model,
session_output,
} = input;
let operation_id = Uuid::new_v4().to_string();
let request_kind = if is_first_message {
AgentRequestKind::SessionStart
} else {
AgentRequestKind::SessionResume { session_output }
};
SessionCommand::Run {
operation_id,
request_kind,
prompt,
turn_metadata: TurnMetadata {
published_upstream_ref,
session_model,
},
}
}
async fn append_reply_status_error(
&self,
services: &AppServices,
session_id: &str,
error: &SessionError,
) {
let status_error = format!("\n[Reply Error] {error}\n");
let Ok(handles) = self.session_handles_or_err(session_id) else {
return;
};
let app_event_tx = services.event_sender();
SessionTaskService::append_session_output(
&handles.output,
services.db(),
&app_event_tx,
&services.session_update_versions(),
session_id,
&status_error,
)
.await;
}
async fn enqueue_reply_command(
&mut self,
services: &AppServices,
output: &Arc<Mutex<String>>,
app_event_tx: &mpsc::UnboundedSender<AppEvent>,
persisted_session_id: &str,
prompt: &TurnPrompt,
command: SessionCommand,
) {
if let Err(error) = self
.enqueue_session_command(services, persisted_session_id, command)
.await
{
self.cleanup_prompt_attachment_files(services, prompt).await;
let error_line = format!("\n[Reply Error] {error}\n");
SessionTaskService::append_session_output(
output,
services.db(),
app_event_tx,
&services.session_update_versions(),
persisted_session_id,
&error_line,
)
.await;
}
}
pub(crate) fn spawn_session_title_generation_task(
app_event_tx: mpsc::UnboundedSender<AppEvent>,
db: db::AppRepositories,
session_id: &str,
folder: &Path,
prompt: &str,
session_model: AgentModel,
tracked_generation: Option<u64>,
) -> tokio::task::JoinHandle<()> {
let folder = folder.to_path_buf();
let prompt = prompt.to_string();
let persisted_session_id = SessionId::from(session_id);
let tracked_completion =
tracked_generation.map(|generation| TitleGenerationTaskCompletion {
generation,
session_id: persisted_session_id.clone(),
});
tokio::spawn(async move {
let Ok(title_generation_prompt) =
SessionManager::session_title_generation_prompt(&prompt)
else {
SessionManager::emit_title_generation_finished_event(
&app_event_tx,
tracked_completion.as_ref(),
);
return;
};
let Some(title_response) = SessionManager::run_title_generation_command(
folder.as_path(),
&title_generation_prompt,
session_model,
)
.await
else {
SessionManager::emit_title_generation_finished_event(
&app_event_tx,
tracked_completion.as_ref(),
);
return;
};
let Some(generated_title) =
SessionManager::parse_generated_session_title(&title_response)
else {
SessionManager::emit_title_generation_finished_event(
&app_event_tx,
tracked_completion.as_ref(),
);
return;
};
if generated_title == prompt {
SessionManager::emit_title_generation_finished_event(
&app_event_tx,
tracked_completion.as_ref(),
);
return;
}
match db
.update_session_title_for_prompt(&persisted_session_id, &prompt, &generated_title)
.await
{
Ok(true) => {
if app_event_tx.send(AppEvent::RefreshSessions).is_err() {
warn!(
session_id = %persisted_session_id,
"failed to refresh sessions after title generation because the app event receiver is closed"
);
}
}
Ok(false) => {}
Err(error) => {
warn!(
session_id = %persisted_session_id,
error = %error,
"failed to persist generated session title"
);
}
}
SessionManager::emit_title_generation_finished_event(
&app_event_tx,
tracked_completion.as_ref(),
);
})
}
fn emit_title_generation_finished_event(
app_event_tx: &mpsc::UnboundedSender<AppEvent>,
tracked_completion: Option<&TitleGenerationTaskCompletion>,
) {
let Some(tracked_completion) = tracked_completion else {
return;
};
if app_event_tx
.send(AppEvent::SessionTitleGenerationFinished {
generation: tracked_completion.generation,
session_id: tracked_completion.session_id.clone(),
})
.is_err()
{
warn!(
session_id = %tracked_completion.session_id,
generation = tracked_completion.generation,
"failed to send session title generation completion event because the app event receiver is closed"
);
}
}
async fn run_title_generation_command(
folder: &Path,
prompt: &str,
model: AgentModel,
) -> Option<String> {
let response = agent::submit_one_shot(agent::OneShotRequest {
child_pid: None,
folder,
model,
prompt,
request_kind: AgentRequestKind::UtilityPrompt,
reasoning_level: ReasoningLevel::default(),
})
.await
.ok()?;
Some(response.to_answer_display_text())
}
fn session_title_generation_prompt(prompt: &str) -> Result<String, SessionError> {
let template = SessionTitleGenerationPromptTemplate { prompt };
template.render().map_err(|error| {
SessionError::Workflow(format!(
"Failed to render `session_title_generation_prompt.md`: {error}"
))
})
}
fn parse_generated_session_title(content: &str) -> Option<String> {
let content = content.trim();
if content.is_empty() {
return None;
}
if let Ok(protocol_response) = agent::protocol::parse_agent_response_strict(content) {
return Self::parse_generated_session_title_from_protocol_response(&protocol_response);
}
let first_line = Self::first_nonempty_line(content)?;
Self::normalize_generated_session_title(first_line)
}
fn parse_generated_session_title_from_protocol_response(
protocol_response: &agent::protocol::AgentResponse,
) -> Option<String> {
for answer in protocol_response.answers() {
if let Some(first_line) = Self::first_nonempty_line(&answer)
&& let Some(parsed_title) = Self::normalize_generated_session_title(first_line)
{
return Some(parsed_title);
}
}
None
}
fn first_nonempty_line(content: &str) -> Option<&str> {
content.lines().find_map(|line| {
let trimmed_line = line.trim();
if trimmed_line.is_empty() {
return None;
}
Some(trimmed_line)
})
}
fn normalize_generated_session_title(candidate: &str) -> Option<String> {
let mut title = candidate.trim().to_string();
if let Some((prefix, remainder)) = title.split_once(':')
&& prefix.trim().eq_ignore_ascii_case("title")
{
title = remainder.trim().to_string();
}
title = title
.trim_matches(|ch| matches!(ch, '"' | '\'' | '`'))
.trim()
.to_string();
if title.is_empty() {
return None;
}
let truncated = title
.chars()
.take(GENERATED_SESSION_TITLE_MAX_CHARACTERS)
.collect::<String>();
Some(truncated)
}
async fn resolve_default_session_model(
&self,
services: &AppServices,
project_id: i64,
) -> AgentModel {
setting::load_default_smart_model_setting(
services,
Some(project_id),
self.default_session_model,
)
.await
}
async fn rollback_failed_session_creation(
&self,
services: &AppServices,
folder: &Path,
repo_root: &Path,
session_id: &str,
worktree_branch: &str,
session_saved: bool,
) {
if session_saved {
if let Err(error) = services.db().delete_session(session_id).await {
warn!(
session_id = session_id,
error = %error,
"failed to roll back persisted session metadata"
);
}
SessionTaskService::remove_session_update_version(
&services.session_update_versions(),
session_id,
);
}
{
let git_client = services.git_client();
let folder = folder.to_path_buf();
let repo_root = repo_root.to_path_buf();
let worktree_branch = worktree_branch.to_string();
if let Err(error) = git_client.remove_worktree(folder).await {
warn!(
session_id = session_id,
error = %error,
"failed to remove worktree while rolling back session creation"
);
}
if let Err(error) = git_client.delete_branch(repo_root, worktree_branch).await {
warn!(
session_id = session_id,
error = %error,
"failed to delete branch while rolling back session creation"
);
}
}
if let Err(error) = services
.fs_client()
.remove_dir_all(folder.to_path_buf())
.await
{
warn!(
session_id = session_id,
error = %error,
"failed to remove session worktree directory while rolling back session creation"
);
}
Self::cleanup_session_temp_directory(services.fs_client(), session_id).await;
}
async fn record_session_creation_activity(services: &AppServices, session_id: &str) {
if let Err(error) = services
.db()
.insert_session_creation_activity_now(session_id)
.await
{
warn!(
session_id = session_id,
error = %error,
"failed to record session creation activity"
);
}
}
pub(crate) async fn append_output_for_session(
&self,
services: &AppServices,
session_id: &str,
output: &str,
) {
let Ok((session, handles)) = self.session_and_handles_or_err(session_id) else {
return;
};
let app_event_tx = services.event_sender();
SessionTaskService::append_session_output(
&handles.output,
services.db(),
&app_event_tx,
&services.session_update_versions(),
&session.id,
output,
)
.await;
}
pub(crate) async fn cleanup_prompt_attachment_files(
&self,
services: &AppServices,
prompt: &TurnPrompt,
) {
Self::cleanup_prompt_attachment_paths(
services.fs_client(),
prompt.local_image_paths().cloned().collect(),
)
.await;
}
pub async fn cancel_session(
&self,
services: &AppServices,
session_id: &str,
) -> Result<(), SessionError> {
let session = self.session_or_err(session_id)?;
if !session.allows_cancel_action() {
return Err(SessionError::Workflow(
"Session must be in review or be an unstarted draft to be canceled".to_string(),
));
}
let branch_name = session_branch(&session.id);
let folder = session.folder.clone();
let has_worktree = services.fs_client().is_dir(folder.clone());
let handles = self.session_handles_or_err(session_id)?;
let status = Arc::clone(&handles.status);
let app_event_tx = services.event_sender();
let status_updated = SessionTaskService::update_status(
&status,
services.clock().as_ref(),
services.db(),
&app_event_tx,
&services.session_update_versions(),
session_id,
Status::Canceled,
)
.await;
if status_updated {
if has_worktree {
let repo_root = services
.git_client()
.main_repo_root(folder.clone())
.await
.ok();
let cleanup_errors = Self::cleanup_session_worktree_resources(
services.fs_client().clone(),
services.git_client(),
folder,
branch_name,
repo_root,
true,
)
.await;
Self::warn_cleanup_errors(session_id, &cleanup_errors);
}
Self::cleanup_session_temp_directory(services.fs_client(), session_id).await;
}
Ok(())
}
#[must_use]
async fn cleanup_session_worktree_resources(
fs_client: Arc<dyn FsClient>,
git_client: Arc<dyn git::GitClient>,
folder: PathBuf,
branch_name: String,
repo_root: Option<PathBuf>,
remove_git_resources: bool,
) -> Vec<String> {
let mut cleanup_errors = Vec::new();
if remove_git_resources {
if let Err(error) = git_client.remove_worktree(folder.clone()).await {
cleanup_errors.push(format!("failed to remove worktree: {error}"));
}
if let Some(repo_root) = repo_root
&& let Err(error) = git_client.delete_branch(repo_root, branch_name).await
{
cleanup_errors.push(format!("failed to delete branch: {error}"));
}
}
if let Err(error) = fs_client.remove_dir_all(folder).await {
cleanup_errors.push(format!("failed to remove worktree directory: {error}"));
}
cleanup_errors
}
fn warn_cleanup_errors(session_id: &str, cleanup_errors: &[String]) {
for cleanup_error in cleanup_errors {
warn!(session_id = session_id, "{cleanup_error}");
}
}
pub(crate) async fn cleanup_prompt_attachment_paths(
fs_client: Arc<dyn FsClient>,
attachment_paths: Vec<PathBuf>,
) {
Self::cleanup_prompt_attachment_paths_in_root(
fs_client,
&prompt_attachment_tmp_root(),
attachment_paths,
)
.await;
}
async fn cleanup_prompt_attachment_paths_in_root(
fs_client: Arc<dyn FsClient>,
managed_tmp_root: &Path,
attachment_paths: Vec<PathBuf>,
) {
if attachment_paths.is_empty() {
return;
}
let image_directory =
managed_prompt_attachment_directory(&attachment_paths, managed_tmp_root);
for attachment_path in attachment_paths {
if is_managed_prompt_attachment_path(&attachment_path, managed_tmp_root)
&& let Err(error) = fs_client.remove_file(attachment_path).await
{
warn!(
error = %error,
"failed to remove managed prompt attachment file"
);
}
}
if let Some(image_directory) = image_directory
&& let Err(error) = fs_client.remove_dir_all(image_directory).await
{
warn!(
error = %error,
"failed to remove managed prompt attachment directory"
);
}
}
async fn cleanup_session_temp_directory(fs_client: Arc<dyn FsClient>, session_id: &str) {
if let Err(error) = fs_client
.remove_dir_all(session_prompt_temp_directory(session_id))
.await
{
warn!(
session_id = session_id,
error = %error,
"failed to remove session prompt temp directory"
);
}
}
}
fn replace_first(haystack: &str, needle: &str, replacement: &str) -> String {
let Some(match_index) = haystack.find(needle) else {
return haystack.to_string();
};
let mut replaced = String::with_capacity(
haystack
.len()
.saturating_sub(needle.len())
.saturating_add(replacement.len()),
);
replaced.push_str(&haystack[..match_index]);
replaced.push_str(replacement);
replaced.push_str(&haystack[match_index + needle.len()..]);
replaced
}
fn session_prompt_temp_directory(session_id: &str) -> PathBuf {
agentty_home().join("tmp").join(session_id)
}
fn prompt_attachment_tmp_root() -> PathBuf {
agentty_home().join("tmp")
}
fn managed_prompt_attachment_directory(
attachment_paths: &[PathBuf],
managed_tmp_root: &Path,
) -> Option<PathBuf> {
let image_directory = attachment_paths.first()?.parent()?.to_path_buf();
if !is_managed_prompt_attachment_directory(&image_directory, managed_tmp_root) {
return None;
}
attachment_paths
.iter()
.all(|attachment_path| {
attachment_path.parent() == Some(image_directory.as_path())
&& is_managed_prompt_attachment_path(attachment_path, managed_tmp_root)
})
.then_some(image_directory)
}
fn is_managed_prompt_attachment_path(path: &Path, managed_tmp_root: &Path) -> bool {
path.parent().is_some_and(|parent| {
is_managed_prompt_attachment_directory(parent, managed_tmp_root)
&& path.starts_with(managed_tmp_root)
})
}
fn is_managed_prompt_attachment_directory(path: &Path, managed_tmp_root: &Path) -> bool {
path.starts_with(managed_tmp_root) && path.ends_with("images")
}
#[cfg(test)]
mod test_support {
use std::sync::Arc;
use super::*;
impl SessionManager {
pub(crate) async fn reply_with_backend(
&mut self,
services: &AppServices,
session_id: &str,
prompt: impl Into<TurnPrompt>,
backend: Arc<dyn agent::AgentBackend>,
session_model: AgentModel,
) {
let prompt = prompt.into();
let channel: Arc<dyn crate::infra::channel::AgentChannel> =
Arc::new(crate::infra::channel::cli::CliAgentChannel::with_backend(
backend,
session_model.kind(),
));
self.worker_service
.test_agent_channels
.insert(session_id.to_string().into(), channel);
self.reply_impl(services, session_id, prompt, session_model)
.await;
}
}
}
#[cfg(test)]
mod tests {
use std::collections::HashMap;
use std::sync::Arc;
use ag_forge as forge;
use ratatui::widgets::TableState;
use tokio::sync::mpsc;
use super::*;
use crate::app::session::{RealClock, SessionDefaults};
use crate::app::{AppEvent, AppServices, SessionState};
use crate::domain::agent::{AgentKind, ReasoningLevel};
use crate::domain::session::{
ForgeKind, ReviewRequestState, ReviewRequestSummary, SessionHandles, SessionSize,
SessionStats,
};
use crate::infra::channel::{TurnPromptAttachment, TurnPromptTextSource};
use crate::infra::db::{self, AppRepositories};
use crate::infra::{app_server, fs};
fn session_manager_with_one_session(session: Session) -> SessionManager {
let mut handles = HashMap::new();
handles.insert(
session.id.clone(),
SessionHandles::new(session.output.clone(), session.status),
);
let state = SessionState::new(
handles,
vec![session],
TableState::default(),
Arc::new(RealClock),
1,
0,
);
SessionManager::new(
SessionDefaults {
model: AgentModel::Gpt54,
},
Arc::new(git::MockGitClient::new()),
state,
Vec::new(),
)
}
fn test_session(prompt: &str, status: Status, title: Option<&str>, output: &str) -> Session {
Session {
base_branch: "main".to_string(),
created_at: 0,
draft_attachments: Vec::new(),
folder: PathBuf::from("/tmp/session"),
follow_up_tasks: Vec::new(),
id: "session-id".into(),
in_progress_started_at: None,
in_progress_total_seconds: 0,
is_draft: false,
model: AgentModel::ClaudeSonnet46,
output: output.to_string(),
project_name: "project".to_string(),
prompt: prompt.to_string(),
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: title.map(ToString::to_string),
updated_at: 0,
}
}
fn mock_app_server() -> Arc<dyn app_server::AppServerClient> {
Arc::new(app_server::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(|_| Box::pin(async { Ok(()) }));
mock_fs_client
.expect_remove_dir_all()
.times(0..)
.returning(|_| Box::pin(async { Ok(()) }));
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(|_| Box::pin(async { Ok(()) }));
mock_fs_client
.expect_exists()
.times(0..)
.returning(|path| path.exists());
mock_fs_client
.expect_is_dir()
.times(0..)
.returning(|path| path.is_dir());
mock_fs_client
}
async fn database_with_session(session: &Session) -> AppRepositories {
let database = AppRepositories::in_memory().await;
let project_id = database
.upsert_project("/tmp/project", Some("main"))
.await
.expect("failed to upsert project");
if session.is_draft {
database
.insert_draft_session(
&session.id,
session.model.as_str(),
&session.base_branch,
&session.status.to_string(),
project_id,
)
.await
.expect("failed to insert draft session");
} else {
database
.insert_session(
&session.id,
session.model.as_str(),
&session.base_branch,
&session.status.to_string(),
project_id,
)
.await
.expect("failed to insert session");
}
database
.update_session_prompt(&session.id, &session.prompt)
.await
.expect("failed to persist session prompt");
if let Some(title) = &session.title {
database
.update_session_title(&session.id, title)
.await
.expect("failed to persist session title");
}
if let Some(review_request) = &session.review_request {
database
.update_session_review_request(&session.id, Some(review_request))
.await
.expect("failed to persist session review request");
}
database
}
fn test_services_with_fs_client(
database: &AppRepositories,
fs_client: Arc<dyn fs::FsClient>,
git_client: Arc<dyn git::GitClient>,
review_request_client: Arc<dyn forge::ReviewRequestClient>,
) -> AppServices {
let (event_tx, _event_rx) = mpsc::unbounded_channel();
AppServices::new(
PathBuf::from("/tmp/agentty-tests"),
Arc::new(crate::app::session::RealClock),
event_tx,
crate::app::service::AppServiceDeps {
app_server_client_override: Some(mock_app_server()),
fs_client,
available_agent_kinds: AgentKind::ALL.to_vec(),
git_client,
repositories: database.clone(),
review_request_client,
},
)
}
fn test_services(
database: &AppRepositories,
git_client: Arc<dyn git::GitClient>,
review_request_client: Arc<dyn forge::ReviewRequestClient>,
) -> AppServices {
test_services_with_fs_client(
database,
Arc::new(create_passthrough_mock_fs_client()),
git_client,
review_request_client,
)
}
fn test_services_with_event_receiver(
database: &AppRepositories,
git_client: Arc<dyn git::GitClient>,
review_request_client: Arc<dyn forge::ReviewRequestClient>,
) -> (AppServices, mpsc::UnboundedReceiver<AppEvent>) {
let (event_tx, event_rx) = mpsc::unbounded_channel();
let services = AppServices::new(
PathBuf::from("/tmp/agentty-tests"),
Arc::new(crate::app::session::RealClock),
event_tx,
crate::app::service::AppServiceDeps {
app_server_client_override: Some(mock_app_server()),
fs_client: Arc::new(create_passthrough_mock_fs_client()),
available_agent_kinds: AgentKind::ALL.to_vec(),
git_client,
repositories: database.clone(),
review_request_client,
},
);
(services, event_rx)
}
fn review_request_summary(display_id: &str) -> ReviewRequestSummary {
ReviewRequestSummary {
display_id: display_id.to_string(),
forge_kind: ForgeKind::GitHub,
source_branch: session_branch("session-id"),
state: ReviewRequestState::Open,
status_summary: Some("Checks pending".to_string()),
target_branch: "main".to_string(),
title: "Add forge review support".to_string(),
web_url: format!(
"https://github.com/agentty-xyz/agentty/pull/{}",
&display_id[1..]
),
}
}
fn github_remote() -> forge::ForgeRemote {
forge::ForgeRemote {
command_working_directory: Some(PathBuf::from("/tmp/session")),
forge_kind: ForgeKind::GitHub,
host: "github.com".to_string(),
namespace: "agentty-xyz".to_string(),
project: "agentty".to_string(),
repo_url: "https://github.com/agentty-xyz/agentty.git".to_string(),
web_url: "https://github.com/agentty-xyz/agentty".to_string(),
}
}
fn expected_create_input() -> forge::CreateReviewRequestInput {
forge::CreateReviewRequestInput {
body: Some("- Keep title in sync".to_string()),
source_branch: session_branch("session-id"),
target_branch: "main".to_string(),
title: "Refine session commit message".to_string(),
}
}
fn expect_published_session_branch(mock_git_client: &mut git::MockGitClient) {
mock_git_client
.expect_push_current_branch()
.times(1)
.returning(|_| Box::pin(async { Ok("origin/wt/session-id".to_string()) }));
mock_git_client.expect_repo_url().times(1).returning(|_| {
Box::pin(async { Ok("https://github.com/agentty-xyz/agentty.git".to_string()) })
});
}
async fn load_persisted_session_row(database: &AppRepositories) -> db::SessionRow {
database
.load_sessions()
.await
.expect("failed to load session rows")
.into_iter()
.find(|row| row.id == "session-id")
.expect("session row should exist")
}
#[tokio::test]
async fn test_set_session_reasoning_level_persists_override_and_emits_event() {
let session = test_session("Prompt", Status::Review, Some("Title"), "");
let database = database_with_session(&session).await;
let mut session_manager = session_manager_with_one_session(session);
let (services, mut event_rx) = test_services_with_event_receiver(
&database,
Arc::new(git::MockGitClient::new()),
Arc::new(forge::MockReviewRequestClient::new()),
);
session_manager
.set_session_reasoning_level(&services, "session-id", Some(ReasoningLevel::High))
.await
.expect("reasoning level update should succeed");
let persisted_reasoning_level = database
.load_session_reasoning_level_override("session-id")
.await
.expect("reasoning override should load");
let emitted_event = event_rx
.try_recv()
.expect("expected reasoning update event");
assert_eq!(persisted_reasoning_level, Some(ReasoningLevel::High));
assert_eq!(
emitted_event,
AppEvent::SessionReasoningLevelUpdated {
reasoning_level_override: Some(ReasoningLevel::High),
session_id: "session-id".into(),
}
);
assert!(event_rx.try_recv().is_err());
}
#[tokio::test]
async fn test_publish_review_request_creates_and_persists_link_when_lookup_misses() {
let session = test_session(
"Implement forge review support",
Status::Review,
Some("Add forge review support"),
"",
);
let database = database_with_session(&session).await;
let mut session_manager = session_manager_with_one_session(session);
let source_branch = session_branch("session-id");
let expected_create_input = expected_create_input();
let remote = github_remote();
let created_summary = review_request_summary("#42");
let mut mock_git_client = git::MockGitClient::new();
expect_published_session_branch(&mut mock_git_client);
mock_git_client
.expect_head_commit_message()
.times(1)
.returning(|_| {
Box::pin(async {
Ok(Some(
"Refine session commit message\n\n- Keep title in sync".to_string(),
))
})
});
let mut mock_review_request_client = forge::MockReviewRequestClient::new();
mock_review_request_client
.expect_detect_remote()
.times(1)
.returning({
let remote = remote.clone();
move |_| Ok(remote.clone())
});
mock_review_request_client
.expect_find_by_source_branch()
.times(1)
.withf({
let remote = remote.clone();
let source_branch = source_branch.clone();
move |candidate_remote, candidate_source_branch| {
candidate_remote == &remote && candidate_source_branch == &source_branch
}
})
.returning(|_, _| Box::pin(async { Ok(None) }));
mock_review_request_client
.expect_create_review_request()
.times(1)
.withf({
let remote = remote.clone();
let expected_create_input = expected_create_input.clone();
move |candidate_remote, candidate_input| {
candidate_remote == &remote && candidate_input == &expected_create_input
}
})
.returning(move |_, _| {
let created_summary = created_summary.clone();
Box::pin(async move { Ok(created_summary) })
});
let services = test_services(
&database,
Arc::new(mock_git_client),
Arc::new(mock_review_request_client),
);
let review_request = session_manager
.publish_review_request(&services, "session-id")
.await
.expect("review request should be created");
let persisted_row = load_persisted_session_row(&database).await;
assert_eq!(review_request.summary.display_id, "#42");
assert_eq!(
session_manager.state.sessions[0].review_request,
Some(review_request.clone())
);
assert_eq!(
persisted_row
.review_request
.as_ref()
.map(|row| row.display_id.as_str()),
Some("#42")
);
assert_eq!(
persisted_row
.review_request
.as_ref()
.map(|row| row.last_refreshed_at),
Some(review_request.last_refreshed_at)
);
assert_eq!(
persisted_row.published_upstream_ref.as_deref(),
Some("origin/wt/session-id")
);
assert_eq!(
session_manager.state.sessions[0]
.published_upstream_ref
.as_deref(),
Some("origin/wt/session-id")
);
}
#[tokio::test]
async fn test_stage_draft_message_preserves_persisted_prompt_when_attachment_write_fails() {
let mut session = test_session("", Status::New, None, "");
session.is_draft = true;
let database = database_with_session(&session).await;
let mut session_manager = session_manager_with_one_session(session);
let mut mock_fs_client = fs::MockFsClient::new();
mock_fs_client
.expect_create_dir_all()
.times(0..)
.returning(|_| Box::pin(async { Ok(()) }));
mock_fs_client
.expect_remove_dir_all()
.times(0..)
.returning(|_| Box::pin(async { Ok(()) }));
mock_fs_client
.expect_read_file()
.times(0..)
.returning(|_| Box::pin(async { Ok(Vec::new()) }));
mock_fs_client
.expect_remove_file()
.times(0..)
.returning(|_| Box::pin(async { Ok(()) }));
mock_fs_client
.expect_exists()
.times(0..)
.returning(|path| path.exists());
mock_fs_client
.expect_is_dir()
.times(0..)
.returning(|path| path.is_dir());
mock_fs_client.expect_write_file().once().returning(|_, _| {
Box::pin(async {
Err(fs::FsError::Io(std::io::Error::other(
"simulated attachment write failure",
)))
})
});
let services = test_services_with_fs_client(
&database,
Arc::new(mock_fs_client),
Arc::new(git::MockGitClient::new()),
Arc::new(forge::MockReviewRequestClient::new()),
);
let prompt = TurnPrompt {
attachments: vec![TurnPromptAttachment {
placeholder: "[Image #1]".to_string(),
local_image_path: PathBuf::from("/tmp/image-1.png"),
}],
text: "Review [Image #1]".to_string(),
text_source: TurnPromptTextSource::UserPrompt,
};
let error = session_manager
.stage_draft_message(&services, "session-id", prompt)
.await
.expect_err("attachment metadata failure should abort draft staging");
let persisted_session = load_persisted_session_row(&database).await;
assert!(matches!(error, SessionError::Fs(_)));
assert!(persisted_session.prompt.is_empty());
assert!(session_manager.sessions[0].prompt.is_empty());
assert!(session_manager.sessions[0].draft_attachments.is_empty());
}
#[tokio::test]
async fn test_ensure_session_worktree_ready_skips_non_draft_sessions() {
let session = test_session("", Status::New, None, "");
let database = database_with_session(&session).await;
let mut session_manager = session_manager_with_one_session(session);
let mut mock_fs_client = fs::MockFsClient::new();
mock_fs_client.expect_is_dir().times(0);
let mut mock_git_client = git::MockGitClient::new();
mock_git_client.expect_create_worktree().times(0);
mock_git_client.expect_find_git_repo_root().times(0);
let services = test_services_with_fs_client(
&database,
Arc::new(mock_fs_client),
Arc::new(mock_git_client),
Arc::new(forge::MockReviewRequestClient::new()),
);
let result = session_manager
.ensure_session_worktree_ready(&services, "session-id")
.await;
assert!(result.is_ok());
}
#[tokio::test]
async fn test_ensure_session_worktree_ready_reuses_existing_draft_worktree() {
let mut session = test_session("", Status::New, None, "");
session.is_draft = true;
let database = database_with_session(&session).await;
let mut session_manager = session_manager_with_one_session(session);
let mut mock_fs_client = fs::MockFsClient::new();
mock_fs_client.expect_is_dir().once().return_const(true);
let mut mock_git_client = git::MockGitClient::new();
mock_git_client.expect_create_worktree().times(0);
mock_git_client.expect_find_git_repo_root().times(0);
let services = test_services_with_fs_client(
&database,
Arc::new(mock_fs_client),
Arc::new(mock_git_client),
Arc::new(forge::MockReviewRequestClient::new()),
);
let result = session_manager
.ensure_session_worktree_ready(&services, "session-id")
.await;
assert!(result.is_ok());
}
#[tokio::test]
async fn test_cleanup_session_worktree_resources_collects_cleanup_errors() {
let mut mock_fs_client = fs::MockFsClient::new();
mock_fs_client
.expect_remove_dir_all()
.once()
.returning(|_| {
Box::pin(async {
Err(fs::FsError::Io(std::io::Error::other(
"simulated directory cleanup failure",
)))
})
});
let mut mock_git_client = git::MockGitClient::new();
mock_git_client
.expect_remove_worktree()
.once()
.returning(|_| {
Box::pin(async {
Err(git::GitError::CommandFailed {
command: "git worktree remove".to_string(),
stderr: "simulated worktree removal failure".to_string(),
})
})
});
mock_git_client
.expect_delete_branch()
.once()
.returning(|_, _| {
Box::pin(async {
Err(git::GitError::CommandFailed {
command: "git branch -D".to_string(),
stderr: "simulated branch deletion failure".to_string(),
})
})
});
let cleanup_errors = SessionManager::cleanup_session_worktree_resources(
Arc::new(mock_fs_client),
Arc::new(mock_git_client),
PathBuf::from("/tmp/session"),
"wt/session-id".to_string(),
Some(PathBuf::from("/tmp/repo")),
true,
)
.await;
assert_eq!(cleanup_errors.len(), 3);
assert!(
cleanup_errors
.iter()
.any(|message| message.contains("failed to remove worktree"))
);
assert!(
cleanup_errors
.iter()
.any(|message| message.contains("failed to delete branch"))
);
assert!(
cleanup_errors
.iter()
.any(|message| message.contains("failed to remove worktree directory"))
);
}
#[tokio::test]
async fn test_publish_review_request_reuses_existing_remote_link_before_create() {
let session = test_session(
"Implement forge review support",
Status::Review,
Some("Add forge review support"),
"",
);
let database = database_with_session(&session).await;
let mut session_manager = session_manager_with_one_session(session);
let source_branch = session_branch("session-id");
let remote = github_remote();
let existing_summary = review_request_summary("#24");
let mut mock_git_client = git::MockGitClient::new();
expect_published_session_branch(&mut mock_git_client);
let mut mock_review_request_client = forge::MockReviewRequestClient::new();
mock_review_request_client
.expect_detect_remote()
.times(1)
.returning({
let remote = remote.clone();
move |_| Ok(remote.clone())
});
mock_review_request_client
.expect_find_by_source_branch()
.times(1)
.withf({
let remote = remote.clone();
let source_branch = source_branch.clone();
move |candidate_remote, candidate_source_branch| {
candidate_remote == &remote && candidate_source_branch == &source_branch
}
})
.returning(move |_, _| {
let existing_summary = existing_summary.clone();
Box::pin(async move { Ok(Some(existing_summary)) })
});
mock_review_request_client
.expect_refresh_review_request()
.times(1)
.withf({
let remote = remote.clone();
move |candidate_remote, display_id| {
candidate_remote == &remote && display_id == "#24"
}
})
.returning(|_, _| Box::pin(async { Ok(review_request_summary("#24")) }));
mock_review_request_client
.expect_create_review_request()
.times(0);
let services = test_services(
&database,
Arc::new(mock_git_client),
Arc::new(mock_review_request_client),
);
let review_request = session_manager
.publish_review_request(&services, "session-id")
.await
.expect("existing remote review request should be reused");
assert_eq!(review_request.summary.display_id, "#24");
assert_eq!(
session_manager.state.sessions[0]
.review_request
.as_ref()
.map(|review_request| review_request.summary.display_id.as_str()),
Some("#24")
);
}
#[tokio::test]
async fn test_publish_review_request_refreshes_stored_link_after_push() {
let mut session = test_session(
"Implement forge review support",
Status::Review,
Some("Add forge review support"),
"",
);
session.review_request = Some(ReviewRequest {
last_refreshed_at: 42,
summary: review_request_summary("#11"),
});
let database = database_with_session(&session).await;
let mut session_manager = session_manager_with_one_session(session);
let remote = github_remote();
let refreshed_summary = review_request_summary("#11");
let mut mock_git_client = git::MockGitClient::new();
expect_published_session_branch(&mut mock_git_client);
let mut mock_review_request_client = forge::MockReviewRequestClient::new();
mock_review_request_client
.expect_detect_remote()
.times(1)
.returning({
let remote = remote.clone();
move |_| Ok(remote.clone())
});
mock_review_request_client
.expect_refresh_review_request()
.times(1)
.withf({
let remote = remote.clone();
move |candidate_remote, display_id| {
candidate_remote == &remote && display_id == "#11"
}
})
.returning(move |_, _| {
let refreshed_summary = refreshed_summary.clone();
Box::pin(async move { Ok(refreshed_summary) })
});
let services = test_services(
&database,
Arc::new(mock_git_client),
Arc::new(mock_review_request_client),
);
let review_request = session_manager
.publish_review_request(&services, "session-id")
.await
.expect("stored review request should be refreshed");
let persisted_row = load_persisted_session_row(&database).await;
assert_eq!(review_request.summary.display_id, "#11");
assert!(review_request.last_refreshed_at >= 42);
assert_eq!(
session_manager.state.sessions[0]
.published_upstream_ref
.as_deref(),
Some("origin/wt/session-id")
);
assert_eq!(
persisted_row
.review_request
.as_ref()
.map(|row| row.display_id.as_str()),
Some("#11")
);
}
#[tokio::test]
async fn test_review_request_web_url_returns_linked_review_request_url() {
let mut session = test_session(
"Implement forge review support",
Status::Done,
Some("Add forge review support"),
"",
);
session.review_request = Some(ReviewRequest {
last_refreshed_at: 42,
summary: review_request_summary("#11"),
});
let session_manager = session_manager_with_one_session(session);
let database = database_with_session(
session_manager
.state
.sessions
.first()
.expect("fixture session should exist"),
)
.await;
let mut mock_review_request_client = forge::MockReviewRequestClient::new();
mock_review_request_client
.expect_review_request_web_url()
.times(1)
.returning(|summary| Ok(summary.web_url.clone()));
let services = test_services(
&database,
Arc::new(git::MockGitClient::new()),
Arc::new(mock_review_request_client),
);
let review_request_url = session_manager
.review_request_web_url(&services, "session-id")
.expect("linked review request URL should be returned");
assert_eq!(
review_request_url,
"https://github.com/agentty-xyz/agentty/pull/11"
);
}
#[test]
fn test_formatted_prompt_output_formats_multiline_prompt_with_continuation_prefix() {
let prompt = TurnPrompt::from_text("first line\n\n\nafter gap".to_string());
let formatted_prompt = SessionManager::formatted_prompt_output(&prompt, false);
assert_eq!(
formatted_prompt,
" › first line\n \n \n after gap\n\n"
);
}
#[test]
fn test_formatted_prompt_output_prepends_newline_for_replies() {
let prompt = TurnPrompt::from_text("reply line".to_string());
let formatted_prompt = SessionManager::formatted_prompt_output(&prompt, true);
assert_eq!(formatted_prompt, "\n › reply line\n\n");
}
#[test]
fn test_formatted_prompt_output_preserves_image_placeholders_in_transcript() {
let prompt = TurnPrompt {
attachments: vec![TurnPromptAttachment {
placeholder: "[Image #1]".to_string(),
local_image_path: PathBuf::from("/tmp/image-1.png"),
}],
text: "Review [Image #1]".to_string(),
text_source: TurnPromptTextSource::UserPrompt,
};
let formatted_prompt = SessionManager::formatted_prompt_output(&prompt, false);
assert_eq!(formatted_prompt, " › Review [Image #1]\n\n");
}
#[test]
fn test_renumbered_prompt_text_rewrites_only_attachment_occurrences() {
let prompt = TurnPrompt {
attachments: vec![TurnPromptAttachment {
placeholder: "[Image #1]".to_string(),
local_image_path: PathBuf::from("/tmp/image-1.png"),
}],
text: "Attach [Image #1] but keep literal [Image #1] text".to_string(),
text_source: TurnPromptTextSource::UserPrompt,
};
let renumbered_prompt = SessionManager::renumbered_prompt_text(&prompt, 2);
assert_eq!(
renumbered_prompt,
"Attach [Image #2] but keep literal [Image #1] text"
);
}
#[tokio::test]
async fn test_cleanup_prompt_attachment_paths_removes_files_and_directory() {
let temp_dir = tempfile::tempdir().expect("temp dir should exist");
let managed_tmp_root = temp_dir.path().join("tmp");
let image_directory = managed_tmp_root.join("session-id").join("images");
std::fs::create_dir_all(&image_directory).expect("image directory should exist");
let first_image = image_directory.join("image-1.png");
let second_image = image_directory.join("image-2.png");
std::fs::write(&first_image, b"png").expect("first image should exist");
std::fs::write(&second_image, b"png").expect("second image should exist");
SessionManager::cleanup_prompt_attachment_paths_in_root(
Arc::new(fs::RealFsClient),
&managed_tmp_root,
vec![first_image.clone(), second_image.clone()],
)
.await;
assert!(!first_image.exists());
assert!(!second_image.exists());
assert!(!image_directory.exists());
}
#[tokio::test]
async fn test_cleanup_prompt_attachment_paths_leaves_unmanaged_files_untouched() {
let temp_dir = tempfile::tempdir().expect("temp dir should exist");
let managed_tmp_root = temp_dir.path().join("tmp");
let image_directory = temp_dir.path().join("user-images");
std::fs::create_dir_all(&image_directory).expect("image directory should exist");
let image_path = image_directory.join("image-1.png");
std::fs::write(&image_path, b"png").expect("image file should exist");
SessionManager::cleanup_prompt_attachment_paths_in_root(
Arc::new(fs::RealFsClient),
&managed_tmp_root,
vec![image_path.clone()],
)
.await;
assert!(image_path.exists());
assert!(image_directory.exists());
}
#[test]
fn test_prepare_reply_context_first_message_sets_title_from_prompt() {
let prompt = "Implement optimistic retry path";
let turn_prompt = TurnPrompt::from_text(prompt.to_string());
let session = test_session("", Status::New, None, "");
let mut session_manager = session_manager_with_one_session(session);
let context = session_manager
.prepare_reply_context(0, &turn_prompt, AgentModel::ClaudeSonnet46, false)
.expect("reply context should be available")
.expect("session should produce reply context");
assert_eq!(context.0, None);
assert!(context.1);
assert_eq!(context.2, "session-id");
assert_eq!(context.3, Some(prompt.to_string()));
assert_eq!(session_manager.sessions[0].prompt, prompt);
assert_eq!(session_manager.sessions[0].title, Some(prompt.to_string()));
}
#[test]
fn test_prepare_reply_context_follow_up_keeps_existing_title() {
let session = test_session(
"Initial prompt",
Status::Review,
Some("Initial prompt"),
"existing output",
);
let mut session_manager = session_manager_with_one_session(session);
let prompt = TurnPrompt::from_text("Follow-up prompt".to_string());
let context = session_manager
.prepare_reply_context(0, &prompt, AgentModel::ClaudeSonnet46, false)
.expect("reply context should be available")
.expect("session should produce reply context");
assert_eq!(context.0, None);
assert!(!context.1);
assert_eq!(context.2, "session-id");
assert_eq!(context.3, None);
assert_eq!(session_manager.sessions[0].prompt, "Initial prompt");
assert_eq!(
session_manager.sessions[0].title,
Some("Initial prompt".to_string())
);
}
#[test]
fn test_prepare_reply_context_returns_workflow_error_when_status_blocks_reply() {
let session = test_session("Initial prompt", Status::InProgress, Some("Title"), "");
let mut session_manager = session_manager_with_one_session(session);
let prompt = TurnPrompt::from_text("Another prompt".to_string());
let result =
session_manager.prepare_reply_context(0, &prompt, AgentModel::ClaudeSonnet46, false);
let error = result.expect_err("in-progress session should block reply");
assert!(
matches!(error, SessionError::Workflow(_)),
"expected SessionError::Workflow, got: {error:?}"
);
}
#[test]
fn test_session_title_generation_prompt_includes_request() {
let request_prompt = "Refactor session lifecycle updates";
let title_prompt = SessionManager::session_title_generation_prompt(request_prompt)
.expect("title generation prompt should render");
assert!(title_prompt.contains("Generate a concise, commit-style title"));
assert!(title_prompt.contains("Describe what the user wants to do"));
assert!(title_prompt.contains("Keep it high-level and intent-focused."));
assert!(title_prompt.contains("Do not include long file names"));
assert!(title_prompt.contains("Return only the title text."));
assert!(title_prompt.contains(request_prompt));
}
#[test]
fn test_parse_generated_session_title_accepts_plain_title() {
let response_content = "Refine session startup flow";
let parsed_title = SessionManager::parse_generated_session_title(response_content);
assert_eq!(
parsed_title,
Some("Refine session startup flow".to_string())
);
}
#[test]
fn test_parse_generated_session_title_accepts_protocol_answer_plain_text() {
let response_content =
r#"{"answer":"Polish Gemini title parsing","questions":[],"summary":null}"#;
let parsed_title = SessionManager::parse_generated_session_title(response_content);
assert_eq!(
parsed_title,
Some("Polish Gemini title parsing".to_string())
);
}
#[test]
fn test_parse_generated_session_title_uses_first_nonempty_line_for_multiline_response() {
let response_content = "Polish Gemini title parsing\nExtra detail that should be ignored";
let parsed_title = SessionManager::parse_generated_session_title(response_content);
assert_eq!(
parsed_title,
Some("Polish Gemini title parsing".to_string())
);
}
#[test]
fn test_parse_generated_session_title_returns_none_for_question_only_protocol_payload() {
let response_content = r#"{"answer":"","questions":[{"text":"Need confirmation?","options":[]}],"summary":null}"#;
let parsed_title = SessionManager::parse_generated_session_title(response_content);
assert_eq!(parsed_title, None);
}
#[test]
fn test_parse_generated_session_title_normalizes_title_prefix() {
let response_content = "Title: \"Polish merge queue behavior\"";
let parsed_title = SessionManager::parse_generated_session_title(response_content);
assert_eq!(
parsed_title,
Some("Polish merge queue behavior".to_string())
);
}
fn session_manager_with_sessions(sessions: Vec<Session>) -> SessionManager {
let mut handles = HashMap::new();
for session in &sessions {
handles.insert(
session.id.clone(),
SessionHandles::new(session.output.clone(), session.status),
);
}
let row_count = i64::try_from(sessions.len()).unwrap_or(0);
let state = SessionState::new(
handles,
sessions,
TableState::default(),
Arc::new(RealClock),
row_count,
0,
);
SessionManager::new(
SessionDefaults {
model: AgentModel::Gpt54,
},
Arc::new(git::MockGitClient::new()),
state,
Vec::new(),
)
}
fn session_with_id(id: &str, status: Status) -> Session {
let mut session = test_session("prompt", status, None, "");
session.id = id.to_string().into();
session
}
#[test]
fn next_starts_at_first_selectable_row_when_no_prior_selection() {
let mut session_manager = session_manager_with_sessions(vec![
session_with_id("session-active", Status::InProgress),
session_with_id("session-archive", Status::Done),
]);
session_manager.next();
assert_eq!(session_manager.state.table_state.selected(), Some(0));
}
#[test]
fn next_advances_selection_to_next_grouped_row() {
let mut session_manager = session_manager_with_sessions(vec![
session_with_id("session-active-1", Status::InProgress),
session_with_id("session-active-2", Status::Review),
]);
session_manager.state.table_state.select(Some(0));
session_manager.next();
assert_eq!(session_manager.state.table_state.selected(), Some(1));
}
#[test]
fn next_wraps_to_first_selectable_row_after_last_row() {
let mut session_manager = session_manager_with_sessions(vec![
session_with_id("session-active", Status::InProgress),
session_with_id("session-archive", Status::Done),
]);
session_manager.state.table_state.select(Some(1));
session_manager.next();
assert_eq!(session_manager.state.table_state.selected(), Some(0));
}
#[test]
fn next_is_no_op_when_no_sessions_present() {
let mut session_manager = session_manager_with_sessions(Vec::new());
session_manager.next();
assert_eq!(session_manager.state.table_state.selected(), None);
}
#[test]
fn previous_starts_at_first_selectable_row_when_no_prior_selection() {
let mut session_manager = session_manager_with_sessions(vec![
session_with_id("session-active", Status::InProgress),
session_with_id("session-archive", Status::Done),
]);
session_manager.previous();
assert_eq!(session_manager.state.table_state.selected(), Some(0));
}
#[test]
fn previous_moves_selection_back_one_grouped_row() {
let mut session_manager = session_manager_with_sessions(vec![
session_with_id("session-active-1", Status::InProgress),
session_with_id("session-active-2", Status::Review),
]);
session_manager.state.table_state.select(Some(1));
session_manager.previous();
assert_eq!(session_manager.state.table_state.selected(), Some(0));
}
#[test]
fn previous_wraps_to_last_selectable_row_when_at_first_row() {
let mut session_manager = session_manager_with_sessions(vec![
session_with_id("session-active", Status::InProgress),
session_with_id("session-archive", Status::Done),
]);
session_manager.state.table_state.select(Some(0));
session_manager.previous();
assert_eq!(session_manager.state.table_state.selected(), Some(1));
}
#[test]
fn previous_is_no_op_when_no_sessions_present() {
let mut session_manager = session_manager_with_sessions(Vec::new());
session_manager.previous();
assert_eq!(session_manager.state.table_state.selected(), None);
}
#[test]
fn selected_session_returns_currently_selected_session_or_none() {
let mut session_manager = session_manager_with_sessions(vec![
session_with_id("session-a", Status::InProgress),
session_with_id("session-b", Status::Review),
]);
assert!(session_manager.selected_session().is_none());
session_manager.state.table_state.select(Some(1));
assert_eq!(
session_manager
.selected_session()
.map(|session| session.id.clone()),
Some("session-b".into())
);
}
#[test]
fn session_at_returns_session_by_index_or_none_for_out_of_range() {
let session_manager = session_manager_with_sessions(vec![
session_with_id("session-a", Status::InProgress),
session_with_id("session-b", Status::Review),
]);
assert_eq!(
session_manager
.session_at(0)
.map(|session| session.id.as_str()),
Some("session-a")
);
assert_eq!(
session_manager
.session_at(1)
.map(|session| session.id.as_str()),
Some("session-b")
);
assert!(session_manager.session_at(99).is_none());
}
#[test]
fn session_id_for_index_returns_owned_id_or_none_for_out_of_range() {
let session_manager =
session_manager_with_sessions(vec![session_with_id("session-a", Status::InProgress)]);
assert_eq!(
session_manager.session_id_for_index(0),
Some("session-a".into())
);
assert!(session_manager.session_id_for_index(1).is_none());
}
#[tokio::test]
async fn set_session_model_persists_new_model_and_clears_conversation_state() {
let mut session = test_session("Prompt", Status::Review, Some("Title"), "");
session.model = AgentModel::ClaudeSonnet46;
let database = database_with_session(&session).await;
database
.update_session_provider_conversation_id("session-id", Some("provider-conv"))
.await
.expect("seed provider conversation id");
database
.update_session_instruction_conversation_id("session-id", Some("instruction-conv"))
.await
.expect("seed instruction conversation id");
let mut session_manager = session_manager_with_one_session(session);
let (services, mut event_rx) = test_services_with_event_receiver(
&database,
Arc::new(git::MockGitClient::new()),
Arc::new(forge::MockReviewRequestClient::new()),
);
session_manager
.set_session_model(&services, "session-id", AgentModel::Gpt54)
.await
.expect("set session model should succeed");
let persisted_model = database
.load_sessions()
.await
.expect("load sessions should succeed")
.into_iter()
.find(|row| row.id == "session-id")
.expect("session row should exist")
.model;
let cleared_provider = database
.get_session_provider_conversation_id("session-id")
.await
.expect("provider id load should succeed");
let cleared_instruction = database
.get_session_instruction_conversation_id("session-id")
.await
.expect("instruction id load should succeed");
let emitted_event = event_rx.try_recv().expect("model event expected");
assert_eq!(persisted_model, AgentModel::Gpt54.as_str());
assert!(cleared_provider.is_none());
assert!(cleared_instruction.is_none());
assert_eq!(
emitted_event,
AppEvent::SessionModelUpdated {
session_id: "session-id".into(),
session_model: AgentModel::Gpt54,
}
);
assert!(session_manager.should_replay_history("session-id"));
}
#[tokio::test]
async fn set_session_model_keeps_conversation_state_when_model_does_not_change() {
let mut session = test_session("Prompt", Status::InProgress, Some("Title"), "");
session.model = AgentModel::ClaudeSonnet46;
let database = database_with_session(&session).await;
database
.update_session_provider_conversation_id("session-id", Some("provider-conv"))
.await
.expect("seed provider conversation id");
let mut session_manager = session_manager_with_one_session(session);
let (services, mut event_rx) = test_services_with_event_receiver(
&database,
Arc::new(git::MockGitClient::new()),
Arc::new(forge::MockReviewRequestClient::new()),
);
session_manager
.set_session_model(&services, "session-id", AgentModel::ClaudeSonnet46)
.await
.expect("set session model should succeed");
let preserved_provider = database
.get_session_provider_conversation_id("session-id")
.await
.expect("provider id load should succeed");
let emitted_event = event_rx.try_recv().expect("model event expected");
assert_eq!(preserved_provider.as_deref(), Some("provider-conv"));
assert_eq!(
emitted_event,
AppEvent::SessionModelUpdated {
session_id: "session-id".into(),
session_model: AgentModel::ClaudeSonnet46,
}
);
assert!(!session_manager.should_replay_history("session-id"));
}
#[tokio::test]
async fn set_session_model_returns_error_for_missing_session() {
let session = test_session("Prompt", Status::Review, Some("Title"), "");
let database = database_with_session(&session).await;
let mut session_manager = session_manager_with_one_session(session);
let services = test_services(
&database,
Arc::new(git::MockGitClient::new()),
Arc::new(forge::MockReviewRequestClient::new()),
);
let result = session_manager
.set_session_model(&services, "missing", AgentModel::Gpt54)
.await;
assert!(
result.is_err(),
"missing session should return SessionError"
);
}
#[test]
fn session_index_for_id_returns_index_or_none_for_unknown_session() {
let session_manager = session_manager_with_sessions(vec![
session_with_id("session-a", Status::InProgress),
session_with_id("session-b", Status::Review),
]);
assert_eq!(session_manager.session_index_for_id("session-a"), Some(0));
assert_eq!(session_manager.session_index_for_id("session-b"), Some(1));
assert!(session_manager.session_index_for_id("missing").is_none());
}
}