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 set_publication_assessment_if_current(
16 session_id: &str,
17 assessment: &mj_core::state::PublicationAssessment,
18) -> Result<bool> {
19 let session_id = session_id.to_owned();
20 let assessment = assessment.clone();
21 submit_database_write("set_publication_assessment_if_current", move |_| {
22 let connection = open(&database_path())?;
23 let updated = connection.execute(
24 "UPDATE sessions SET publication_json = ?2
25 WHERE session_id = ?1 AND state = 'stopped'
26 AND EXISTS (SELECT 1 FROM session_checkpoints c
27 WHERE c.session_id = ?1 AND c.sha256 = ?3)",
28 params![
29 session_id,
30 serde_json::to_string(&assessment)?,
31 assessment.checkpoint_sha256
32 ],
33 )?;
34 Ok(updated == 1)
35 })
36}
37
38pub fn save_subagent_session(
40 session: &SessionRecord,
41 subagent: &mj_core::subagent::SubagentRecord,
42) -> Result<()> {
43 let session = session.clone();
44 let subagent = subagent.clone();
45 submit_database_write("save_subagent_session", move |_| {
46 let mut connection = open(&database_path())?;
47 let tx = connection.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
48 insert_session(&tx, &session)?;
49 tx.execute(
50 "INSERT INTO subagent_sessions(
51 child_session_id, parent_session_id, request_key, record_json
52 ) VALUES (?1, ?2, ?3, ?4)",
53 params![
54 subagent.child_session_id,
55 subagent.parent_session_id,
56 subagent.request_key,
57 serde_json::to_string(&subagent)?,
58 ],
59 )?;
60 tx.commit()?;
61 Ok(())
62 })
63}
64
65pub fn mark_subagent_turn_noticed(child_session_id: &str, turn: u64) -> Result<()> {
67 let child_session_id = child_session_id.to_owned();
68 submit_database_write("mark_subagent_turn_noticed", move |_| {
69 let mut relation = load_subagent(&child_session_id)?
70 .with_context(|| format!("unknown sub-agent session {child_session_id}"))?;
71 relation.noticed_turn = Some(turn);
72 let json = serde_json::to_string(&relation)?;
73 let connection = open(&database_path())?;
74 connection.execute(
75 "UPDATE subagent_sessions SET record_json = ?2 WHERE child_session_id = ?1",
76 params![child_session_id, json],
77 )?;
78 Ok(())
79 })
80}
81
82pub fn load_subagent(child_session_id: &str) -> Result<Option<mj_core::subagent::SubagentRecord>> {
83 let connection = open_reader(&database_path())?;
84 connection
85 .query_row(
86 "SELECT record_json FROM subagent_sessions WHERE child_session_id = ?1",
87 [child_session_id],
88 |row| row.get::<_, String>(0),
89 )
90 .optional()?
91 .map(|json| serde_json::from_str(&json).context("decode sub-agent record"))
92 .transpose()
93}
94
95pub fn list_subagents(parent_session_id: &str) -> Result<Vec<mj_core::subagent::SubagentRecord>> {
96 let connection = open_reader(&database_path())?;
97 let mut statement = connection.prepare(
98 "SELECT record_json FROM subagent_sessions
99 WHERE parent_session_id = ?1 ORDER BY rowid",
100 )?;
101 statement
102 .query_map([parent_session_id], |row| row.get::<_, String>(0))?
103 .map(|row| serde_json::from_str(&row?).context("decode sub-agent record"))
104 .collect()
105}
106
107pub fn load_subagent_report(child_session_id: &str) -> Result<mj_core::subagent::SubagentReport> {
110 load_subagent_report_from(&database_path(), child_session_id)
111}
112
113pub(super) fn load_subagent_report_from(
114 path: &Path,
115 child_session_id: &str,
116) -> Result<mj_core::subagent::SubagentReport> {
117 let connection = open_reader(path)?;
118 let row = connection
119 .query_row(
120 "SELECT handback_command_id, handback_message, handback_recorded_at_ms,
121 reminder_command_id, reminder_for_command_id, reminder_sent_at_ms,
122 reminder_failed_for_command_id, awaited_ordinal
123 FROM subagent_handbacks WHERE child_session_id = ?1",
124 [child_session_id],
125 |row| {
126 Ok((
127 row.get::<_, Option<String>>(0)?,
128 row.get::<_, Option<String>>(1)?,
129 row.get::<_, Option<i64>>(2)?,
130 row.get::<_, Option<String>>(3)?,
131 row.get::<_, Option<String>>(4)?,
132 row.get::<_, Option<i64>>(5)?,
133 row.get::<_, Option<String>>(6)?,
134 row.get::<_, Option<i64>>(7)?,
135 ))
136 },
137 )
138 .optional()?;
139 let Some((
140 handback_command,
141 handback_message,
142 handback_at,
143 reminder_command,
144 reminder_for,
145 reminder_at,
146 reminder_failed_for,
147 awaited_ordinal,
148 )) = row
149 else {
150 return Ok(mj_core::subagent::SubagentReport::default());
151 };
152 Ok(mj_core::subagent::SubagentReport {
153 handback: match (handback_command, handback_message, handback_at) {
154 (Some(command_id), Some(message), Some(recorded_at_ms)) => {
155 Some(mj_core::subagent::SubagentHandback {
156 command_id,
157 message,
158 recorded_at_ms,
159 })
160 }
161 _ => None,
162 },
163 reminder: match (reminder_command, reminder_for, reminder_at) {
164 (Some(command_id), Some(for_command_id), Some(sent_at_ms)) => {
165 Some(mj_core::subagent::HandbackReminder {
166 command_id,
167 for_command_id,
168 sent_at_ms,
169 })
170 }
171 _ => None,
172 },
173 reminder_failed_for,
174 awaited_ordinal: awaited_ordinal.and_then(|ordinal| u64::try_from(ordinal).ok()),
175 })
176}
177
178pub fn record_subagent_prompt(child_session_id: &str, ordinal: u64) -> Result<()> {
181 let child_session_id = child_session_id.to_owned();
182 submit_database_write("record_subagent_prompt", move |_| {
183 record_subagent_prompt_to(&database_path(), &child_session_id, ordinal)
184 })
185}
186
187pub(super) fn record_subagent_prompt_to(
188 path: &Path,
189 child_session_id: &str,
190 ordinal: u64,
191) -> Result<()> {
192 let ordinal = i64::try_from(ordinal).context("prompt ordinal exceeds the store's range")?;
193 open(path)?.execute(
194 "INSERT INTO subagent_handbacks(child_session_id, awaited_ordinal)
195 SELECT ?1, ?2 WHERE EXISTS (
196 SELECT 1 FROM subagent_sessions WHERE child_session_id = ?1
197 )
198 ON CONFLICT(child_session_id) DO UPDATE SET
199 awaited_ordinal = max(coalesce(awaited_ordinal, 0), excluded.awaited_ordinal)",
200 params![child_session_id, ordinal],
201 )?;
202 Ok(())
203}
204
205pub fn record_subagent_handback(
210 child_session_id: &str,
211 handback: &mj_core::subagent::SubagentHandback,
212) -> Result<bool> {
213 let child_session_id = child_session_id.to_owned();
214 let handback = handback.clone();
215 submit_database_write("record_subagent_handback", move |_| {
216 record_subagent_handback_to(&database_path(), &child_session_id, &handback)
217 })
218}
219
220pub(super) fn record_subagent_handback_to(
221 path: &Path,
222 child_session_id: &str,
223 handback: &mj_core::subagent::SubagentHandback,
224) -> Result<bool> {
225 let mut connection = open(path)?;
226 let tx = connection.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
227 let recorded_for: Option<String> = tx
228 .query_row(
229 "SELECT handback_command_id FROM subagent_handbacks WHERE child_session_id = ?1",
230 [child_session_id],
231 |row| row.get(0),
232 )
233 .optional()?
234 .flatten();
235 if recorded_for.as_deref() == Some(handback.command_id.as_str()) {
236 return Ok(false);
237 }
238 tx.execute(
239 "INSERT INTO subagent_handbacks(
240 child_session_id, handback_command_id, handback_message, handback_recorded_at_ms
241 ) VALUES (?1, ?2, ?3, ?4)
242 ON CONFLICT(child_session_id) DO UPDATE SET
243 handback_command_id = excluded.handback_command_id,
244 handback_message = excluded.handback_message,
245 handback_recorded_at_ms = excluded.handback_recorded_at_ms",
246 params![
247 child_session_id,
248 handback.command_id,
249 handback.message,
250 handback.recorded_at_ms
251 ],
252 )?;
253 tx.commit()?;
254 Ok(true)
255}
256
257pub fn record_handback_reminder(
259 child_session_id: &str,
260 reminder: &mj_core::subagent::HandbackReminder,
261) -> Result<()> {
262 let child_session_id = child_session_id.to_owned();
263 let reminder = reminder.clone();
264 submit_database_write("record_handback_reminder", move |_| {
265 open(&database_path())?.execute(
266 "INSERT INTO subagent_handbacks(
267 child_session_id, reminder_command_id, reminder_for_command_id, reminder_sent_at_ms
268 ) VALUES (?1, ?2, ?3, ?4)
269 ON CONFLICT(child_session_id) DO UPDATE SET
270 reminder_command_id = excluded.reminder_command_id,
271 reminder_for_command_id = excluded.reminder_for_command_id,
272 reminder_sent_at_ms = excluded.reminder_sent_at_ms",
273 params![
274 child_session_id,
275 reminder.command_id,
276 reminder.for_command_id,
277 reminder.sent_at_ms
278 ],
279 )?;
280 Ok(())
281 })
282}
283
284pub fn record_handback_reminder_failed(child_session_id: &str, for_command_id: &str) -> Result<()> {
287 let child_session_id = child_session_id.to_owned();
288 let for_command_id = for_command_id.to_owned();
289 submit_database_write("record_handback_reminder_failed", move |_| {
290 open(&database_path())?.execute(
291 "INSERT INTO subagent_handbacks(child_session_id, reminder_failed_for_command_id)
292 VALUES (?1, ?2)
293 ON CONFLICT(child_session_id) DO UPDATE SET
294 reminder_failed_for_command_id = excluded.reminder_failed_for_command_id",
295 params![child_session_id, for_command_id],
296 )?;
297 Ok(())
298 })
299}
300
301pub fn lookup_subagent_request(
302 parent_session_id: &str,
303 request_key: &str,
304) -> Result<Option<mj_core::subagent::SubagentRecord>> {
305 let connection = open_reader(&database_path())?;
306 connection
307 .query_row(
308 "SELECT record_json FROM subagent_sessions
309 WHERE parent_session_id = ?1 AND request_key = ?2",
310 params![parent_session_id, request_key],
311 |row| row.get::<_, String>(0),
312 )
313 .optional()?
314 .map(|json| serde_json::from_str(&json).context("decode sub-agent record"))
315 .transpose()
316}
317
318pub fn save_session_with_container_size(
321 session: &SessionRecord,
322 host: &str,
323 size: HostContainerSize,
324) -> Result<()> {
325 let session = session.clone();
326 let host = host.to_owned();
327 submit_database_write("save_session_with_container_size", move |_| {
328 save_session_with_container_size_to(&database_path(), &session, Some((&host, size)))
329 })
330}
331
332pub fn save_lifecycle_session(session: &SessionRecord) -> Result<()> {
336 let session = session.clone();
337 submit_database_write("save_lifecycle_session", move |_| {
338 save_lifecycle_session_to(&database_path(), &session)
339 })
340}
341
342pub fn save_checkpointed_session(session: &SessionRecord) -> Result<()> {
345 let session = session.clone();
346 submit_database_write("save_checkpointed_session", move |_| {
347 save_checkpointed_session_to(&database_path(), &session)
348 })
349}
350
351pub fn recover_interrupted_checkpointing_sessions(updated_at: &str) -> Result<usize> {
355 let updated_at = updated_at.to_owned();
356 submit_database_write("recover_interrupted_checkpointing_sessions", move |_| {
357 recover_interrupted_checkpointing_sessions_to(&database_path(), &updated_at)
358 })
359}
360
361pub fn set_session_title_override(session_id: &str, title: &str, updated_at: &str) -> Result<()> {
364 let session_id = session_id.to_owned();
365 let title = title.to_owned();
366 let updated_at = updated_at.to_owned();
367 submit_database_write("set_session_title_override", move |_| {
368 set_session_title_override_to(&database_path(), &session_id, &title, &updated_at)
369 })
370}
371
372pub fn rename_profile_references(old_id: &str, new_id: &str) -> Result<usize> {
376 rename_session_reference("last_profile", old_id, new_id)
377}
378
379pub fn rename_target_references(old_id: &str, new_id: &str) -> Result<usize> {
382 rename_session_reference("target_template_id", old_id, new_id)
383}
384
385pub(super) fn rename_session_reference(
386 column: &'static str,
387 old_id: &str,
388 new_id: &str,
389) -> Result<usize> {
390 ensure!(
391 matches!(column, "last_profile" | "target_template_id"),
392 "unsupported session reference column"
393 );
394 let old_id = old_id.to_owned();
395 let new_id = new_id.to_owned();
396 submit_database_write("rename_session_reference", move |_| {
397 rename_session_reference_at(&database_path(), column, &old_id, &new_id)
398 })
399}
400
401pub(super) fn rename_session_reference_at(
402 path: &Path,
403 column: &str,
404 old_id: &str,
405 new_id: &str,
406) -> Result<usize> {
407 let mut connection = open(path)?;
408 let tx = connection.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
409 let changed = tx.execute(
410 &format!("UPDATE sessions SET {column} = ?2 WHERE {column} = ?1"),
411 params![old_id, new_id],
412 )?;
413 tx.commit()?;
414 Ok(changed)
415}
416
417pub fn set_session_archived(session_id: &str, archived: bool) -> Result<()> {
421 let session_id = session_id.to_owned();
422 submit_database_write("set_session_archived", move |_| {
423 set_session_archived_to(&database_path(), &session_id, archived)
424 })
425}
426
427pub fn mark_session_target_missing(
432 session_id: &str,
433 detail: &str,
434 updated_at: &str,
435) -> Result<Option<SessionState>> {
436 let session_id = session_id.to_owned();
437 let detail = detail.to_owned();
438 let updated_at = updated_at.to_owned();
439 submit_database_write("mark_session_target_missing", move |_| {
440 mark_session_target_missing_to(&database_path(), &session_id, &detail, &updated_at)
441 })
442}
443
444pub(super) fn mark_session_target_missing_to(
445 path: &Path,
446 session_id: &str,
447 detail: &str,
448 updated_at: &str,
449) -> Result<Option<SessionState>> {
450 mark_session_target_missing_if_current_to(path, session_id, detail, updated_at, None)
451}
452
453pub fn mark_session_target_missing_if_current(
456 session_id: &str,
457 detail: &str,
458 updated_at: &str,
459 observed_updated_at: &str,
460) -> Result<Option<SessionState>> {
461 let session_id = session_id.to_owned();
462 let detail = detail.to_owned();
463 let updated_at = updated_at.to_owned();
464 let observed_updated_at = observed_updated_at.to_owned();
465 submit_database_write("mark_session_target_missing_if_current", move |_| {
466 mark_session_target_missing_if_current_to(
467 &database_path(),
468 &session_id,
469 &detail,
470 &updated_at,
471 Some(&observed_updated_at),
472 )
473 })
474}
475
476pub(super) fn mark_session_target_missing_if_current_to(
477 path: &Path,
478 session_id: &str,
479 detail: &str,
480 updated_at: &str,
481 observed_updated_at: Option<&str>,
482) -> Result<Option<SessionState>> {
483 let mut connection = open(path)?;
484 let tx = connection.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
485 let changed = tx.execute(
486 "UPDATE sessions
487 SET state = CASE
488 WHEN EXISTS(
489 SELECT 1 FROM session_checkpoints
490 WHERE session_checkpoints.session_id = sessions.session_id
491 ) THEN 'error'
492 ELSE 'lost'
493 END,
494 last_error = ?2,
495 updated_at = ?3
496 WHERE session_id = ?1
497 AND (?4 IS NULL OR updated_at = ?4)
498 AND state IN ('provisioning', 'running', 'disconnected', 'error')",
499 params![session_id, detail, updated_at, observed_updated_at],
500 )?;
501 ensure!(changed <= 1, "updated {changed} sessions for {session_id}");
502 let state = if changed == 1 {
503 let stored: String = tx.query_row(
504 "SELECT state FROM sessions WHERE session_id = ?1",
505 [session_id],
506 |row| row.get(0),
507 )?;
508 Some(stored_session_state(&stored))
509 } else {
510 None
511 };
512 tx.commit()?;
513 Ok(state)
514}
515
516pub(super) fn set_session_archived_to(path: &Path, session_id: &str, archived: bool) -> Result<()> {
517 let connection = open(path)?;
518 let changed = connection.execute(
519 "UPDATE sessions SET archived = ?2 WHERE session_id = ?1",
520 params![session_id, archived],
521 )?;
522 if changed != 1 {
523 bail!("unknown session {session_id}");
524 }
525 Ok(())
526}
527
528pub fn hidden_native_sessions() -> Result<BTreeSet<(mj_core::config::HarnessKind, String)>> {
531 hidden_native_sessions_from(&database_path())
532}
533
534pub(super) fn hidden_native_sessions_from(
535 path: &Path,
536) -> Result<BTreeSet<(mj_core::config::HarnessKind, String)>> {
537 let connection = open_reader(path)?;
538 let mut statement =
539 connection.prepare("SELECT harness_kind, native_session_id FROM hidden_native_sessions")?;
540 let rows = statement.query_map([], |row| {
541 Ok((row.get::<_, String>(0)?, row.get::<_, String>(1)?))
542 })?;
543 let mut hidden = BTreeSet::new();
544 for row in rows {
545 let (harness, native_session_id) = row?;
546 match harness.parse::<mj_core::config::HarnessKind>() {
549 Ok(harness) => {
550 hidden.insert((harness, native_session_id));
551 }
552 Err(_) => tracing::warn!(
553 harness = %harness,
554 "ignoring a hidden native session for a harness that is no longer supported"
555 ),
556 }
557 }
558 Ok(hidden)
559}
560
561pub fn set_native_session_hidden(
563 harness: mj_core::config::HarnessKind,
564 native_session_id: &str,
565 hidden: bool,
566) -> Result<()> {
567 let native_session_id = native_session_id.to_owned();
568 submit_database_write("set_native_session_hidden", move |_| {
569 set_native_session_hidden_to(&database_path(), harness, &native_session_id, hidden)
570 })
571}
572
573pub(super) fn set_native_session_hidden_to(
574 path: &Path,
575 harness: mj_core::config::HarnessKind,
576 native_session_id: &str,
577 hidden: bool,
578) -> Result<()> {
579 if native_session_id.trim().is_empty() {
580 bail!("native session id is empty");
581 }
582 let connection = open(path)?;
583 if hidden {
584 connection.execute(
585 "INSERT INTO hidden_native_sessions(harness_kind, native_session_id, hidden_at)
586 VALUES (?1, ?2, ?3)
587 ON CONFLICT(harness_kind, native_session_id) DO NOTHING",
588 params![harness.id(), native_session_id, Utc::now().to_rfc3339()],
589 )?;
590 } else {
591 connection.execute(
592 "DELETE FROM hidden_native_sessions
593 WHERE harness_kind = ?1 AND native_session_id = ?2",
594 params![harness.id(), native_session_id],
595 )?;
596 }
597 Ok(())
598}
599
600pub fn set_session_container_settings(
604 session_id: &str,
605 cpus: Option<&str>,
606 memory: Option<&str>,
607 mounts: &[AdditionalMount],
608 updated_at: &str,
609) -> Result<()> {
610 let session_id = session_id.to_owned();
611 let cpus = cpus.map(str::to_owned);
612 let memory = memory.map(str::to_owned);
613 let mounts = mounts.to_vec();
614 let updated_at = updated_at.to_owned();
615 submit_database_write("set_session_container_settings", move |_| {
616 set_session_container_settings_to(
617 &database_path(),
618 &session_id,
619 cpus.as_deref(),
620 memory.as_deref(),
621 &mounts,
622 &updated_at,
623 )
624 })
625}
626
627pub(super) fn set_session_container_settings_to(
628 path: &Path,
629 session_id: &str,
630 cpus: Option<&str>,
631 memory: Option<&str>,
632 mounts: &[AdditionalMount],
633 updated_at: &str,
634) -> Result<()> {
635 if updated_at.trim().is_empty() {
636 bail!("session update timestamp is empty");
637 }
638 crate::targets::validate_additional_mounts(mounts)?;
639 let mut connection = open(path)?;
640 let tx = connection.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
641 let changed = tx.execute(
642 "UPDATE sessions
643 SET container_cpus = ?2, container_memory = ?3, updated_at = ?4
644 WHERE session_id = ?1",
645 params![session_id, cpus, memory, updated_at],
646 )?;
647 if changed != 1 {
648 bail!("unknown session {session_id}");
649 }
650 replace_mounts(&tx, session_id, mounts)?;
651 tx.commit()?;
652 Ok(())
653}
654
655pub(super) fn set_session_title_override_to(
656 path: &Path,
657 session_id: &str,
658 title: &str,
659 updated_at: &str,
660) -> Result<()> {
661 if title.trim().is_empty() {
662 bail!("session title is empty");
663 }
664 if updated_at.trim().is_empty() {
665 bail!("session update timestamp is empty");
666 }
667 let connection = open(path)?;
668 let changed = connection.execute(
669 "UPDATE sessions
670 SET session_title_override = ?2, updated_at = ?3
671 WHERE session_id = ?1",
672 params![session_id, title, updated_at],
673 )?;
674 if changed != 1 {
675 bail!("unknown session {session_id}");
676 }
677 Ok(())
678}
679
680pub fn set_session_acp_title(session_id: &str, title: Option<&str>) -> Result<()> {
683 let session_id = session_id.to_owned();
684 let title = title.map(str::to_owned);
685 submit_database_write("set_session_acp_title", move |_| {
686 set_session_acp_title_to(&database_path(), &session_id, title.as_deref())
687 })
688}
689
690pub(super) fn set_session_acp_title_to(
691 path: &Path,
692 session_id: &str,
693 title: Option<&str>,
694) -> Result<()> {
695 if title.is_some_and(|title| title.trim().is_empty()) {
696 bail!("ACP session title is empty");
697 }
698 let title = title.and_then(mj_core::state::normalize_session_title);
699 let connection = open(path)?;
700 let changed = connection.execute(
701 "UPDATE sessions SET acp_session_title = ?2 WHERE session_id = ?1",
702 params![session_id, title],
703 )?;
704 if changed != 1 {
705 bail!("unknown session {session_id}");
706 }
707 Ok(())
708}
709
710pub fn mark_session_worker_connected(
713 session_id: &str,
714 native_session_id: Option<&str>,
715 updated_at: &str,
716) -> Result<()> {
717 let session_id = session_id.to_owned();
718 let native_session_id = native_session_id.map(str::to_owned);
719 let updated_at = updated_at.to_owned();
720 submit_database_write("mark_session_worker_connected", move |_| {
721 mark_session_worker_connected_to(
722 &database_path(),
723 &session_id,
724 native_session_id.as_deref(),
725 &updated_at,
726 )
727 })
728}
729
730pub fn adopt_native_session_id(session_id: &str, native_session_id: &str) -> Result<()> {
734 let session_id = session_id.to_owned();
735 let native_session_id = native_session_id.to_owned();
736 submit_database_write("adopt_native_session_id", move |_| {
737 adopt_native_session_id_to(&database_path(), &session_id, &native_session_id)
738 })
739}
740
741pub(super) fn adopt_native_session_id_to(
742 path: &Path,
743 session_id: &str,
744 native_session_id: &str,
745) -> Result<()> {
746 let connection = open(path)?;
747 let changed = connection.execute(
748 "UPDATE sessions SET native_session_id = ?2 WHERE session_id = ?1",
749 params![session_id, native_session_id],
750 )?;
751 if changed != 1 {
752 bail!("unknown session {session_id}");
753 }
754 Ok(())
755}
756
757pub(super) fn mark_session_worker_connected_to(
758 path: &Path,
759 session_id: &str,
760 native_session_id: Option<&str>,
761 updated_at: &str,
762) -> Result<()> {
763 if updated_at.trim().is_empty() {
764 bail!("worker connection timestamp is empty");
765 }
766 let connection = open(path)?;
767 let changed = connection.execute(
768 "UPDATE sessions
769 SET state = 'running',
770 native_session_id = coalesce(?2, native_session_id),
771 updated_at = ?3,
772 last_error = NULL
773 WHERE session_id = ?1",
774 params![session_id, native_session_id, updated_at],
775 )?;
776 if changed != 1 {
777 bail!("unknown session {session_id}");
778 }
779 Ok(())
780}
781
782pub(super) fn recover_interrupted_checkpointing_sessions_to(
783 path: &Path,
784 updated_at: &str,
785) -> Result<usize> {
786 if updated_at.trim().is_empty() {
787 bail!("checkpoint recovery timestamp is empty");
788 }
789 let connection = open(path)?;
790 connection
791 .execute(
792 "UPDATE sessions
793 SET state = 'running', updated_at = ?1, last_checkpoint_error = ?2
794 WHERE state = 'checkpointing'",
795 params![
796 updated_at,
797 "checkpointing was interrupted by a controller restart; the target was left running"
798 ],
799 )
800 .context("recover interrupted checkpointing sessions")
801}
802
803pub(super) fn save_session_to(path: &Path, session: &SessionRecord) -> Result<()> {
804 save_session_with_container_size_to(path, session, None)
805}
806
807pub(super) fn save_session_with_container_size_to(
808 path: &Path,
809 session: &SessionRecord,
810 container_size: Option<(&str, HostContainerSize)>,
811) -> Result<()> {
812 validate_session_record(session)?;
813
814 let mut connection = open(path)?;
815 let tx = connection.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
816 if let Some(existing_bundle) = tx
817 .query_row(
818 "SELECT bundle_id FROM session_contexts WHERE session_id = ?1",
819 [session.id.as_str()],
820 |row| row.get::<_, String>(0),
821 )
822 .optional()?
823 && existing_bundle != session.bundle_id
824 {
825 bail!(
826 "session {} was already associated with bundle {}, not {}",
827 session.id,
828 existing_bundle,
829 session.bundle_id
830 );
831 }
832 let mut session = session.clone();
833 let moving: bool = tx.query_row(
834 "SELECT EXISTS(SELECT 1 FROM session_moves WHERE session_id=?1
835 AND json_extract(operation_json, '$.phase') IN ('preparing','closing_source','resuming_destination','starting_queue'))",
836 [&session.id], |row| row.get(0),
837 )?;
838 if moving {
839 let (draft, title, acp_title, viewed, archived) = tx.query_row(
843 "SELECT draft_input, session_title_override, acp_session_title, viewed_through_event_ordinal, archived
844 FROM sessions WHERE session_id=?1", [&session.id], |row| Ok((
845 row.get::<_, String>(0)?, row.get::<_, Option<String>>(1)?, row.get::<_, Option<String>>(2)?,
846 row.get::<_, u64>(3)?, row.get::<_, bool>(4)?,
847 )),
848 )?;
849 session.draft_input = draft;
850 session.session_title_override = title;
851 session.acp_session_title = acp_title;
852 session.viewed_through_event_ordinal = viewed;
853 session.archived = archived;
854 }
855 insert_session(&tx, &session)?;
856 if let Some((host, size)) = container_size {
857 write_host_container_size(&tx, host, size)?;
858 }
859 tx.commit()?;
860 Ok(())
861}
862
863pub(super) fn save_lifecycle_session_to(path: &Path, session: &SessionRecord) -> Result<()> {
864 validate_session_record(session)?;
865
866 let mut connection = open(path)?;
867 let tx = connection.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
868 update_lifecycle_fields(&tx, session)?;
869 tx.commit()?;
870 Ok(())
871}
872
873pub(super) fn save_checkpointed_session_to(path: &Path, session: &SessionRecord) -> Result<()> {
874 validate_session_record(session)?;
875
876 let mut connection = open(path)?;
877 let tx = connection.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
878 update_lifecycle_fields(&tx, session)?;
879 tx.execute(
880 "UPDATE sessions SET native_session_id = ?2 WHERE session_id = ?1",
881 params![session.id, session.native_session_id],
882 )?;
883 replace_checkpoint(&tx, session)?;
884 tx.commit()?;
885 Ok(())
886}
887
888pub(super) fn validate_session_record(session: &SessionRecord) -> Result<()> {
889 let mut validation = State::default();
890 validation
891 .sessions
892 .insert(session.id.clone(), session.clone());
893 validation.validate()
894}
895
896pub fn delete_session(session_id: &str) -> Result<()> {
899 let session_id = session_id.to_owned();
900 submit_database_write("delete_session", move |_| {
901 delete_session_from(&database_path(), &session_id)
902 })
903}
904
905pub(super) fn delete_session_from(path: &Path, session_id: &str) -> Result<()> {
906 let connection = open(path)?;
907 connection.execute("DELETE FROM sessions WHERE session_id = ?1", [session_id])?;
908 Ok(())
909}
910
911pub fn set_session_draft_input(session_id: &str, draft: &str) -> Result<()> {
915 let session_id = session_id.to_owned();
916 let draft = draft.to_owned();
917 submit_database_write("set_session_draft_input", move |_| {
918 set_session_draft_input_at(&database_path(), &session_id, &draft)
919 })
920}
921
922pub(super) fn set_session_draft_input_at(path: &Path, session_id: &str, draft: &str) -> Result<()> {
923 let mut connection = open(path)?;
924 let tx = connection.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
925 let updated = tx.execute(
926 "UPDATE sessions SET draft_input = ?2 WHERE session_id = ?1",
927 params![session_id, draft],
928 )?;
929 ensure!(updated == 1, "unknown session {session_id}");
930 tx.commit()?;
931 Ok(())
932}
933
934pub fn clear_session_draft_input_if_matches(session_id: &str, expected: &str) -> Result<()> {
936 let session_id = session_id.to_owned();
937 let expected = expected.to_owned();
938 submit_database_write("clear_session_draft_input_if_matches", move |connection| {
939 connection.execute(
940 "UPDATE sessions SET draft_input = '' WHERE session_id = ?1 AND draft_input = ?2",
941 params![session_id, expected],
942 )?;
943 Ok(())
944 })
945}
946
947pub fn record_recovery_success(
948 session_id: &str,
949 native_session_id: &str,
950 checkpoint: &CheckpointMetadata,
951) -> Result<()> {
952 let session_id = session_id.to_owned();
953 let native_session_id = native_session_id.to_owned();
954 let checkpoint = checkpoint.clone();
955 submit_database_write("record_recovery_success", move |_| {
956 record_recovery_success_to(
957 &database_path(),
958 &session_id,
959 &native_session_id,
960 &checkpoint,
961 )
962 })
963}
964
965pub(super) fn record_recovery_success_to(
966 path: &Path,
967 session_id: &str,
968 native_session_id: &str,
969 checkpoint: &CheckpointMetadata,
970) -> Result<()> {
971 let mut connection = open(path)?;
972 let tx = connection.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
973 let changed = tx.execute(
974 "UPDATE sessions
975 SET native_session_id = ?2, last_checkpoint_error = NULL
976 WHERE session_id = ?1",
977 params![session_id, native_session_id],
978 )?;
979 if changed != 1 {
980 bail!("unknown session {session_id}");
981 }
982 tx.execute(
983 "INSERT INTO session_checkpoints(
984 session_id, archive_path, sha256, created_at, event_frontier
985 ) VALUES (?1,?2,?3,?4,?5)
986 ON CONFLICT(session_id) DO UPDATE SET
987 archive_path = excluded.archive_path,
988 sha256 = excluded.sha256,
989 created_at = excluded.created_at,
990 event_frontier = excluded.event_frontier",
991 params![
992 session_id,
993 path_to_blob(&checkpoint.archive_path),
994 checkpoint.sha256,
995 checkpoint.created_at,
996 checkpoint.event_frontier,
997 ],
998 )?;
999 tx.commit()?;
1000 Ok(())
1001}
1002
1003pub fn record_recovery_failure(session_id: &str, detail: &str) -> Result<()> {
1004 let session_id = session_id.to_owned();
1005 let detail = detail.to_owned();
1006 submit_database_write("record_recovery_failure", move |_| {
1007 record_recovery_failure_to(&database_path(), &session_id, &detail)
1008 })
1009}
1010
1011pub(super) fn record_recovery_failure_to(
1012 path: &Path,
1013 session_id: &str,
1014 detail: &str,
1015) -> Result<()> {
1016 let connection = open(path)?;
1017 let changed = connection.execute(
1018 "UPDATE sessions SET last_checkpoint_error = ?2 WHERE session_id = ?1",
1019 params![session_id, detail],
1020 )?;
1021 if changed != 1 {
1022 bail!("unknown session {session_id}");
1023 }
1024 Ok(())
1025}
1026
1027pub fn rebind_session_bundle(session_id: &str, bundle_id: &str) -> Result<()> {
1034 let session_id = session_id.to_owned();
1035 let bundle_id = bundle_id.to_owned();
1036 submit_database_write("rebind_session_bundle", move |_| {
1037 rebind_session_bundle_to(&database_path(), &session_id, &bundle_id)
1038 })
1039}
1040
1041pub(super) fn rebind_session_bundle_to(
1042 path: &Path,
1043 session_id: &str,
1044 bundle_id: &str,
1045) -> Result<()> {
1046 let mut connection = open(path)?;
1047 let tx = connection.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
1048 let changed = tx.execute(
1049 "UPDATE session_contexts SET bundle_id = ?2 WHERE session_id = ?1",
1050 params![session_id, bundle_id],
1051 )?;
1052 if changed == 0 {
1053 tx.execute(
1054 "INSERT INTO session_contexts(session_id, bundle_id, created_at) VALUES (?1, ?2, ?3)",
1055 params![session_id, bundle_id, Utc::now().to_rfc3339()],
1056 )?;
1057 }
1058 tx.commit()?;
1059 Ok(())
1060}