1use super::*;
2
3pub fn save_session(session: &SessionRecord) -> Result<()> {
7 let session = session.clone();
8 submit_database_write("save_session", move |_| {
9 save_session_to(&database_path(), &session)
10 })
11}
12
13pub fn save_new_session(
14 session: &SessionRecord,
15 container_size: Option<(String, HostContainerSize)>,
16) -> Result<()> {
17 let session = session.clone();
18 submit_database_write("save_new_session", move |_| {
19 save_new_session_to(&database_path(), &session, container_size)
20 })
21}
22
23pub(super) fn save_new_session_to(
24 path: &Path,
25 session: &SessionRecord,
26 container_size: Option<(String, HostContainerSize)>,
27) -> Result<()> {
28 let mut connection = open(path)?;
29 let tx = connection.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
30 validate_session_record(session)?;
31 insert_session(&tx, session)?;
32 if let Some((host, size)) = container_size {
33 write_host_container_size(&tx, &host, size)?;
34 }
35 if session.harness_kind.supports_delegation_tools() {
36 tx.execute("INSERT INTO subagent_preference(singleton, policy) VALUES(1, ?1) ON CONFLICT(singleton) DO UPDATE SET policy = excluded.policy", [serde_json::to_string(&session.subagents.clone().unwrap_or_default())?])?;
37 }
38 tx.commit()?;
39 Ok(())
40}
41
42pub fn set_publication_assessment_if_current(
45 session_id: &str,
46 assessment: &mj_core::state::PublicationAssessment,
47) -> Result<bool> {
48 let session_id = session_id.to_owned();
49 let assessment = assessment.clone();
50 submit_database_write("set_publication_assessment_if_current", move |_| {
51 let connection = open(&database_path())?;
52 let updated = connection.execute(
53 "UPDATE sessions SET publication_json = ?2
54 WHERE session_id = ?1 AND state = 'stopped'
55 AND EXISTS (SELECT 1 FROM session_checkpoints c
56 WHERE c.session_id = ?1 AND c.sha256 = ?3)",
57 params![
58 session_id,
59 serde_json::to_string(&assessment)?,
60 assessment.checkpoint_sha256
61 ],
62 )?;
63 Ok(updated == 1)
64 })
65}
66
67pub fn save_subagent_session(
69 session: &SessionRecord,
70 subagent: &mj_core::subagent::SubagentRecord,
71) -> Result<()> {
72 let session = session.clone();
73 let subagent = subagent.clone();
74 submit_database_write("save_subagent_session", move |_| {
75 let mut connection = open(&database_path())?;
76 let tx = connection.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
77 insert_session(&tx, &session)?;
78 super::usage::record_child_accounting(&tx, &subagent)?;
79 tx.execute(
80 "INSERT INTO subagent_sessions(
81 child_session_id, parent_session_id, request_key, record_json
82 ) VALUES (?1, ?2, ?3, ?4)",
83 params![
84 subagent.child_session_id,
85 subagent.parent_session_id,
86 subagent.request_key,
87 serde_json::to_string(&subagent)?,
88 ],
89 )?;
90 tx.commit()?;
91 Ok(())
92 })
93}
94
95pub fn mark_subagent_turn_noticed(child_session_id: &str, turn: u64) -> Result<()> {
97 let child_session_id = child_session_id.to_owned();
98 submit_database_write("mark_subagent_turn_noticed", move |_| {
99 let mut relation = load_subagent(&child_session_id)?
100 .with_context(|| format!("unknown sub-agent session {child_session_id}"))?;
101 relation.noticed_turn = Some(turn);
102 let json = serde_json::to_string(&relation)?;
103 let connection = open(&database_path())?;
104 connection.execute(
105 "UPDATE subagent_sessions SET record_json = ?2 WHERE child_session_id = ?1",
106 params![child_session_id, json],
107 )?;
108 Ok(())
109 })
110}
111
112pub fn load_subagent(child_session_id: &str) -> Result<Option<mj_core::subagent::SubagentRecord>> {
113 let connection = open_reader(&database_path())?;
114 connection
115 .query_row(
116 "SELECT record_json FROM subagent_sessions WHERE child_session_id = ?1",
117 [child_session_id],
118 |row| row.get::<_, String>(0),
119 )
120 .optional()?
121 .map(|json| serde_json::from_str(&json).context("decode sub-agent record"))
122 .transpose()
123}
124
125pub fn list_subagents(parent_session_id: &str) -> Result<Vec<mj_core::subagent::SubagentRecord>> {
126 let connection = open_reader(&database_path())?;
127 let mut statement = connection.prepare(
128 "SELECT record_json FROM subagent_sessions
129 WHERE parent_session_id = ?1 ORDER BY rowid",
130 )?;
131 statement
132 .query_map([parent_session_id], |row| row.get::<_, String>(0))?
133 .map(|row| serde_json::from_str(&row?).context("decode sub-agent record"))
134 .collect()
135}
136
137pub fn record_stopped_subagents(
141 parent_session_id: &str,
142 stopped: &[mj_core::subagent::StoppedSubagent],
143) -> Result<()> {
144 let parent_session_id = parent_session_id.to_owned();
145 let stopped = stopped.to_vec();
146 submit_database_write("record_stopped_subagents", move |_| {
147 record_stopped_subagents_to(&database_path(), &parent_session_id, &stopped)
148 })
149}
150
151pub(super) fn record_stopped_subagents_to(
152 path: &Path,
153 parent_session_id: &str,
154 stopped: &[mj_core::subagent::StoppedSubagent],
155) -> Result<()> {
156 let mut connection = open(path)?;
157 let tx = connection.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
158 for child in stopped {
159 tx.execute(
160 "INSERT INTO stopped_subagents(parent_session_id, child_session_id, record_json)
161 VALUES (?1, ?2, ?3)
162 ON CONFLICT(parent_session_id, child_session_id) DO UPDATE SET
163 record_json = excluded.record_json",
164 params![
165 parent_session_id,
166 child.child_session_id,
167 serde_json::to_string(child)?
168 ],
169 )?;
170 }
171 tx.commit()?;
172 Ok(())
173}
174
175pub fn load_stopped_subagents(
178 parent_session_id: &str,
179) -> Result<Vec<mj_core::subagent::StoppedSubagent>> {
180 load_stopped_subagents_from(&database_path(), parent_session_id)
181}
182
183pub(super) fn load_stopped_subagents_from(
184 path: &Path,
185 parent_session_id: &str,
186) -> Result<Vec<mj_core::subagent::StoppedSubagent>> {
187 let connection = open_reader(path)?;
188 let mut statement = connection.prepare(
189 "SELECT record_json FROM stopped_subagents
190 WHERE parent_session_id = ?1 ORDER BY rowid",
191 )?;
192 statement
193 .query_map([parent_session_id], |row| row.get::<_, String>(0))?
194 .map(|row| serde_json::from_str(&row?).context("decode stopped sub-agent record"))
195 .collect()
196}
197
198pub fn clear_stopped_subagents(
202 parent_session_id: &str,
203 child_session_ids: &[String],
204) -> Result<()> {
205 let parent_session_id = parent_session_id.to_owned();
206 let child_session_ids = child_session_ids.to_vec();
207 submit_database_write("clear_stopped_subagents", move |_| {
208 clear_stopped_subagents_from(&database_path(), &parent_session_id, &child_session_ids)
209 })
210}
211
212pub(super) fn clear_stopped_subagents_from(
213 path: &Path,
214 parent_session_id: &str,
215 child_session_ids: &[String],
216) -> Result<()> {
217 let mut connection = open(path)?;
218 let tx = connection.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
219 for child_session_id in child_session_ids {
220 tx.execute(
221 "DELETE FROM stopped_subagents
222 WHERE parent_session_id = ?1 AND child_session_id = ?2",
223 params![parent_session_id, child_session_id],
224 )?;
225 }
226 tx.commit()?;
227 Ok(())
228}
229
230pub fn load_session_state(session_id: &str) -> Result<Option<SessionState>> {
235 let connection = open_reader(&database_path())?;
236 let stored = connection
237 .query_row(
238 "SELECT state FROM sessions WHERE session_id = ?1",
239 [session_id],
240 |row| row.get::<_, String>(0),
241 )
242 .optional()?;
243 Ok(stored.as_deref().map(stored_session_state))
244}
245
246pub fn load_subagent_report(child_session_id: &str) -> Result<mj_core::subagent::SubagentReport> {
247 load_subagent_report_from(&database_path(), child_session_id)
248}
249
250pub(super) fn load_subagent_report_from(
251 path: &Path,
252 child_session_id: &str,
253) -> Result<mj_core::subagent::SubagentReport> {
254 Ok(load_subagent_report_with(&*open_reader(path)?, child_session_id)?.unwrap_or_default())
255}
256
257pub(super) fn load_subagent_report_with(
259 connection: &Connection,
260 child_session_id: &str,
261) -> Result<Option<mj_core::subagent::SubagentReport>> {
262 let row = connection
263 .prepare_cached(
264 "SELECT handback_command_id, handback_message, handback_recorded_at_ms,
265 reminder_command_id, reminder_for_command_id, reminder_sent_at_ms,
266 reminder_failed_for_command_id, awaited_ordinal, report_dir
267 FROM subagent_handbacks WHERE child_session_id = ?1",
268 )?
269 .query_row([child_session_id], |row| {
270 Ok((
271 row.get::<_, Option<String>>(0)?,
272 row.get::<_, Option<String>>(1)?,
273 row.get::<_, Option<i64>>(2)?,
274 row.get::<_, Option<String>>(3)?,
275 row.get::<_, Option<String>>(4)?,
276 row.get::<_, Option<i64>>(5)?,
277 row.get::<_, Option<String>>(6)?,
278 row.get::<_, Option<i64>>(7)?,
279 row.get::<_, Option<String>>(8)?,
280 ))
281 })
282 .optional()?;
283 let Some((
284 handback_command,
285 handback_message,
286 handback_at,
287 reminder_command,
288 reminder_for,
289 reminder_at,
290 reminder_failed_for,
291 awaited_ordinal,
292 report_dir,
293 )) = row
294 else {
295 return Ok(None);
296 };
297 Ok(Some(mj_core::subagent::SubagentReport {
298 handback: match (handback_command, handback_message, handback_at) {
299 (Some(command_id), Some(message), Some(recorded_at_ms)) => {
300 Some(mj_core::subagent::SubagentHandback {
301 command_id,
302 message,
303 recorded_at_ms,
304 })
305 }
306 _ => None,
307 },
308 reminder: match (reminder_command, reminder_for, reminder_at) {
309 (Some(command_id), Some(for_command_id), Some(sent_at_ms)) => {
310 Some(mj_core::subagent::HandbackReminder {
311 command_id,
312 for_command_id,
313 sent_at_ms,
314 })
315 }
316 _ => None,
317 },
318 reminder_failed_for,
319 awaited_ordinal: awaited_ordinal.and_then(|ordinal| u64::try_from(ordinal).ok()),
320 report_dir,
321 }))
322}
323
324pub fn record_subagent_report_dir(child_session_id: &str, report_dir: &str) -> Result<()> {
327 let child_session_id = child_session_id.to_owned();
328 let report_dir = report_dir.to_owned();
329 submit_database_write("record_subagent_report_dir", move |_| {
330 record_subagent_report_dir_to(&database_path(), &child_session_id, &report_dir)
331 })
332}
333
334pub(super) fn record_subagent_report_dir_to(
335 path: &Path,
336 child_session_id: &str,
337 report_dir: &str,
338) -> Result<()> {
339 open(path)?.execute(
340 "INSERT INTO subagent_handbacks(child_session_id, report_dir)
341 SELECT ?1, ?2 WHERE EXISTS (
342 SELECT 1 FROM subagent_sessions WHERE child_session_id = ?1
343 )
344 ON CONFLICT(child_session_id) DO UPDATE SET report_dir = excluded.report_dir",
345 params![child_session_id, report_dir],
346 )?;
347 Ok(())
348}
349
350pub fn record_subagent_prompt(child_session_id: &str, ordinal: u64) -> Result<()> {
353 let child_session_id = child_session_id.to_owned();
354 submit_database_write("record_subagent_prompt", move |_| {
355 record_subagent_prompt_to(&database_path(), &child_session_id, ordinal)
356 })
357}
358
359pub(super) fn record_subagent_prompt_to(
360 path: &Path,
361 child_session_id: &str,
362 ordinal: u64,
363) -> Result<()> {
364 let ordinal = i64::try_from(ordinal).context("prompt ordinal exceeds the store's range")?;
365 open(path)?.execute(
366 "INSERT INTO subagent_handbacks(child_session_id, awaited_ordinal)
367 SELECT ?1, ?2 WHERE EXISTS (
368 SELECT 1 FROM subagent_sessions WHERE child_session_id = ?1
369 )
370 ON CONFLICT(child_session_id) DO UPDATE SET
371 awaited_ordinal = max(coalesce(awaited_ordinal, 0), excluded.awaited_ordinal)",
372 params![child_session_id, ordinal],
373 )?;
374 Ok(())
375}
376
377pub fn record_subagent_handback(
382 child_session_id: &str,
383 handback: &mj_core::subagent::SubagentHandback,
384) -> Result<bool> {
385 let child_session_id = child_session_id.to_owned();
386 let handback = handback.clone();
387 submit_database_write("record_subagent_handback", move |_| {
388 record_subagent_handback_to(&database_path(), &child_session_id, &handback)
389 })
390}
391
392pub(super) fn record_subagent_handback_to(
393 path: &Path,
394 child_session_id: &str,
395 handback: &mj_core::subagent::SubagentHandback,
396) -> Result<bool> {
397 let mut connection = open(path)?;
398 let tx = connection.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
399 let recorded_for: Option<String> = tx
400 .query_row(
401 "SELECT handback_command_id FROM subagent_handbacks WHERE child_session_id = ?1",
402 [child_session_id],
403 |row| row.get(0),
404 )
405 .optional()?
406 .flatten();
407 if recorded_for.as_deref() == Some(handback.command_id.as_str()) {
408 return Ok(false);
409 }
410 tx.execute(
411 "INSERT INTO subagent_handbacks(
412 child_session_id, handback_command_id, handback_message, handback_recorded_at_ms
413 ) VALUES (?1, ?2, ?3, ?4)
414 ON CONFLICT(child_session_id) DO UPDATE SET
415 handback_command_id = excluded.handback_command_id,
416 handback_message = excluded.handback_message,
417 handback_recorded_at_ms = excluded.handback_recorded_at_ms",
418 params![
419 child_session_id,
420 handback.command_id,
421 handback.message,
422 handback.recorded_at_ms
423 ],
424 )?;
425 tx.commit()?;
426 Ok(true)
427}
428
429pub fn record_handback_reminder(
431 child_session_id: &str,
432 reminder: &mj_core::subagent::HandbackReminder,
433) -> Result<()> {
434 let child_session_id = child_session_id.to_owned();
435 let reminder = reminder.clone();
436 submit_database_write("record_handback_reminder", move |_| {
437 open(&database_path())?.execute(
438 "INSERT INTO subagent_handbacks(
439 child_session_id, reminder_command_id, reminder_for_command_id, reminder_sent_at_ms
440 ) VALUES (?1, ?2, ?3, ?4)
441 ON CONFLICT(child_session_id) DO UPDATE SET
442 reminder_command_id = excluded.reminder_command_id,
443 reminder_for_command_id = excluded.reminder_for_command_id,
444 reminder_sent_at_ms = excluded.reminder_sent_at_ms",
445 params![
446 child_session_id,
447 reminder.command_id,
448 reminder.for_command_id,
449 reminder.sent_at_ms
450 ],
451 )?;
452 Ok(())
453 })
454}
455
456pub fn record_handback_reminder_failed(child_session_id: &str, for_command_id: &str) -> Result<()> {
459 let child_session_id = child_session_id.to_owned();
460 let for_command_id = for_command_id.to_owned();
461 submit_database_write("record_handback_reminder_failed", move |_| {
462 open(&database_path())?.execute(
463 "INSERT INTO subagent_handbacks(child_session_id, reminder_failed_for_command_id)
464 VALUES (?1, ?2)
465 ON CONFLICT(child_session_id) DO UPDATE SET
466 reminder_failed_for_command_id = excluded.reminder_failed_for_command_id",
467 params![child_session_id, for_command_id],
468 )?;
469 Ok(())
470 })
471}
472
473pub fn lookup_subagent_request(
474 parent_session_id: &str,
475 request_key: &str,
476) -> Result<Option<mj_core::subagent::SubagentRecord>> {
477 let connection = open_reader(&database_path())?;
478 connection
479 .query_row(
480 "SELECT record_json FROM subagent_sessions
481 WHERE parent_session_id = ?1 AND request_key = ?2",
482 params![parent_session_id, request_key],
483 |row| row.get::<_, String>(0),
484 )
485 .optional()?
486 .map(|json| serde_json::from_str(&json).context("decode sub-agent record"))
487 .transpose()
488}
489
490pub fn save_session_with_container_size(
493 session: &SessionRecord,
494 host: &str,
495 size: HostContainerSize,
496) -> Result<()> {
497 let session = session.clone();
498 let host = host.to_owned();
499 submit_database_write("save_session_with_container_size", move |_| {
500 save_session_with_container_size_to(&database_path(), &session, Some((&host, size)))
501 })
502}
503
504pub fn save_resumed_session(
507 session: &SessionRecord,
508 container_size: Option<(&str, HostContainerSize)>,
509) -> Result<()> {
510 let session = session.clone();
511 let container_size = container_size.map(|(host, size)| (host.to_owned(), size));
512 submit_database_write("save_resumed_session", move |_| {
513 save_resumed_session_to(
514 &database_path(),
515 &session,
516 container_size
517 .as_ref()
518 .map(|(host, size)| (host.as_str(), *size)),
519 )
520 })
521}
522
523pub(super) fn save_resumed_session_to(
524 path: &Path,
525 session: &SessionRecord,
526 container_size: Option<(&str, HostContainerSize)>,
527) -> Result<()> {
528 validate_session_record(session)?;
529 let mut connection = open(path)?;
530 let tx = connection.transaction()?;
531 let (bundle, workspace): (String, String) = tx.query_row(
532 "SELECT bundle_id, workspace_id FROM session_contexts WHERE session_id = ?1",
533 [&session.id],
534 |row| Ok((row.get(0)?, row.get(1)?)),
535 )?;
536 ensure!(
537 bundle == session.bundle_id && workspace == session.workspace_id,
538 "session {} context changed before resume publication",
539 session.id
540 );
541 update_lifecycle_fields(&tx, session)?;
542 tx.execute(
543 "UPDATE sessions SET native_session_id = ?2, container_cpus = ?3,
544 container_memory = ?4, container_workspace = ?5,
545 create_managed_worktree = ?6, launch_base = ?7, launch_branch = ?8,
546 checkout_json = ?9
547 WHERE session_id = ?1",
548 params![
549 session.id,
550 session.native_session_id,
551 session.container_cpus,
552 session.container_memory,
553 session
554 .container_workspace
555 .as_ref()
556 .map(|path| path.to_string_lossy().into_owned()),
557 session.create_managed_worktree,
558 session.launch_base,
559 session.launch_branch,
560 session
561 .checkout
562 .as_ref()
563 .map(serde_json::to_string)
564 .transpose()?,
565 ],
566 )?;
567 replace_mounts(&tx, &session.id, &session.additional_mounts)?;
568 replace_checkpoint(&tx, session)?;
569 if let Some((host, size)) = container_size {
570 write_host_container_size(&tx, host, size)?;
571 }
572 tx.commit()?;
573 Ok(())
574}
575
576pub fn save_lifecycle_session(session: &SessionRecord) -> Result<()> {
580 let session = session.clone();
581 submit_database_write("save_lifecycle_session", move |_| {
582 save_lifecycle_session_to(&database_path(), &session)
583 })
584}
585
586pub fn save_checkpointed_session(session: &SessionRecord) -> Result<()> {
589 let session = session.clone();
590 submit_database_write("save_checkpointed_session", move |_| {
591 save_checkpointed_session_to(&database_path(), &session)
592 })
593}
594
595pub fn recover_interrupted_checkpointing_sessions(updated_at: &str) -> Result<usize> {
599 let updated_at = updated_at.to_owned();
600 submit_database_write("recover_interrupted_checkpointing_sessions", move |_| {
601 recover_interrupted_checkpointing_sessions_to(&database_path(), &updated_at)
602 })
603}
604
605pub fn set_session_title_override(session_id: &str, title: &str, updated_at: &str) -> Result<()> {
608 let session_id = session_id.to_owned();
609 let title = title.to_owned();
610 let updated_at = updated_at.to_owned();
611 submit_database_write("set_session_title_override", move |_| {
612 set_session_title_override_to(&database_path(), &session_id, &title, &updated_at)
613 })
614}
615
616pub fn rename_profile_references(old_id: &str, new_id: &str) -> Result<usize> {
620 rename_session_reference("last_profile", old_id, new_id)
621}
622
623pub fn rename_target_references(old_id: &str, new_id: &str) -> Result<usize> {
626 rename_session_reference("target_template_id", old_id, new_id)
627}
628
629pub(super) fn rename_session_reference(
630 column: &'static str,
631 old_id: &str,
632 new_id: &str,
633) -> Result<usize> {
634 ensure!(
635 matches!(column, "last_profile" | "target_template_id"),
636 "unsupported session reference column"
637 );
638 let old_id = old_id.to_owned();
639 let new_id = new_id.to_owned();
640 submit_database_write("rename_session_reference", move |_| {
641 rename_session_reference_at(&database_path(), column, &old_id, &new_id)
642 })
643}
644
645pub(super) fn rename_session_reference_at(
646 path: &Path,
647 column: &str,
648 old_id: &str,
649 new_id: &str,
650) -> Result<usize> {
651 let mut connection = open(path)?;
652 let tx = connection.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
653 let changed = tx.execute(
654 &format!("UPDATE sessions SET {column} = ?2 WHERE {column} = ?1"),
655 params![old_id, new_id],
656 )?;
657 tx.commit()?;
658 Ok(changed)
659}
660
661pub fn set_session_archived(session_id: &str, archived: bool) -> Result<()> {
665 let session_id = session_id.to_owned();
666 submit_database_write("set_session_archived", move |_| {
667 set_session_archived_to(&database_path(), &session_id, archived)
668 })
669}
670
671pub fn mark_session_target_missing(
676 session_id: &str,
677 detail: &str,
678 updated_at: &str,
679) -> Result<Option<SessionState>> {
680 let session_id = session_id.to_owned();
681 let detail = detail.to_owned();
682 let updated_at = updated_at.to_owned();
683 submit_database_write("mark_session_target_missing", move |_| {
684 mark_session_target_missing_to(&database_path(), &session_id, &detail, &updated_at)
685 })
686}
687
688pub(super) fn mark_session_target_missing_to(
689 path: &Path,
690 session_id: &str,
691 detail: &str,
692 updated_at: &str,
693) -> Result<Option<SessionState>> {
694 mark_session_target_missing_if_current_to(path, session_id, detail, updated_at, None)
695}
696
697pub fn mark_session_target_missing_if_current(
700 session_id: &str,
701 detail: &str,
702 updated_at: &str,
703 observed_updated_at: &str,
704) -> Result<Option<SessionState>> {
705 let session_id = session_id.to_owned();
706 let detail = detail.to_owned();
707 let updated_at = updated_at.to_owned();
708 let observed_updated_at = observed_updated_at.to_owned();
709 submit_database_write("mark_session_target_missing_if_current", move |_| {
710 mark_session_target_missing_if_current_to(
711 &database_path(),
712 &session_id,
713 &detail,
714 &updated_at,
715 Some(&observed_updated_at),
716 )
717 })
718}
719
720pub(super) fn mark_session_target_missing_if_current_to(
721 path: &Path,
722 session_id: &str,
723 detail: &str,
724 updated_at: &str,
725 observed_updated_at: Option<&str>,
726) -> Result<Option<SessionState>> {
727 let mut connection = open(path)?;
728 let tx = connection.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
729 let previous_error = super::events::previous_session_error(&tx, session_id)?;
730 let changed = tx.execute(
731 "UPDATE sessions
732 SET state = CASE
733 WHEN EXISTS(
734 SELECT 1 FROM session_checkpoints
735 WHERE session_checkpoints.session_id = sessions.session_id
736 ) THEN 'error'
737 ELSE 'lost'
738 END,
739 last_error = ?2,
740 updated_at = ?3
741 WHERE session_id = ?1
742 AND (?4 IS NULL OR updated_at = ?4)
743 AND state IN ('provisioning', 'running', 'disconnected', 'error')",
744 params![session_id, detail, updated_at, observed_updated_at],
745 )?;
746 ensure!(changed <= 1, "updated {changed} sessions for {session_id}");
747 let state = if changed == 1 {
748 let stored: String = tx.query_row(
749 "SELECT state FROM sessions WHERE session_id = ?1",
750 [session_id],
751 |row| row.get(0),
752 )?;
753 Some(stored_session_state(&stored))
754 } else {
755 None
756 };
757 if changed == 1 && previous_error.as_deref() != Some(detail) {
758 super::events::insert_api_event(
759 &tx,
760 session_id,
761 Utc::now().timestamp_millis(),
762 &ApiEventData::SessionFault {
763 reason: mj_core::event_outcome::OutcomeReason::RuntimeUnavailable,
764 message: detail.into(),
765 command_id: None,
766 },
767 )?;
768 }
769 tx.commit()?;
770 Ok(state)
771}
772
773pub(super) fn set_session_archived_to(path: &Path, session_id: &str, archived: bool) -> Result<()> {
774 let connection = open(path)?;
775 let changed = connection.execute(
776 "UPDATE sessions SET archived = ?2 WHERE session_id = ?1",
777 params![session_id, archived],
778 )?;
779 if changed != 1 {
780 bail!("unknown session {session_id}");
781 }
782 Ok(())
783}
784
785pub fn hidden_native_sessions() -> Result<BTreeSet<(mj_core::config::HarnessKind, String)>> {
788 hidden_native_sessions_from(&database_path())
789}
790
791pub(super) fn hidden_native_sessions_from(
792 path: &Path,
793) -> Result<BTreeSet<(mj_core::config::HarnessKind, String)>> {
794 let connection = open_reader(path)?;
795 let mut statement =
796 connection.prepare("SELECT harness_kind, native_session_id FROM hidden_native_sessions")?;
797 let rows = statement.query_map([], |row| {
798 Ok((row.get::<_, String>(0)?, row.get::<_, String>(1)?))
799 })?;
800 let mut hidden = BTreeSet::new();
801 for row in rows {
802 let (harness, native_session_id) = row?;
803 match harness.parse::<mj_core::config::HarnessKind>() {
806 Ok(harness) => {
807 hidden.insert((harness, native_session_id));
808 }
809 Err(_) => tracing::warn!(
810 harness = %harness,
811 "ignoring a hidden native session for a harness that is no longer supported"
812 ),
813 }
814 }
815 Ok(hidden)
816}
817
818pub fn set_native_session_hidden(
820 harness: mj_core::config::HarnessKind,
821 native_session_id: &str,
822 hidden: bool,
823) -> Result<()> {
824 let native_session_id = native_session_id.to_owned();
825 submit_database_write("set_native_session_hidden", move |_| {
826 set_native_session_hidden_to(&database_path(), harness, &native_session_id, hidden)
827 })
828}
829
830pub(super) fn set_native_session_hidden_to(
831 path: &Path,
832 harness: mj_core::config::HarnessKind,
833 native_session_id: &str,
834 hidden: bool,
835) -> Result<()> {
836 if native_session_id.trim().is_empty() {
837 bail!("native session id is empty");
838 }
839 let connection = open(path)?;
840 if hidden {
841 connection.execute(
842 "INSERT INTO hidden_native_sessions(harness_kind, native_session_id, hidden_at)
843 VALUES (?1, ?2, ?3)
844 ON CONFLICT(harness_kind, native_session_id) DO NOTHING",
845 params![harness.id(), native_session_id, Utc::now().to_rfc3339()],
846 )?;
847 } else {
848 connection.execute(
849 "DELETE FROM hidden_native_sessions
850 WHERE harness_kind = ?1 AND native_session_id = ?2",
851 params![harness.id(), native_session_id],
852 )?;
853 }
854 Ok(())
855}
856
857pub fn set_session_container_settings(
861 session_id: &str,
862 cpus: Option<&str>,
863 memory: Option<&str>,
864 mounts: &[AdditionalMount],
865 updated_at: &str,
866) -> Result<()> {
867 let session_id = session_id.to_owned();
868 let cpus = cpus.map(str::to_owned);
869 let memory = memory.map(str::to_owned);
870 let mounts = mounts.to_vec();
871 let updated_at = updated_at.to_owned();
872 submit_database_write("set_session_container_settings", move |_| {
873 set_session_container_settings_to(
874 &database_path(),
875 &session_id,
876 cpus.as_deref(),
877 memory.as_deref(),
878 &mounts,
879 &updated_at,
880 )
881 })
882}
883
884pub(super) fn set_session_container_settings_to(
885 path: &Path,
886 session_id: &str,
887 cpus: Option<&str>,
888 memory: Option<&str>,
889 mounts: &[AdditionalMount],
890 updated_at: &str,
891) -> Result<()> {
892 if updated_at.trim().is_empty() {
893 bail!("session update timestamp is empty");
894 }
895 crate::targets::validate_additional_mounts(mounts)?;
896 let mut connection = open(path)?;
897 let tx = connection.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
898 let changed = tx.execute(
899 "UPDATE sessions
900 SET container_cpus = ?2, container_memory = ?3, updated_at = ?4
901 WHERE session_id = ?1",
902 params![session_id, cpus, memory, updated_at],
903 )?;
904 if changed != 1 {
905 bail!("unknown session {session_id}");
906 }
907 replace_mounts(&tx, session_id, mounts)?;
908 tx.commit()?;
909 Ok(())
910}
911
912pub(super) fn set_session_title_override_to(
913 path: &Path,
914 session_id: &str,
915 title: &str,
916 updated_at: &str,
917) -> Result<()> {
918 if title.trim().is_empty() {
919 bail!("session title is empty");
920 }
921 if updated_at.trim().is_empty() {
922 bail!("session update timestamp is empty");
923 }
924 let connection = open(path)?;
925 let changed = connection.execute(
926 "UPDATE sessions
927 SET session_title_override = ?2, updated_at = ?3
928 WHERE session_id = ?1",
929 params![session_id, title, updated_at],
930 )?;
931 if changed != 1 {
932 bail!("unknown session {session_id}");
933 }
934 Ok(())
935}
936
937pub fn set_session_acp_title(session_id: &str, title: Option<&str>) -> Result<()> {
940 let session_id = session_id.to_owned();
941 let title = title.map(str::to_owned);
942 submit_database_write("set_session_acp_title", move |_| {
943 set_session_acp_title_to(&database_path(), &session_id, title.as_deref())
944 })
945}
946
947pub(super) fn set_session_acp_title_to(
948 path: &Path,
949 session_id: &str,
950 title: Option<&str>,
951) -> Result<()> {
952 if title.is_some_and(|title| title.trim().is_empty()) {
953 bail!("ACP session title is empty");
954 }
955 let title = title.and_then(mj_core::state::normalize_session_title);
956 let connection = open(path)?;
957 let changed = connection.execute(
958 "UPDATE sessions SET acp_session_title = ?2 WHERE session_id = ?1",
959 params![session_id, title],
960 )?;
961 if changed != 1 {
962 bail!("unknown session {session_id}");
963 }
964 Ok(())
965}
966
967pub fn mark_session_worker_connected(
970 session_id: &str,
971 native_session_id: Option<&str>,
972 updated_at: &str,
973) -> Result<()> {
974 let session_id = session_id.to_owned();
975 let native_session_id = native_session_id.map(str::to_owned);
976 let updated_at = updated_at.to_owned();
977 submit_database_write("mark_session_worker_connected", move |_| {
978 mark_session_worker_connected_to(
979 &database_path(),
980 &session_id,
981 native_session_id.as_deref(),
982 &updated_at,
983 )
984 })
985}
986
987pub fn adopt_native_session_id(session_id: &str, native_session_id: &str) -> Result<()> {
991 let session_id = session_id.to_owned();
992 let native_session_id = native_session_id.to_owned();
993 submit_database_write("adopt_native_session_id", move |_| {
994 adopt_native_session_id_to(&database_path(), &session_id, &native_session_id)
995 })
996}
997
998pub(super) fn adopt_native_session_id_to(
999 path: &Path,
1000 session_id: &str,
1001 native_session_id: &str,
1002) -> Result<()> {
1003 let connection = open(path)?;
1004 let changed = connection.execute(
1005 "UPDATE sessions SET native_session_id = ?2 WHERE session_id = ?1",
1006 params![session_id, native_session_id],
1007 )?;
1008 if changed != 1 {
1009 bail!("unknown session {session_id}");
1010 }
1011 Ok(())
1012}
1013
1014pub(super) fn mark_session_worker_connected_to(
1015 path: &Path,
1016 session_id: &str,
1017 native_session_id: Option<&str>,
1018 updated_at: &str,
1019) -> Result<()> {
1020 if updated_at.trim().is_empty() {
1021 bail!("worker connection timestamp is empty");
1022 }
1023 let connection = open(path)?;
1024 let changed = connection.execute(
1025 "UPDATE sessions
1026 SET state = 'running',
1027 native_session_id = coalesce(?2, native_session_id),
1028 updated_at = ?3,
1029 last_error = NULL
1030 WHERE session_id = ?1",
1031 params![session_id, native_session_id, updated_at],
1032 )?;
1033 if changed != 1 {
1034 bail!("unknown session {session_id}");
1035 }
1036 Ok(())
1037}
1038
1039pub(super) fn recover_interrupted_checkpointing_sessions_to(
1040 path: &Path,
1041 updated_at: &str,
1042) -> Result<usize> {
1043 ensure!(
1044 !updated_at.trim().is_empty(),
1045 "checkpoint recovery timestamp is empty"
1046 );
1047 let mut connection = open(path)?;
1048 let tx = connection.transaction()?;
1049 let ids = {
1050 let mut query = tx.prepare("SELECT session_id FROM checkpoint_operations")?;
1051 query
1052 .query_map([], |row| row.get::<_, String>(0))?
1053 .collect::<rusqlite::Result<Vec<_>>>()?
1054 };
1055 for id in ids {
1056 super::events::finish_checkpoint_operation(
1057 &tx,
1058 &id,
1059 mj_core::event_outcome::CommandResultKind::Failed,
1060 Some(mj_core::event_outcome::OutcomeReason::ControllerRestarted),
1061 Some("Checkpoint operation interrupted by controller restart".into()),
1062 )?;
1063 }
1064 let changed = tx.execute("UPDATE sessions SET state = 'running', updated_at = ?1, last_checkpoint_error = ?2 WHERE state = 'checkpointing'",
1065 params![updated_at, "checkpointing was interrupted by a controller restart; the target was left running"])?;
1066 tx.commit()?;
1067 Ok(changed)
1068}
1069
1070pub(super) fn save_session_to(path: &Path, session: &SessionRecord) -> Result<()> {
1071 save_session_with_container_size_to(path, session, None)
1072}
1073
1074pub(super) fn save_session_with_container_size_to(
1075 path: &Path,
1076 session: &SessionRecord,
1077 container_size: Option<(&str, HostContainerSize)>,
1078) -> Result<()> {
1079 validate_session_record(session)?;
1080
1081 let mut connection = open(path)?;
1082 let tx = connection.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
1083 if let Some(existing_bundle) = tx
1084 .query_row(
1085 "SELECT bundle_id FROM session_contexts WHERE session_id = ?1",
1086 [session.id.as_str()],
1087 |row| row.get::<_, String>(0),
1088 )
1089 .optional()?
1090 && existing_bundle
1091 != super::projects::session_project_id(&tx, &session.id, &session.bundle_id)?
1092 {
1093 bail!(
1094 "session {} was already associated with bundle {}, not {}",
1095 session.id,
1096 existing_bundle,
1097 session.bundle_id
1098 );
1099 }
1100 let mut session = session.clone();
1101 let moving: bool = tx.query_row(
1102 "SELECT EXISTS(SELECT 1 FROM session_moves WHERE session_id=?1
1103 AND json_extract(operation_json, '$.phase') IN ('preparing','closing_source','resuming_destination','starting_queue'))",
1104 [&session.id], |row| row.get(0),
1105 )?;
1106 if moving {
1107 let (draft, title, acp_title, viewed, archived) = tx.query_row(
1111 "SELECT draft_input, session_title_override, acp_session_title, viewed_through_event_ordinal, archived
1112 FROM sessions WHERE session_id=?1", [&session.id], |row| Ok((
1113 row.get::<_, String>(0)?, row.get::<_, Option<String>>(1)?, row.get::<_, Option<String>>(2)?,
1114 row.get::<_, u64>(3)?, row.get::<_, bool>(4)?,
1115 )),
1116 )?;
1117 session.draft_input = draft;
1118 session.session_title_override = title;
1119 session.acp_session_title = acp_title;
1120 session.viewed_through_event_ordinal = viewed;
1121 session.archived = archived;
1122 }
1123 insert_session(&tx, &session)?;
1124 if let Some((host, size)) = container_size {
1125 write_host_container_size(&tx, host, size)?;
1126 }
1127 tx.commit()?;
1128 Ok(())
1129}
1130
1131pub(super) fn save_lifecycle_session_to(path: &Path, session: &SessionRecord) -> Result<()> {
1132 validate_session_record(session)?;
1133
1134 let mut connection = open(path)?;
1135 let tx = connection.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
1136 update_lifecycle_fields(&tx, session)?;
1137 tx.commit()?;
1138 Ok(())
1139}
1140
1141pub(crate) fn save_startup_cleanup_outcome(session: &SessionRecord, archive: bool) -> Result<()> {
1143 let session = session.clone();
1144 submit_database_write("save_startup_cleanup_outcome", move |_| {
1145 validate_session_record(&session)?;
1146 let mut connection = open(&database_path())?;
1147 let tx = connection.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
1148 update_lifecycle_fields(&tx, &session)?;
1149 if archive {
1150 tx.execute(
1151 "UPDATE sessions SET archived=1 WHERE session_id=?1",
1152 [&session.id],
1153 )?;
1154 }
1155 tx.commit()?;
1156 Ok(())
1157 })
1158}
1159
1160pub(super) fn save_checkpointed_session_to(path: &Path, session: &SessionRecord) -> Result<()> {
1161 validate_session_record(session)?;
1162
1163 let mut connection = open(path)?;
1164 let tx = connection.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
1165 update_lifecycle_fields(&tx, session)?;
1166 tx.execute(
1167 "UPDATE sessions SET native_session_id = ?2 WHERE session_id = ?1",
1168 params![session.id, session.native_session_id],
1169 )?;
1170 replace_checkpoint(&tx, session)?;
1171 tx.commit()?;
1172 Ok(())
1173}
1174
1175pub(super) fn validate_session_record(session: &SessionRecord) -> Result<()> {
1176 let mut validation = State::default();
1177 validation
1178 .sessions
1179 .insert(session.id.clone(), session.clone());
1180 validation.validate()
1181}
1182
1183pub fn delete_session(session_id: &str) -> Result<()> {
1186 let session_id = session_id.to_owned();
1187 submit_database_write("delete_session", move |_| {
1188 delete_session_from(&database_path(), &session_id)
1189 })
1190}
1191
1192pub(super) fn delete_session_from(path: &Path, session_id: &str) -> Result<()> {
1193 let connection = open(path)?;
1194 connection.execute("DELETE FROM sessions WHERE session_id = ?1", [session_id])?;
1195 Ok(())
1196}
1197
1198pub fn set_session_draft_input(session_id: &str, draft: &str) -> Result<()> {
1202 let session_id = session_id.to_owned();
1203 let draft = draft.to_owned();
1204 submit_database_write("set_session_draft_input", move |_| {
1205 set_session_draft_input_at(&database_path(), &session_id, &draft)
1206 })
1207}
1208
1209pub(super) fn set_session_draft_input_at(path: &Path, session_id: &str, draft: &str) -> Result<()> {
1210 let mut connection = open(path)?;
1211 let tx = connection.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
1212 let updated = tx.execute(
1213 "UPDATE sessions SET draft_input = ?2 WHERE session_id = ?1",
1214 params![session_id, draft],
1215 )?;
1216 ensure!(updated == 1, "unknown session {session_id}");
1217 tx.commit()?;
1218 Ok(())
1219}
1220
1221pub fn append_session_draft_input(session_id: &str, text: &str) -> Result<()> {
1224 let session_id = session_id.to_owned();
1225 let text = text.to_owned();
1226 submit_database_write("append_session_draft_input", move |connection| {
1227 let updated = connection.execute(
1228 "UPDATE sessions SET draft_input = CASE WHEN ?2 = '' THEN draft_input
1229 WHEN draft_input = '' THEN ?2 ELSE draft_input || char(10) || char(10) || ?2 END
1230 WHERE session_id = ?1",
1231 params![session_id, text],
1232 )?;
1233 ensure!(updated == 1, "unknown session {session_id}");
1234 Ok(())
1235 })
1236}
1237
1238pub fn clear_session_draft_input_if_matches(session_id: &str, expected: &str) -> Result<()> {
1240 let session_id = session_id.to_owned();
1241 let expected = expected.to_owned();
1242 submit_database_write("clear_session_draft_input_if_matches", move |connection| {
1243 connection.execute(
1244 "UPDATE sessions SET draft_input = '' WHERE session_id = ?1 AND draft_input = ?2",
1245 params![session_id, expected],
1246 )?;
1247 Ok(())
1248 })
1249}
1250
1251pub fn record_recovery_success(
1252 session_id: &str,
1253 native_session_id: &str,
1254 checkpoint: &CheckpointMetadata,
1255) -> Result<()> {
1256 let session_id = session_id.to_owned();
1257 let native_session_id = native_session_id.to_owned();
1258 let checkpoint = checkpoint.clone();
1259 submit_database_write("record_recovery_success", move |_| {
1260 record_recovery_success_to(
1261 &database_path(),
1262 &session_id,
1263 &native_session_id,
1264 &checkpoint,
1265 )
1266 })
1267}
1268
1269pub(super) fn record_recovery_success_to(
1270 path: &Path,
1271 session_id: &str,
1272 native_session_id: &str,
1273 checkpoint: &CheckpointMetadata,
1274) -> Result<()> {
1275 let mut connection = open(path)?;
1276 let tx = connection.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
1277 let changed = tx.execute(
1278 "UPDATE sessions
1279 SET native_session_id = ?2, last_checkpoint_error = NULL
1280 WHERE session_id = ?1",
1281 params![session_id, native_session_id],
1282 )?;
1283 if changed != 1 {
1284 bail!("unknown session {session_id}");
1285 }
1286 tx.execute(
1287 "INSERT INTO session_checkpoints(
1288 session_id, archive_path, sha256, created_at, event_frontier
1289 ) VALUES (?1,?2,?3,?4,?5)
1290 ON CONFLICT(session_id) DO UPDATE SET
1291 archive_path = excluded.archive_path,
1292 sha256 = excluded.sha256,
1293 created_at = excluded.created_at,
1294 event_frontier = excluded.event_frontier",
1295 params![
1296 session_id,
1297 path_to_blob(&checkpoint.archive_path),
1298 checkpoint.sha256,
1299 checkpoint.created_at,
1300 checkpoint.event_frontier,
1301 ],
1302 )?;
1303 tx.commit()?;
1304 Ok(())
1305}
1306
1307pub fn record_recovery_success_if_current(
1308 session_id: &str,
1309 expected_target: &TargetLocator,
1310 expected_checkpoint: Option<&CheckpointMetadata>,
1311 native_session_id: &str,
1312 checkpoint: &CheckpointMetadata,
1313) -> Result<bool> {
1314 let session_id = session_id.to_owned();
1315 let expected_target = expected_target.clone();
1316 let expected_checkpoint = expected_checkpoint.cloned();
1317 let native_session_id = native_session_id.to_owned();
1318 let checkpoint = checkpoint.clone();
1319 submit_database_write("record_recovery_success_if_current", move |connection| {
1320 record_recovery_success_if_current_with(
1321 connection,
1322 &session_id,
1323 &expected_target,
1324 expected_checkpoint.as_ref(),
1325 &native_session_id,
1326 &checkpoint,
1327 )
1328 })
1329}
1330
1331fn record_recovery_success_if_current_with(
1332 connection: &mut Connection,
1333 session_id: &str,
1334 expected_target: &TargetLocator,
1335 expected_checkpoint: Option<&CheckpointMetadata>,
1336 native_session_id: &str,
1337 checkpoint: &CheckpointMetadata,
1338) -> Result<bool> {
1339 let tx = connection.transaction()?;
1340 let Some(mut current) = load_session_with(&tx, session_id)? else {
1341 return Ok(false);
1342 };
1343 if current.target.as_ref() != Some(expected_target)
1344 || current.checkpoint.as_ref() != expected_checkpoint
1345 {
1346 return Ok(false);
1347 }
1348 tx.execute(
1349 "UPDATE sessions SET native_session_id = ?2, last_checkpoint_error = NULL
1350 WHERE session_id = ?1",
1351 params![session_id, native_session_id],
1352 )?;
1353 current.checkpoint = Some(checkpoint.clone());
1354 replace_checkpoint(&tx, ¤t)?;
1355 tx.commit()?;
1356 Ok(true)
1357}
1358
1359pub fn record_recovery_failure(session_id: &str, detail: &str) -> Result<()> {
1360 let session_id = session_id.to_owned();
1361 let detail = detail.to_owned();
1362 submit_database_write("record_recovery_failure", move |_| {
1363 record_recovery_failure_to(&database_path(), &session_id, &detail)
1364 })
1365}
1366
1367pub fn record_recovery_failure_if_current(
1370 session_id: &str,
1371 expected_target: &TargetLocator,
1372 expected_checkpoint: Option<&CheckpointMetadata>,
1373 detail: &str,
1374) -> Result<bool> {
1375 let session_id = session_id.to_owned();
1376 let expected_target = expected_target.clone();
1377 let expected_checkpoint = expected_checkpoint.cloned();
1378 let detail = detail.to_owned();
1379 submit_database_write("record_recovery_failure_if_current", move |connection| {
1380 record_recovery_failure_if_current_with(
1381 connection,
1382 &session_id,
1383 &expected_target,
1384 expected_checkpoint.as_ref(),
1385 &detail,
1386 )
1387 })
1388}
1389
1390fn record_recovery_failure_if_current_with(
1391 connection: &mut Connection,
1392 session_id: &str,
1393 expected_target: &TargetLocator,
1394 expected_checkpoint: Option<&CheckpointMetadata>,
1395 detail: &str,
1396) -> Result<bool> {
1397 let tx = connection.transaction()?;
1398 let Some(current) = load_session_with(&tx, session_id)? else {
1399 return Ok(false);
1400 };
1401 if current.target.as_ref() != Some(expected_target)
1402 || current.checkpoint.as_ref() != expected_checkpoint
1403 {
1404 return Ok(false);
1405 }
1406 tx.execute(
1407 "UPDATE sessions SET last_checkpoint_error = ?2 WHERE session_id = ?1",
1408 params![session_id, detail],
1409 )?;
1410 tx.commit()?;
1411 Ok(true)
1412}
1413
1414pub(super) fn record_recovery_failure_to(
1415 path: &Path,
1416 session_id: &str,
1417 detail: &str,
1418) -> Result<()> {
1419 let connection = open(path)?;
1420 let changed = connection.execute(
1421 "UPDATE sessions SET last_checkpoint_error = ?2 WHERE session_id = ?1",
1422 params![session_id, detail],
1423 )?;
1424 if changed != 1 {
1425 bail!("unknown session {session_id}");
1426 }
1427 Ok(())
1428}
1429
1430pub fn rebind_session_bundle(session_id: &str, bundle_id: &str) -> Result<()> {
1437 let session_id = session_id.to_owned();
1438 let bundle_id = bundle_id.to_owned();
1439 submit_database_write("rebind_session_bundle", move |_| {
1440 rebind_session_bundle_to(&database_path(), &session_id, &bundle_id)
1441 })
1442}
1443
1444pub(super) fn rebind_session_bundle_to(
1445 path: &Path,
1446 session_id: &str,
1447 bundle_id: &str,
1448) -> Result<()> {
1449 let mut connection = open(path)?;
1450 let tx = connection.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
1451 let canonical_bundle = super::projects::canonical_project_id(&tx, bundle_id)?;
1452 tx.execute("INSERT OR IGNORE INTO project_session_aliases(session_id,bundle_id) SELECT session_id,bundle_id FROM session_contexts WHERE session_id=?1",[session_id])?;
1453 let changed = tx.execute(
1454 "UPDATE session_contexts SET bundle_id = ?2 WHERE session_id = ?1",
1455 params![session_id, canonical_bundle],
1456 )?;
1457 if changed == 0 {
1458 tx.execute(
1459 "INSERT INTO session_contexts(session_id, bundle_id, created_at) VALUES (?1, ?2, ?3)",
1460 params![session_id, canonical_bundle, Utc::now().to_rfc3339()],
1461 )?;
1462 }
1463 tx.commit()?;
1464 Ok(())
1465}
1466
1467#[cfg(test)]
1468mod ownership_tests {
1469 use super::*;
1470
1471 #[test]
1472 fn resume_and_rollback_preserve_edits_committed_after_their_snapshot() {
1473 let directory = tempfile::tempdir().unwrap();
1474 let path = directory.path().join("resume.sqlite3");
1475 let original = super::super::tests::session("resume-owner", "project");
1476 save_session_to(&path, &original).unwrap();
1477 let connection = open(&path).unwrap();
1478 connection
1479 .execute(
1480 "UPDATE sessions SET session_title_override = 'renamed during resume',
1481 acp_session_title = 'new worker title', draft_input = 'new draft',
1482 archived = 1, viewed_through_event_ordinal = 99 WHERE session_id = ?1",
1483 [&original.id],
1484 )
1485 .unwrap();
1486 let mut provisioning = original.clone();
1487 provisioning.state = SessionState::Provisioning;
1488 provisioning.target = None;
1489 provisioning.native_session_id = Some("resumed-native".into());
1490 provisioning.additional_mounts.clear();
1491 save_resumed_session_to(&path, &provisioning, None).unwrap();
1492 let current = load_session_with(&connection, &original.id)
1493 .unwrap()
1494 .unwrap();
1495 assert_eq!(current.state, SessionState::Provisioning);
1496 assert_eq!(current.native_session_id, provisioning.native_session_id);
1497 assert!(current.additional_mounts.is_empty());
1498 assert_eq!(
1499 current.session_title_override.as_deref(),
1500 Some("renamed during resume")
1501 );
1502 assert_eq!(
1503 current.acp_session_title.as_deref(),
1504 Some("new worker title")
1505 );
1506 assert_eq!(current.draft_input, "new draft");
1507 assert!(current.archived);
1508 assert_eq!(current.viewed_through_event_ordinal, 99);
1509
1510 save_resumed_session_to(&path, &original, None).unwrap();
1512 let rolled_back = load_session_with(&connection, &original.id)
1513 .unwrap()
1514 .unwrap();
1515 assert_eq!(rolled_back.state, original.state);
1516 assert_eq!(rolled_back.additional_mounts, original.additional_mounts);
1517 assert_eq!(
1518 rolled_back.session_title_override,
1519 current.session_title_override
1520 );
1521 assert_eq!(rolled_back.acp_session_title, current.acp_session_title);
1522 assert_eq!(rolled_back.draft_input, current.draft_input);
1523 assert_eq!(rolled_back.archived, current.archived);
1524 assert_eq!(rolled_back.viewed_through_event_ordinal, 99);
1525 }
1526
1527 #[test]
1528 fn recovery_completion_cannot_replace_a_newer_checkpoint_or_target() {
1529 let directory = tempfile::tempdir().unwrap();
1530 let path = directory.path().join("recovery.sqlite3");
1531 let original = super::super::tests::session("recovery-owner", "project");
1532 save_session_to(&path, &original).unwrap();
1533 let target = original.target.as_ref().unwrap();
1534 let previous = original.checkpoint.as_ref().unwrap();
1535 let mut next = previous.clone();
1536 next.sha256 = "b".repeat(64);
1537 next.event_frontier += 1;
1538 let mut connection = open(&path).unwrap();
1539 assert!(
1540 record_recovery_success_if_current_with(
1541 &mut connection,
1542 &original.id,
1543 target,
1544 Some(previous),
1545 "new-native",
1546 &next,
1547 )
1548 .unwrap()
1549 );
1550 assert!(
1551 !record_recovery_failure_if_current_with(
1552 &mut connection,
1553 &original.id,
1554 target,
1555 Some(previous),
1556 "old failure",
1557 )
1558 .unwrap()
1559 );
1560 assert!(
1561 !record_recovery_success_if_current_with(
1562 &mut connection,
1563 &original.id,
1564 target,
1565 Some(previous),
1566 "old-native",
1567 previous,
1568 )
1569 .unwrap()
1570 );
1571 let current = load_session_with(&connection, &original.id)
1572 .unwrap()
1573 .unwrap();
1574 assert_eq!(current.checkpoint.as_ref(), Some(&next));
1575 assert_eq!(current.native_session_id.as_deref(), Some("new-native"));
1576 assert_eq!(current.last_checkpoint_error, None);
1577
1578 let mut replacement = current.clone();
1579 replacement.target = None;
1580 replacement.state = SessionState::Stopped;
1581 save_resumed_session_to(&path, &replacement, None).unwrap();
1582 assert!(
1583 !record_recovery_failure_if_current_with(
1584 &mut connection,
1585 &original.id,
1586 target,
1587 Some(&next),
1588 "old target failure",
1589 )
1590 .unwrap()
1591 );
1592 assert!(
1593 !record_recovery_success_if_current_with(
1594 &mut connection,
1595 &original.id,
1596 target,
1597 Some(&next),
1598 "old-native",
1599 previous,
1600 )
1601 .unwrap()
1602 );
1603 }
1604
1605 #[test]
1606 fn current_recovery_failure_settles_without_resurrecting_a_removed_session() {
1607 let directory = tempfile::tempdir().unwrap();
1608 let path = directory.path().join("recovery.sqlite3");
1609 let original = super::super::tests::session("recovery-owner", "project");
1610 save_session_to(&path, &original).unwrap();
1611 let mut connection = open(&path).unwrap();
1612 assert!(
1613 record_recovery_failure_if_current_with(
1614 &mut connection,
1615 &original.id,
1616 original.target.as_ref().unwrap(),
1617 original.checkpoint.as_ref(),
1618 "current failure",
1619 )
1620 .unwrap()
1621 );
1622 assert_eq!(
1623 load_session_with(&connection, &original.id)
1624 .unwrap()
1625 .unwrap()
1626 .last_checkpoint_error
1627 .as_deref(),
1628 Some("current failure")
1629 );
1630 connection
1631 .execute("DELETE FROM sessions WHERE session_id = ?1", [&original.id])
1632 .unwrap();
1633 assert!(
1634 !record_recovery_failure_if_current_with(
1635 &mut connection,
1636 &original.id,
1637 original.target.as_ref().unwrap(),
1638 original.checkpoint.as_ref(),
1639 "late failure",
1640 )
1641 .unwrap()
1642 );
1643 assert!(
1644 load_session_with(&connection, &original.id)
1645 .unwrap()
1646 .is_none()
1647 );
1648 assert!(save_resumed_session_to(&path, &original, None).is_err());
1649 }
1650}