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 background_work_settled<'a>(
62 &'a self,
63 _session_id: &'a str,
64 _deadline: tokio::time::Instant,
65 ) -> mj_client::session::BoxFuture<'a, ()> {
66 Box::pin(async {})
67 }
68}
69
70#[derive(Default)]
76pub struct ControllerEnvironment {
77 pub background: Option<Arc<crate::recovery_gate::RecoveryGate>>,
80}
81
82impl ReviewEnvironment for ControllerEnvironment {
83 fn check(&self, session_id: &str, profile: &str) -> Result<(), String> {
84 let controller =
85 crate::controller::Controller::load().map_err(|error| format!("{error:#}"))?;
86 let Some(reviewer) = controller.config.profiles.get(profile) else {
87 return Err(format!(
88 "turn review needs a reviewer: [review] profile {profile:?} is not a profile in config.toml"
89 ));
90 };
91 if !reviewer.enabled {
92 return Err(format!(
93 "turn review needs an enabled reviewer: [review] profile {profile:?} is disabled"
94 ));
95 }
96 validate_reviewer_assignment(
97 session_id,
98 controller.state.sessions.get(session_id),
99 profile,
100 )
101 }
102
103 fn is_subagent(&self, session_id: &str) -> bool {
104 crate::database::load_subagent(session_id).is_ok_and(|record| record.is_some())
105 }
106
107 fn resolve<'a>(
108 &'a self,
109 handle: ManagedSessionHandle,
110 config: ReviewConfig,
111 cancelled: Arc<std::sync::atomic::AtomicBool>,
112 ) -> mj_client::session::BoxFuture<
113 'a,
114 Result<mj_core::review::settings::ResolvedReviewSettings, String>,
115 > {
116 Box::pin(async move {
117 let specialists = config.tier == ReviewTier::Extended;
118 crate::review_selection::resolve(handle, Some(config), specialists, cancelled)
119 .await
120 .map_err(|e| format!("{e:#}"))
121 })
122 }
123
124 fn stage(
125 &self,
126 session_id: &str,
127 profile: &str,
128 generation: u64,
129 mcp_servers: &[mj_core::worker_launch::ReviewMcpServer],
130 dispatch_tool: bool,
131 ) -> Result<mj_core::worker_launch::ReviewerLaunchConfig, String> {
132 let controller =
133 crate::controller::Controller::load().map_err(|error| format!("{error:#}"))?;
134 controller
135 .stage_reviewer_profile_with_mcp(
136 session_id,
137 profile,
138 generation,
139 mcp_servers,
140 dispatch_tool,
141 )
142 .map_err(|error| format!("{error:#}"))
143 }
144
145 fn load_state(&self, session_id: &str) -> Result<TurnReviewState, String> {
146 crate::database::turn_review_state(session_id).map_err(|error| format!("{error:#}"))
147 }
148
149 fn save_state(&self, session_id: &str, state: &TurnReviewState) -> Result<(), String> {
150 crate::database::save_turn_review_state(session_id, state)
151 .map_err(|error| format!("{error:#}"))
152 }
153
154 fn clear_interrupted(&self) -> Result<Vec<String>, String> {
155 crate::database::clear_interrupted_turn_reviews().map_err(|error| format!("{error:#}"))
156 }
157
158 fn background_work_settled<'a>(
159 &'a self,
160 session_id: &'a str,
161 deadline: tokio::time::Instant,
162 ) -> mj_client::session::BoxFuture<'a, ()> {
163 Box::pin(async move {
164 let Some(gate) = &self.background else {
165 return;
166 };
167 let mut busy = gate.subscribe();
168 let _ = tokio::time::timeout_at(
171 deadline,
172 busy.wait_for(|sessions| !sessions.contains(session_id)),
173 )
174 .await;
175 })
176 }
177}
178
179pub(crate) fn validate_reviewer_assignment(
180 session_id: &str,
181 session: Option<&mj_core::state::SessionRecord>,
182 _profile: &str,
183) -> Result<(), String> {
184 let Some(session) = session else {
185 return Err(format!(
186 "session {session_id:?} is not in the controller store"
187 ));
188 };
189 if session.archived {
190 return Err("this session is archived".to_owned());
191 }
192 Ok(())
193}