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