use std::collections::HashMap;
use std::future::Future;
use std::path::{Path, PathBuf};
use std::pin::Pin;
use std::sync::Arc;
use std::sync::atomic::{AtomicBool, Ordering};
use std::time::Duration;
use ag_forge::ReviewRequestClient;
use askama::Template;
use tokio::sync::mpsc;
use super::core::SyncReviewRequestTaskResult;
use crate::app::error::AppError;
use crate::app::session_state::SessionGitStatus;
use crate::app::{AppEvent, UpdateStatus, session};
use crate::domain::agent::{AgentKind, AgentModel, ReasoningLevel};
use crate::domain::session::SessionId;
use crate::infra::agent;
use crate::infra::git::GitClient;
use crate::infra::review_comment_cache::ReviewCommentCache;
use crate::version;
pub(super) struct TaskService;
#[derive(Clone, Debug, Eq, PartialEq)]
pub(crate) struct SessionGitStatusTarget {
pub(crate) base_branch: String,
pub(crate) branch_name: String,
pub(crate) session_id: SessionId,
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub(crate) struct ReviewRequestSyncTarget {
pub(crate) folder: PathBuf,
pub(crate) linked_review_request: Option<crate::domain::session::ReviewRequest>,
pub(crate) published_upstream_ref: Option<String>,
pub(crate) session_id: SessionId,
}
pub(super) struct ReviewAssistTaskInput {
pub(super) app_event_tx: mpsc::UnboundedSender<AppEvent>,
pub(super) diff_hash: u64,
pub(super) review_diff: String,
pub(super) review_model: AgentModel,
pub(super) session_folder: PathBuf,
pub(super) session_id: SessionId,
pub(super) session_summary: Option<String>,
}
#[derive(Template)]
#[template(path = "review_assist_prompt.md", escape = "none")]
struct ReviewAssistPromptTemplate<'a> {
fenced_diff: &'a str,
session_summary: &'a str,
}
const REVIEW_REQUEST_SYNC_INTERVAL_SECONDS: u64 = 60;
impl TaskService {
pub(super) async fn load_agent_availability(
availability_probe: Arc<dyn agent::AgentAvailabilityProbe>,
) -> Vec<AgentKind> {
tokio::task::spawn_blocking(move || availability_probe.available_agent_kinds())
.await
.unwrap_or_else(|_| AgentKind::ALL.to_vec())
}
pub(super) fn spawn_requested_reviews_task(
generation: u64,
project_id: i64,
working_dir: PathBuf,
app_event_tx: mpsc::UnboundedSender<AppEvent>,
git_client: Arc<dyn GitClient>,
review_request_client: Arc<dyn ReviewRequestClient>,
) {
tokio::spawn(async move {
let result = load_requested_reviews(
working_dir,
git_client.as_ref(),
review_request_client.as_ref(),
)
.await;
let _ = app_event_tx.send(AppEvent::RequestedReviewsLoaded {
generation,
project_id,
result,
});
});
}
pub(super) fn spawn_git_status_task(
working_dir: &Path,
project_branch_name: String,
session_git_status_targets: Vec<SessionGitStatusTarget>,
cancel: Arc<AtomicBool>,
app_event_tx: mpsc::UnboundedSender<AppEvent>,
git_client: Arc<dyn GitClient>,
) {
let dir = working_dir.to_path_buf();
tokio::spawn(async move {
let repo_root = git_client
.find_git_repo_root(dir.clone())
.await
.unwrap_or(dir);
loop {
if cancel.load(Ordering::Relaxed) {
break;
}
{
let root = repo_root.clone();
let _ = git_client.fetch_remote(root).await;
}
let branch_tracking_statuses = {
let root = repo_root.clone();
git_client
.branch_tracking_statuses(root)
.await
.unwrap_or_default()
};
let status = branch_tracking_statuses
.get(&project_branch_name)
.copied()
.flatten();
let session_git_statuses = Self::session_git_statuses(
&branch_tracking_statuses,
&repo_root,
&session_git_status_targets,
git_client.as_ref(),
);
let session_git_statuses = session_git_statuses.await;
if cancel.load(Ordering::Relaxed) {
break;
}
let _ = app_event_tx.send(AppEvent::GitStatusUpdated {
session_statuses: session_git_statuses,
status,
});
for _ in 0..30 {
if cancel.load(Ordering::Relaxed) {
return;
}
tokio::time::sleep(Duration::from_secs(1)).await;
}
}
});
}
async fn session_git_statuses(
branch_tracking_statuses: &HashMap<String, Option<(u32, u32)>>,
repo_root: &Path,
session_git_status_targets: &[SessionGitStatusTarget],
git_client: &dyn GitClient,
) -> HashMap<SessionId, SessionGitStatus> {
let mut session_git_statuses = HashMap::with_capacity(session_git_status_targets.len());
for session_git_status_target in session_git_status_targets {
let base_status = git_client
.get_ref_ahead_behind(
repo_root.to_path_buf(),
session_git_status_target.branch_name.clone(),
session_git_status_target.base_branch.clone(),
)
.await
.ok();
let remote_status = branch_tracking_statuses
.get(&session_git_status_target.branch_name)
.copied()
.flatten();
session_git_statuses.insert(
session_git_status_target.session_id.clone(),
SessionGitStatus {
base_status,
remote_status,
},
);
}
session_git_statuses
}
pub(super) fn spawn_review_request_status_task(
review_request_sync_targets: Vec<ReviewRequestSyncTarget>,
cancel: Arc<AtomicBool>,
app_event_tx: mpsc::UnboundedSender<AppEvent>,
git_client: Arc<dyn GitClient>,
review_request_client: Arc<dyn ReviewRequestClient>,
review_comment_cache: ReviewCommentCache,
) {
tokio::spawn(async move {
loop {
if cancel.load(Ordering::Relaxed) {
break;
}
for review_request_sync_target in &review_request_sync_targets {
if cancel.load(Ordering::Relaxed) {
return;
}
let result = sync_review_request_status(
review_request_sync_target.folder.clone(),
git_client.as_ref(),
review_request_sync_target.linked_review_request.clone(),
review_request_sync_target.published_upstream_ref.clone(),
review_request_client.as_ref(),
)
.await;
if cancel.load(Ordering::Relaxed) {
return;
}
let review_comments_sync_action = review_comments_sync_action(
review_request_sync_target.linked_review_request.as_ref(),
result.as_ref().ok(),
);
let _ = app_event_tx.send(AppEvent::ReviewRequestStatusUpdated {
result,
session_id: review_request_sync_target.session_id.clone(),
});
match review_comments_sync_action {
ReviewCommentsSyncAction::Fetch(display_id) => {
sync_review_comments_for_target(
review_request_sync_target,
display_id,
git_client.as_ref(),
review_request_client.as_ref(),
&review_comment_cache,
&app_event_tx,
)
.await;
}
ReviewCommentsSyncAction::Forget => {
forget_review_comments_for_session(
&review_request_sync_target.session_id,
&review_comment_cache,
&app_event_tx,
);
}
ReviewCommentsSyncAction::Skip => {}
}
}
for _ in 0..REVIEW_REQUEST_SYNC_INTERVAL_SECONDS {
if cancel.load(Ordering::Relaxed) {
return;
}
tokio::time::sleep(Duration::from_secs(1)).await;
}
}
});
}
pub(super) fn spawn_version_check_task(
app_event_tx: &mpsc::UnboundedSender<AppEvent>,
auto_update: bool,
) {
#[cfg(test)]
{
let _ = auto_update;
let _ = app_event_tx.send(Self::version_availability_event(None));
}
#[cfg(not(test))]
let app_event_tx = app_event_tx.clone();
#[cfg(not(test))]
tokio::spawn(async move {
let latest_version_tag = version::latest_npm_version_tag().await;
let version_event = Self::version_availability_event(latest_version_tag);
let newer_version = match &version_event {
AppEvent::VersionAvailabilityUpdated {
latest_available_version: Some(version),
} => Some(version.clone()),
_ => None,
};
let _ = app_event_tx.send(version_event);
if let Some(newer_version) = newer_version
&& auto_update
{
Self::run_background_update(&app_event_tx, &newer_version).await;
}
});
}
#[cfg(not(test))]
async fn run_background_update(
app_event_tx: &mpsc::UnboundedSender<AppEvent>,
newer_version: &str,
) {
let _ = app_event_tx.send(AppEvent::UpdateStatusChanged {
update_status: UpdateStatus::InProgress {
version: newer_version.to_string(),
},
});
let update_result = tokio::task::spawn_blocking(move || {
let update_runner = version::RealUpdateRunner;
version::run_npm_update_sync(&update_runner)
})
.await;
let update_status = match update_result {
Ok(Ok(_)) => UpdateStatus::Complete {
version: newer_version.to_string(),
},
Ok(Err(_)) | Err(_) => UpdateStatus::Failed {
version: newer_version.to_string(),
},
};
let _ = app_event_tx.send(AppEvent::UpdateStatusChanged { update_status });
}
pub(super) fn spawn_review_assist_task(input: ReviewAssistTaskInput) {
let ReviewAssistTaskInput {
app_event_tx,
diff_hash,
review_diff,
review_model,
session_folder,
session_id,
session_summary,
} = input;
tokio::spawn(async move {
let review_result = Self::review_assist_text(
&session_folder,
review_model,
&review_diff,
session_summary.as_deref(),
)
.await;
let app_event = Self::review_app_event(diff_hash, review_result, session_id);
let _ = app_event_tx.send(app_event);
});
}
async fn review_assist_text(
session_folder: &Path,
review_model: AgentModel,
review_diff: &str,
session_summary: Option<&str>,
) -> Result<String, AppError> {
Self::review_assist_text_with_submitter(
session_folder,
review_model,
review_diff,
session_summary,
|review_folder, review_model, review_prompt| {
Box::pin(async move {
agent::submit_one_shot(agent::OneShotRequest {
child_pid: None,
folder: review_folder,
model: review_model,
prompt: review_prompt,
request_kind: crate::infra::channel::AgentRequestKind::UtilityPrompt,
reasoning_level: ReasoningLevel::default(),
})
.await
})
},
)
.await
}
fn version_availability_event(latest_version_tag: Option<String>) -> AppEvent {
let latest_available_version = latest_version_tag.filter(|latest_version| {
version::is_newer_than_current_version(env!("CARGO_PKG_VERSION"), latest_version)
});
AppEvent::VersionAvailabilityUpdated {
latest_available_version,
}
}
async fn review_assist_text_with_submitter<Submitter>(
session_folder: &Path,
review_model: AgentModel,
review_diff: &str,
session_summary: Option<&str>,
submitter: Submitter,
) -> Result<String, AppError>
where
Submitter: for<'submit> FnOnce(
&'submit Path,
AgentModel,
&'submit str,
) -> Pin<
Box<dyn Future<Output = Result<agent::AgentResponse, String>> + Send + 'submit>,
>,
{
let review_prompt = Self::review_assist_prompt(review_diff, session_summary)?;
let agent_response = submitter(session_folder, review_model, &review_prompt)
.await
.map_err(AppError::Workflow)?;
Self::review_output_text(&agent_response)
}
fn review_app_event(
diff_hash: u64,
review_result: Result<String, AppError>,
session_id: SessionId,
) -> AppEvent {
match review_result {
Ok(review_text) => AppEvent::ReviewPrepared {
diff_hash,
review_text,
session_id,
},
Err(error) => AppEvent::ReviewPreparationFailed {
diff_hash,
error: error.to_string(),
session_id,
},
}
}
fn review_output_text(agent_response: &agent::AgentResponse) -> Result<String, AppError> {
let review_text = agent_response.to_display_text();
let review_text = review_text.trim();
if review_text.is_empty() {
return Err(AppError::Workflow(
"Review assist returned empty output".to_string(),
));
}
Ok(review_text.to_string())
}
fn review_assist_prompt(
review_diff: &str,
session_summary: Option<&str>,
) -> Result<String, AppError> {
let trimmed_diff = review_diff.trim();
let fence = agent::diff_fence(trimmed_diff);
let fenced_diff = format!("{fence}diff\n{trimmed_diff}\n{fence}");
let template = ReviewAssistPromptTemplate {
fenced_diff: &fenced_diff,
session_summary: session_summary.map_or("", str::trim),
};
template.render().map_err(|error| {
AppError::Workflow(format!(
"Failed to render `review_assist_prompt.md`: {error}"
))
})
}
}
#[derive(Clone, Debug, Eq, PartialEq)]
enum ReviewCommentsSyncAction {
Fetch(String),
Forget,
Skip,
}
fn review_comments_sync_action(
linked_review_request: Option<&crate::domain::session::ReviewRequest>,
sync_result: Option<&SyncReviewRequestTaskResult>,
) -> ReviewCommentsSyncAction {
let observed_forge_kind = sync_result
.and_then(|result| result.summary.as_ref())
.map(|summary| summary.forge_kind)
.or_else(|| linked_review_request.map(|linked| linked.summary.forge_kind));
if let Some(forge_kind) = observed_forge_kind
&& !forge_kind.supports_review_comments_preview()
{
return ReviewCommentsSyncAction::Skip;
}
if let Some(result) = sync_result {
return match &result.outcome {
session::SyncReviewRequestOutcome::Open { display_id, .. } => {
ReviewCommentsSyncAction::Fetch(display_id.clone())
}
_ => ReviewCommentsSyncAction::Forget,
};
}
let Some(linked) = linked_review_request else {
return ReviewCommentsSyncAction::Skip;
};
if linked.summary.state == ag_forge::ReviewRequestState::Open {
return ReviewCommentsSyncAction::Fetch(linked.summary.display_id.clone());
}
ReviewCommentsSyncAction::Forget
}
fn forget_review_comments_for_session(
session_id: &SessionId,
review_comment_cache: &ReviewCommentCache,
app_event_tx: &mpsc::UnboundedSender<AppEvent>,
) {
if review_comment_cache.forget(session_id) {
let _ = app_event_tx.send(AppEvent::ReviewCommentsUpdated {
session_id: session_id.clone(),
});
}
}
async fn sync_review_comments_for_target(
review_request_sync_target: &ReviewRequestSyncTarget,
display_id: String,
git_client: &dyn GitClient,
review_request_client: &dyn ReviewRequestClient,
review_comment_cache: &ReviewCommentCache,
app_event_tx: &mpsc::UnboundedSender<AppEvent>,
) {
let result = fetch_review_comment_snapshot(
review_request_sync_target.folder.clone(),
display_id,
git_client,
review_request_client,
)
.await;
let session_id = review_request_sync_target.session_id.clone();
match result {
Ok(snapshot) => {
if review_comment_cache.record_snapshot(session_id.clone(), snapshot) {
let _ = app_event_tx.send(AppEvent::ReviewCommentsUpdated { session_id });
}
}
Err(error) => {
tracing::warn!(
session_id = %session_id,
"failed to fetch review comments: {error}",
);
}
}
}
async fn fetch_review_comment_snapshot(
folder: PathBuf,
display_id: String,
git_client: &dyn GitClient,
review_request_client: &dyn ReviewRequestClient,
) -> Result<ag_forge::ReviewCommentSnapshot, String> {
let repo_url = git_client
.repo_url(folder.clone())
.await
.map_err(|error| format!("Failed to resolve repository remote: {error}"))?;
let remote = review_request_client
.detect_remote(repo_url)
.map(|remote| remote.with_command_working_directory(folder))
.map_err(|error| error.detail_message())?;
review_request_client
.fetch_review_comment_snapshot(remote, display_id)
.await
.map_err(|error| error.detail_message())
}
async fn load_requested_reviews(
working_dir: PathBuf,
git_client: &dyn GitClient,
review_request_client: &dyn ReviewRequestClient,
) -> Result<Vec<ag_forge::RequestedReview>, String> {
let repo_url = git_client
.repo_url(working_dir.clone())
.await
.map_err(|error| format!("Failed to resolve repository remote: {error}"))?;
let remote = review_request_client
.detect_remote(repo_url)
.map(|remote| remote.with_command_working_directory(working_dir))
.map_err(|error| error.detail_message())?;
review_request_client
.list_requested_reviews(remote)
.await
.map_err(|error| error.detail_message())
}
async fn sync_review_request_status(
folder: PathBuf,
git_client: &dyn GitClient,
linked_review_request: Option<crate::domain::session::ReviewRequest>,
published_upstream_ref: Option<String>,
review_request_client: &dyn ReviewRequestClient,
) -> Result<SyncReviewRequestTaskResult, String> {
let repo_url = git_client
.repo_url(folder.clone())
.await
.map_err(|error| format!("Failed to resolve repository remote: {error}"))?;
let remote = review_request_client
.detect_remote(repo_url)
.map(|remote| remote.with_command_working_directory(folder))
.map_err(|error| error.detail_message())?;
if let Some(review_request) = linked_review_request {
let refreshed_summary = review_request_client
.refresh_review_request(remote, review_request.summary.display_id)
.await
.map_err(|error| error.detail_message())?;
return Ok(sync_task_result_from_summary(refreshed_summary));
}
let upstream_ref = published_upstream_ref
.ok_or_else(|| "Session branch has not been published yet".to_string())?;
let source_branch = session::remote_branch_name_from_upstream_ref(&upstream_ref);
let found_summary = review_request_client
.find_by_source_branch(remote, source_branch)
.await
.map_err(|error| error.detail_message())?;
match found_summary {
Some(summary) => Ok(sync_task_result_from_summary(summary)),
None => Ok(SyncReviewRequestTaskResult {
outcome: session::SyncReviewRequestOutcome::NoReviewRequest,
summary: None,
}),
}
}
fn sync_task_result_from_summary(
summary: crate::domain::session::ReviewRequestSummary,
) -> SyncReviewRequestTaskResult {
let display_id = summary.display_id.clone();
let outcome = match summary.state {
crate::domain::session::ReviewRequestState::Open => {
session::SyncReviewRequestOutcome::Open {
display_id,
status_summary: summary.status_summary.clone(),
}
}
crate::domain::session::ReviewRequestState::Merged => {
session::SyncReviewRequestOutcome::Merged { display_id }
}
crate::domain::session::ReviewRequestState::Closed => {
session::SyncReviewRequestOutcome::Closed { display_id }
}
};
SyncReviewRequestTaskResult {
outcome,
summary: Some(summary),
}
}
#[cfg(test)]
mod tests {
use std::path::{Path, PathBuf};
use ag_forge::{ForgeKind, MockReviewRequestClient, ReviewRequestState, ReviewRequestSummary};
use super::*;
use crate::infra::agent::protocol::AgentResponse;
use crate::infra::git::{GitError, MockGitClient};
#[tokio::test]
async fn spawn_version_check_task_emits_none_update_in_tests() {
let (app_event_tx, mut app_event_rx) = mpsc::unbounded_channel();
TaskService::spawn_version_check_task(&app_event_tx, true);
let app_event = tokio::time::timeout(Duration::from_secs(1), app_event_rx.recv())
.await
.expect("timed out waiting for version-check event")
.expect("version-check task should emit one event");
assert_eq!(
app_event,
AppEvent::VersionAvailabilityUpdated {
latest_available_version: None,
}
);
}
#[tokio::test]
async fn spawn_version_check_task_with_no_update_emits_version_event() {
let (app_event_tx, mut app_event_rx) = mpsc::unbounded_channel();
TaskService::spawn_version_check_task(&app_event_tx, false);
let app_event = tokio::time::timeout(Duration::from_secs(1), app_event_rx.recv())
.await
.expect("timed out waiting for version-check event")
.expect("version-check task should emit one event");
assert_eq!(
app_event,
AppEvent::VersionAvailabilityUpdated {
latest_available_version: None,
}
);
}
#[test]
fn version_availability_event_keeps_newer_version_tags() {
let latest_version_tag = Some("v999.0.0".to_string());
let app_event = TaskService::version_availability_event(latest_version_tag);
assert_eq!(
app_event,
AppEvent::VersionAvailabilityUpdated {
latest_available_version: Some("v999.0.0".to_string()),
}
);
}
#[test]
fn version_availability_event_ignores_current_version_tag() {
let latest_version_tag = Some(format!("v{}", env!("CARGO_PKG_VERSION")));
let app_event = TaskService::version_availability_event(latest_version_tag);
assert_eq!(
app_event,
AppEvent::VersionAvailabilityUpdated {
latest_available_version: None,
}
);
}
fn test_review_request_summary(
display_id: &str,
state: ReviewRequestState,
) -> ReviewRequestSummary {
test_review_request_summary_with_forge(display_id, state, ForgeKind::GitHub)
}
fn test_review_request_summary_with_forge(
display_id: &str,
state: ReviewRequestState,
forge_kind: ForgeKind,
) -> ReviewRequestSummary {
ReviewRequestSummary {
display_id: display_id.to_string(),
forge_kind,
source_branch: "wt/session-id".to_string(),
state,
status_summary: None,
target_branch: "main".to_string(),
title: "feat".to_string(),
web_url: String::new(),
}
}
#[tokio::test]
async fn review_assist_text_with_submitter_returns_workflow_error_on_submit_failure() {
let session_folder = Path::new("/tmp/review-assist-submit-error");
let review_model = AgentModel::ClaudeSonnet46;
let review_diff = "diff --git a/src/lib.rs b/src/lib.rs";
let result = TaskService::review_assist_text_with_submitter(
session_folder,
review_model,
review_diff,
None,
|_, _, _| Box::pin(async { Err("submit failed".to_string()) }),
)
.await;
let error = result.expect_err("submit failure should be returned");
assert!(
matches!(error, AppError::Workflow(_)),
"expected AppError::Workflow, got: {error:?}"
);
assert_eq!(error.to_string(), "submit failed");
}
#[tokio::test]
async fn session_git_statuses_collects_all_target_statuses() {
let repo_root = Path::new("/tmp/task-service-session-statuses");
let branch_tracking_statuses = HashMap::from([
("wt/session-a".to_string(), Some((7, 0))),
("wt/session-b".to_string(), Some((0, 4))),
]);
let session_git_status_targets = vec![
SessionGitStatusTarget {
base_branch: "main".to_string(),
branch_name: "wt/session-a".to_string(),
session_id: "session-a".into(),
},
SessionGitStatusTarget {
base_branch: "develop".to_string(),
branch_name: "wt/session-b".to_string(),
session_id: "session-b".into(),
},
];
let mut mock_git_client = MockGitClient::new();
mock_git_client
.expect_get_ref_ahead_behind()
.times(2)
.returning(|_, left_ref, right_ref| {
Box::pin(async move {
match (left_ref.as_str(), right_ref.as_str()) {
("wt/session-a", "main") => Ok((2, 1)),
("wt/session-b", "develop") => Ok((0, 0)),
_ => Err(GitError::OutputParse("unexpected ref pair".to_string())),
}
})
});
let statuses = TaskService::session_git_statuses(
&branch_tracking_statuses,
repo_root,
&session_git_status_targets,
&mock_git_client,
)
.await;
assert_eq!(
statuses.get("session-a"),
Some(&SessionGitStatus {
base_status: Some((2, 1)),
remote_status: Some((7, 0)),
})
);
assert_eq!(
statuses.get("session-b"),
Some(&SessionGitStatus {
base_status: Some((0, 0)),
remote_status: Some((0, 4)),
})
);
}
#[tokio::test]
async fn session_git_statuses_keeps_failed_targets_as_none() {
let repo_root = Path::new("/tmp/task-service-session-statuses-error");
let branch_tracking_statuses = HashMap::new();
let session_git_status_targets = vec![SessionGitStatusTarget {
base_branch: "main".to_string(),
branch_name: "wt/session-a".to_string(),
session_id: "session-a".into(),
}];
let mut mock_git_client = MockGitClient::new();
mock_git_client
.expect_get_ref_ahead_behind()
.once()
.returning(|_, _, _| {
Box::pin(async {
Err(GitError::OutputParse(
"failed to compare session branch".to_string(),
))
})
});
let statuses = TaskService::session_git_statuses(
&branch_tracking_statuses,
repo_root,
&session_git_status_targets,
&mock_git_client,
)
.await;
assert_eq!(
statuses.get("session-a"),
Some(&SessionGitStatus {
base_status: None,
remote_status: None,
})
);
}
#[test]
fn sync_task_result_from_open_summary_maps_open_outcome() {
let mut summary = test_review_request_summary("#42", ReviewRequestState::Open);
summary.status_summary = Some("Checks passing".to_string());
let result = sync_task_result_from_summary(summary);
assert_eq!(
result.outcome,
session::SyncReviewRequestOutcome::Open {
display_id: "#42".to_string(),
status_summary: Some("Checks passing".to_string()),
}
);
assert!(result.summary.is_some());
}
#[tokio::test]
async fn sync_review_request_status_attaches_worktree_to_detected_remote() {
let folder = PathBuf::from("/tmp/session-worktree");
let linked_review_request = crate::domain::session::ReviewRequest {
last_refreshed_at: 42,
summary: test_review_request_summary("#42", ReviewRequestState::Open),
};
let expected_remote = ag_forge::ForgeRemote {
command_working_directory: Some(folder.clone()),
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(),
};
let expected_summary = test_review_request_summary("#42", ReviewRequestState::Merged);
let mut mock_git_client = MockGitClient::new();
mock_git_client
.expect_repo_url()
.once()
.withf({
let folder = folder.clone();
move |candidate_folder| candidate_folder == &folder
})
.returning(|_| {
Box::pin(async { Ok("https://github.com/agentty-xyz/agentty.git".to_string()) })
});
let mut mock_review_request_client = MockReviewRequestClient::new();
mock_review_request_client
.expect_detect_remote()
.once()
.withf(|repo_url| repo_url == "https://github.com/agentty-xyz/agentty.git")
.returning(|_| {
Ok(ag_forge::ForgeRemote {
command_working_directory: None,
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(),
})
});
mock_review_request_client
.expect_refresh_review_request()
.once()
.withf({
let expected_remote = expected_remote.clone();
move |candidate_remote, display_id| {
candidate_remote == &expected_remote && display_id == "#42"
}
})
.returning({
let expected_summary = expected_summary.clone();
move |_, _| {
let expected_summary = expected_summary.clone();
Box::pin(async move { Ok(expected_summary) })
}
});
let result = sync_review_request_status(
folder,
&mock_git_client,
Some(linked_review_request),
None,
&mock_review_request_client,
)
.await
.expect("sync should succeed");
assert_eq!(
result.outcome,
session::SyncReviewRequestOutcome::Merged {
display_id: "#42".to_string(),
}
);
}
#[test]
fn sync_task_result_from_merged_summary_maps_merged_outcome() {
let summary = test_review_request_summary("#99", ReviewRequestState::Merged);
let result = sync_task_result_from_summary(summary);
assert_eq!(
result.outcome,
session::SyncReviewRequestOutcome::Merged {
display_id: "#99".to_string(),
}
);
assert!(result.summary.is_some());
}
#[test]
fn sync_task_result_from_closed_summary_maps_closed_outcome() {
let summary = test_review_request_summary("#7", ReviewRequestState::Closed);
let result = sync_task_result_from_summary(summary);
assert_eq!(
result.outcome,
session::SyncReviewRequestOutcome::Closed {
display_id: "#7".to_string(),
}
);
assert!(result.summary.is_some());
}
#[test]
fn review_app_event_maps_successful_review_output() {
let diff_hash = 7;
let review_result = Ok("Flagged one missing error branch.".to_string());
let session_id = "session-7".to_string();
let app_event = TaskService::review_app_event(diff_hash, review_result, session_id.into());
assert_eq!(
app_event,
AppEvent::ReviewPrepared {
diff_hash: 7,
review_text: "Flagged one missing error branch.".to_string(),
session_id: "session-7".into(),
}
);
}
#[test]
fn review_app_event_maps_failure_output() {
let diff_hash = 9;
let review_result = Err(AppError::Workflow("empty response".to_string()));
let session_id = "session-9".to_string();
let app_event = TaskService::review_app_event(diff_hash, review_result, session_id.into());
assert_eq!(
app_event,
AppEvent::ReviewPreparationFailed {
diff_hash: 9,
error: "empty response".to_string(),
session_id: "session-9".into(),
}
);
}
#[test]
fn review_output_text_trims_agent_response_text() {
let agent_response = AgentResponse::plain(" Review looks good. \n");
let review_text = TaskService::review_output_text(&agent_response)
.expect("non-empty output should be accepted");
assert_eq!(review_text, "Review looks good.");
}
#[test]
fn review_output_text_rejects_blank_agent_response_text() {
let agent_response = AgentResponse::plain(" \n\t ");
let result = TaskService::review_output_text(&agent_response);
let error = result.expect_err("blank output should be rejected");
assert!(
matches!(error, AppError::Workflow(_)),
"expected AppError::Workflow, got: {error:?}"
);
assert_eq!(error.to_string(), "Review assist returned empty output");
}
#[test]
fn test_review_assist_prompt_enforces_read_only_constraints() {
let review_diff = "diff --git a/src/lib.rs b/src/lib.rs";
let session_summary = Some("Refactor parser error mapping.");
let prompt = TaskService::review_assist_prompt(review_diff, session_summary)
.expect("review prompt should render");
assert!(prompt.contains("You are in read-only review mode."));
assert!(prompt.contains("Do not create, modify, rename, or delete files."));
assert!(prompt.contains("You may browse the internet when needed."));
assert!(prompt.contains("You may run non-editing CLI commands"));
let fenced_diff = format!("```diff\n{review_diff}\n```");
assert!(
prompt.contains(&fenced_diff),
"review prompt must wrap the diff in a ```diff``` fence so `@`-prefixed decorator \
tokens are not misread as file mentions"
);
assert!(prompt.contains("`@`-prefixed tokens inside the diff"));
}
#[test]
fn test_review_assist_prompt_escapes_triple_backtick_fence_in_diff() {
let review_diff = concat!(
"diff --git a/notes.md b/notes.md\n",
"+```\n",
"+example fenced block\n",
"+```\n",
);
let session_summary: Option<&str> = None;
let prompt = TaskService::review_assist_prompt(review_diff, session_summary)
.expect("review prompt should render");
assert!(
prompt.contains("````diff\n"),
"outer fence must be longer than the longest backtick run in the diff to preserve \
prompt boundaries"
);
let matches = prompt.matches("\n````").count();
assert!(
matches >= 2,
"prompt must contain an opening and closing 4-backtick fence, got {matches} \
occurrences"
);
assert!(prompt.contains("+```\n"));
}
#[test]
fn test_structured_agent_response_is_unwrapped_to_display_text() {
let structured_json = r#"{"answer":"Review looks good.","questions":[],"summary":null}"#;
let agent_response = agent::protocol::parse_agent_response_strict(structured_json)
.expect("structured response should parse");
let display_text = agent_response.to_display_text();
assert_eq!(display_text.trim(), "Review looks good.");
}
#[test]
fn update_status_changed_event_roundtrips_all_variants() {
let in_progress = AppEvent::UpdateStatusChanged {
update_status: UpdateStatus::InProgress {
version: "v1.0.0".to_string(),
},
};
let complete = AppEvent::UpdateStatusChanged {
update_status: UpdateStatus::Complete {
version: "v1.0.0".to_string(),
},
};
let failed = AppEvent::UpdateStatusChanged {
update_status: UpdateStatus::Failed {
version: "v1.0.0".to_string(),
},
};
assert_ne!(in_progress, complete);
assert_ne!(complete, failed);
assert_ne!(in_progress, failed);
}
fn linked_review_request(
display_id: &str,
state: ReviewRequestState,
) -> crate::domain::session::ReviewRequest {
crate::domain::session::ReviewRequest {
last_refreshed_at: 0,
summary: test_review_request_summary(display_id, state),
}
}
#[test]
fn review_comments_sync_action_fetches_on_open_sync_outcome() {
let linked = linked_review_request("#1", ReviewRequestState::Open);
let sync_result = SyncReviewRequestTaskResult {
outcome: session::SyncReviewRequestOutcome::Open {
display_id: "#42".to_string(),
status_summary: None,
},
summary: None,
};
let action = review_comments_sync_action(Some(&linked), Some(&sync_result));
assert_eq!(action, ReviewCommentsSyncAction::Fetch("#42".to_string()));
}
#[test]
fn review_comments_sync_action_forgets_on_successful_merged_sync() {
let linked = linked_review_request("#7", ReviewRequestState::Open);
let sync_result = SyncReviewRequestTaskResult {
outcome: session::SyncReviewRequestOutcome::Merged {
display_id: "#7".to_string(),
},
summary: None,
};
let action = review_comments_sync_action(Some(&linked), Some(&sync_result));
assert_eq!(action, ReviewCommentsSyncAction::Forget);
}
#[test]
fn review_comments_sync_action_forgets_on_no_review_request_sync_outcome() {
let linked = linked_review_request("#21", ReviewRequestState::Open);
let sync_result = SyncReviewRequestTaskResult {
outcome: session::SyncReviewRequestOutcome::NoReviewRequest,
summary: None,
};
let action = review_comments_sync_action(Some(&linked), Some(&sync_result));
assert_eq!(action, ReviewCommentsSyncAction::Forget);
}
#[test]
fn review_comments_sync_action_falls_back_to_linked_open_request_on_sync_failure() {
let linked = linked_review_request("#11", ReviewRequestState::Open);
let action = review_comments_sync_action(Some(&linked), None);
assert_eq!(action, ReviewCommentsSyncAction::Fetch("#11".to_string()));
}
#[test]
fn review_comments_sync_action_forgets_on_failure_when_linked_is_not_open() {
let linked = linked_review_request("#13", ReviewRequestState::Closed);
let action = review_comments_sync_action(Some(&linked), None);
assert_eq!(action, ReviewCommentsSyncAction::Forget);
}
#[test]
fn review_comments_sync_action_fetches_on_gitlab_open_sync_outcome() {
let linked = crate::domain::session::ReviewRequest {
last_refreshed_at: 0,
summary: test_review_request_summary_with_forge(
"!42",
ReviewRequestState::Open,
ForgeKind::GitLab,
),
};
let sync_result = SyncReviewRequestTaskResult {
outcome: session::SyncReviewRequestOutcome::Open {
display_id: "!42".to_string(),
status_summary: None,
},
summary: Some(test_review_request_summary_with_forge(
"!42",
ReviewRequestState::Open,
ForgeKind::GitLab,
)),
};
let action = review_comments_sync_action(Some(&linked), Some(&sync_result));
assert_eq!(action, ReviewCommentsSyncAction::Fetch("!42".to_string()));
}
#[test]
fn review_comments_sync_action_fetches_on_gitlab_sync_failure() {
let linked = crate::domain::session::ReviewRequest {
last_refreshed_at: 0,
summary: test_review_request_summary_with_forge(
"!11",
ReviewRequestState::Open,
ForgeKind::GitLab,
),
};
let action = review_comments_sync_action(Some(&linked), None);
assert_eq!(action, ReviewCommentsSyncAction::Fetch("!11".to_string()));
}
#[test]
fn review_comments_sync_action_skips_on_failure_without_linked_request() {
let action = review_comments_sync_action(None, None);
assert_eq!(action, ReviewCommentsSyncAction::Skip);
}
#[tokio::test]
async fn forget_review_comments_for_session_emits_update_when_entry_removed() {
let session_id: SessionId = "session-forget".into();
let review_comment_cache = ReviewCommentCache::default();
review_comment_cache.record_snapshot(
session_id.clone(),
ag_forge::ReviewCommentSnapshot::default(),
);
let (app_event_tx, mut app_event_rx) = mpsc::unbounded_channel();
forget_review_comments_for_session(&session_id, &review_comment_cache, &app_event_tx);
let event = app_event_rx
.try_recv()
.expect("cache removal should emit a refresh event");
assert!(matches!(
event,
AppEvent::ReviewCommentsUpdated { ref session_id } if session_id.as_str() == "session-forget"
));
}
#[tokio::test]
async fn forget_review_comments_for_session_stays_silent_when_entry_missing() {
let session_id: SessionId = "session-missing".into();
let review_comment_cache = ReviewCommentCache::default();
let (app_event_tx, mut app_event_rx) = mpsc::unbounded_channel();
forget_review_comments_for_session(&session_id, &review_comment_cache, &app_event_tx);
assert!(
app_event_rx.try_recv().is_err(),
"missing cache entry should not trigger an event"
);
}
}