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