Skip to main content

mj_controller/database/
reviews.rs

1use super::*;
2
3/// Replace one host's remembered mount sources with exactly this list, so the
4/// dashboard can forget a directory the user no longer wants suggested.
5/// What each workspace last chose for a second opinion.
6///
7/// The selection is remembered so a repeat review does not ask again, and it
8/// is workspace scoped because a reviewer that suits one project rarely suits
9/// the next. Values are validated against what the harness advertises now
10/// before they are used, so a retired profile is harmless here.
11pub fn reviewer_defaults() -> Result<mj_core::second_opinion::ReviewerDefaults> {
12    reviewer_defaults_in(&database_path())
13}
14
15pub(super) fn reviewer_defaults_in(
16    path: &Path,
17) -> Result<mj_core::second_opinion::ReviewerDefaults> {
18    let connection = open_reader(path)?;
19    let mut statement = connection.prepare(
20        "SELECT workspace_id, profile_id, model, effort FROM second_opinion_defaults
21         ORDER BY workspace_id, profile_id, model",
22    )?;
23    let mut defaults = mj_core::second_opinion::ReviewerDefaults::default();
24    let mut rows = statement.query([])?;
25    while let Some(row) = rows.next()? {
26        let workspace_id: String = row.get(0)?;
27        let profile_id: String = row.get(1)?;
28        let model: String = row.get(2)?;
29        let effort: String = row.get(3)?;
30        defaults.restore(&workspace_id, &profile_id, &model, &effort);
31    }
32    Ok(defaults)
33}
34
35/// Record one confirmed selection.
36pub fn remember_reviewer_selection(
37    workspace_id: &str,
38    selection: &mj_core::second_opinion::ReviewerSelection,
39) -> Result<()> {
40    let workspace_id = workspace_id.to_owned();
41    let selection = selection.clone();
42    submit_database_write("remember_reviewer_selection", move |_| {
43        remember_reviewer_selection_in(&database_path(), &workspace_id, &selection)
44    })
45}
46
47pub(super) fn remember_reviewer_selection_in(
48    path: &Path,
49    workspace_id: &str,
50    selection: &mj_core::second_opinion::ReviewerSelection,
51) -> Result<()> {
52    ensure!(
53        !workspace_id.trim().is_empty(),
54        "second-opinion defaults need a workspace"
55    );
56    let mut connection = open(path)?;
57    let tx = connection.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
58    let (profile_id, model, effort) = selection.stored_values();
59    // One profile is the workspace's reviewer at a time, so the rows for the
60    // others stop being the remembered choice rather than accumulating.
61    tx.execute(
62        "DELETE FROM second_opinion_defaults WHERE workspace_id = ?1 AND profile_id <> ?2",
63        params![workspace_id, profile_id],
64    )?;
65    tx.execute(
66        "INSERT INTO second_opinion_defaults(workspace_id, profile_id, model, effort)
67         VALUES (?1, ?2, ?3, ?4)
68         ON CONFLICT(workspace_id, profile_id, model) DO UPDATE SET effort = excluded.effort",
69        params![workspace_id, profile_id, model, effort],
70    )?;
71    tx.commit()?;
72    Ok(())
73}
74
75/// The open review for `session_id`, if the session has one.
76pub fn active_review(session_id: &str) -> Result<Option<StoredReview>> {
77    active_review_in(&database_path(), session_id)
78}
79
80pub(super) fn active_review_in(path: &Path, session_id: &str) -> Result<Option<StoredReview>> {
81    let connection = open_reader(path)?;
82    let row = connection
83        .query_row(
84            "SELECT workflow, generation, context_baseline, native_lost, reviewer_transcript
85             FROM second_opinion_reviews WHERE session_id = ?1",
86            [session_id],
87            |row| {
88                Ok((
89                    row.get::<_, String>(0)?,
90                    row.get::<_, i64>(1)?,
91                    row.get::<_, i64>(2)?,
92                    row.get::<_, i64>(3)?,
93                    row.get::<_, String>(4)?,
94                ))
95            },
96        )
97        .optional()?;
98    let Some((workflow, generation, baseline, native_lost, transcript)) = row else {
99        return Ok(None);
100    };
101    Ok(Some(StoredReview {
102        workflow: serde_json::from_str(&workflow).context("parse the stored review workflow")?,
103        generation: u64::try_from(generation).unwrap_or_default(),
104        context_baseline: u64::try_from(baseline).unwrap_or_default(),
105        native_lost: native_lost != 0,
106        reviewer_transcript: serde_json::from_str(&transcript)
107            .context("parse the stored reviewer transcript")?,
108    }))
109}
110
111/// Records the open review, replacing any earlier one for this session.
112pub fn save_active_review(session_id: &str, review: &StoredReview) -> Result<()> {
113    let session_id = session_id.to_owned();
114    let review = review.clone();
115    submit_database_write("save_active_review", move |_| {
116        save_active_review_in(&database_path(), &session_id, &review)
117    })
118}
119
120pub(super) fn save_active_review_in(
121    path: &Path,
122    session_id: &str,
123    review: &StoredReview,
124) -> Result<()> {
125    let connection = open(path)?;
126    connection.execute(
127        "INSERT INTO second_opinion_reviews(
128             session_id, workflow, generation, context_baseline, native_lost,
129             reviewer_transcript
130         ) VALUES (?1, ?2, ?3, ?4, ?5, ?6)
131         ON CONFLICT(session_id) DO UPDATE SET
132             workflow = excluded.workflow,
133             generation = excluded.generation,
134             context_baseline = excluded.context_baseline,
135             native_lost = excluded.native_lost,
136             reviewer_transcript = excluded.reviewer_transcript",
137        params![
138            session_id,
139            serde_json::to_string(&review.workflow)?,
140            i64::try_from(review.generation).unwrap_or(i64::MAX),
141            i64::try_from(review.context_baseline).unwrap_or(i64::MAX),
142            i64::from(review.native_lost),
143            serde_json::to_string(&review.reviewer_transcript)?,
144        ],
145    )?;
146    Ok(())
147}
148
149/// Forgets the open review once it has finished.
150pub fn clear_active_review(session_id: &str) -> Result<()> {
151    let session_id = session_id.to_owned();
152    submit_database_write("clear_active_review", move |_| {
153        clear_active_review_in(&database_path(), &session_id)
154    })
155}
156
157pub(super) fn clear_active_review_in(path: &Path, session_id: &str) -> Result<()> {
158    let connection = open(path)?;
159    connection.execute(
160        "DELETE FROM second_opinion_reviews WHERE session_id = ?1",
161        [session_id],
162    )?;
163    Ok(())
164}
165
166/// How far `session_id` has been reviewed, or a fresh state when it has never
167/// been reviewed.
168pub fn turn_review_state(session_id: &str) -> Result<TurnReviewState> {
169    turn_review_state_in(&database_path(), session_id)
170}
171
172pub(super) fn turn_review_state_in(path: &Path, session_id: &str) -> Result<TurnReviewState> {
173    let connection = open_reader(path)?;
174    let row = connection
175        .query_row(
176            "SELECT baselines, reviewed_through_ordinal, prior_review, active,
177                    pending_forward
178             FROM turn_review_state WHERE session_id = ?1",
179            [session_id],
180            |row| {
181                Ok((
182                    row.get::<_, String>(0)?,
183                    row.get::<_, i64>(1)?,
184                    row.get::<_, Option<String>>(2)?,
185                    row.get::<_, Option<String>>(3)?,
186                    row.get::<_, Option<String>>(4)?,
187                ))
188            },
189        )
190        .optional()?;
191    let Some((baselines, ordinal, prior, active, pending_forward)) = row else {
192        return Ok(TurnReviewState::default());
193    };
194    Ok(TurnReviewState {
195        baselines: serde_json::from_str(&baselines).context("parse the stored review baselines")?,
196        reviewed_through_ordinal: u64::try_from(ordinal).unwrap_or_default(),
197        prior_review: prior
198            .map(|prior| serde_json::from_str(&prior))
199            .transpose()
200            .context("parse the stored prior review")?,
201        active,
202        pending_forward: pending_forward
203            .map(|pending| serde_json::from_str(&pending))
204            .transpose()
205            .context("parse the stored pending review handoff")?,
206    })
207}
208
209/// Records how far a session has been reviewed.
210pub fn save_turn_review_state(session_id: &str, state: &TurnReviewState) -> Result<()> {
211    let session_id = session_id.to_owned();
212    let state = state.clone();
213    submit_database_write("save_turn_review_state", move |_| {
214        save_turn_review_state_in(&database_path(), &session_id, &state)
215    })
216}
217
218pub(super) fn save_turn_review_state_in(
219    path: &Path,
220    session_id: &str,
221    state: &TurnReviewState,
222) -> Result<()> {
223    let connection = open(path)?;
224    connection.execute(
225        "INSERT INTO turn_review_state(
226             session_id, baselines, reviewed_through_ordinal, prior_review, active,
227             pending_forward
228         ) VALUES (?1, ?2, ?3, ?4, ?5, ?6)
229         ON CONFLICT(session_id) DO UPDATE SET
230             baselines = excluded.baselines,
231             reviewed_through_ordinal = excluded.reviewed_through_ordinal,
232             prior_review = excluded.prior_review,
233             active = excluded.active,
234             pending_forward = excluded.pending_forward",
235        params![
236            session_id,
237            serde_json::to_string(&state.baselines)?,
238            i64::try_from(state.reviewed_through_ordinal).unwrap_or(i64::MAX),
239            state
240                .prior_review
241                .as_ref()
242                .map(serde_json::to_string)
243                .transpose()?,
244            state.active,
245            state
246                .pending_forward
247                .as_ref()
248                .map(serde_json::to_string)
249                .transpose()?,
250        ],
251    )?;
252    Ok(())
253}
254
255/// Clears every session's in-flight review flag.
256///
257/// A review that was running when the daemon stopped is not resumed: the
258/// baseline never advanced, so the next review covers the same change, and
259/// half a multi-agent fan-out is not worth rebuilding. A pending corrective
260/// handoff is returned as well so the host can retry its exact command id.
261/// Baselines are deliberately left alone, which is what makes interruption
262/// lossless. Returns the sessions whose review or handoff was interrupted.
263pub fn clear_interrupted_turn_reviews() -> Result<Vec<String>> {
264    submit_database_write("clear_interrupted_turn_reviews", move |_| {
265        clear_interrupted_turn_reviews_in(&database_path())
266    })
267}
268
269pub(super) fn clear_interrupted_turn_reviews_in(path: &Path) -> Result<Vec<String>> {
270    let connection = open(path)?;
271    let interrupted = {
272        let mut statement = connection.prepare(
273            "SELECT session_id FROM turn_review_state
274                 WHERE active IS NOT NULL OR pending_forward IS NOT NULL",
275        )?;
276        let mut rows = statement.query([])?;
277        let mut interrupted = Vec::new();
278        while let Some(row) = rows.next()? {
279            interrupted.push(row.get::<_, String>(0)?);
280        }
281        interrupted
282    };
283    connection.execute(
284        "UPDATE turn_review_state SET active = NULL WHERE active IS NOT NULL",
285        [],
286    )?;
287    Ok(interrupted)
288}
289
290/// Marks this session's reviewer conversation as no longer continuable, and
291/// reports the generation a future review must start under.
292///
293/// Losing the target takes the reviewer's native session with it. The
294/// materialized transcript is kept for reference, but the next review is a new
295/// conversation, so it runs under a new generation.
296pub fn lose_reviewer_continuity(session_id: &str) -> Result<u64> {
297    let session_id = session_id.to_owned();
298    submit_database_write("lose_reviewer_continuity", move |_| {
299        lose_reviewer_continuity_in(&database_path(), &session_id)
300    })
301}
302
303pub(super) fn lose_reviewer_continuity_in(path: &Path, session_id: &str) -> Result<u64> {
304    let Some(mut review) = active_review_in(path, session_id)? else {
305        return Ok(0);
306    };
307    if review.native_lost {
308        return Ok(review.generation);
309    }
310    review.native_lost = true;
311    review.generation = review.generation.saturating_add(1);
312    save_active_review_in(path, session_id, &review)?;
313    Ok(review.generation)
314}