Skip to main content

mj_controller/review_host/
environment.rs

1use super::*;
2
3/// Everything a review needs from the controller: whether it can review this
4/// session at all, and a staged reviewer profile to launch a role from.
5///
6/// It is a trait so the host's own tests can drive a whole review without a
7/// container, a harness, or the developer's own `config.toml`. The daemon
8/// installs [`ControllerEnvironment`], which loads the real controller.
9pub trait ReviewEnvironment: Send + Sync {
10    /// Refuses, with a sentence for a person, when this session cannot be
11    /// reviewed under `profile`.
12    fn check(&self, session_id: &str, profile: &str) -> Result<(), String>;
13
14    /// Whether this session is a Mjolnir-managed sub-agent. Its changes are
15    /// reviewed through the parent's turn, so it is never reviewed on its
16    /// own. Blocking: it reads the controller's database.
17    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    /// Stages the reviewer profile for one role and describes how to launch
32    /// it. Blocking: it copies a profile onto the session's target.
33    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    /// How far this session has been reviewed. Blocking: it reads the
43    /// controller's database.
44    fn load_state(&self, session_id: &str) -> Result<TurnReviewState, String>;
45
46    /// Records how far this session has been reviewed. Blocking: the host
47    /// routes it through its ordered persistence lane rather than calling it
48    /// on the Tokio task that owns review state.
49    fn save_state(&self, session_id: &str, state: &TurnReviewState) -> Result<(), String>;
50
51    /// Clears the in-flight flag of every review a restart interrupted, and
52    /// reports whose they were. Baselines are deliberately left alone: the
53    /// interrupted review never advanced one, so the next review covers the
54    /// same change and nothing is lost.
55    fn clear_interrupted(&self) -> Result<Vec<String>, String>;
56
57    /// Holds the session against background work for as long as a review
58    /// needs its worker: waits until no recovery copy or worker upgrade is
59    /// running (up to `deadline`), then keeps new ones from starting until
60    /// the returned hold is dropped.
61    ///
62    /// A recovery copy takes the session's lease, which cancels reviewer
63    /// actions in flight (R4-9), and a checkpoint or worker replacement makes
64    /// the worker refuse reviewer work. Both are started by the same finished
65    /// turn that starts the review, so waiting for them to settle was not
66    /// enough: one could start a moment later, while the review chose its
67    /// reviewer (Series 34 #3444, 2026-10-03). The hold closes that window.
68    /// An environment with no background work holds nothing.
69    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
78/// Keeps background work off a session while a review needs its worker.
79/// Dropping it lets recovery copies and worker upgrades start again.
80pub type BackgroundHold = Box<dyn Send + Sync>;
81
82/// The production environment: the controller as it is on disk right now.
83///
84/// It is reloaded per call rather than held, because a review is rare and the
85/// answer must reflect the config as it stands when the review starts -- the
86/// daemon reloads config.toml every 500 ms for the same reason.
87#[derive(Default)]
88pub struct ControllerEnvironment {
89    /// The daemon's gate for recovery copies and worker upgrades, which a
90    /// review waits behind instead of racing for the session's lease.
91    pub background: Option<Arc<crate::recovery_gate::RecoveryGate>>,
92    /// The daemon's profile catalog, once it is installed. A pinned review
93    /// model is matched against it before any profile is staged.
94    pub(crate) catalog: Option<SharedProfileCatalog>,
95}
96
97/// The daemon installs its profile catalog after the review host starts.
98pub(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            // A deadline that passes leaves the review to try anyway; the
190            // running work's own refusal then says what held the session.
191            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}