use super::*;
pub trait ReviewEnvironment: Send + Sync {
fn check(&self, session_id: &str, profile: &str) -> Result<(), String>;
fn is_subagent(&self, _session_id: &str) -> bool {
false
}
fn resolve<'a>(
&'a self,
handle: ManagedSessionHandle,
config: ReviewConfig,
cancelled: Arc<std::sync::atomic::AtomicBool>,
) -> mj_client::session::BoxFuture<
'a,
Result<mj_core::review::settings::ResolvedReviewSettings, String>,
>;
fn stage(
&self,
session_id: &str,
profile: &str,
generation: u64,
mcp_servers: &[mj_core::worker_launch::ReviewMcpServer],
dispatch_tool: bool,
) -> Result<mj_core::worker_launch::ReviewerLaunchConfig, String>;
fn load_state(&self, session_id: &str) -> Result<TurnReviewState, String>;
fn primary_prompt_accepted(
&self,
_session_id: &str,
_command_id: &str,
) -> Result<bool, String> {
Ok(false)
}
fn save_state(&self, session_id: &str, state: &TurnReviewState) -> Result<(), String>;
fn recoverable_reviews(&self) -> Result<Vec<String>, String>;
fn background_work_settled<'a>(
&'a self,
_session_id: &'a str,
_deadline: tokio::time::Instant,
) -> mj_client::session::BoxFuture<'a, ()> {
Box::pin(async {})
}
}
#[derive(Default)]
pub struct ControllerEnvironment {
pub background: Option<Arc<crate::recovery_gate::RecoveryGate>>,
}
impl ReviewEnvironment for ControllerEnvironment {
fn check(&self, session_id: &str, profile: &str) -> Result<(), String> {
let controller =
crate::controller::Controller::load().map_err(|error| format!("{error:#}"))?;
let Some(reviewer) = controller.config.profiles.get(profile) else {
return Err(format!(
"turn review needs a reviewer: [review] profile {profile:?} is not a profile in config.toml"
));
};
if !reviewer.enabled {
return Err(format!(
"turn review needs an enabled reviewer: [review] profile {profile:?} is disabled"
));
}
validate_reviewer_assignment(
session_id,
controller.state.sessions.get(session_id),
profile,
)
}
fn is_subagent(&self, session_id: &str) -> bool {
crate::database::load_subagent(session_id).is_ok_and(|record| record.is_some())
}
fn resolve<'a>(
&'a self,
handle: ManagedSessionHandle,
config: ReviewConfig,
cancelled: Arc<std::sync::atomic::AtomicBool>,
) -> mj_client::session::BoxFuture<
'a,
Result<mj_core::review::settings::ResolvedReviewSettings, String>,
> {
Box::pin(async move {
let specialists = config.tier == ReviewTier::Extended;
crate::review_selection::resolve(handle, Some(config), specialists, cancelled)
.await
.map_err(|e| format!("{e:#}"))
})
}
fn stage(
&self,
session_id: &str,
profile: &str,
generation: u64,
mcp_servers: &[mj_core::worker_launch::ReviewMcpServer],
dispatch_tool: bool,
) -> Result<mj_core::worker_launch::ReviewerLaunchConfig, String> {
let controller =
crate::controller::Controller::load().map_err(|error| format!("{error:#}"))?;
controller
.stage_reviewer_profile_with_mcp(
session_id,
profile,
generation,
mcp_servers,
dispatch_tool,
)
.map_err(|error| format!("{error:#}"))
}
fn load_state(&self, session_id: &str) -> Result<TurnReviewState, String> {
crate::database::turn_review_state(session_id).map_err(|error| format!("{error:#}"))
}
fn save_state(&self, session_id: &str, state: &TurnReviewState) -> Result<(), String> {
crate::database::save_turn_review_state(session_id, state)
.map_err(|error| format!("{error:#}"))
}
fn primary_prompt_accepted(&self, session_id: &str, command_id: &str) -> Result<bool, String> {
crate::database::load_prompt_acceptance(session_id, command_id)
.map(|accepted| accepted.is_some())
.map_err(|error| format!("reconcile review handoff: {error:#}"))
}
fn recoverable_reviews(&self) -> Result<Vec<String>, String> {
crate::database::recoverable_turn_reviews().map_err(|error| format!("{error:#}"))
}
fn background_work_settled<'a>(
&'a self,
session_id: &'a str,
deadline: tokio::time::Instant,
) -> mj_client::session::BoxFuture<'a, ()> {
Box::pin(async move {
let Some(gate) = &self.background else {
return;
};
let mut busy = gate.subscribe();
let _ = tokio::time::timeout_at(
deadline,
busy.wait_for(|sessions| !sessions.contains(session_id)),
)
.await;
})
}
}
pub(crate) fn validate_reviewer_assignment(
session_id: &str,
session: Option<&mj_core::state::SessionRecord>,
_profile: &str,
) -> Result<(), String> {
let Some(session) = session else {
return Err(format!(
"session {session_id:?} is not in the controller store"
));
};
if session.archived {
return Err("this session is archived".to_owned());
}
Ok(())
}