use std::collections::hash_map::Entry;
use std::collections::{HashMap, HashSet};
use std::path::PathBuf;
use app::branch_publish::{
BranchPublishActionUpdate, BranchPublishTaskResult, BranchPublishTaskSuccess,
branch_publish_loading_label as branch_publish_loading_label_text,
branch_publish_loading_message as branch_publish_loading_message_text,
branch_publish_loading_title as branch_publish_loading_title_text,
branch_publish_success_title as branch_publish_success_title_text,
detected_forge_kind_from_git_push_error, git_push_authentication_message,
is_git_push_authentication_error,
pull_request_publish_success_message as pull_request_publish_success_message_text,
};
use app::reducer::AppEventReducer;
use app::review::{
FocusedReviewPersistence, ReviewUpdate, apply_review_updates, auto_start_reviews,
};
use super::state::{App, SyncPopupContext, SyncReviewRequestTaskResult, UpdateStatus};
use crate::app;
use crate::app::session::{
SessionTaskService, SyncMainOutcome, SyncSessionStartError, TurnAppliedState,
};
use crate::app::session_state::SessionGitStatus;
use crate::domain::file_entry::FileEntry;
use crate::domain::input::InputState;
use crate::domain::session::{
PublishBranchAction, PublishedBranchSyncStatus, SessionId, SessionSize, Status,
};
use crate::domain::transcript_notice::TranscriptNotice;
use crate::runtime::mode::{at_mention, question, sync_blocked};
use crate::ui::state::app_mode::{AppMode, ConfirmationViewMode, QuestionFocus};
use crate::ui::state::prompt::PromptAtMentionState;
#[derive(Clone, Debug, Eq, PartialEq)]
pub(crate) enum AppEvent {
AtMentionEntriesLoaded {
entries: Vec<FileEntry>,
session_id: SessionId,
},
GitStatusUpdated {
session_statuses: HashMap<SessionId, SessionGitStatus>,
status: Option<(u32, u32)>,
},
VersionAvailabilityUpdated {
latest_available_version: Option<String>,
},
UpdateStatusChanged { update_status: UpdateStatus },
SessionModelUpdated {
session_id: SessionId,
session_model: crate::domain::agent::AgentModel,
},
SessionReasoningLevelUpdated {
reasoning_level_override: Option<crate::domain::agent::ReasoningLevel>,
session_id: SessionId,
},
RefreshSessions,
RefreshGitStatus,
RequestedReviewsLoaded {
generation: u64,
project_id: i64,
result: Result<Vec<ag_forge::RequestedReview>, String>,
},
SessionProgressUpdated {
progress_message: Option<String>,
session_id: SessionId,
},
SyncMainCompleted {
result: Result<SyncMainOutcome, SyncSessionStartError>,
},
SessionSizeUpdated {
added_lines: u64,
deleted_lines: u64,
session_id: SessionId,
session_size: SessionSize,
},
SessionTitleGenerationFinished {
generation: u64,
session_id: SessionId,
},
BranchPublishActionCompleted {
restore_view: ConfirmationViewMode,
result: Box<BranchPublishTaskResult>,
session_id: SessionId,
},
ReviewPrepared {
diff_hash: u64,
review_text: String,
session_id: SessionId,
},
ReviewPreparationFailed {
diff_hash: u64,
error: String,
session_id: SessionId,
},
SessionUpdated { session_id: SessionId, version: u64 },
AgentResponseReceived {
session_id: SessionId,
turn_applied_state: TurnAppliedState,
},
SessionWorkflowNoticeUpdated {
notice: String,
session_id: SessionId,
},
PublishedBranchSyncUpdated {
session_id: SessionId,
sync_operation_id: String,
sync_status: PublishedBranchSyncStatus,
},
ReviewRequestStatusUpdated {
result: Result<SyncReviewRequestTaskResult, String>,
session_id: SessionId,
},
ReviewCommentsUpdated { session_id: SessionId },
}
#[derive(Default)]
pub(super) struct AppEventBatch {
pub(super) applied_turns: HashMap<SessionId, TurnAppliedState>,
pub(super) at_mention_entries_updates: HashMap<SessionId, Vec<FileEntry>>,
pub(super) branch_publish_action_update: Option<BranchPublishActionUpdate>,
pub(super) git_status_update: Option<GitStatusBatchUpdate>,
pub(super) latest_available_version_update: Option<LatestAvailableVersionUpdate>,
pub(super) published_branch_sync_updates: Vec<(SessionId, PublishedBranchSyncUpdate)>,
pub(super) review_updates: HashMap<SessionId, ReviewUpdate>,
pub(super) session_git_status_updates: HashMap<SessionId, SessionGitStatus>,
pub(super) session_ids: HashSet<SessionId>,
pub(super) session_update_versions: HashMap<SessionId, u64>,
pub(super) session_model_updates: HashMap<SessionId, crate::domain::agent::AgentModel>,
pub(super) session_reasoning_level_updates:
HashMap<SessionId, Option<crate::domain::agent::ReasoningLevel>>,
pub(super) session_progress_updates: HashMap<SessionId, Option<String>>,
pub(super) session_size_updates: HashMap<SessionId, (u64, u64, SessionSize)>,
pub(super) session_title_generation_finished: HashMap<SessionId, u64>,
pub(super) session_workflow_notice_updates: HashMap<SessionId, Vec<String>>,
pub(super) should_refresh_git_status: bool,
pub(super) should_force_reload: bool,
pub(super) review_request_status_updates: Vec<ReviewRequestStatusUpdate>,
pub(super) review_comment_session_ids: HashSet<SessionId>,
pub(super) requested_reviews:
Option<(u64, i64, Result<Vec<ag_forge::RequestedReview>, String>)>,
pub(super) sync_main_result: Option<Result<SyncMainOutcome, SyncSessionStartError>>,
pub(super) update_status: Option<UpdateStatus>,
}
pub(super) struct GitStatusBatchUpdate {
status: Option<(u32, u32)>,
}
pub(super) struct LatestAvailableVersionUpdate {
latest_available_version: Option<String>,
}
pub(super) struct PublishedBranchSyncUpdate {
sync_operation_id: String,
sync_status: PublishedBranchSyncStatus,
}
pub(super) struct ReviewRequestStatusUpdate {
pub(super) result: Result<SyncReviewRequestTaskResult, String>,
pub(super) session_id: SessionId,
}
impl AppEventBatch {
pub(super) fn collect_event(&mut self, event: AppEvent) {
match event {
AppEvent::AtMentionEntriesLoaded {
entries,
session_id,
} => self.collect_at_mention_entries_loaded(session_id, entries),
AppEvent::GitStatusUpdated {
session_statuses,
status,
} => self.collect_git_status_updated(session_statuses, status),
AppEvent::VersionAvailabilityUpdated {
latest_available_version,
} => self.collect_version_availability_updated(latest_available_version),
AppEvent::UpdateStatusChanged { update_status } => {
self.collect_update_status_changed(update_status);
}
AppEvent::SessionModelUpdated {
session_id,
session_model,
} => self.collect_session_model_updated(session_id, session_model),
AppEvent::SessionReasoningLevelUpdated {
reasoning_level_override,
session_id,
} => self.collect_session_reasoning_level_updated(session_id, reasoning_level_override),
AppEvent::RefreshSessions => self.collect_refresh_sessions(),
AppEvent::RefreshGitStatus => self.collect_refresh_git_status(),
AppEvent::RequestedReviewsLoaded {
generation,
project_id,
result,
} => self.collect_requested_reviews_loaded(generation, project_id, result),
AppEvent::SessionProgressUpdated {
progress_message,
session_id,
} => {
self.session_progress_updates
.insert(session_id, progress_message);
}
AppEvent::SyncMainCompleted { result } => self.collect_sync_main_completed(result),
AppEvent::SessionSizeUpdated {
added_lines,
deleted_lines,
session_id,
session_size,
} => {
self.session_size_updates
.insert(session_id, (added_lines, deleted_lines, session_size));
}
AppEvent::SessionTitleGenerationFinished {
generation,
session_id,
} => {
self.session_title_generation_finished
.insert(session_id, generation);
}
AppEvent::BranchPublishActionCompleted {
restore_view,
result,
session_id,
} => self.collect_branch_publish_action_completed(restore_view, *result, session_id),
AppEvent::ReviewPrepared {
diff_hash,
review_text,
session_id,
} => self.collect_review_prepared(diff_hash, review_text, session_id),
AppEvent::ReviewPreparationFailed {
diff_hash,
error,
session_id,
} => self.collect_review_preparation_failed(diff_hash, error, session_id),
AppEvent::SessionUpdated {
session_id,
version,
} => self.collect_session_updated(session_id, version),
AppEvent::AgentResponseReceived {
session_id,
turn_applied_state,
} => self.collect_agent_response_received(session_id, turn_applied_state),
AppEvent::SessionWorkflowNoticeUpdated { notice, session_id } => {
self.collect_session_workflow_notice_updated(session_id, notice);
}
AppEvent::PublishedBranchSyncUpdated {
session_id,
sync_operation_id,
sync_status,
} => self.collect_published_branch_sync_updated(
session_id,
sync_operation_id,
sync_status,
),
AppEvent::ReviewRequestStatusUpdated { result, session_id } => {
self.collect_review_request_status_updated(result, session_id);
}
AppEvent::ReviewCommentsUpdated { session_id } => {
self.collect_review_comments_updated(session_id);
}
}
}
fn collect_requested_reviews_loaded(
&mut self,
generation: u64,
project_id: i64,
result: Result<Vec<ag_forge::RequestedReview>, String>,
) {
if self
.requested_reviews
.as_ref()
.is_none_or(|(batched_generation, _, _)| generation >= *batched_generation)
{
self.requested_reviews = Some((generation, project_id, result));
}
}
fn collect_session_model_updated(
&mut self,
session_id: SessionId,
session_model: crate::domain::agent::AgentModel,
) {
self.session_model_updates.insert(session_id, session_model);
}
fn collect_session_reasoning_level_updated(
&mut self,
session_id: SessionId,
reasoning_level_override: Option<crate::domain::agent::ReasoningLevel>,
) {
self.session_reasoning_level_updates
.insert(session_id, reasoning_level_override);
}
fn collect_session_workflow_notice_updated(&mut self, session_id: SessionId, notice: String) {
self.session_ids.insert(session_id.clone());
self.session_workflow_notice_updates
.entry(session_id)
.or_default()
.push(notice);
}
fn collect_review_comments_updated(&mut self, session_id: SessionId) {
self.review_comment_session_ids.insert(session_id);
}
fn collect_at_mention_entries_loaded(
&mut self,
session_id: SessionId,
entries: Vec<FileEntry>,
) {
self.at_mention_entries_updates.insert(session_id, entries);
}
fn collect_update_status_changed(&mut self, update_status: UpdateStatus) {
self.update_status = Some(update_status);
}
fn collect_refresh_sessions(&mut self) {
self.should_force_reload = true;
}
fn collect_refresh_git_status(&mut self) {
self.should_refresh_git_status = true;
}
fn collect_git_status_updated(
&mut self,
session_statuses: HashMap<SessionId, SessionGitStatus>,
status: Option<(u32, u32)>,
) {
self.git_status_update = Some(GitStatusBatchUpdate { status });
self.session_git_status_updates = session_statuses;
}
fn collect_version_availability_updated(&mut self, latest_available_version: Option<String>) {
self.latest_available_version_update = Some(LatestAvailableVersionUpdate {
latest_available_version,
});
}
fn collect_sync_main_completed(
&mut self,
result: Result<SyncMainOutcome, SyncSessionStartError>,
) {
if result.is_ok() {
self.should_refresh_git_status = true;
}
self.sync_main_result = Some(result);
}
fn collect_branch_publish_action_completed(
&mut self,
restore_view: ConfirmationViewMode,
result: BranchPublishTaskResult,
session_id: SessionId,
) {
if result.is_ok() {
self.should_refresh_git_status = true;
}
self.branch_publish_action_update = Some(BranchPublishActionUpdate {
restore_view,
result,
session_id,
});
}
fn collect_review_prepared(
&mut self,
diff_hash: u64,
review_text: String,
session_id: SessionId,
) {
self.review_updates.insert(
session_id,
ReviewUpdate {
diff_hash,
result: Ok(review_text),
},
);
}
fn collect_review_preparation_failed(
&mut self,
diff_hash: u64,
error: String,
session_id: SessionId,
) {
self.review_updates.insert(
session_id,
ReviewUpdate {
diff_hash,
result: Err(error),
},
);
}
fn collect_published_branch_sync_updated(
&mut self,
session_id: SessionId,
sync_operation_id: String,
sync_status: PublishedBranchSyncStatus,
) {
if matches!(
sync_status,
PublishedBranchSyncStatus::Idle | PublishedBranchSyncStatus::Succeeded
) {
self.should_refresh_git_status = true;
}
self.published_branch_sync_updates.push((
session_id,
PublishedBranchSyncUpdate {
sync_operation_id,
sync_status,
},
));
}
fn collect_review_request_status_updated(
&mut self,
result: Result<SyncReviewRequestTaskResult, String>,
session_id: SessionId,
) {
self.review_request_status_updates
.push(ReviewRequestStatusUpdate { result, session_id });
}
fn collect_session_updated(&mut self, session_id: SessionId, version: u64) {
self.session_ids.insert(session_id.clone());
self.session_update_versions.insert(session_id, version);
}
fn collect_agent_response_received(
&mut self,
session_id: SessionId,
turn_applied_state: TurnAppliedState,
) {
self.session_ids.insert(session_id.clone());
match self.applied_turns.entry(session_id) {
Entry::Occupied(mut occupied_entry) => {
occupied_entry.get_mut().merge_newer(turn_applied_state);
}
Entry::Vacant(vacant_entry) => {
vacant_entry.insert(turn_applied_state);
}
}
}
}
impl App {
pub(crate) async fn apply_app_events(&mut self, first_event: AppEvent) {
let drained_events = AppEventReducer::drain(&mut self.event_rx, first_event);
let mut event_batch = AppEventBatch::default();
for event in drained_events {
event_batch.collect_event(event);
}
self.apply_app_event_batch(event_batch).await;
}
pub(crate) async fn process_pending_app_events(&mut self) {
let Ok(first_event) = self.event_rx.try_recv() else {
return;
};
self.apply_app_events(first_event).await;
}
pub(crate) async fn next_app_event(&mut self) -> Option<AppEvent> {
self.event_rx.recv().await
}
async fn apply_app_event_batch(&mut self, mut event_batch: AppEventBatch) {
let mut should_mark_dirty = Self::app_event_batch_changes_observable_state(&event_batch);
let previous_session_states = self.previous_session_states(&event_batch.session_ids);
should_mark_dirty |=
self.update_session_redraw_versions(&event_batch.session_update_versions);
self.apply_batch_runtime_updates(&mut event_batch).await;
self.apply_batch_session_snapshot_updates(&mut event_batch);
let focused_review_persistence = apply_review_updates(
&mut self.review_cache,
&mut self.mode,
self.sessions.state_mut(),
event_batch.review_updates,
);
self.persist_focused_review_updates(focused_review_persistence)
.await;
if let Some(branch_publish_action_update) = event_batch.branch_publish_action_update {
self.apply_branch_publish_action_update(branch_publish_action_update);
}
for review_request_status_update in event_batch.review_request_status_updates {
self.apply_review_request_status_update(review_request_status_update)
.await;
}
self.invalidate_diff_scroll_cache_for_review_comments(
&event_batch.review_comment_session_ids,
);
self.apply_session_progress_updates(std::mem::take(
&mut event_batch.session_progress_updates,
));
for (session_id, turn_applied_state) in event_batch.applied_turns {
self.apply_agent_response_received(&session_id, &turn_applied_state);
}
for (session_id, notices) in
std::mem::take(&mut event_batch.session_workflow_notice_updates)
{
for notice in notices {
self.sessions.append_workflow_notice(&session_id, notice);
}
}
for (session_id, sync_update) in event_batch.published_branch_sync_updates {
self.apply_published_branch_sync_update(&session_id, sync_update);
}
self.sync_touched_sessions(&event_batch.session_ids);
auto_start_reviews(
&mut self.review_cache,
&event_batch.session_ids,
self.sessions.state_mut(),
&mut self.mode,
self.services.git_client(),
self.services.event_sender(),
self.settings.default_review_model,
)
.await;
if let Some(sync_main_result) = event_batch.sync_main_result {
let sync_popup_context = self.sync_popup_context();
self.mode = Self::sync_main_popup_mode(sync_main_result, &sync_popup_context);
}
self.handle_merge_queue_progress(&event_batch.session_ids, &previous_session_states)
.await;
self.retain_valid_session_progress_messages();
self.sessions.retain_active_prompt_outputs();
if should_mark_dirty {
self.mark_dirty();
}
}
fn invalidate_diff_scroll_cache_for_review_comments(
&mut self,
review_comment_session_ids: &HashSet<SessionId>,
) {
let AppMode::Diff {
scroll_cache,
session_id,
..
} = &mut self.mode
else {
return;
};
if review_comment_session_ids.contains(session_id) {
*scroll_cache = None;
}
}
async fn apply_batch_runtime_updates(&mut self, event_batch: &mut AppEventBatch) {
if event_batch.should_force_reload {
self.refresh_sessions_now().await;
self.reload_projects().await;
}
if event_batch.should_refresh_git_status {
self.restart_git_status_task();
}
if let Some(git_status_update) = &event_batch.git_status_update {
self.projects.set_git_status(git_status_update.status);
self.sessions
.replace_session_git_statuses(event_batch.session_git_status_updates.clone());
}
if let Some((generation, project_id, result)) = event_batch.requested_reviews.take()
&& project_id == self.projects.active_project_id()
&& self
.requested_reviews
.matches_loading_request(project_id, generation)
{
self.requested_reviews = match result {
Ok(items) => app::RequestedReviewState::Loaded { items, project_id },
Err(message) => app::RequestedReviewState::Failed {
message,
project_id,
},
};
}
self.apply_status_bar_updates(
event_batch.latest_available_version_update.as_ref(),
event_batch.update_status.take(),
);
}
fn sync_touched_sessions(&mut self, session_ids: &HashSet<SessionId>) {
for session_id in session_ids {
self.sessions.sync_session_from_handle(session_id);
}
self.sessions.clear_terminal_session_workers(session_ids);
}
fn apply_status_bar_updates(
&mut self,
latest_available_version_update: Option<&LatestAvailableVersionUpdate>,
update_status: Option<UpdateStatus>,
) {
if let Some(latest_available_version_update) = latest_available_version_update {
self.latest_available_version
.clone_from(&latest_available_version_update.latest_available_version);
}
if let Some(update_status) = update_status {
self.update_status = Some(update_status);
}
}
fn previous_session_states(
&self,
session_ids: &HashSet<SessionId>,
) -> HashMap<SessionId, Status> {
session_ids
.iter()
.filter_map(|session_id| {
self.sessions
.sessions()
.iter()
.find(|session| session.id == *session_id)
.map(|session| (session_id.clone(), session.status))
})
.collect()
}
fn app_event_batch_changes_observable_state(event_batch: &AppEventBatch) -> bool {
event_batch.should_force_reload
|| event_batch.git_status_update.is_some()
|| event_batch.latest_available_version_update.is_some()
|| event_batch.update_status.is_some()
|| !event_batch.applied_turns.is_empty()
|| !event_batch.at_mention_entries_updates.is_empty()
|| event_batch.branch_publish_action_update.is_some()
|| !event_batch.published_branch_sync_updates.is_empty()
|| !event_batch.review_request_status_updates.is_empty()
|| !event_batch.review_comment_session_ids.is_empty()
|| event_batch.requested_reviews.is_some()
|| !event_batch.review_updates.is_empty()
|| !event_batch.session_model_updates.is_empty()
|| !event_batch.session_progress_updates.is_empty()
|| !event_batch.session_reasoning_level_updates.is_empty()
|| !event_batch.session_size_updates.is_empty()
|| !event_batch.session_title_generation_finished.is_empty()
|| !event_batch.session_workflow_notice_updates.is_empty()
|| event_batch.sync_main_result.is_some()
}
fn update_session_redraw_versions(
&mut self,
session_update_versions: &HashMap<SessionId, u64>,
) -> bool {
let mut did_change = false;
for (session_id, version) in session_update_versions {
let previous_version = self
.last_seen_session_update_versions
.insert(session_id.clone(), *version);
if previous_version != Some(*version) {
did_change = true;
}
}
did_change
}
fn apply_batch_session_snapshot_updates(&mut self, event_batch: &mut AppEventBatch) {
for (session_id, session_model) in std::mem::take(&mut event_batch.session_model_updates) {
self.sessions
.apply_session_model_updated(&session_id, session_model);
}
for (session_id, reasoning_level_override) in
std::mem::take(&mut event_batch.session_reasoning_level_updates)
{
self.sessions
.apply_session_reasoning_level_updated(&session_id, reasoning_level_override);
}
for (session_id, (added_lines, deleted_lines, session_size)) in
std::mem::take(&mut event_batch.session_size_updates)
{
self.sessions.apply_session_size_updated(
&session_id,
added_lines,
deleted_lines,
session_size,
);
}
for (session_id, generation) in
std::mem::take(&mut event_batch.session_title_generation_finished)
{
self.sessions
.clear_title_generation_task_if_matches(&session_id, generation);
}
for (session_id, entries) in std::mem::take(&mut event_batch.at_mention_entries_updates) {
self.sessions.set_at_mention_index_for_root(
self.at_mention_lookup_root(&session_id),
entries.clone(),
);
self.apply_prompt_at_mention_entries(&session_id, entries);
}
}
fn apply_session_progress_updates(
&mut self,
session_progress_updates: HashMap<SessionId, Option<String>>,
) {
for (session_id, progress_message) in session_progress_updates {
if let Some(progress_message) = progress_message {
self.session_progress_messages
.insert(session_id, progress_message);
} else {
self.session_progress_messages.remove(&session_id);
}
}
}
fn apply_agent_response_received(
&mut self,
session_id: &str,
turn_applied_state: &TurnAppliedState,
) {
if !self
.sessions
.sessions()
.iter()
.any(|session| session.id == session_id)
{
return;
}
self.sessions
.apply_turn_applied_state(session_id, turn_applied_state);
let questions = turn_applied_state.questions.clone();
if questions.is_empty() {
return;
}
if self.is_viewing_session(session_id) {
let (review_status_message, review_text) = self.question_mode_review_state(session_id);
self.mode = AppMode::Question {
at_mention_state: None,
selected_option_index: question::default_option_index(&questions, 0),
session_id: session_id.into(),
questions,
review_status_message,
review_text,
responses: Vec::new(),
current_index: 0,
focus: QuestionFocus::Answer,
input: InputState::default(),
scroll_offset: None,
};
}
}
fn is_viewing_session(&self, session_id: &str) -> bool {
match &self.mode {
AppMode::View {
session_id: view_id,
..
}
| AppMode::Prompt {
session_id: view_id,
..
}
| AppMode::Diff {
session_id: view_id,
..
}
| AppMode::Question {
session_id: view_id,
..
}
| AppMode::OpenCommandSelector {
restore_view:
ConfirmationViewMode {
session_id: view_id,
..
},
..
}
| AppMode::PublishBranchInput {
restore_view:
ConfirmationViewMode {
session_id: view_id,
..
},
..
}
| AppMode::ViewInfoPopup {
restore_view:
ConfirmationViewMode {
session_id: view_id,
..
},
..
} => view_id == session_id,
AppMode::List
| AppMode::SessionCreation { .. }
| AppMode::Confirmation { .. }
| AppMode::SyncBlockedPopup { .. }
| AppMode::Help { .. } => false,
}
}
fn question_mode_review_state(&self, session_id: &str) -> (Option<String>, Option<String>) {
match &self.mode {
AppMode::View {
review_status_message,
review_text,
..
}
| AppMode::Prompt {
review_status_message,
review_text,
..
}
| AppMode::Question {
review_status_message,
review_text,
..
} => (review_status_message.clone(), review_text.clone()),
AppMode::OpenCommandSelector { restore_view, .. }
| AppMode::PublishBranchInput { restore_view, .. }
| AppMode::ViewInfoPopup { restore_view, .. } => (
restore_view.review_status_message.clone(),
restore_view.review_text.clone(),
),
AppMode::Diff {
session_id: diff_session_id,
..
} if diff_session_id == session_id => self.review_view_state(session_id),
AppMode::List
| AppMode::SessionCreation { .. }
| AppMode::Confirmation { .. }
| AppMode::SyncBlockedPopup { .. }
| AppMode::Diff { .. }
| AppMode::Help { .. } => (None, None),
}
}
fn apply_published_branch_sync_update(
&mut self,
session_id: &str,
sync_update: PublishedBranchSyncUpdate,
) {
let PublishedBranchSyncUpdate {
sync_operation_id,
sync_status,
} = sync_update;
match sync_status {
PublishedBranchSyncStatus::InProgress => {
self.sessions
.start_published_branch_sync(session_id, sync_operation_id);
}
PublishedBranchSyncStatus::Idle
| PublishedBranchSyncStatus::Succeeded
| PublishedBranchSyncStatus::Failed => {
self.sessions.finish_published_branch_sync(
session_id,
&sync_operation_id,
sync_status,
);
}
}
}
fn at_mention_lookup_root(&self, session_id: &str) -> PathBuf {
let project_working_dir = self.working_dir().to_path_buf();
self.sessions.session_for_id(session_id).map_or_else(
|| project_working_dir.clone(),
|session| {
let project_working_dir = project_working_dir.clone();
let session_folder = session.folder.clone();
let has_session_folder = self.services.fs_client().is_dir(session_folder.clone());
at_mention::lookup_root(
project_working_dir,
Some(session_folder),
has_session_folder,
)
},
)
}
fn apply_prompt_at_mention_entries(&mut self, session_id: &str, entries: Vec<FileEntry>) {
let (at_mention_state, has_query) = match &mut self.mode {
AppMode::Prompt {
at_mention_state,
input,
session_id: mode_session_id,
..
} if mode_session_id == session_id => {
(at_mention_state, input.at_mention_query().is_some())
}
AppMode::Question {
at_mention_state,
input,
session_id: mode_session_id,
..
} if mode_session_id == session_id => {
(at_mention_state, input.at_mention_query().is_some())
}
_ => return,
};
if !has_query {
return;
}
if let Some(state) = at_mention_state.as_mut() {
state.all_entries = entries;
state.selected_index = 0;
return;
}
*at_mention_state = Some(PromptAtMentionState::new(entries));
}
#[cfg(test)]
pub(super) fn apply_review_update(
&mut self,
session_id: &str,
review_update: app::review::ReviewUpdate,
) {
let mut review_updates = HashMap::new();
review_updates.insert(SessionId::from(session_id), review_update);
apply_review_updates(
&mut self.review_cache,
&mut self.mode,
self.sessions.state_mut(),
review_updates,
);
}
async fn persist_focused_review_updates(
&self,
focused_review_persistence: Vec<FocusedReviewPersistence>,
) {
for persistence_update in focused_review_persistence {
let diff_hash = persistence_update
.diff_hash
.map(|diff_hash| diff_hash.to_string());
let _ = self
.services
.db()
.sessions()
.update_session_focused_review(
persistence_update.session_id.as_str(),
diff_hash,
persistence_update.text,
)
.await;
}
}
#[cfg(test)]
pub(super) async fn auto_start_reviews(&mut self, session_ids: &HashSet<SessionId>) {
auto_start_reviews(
&mut self.review_cache,
session_ids,
self.sessions.state_mut(),
&mut self.mode,
self.services.git_client(),
self.services.event_sender(),
self.settings.default_review_model,
)
.await;
}
pub(super) fn apply_branch_publish_action_update(
&mut self,
branch_publish_action_update: BranchPublishActionUpdate,
) {
let BranchPublishActionUpdate {
restore_view,
result,
session_id,
} = branch_publish_action_update;
let popup_mode = match result {
Ok(BranchPublishTaskSuccess::Pushed {
branch_name,
review_request_creation,
upstream_reference,
}) => {
self.sessions
.apply_published_upstream_ref(&session_id, upstream_reference);
Self::view_info_popup_mode(
Self::branch_publish_success_title(PublishBranchAction::Push),
Self::branch_publish_success_message(
&branch_name,
review_request_creation.as_ref(),
),
false,
String::new(),
restore_view,
)
}
Ok(BranchPublishTaskSuccess::PullRequestPublished {
branch_name,
review_request,
upstream_reference,
}) => {
self.sessions
.apply_published_upstream_ref(&session_id, upstream_reference);
self.sessions
.apply_review_request(&session_id, review_request.clone());
Self::view_info_popup_mode(
Self::review_request_publish_success_title(&review_request),
Self::pull_request_publish_success_message(&branch_name, &review_request),
false,
String::new(),
restore_view,
)
}
Err(failure) => Self::view_info_popup_mode(
failure.title,
failure.message,
false,
String::new(),
restore_view,
),
};
self.mode = popup_mode;
}
pub(super) async fn apply_review_request_status_update(
&mut self,
review_request_status_update: ReviewRequestStatusUpdate,
) {
let ReviewRequestStatusUpdate { result, session_id } = review_request_status_update;
let Ok(task_result) = result else {
return;
};
if let Some(summary) = task_result.summary {
let _ = self
.sessions
.store_review_request_summary(&self.services, &session_id, summary)
.await;
}
match task_result.outcome {
crate::app::session::SyncReviewRequestOutcome::Merged {
session_head_hash, ..
} => {
if let Some(warning) = self
.complete_externally_merged_session(&session_id, session_head_hash)
.await
{
self.append_output_for_session(
&session_id,
&TranscriptNotice::ReviewRequestSyncWarning.format(warning),
)
.await;
}
}
crate::app::session::SyncReviewRequestOutcome::Closed { .. } => {
self.cancel_externally_closed_session(&session_id).await;
}
crate::app::session::SyncReviewRequestOutcome::Open { .. }
| crate::app::session::SyncReviewRequestOutcome::NoReviewRequest => {}
}
}
async fn complete_externally_merged_session(
&self,
session_id: &str,
session_head_hash: Option<String>,
) -> Option<String> {
let Ok(session) = self.sessions.session_or_err(session_id) else {
return None;
};
let Ok(handles) = self.sessions.session_handles_or_err(session_id) else {
return None;
};
let mut warnings = Vec::new();
let commit_hash_persistence_error = if let Some(session_head_hash) = session_head_hash {
self.services
.db()
.sessions()
.update_session_merged_commit_hash(session_id, Some(session_head_hash))
.await
.err()
} else {
None
};
if let Some(error) = commit_hash_persistence_error {
warnings.push(format!("Merged commit hash persistence failed: {error}"));
}
let folder = session.folder.clone();
let base_branch = session.base_branch.clone();
let source_branch = crate::app::session::session_branch(session_id);
if let Err(error) = crate::app::session::SessionManager::cleanup_merged_session_worktree(
folder,
self.services.fs_client(),
self.services.git_client(),
source_branch,
None,
)
.await
{
warnings.push(format!("Worktree cleanup failed: {error}"));
}
if let Err(error) =
crate::app::session::SessionManager::restack_child_sessions_after_parent_merge(
self.services.db(),
session_id,
&base_branch,
)
.await
{
warnings.push(format!("Stacked child restack failed: {error}"));
}
let app_event_tx = self.services.event_sender();
SessionTaskService::update_status(
handles.status.as_ref(),
self.services.clock().as_ref(),
self.services.db(),
&app_event_tx,
&self.services.session_update_versions(),
session_id,
Status::Done,
)
.await;
(!warnings.is_empty()).then(|| warnings.join("\n"))
}
async fn cancel_externally_closed_session(&self, session_id: &str) {
let Ok(handles) = self.sessions.session_handles_or_err(session_id) else {
return;
};
let app_event_tx = self.services.event_sender();
let _ = SessionTaskService::update_status(
handles.status.as_ref(),
self.services.clock().as_ref(),
self.services.db(),
&app_event_tx,
&self.services.session_update_versions(),
session_id,
Status::Canceled,
)
.await;
self.sessions
.cancel_stacked_child_sessions(&self.services, session_id)
.await;
}
pub(super) fn view_info_popup_mode(
title: String,
message: String,
is_loading: bool,
loading_label: String,
restore_view: ConfirmationViewMode,
) -> AppMode {
AppMode::ViewInfoPopup {
is_loading,
loading_label,
message,
restore_view,
title,
}
}
pub(super) fn branch_publish_loading_title(
publish_branch_action: PublishBranchAction,
) -> String {
branch_publish_loading_title_text(publish_branch_action)
}
pub(super) fn branch_publish_loading_message(
publish_branch_action: PublishBranchAction,
remote_branch_name: Option<&str>,
) -> String {
branch_publish_loading_message_text(publish_branch_action, remote_branch_name)
}
pub(super) fn branch_publish_loading_label(
publish_branch_action: PublishBranchAction,
) -> String {
branch_publish_loading_label_text(publish_branch_action)
}
pub(super) fn branch_publish_success_title(
publish_branch_action: PublishBranchAction,
) -> String {
branch_publish_success_title_text(publish_branch_action)
}
pub(super) fn branch_publish_success_message(
branch_name: &str,
review_request_creation: Option<&crate::app::branch_publish::ReviewRequestCreationInfo>,
) -> String {
crate::app::branch_publish::branch_push_success_message(
branch_name,
review_request_creation,
)
}
pub(super) fn review_request_publish_success_title(
review_request: &crate::domain::session::ReviewRequest,
) -> String {
crate::app::branch_publish::review_request_publish_success_title(review_request)
}
pub(super) fn pull_request_publish_success_message(
branch_name: &str,
review_request: &crate::domain::session::ReviewRequest,
) -> String {
pull_request_publish_success_message_text(branch_name, review_request)
}
pub(super) fn sync_main_popup_mode(
sync_main_result: Result<SyncMainOutcome, SyncSessionStartError>,
sync_popup_context: &SyncPopupContext,
) -> AppMode {
match sync_main_result {
Ok(sync_main_outcome) => AppMode::SyncBlockedPopup {
project_name: Some(sync_popup_context.project_name.clone()),
default_branch: Some(sync_popup_context.default_branch.clone()),
is_loading: false,
message: Self::sync_success_message(&sync_main_outcome),
title: "Sync complete".to_string(),
},
Err(sync_error @ SyncSessionStartError::MainHasUncommittedChanges { .. }) => {
AppMode::SyncBlockedPopup {
project_name: Some(sync_popup_context.project_name.clone()),
default_branch: Some(sync_popup_context.default_branch.clone()),
is_loading: false,
message: sync_error.detail_message(),
title: "Sync blocked".to_string(),
}
}
Err(sync_error @ SyncSessionStartError::Other(_)) => AppMode::SyncBlockedPopup {
project_name: Some(sync_popup_context.project_name.clone()),
default_branch: Some(sync_popup_context.default_branch.clone()),
is_loading: false,
message: Self::sync_failure_message(&sync_error),
title: "Sync failed".to_string(),
},
}
}
fn sync_success_message(sync_main_outcome: &SyncMainOutcome) -> String {
let pulled_summary = Self::sync_commit_summary("pulled", sync_main_outcome.pulled_commits);
let pulled_titles =
Self::sync_pulled_commit_titles_summary(&sync_main_outcome.pulled_commit_titles);
let pushed_titles =
Self::sync_pushed_commit_titles_summary(&sync_main_outcome.pushed_commit_titles);
let pushed_summary = Self::sync_commit_summary("pushed", sync_main_outcome.pushed_commits);
let conflict_summary =
Self::sync_conflict_summary(&sync_main_outcome.resolved_conflict_files);
sync_blocked::format_sync_success_message(
&pulled_summary,
&pulled_titles,
&pushed_summary,
&pushed_titles,
&conflict_summary,
)
}
fn sync_pulled_commit_titles_summary(pulled_commit_titles: &[String]) -> String {
if pulled_commit_titles.is_empty() {
return String::new();
}
pulled_commit_titles
.iter()
.map(|title| format!(" - {title}"))
.collect::<Vec<String>>()
.join("\n")
}
fn sync_pushed_commit_titles_summary(pushed_commit_titles: &[String]) -> String {
if pushed_commit_titles.is_empty() {
return String::new();
}
pushed_commit_titles
.iter()
.map(|title| format!(" - {title}"))
.collect::<Vec<String>>()
.join("\n")
}
fn sync_failure_message(sync_error: &SyncSessionStartError) -> String {
let detail_message = sync_error.detail_message();
if !is_git_push_authentication_error(&detail_message) {
return detail_message;
}
git_push_authentication_message(
detected_forge_kind_from_git_push_error(&detail_message),
"run sync again",
)
}
fn sync_commit_summary(direction: &str, commit_count: Option<u32>) -> String {
match commit_count {
Some(1) => format!("1 commit {direction}"),
Some(commit_count) => format!("{commit_count} commits {direction}"),
None => format!("commits {direction}: unknown"),
}
}
fn sync_conflict_summary(resolved_conflict_files: &[String]) -> String {
if resolved_conflict_files.is_empty() {
return "no conflicts fixed".to_string();
}
format!("conflicts fixed: {}", resolved_conflict_files.join(", "))
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn test_requested_review_batch_keeps_newer_generation_when_stale_event_arrives_later() {
let mut event_batch = AppEventBatch::default();
let newer_event = AppEvent::RequestedReviewsLoaded {
generation: 2,
project_id: 42,
result: Ok(Vec::new()),
};
let stale_event = AppEvent::RequestedReviewsLoaded {
generation: 1,
project_id: 42,
result: Err("stale failure".to_string()),
};
event_batch.collect_event(newer_event);
event_batch.collect_event(stale_event);
let (generation, project_id, result) = event_batch
.requested_reviews
.expect("newer requested-review event should be retained");
assert_eq!(generation, 2);
assert_eq!(project_id, 42);
assert_eq!(result.expect("newer result should be successful").len(), 0);
}
}