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