use std::future::Future;
use std::path::{Path, PathBuf};
use std::pin::Pin;
use std::sync::Arc;
use ag_forge::ReviewRequestClient;
use askama::Template;
use tokio::sync::mpsc;
use crate::app::error::AppError;
use crate::app::{AppEvent, UpdateStatus};
use crate::domain::agent::{AgentKind, AgentModel, ReasoningLevel};
use crate::domain::session::SessionId;
use crate::infra::agent;
use crate::infra::git::GitClient;
use crate::version;
pub(super) struct TaskService;
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,
}
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_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}"
))
})
}
}
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())
}
#[cfg(test)]
mod tests {
use std::path::Path;
use std::time::Duration;
use super::*;
use crate::infra::agent::protocol::AgentResponse;
#[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,
}
);
}
#[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");
}
#[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("put the Markdown review body in `answer`"));
assert!(prompt.contains("leave `questions` empty"));
assert!(prompt.contains("set `summary` to null"));
assert!(!prompt.contains("Return Markdown only."));
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"));
assert!(prompt.contains("Use the surrounding Agentty protocol for file-reference"));
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);
}
}