use super::*;
pub(super) fn prompt_command(prompt: String) -> RelayCommand {
RelayCommand::Prompt {
prompt: vec![agent_client_protocol::schema::v1::ContentBlock::Text(
agent_client_protocol::schema::v1::TextContent::new(prompt),
)],
}
}
pub(super) async fn reviewer_action(
control: &SessionManagerControl,
session_id: &str,
role: Option<String>,
action: ReviewerAction,
) -> Result<ReviewerOutcome, String> {
let handle: ManagedSessionHandle = control
.session(session_id.to_owned())
.await
.map_err(|error| format!("{error:#}"))?;
let deadline = tokio::time::Instant::now() + LIFECYCLE_WAIT;
let mut pause = LEASE_RETRY_PAUSE;
loop {
match handle
.reviewer_as(role.clone(), action.clone())
.await
.map_err(|error| format!("{error:#}"))
{
Err(reason)
if preempted_by_lifecycle(&reason)
&& tokio::time::Instant::now() + pause < deadline =>
{
tracing::info!(
%session_id,
operation = action.operation_name(),
%reason,
"a running review waits for another operation on the session"
);
tokio::time::sleep(pause).await;
pause = (pause * 2).min(LIFECYCLE_RETRY_PAUSE_MAX);
}
outcome => return outcome,
}
}
}
pub(super) const LIFECYCLE_WAIT: std::time::Duration = std::time::Duration::from_secs(300);
const LIFECYCLE_RETRY_PAUSE_MAX: std::time::Duration = std::time::Duration::from_secs(5);
pub(super) async fn launch_role(
control: &SessionManagerControl,
environment: &Arc<dyn ReviewEnvironment>,
session_id: &str,
role: &str,
reviewer: &ReviewerIdentity,
generation: u64,
) -> Result<(), String> {
let staged = {
let session_id = session_id.to_owned();
let profile = reviewer.profile.clone();
let environment = environment.clone();
tokio::task::spawn_blocking(move || environment.stage(&session_id, &profile, generation))
.await
.map_err(|error| format!("staging the reviewer stopped: {error}"))??
};
let mut config = staged;
config.model = reviewer.main.model.clone();
config.effort = reviewer.main.effort.clone();
config.fast_mode = reviewer.main.fast_mode.then_some(true);
match reviewer_action(
control,
session_id,
Some(role.to_owned()),
ReviewerAction::Start {
config: Box::new(config),
},
)
.await
{
Ok(ReviewerOutcome::Started(_)) => Ok(()),
other => Err(unexpected(other)),
}
}
pub(super) async fn prepare(
control: &SessionManagerControl,
environment: &Arc<dyn ReviewEnvironment>,
session_id: &str,
config: ReviewConfig,
cancelled: Arc<std::sync::atomic::AtomicBool>,
) -> Result<Prepared, StartRefusal> {
let child = {
let environment = environment.clone();
let session = session_id.to_owned();
tokio::task::spawn_blocking(move || environment.is_subagent(&session))
.await
.unwrap_or(false)
};
if child {
return Err(StartRefusal(SUBAGENT_REFUSAL.to_owned()));
}
let handle = control
.session(session_id.to_owned())
.await
.map_err(|error| StartRefusal(format!("{error:#}")))?;
match handle.reviewer(ReviewerAction::Status).await {
Ok(ReviewerOutcome::Status(state)) => match &state.active_prompt {
Some(prompt) if is_turn_review_command(&prompt.command_id) => {
stop_leftover_review(&handle).await.map_err(|error| {
StartRefusal(format!(
"the review left running when Mjolnir restarted could not be \
stopped: {error}"
))
})?;
}
Some(prompt) => return Err(StartRefusal(busy_reviewer(&prompt.command_id))),
None => {}
},
Ok(_) => {}
Err(error) => return Err(StartRefusal(format!("{error:#}"))),
}
let view = handle.view();
if !view.connected {
return Err(StartRefusal("this session is not connected".to_owned()));
}
let Some(snapshot) = view.snapshot else {
return Err(StartRefusal(
"this session has no transcript yet".to_owned(),
));
};
if !matches!(
snapshot.materialized.execution,
MaterializedExecutionState::Idle
) {
return Err(StartRefusal(
"a review runs between turns; this one is still working".to_owned(),
));
}
if !snapshot.materialized.queued_prompts.is_empty() {
return Err(StartRefusal(
"prompts are queued; the review waits for them".to_owned(),
));
}
let state = {
let session = session_id.to_owned();
let environment = environment.clone();
tokio::task::spawn_blocking(move || environment.load_state(&session))
.await
.map_err(|e| StartRefusal(format!("preparing review: {e}")))?
.map_err(StartRefusal)?
};
let deadline = tokio::time::Instant::now() + BACKGROUND_WORK_WAIT;
let background = environment.hold_background_work(session_id, deadline).await;
let captured = {
let handle = &handle;
let baselines = &state.baselines;
after_background_work(session_id, deadline, &cancelled, move || async move {
handle
.reviewer(ReviewerAction::CaptureDelta {
baselines: baselines.clone(),
})
.await
.map_err(|error| format!("{error:#}"))
})
.await
};
let captured = match captured {
Ok(ReviewerOutcome::Delta { repositories }) => Some(repositories),
_ => None,
};
if cancelled.load(std::sync::atomic::Ordering::Acquire) {
return Err(StartRefusal("review preparation cancelled".into()));
}
if captured
.as_deref()
.is_some_and(|deltas| !mj_review::delta::has_changes(deltas))
{
return Ok(Prepared {
state,
reviewer: ReviewerIdentity::default(),
materialized: Box::new(snapshot.materialized),
resume_forward: None,
captured,
background,
});
}
let reviewer = after_background_work(session_id, deadline, &cancelled, || {
environment.resolve(handle.clone(), config.clone(), cancelled.clone())
})
.await
.map_err(StartRefusal)?;
if cancelled.load(std::sync::atomic::Ordering::Acquire) {
return Err(StartRefusal("review preparation cancelled".into()));
}
let profile = reviewer.profile.clone();
let session = session_id.to_owned();
let environment = environment.clone();
tokio::task::spawn_blocking(move || environment.check(&session, &profile))
.await
.map_err(|e| StartRefusal(format!("preparing review: {e}")))?
.map_err(StartRefusal)?;
Ok(Prepared {
state,
reviewer: reviewer.clone(),
materialized: Box::new(snapshot.materialized),
resume_forward: None,
captured,
background,
})
}
pub(super) async fn prepare_recovery(
control: &SessionManagerControl,
environment: &Arc<dyn ReviewEnvironment>,
session_id: &str,
) -> Result<Option<Prepared>, String> {
let session = session_id.to_owned();
let environment = environment.clone();
let state = tokio::task::spawn_blocking(move || environment.load_state(&session))
.await
.map_err(|error| format!("loading the pending review handoff stopped: {error}"))??;
let handle = control
.session(session_id.to_owned())
.await
.map_err(|error| format!("{error:#}"))?;
if let Err(error) = stop_leftover_review(&handle).await {
tracing::warn!(%session_id, %error, "could not stop the interrupted review's reviewers");
}
let Some(pending) = state.pending_forward.clone() else {
return Ok(None);
};
let view = handle.view();
if !view.connected {
return Err("the primary session is not connected".to_owned());
}
let Some(snapshot) = view.snapshot else {
return Err("the primary session has no transcript yet".to_owned());
};
if !matches!(
snapshot.materialized.execution,
MaterializedExecutionState::Idle
) {
return Err("the primary session is still working".to_owned());
}
if !snapshot.materialized.queued_prompts.is_empty() {
return Err("prompts are queued; the pending handoff waits for them".to_owned());
}
Ok(Some(Prepared {
state,
reviewer: ReviewerIdentity::default(),
materialized: Box::new(snapshot.materialized),
resume_forward: Some(pending),
captured: None,
background: None,
}))
}
fn busy_reviewer(command_id: &str) -> String {
if command_id.starts_with(mj_core::second_opinion::COMMAND_ID_PREFIX) {
"the reviewer is busy with a second opinion".to_owned()
} else {
"the reviewer is busy".to_owned()
}
}
fn is_turn_review_command(command_id: &str) -> bool {
command_id.starts_with(mj_core::review::driver::COMMAND_ID_PREFIX)
}
const LEGACY_ROLES: [&str; 9] = [
"validator",
"intent",
"supervisor",
"control_flow",
"duplication",
"error_handling",
"dead_code",
"tests",
"contracts",
];
pub(super) async fn stop_leftover_review(handle: &ManagedSessionHandle) -> Result<(), String> {
use mj_core::review::driver::REVIEWER_ROLE;
let mut roles = LEGACY_ROLES.to_vec();
let status = handle
.reviewer_as(Some(REVIEWER_ROLE.to_owned()), ReviewerAction::Status)
.await
.map_err(|error| format!("{error:#}"))?;
if let ReviewerOutcome::Status(state) = status
&& state
.active_prompt
.as_ref()
.is_some_and(|prompt| is_turn_review_command(&prompt.command_id))
{
roles.insert(0, REVIEWER_ROLE);
}
let mut failures = Vec::new();
for role in roles {
if let Err(error) = handle
.reviewer_as(Some(role.to_owned()), ReviewerAction::Pause)
.await
{
if error.to_string().contains("no longer supported") {
continue;
}
failures.push(format!("{role}: {error:#}"));
}
}
if failures.is_empty() {
Ok(())
} else {
Err(failures.join("; "))
}
}
pub(super) fn quoted_question(message: &str) -> String {
const LIMIT: usize = 160;
let line = message.lines().next().unwrap_or_default().trim();
if line.chars().count() > LIMIT {
format!("\"{}…\"", line.chars().take(LIMIT).collect::<String>())
} else {
format!("\"{line}\"")
}
}
pub(super) const BACKGROUND_WORK_WAIT: std::time::Duration = std::time::Duration::from_secs(120);
const LEASE_RETRY_PAUSE: std::time::Duration = std::time::Duration::from_millis(500);
pub(super) use crate::review_selection::preempted_by_lifecycle;
async fn after_background_work<T, Step, Attempt>(
session_id: &str,
deadline: tokio::time::Instant,
cancelled: &std::sync::atomic::AtomicBool,
mut step: Step,
) -> Result<T, String>
where
Step: FnMut() -> Attempt,
Attempt: std::future::Future<Output = Result<T, String>>,
{
let mut attempt = 0_u32;
loop {
if attempt > 0 {
tokio::time::sleep(LEASE_RETRY_PAUSE).await;
}
attempt += 1;
match step().await {
Err(reason)
if preempted_by_lifecycle(&reason)
&& tokio::time::Instant::now() < deadline
&& !cancelled.load(std::sync::atomic::Ordering::Acquire) =>
{
tracing::info!(
%session_id,
attempt,
%reason,
"turn review waits for another operation on the session"
);
}
outcome => return outcome,
}
}
}
pub(super) const SUBAGENT_REFUSAL: &str =
"a sub-agent's changes are reviewed with its parent's turn";
#[must_use]
pub fn start_refusal_notice(reason: &str) -> String {
let reason = if preempted_by_lifecycle(reason) {
"another operation was using the session"
} else {
reason
};
format!("Turn review did not start: {reason}. The next review covers these changes.")
}
#[must_use]
pub fn resolution_notice(
phase: &TurnReviewPhase,
last_verdict: Option<&ReviewVerdict>,
) -> Option<String> {
let TurnReviewPhase::Resolved(resolution) = phase else {
return None;
};
Some(match resolution {
Resolution::Forwarded => "Review findings sent to the agent".to_owned(),
Resolution::Dismissed => match last_verdict {
Some(ReviewVerdict::Clean) => "Review complete: no material findings".to_owned(),
Some(ReviewVerdict::Failed { .. }) => {
"Review failed; the change stays unreviewed".to_owned()
}
_ => "Review dismissed".to_owned(),
},
Resolution::Cancelled => match last_verdict {
Some(ReviewVerdict::Failed { .. }) => {
"Review failed; the change stays unreviewed".to_owned()
}
_ => "Review cancelled".to_owned(),
},
Resolution::NothingToReview => "Nothing to review: the turn changed no files".to_owned(),
Resolution::CoverageStarted => {
"Review coverage starts here; the next completed turn is reviewed".to_owned()
}
})
}
pub(super) fn seed_from_session(
session: &MaterializedSession,
state: &TurnReviewState,
_trigger: &str,
) -> TurnReviewSeed {
let mut task = String::new();
let mut user_messages = Vec::new();
let context_start = session
.transcript
.iter()
.filter(|item| mj_core::archive::is_context_boundary(&item.stable_id))
.map(|item| item.position)
.max()
.unwrap_or(0);
for item in session
.transcript
.iter()
.filter(|item| context_start == 0 || item.position > context_start)
{
let mj_core::state::TranscriptBody::User { content } = &item.body else {
continue;
};
let text = mj_core::transcript::materialized_content_text(content);
let text = text.trim();
if text.is_empty() {
continue;
}
if mj_core::second_opinion::is_control_origin_prompt(text)
|| mj_core::continuation::is_generated_prompt(
item.stable_id
.strip_prefix("user:")
.unwrap_or(&item.stable_id),
)
{
continue;
}
task = text.to_owned();
user_messages.push(UserMessage::prompt(text));
}
TurnReviewSeed {
task,
user_messages,
baselines: state.baselines.clone(),
through_ordinal: session.applied_event_ordinal,
prior_review: state.prior_review.clone(),
}
}