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 session = super::super::tests::session("move-reopen", "project");
244 save_session_to(&path, &session).unwrap();
245 let mut intent = operation(&session);
246 for phase in [
247 MovePhase::Preparing,
248 MovePhase::ClosingSource,
249 MovePhase::ResumingDestination,
250 MovePhase::StartingQueue,
251 MovePhase::Failed,
252 MovePhase::Completed,
253 ] {
254 intent.phase = phase;
255 if phase == MovePhase::StartingQueue {
256 intent.queue_admission_started = true;
257 intent.destination_target = session.target.clone();
258 intent.destination_store_id = Some("durable-destination".into());
259 }
260 if phase == MovePhase::Completed {
261 intent.queue_admission_finished = true;
262 }
263 let connection = open(&path).unwrap();
264 save_move_operation_with(&connection, &intent).unwrap();
265 drop(connection);
266 let reopened = open_reader(&path).unwrap();
267 let restored = load_move_operation_with(&reopened, &session.id)
268 .unwrap()
269 .unwrap();
270 assert_eq!(restored, intent);
271 assert_eq!(restored.source_target, session.target);
272 assert_eq!(restored.retains_checkpoint(), phase != MovePhase::Completed);
273 }
274 }
275
276 #[test]
277 fn bulk_load_skips_a_move_intent_whose_harness_no_longer_decodes() {
278 let directory = tempfile::tempdir().unwrap();
279 let path = directory.path().join("mj.sqlite3");
280 let good = super::super::tests::session("move-good", "project");
281 let stale_session = super::super::tests::session("move-removed-harness", "project");
282 save_session_to(&path, &good).unwrap();
283 save_session_to(&path, &stale_session).unwrap();
284 let connection = open(&path).unwrap();
285 save_move_operation_with(&connection, &operation(&good)).unwrap();
286 let mut stale_operation = operation(&stale_session);
287 stale_operation.operation_id = "move-two".into();
288 save_move_operation_with(&connection, &stale_operation).unwrap();
289 let rewritten = connection
293 .execute(
294 "UPDATE session_moves
295 SET operation_json = replace(operation_json, '\"harness_kind\":\"codex\"', '\"harness_kind\":\"zcode\"')
296 WHERE session_id = 'move-removed-harness'",
297 [],
298 )
299 .unwrap();
300 assert_eq!(
301 rewritten, 1,
302 "the test session must store a codex harness to rewrite"
303 );
304 let loaded = load_move_operations_with(&connection).unwrap();
305 assert_eq!(
306 loaded.len(),
307 1,
308 "the undecodable row must be skipped, not fail the load"
309 );
310 assert_eq!(loaded[0].selection.session_id, good.id);
311 }
312
313 #[test]
314 fn in_place_intent_round_trips_and_a_legacy_row_without_it_reads_as_a_fresh_environment() {
315 let directory = tempfile::tempdir().unwrap();
316 let path = directory.path().join("mj.sqlite3");
317 let session = super::super::tests::session("move-in-place", "project");
318 save_session_to(&path, &session).unwrap();
319 let connection = open(&path).unwrap();
320 let mut intent = operation(&session);
321 intent.in_place = true;
322 save_move_operation_with(&connection, &intent).unwrap();
323 let stored: String = connection
324 .query_row(
325 "SELECT operation_json FROM session_moves WHERE session_id=?1",
326 [&session.id],
327 |row| row.get(0),
328 )
329 .unwrap();
330 assert!(
331 stored.contains("\"in_place\":true"),
332 "the in-place choice must be durable: {stored}"
333 );
334 let restored = load_move_operation_with(&connection, &session.id)
335 .unwrap()
336 .unwrap();
337 assert_eq!(restored, intent);
338 let rewritten = connection
341 .execute(
342 "UPDATE session_moves
343 SET operation_json = replace(operation_json, '\"in_place\":true,', '')
344 WHERE session_id = ?1",
345 [&session.id],
346 )
347 .unwrap();
348 assert_eq!(rewritten, 1);
349 let legacy = load_move_operation_with(&connection, &session.id)
350 .unwrap()
351 .expect("a row without in_place must still decode");
352 assert!(!legacy.in_place);
353 }
354
355 #[test]
356 fn per_session_load_treats_an_intent_whose_harness_no_longer_decodes_as_absent() {
357 let directory = tempfile::tempdir().unwrap();
361 let path = directory.path().join("mj.sqlite3");
362 let stale_session = super::super::tests::session("move-removed-harness", "project");
363 save_session_to(&path, &stale_session).unwrap();
364 let connection = open(&path).unwrap();
365 let mut stale_operation = operation(&stale_session);
366 stale_operation.phase = MovePhase::Completed;
367 save_move_operation_with(&connection, &stale_operation).unwrap();
368 let rewritten = connection
369 .execute(
370 "UPDATE session_moves
371 SET operation_json = replace(operation_json, '\"harness_kind\":\"codex\"', '\"harness_kind\":\"zcode\"')
372 WHERE session_id = 'move-removed-harness'",
373 [],
374 )
375 .unwrap();
376 assert_eq!(
377 rewritten, 1,
378 "the test session must store a codex harness to rewrite"
379 );
380 let loaded = load_move_operation_with(&connection, &stale_session.id)
381 .expect("an undecodable intent must not fail the read");
382 assert!(
383 loaded.is_none(),
384 "the undecodable intent is treated as absent"
385 );
386 }
387
388 #[test]
389 fn reaping_deletes_a_finished_move_only_once_its_checkpoint_archive_is_gone() {
390 let directory = tempfile::tempdir().unwrap();
391 let path = directory.path().join("mj.sqlite3");
392 let present = directory.path().join("present.hel.zip");
393 std::fs::write(&present, b"archive").unwrap();
394 let gone = directory.path().join("gone.hel.zip");
395 let cases = [
398 (
399 "completed-retained",
400 MovePhase::Completed,
401 true,
402 false,
403 true,
404 ),
405 ("completed-gone", MovePhase::Completed, false, false, false),
406 (
407 "completed-admitting",
408 MovePhase::Completed,
409 false,
410 true,
411 true,
412 ),
413 (
414 "cancelled-retained",
415 MovePhase::Cancelled,
416 true,
417 false,
418 true,
419 ),
420 ("cancelled-gone", MovePhase::Cancelled, false, false, false),
421 ("failed-gone", MovePhase::Failed, false, false, true),
422 (
423 "running-gone",
424 MovePhase::ResumingDestination,
425 false,
426 false,
427 true,
428 ),
429 ];
430 for (session_id, phase, archive_present, mid_admission, _) in cases {
431 let session = super::super::tests::session(session_id, "project");
432 save_session_to(&path, &session).unwrap();
433 let connection = open(&path).unwrap();
434 let mut intent = operation(&session);
435 intent.operation_id = format!("{session_id}-operation");
436 intent.phase = phase;
437 intent.queue_admission_started = mid_admission;
438 intent.checkpoint = Some(CheckpointMetadata {
439 archive_path: if archive_present {
440 present.clone()
441 } else {
442 gone.clone()
443 },
444 sha256: "b".repeat(64),
445 created_at: session.created_at.clone(),
446 event_frontier: 6,
447 });
448 save_move_operation_with(&connection, &intent).unwrap();
449 }
450 let connection = open(&path).unwrap();
451 let reaped = reap_finished_move_intents_with(&connection).unwrap();
452 assert_eq!(
453 reaped,
454 cases.iter().filter(|case| !case.4).count(),
455 "only the finished moves whose archive is gone are reaped"
456 );
457 for (session_id, _, _, _, survives) in cases {
458 assert_eq!(
459 load_move_operation_with(&connection, session_id)
460 .unwrap()
461 .is_some(),
462 survives,
463 "{session_id} row survival"
464 );
465 }
466 assert_eq!(
467 reap_finished_move_intents_with(&connection).unwrap(),
468 0,
469 "a second sweep finds nothing left to reap"
470 );
471 }
472
473 #[test]
474 fn concurrent_phase_save_cannot_erase_durable_cancellation() {
475 let directory = tempfile::tempdir().unwrap();
476 let path = directory.path().join("mj.sqlite3");
477 let session = super::super::tests::session("move-cancel", "project");
478 save_session_to(&path, &session).unwrap();
479 let connection = open(&path).unwrap();
480 let mut intent = operation(&session);
481 save_move_operation_with(&connection, &intent).unwrap();
482 let mut cancelled = intent.clone();
483 cancelled.cancellation_requested = true;
484 save_move_operation_with(&connection, &cancelled).unwrap();
485 intent.phase = MovePhase::ClosingSource;
486 save_move_operation_with(&connection, &intent).unwrap();
487 let restored = load_move_operation_with(&connection, &session.id)
488 .unwrap()
489 .unwrap();
490 assert!(restored.cancellation_requested);
491 assert_eq!(restored.phase, MovePhase::ClosingSource);
492 intent.operation_id = "explicit-new-operation".into();
493 save_move_operation_with(&connection, &intent).unwrap();
494 assert!(
495 !load_move_operation_with(&connection, &session.id)
496 .unwrap()
497 .unwrap()
498 .cancellation_requested
499 );
500 }
501
502 #[test]
503 fn cancelling_partial_queue_admission_does_not_release_its_archive() {
504 let session = super::super::tests::session("move-queue", "project");
505 let mut intent = operation(&session);
506 intent.phase = MovePhase::Cancelled;
507 intent.queue_admission_started = true;
508 assert!(intent.retains_checkpoint());
509 intent.queue_admission_finished = true;
510 assert!(!intent.retains_checkpoint());
511 }
512
513 #[test]
514 fn destination_record_install_keeps_drafts_and_titles_edited_during_move() {
515 let directory = tempfile::tempdir().unwrap();
516 let path = directory.path().join("mj.sqlite3");
517 let mut stale = super::super::tests::session("move-draft", "project");
518 save_session_to(&path, &stale).unwrap();
519 let connection = open(&path).unwrap();
520 save_move_operation_with(&connection, &operation(&stale)).unwrap();
521 connection.execute("UPDATE sessions SET draft_input='keep this draft', session_title_override='new title' WHERE session_id=?1", [&stale.id]).unwrap();
522 stale.last_profile = "destination".into();
523 stale.state = SessionState::Provisioning;
524 save_session_to(&path, &stale).unwrap();
525 let current = load_state_from(&path)
526 .unwrap()
527 .sessions
528 .remove(&stale.id)
529 .unwrap();
530 assert_eq!(current.draft_input, "keep this draft");
531 assert_eq!(current.session_title_override.as_deref(), Some("new title"));
532 assert_eq!(current.last_profile, "destination");
533 }
534}