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    ) -> Result<mj_core::worker_launch::ReviewerLaunchConfig, String>;
39
40    /// How far this session has been reviewed. Blocking: it reads the
41    /// controller's database.
42    fn load_state(&self, session_id: &str) -> Result<TurnReviewState, String>;
43
44    /// Records how far this session has been reviewed. Blocking: the host
45    /// routes it through its ordered persistence lane rather than calling it
46    /// on the Tokio task that owns review state.
47    fn save_state(&self, session_id: &str, state: &TurnReviewState) -> Result<(), String>;
48
49    /// Clears the in-flight flag of every review a restart interrupted, and
50    /// reports whose they were. Baselines are deliberately left alone: the
51    /// interrupted review never advanced one, so the next review covers the
52    /// same change and nothing is lost.
53    fn clear_interrupted(&self) -> Result<Vec<String>, String>;
54
55    /// Holds the session against background work for as long as a review
56    /// needs its worker: waits until no recovery copy or worker upgrade is
57    /// running (up to `deadline`), then keeps new ones from starting until
58    /// the returned hold is dropped.
59    ///
60    /// A recovery copy takes the session's lease, which cancels reviewer
61    /// actions in flight (R4-9), and a checkpoint or worker replacement makes
62    /// the worker refuse reviewer work. Both are started by the same finished
63    /// turn that starts the review, so waiting for them to settle was not
64    /// enough: one could start a moment later, while the review chose its
65    /// reviewer (Series 34 #3444, 2026-10-03). The hold closes that window.
66    /// An environment with no background work holds nothing.
67    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
76/// Keeps background work off a session while a review needs its worker.
77/// Dropping it lets recovery copies and worker upgrades start again.
78pub type BackgroundHold = Box<dyn Send + Sync>;
79
80/// The production environment: the controller as it is on disk right now.
81///
82/// It is reloaded per call rather than held, because a review is rare and the
83/// answer must reflect the config as it stands when the review starts -- the
84/// daemon reloads config.toml every 500 ms for the same reason.
85#[derive(Default)]
86pub struct ControllerEnvironment {
87    /// The daemon's gate for recovery copies and worker upgrades, which a
88    /// review waits behind instead of racing for the session's lease.
89    pub background: Option<Arc<crate::recovery_gate::RecoveryGate>>,
90    /// The daemon's profile catalog, once it is installed. A pinned review
91    /// model is matched against it before any profile is staged.
92    pub(crate) catalog: Option<SharedProfileCatalog>,
93}
94
95/// The daemon installs its profile catalog after the review host starts.
96pub(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            // A deadline that passes leaves the review to try anyway; the
179            // running work's own refusal then says what held the session.
180            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}