Skip to main content

mj_controller/database/
session_move.rs

1//! Move intent uses the same guarded writer as session lifecycle transitions.
2
3use super::*;
4use mj_core::state::MoveOperation;
5
6pub fn save_move_operation(operation: &MoveOperation) -> Result<()> {
7    let operation = operation.clone();
8    submit_database_write("save_move_operation", move |connection| {
9        save_move_operation_with(connection, &operation)
10    })
11}
12
13pub(super) fn save_move_operation_with(
14    connection: &Connection,
15    operation: &MoveOperation,
16) -> Result<()> {
17    connection.execute(
18        "INSERT INTO session_moves(session_id, operation_id, operation_json) VALUES (?1, ?2, ?3)
19         ON CONFLICT(session_id) DO UPDATE SET operation_id=excluded.operation_id,
20             operation_json=CASE WHEN session_moves.operation_id=excluded.operation_id
21                 AND json_extract(session_moves.operation_json, '$.cancellation_requested')=1
22                 THEN json_set(excluded.operation_json, '$.cancellation_requested', json('true'))
23                 ELSE excluded.operation_json END",
24        params![
25            operation.selection.session_id,
26            operation.operation_id,
27            serde_json::to_string(operation)?
28        ],
29    )?;
30    Ok(())
31}
32
33pub fn load_move_operation(session_id: &str) -> Result<Option<MoveOperation>> {
34    let connection = open_reader(&database_path())?;
35    load_move_operation_with(&connection, session_id)
36}
37
38pub(super) fn load_move_operation_with(
39    connection: &Connection,
40    session_id: &str,
41) -> Result<Option<MoveOperation>> {
42    let json: Option<String> = connection
43        .query_row(
44            "SELECT operation_json FROM session_moves WHERE session_id=?1",
45            [session_id],
46            |row| row.get(0),
47        )
48        .optional()?;
49    Ok(json.and_then(|json| decode_move_operation(session_id, &json)))
50}
51
52/// Decode one stored move intent, or `None` when it no longer decodes.
53///
54/// A move intent recorded before a harness was removed keeps that harness in
55/// its recovery snapshot, so it fails to decode under a binary that no longer
56/// knows the harness. Nothing deletes a completed intent, so such rows linger
57/// indefinitely (BrokkAi/mjolnir#1026). Every reader asks the same question,
58/// "is there a move to act on for this session?", and a move whose snapshot
59/// names a removed harness cannot be acted on, so all readers treat the row as
60/// absent with a warning rather than failing the operation that asked, whether
61/// that is daemon startup, a stop, a checkpoint, or a new move.
62fn decode_move_operation(session_id: &str, json: &str) -> Option<MoveOperation> {
63    match serde_json::from_str::<MoveOperation>(json) {
64        Ok(operation) => Some(operation),
65        Err(error) => {
66            tracing::warn!(
67                session_id,
68                %error,
69                "durable move intent no longer decodes; treating it as absent (its harness may have been removed)"
70            );
71            None
72        }
73    }
74}
75
76pub fn load_move_operations() -> Result<Vec<MoveOperation>> {
77    let connection = open_reader(&database_path())?;
78    load_move_operations_with(&connection)
79}
80
81pub(super) fn load_move_operations_with(connection: &Connection) -> Result<Vec<MoveOperation>> {
82    let mut statement = connection
83        .prepare("SELECT session_id, operation_json FROM session_moves ORDER BY session_id")?;
84    // The daemon loads every intent at startup, so one undecodable row would
85    // otherwise stop the daemon from starting; `decode_move_operation` says why
86    // such a row is skipped rather than fatal.
87    let rows = statement.query_map([], |row| {
88        Ok((row.get::<_, String>(0)?, row.get::<_, String>(1)?))
89    })?;
90    let mut operations = Vec::new();
91    for row in rows {
92        let (session_id, json) = row?;
93        if let Some(operation) = decode_move_operation(&session_id, &json) {
94            operations.push(operation);
95        }
96    }
97    Ok(operations)
98}
99
100pub fn move_checkpoint_is_retained(path: &Path) -> Result<bool> {
101    Ok(load_move_operations()?.iter().any(|operation| {
102        operation.retains_checkpoint()
103            && operation
104                .checkpoint
105                .as_ref()
106                .is_some_and(|checkpoint| checkpoint.archive_path == path)
107    }))
108}
109
110/// Delete the rows of finished moves that can no longer act, and report how
111/// many went away.
112///
113/// A finished move keeps its row only for what still reads it: the resume
114/// dialog's "move needs recovery" offer and the viewer's retry projection, both
115/// of which restore a destination from the move's own checkpoint archive. Once
116/// that archive is gone the offer cannot do anything, so the row is litter every
117/// reader has to skip, and a lingering row is what wedged daemon startup once
118/// its harness was removed (BrokkAi/mjolnir#1026).
119///
120/// [`MoveOperation::retains_checkpoint`] already draws the line: it is false
121/// only for a completed or cancelled move that is not mid queue admission. A
122/// move that is preparing, closing, resuming, starting a queue, failed, or
123/// holding a partly admitted queue keeps its row however its archive looks.
124/// Deleting rows needs no migration.
125pub fn reap_finished_move_intents() -> Result<usize> {
126    submit_database_write("reap_finished_move_intents", |connection| {
127        reap_finished_move_intents_with(connection)
128    })
129}
130
131pub(super) fn reap_finished_move_intents_with(connection: &Connection) -> Result<usize> {
132    let mut reaped = 0;
133    for operation in load_move_operations_with(connection)? {
134        if operation.retains_checkpoint()
135            || operation
136                .checkpoint
137                .as_ref()
138                .is_some_and(|checkpoint| checkpoint.archive_path.exists())
139        {
140            continue;
141        }
142        let removed = connection.execute(
143            "DELETE FROM session_moves WHERE session_id=?1 AND operation_id=?2",
144            params![operation.selection.session_id, operation.operation_id],
145        )?;
146        if removed > 0 {
147            tracing::debug!(
148                session_id = operation.selection.session_id,
149                operation_id = operation.operation_id,
150                phase = ?operation.phase,
151                "reaped the durable row of a finished move whose checkpoint is gone"
152            );
153            reaped += removed;
154        }
155    }
156    Ok(reaped)
157}
158
159pub fn request_move_cancellation(session_id: &str) -> Result<()> {
160    let session_id = session_id.to_owned();
161    submit_database_write("request_move_cancellation", move |connection| {
162        connection.execute(
163            "UPDATE session_moves SET operation_json=json_set(operation_json, '$.cancellation_requested', json('true'))
164             WHERE session_id=?1 AND json_extract(operation_json, '$.phase') IN ('preparing','closing_source','resuming_destination','starting_queue')",
165            [session_id],
166        )?;
167        Ok(())
168    })
169}
170
171pub fn clear_move_cancellation_for_retry(session_id: &str) -> Result<()> {
172    let session_id = session_id.to_owned();
173    submit_database_write("clear_move_cancellation_for_retry", move |connection| {
174        connection.execute(
175            "UPDATE session_moves SET operation_json=json_set(operation_json, '$.cancellation_requested', json('false')) WHERE session_id=?1",
176            [session_id],
177        )?;
178        Ok(())
179    })
180}
181
182/// Read the confirmation data without loading the conversation history.
183pub fn move_pending_work(session_id: &str) -> Result<(bool, Vec<MaterializedQueuedPrompt>)> {
184    let connection = open_reader(&database_path())?;
185    let running: Option<String> = connection
186        .query_row(
187            "SELECT execution_state FROM materialized_sessions WHERE session_id=?1",
188            [session_id],
189            |row| row.get(0),
190        )
191        .optional()?;
192    Ok((
193        running.as_deref() == Some("running"),
194        read_materialized_queued_prompts(&connection, session_id)?,
195    ))
196}
197
198#[cfg(test)]
199mod tests {
200    use super::*;
201    use mj_core::state::{MovePhase, MoveSelection, ResumeQueueDisposition};
202
203    fn operation(session: &SessionRecord) -> MoveOperation {
204        MoveOperation {
205            in_place: false,
206            source_checkpoint_only: false,
207            operation_id: "move-one".into(),
208            selection: MoveSelection {
209                clear_resource_allocation: false,
210                session_id: session.id.clone(),
211                profile_id: Some("destination".into()),
212                target_template_id: Some("local".into()),
213                additional_mounts: Some(Vec::new()),
214                resource_allocation: None,
215            },
216            source_profile_id: session.last_profile.clone(),
217            source_target_template_id: session.target_template_id.clone(),
218            source_target: session.target.clone(),
219            source_native_session_id: session.native_session_id.clone(),
220            source_additional_mounts: session.additional_mounts.clone(),
221            source_resource_allocation: session.resource_allocation.clone(),
222            destination_target: None,
223            destination_native_session_id: None,
224            destination_store_id: None,
225            configuration_fingerprint: "fingerprint".into(),
226            checkpoint: session.checkpoint.clone(),
227            recovery_session: Some(session.clone()),
228            queue: ResumeQueueDisposition::Start,
229            phase: MovePhase::Preparing,
230            queue_admission_started: false,
231            queue_admission_finished: false,
232            cancellation_requested: false,
233            created_at: session.created_at.clone(),
234            updated_at: session.updated_at.clone(),
235            error: None,
236        }
237    }
238
239    #[test]
240    fn move_boundaries_survive_database_reopen_and_retain_the_source_locator() {
241        let directory = tempfile::tempdir().unwrap();
242        let path = directory.path().join("mj.sqlite3");
243        let mut session = super::super::tests::session("move-reopen", "project");
244        let template: mj_core::config::TargetTemplate = serde_json::from_str(
245            r#"{"kind":"ssh-podman","host":"original.test","image":"test","user":"builder"}"#,
246        )
247        .unwrap();
248        session.target_runtime = Some((&template).into());
249        session.target = Some(mj_core::state::TargetLocator::SshPodman {
250            host: "original.test".into(),
251            container_id: "source-container".into(),
252            workspace_storage: Default::default(),
253            borrowed_from: None,
254        });
255        save_session_to(&path, &session).unwrap();
256        let mut intent = operation(&session);
257        for phase in [
258            MovePhase::Preparing,
259            MovePhase::ClosingSource,
260            MovePhase::ResumingDestination,
261            MovePhase::StartingQueue,
262            MovePhase::Failed,
263            MovePhase::Completed,
264        ] {
265            intent.phase = phase;
266            if phase == MovePhase::StartingQueue {
267                intent.queue_admission_started = true;
268                intent.destination_target = session.target.clone();
269                intent.destination_store_id = Some("durable-destination".into());
270            }
271            if phase == MovePhase::Completed {
272                intent.queue_admission_finished = true;
273            }
274            let connection = open(&path).unwrap();
275            save_move_operation_with(&connection, &intent).unwrap();
276            drop(connection);
277            let reopened = open_reader(&path).unwrap();
278            let restored = load_move_operation_with(&reopened, &session.id)
279                .unwrap()
280                .unwrap();
281            assert_eq!(restored, intent);
282            assert_eq!(restored.source_target, session.target);
283            assert_eq!(restored.retains_checkpoint(), phase != MovePhase::Completed);
284        }
285    }
286
287    #[test]
288    fn bulk_load_skips_a_move_intent_whose_harness_no_longer_decodes() {
289        let directory = tempfile::tempdir().unwrap();
290        let path = directory.path().join("mj.sqlite3");
291        let good = super::super::tests::session("move-good", "project");
292        let stale_session = super::super::tests::session("move-removed-harness", "project");
293        save_session_to(&path, &good).unwrap();
294        save_session_to(&path, &stale_session).unwrap();
295        let connection = open(&path).unwrap();
296        save_move_operation_with(&connection, &operation(&good)).unwrap();
297        let mut stale_operation = operation(&stale_session);
298        stale_operation.operation_id = "move-two".into();
299        save_move_operation_with(&connection, &stale_operation).unwrap();
300        // Simulate a row recorded before a harness was removed: rewrite the
301        // stored recovery snapshot to name a harness the current binary no
302        // longer knows, so the row fails to decode.
303        let rewritten = connection
304            .execute(
305                "UPDATE session_moves
306                 SET operation_json = replace(operation_json, '\"harness_kind\":\"codex\"', '\"harness_kind\":\"zcode\"')
307                 WHERE session_id = 'move-removed-harness'",
308                [],
309            )
310            .unwrap();
311        assert_eq!(
312            rewritten, 1,
313            "the test session must store a codex harness to rewrite"
314        );
315        let loaded = load_move_operations_with(&connection).unwrap();
316        assert_eq!(
317            loaded.len(),
318            1,
319            "the undecodable row must be skipped, not fail the load"
320        );
321        assert_eq!(loaded[0].selection.session_id, good.id);
322    }
323
324    #[test]
325    fn in_place_intent_round_trips_and_a_legacy_row_without_it_reads_as_a_fresh_environment() {
326        let directory = tempfile::tempdir().unwrap();
327        let path = directory.path().join("mj.sqlite3");
328        let session = super::super::tests::session("move-in-place", "project");
329        save_session_to(&path, &session).unwrap();
330        let connection = open(&path).unwrap();
331        let mut intent = operation(&session);
332        intent.in_place = true;
333        save_move_operation_with(&connection, &intent).unwrap();
334        let stored: String = connection
335            .query_row(
336                "SELECT operation_json FROM session_moves WHERE session_id=?1",
337                [&session.id],
338                |row| row.get(0),
339            )
340            .unwrap();
341        assert!(
342            stored.contains("\"in_place\":true"),
343            "the in-place choice must be durable: {stored}"
344        );
345        let restored = load_move_operation_with(&connection, &session.id)
346            .unwrap()
347            .unwrap();
348        assert_eq!(restored, intent);
349        // A row written before this field existed must still decode, as the
350        // full fresh-environment move it was.
351        let rewritten = connection
352            .execute(
353                "UPDATE session_moves
354                 SET operation_json = replace(operation_json, '\"in_place\":true,', '')
355                 WHERE session_id = ?1",
356                [&session.id],
357            )
358            .unwrap();
359        assert_eq!(rewritten, 1);
360        let legacy = load_move_operation_with(&connection, &session.id)
361            .unwrap()
362            .expect("a row without in_place must still decode");
363        assert!(!legacy.in_place);
364    }
365
366    #[test]
367    fn per_session_load_treats_an_intent_whose_harness_no_longer_decodes_as_absent() {
368        // The stop, checkpoint, recovery, and new-move paths each read the one
369        // intent for their session. A completed move whose recovery snapshot
370        // names a removed harness must not fail those operations.
371        let directory = tempfile::tempdir().unwrap();
372        let path = directory.path().join("mj.sqlite3");
373        let stale_session = super::super::tests::session("move-removed-harness", "project");
374        save_session_to(&path, &stale_session).unwrap();
375        let connection = open(&path).unwrap();
376        let mut stale_operation = operation(&stale_session);
377        stale_operation.phase = MovePhase::Completed;
378        save_move_operation_with(&connection, &stale_operation).unwrap();
379        let rewritten = connection
380            .execute(
381                "UPDATE session_moves
382                 SET operation_json = replace(operation_json, '\"harness_kind\":\"codex\"', '\"harness_kind\":\"zcode\"')
383                 WHERE session_id = 'move-removed-harness'",
384                [],
385            )
386            .unwrap();
387        assert_eq!(
388            rewritten, 1,
389            "the test session must store a codex harness to rewrite"
390        );
391        let loaded = load_move_operation_with(&connection, &stale_session.id)
392            .expect("an undecodable intent must not fail the read");
393        assert!(
394            loaded.is_none(),
395            "the undecodable intent is treated as absent"
396        );
397    }
398
399    #[test]
400    fn reaping_deletes_a_finished_move_only_once_its_checkpoint_archive_is_gone() {
401        let directory = tempfile::tempdir().unwrap();
402        let path = directory.path().join("mj.sqlite3");
403        let present = directory.path().join("present.hel.zip");
404        std::fs::write(&present, b"archive").unwrap();
405        let gone = directory.path().join("gone.hel.zip");
406        // session id, phase, whether its archive is still on disk, whether its
407        // queue admission is half done, and whether the row must survive.
408        let cases = [
409            (
410                "completed-retained",
411                MovePhase::Completed,
412                true,
413                false,
414                true,
415            ),
416            ("completed-gone", MovePhase::Completed, false, false, false),
417            (
418                "completed-admitting",
419                MovePhase::Completed,
420                false,
421                true,
422                true,
423            ),
424            (
425                "cancelled-retained",
426                MovePhase::Cancelled,
427                true,
428                false,
429                true,
430            ),
431            ("cancelled-gone", MovePhase::Cancelled, false, false, false),
432            ("failed-gone", MovePhase::Failed, false, false, true),
433            (
434                "running-gone",
435                MovePhase::ResumingDestination,
436                false,
437                false,
438                true,
439            ),
440        ];
441        for (session_id, phase, archive_present, mid_admission, _) in cases {
442            let session = super::super::tests::session(session_id, "project");
443            save_session_to(&path, &session).unwrap();
444            let connection = open(&path).unwrap();
445            let mut intent = operation(&session);
446            intent.operation_id = format!("{session_id}-operation");
447            intent.phase = phase;
448            intent.queue_admission_started = mid_admission;
449            intent.checkpoint = Some(CheckpointMetadata {
450                archive_path: if archive_present {
451                    present.clone()
452                } else {
453                    gone.clone()
454                },
455                sha256: "b".repeat(64),
456                created_at: session.created_at.clone(),
457                event_frontier: 6,
458            });
459            save_move_operation_with(&connection, &intent).unwrap();
460        }
461        let connection = open(&path).unwrap();
462        let reaped = reap_finished_move_intents_with(&connection).unwrap();
463        assert_eq!(
464            reaped,
465            cases.iter().filter(|case| !case.4).count(),
466            "only the finished moves whose archive is gone are reaped"
467        );
468        for (session_id, _, _, _, survives) in cases {
469            assert_eq!(
470                load_move_operation_with(&connection, session_id)
471                    .unwrap()
472                    .is_some(),
473                survives,
474                "{session_id} row survival"
475            );
476        }
477        assert_eq!(
478            reap_finished_move_intents_with(&connection).unwrap(),
479            0,
480            "a second sweep finds nothing left to reap"
481        );
482    }
483
484    #[test]
485    fn concurrent_phase_save_cannot_erase_durable_cancellation() {
486        let directory = tempfile::tempdir().unwrap();
487        let path = directory.path().join("mj.sqlite3");
488        let session = super::super::tests::session("move-cancel", "project");
489        save_session_to(&path, &session).unwrap();
490        let connection = open(&path).unwrap();
491        let mut intent = operation(&session);
492        save_move_operation_with(&connection, &intent).unwrap();
493        let mut cancelled = intent.clone();
494        cancelled.cancellation_requested = true;
495        save_move_operation_with(&connection, &cancelled).unwrap();
496        intent.phase = MovePhase::ClosingSource;
497        save_move_operation_with(&connection, &intent).unwrap();
498        let restored = load_move_operation_with(&connection, &session.id)
499            .unwrap()
500            .unwrap();
501        assert!(restored.cancellation_requested);
502        assert_eq!(restored.phase, MovePhase::ClosingSource);
503        intent.operation_id = "explicit-new-operation".into();
504        save_move_operation_with(&connection, &intent).unwrap();
505        assert!(
506            !load_move_operation_with(&connection, &session.id)
507                .unwrap()
508                .unwrap()
509                .cancellation_requested
510        );
511    }
512
513    #[test]
514    fn cancelling_partial_queue_admission_does_not_release_its_archive() {
515        let session = super::super::tests::session("move-queue", "project");
516        let mut intent = operation(&session);
517        intent.phase = MovePhase::Cancelled;
518        intent.queue_admission_started = true;
519        assert!(intent.retains_checkpoint());
520        intent.queue_admission_finished = true;
521        assert!(!intent.retains_checkpoint());
522    }
523
524    #[test]
525    fn destination_record_install_keeps_drafts_and_titles_edited_during_move() {
526        let directory = tempfile::tempdir().unwrap();
527        let path = directory.path().join("mj.sqlite3");
528        let mut stale = super::super::tests::session("move-draft", "project");
529        save_session_to(&path, &stale).unwrap();
530        let connection = open(&path).unwrap();
531        save_move_operation_with(&connection, &operation(&stale)).unwrap();
532        connection.execute("UPDATE sessions SET draft_input='keep this draft', session_title_override='new title' WHERE session_id=?1", [&stale.id]).unwrap();
533        stale.last_profile = "destination".into();
534        stale.state = SessionState::Provisioning;
535        save_session_to(&path, &stale).unwrap();
536        let current = load_state_from(&path)
537            .unwrap()
538            .sessions
539            .remove(&stale.id)
540            .unwrap();
541        assert_eq!(current.draft_input, "keep this draft");
542        assert_eq!(current.session_title_override.as_deref(), Some("new title"));
543        assert_eq!(current.last_profile, "destination");
544    }
545}