mj_controller/review_host/
environment.rs1use super::*;
2
3pub trait ReviewEnvironment: Send + Sync {
10 fn check(&self, session_id: &str, profile: &str) -> Result<(), String>;
13
14 fn is_subagent(&self, _session_id: &str) -> bool {
18 false
19 }
20
21 fn resolve<'a>(
22 &'a self,
23 handle: ManagedSessionHandle,
24 config: ReviewConfig,
25 cancelled: Arc<std::sync::atomic::AtomicBool>,
26 ) -> mj_client::session::BoxFuture<
27 'a,
28 Result<mj_core::review::settings::ResolvedReviewSettings, String>,
29 >;
30
31 fn stage(
34 &self,
35 session_id: &str,
36 profile: &str,
37 generation: u64,
38 mcp_servers: &[mj_core::worker_launch::ReviewMcpServer],
39 dispatch_tool: bool,
40 ) -> Result<mj_core::worker_launch::ReviewerLaunchConfig, String>;
41
42 fn load_state(&self, session_id: &str) -> Result<TurnReviewState, String>;
45
46 fn save_state(&self, session_id: &str, state: &TurnReviewState) -> Result<(), String>;
50
51 fn clear_interrupted(&self) -> Result<Vec<String>, String>;
56
57 fn hold_background_work<'a>(
70 &'a self,
71 _session_id: &'a str,
72 _deadline: tokio::time::Instant,
73 ) -> mj_client::session::BoxFuture<'a, Option<BackgroundHold>> {
74 Box::pin(async { None })
75 }
76}
77
78pub type BackgroundHold = Box<dyn Send + Sync>;
81
82#[derive(Default)]
88pub struct ControllerEnvironment {
89 pub background: Option<Arc<crate::recovery_gate::RecoveryGate>>,
92 pub(crate) catalog: Option<SharedProfileCatalog>,
95}
96
97pub(crate) type SharedProfileCatalog =
99 Arc<std::sync::OnceLock<Arc<crate::server_runtime::profile_catalog::ProfileCatalog>>>;
100
101impl ReviewEnvironment for ControllerEnvironment {
102 fn check(&self, session_id: &str, profile: &str) -> Result<(), String> {
103 let controller =
104 crate::controller::Controller::load().map_err(|error| format!("{error:#}"))?;
105 let Some(reviewer) = controller.config.profiles.get(profile) else {
106 return Err(format!(
107 "turn review needs a reviewer: [review] profile {profile:?} is not a profile in config.toml"
108 ));
109 };
110 if !reviewer.enabled {
111 return Err(format!(
112 "turn review needs an enabled reviewer: [review] profile {profile:?} is disabled"
113 ));
114 }
115 validate_reviewer_assignment(
116 session_id,
117 controller.state.sessions.get(session_id),
118 profile,
119 )
120 }
121
122 fn is_subagent(&self, session_id: &str) -> bool {
123 crate::database::load_subagent(session_id).is_ok_and(|record| record.is_some())
124 }
125
126 fn resolve<'a>(
127 &'a self,
128 handle: ManagedSessionHandle,
129 config: ReviewConfig,
130 cancelled: Arc<std::sync::atomic::AtomicBool>,
131 ) -> mj_client::session::BoxFuture<
132 'a,
133 Result<mj_core::review::settings::ResolvedReviewSettings, String>,
134 > {
135 Box::pin(async move {
136 let specialists = config.tier == ReviewTier::Extended;
137 let offered = self
138 .catalog
139 .as_ref()
140 .and_then(|catalog| catalog.get())
141 .map(|catalog| catalog.snapshot());
142 crate::review_selection::resolve(handle, Some(config), specialists, cancelled, offered)
143 .await
144 .map_err(|e| format!("{e:#}"))
145 })
146 }
147
148 fn stage(
149 &self,
150 session_id: &str,
151 profile: &str,
152 generation: u64,
153 mcp_servers: &[mj_core::worker_launch::ReviewMcpServer],
154 dispatch_tool: bool,
155 ) -> Result<mj_core::worker_launch::ReviewerLaunchConfig, String> {
156 let controller =
157 crate::controller::Controller::load().map_err(|error| format!("{error:#}"))?;
158 controller
159 .stage_reviewer_profile_with_mcp(
160 session_id,
161 profile,
162 generation,
163 mcp_servers,
164 dispatch_tool,
165 )
166 .map_err(|error| format!("{error:#}"))
167 }
168
169 fn load_state(&self, session_id: &str) -> Result<TurnReviewState, String> {
170 crate::database::turn_review_state(session_id).map_err(|error| format!("{error:#}"))
171 }
172
173 fn save_state(&self, session_id: &str, state: &TurnReviewState) -> Result<(), String> {
174 crate::database::save_turn_review_state(session_id, state)
175 .map_err(|error| format!("{error:#}"))
176 }
177
178 fn clear_interrupted(&self) -> Result<Vec<String>, String> {
179 crate::database::clear_interrupted_turn_reviews().map_err(|error| format!("{error:#}"))
180 }
181
182 fn hold_background_work<'a>(
183 &'a self,
184 session_id: &'a str,
185 deadline: tokio::time::Instant,
186 ) -> mj_client::session::BoxFuture<'a, Option<BackgroundHold>> {
187 Box::pin(async move {
188 let gate = self.background.as_ref()?;
189 let hold: BackgroundHold = Box::new(gate.reserve_when_idle(session_id, deadline).await);
192 Some(hold)
193 })
194 }
195}
196
197pub(crate) fn validate_reviewer_assignment(
198 session_id: &str,
199 session: Option<&mj_core::state::SessionRecord>,
200 _profile: &str,
201) -> Result<(), String> {
202 let Some(session) = session else {
203 return Err(format!(
204 "session {session_id:?} is not in the controller store"
205 ));
206 };
207 if session.archived {
208 return Err("this session is archived".to_owned());
209 }
210 Ok(())
211}