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 retain_move_source(operation: &MoveOperation) -> Result<()> {
7    let operation = operation.clone();
8    submit_database_write("retain_move_source", move |connection| {
9        let transfer = operation
10            .workspace_transfer
11            .as_ref()
12            .context("Move transfer missing")?;
13        connection.execute(
14            "INSERT INTO retained_move_sources(operation_id, session_id, source_json, exclusions_json, created_at)
15             VALUES (?1, ?2, ?3, ?4, ?5) ON CONFLICT(operation_id) DO NOTHING",
16            params![operation.operation_id, operation.selection.session_id,
17                serde_json::to_string(&transfer.source)?, serde_json::to_string(&operation.selection.workspace.exclusions)?, operation.created_at],
18        )?;
19        Ok(())
20    })
21}
22
23pub fn retained_move_sources(
24    session_id: &str,
25) -> Result<Vec<mj_core::move_workspace::RetainedMoveSource>> {
26    let connection = open_reader(&database_path())?;
27    let mut query = connection.prepare("SELECT operation_id, source_json, exclusions_json, created_at FROM retained_move_sources WHERE session_id=?1 ORDER BY created_at")?;
28    let rows = query.query_map([session_id], |row| {
29        Ok((
30            row.get::<_, String>(0)?,
31            row.get::<_, String>(1)?,
32            row.get::<_, String>(2)?,
33            row.get::<_, String>(3)?,
34        ))
35    })?;
36    rows.map(|row| {
37        let (operation_id, source, exclusions, created_at) = row?;
38        Ok(mj_core::move_workspace::RetainedMoveSource {
39            operation_id,
40            session_id: session_id.into(),
41            source: serde_json::from_str(&source)?,
42            exclusions: serde_json::from_str(&exclusions)?,
43            created_at,
44        })
45    })
46    .collect()
47}
48
49pub fn forget_retained_move_source(operation_id: &str) -> Result<()> {
50    let operation_id = operation_id.to_owned();
51    submit_database_write("forget_retained_move_source", move |connection| {
52        connection.execute(
53            "DELETE FROM retained_move_sources WHERE operation_id=?1",
54            [operation_id],
55        )?;
56        Ok(())
57    })
58}
59
60pub fn save_move_operation(operation: &MoveOperation) -> Result<()> {
61    let operation = operation.clone();
62    submit_database_write("save_move_operation", move |connection| {
63        save_move_operation_with(connection, &operation)
64    })
65}
66
67/// Publish destination adoption and ownership together, without changing client fields.
68pub fn adopt_move_destination(operation: &MoveOperation, session: &SessionRecord) -> Result<()> {
69    let operation = operation.clone();
70    let session = session.clone();
71    submit_database_write("adopt_move_destination", move |connection| {
72        let tx = connection.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
73        update_lifecycle_fields(&tx, &session)?;
74        save_move_operation_with(&tx, &operation)?;
75        tx.commit()?;
76        Ok(())
77    })
78}
79
80/// Publish a Move result and its session message as one committed outcome.
81/// Update only the message fields; other owners may have changed the session.
82pub fn save_move_outcome(operation: &MoveOperation, last_error: Option<&str>) -> Result<()> {
83    let operation = operation.clone();
84    let last_error = last_error.map(str::to_owned);
85    submit_database_write("save_move_outcome", move |connection| {
86        save_move_outcome_with(connection, &operation, last_error.as_deref())
87    })
88}
89
90fn save_move_outcome_with(
91    connection: &mut Connection,
92    operation: &MoveOperation,
93    last_error: Option<&str>,
94) -> Result<()> {
95    let tx = connection.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
96    save_move_operation_with(&tx, operation)?;
97    let updated = tx.execute(
98        "UPDATE sessions SET last_error=?2, updated_at=?3 WHERE session_id=?1",
99        params![
100            operation.selection.session_id,
101            last_error,
102            operation.updated_at
103        ],
104    )?;
105    anyhow::ensure!(updated == 1, "Move outcome session is missing");
106    tx.commit()?;
107    Ok(())
108}
109
110pub(super) fn save_move_operation_with(
111    connection: &Connection,
112    operation: &MoveOperation,
113) -> Result<()> {
114    connection.execute(
115        "INSERT INTO session_moves(session_id, operation_id, operation_json) VALUES (?1, ?2, ?3)
116         ON CONFLICT(session_id) DO UPDATE SET operation_id=excluded.operation_id,
117             operation_json=CASE WHEN session_moves.operation_id=excluded.operation_id
118                 AND json_extract(session_moves.operation_json, '$.cancellation_requested')=1
119                 THEN json_set(excluded.operation_json, '$.cancellation_requested', json('true'))
120                 ELSE excluded.operation_json END",
121        params![
122            operation.selection.session_id,
123            operation.operation_id,
124            serde_json::to_string(operation)?
125        ],
126    )?;
127    Ok(())
128}
129
130pub fn load_move_operation(session_id: &str) -> Result<Option<MoveOperation>> {
131    if let Some(committed) = committed_state()? {
132        return Ok(committed.moves.get(session_id).cloned());
133    }
134    let connection = open_reader(&database_path())?;
135    load_move_operation_with(&connection, session_id)
136}
137
138pub(super) fn load_move_operation_with(
139    connection: &Connection,
140    session_id: &str,
141) -> Result<Option<MoveOperation>> {
142    let json: Option<String> = connection
143        .query_row(
144            "SELECT operation_json FROM session_moves WHERE session_id=?1",
145            [session_id],
146            |row| row.get(0),
147        )
148        .optional()?;
149    Ok(json.and_then(|json| decode_move_operation(session_id, &json)))
150}
151
152/// Decode one stored move intent, or `None` when it no longer decodes.
153///
154/// A move intent recorded before a harness was removed keeps that harness in
155/// its recovery snapshot, so it fails to decode under a binary that no longer
156/// knows the harness. Nothing deletes a completed intent, so such rows linger
157/// indefinitely (BrokkAi/mjolnir#1026). Every reader asks the same question,
158/// "is there a move to act on for this session?", and a move whose snapshot
159/// names a removed harness cannot be acted on, so all readers treat the row as
160/// absent with a warning rather than failing the operation that asked, whether
161/// that is daemon startup, a stop, a checkpoint, or a new move.
162fn decode_move_operation(session_id: &str, json: &str) -> Option<MoveOperation> {
163    match serde_json::from_str::<MoveOperation>(json) {
164        Ok(operation) => Some(operation),
165        Err(error) => {
166            tracing::warn!(
167                session_id,
168                %error,
169                "durable move intent no longer decodes; treating it as absent (its harness may have been removed)"
170            );
171            None
172        }
173    }
174}
175
176pub fn load_move_operations() -> Result<Vec<MoveOperation>> {
177    if let Some(committed) = committed_state()? {
178        return Ok(committed.moves.values().cloned().collect());
179    }
180    let connection = open_reader(&database_path())?;
181    load_move_operations_with(&connection)
182}
183
184pub(super) fn load_move_operations_with(connection: &Connection) -> Result<Vec<MoveOperation>> {
185    let mut statement = connection
186        .prepare("SELECT session_id, operation_json FROM session_moves ORDER BY session_id")?;
187    // The daemon loads every intent at startup, so one undecodable row would
188    // otherwise stop the daemon from starting; `decode_move_operation` says why
189    // such a row is skipped rather than fatal.
190    let rows = statement.query_map([], |row| {
191        Ok((row.get::<_, String>(0)?, row.get::<_, String>(1)?))
192    })?;
193    let mut operations = Vec::new();
194    for row in rows {
195        let (session_id, json) = row?;
196        if let Some(operation) = decode_move_operation(&session_id, &json) {
197            operations.push(operation);
198        }
199    }
200    Ok(operations)
201}
202
203pub fn move_checkpoint_is_retained(path: &Path) -> Result<bool> {
204    Ok(load_move_operations()?.iter().any(|operation| {
205        operation
206            .retained_archives()
207            .any(|checkpoint| checkpoint.archive_path == path)
208    }))
209}
210
211/// Delete the rows of finished moves that can no longer act, and report how
212/// many went away.
213///
214/// A finished move keeps its row only for what still reads it: the resume
215/// dialog's "move needs recovery" offer and the viewer's retry projection, both
216/// of which restore a destination from the move's own checkpoint archive. Once
217/// that archive is gone the offer cannot do anything, so the row is litter every
218/// reader has to skip, and a lingering row is what wedged daemon startup once
219/// its harness was removed (BrokkAi/mjolnir#1026).
220///
221/// [`MoveOperation::retains_checkpoint`] already draws the line: it is false
222/// only for a completed or cancelled move that is not mid queue admission. A
223/// move that is preparing, closing, resuming, starting a queue, failed, or
224/// holding a partly admitted queue keeps its row however its archive looks.
225/// Deleting rows needs no migration.
226pub fn reap_finished_move_intents() -> Result<usize> {
227    submit_database_write("reap_finished_move_intents", |connection| {
228        reap_finished_move_intents_with(connection)
229    })
230}
231
232pub(super) fn reap_finished_move_intents_with(connection: &Connection) -> Result<usize> {
233    let mut reaped = 0;
234    for operation in load_move_operations_with(connection)? {
235        if (operation.phase != mj_core::state::MovePhase::Completed
236            && operation
237                .prepared_destination
238                .as_ref()
239                .is_some_and(|destination| destination.owns_resource()))
240            || operation.retains_checkpoint()
241            || operation
242                .restore_artifact()
243                .is_some_and(|checkpoint| checkpoint.archive_path.exists())
244        {
245            continue;
246        }
247        let removed = connection.execute(
248            "DELETE FROM session_moves WHERE session_id=?1 AND operation_id=?2",
249            params![operation.selection.session_id, operation.operation_id],
250        )?;
251        if removed > 0 {
252            tracing::debug!(
253                session_id = operation.selection.session_id,
254                operation_id = operation.operation_id,
255                phase = ?operation.phase,
256                "reaped the durable row of a finished move whose checkpoint is gone"
257            );
258            reaped += removed;
259        }
260    }
261    Ok(reaped)
262}
263
264pub fn request_move_cancellation(session_id: &str) -> Result<()> {
265    let session_id = session_id.to_owned();
266    submit_database_write("request_move_cancellation", move |connection| {
267        connection.execute(
268            "UPDATE session_moves SET operation_json=json_set(operation_json, '$.cancellation_requested', json('true'))
269             WHERE session_id=?1 AND json_extract(operation_json, '$.phase') IN ('preparing','closing_source','resuming_destination','starting_queue')",
270            [session_id],
271        )?;
272        Ok(())
273    })
274}
275
276pub fn clear_move_cancellation_for_retry(session_id: &str) -> Result<()> {
277    let session_id = session_id.to_owned();
278    submit_database_write("clear_move_cancellation_for_retry", move |connection| {
279        connection.execute(
280            "UPDATE session_moves SET operation_json=json_set(operation_json, '$.cancellation_requested', json('false')) WHERE session_id=?1",
281            [session_id],
282        )?;
283        Ok(())
284    })
285}
286
287/// Read the confirmation data without loading the conversation history.
288pub fn move_pending_work(session_id: &str) -> Result<(bool, Vec<MaterializedQueuedPrompt>)> {
289    let connection = open_reader(&database_path())?;
290    let running: Option<String> = connection
291        .query_row(
292            "SELECT execution_state FROM materialized_sessions WHERE session_id=?1",
293            [session_id],
294            |row| row.get(0),
295        )
296        .optional()?;
297    Ok((
298        running.as_deref() == Some("running"),
299        read_materialized_queued_prompts(&connection, session_id)?,
300    ))
301}
302
303#[cfg(test)]
304mod tests {
305    use super::*;
306    use mj_core::state::{MovePhase, MoveSelection, ResumeQueueDisposition};
307
308    fn operation(session: &SessionRecord) -> MoveOperation {
309        MoveOperation {
310            prepared_destination: None,
311            accepted_preparation: None,
312            acknowledge_interruption: false,
313            workspace_transfer: None,
314            handoff: None,
315            in_place: false,
316            source_checkpoint_only: false,
317            operation_id: "move-one".into(),
318            selection: MoveSelection {
319                subagents: None,
320                workspace: Default::default(),
321                clear_resource_allocation: false,
322                session_id: session.id.clone(),
323                profile_id: Some("destination".into()),
324                target_template_id: Some("local".into()),
325                additional_mounts: Some(Vec::new()),
326                resource_allocation: None,
327            },
328            source_profile_id: session.last_profile.clone(),
329            source_target_template_id: session.target_template_id.clone(),
330            source_target: session.target.clone(),
331            source_native_session_id: session.native_session_id.clone(),
332            source_additional_mounts: session.additional_mounts.clone(),
333            source_resource_allocation: session.resource_allocation.clone(),
334            destination_target: None,
335            destination_native_session_id: None,
336            destination_store_id: None,
337            configuration_fingerprint: "fingerprint".into(),
338            checkpoint: session.checkpoint.clone(),
339            recovery_session: Some(session.clone()),
340            queue: ResumeQueueDisposition::Start,
341            phase: MovePhase::Preparing,
342            queue_admission_started: false,
343            queue_admission_finished: false,
344            cancellation_requested: false,
345            created_at: session.created_at.clone(),
346            updated_at: session.updated_at.clone(),
347            error: None,
348        }
349    }
350
351    #[test]
352    fn move_outcome_commits_the_message_with_the_result_and_rolls_back_both_on_failure() {
353        let directory = tempfile::tempdir().unwrap();
354        let path = directory.path().join("controller.sqlite");
355        let mut session = super::super::tests::session("moving", "project");
356        session.draft_input = "a concurrent draft".into();
357        save_session_to(&path, &session).unwrap();
358        let mut connection = open(&path).unwrap();
359        let mut intent = operation(&session);
360        save_move_operation_with(&connection, &intent).unwrap();
361        connection
362            .execute_batch(
363                "CREATE TRIGGER refuse_move_message BEFORE UPDATE OF last_error ON sessions
364             BEGIN SELECT RAISE(ABORT, 'message write failed'); END;",
365            )
366            .unwrap();
367        intent.phase = MovePhase::Failed;
368        intent.error = Some("private diagnostic".into());
369        assert!(save_move_outcome_with(&mut connection, &intent, Some("public message")).is_err());
370        assert_eq!(
371            load_move_operation_with(&connection, &session.id)
372                .unwrap()
373                .unwrap()
374                .phase,
375            MovePhase::Preparing
376        );
377        assert_eq!(
378            load_state_from(&path).unwrap().sessions[&session.id].last_error,
379            None
380        );
381        connection
382            .execute_batch("DROP TRIGGER refuse_move_message")
383            .unwrap();
384        save_move_outcome_with(&mut connection, &intent, Some("public message")).unwrap();
385        drop(connection);
386        let restored = load_state_from(&path).unwrap();
387        assert_eq!(
388            restored.sessions[&session.id].last_error.as_deref(),
389            Some("public message")
390        );
391        assert_eq!(
392            restored.sessions[&session.id].draft_input,
393            "a concurrent draft"
394        );
395        assert_eq!(
396            load_move_operation_with(&open_reader(&path).unwrap(), &session.id)
397                .unwrap()
398                .unwrap(),
399            intent
400        );
401    }
402
403    #[test]
404    fn committed_moves_follow_phase_changes_and_cascading_deletion() {
405        let directory = tempfile::tempdir().unwrap();
406        let path = directory.path().join("controller.sqlite");
407        let session = super::super::tests::session("moving", "project");
408        save_session_to(&path, &session).unwrap();
409        let intent = operation(&session);
410        save_move_operation_with(&open(&path).unwrap(), &intent).unwrap();
411        let writer = start_database_writer_at(&path, false).unwrap();
412        let held = writer.writer.committed_state().unwrap();
413        assert_eq!(held.moves["moving"], intent);
414        let mut failed = intent.clone();
415        failed.phase = MovePhase::Failed;
416        writer
417            .writer
418            .execute("fail move", move |connection| {
419                save_move_operation_with(connection, &failed)
420            })
421            .unwrap();
422        assert_eq!(
423            writer.writer.committed_state().unwrap().moves["moving"].phase,
424            MovePhase::Failed
425        );
426        assert_eq!(held.moves["moving"].phase, MovePhase::Preparing);
427        writer
428            .writer
429            .execute("delete moving session", |connection| {
430                connection.execute("DELETE FROM sessions WHERE session_id='moving'", [])?;
431                Ok(())
432            })
433            .unwrap();
434        assert!(writer.writer.committed_state().unwrap().moves.is_empty());
435    }
436
437    #[test]
438    fn move_boundaries_survive_database_reopen_and_retain_the_source_locator() {
439        let directory = tempfile::tempdir().unwrap();
440        let path = directory.path().join("mj.sqlite3");
441        let mut session = super::super::tests::session("move-reopen", "project");
442        let template: mj_core::config::TargetTemplate = serde_json::from_str(
443            r#"{"kind":"ssh-podman","host":"original.test","image":"test","user":"builder"}"#,
444        )
445        .unwrap();
446        session.target_runtime = Some((&template).into());
447        session.target = Some(mj_core::state::TargetLocator::SshPodman {
448            host: "original.test".into(),
449            container_id: "source-container".into(),
450            workspace_storage: Default::default(),
451            borrowed_from: None,
452        });
453        save_session_to(&path, &session).unwrap();
454        let mut intent = operation(&session);
455        for phase in [
456            MovePhase::Preparing,
457            MovePhase::ClosingSource,
458            MovePhase::ResumingDestination,
459            MovePhase::StartingQueue,
460            MovePhase::Failed,
461            MovePhase::Completed,
462        ] {
463            intent.phase = phase;
464            if phase == MovePhase::StartingQueue {
465                intent.queue_admission_started = true;
466                intent.destination_target = session.target.clone();
467                intent.destination_store_id = Some("durable-destination".into());
468            }
469            if phase == MovePhase::Completed {
470                intent.queue_admission_finished = true;
471            }
472            let connection = open(&path).unwrap();
473            save_move_operation_with(&connection, &intent).unwrap();
474            drop(connection);
475            let reopened = open_reader(&path).unwrap();
476            let restored = load_move_operation_with(&reopened, &session.id)
477                .unwrap()
478                .unwrap();
479            assert_eq!(restored, intent);
480            assert_eq!(restored.source_target, session.target);
481            assert_eq!(restored.retains_checkpoint(), phase != MovePhase::Completed);
482        }
483    }
484
485    #[test]
486    fn bulk_load_skips_a_move_intent_whose_harness_no_longer_decodes() {
487        let directory = tempfile::tempdir().unwrap();
488        let path = directory.path().join("mj.sqlite3");
489        let good = super::super::tests::session("move-good", "project");
490        let stale_session = super::super::tests::session("move-removed-harness", "project");
491        save_session_to(&path, &good).unwrap();
492        save_session_to(&path, &stale_session).unwrap();
493        let connection = open(&path).unwrap();
494        save_move_operation_with(&connection, &operation(&good)).unwrap();
495        let mut stale_operation = operation(&stale_session);
496        stale_operation.operation_id = "move-two".into();
497        save_move_operation_with(&connection, &stale_operation).unwrap();
498        // Simulate a row recorded before a harness was removed: rewrite the
499        // stored recovery snapshot to name a harness the current binary no
500        // longer knows, so the row fails to decode.
501        let rewritten = connection
502            .execute(
503                "UPDATE session_moves
504                 SET operation_json = replace(operation_json, '\"harness_kind\":\"codex\"', '\"harness_kind\":\"zcode\"')
505                 WHERE session_id = 'move-removed-harness'",
506                [],
507            )
508            .unwrap();
509        assert_eq!(
510            rewritten, 1,
511            "the test session must store a codex harness to rewrite"
512        );
513        let loaded = load_move_operations_with(&connection).unwrap();
514        assert_eq!(
515            loaded.len(),
516            1,
517            "the undecodable row must be skipped, not fail the load"
518        );
519        assert_eq!(loaded[0].selection.session_id, good.id);
520    }
521
522    #[test]
523    fn in_place_intent_round_trips_and_a_legacy_row_without_it_reads_as_a_fresh_environment() {
524        let directory = tempfile::tempdir().unwrap();
525        let path = directory.path().join("mj.sqlite3");
526        let session = super::super::tests::session("move-in-place", "project");
527        save_session_to(&path, &session).unwrap();
528        let connection = open(&path).unwrap();
529        let mut intent = operation(&session);
530        intent.in_place = true;
531        save_move_operation_with(&connection, &intent).unwrap();
532        let stored: String = connection
533            .query_row(
534                "SELECT operation_json FROM session_moves WHERE session_id=?1",
535                [&session.id],
536                |row| row.get(0),
537            )
538            .unwrap();
539        assert!(
540            stored.contains("\"in_place\":true"),
541            "the in-place choice must be durable: {stored}"
542        );
543        let restored = load_move_operation_with(&connection, &session.id)
544            .unwrap()
545            .unwrap();
546        assert_eq!(restored, intent);
547        // A row written before this field existed must still decode, as the
548        // full fresh-environment move it was.
549        let rewritten = connection
550            .execute(
551                "UPDATE session_moves
552                 SET operation_json = replace(operation_json, '\"in_place\":true,', '')
553                 WHERE session_id = ?1",
554                [&session.id],
555            )
556            .unwrap();
557        assert_eq!(rewritten, 1);
558        let legacy = load_move_operation_with(&connection, &session.id)
559            .unwrap()
560            .expect("a row without in_place must still decode");
561        assert!(!legacy.in_place);
562    }
563
564    #[test]
565    fn per_session_load_treats_an_intent_whose_harness_no_longer_decodes_as_absent() {
566        // The stop, checkpoint, recovery, and new-move paths each read the one
567        // intent for their session. A completed move whose recovery snapshot
568        // names a removed harness must not fail those operations.
569        let directory = tempfile::tempdir().unwrap();
570        let path = directory.path().join("mj.sqlite3");
571        let stale_session = super::super::tests::session("move-removed-harness", "project");
572        save_session_to(&path, &stale_session).unwrap();
573        let connection = open(&path).unwrap();
574        let mut stale_operation = operation(&stale_session);
575        stale_operation.phase = MovePhase::Completed;
576        save_move_operation_with(&connection, &stale_operation).unwrap();
577        let rewritten = connection
578            .execute(
579                "UPDATE session_moves
580                 SET operation_json = replace(operation_json, '\"harness_kind\":\"codex\"', '\"harness_kind\":\"zcode\"')
581                 WHERE session_id = 'move-removed-harness'",
582                [],
583            )
584            .unwrap();
585        assert_eq!(
586            rewritten, 1,
587            "the test session must store a codex harness to rewrite"
588        );
589        let loaded = load_move_operation_with(&connection, &stale_session.id)
590            .expect("an undecodable intent must not fail the read");
591        assert!(
592            loaded.is_none(),
593            "the undecodable intent is treated as absent"
594        );
595    }
596
597    #[test]
598    fn reaping_deletes_a_finished_move_only_once_its_checkpoint_archive_is_gone() {
599        let directory = tempfile::tempdir().unwrap();
600        let path = directory.path().join("mj.sqlite3");
601        let present = directory.path().join("present.hel.zip");
602        std::fs::write(&present, b"archive").unwrap();
603        let gone = directory.path().join("gone.hel.zip");
604        // session id, phase, whether its archive is still on disk, whether its
605        // queue admission is half done, and whether the row must survive.
606        let cases = [
607            (
608                "completed-retained",
609                MovePhase::Completed,
610                true,
611                false,
612                true,
613            ),
614            ("completed-gone", MovePhase::Completed, false, false, false),
615            (
616                "completed-admitting",
617                MovePhase::Completed,
618                false,
619                true,
620                true,
621            ),
622            (
623                "cancelled-retained",
624                MovePhase::Cancelled,
625                true,
626                false,
627                true,
628            ),
629            ("cancelled-gone", MovePhase::Cancelled, false, false, false),
630            ("failed-gone", MovePhase::Failed, false, false, true),
631            (
632                "running-gone",
633                MovePhase::ResumingDestination,
634                false,
635                false,
636                true,
637            ),
638        ];
639        for (session_id, phase, archive_present, mid_admission, _) in cases {
640            let session = super::super::tests::session(session_id, "project");
641            save_session_to(&path, &session).unwrap();
642            let connection = open(&path).unwrap();
643            let mut intent = operation(&session);
644            intent.operation_id = format!("{session_id}-operation");
645            intent.phase = phase;
646            intent.queue_admission_started = mid_admission;
647            intent.checkpoint = Some(CheckpointMetadata {
648                archive_path: if archive_present {
649                    present.clone()
650                } else {
651                    gone.clone()
652                },
653                sha256: "b".repeat(64),
654                created_at: session.created_at.clone(),
655                event_frontier: 6,
656            });
657            save_move_operation_with(&connection, &intent).unwrap();
658        }
659        let connection = open(&path).unwrap();
660        let reaped = reap_finished_move_intents_with(&connection).unwrap();
661        assert_eq!(
662            reaped,
663            cases.iter().filter(|case| !case.4).count(),
664            "only the finished moves whose archive is gone are reaped"
665        );
666        for (session_id, _, _, _, survives) in cases {
667            assert_eq!(
668                load_move_operation_with(&connection, session_id)
669                    .unwrap()
670                    .is_some(),
671                survives,
672                "{session_id} row survival"
673            );
674        }
675        assert_eq!(
676            reap_finished_move_intents_with(&connection).unwrap(),
677            0,
678            "a second sweep finds nothing left to reap"
679        );
680    }
681
682    #[test]
683    fn concurrent_phase_save_cannot_erase_durable_cancellation() {
684        let directory = tempfile::tempdir().unwrap();
685        let path = directory.path().join("mj.sqlite3");
686        let session = super::super::tests::session("move-cancel", "project");
687        save_session_to(&path, &session).unwrap();
688        let connection = open(&path).unwrap();
689        let mut intent = operation(&session);
690        save_move_operation_with(&connection, &intent).unwrap();
691        let mut cancelled = intent.clone();
692        cancelled.cancellation_requested = true;
693        save_move_operation_with(&connection, &cancelled).unwrap();
694        intent.phase = MovePhase::ClosingSource;
695        save_move_operation_with(&connection, &intent).unwrap();
696        let restored = load_move_operation_with(&connection, &session.id)
697            .unwrap()
698            .unwrap();
699        assert!(restored.cancellation_requested);
700        assert_eq!(restored.phase, MovePhase::ClosingSource);
701        intent.operation_id = "explicit-new-operation".into();
702        save_move_operation_with(&connection, &intent).unwrap();
703        assert!(
704            !load_move_operation_with(&connection, &session.id)
705                .unwrap()
706                .unwrap()
707                .cancellation_requested
708        );
709    }
710
711    #[test]
712    fn cancelling_partial_queue_admission_does_not_release_its_archive() {
713        let session = super::super::tests::session("move-queue", "project");
714        let mut intent = operation(&session);
715        intent.phase = MovePhase::Cancelled;
716        intent.queue_admission_started = true;
717        assert!(intent.retains_checkpoint());
718        intent.queue_admission_finished = true;
719        assert!(!intent.retains_checkpoint());
720    }
721
722    #[test]
723    fn destination_record_install_keeps_drafts_and_titles_edited_during_move() {
724        let directory = tempfile::tempdir().unwrap();
725        let path = directory.path().join("mj.sqlite3");
726        let mut stale = super::super::tests::session("move-draft", "project");
727        save_session_to(&path, &stale).unwrap();
728        let connection = open(&path).unwrap();
729        save_move_operation_with(&connection, &operation(&stale)).unwrap();
730        connection.execute("UPDATE sessions SET draft_input='keep this draft', session_title_override='new title' WHERE session_id=?1", [&stale.id]).unwrap();
731        stale.last_profile = "destination".into();
732        stale.state = SessionState::Provisioning;
733        save_session_to(&path, &stale).unwrap();
734        let current = load_state_from(&path)
735            .unwrap()
736            .sessions
737            .remove(&stale.id)
738            .unwrap();
739        assert_eq!(current.draft_input, "keep this draft");
740        assert_eq!(current.session_title_override.as_deref(), Some("new title"));
741        assert_eq!(current.last_profile, "destination");
742    }
743}