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 ) -> Result<mj_core::worker_launch::ReviewerLaunchConfig, String>;
39
40 fn load_state(&self, session_id: &str) -> Result<TurnReviewState, String>;
43
44 fn save_state(&self, session_id: &str, state: &TurnReviewState) -> Result<(), String>;
48
49 fn clear_interrupted(&self) -> Result<Vec<String>, String>;
54
55 fn hold_background_work<'a>(
68 &'a self,
69 _session_id: &'a str,
70 _deadline: tokio::time::Instant,
71 ) -> mj_client::session::BoxFuture<'a, Option<BackgroundHold>> {
72 Box::pin(async { None })
73 }
74}
75
76pub type BackgroundHold = Box<dyn Send + Sync>;
79
80#[derive(Default)]
86pub struct ControllerEnvironment {
87 pub background: Option<Arc<crate::recovery_gate::RecoveryGate>>,
90 pub(crate) catalog: Option<SharedProfileCatalog>,
93}
94
95pub(crate) type SharedProfileCatalog =
97 Arc<std::sync::OnceLock<Arc<crate::server_runtime::profile_catalog::ProfileCatalog>>>;
98
99impl ReviewEnvironment for ControllerEnvironment {
100 fn check(&self, session_id: &str, profile: &str) -> Result<(), String> {
101 let controller =
102 crate::controller::Controller::load().map_err(|error| format!("{error:#}"))?;
103 let Some(reviewer) = controller.config.profiles.get(profile) else {
104 return Err(format!(
105 "turn review needs a reviewer: [review] profile {profile:?} is not a profile in config.toml"
106 ));
107 };
108 if !reviewer.enabled {
109 return Err(format!(
110 "turn review needs an enabled reviewer: [review] profile {profile:?} is disabled"
111 ));
112 }
113 validate_reviewer_assignment(
114 session_id,
115 controller.state.sessions.get(session_id),
116 profile,
117 )
118 }
119
120 fn is_subagent(&self, session_id: &str) -> bool {
121 crate::database::load_subagent(session_id).is_ok_and(|record| record.is_some())
122 }
123
124 fn resolve<'a>(
125 &'a self,
126 handle: ManagedSessionHandle,
127 config: ReviewConfig,
128 cancelled: Arc<std::sync::atomic::AtomicBool>,
129 ) -> mj_client::session::BoxFuture<
130 'a,
131 Result<mj_core::review::settings::ResolvedReviewSettings, String>,
132 > {
133 Box::pin(async move {
134 let offered = self
135 .catalog
136 .as_ref()
137 .and_then(|catalog| catalog.get())
138 .map(|catalog| catalog.snapshot());
139 crate::review_selection::resolve(handle, Some(config), cancelled, offered)
140 .await
141 .map_err(|e| format!("{e:#}"))
142 })
143 }
144
145 fn stage(
146 &self,
147 session_id: &str,
148 profile: &str,
149 generation: u64,
150 ) -> Result<mj_core::worker_launch::ReviewerLaunchConfig, String> {
151 let controller =
152 crate::controller::Controller::load().map_err(|error| format!("{error:#}"))?;
153 controller
154 .stage_reviewer_profile(session_id, profile, generation)
155 .map_err(|error| format!("{error:#}"))
156 }
157
158 fn load_state(&self, session_id: &str) -> Result<TurnReviewState, String> {
159 crate::database::turn_review_state(session_id).map_err(|error| format!("{error:#}"))
160 }
161
162 fn save_state(&self, session_id: &str, state: &TurnReviewState) -> Result<(), String> {
163 crate::database::save_turn_review_state(session_id, state)
164 .map_err(|error| format!("{error:#}"))
165 }
166
167 fn clear_interrupted(&self) -> Result<Vec<String>, String> {
168 crate::database::clear_interrupted_turn_reviews().map_err(|error| format!("{error:#}"))
169 }
170
171 fn hold_background_work<'a>(
172 &'a self,
173 session_id: &'a str,
174 deadline: tokio::time::Instant,
175 ) -> mj_client::session::BoxFuture<'a, Option<BackgroundHold>> {
176 Box::pin(async move {
177 let gate = self.background.as_ref()?;
178 let hold: BackgroundHold = Box::new(gate.reserve_when_idle(session_id, deadline).await);
181 Some(hold)
182 })
183 }
184}
185
186pub(crate) fn validate_reviewer_assignment(
187 session_id: &str,
188 session: Option<&mj_core::state::SessionRecord>,
189 _profile: &str,
190) -> Result<(), String> {
191 let Some(session) = session else {
192 return Err(format!(
193 "session {session_id:?} is not in the controller store"
194 ));
195 };
196 if session.archived {
197 return Err("this session is archived".to_owned());
198 }
199 Ok(())
200}