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,
) -> Result<mj_core::worker_launch::ReviewerLaunchConfig, String>;
fn load_state(&self, session_id: &str) -> Result<TurnReviewState, String>;
fn save_state(&self, session_id: &str, state: &TurnReviewState) -> Result<(), String>;
fn clear_interrupted(&self) -> Result<Vec<String>, String>;
fn hold_background_work<'a>(
&'a self,
_session_id: &'a str,
_deadline: tokio::time::Instant,
) -> mj_client::session::BoxFuture<'a, Option<BackgroundHold>> {
Box::pin(async { None })
}
}
pub type BackgroundHold = Box<dyn Send + Sync>;
#[derive(Default)]
pub struct ControllerEnvironment {
pub background: Option<Arc<crate::recovery_gate::RecoveryGate>>,
pub(crate) catalog: Option<SharedProfileCatalog>,
}
pub(crate) type SharedProfileCatalog =
Arc<std::sync::OnceLock<Arc<crate::server_runtime::profile_catalog::ProfileCatalog>>>;
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 offered = self
.catalog
.as_ref()
.and_then(|catalog| catalog.get())
.map(|catalog| catalog.snapshot());
crate::review_selection::resolve(handle, Some(config), cancelled, offered)
.await
.map_err(|e| format!("{e:#}"))
})
}
fn stage(
&self,
session_id: &str,
profile: &str,
generation: u64,
) -> Result<mj_core::worker_launch::ReviewerLaunchConfig, String> {
let controller =
crate::controller::Controller::load().map_err(|error| format!("{error:#}"))?;
controller
.stage_reviewer_profile(session_id, profile, generation)
.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 clear_interrupted(&self) -> Result<Vec<String>, String> {
crate::database::clear_interrupted_turn_reviews().map_err(|error| format!("{error:#}"))
}
fn hold_background_work<'a>(
&'a self,
session_id: &'a str,
deadline: tokio::time::Instant,
) -> mj_client::session::BoxFuture<'a, Option<BackgroundHold>> {
Box::pin(async move {
let gate = self.background.as_ref()?;
let hold: BackgroundHold = Box::new(gate.reserve_when_idle(session_id, deadline).await);
Some(hold)
})
}
}
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(())
}