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    /// Waits until no background work holds this session, or until
58    /// `deadline`. The automatic recovery copy a finished turn starts takes
59    /// the session's lease, and a lease cancels reviewer actions in flight
60    /// (R4-9). An environment with no background work returns at once.
61    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/// The production environment: the controller as it is on disk right now.
71///
72/// It is reloaded per call rather than held, because a review is rare and the
73/// answer must reflect the config as it stands when the review starts -- the
74/// daemon reloads config.toml every 500 ms for the same reason.
75#[derive(Default)]
76pub struct ControllerEnvironment {
77    /// The daemon's gate for recovery copies and worker upgrades, which a
78    /// review waits behind instead of racing for the session's lease.
79    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            // A deadline that passes leaves the review to try anyway; its
169            // own refusal then says what held the session.
170            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}