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 load_subagent_parent_worktree(
99 child_session_id: &str,
100) -> Result<Option<mj_core::state::ManagedWorktree>> {
101 let connection = open_reader(&database_path())?;
102 connection
103 .query_row(
104 "SELECT s.managed_worktree FROM subagent_sessions r
105 JOIN sessions s ON s.session_id = r.parent_session_id
106 WHERE r.child_session_id = ?1",
107 [child_session_id],
108 |row| row.get::<_, Option<String>>(0),
109 )
110 .optional()?
111 .flatten()
112 .map(|json| serde_json::from_str(&json).context("decode the parent's managed checkout"))
113 .transpose()
114}
115
116pub fn list_subagents(parent_session_id: &str) -> Result<Vec<mj_core::subagent::SubagentRecord>> {
117 let connection = open_reader(&database_path())?;
118 let mut statement = connection.prepare(
119 "SELECT record_json FROM subagent_sessions
120 WHERE parent_session_id = ?1 ORDER BY rowid",
121 )?;
122 statement
123 .query_map([parent_session_id], |row| row.get::<_, String>(0))?
124 .map(|row| serde_json::from_str(&row?).context("decode sub-agent record"))
125 .collect()
126}
127
128pub fn record_stopped_subagents(
132 parent_session_id: &str,
133 stopped: &[mj_core::subagent::StoppedSubagent],
134) -> Result<()> {
135 let parent_session_id = parent_session_id.to_owned();
136 let stopped = stopped.to_vec();
137 submit_database_write("record_stopped_subagents", move |_| {
138 record_stopped_subagents_to(&database_path(), &parent_session_id, &stopped)
139 })
140}
141
142pub(super) fn record_stopped_subagents_to(
143 path: &Path,
144 parent_session_id: &str,
145 stopped: &[mj_core::subagent::StoppedSubagent],
146) -> Result<()> {
147 let mut connection = open(path)?;
148 let tx = connection.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
149 for child in stopped {
150 tx.execute(
151 "INSERT INTO stopped_subagents(parent_session_id, child_session_id, record_json)
152 VALUES (?1, ?2, ?3)
153 ON CONFLICT(parent_session_id, child_session_id) DO UPDATE SET
154 record_json = excluded.record_json",
155 params![
156 parent_session_id,
157 child.child_session_id,
158 serde_json::to_string(child)?
159 ],
160 )?;
161 }
162 tx.commit()?;
163 Ok(())
164}
165
166pub fn load_stopped_subagents(
169 parent_session_id: &str,
170) -> Result<Vec<mj_core::subagent::StoppedSubagent>> {
171 load_stopped_subagents_from(&database_path(), parent_session_id)
172}
173
174pub(super) fn load_stopped_subagents_from(
175 path: &Path,
176 parent_session_id: &str,
177) -> Result<Vec<mj_core::subagent::StoppedSubagent>> {
178 let connection = open_reader(path)?;
179 let mut statement = connection.prepare(
180 "SELECT record_json FROM stopped_subagents
181 WHERE parent_session_id = ?1 ORDER BY rowid",
182 )?;
183 statement
184 .query_map([parent_session_id], |row| row.get::<_, String>(0))?
185 .map(|row| serde_json::from_str(&row?).context("decode stopped sub-agent record"))
186 .collect()
187}
188
189pub fn clear_stopped_subagents(
193 parent_session_id: &str,
194 child_session_ids: &[String],
195) -> Result<()> {
196 let parent_session_id = parent_session_id.to_owned();
197 let child_session_ids = child_session_ids.to_vec();
198 submit_database_write("clear_stopped_subagents", move |_| {
199 clear_stopped_subagents_from(&database_path(), &parent_session_id, &child_session_ids)
200 })
201}
202
203pub(super) fn clear_stopped_subagents_from(
204 path: &Path,
205 parent_session_id: &str,
206 child_session_ids: &[String],
207) -> Result<()> {
208 let mut connection = open(path)?;
209 let tx = connection.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
210 for child_session_id in child_session_ids {
211 tx.execute(
212 "DELETE FROM stopped_subagents
213 WHERE parent_session_id = ?1 AND child_session_id = ?2",
214 params![parent_session_id, child_session_id],
215 )?;
216 }
217 tx.commit()?;
218 Ok(())
219}
220
221pub fn load_session_state(session_id: &str) -> Result<Option<SessionState>> {
226 let connection = open_reader(&database_path())?;
227 let stored = connection
228 .query_row(
229 "SELECT state FROM sessions WHERE session_id = ?1",
230 [session_id],
231 |row| row.get::<_, String>(0),
232 )
233 .optional()?;
234 Ok(stored.as_deref().map(stored_session_state))
235}
236
237pub fn load_subagent_report(child_session_id: &str) -> Result<mj_core::subagent::SubagentReport> {
238 load_subagent_report_from(&database_path(), child_session_id)
239}
240
241pub(super) fn load_subagent_report_from(
242 path: &Path,
243 child_session_id: &str,
244) -> Result<mj_core::subagent::SubagentReport> {
245 let connection = open_reader(path)?;
246 let row = connection
247 .query_row(
248 "SELECT handback_command_id, handback_message, handback_recorded_at_ms,
249 reminder_command_id, reminder_for_command_id, reminder_sent_at_ms,
250 reminder_failed_for_command_id, awaited_ordinal, report_dir
251 FROM subagent_handbacks WHERE child_session_id = ?1",
252 [child_session_id],
253 |row| {
254 Ok((
255 row.get::<_, Option<String>>(0)?,
256 row.get::<_, Option<String>>(1)?,
257 row.get::<_, Option<i64>>(2)?,
258 row.get::<_, Option<String>>(3)?,
259 row.get::<_, Option<String>>(4)?,
260 row.get::<_, Option<i64>>(5)?,
261 row.get::<_, Option<String>>(6)?,
262 row.get::<_, Option<i64>>(7)?,
263 row.get::<_, Option<String>>(8)?,
264 ))
265 },
266 )
267 .optional()?;
268 let Some((
269 handback_command,
270 handback_message,
271 handback_at,
272 reminder_command,
273 reminder_for,
274 reminder_at,
275 reminder_failed_for,
276 awaited_ordinal,
277 report_dir,
278 )) = row
279 else {
280 return Ok(mj_core::subagent::SubagentReport::default());
281 };
282 Ok(mj_core::subagent::SubagentReport {
283 handback: match (handback_command, handback_message, handback_at) {
284 (Some(command_id), Some(message), Some(recorded_at_ms)) => {
285 Some(mj_core::subagent::SubagentHandback {
286 command_id,
287 message,
288 recorded_at_ms,
289 })
290 }
291 _ => None,
292 },
293 reminder: match (reminder_command, reminder_for, reminder_at) {
294 (Some(command_id), Some(for_command_id), Some(sent_at_ms)) => {
295 Some(mj_core::subagent::HandbackReminder {
296 command_id,
297 for_command_id,
298 sent_at_ms,
299 })
300 }
301 _ => None,
302 },
303 reminder_failed_for,
304 awaited_ordinal: awaited_ordinal.and_then(|ordinal| u64::try_from(ordinal).ok()),
305 report_dir,
306 })
307}
308
309pub fn record_subagent_report_dir(child_session_id: &str, report_dir: &str) -> Result<()> {
312 let child_session_id = child_session_id.to_owned();
313 let report_dir = report_dir.to_owned();
314 submit_database_write("record_subagent_report_dir", move |_| {
315 record_subagent_report_dir_to(&database_path(), &child_session_id, &report_dir)
316 })
317}
318
319pub(super) fn record_subagent_report_dir_to(
320 path: &Path,
321 child_session_id: &str,
322 report_dir: &str,
323) -> Result<()> {
324 open(path)?.execute(
325 "INSERT INTO subagent_handbacks(child_session_id, report_dir)
326 SELECT ?1, ?2 WHERE EXISTS (
327 SELECT 1 FROM subagent_sessions WHERE child_session_id = ?1
328 )
329 ON CONFLICT(child_session_id) DO UPDATE SET report_dir = excluded.report_dir",
330 params![child_session_id, report_dir],
331 )?;
332 Ok(())
333}
334
335pub fn record_subagent_prompt(child_session_id: &str, ordinal: u64) -> Result<()> {
338 let child_session_id = child_session_id.to_owned();
339 submit_database_write("record_subagent_prompt", move |_| {
340 record_subagent_prompt_to(&database_path(), &child_session_id, ordinal)
341 })
342}
343
344pub(super) fn record_subagent_prompt_to(
345 path: &Path,
346 child_session_id: &str,
347 ordinal: u64,
348) -> Result<()> {
349 let ordinal = i64::try_from(ordinal).context("prompt ordinal exceeds the store's range")?;
350 open(path)?.execute(
351 "INSERT INTO subagent_handbacks(child_session_id, awaited_ordinal)
352 SELECT ?1, ?2 WHERE EXISTS (
353 SELECT 1 FROM subagent_sessions WHERE child_session_id = ?1
354 )
355 ON CONFLICT(child_session_id) DO UPDATE SET
356 awaited_ordinal = max(coalesce(awaited_ordinal, 0), excluded.awaited_ordinal)",
357 params![child_session_id, ordinal],
358 )?;
359 Ok(())
360}
361
362pub fn record_subagent_handback(
367 child_session_id: &str,
368 handback: &mj_core::subagent::SubagentHandback,
369) -> Result<bool> {
370 let child_session_id = child_session_id.to_owned();
371 let handback = handback.clone();
372 submit_database_write("record_subagent_handback", move |_| {
373 record_subagent_handback_to(&database_path(), &child_session_id, &handback)
374 })
375}
376
377pub(super) fn record_subagent_handback_to(
378 path: &Path,
379 child_session_id: &str,
380 handback: &mj_core::subagent::SubagentHandback,
381) -> Result<bool> {
382 let mut connection = open(path)?;
383 let tx = connection.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
384 let recorded_for: Option<String> = tx
385 .query_row(
386 "SELECT handback_command_id FROM subagent_handbacks WHERE child_session_id = ?1",
387 [child_session_id],
388 |row| row.get(0),
389 )
390 .optional()?
391 .flatten();
392 if recorded_for.as_deref() == Some(handback.command_id.as_str()) {
393 return Ok(false);
394 }
395 tx.execute(
396 "INSERT INTO subagent_handbacks(
397 child_session_id, handback_command_id, handback_message, handback_recorded_at_ms
398 ) VALUES (?1, ?2, ?3, ?4)
399 ON CONFLICT(child_session_id) DO UPDATE SET
400 handback_command_id = excluded.handback_command_id,
401 handback_message = excluded.handback_message,
402 handback_recorded_at_ms = excluded.handback_recorded_at_ms",
403 params![
404 child_session_id,
405 handback.command_id,
406 handback.message,
407 handback.recorded_at_ms
408 ],
409 )?;
410 tx.commit()?;
411 Ok(true)
412}
413
414pub fn record_handback_reminder(
416 child_session_id: &str,
417 reminder: &mj_core::subagent::HandbackReminder,
418) -> Result<()> {
419 let child_session_id = child_session_id.to_owned();
420 let reminder = reminder.clone();
421 submit_database_write("record_handback_reminder", move |_| {
422 open(&database_path())?.execute(
423 "INSERT INTO subagent_handbacks(
424 child_session_id, reminder_command_id, reminder_for_command_id, reminder_sent_at_ms
425 ) VALUES (?1, ?2, ?3, ?4)
426 ON CONFLICT(child_session_id) DO UPDATE SET
427 reminder_command_id = excluded.reminder_command_id,
428 reminder_for_command_id = excluded.reminder_for_command_id,
429 reminder_sent_at_ms = excluded.reminder_sent_at_ms",
430 params![
431 child_session_id,
432 reminder.command_id,
433 reminder.for_command_id,
434 reminder.sent_at_ms
435 ],
436 )?;
437 Ok(())
438 })
439}
440
441pub fn record_handback_reminder_failed(child_session_id: &str, for_command_id: &str) -> Result<()> {
444 let child_session_id = child_session_id.to_owned();
445 let for_command_id = for_command_id.to_owned();
446 submit_database_write("record_handback_reminder_failed", move |_| {
447 open(&database_path())?.execute(
448 "INSERT INTO subagent_handbacks(child_session_id, reminder_failed_for_command_id)
449 VALUES (?1, ?2)
450 ON CONFLICT(child_session_id) DO UPDATE SET
451 reminder_failed_for_command_id = excluded.reminder_failed_for_command_id",
452 params![child_session_id, for_command_id],
453 )?;
454 Ok(())
455 })
456}
457
458pub fn lookup_subagent_request(
459 parent_session_id: &str,
460 request_key: &str,
461) -> Result<Option<mj_core::subagent::SubagentRecord>> {
462 let connection = open_reader(&database_path())?;
463 connection
464 .query_row(
465 "SELECT record_json FROM subagent_sessions
466 WHERE parent_session_id = ?1 AND request_key = ?2",
467 params![parent_session_id, request_key],
468 |row| row.get::<_, String>(0),
469 )
470 .optional()?
471 .map(|json| serde_json::from_str(&json).context("decode sub-agent record"))
472 .transpose()
473}
474
475pub fn save_session_with_container_size(
478 session: &SessionRecord,
479 host: &str,
480 size: HostContainerSize,
481) -> Result<()> {
482 let session = session.clone();
483 let host = host.to_owned();
484 submit_database_write("save_session_with_container_size", move |_| {
485 save_session_with_container_size_to(&database_path(), &session, Some((&host, size)))
486 })
487}
488
489pub fn save_lifecycle_session(session: &SessionRecord) -> Result<()> {
493 let session = session.clone();
494 submit_database_write("save_lifecycle_session", move |_| {
495 save_lifecycle_session_to(&database_path(), &session)
496 })
497}
498
499pub fn save_checkpointed_session(session: &SessionRecord) -> Result<()> {
502 let session = session.clone();
503 submit_database_write("save_checkpointed_session", move |_| {
504 save_checkpointed_session_to(&database_path(), &session)
505 })
506}
507
508pub fn recover_interrupted_checkpointing_sessions(updated_at: &str) -> Result<usize> {
512 let updated_at = updated_at.to_owned();
513 submit_database_write("recover_interrupted_checkpointing_sessions", move |_| {
514 recover_interrupted_checkpointing_sessions_to(&database_path(), &updated_at)
515 })
516}
517
518pub fn set_session_title_override(session_id: &str, title: &str, updated_at: &str) -> Result<()> {
521 let session_id = session_id.to_owned();
522 let title = title.to_owned();
523 let updated_at = updated_at.to_owned();
524 submit_database_write("set_session_title_override", move |_| {
525 set_session_title_override_to(&database_path(), &session_id, &title, &updated_at)
526 })
527}
528
529pub fn rename_profile_references(old_id: &str, new_id: &str) -> Result<usize> {
533 rename_session_reference("last_profile", old_id, new_id)
534}
535
536pub fn rename_target_references(old_id: &str, new_id: &str) -> Result<usize> {
539 rename_session_reference("target_template_id", old_id, new_id)
540}
541
542pub(super) fn rename_session_reference(
543 column: &'static str,
544 old_id: &str,
545 new_id: &str,
546) -> Result<usize> {
547 ensure!(
548 matches!(column, "last_profile" | "target_template_id"),
549 "unsupported session reference column"
550 );
551 let old_id = old_id.to_owned();
552 let new_id = new_id.to_owned();
553 submit_database_write("rename_session_reference", move |_| {
554 rename_session_reference_at(&database_path(), column, &old_id, &new_id)
555 })
556}
557
558pub(super) fn rename_session_reference_at(
559 path: &Path,
560 column: &str,
561 old_id: &str,
562 new_id: &str,
563) -> Result<usize> {
564 let mut connection = open(path)?;
565 let tx = connection.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
566 let changed = tx.execute(
567 &format!("UPDATE sessions SET {column} = ?2 WHERE {column} = ?1"),
568 params![old_id, new_id],
569 )?;
570 tx.commit()?;
571 Ok(changed)
572}
573
574pub fn set_session_archived(session_id: &str, archived: bool) -> Result<()> {
578 let session_id = session_id.to_owned();
579 submit_database_write("set_session_archived", move |_| {
580 set_session_archived_to(&database_path(), &session_id, archived)
581 })
582}
583
584pub fn mark_session_target_missing(
589 session_id: &str,
590 detail: &str,
591 updated_at: &str,
592) -> Result<Option<SessionState>> {
593 let session_id = session_id.to_owned();
594 let detail = detail.to_owned();
595 let updated_at = updated_at.to_owned();
596 submit_database_write("mark_session_target_missing", move |_| {
597 mark_session_target_missing_to(&database_path(), &session_id, &detail, &updated_at)
598 })
599}
600
601pub(super) fn mark_session_target_missing_to(
602 path: &Path,
603 session_id: &str,
604 detail: &str,
605 updated_at: &str,
606) -> Result<Option<SessionState>> {
607 mark_session_target_missing_if_current_to(path, session_id, detail, updated_at, None)
608}
609
610pub fn mark_session_target_missing_if_current(
613 session_id: &str,
614 detail: &str,
615 updated_at: &str,
616 observed_updated_at: &str,
617) -> Result<Option<SessionState>> {
618 let session_id = session_id.to_owned();
619 let detail = detail.to_owned();
620 let updated_at = updated_at.to_owned();
621 let observed_updated_at = observed_updated_at.to_owned();
622 submit_database_write("mark_session_target_missing_if_current", move |_| {
623 mark_session_target_missing_if_current_to(
624 &database_path(),
625 &session_id,
626 &detail,
627 &updated_at,
628 Some(&observed_updated_at),
629 )
630 })
631}
632
633pub(super) fn mark_session_target_missing_if_current_to(
634 path: &Path,
635 session_id: &str,
636 detail: &str,
637 updated_at: &str,
638 observed_updated_at: Option<&str>,
639) -> Result<Option<SessionState>> {
640 let mut connection = open(path)?;
641 let tx = connection.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
642 let changed = tx.execute(
643 "UPDATE sessions
644 SET state = CASE
645 WHEN EXISTS(
646 SELECT 1 FROM session_checkpoints
647 WHERE session_checkpoints.session_id = sessions.session_id
648 ) THEN 'error'
649 ELSE 'lost'
650 END,
651 last_error = ?2,
652 updated_at = ?3
653 WHERE session_id = ?1
654 AND (?4 IS NULL OR updated_at = ?4)
655 AND state IN ('provisioning', 'running', 'disconnected', 'error')",
656 params![session_id, detail, updated_at, observed_updated_at],
657 )?;
658 ensure!(changed <= 1, "updated {changed} sessions for {session_id}");
659 let state = if changed == 1 {
660 let stored: String = tx.query_row(
661 "SELECT state FROM sessions WHERE session_id = ?1",
662 [session_id],
663 |row| row.get(0),
664 )?;
665 Some(stored_session_state(&stored))
666 } else {
667 None
668 };
669 tx.commit()?;
670 Ok(state)
671}
672
673pub(super) fn set_session_archived_to(path: &Path, session_id: &str, archived: bool) -> Result<()> {
674 let connection = open(path)?;
675 let changed = connection.execute(
676 "UPDATE sessions SET archived = ?2 WHERE session_id = ?1",
677 params![session_id, archived],
678 )?;
679 if changed != 1 {
680 bail!("unknown session {session_id}");
681 }
682 Ok(())
683}
684
685pub fn hidden_native_sessions() -> Result<BTreeSet<(mj_core::config::HarnessKind, String)>> {
688 hidden_native_sessions_from(&database_path())
689}
690
691pub(super) fn hidden_native_sessions_from(
692 path: &Path,
693) -> Result<BTreeSet<(mj_core::config::HarnessKind, String)>> {
694 let connection = open_reader(path)?;
695 let mut statement =
696 connection.prepare("SELECT harness_kind, native_session_id FROM hidden_native_sessions")?;
697 let rows = statement.query_map([], |row| {
698 Ok((row.get::<_, String>(0)?, row.get::<_, String>(1)?))
699 })?;
700 let mut hidden = BTreeSet::new();
701 for row in rows {
702 let (harness, native_session_id) = row?;
703 match harness.parse::<mj_core::config::HarnessKind>() {
706 Ok(harness) => {
707 hidden.insert((harness, native_session_id));
708 }
709 Err(_) => tracing::warn!(
710 harness = %harness,
711 "ignoring a hidden native session for a harness that is no longer supported"
712 ),
713 }
714 }
715 Ok(hidden)
716}
717
718pub fn set_native_session_hidden(
720 harness: mj_core::config::HarnessKind,
721 native_session_id: &str,
722 hidden: bool,
723) -> Result<()> {
724 let native_session_id = native_session_id.to_owned();
725 submit_database_write("set_native_session_hidden", move |_| {
726 set_native_session_hidden_to(&database_path(), harness, &native_session_id, hidden)
727 })
728}
729
730pub(super) fn set_native_session_hidden_to(
731 path: &Path,
732 harness: mj_core::config::HarnessKind,
733 native_session_id: &str,
734 hidden: bool,
735) -> Result<()> {
736 if native_session_id.trim().is_empty() {
737 bail!("native session id is empty");
738 }
739 let connection = open(path)?;
740 if hidden {
741 connection.execute(
742 "INSERT INTO hidden_native_sessions(harness_kind, native_session_id, hidden_at)
743 VALUES (?1, ?2, ?3)
744 ON CONFLICT(harness_kind, native_session_id) DO NOTHING",
745 params![harness.id(), native_session_id, Utc::now().to_rfc3339()],
746 )?;
747 } else {
748 connection.execute(
749 "DELETE FROM hidden_native_sessions
750 WHERE harness_kind = ?1 AND native_session_id = ?2",
751 params![harness.id(), native_session_id],
752 )?;
753 }
754 Ok(())
755}
756
757pub fn set_session_container_settings(
761 session_id: &str,
762 cpus: Option<&str>,
763 memory: Option<&str>,
764 mounts: &[AdditionalMount],
765 updated_at: &str,
766) -> Result<()> {
767 let session_id = session_id.to_owned();
768 let cpus = cpus.map(str::to_owned);
769 let memory = memory.map(str::to_owned);
770 let mounts = mounts.to_vec();
771 let updated_at = updated_at.to_owned();
772 submit_database_write("set_session_container_settings", move |_| {
773 set_session_container_settings_to(
774 &database_path(),
775 &session_id,
776 cpus.as_deref(),
777 memory.as_deref(),
778 &mounts,
779 &updated_at,
780 )
781 })
782}
783
784pub(super) fn set_session_container_settings_to(
785 path: &Path,
786 session_id: &str,
787 cpus: Option<&str>,
788 memory: Option<&str>,
789 mounts: &[AdditionalMount],
790 updated_at: &str,
791) -> Result<()> {
792 if updated_at.trim().is_empty() {
793 bail!("session update timestamp is empty");
794 }
795 crate::targets::validate_additional_mounts(mounts)?;
796 let mut connection = open(path)?;
797 let tx = connection.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
798 let changed = tx.execute(
799 "UPDATE sessions
800 SET container_cpus = ?2, container_memory = ?3, updated_at = ?4
801 WHERE session_id = ?1",
802 params![session_id, cpus, memory, updated_at],
803 )?;
804 if changed != 1 {
805 bail!("unknown session {session_id}");
806 }
807 replace_mounts(&tx, session_id, mounts)?;
808 tx.commit()?;
809 Ok(())
810}
811
812pub(super) fn set_session_title_override_to(
813 path: &Path,
814 session_id: &str,
815 title: &str,
816 updated_at: &str,
817) -> Result<()> {
818 if title.trim().is_empty() {
819 bail!("session title is empty");
820 }
821 if updated_at.trim().is_empty() {
822 bail!("session update timestamp is empty");
823 }
824 let connection = open(path)?;
825 let changed = connection.execute(
826 "UPDATE sessions
827 SET session_title_override = ?2, updated_at = ?3
828 WHERE session_id = ?1",
829 params![session_id, title, updated_at],
830 )?;
831 if changed != 1 {
832 bail!("unknown session {session_id}");
833 }
834 Ok(())
835}
836
837pub fn set_session_acp_title(session_id: &str, title: Option<&str>) -> Result<()> {
840 let session_id = session_id.to_owned();
841 let title = title.map(str::to_owned);
842 submit_database_write("set_session_acp_title", move |_| {
843 set_session_acp_title_to(&database_path(), &session_id, title.as_deref())
844 })
845}
846
847pub(super) fn set_session_acp_title_to(
848 path: &Path,
849 session_id: &str,
850 title: Option<&str>,
851) -> Result<()> {
852 if title.is_some_and(|title| title.trim().is_empty()) {
853 bail!("ACP session title is empty");
854 }
855 let title = title.and_then(mj_core::state::normalize_session_title);
856 let connection = open(path)?;
857 let changed = connection.execute(
858 "UPDATE sessions SET acp_session_title = ?2 WHERE session_id = ?1",
859 params![session_id, title],
860 )?;
861 if changed != 1 {
862 bail!("unknown session {session_id}");
863 }
864 Ok(())
865}
866
867pub fn mark_session_worker_connected(
870 session_id: &str,
871 native_session_id: Option<&str>,
872 updated_at: &str,
873) -> Result<()> {
874 let session_id = session_id.to_owned();
875 let native_session_id = native_session_id.map(str::to_owned);
876 let updated_at = updated_at.to_owned();
877 submit_database_write("mark_session_worker_connected", move |_| {
878 mark_session_worker_connected_to(
879 &database_path(),
880 &session_id,
881 native_session_id.as_deref(),
882 &updated_at,
883 )
884 })
885}
886
887pub fn adopt_native_session_id(session_id: &str, native_session_id: &str) -> Result<()> {
891 let session_id = session_id.to_owned();
892 let native_session_id = native_session_id.to_owned();
893 submit_database_write("adopt_native_session_id", move |_| {
894 adopt_native_session_id_to(&database_path(), &session_id, &native_session_id)
895 })
896}
897
898pub(super) fn adopt_native_session_id_to(
899 path: &Path,
900 session_id: &str,
901 native_session_id: &str,
902) -> Result<()> {
903 let connection = open(path)?;
904 let changed = connection.execute(
905 "UPDATE sessions SET native_session_id = ?2 WHERE session_id = ?1",
906 params![session_id, native_session_id],
907 )?;
908 if changed != 1 {
909 bail!("unknown session {session_id}");
910 }
911 Ok(())
912}
913
914pub(super) fn mark_session_worker_connected_to(
915 path: &Path,
916 session_id: &str,
917 native_session_id: Option<&str>,
918 updated_at: &str,
919) -> Result<()> {
920 if updated_at.trim().is_empty() {
921 bail!("worker connection timestamp is empty");
922 }
923 let connection = open(path)?;
924 let changed = connection.execute(
925 "UPDATE sessions
926 SET state = 'running',
927 native_session_id = coalesce(?2, native_session_id),
928 updated_at = ?3,
929 last_error = NULL
930 WHERE session_id = ?1",
931 params![session_id, native_session_id, updated_at],
932 )?;
933 if changed != 1 {
934 bail!("unknown session {session_id}");
935 }
936 Ok(())
937}
938
939pub(super) fn recover_interrupted_checkpointing_sessions_to(
940 path: &Path,
941 updated_at: &str,
942) -> Result<usize> {
943 if updated_at.trim().is_empty() {
944 bail!("checkpoint recovery timestamp is empty");
945 }
946 let connection = open(path)?;
947 connection
948 .execute(
949 "UPDATE sessions
950 SET state = 'running', updated_at = ?1, last_checkpoint_error = ?2
951 WHERE state = 'checkpointing'",
952 params![
953 updated_at,
954 "checkpointing was interrupted by a controller restart; the target was left running"
955 ],
956 )
957 .context("recover interrupted checkpointing sessions")
958}
959
960pub(super) fn save_session_to(path: &Path, session: &SessionRecord) -> Result<()> {
961 save_session_with_container_size_to(path, session, None)
962}
963
964pub(super) fn save_session_with_container_size_to(
965 path: &Path,
966 session: &SessionRecord,
967 container_size: Option<(&str, HostContainerSize)>,
968) -> Result<()> {
969 validate_session_record(session)?;
970
971 let mut connection = open(path)?;
972 let tx = connection.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
973 if let Some(existing_bundle) = tx
974 .query_row(
975 "SELECT bundle_id FROM session_contexts WHERE session_id = ?1",
976 [session.id.as_str()],
977 |row| row.get::<_, String>(0),
978 )
979 .optional()?
980 && existing_bundle != session.bundle_id
981 {
982 bail!(
983 "session {} was already associated with bundle {}, not {}",
984 session.id,
985 existing_bundle,
986 session.bundle_id
987 );
988 }
989 let mut session = session.clone();
990 let moving: bool = tx.query_row(
991 "SELECT EXISTS(SELECT 1 FROM session_moves WHERE session_id=?1
992 AND json_extract(operation_json, '$.phase') IN ('preparing','closing_source','resuming_destination','starting_queue'))",
993 [&session.id], |row| row.get(0),
994 )?;
995 if moving {
996 let (draft, title, acp_title, viewed, archived) = tx.query_row(
1000 "SELECT draft_input, session_title_override, acp_session_title, viewed_through_event_ordinal, archived
1001 FROM sessions WHERE session_id=?1", [&session.id], |row| Ok((
1002 row.get::<_, String>(0)?, row.get::<_, Option<String>>(1)?, row.get::<_, Option<String>>(2)?,
1003 row.get::<_, u64>(3)?, row.get::<_, bool>(4)?,
1004 )),
1005 )?;
1006 session.draft_input = draft;
1007 session.session_title_override = title;
1008 session.acp_session_title = acp_title;
1009 session.viewed_through_event_ordinal = viewed;
1010 session.archived = archived;
1011 }
1012 insert_session(&tx, &session)?;
1013 if let Some((host, size)) = container_size {
1014 write_host_container_size(&tx, host, size)?;
1015 }
1016 tx.commit()?;
1017 Ok(())
1018}
1019
1020pub(super) fn save_lifecycle_session_to(path: &Path, session: &SessionRecord) -> Result<()> {
1021 validate_session_record(session)?;
1022
1023 let mut connection = open(path)?;
1024 let tx = connection.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
1025 update_lifecycle_fields(&tx, session)?;
1026 tx.commit()?;
1027 Ok(())
1028}
1029
1030pub(super) fn save_checkpointed_session_to(path: &Path, session: &SessionRecord) -> Result<()> {
1031 validate_session_record(session)?;
1032
1033 let mut connection = open(path)?;
1034 let tx = connection.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
1035 update_lifecycle_fields(&tx, session)?;
1036 tx.execute(
1037 "UPDATE sessions SET native_session_id = ?2 WHERE session_id = ?1",
1038 params![session.id, session.native_session_id],
1039 )?;
1040 replace_checkpoint(&tx, session)?;
1041 tx.commit()?;
1042 Ok(())
1043}
1044
1045pub(super) fn validate_session_record(session: &SessionRecord) -> Result<()> {
1046 let mut validation = State::default();
1047 validation
1048 .sessions
1049 .insert(session.id.clone(), session.clone());
1050 validation.validate()
1051}
1052
1053pub fn delete_session(session_id: &str) -> Result<()> {
1056 let session_id = session_id.to_owned();
1057 submit_database_write("delete_session", move |_| {
1058 delete_session_from(&database_path(), &session_id)
1059 })
1060}
1061
1062pub(super) fn delete_session_from(path: &Path, session_id: &str) -> Result<()> {
1063 let connection = open(path)?;
1064 connection.execute("DELETE FROM sessions WHERE session_id = ?1", [session_id])?;
1065 Ok(())
1066}
1067
1068pub fn set_session_draft_input(session_id: &str, draft: &str) -> Result<()> {
1072 let session_id = session_id.to_owned();
1073 let draft = draft.to_owned();
1074 submit_database_write("set_session_draft_input", move |_| {
1075 set_session_draft_input_at(&database_path(), &session_id, &draft)
1076 })
1077}
1078
1079pub(super) fn set_session_draft_input_at(path: &Path, session_id: &str, draft: &str) -> Result<()> {
1080 let mut connection = open(path)?;
1081 let tx = connection.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
1082 let updated = tx.execute(
1083 "UPDATE sessions SET draft_input = ?2 WHERE session_id = ?1",
1084 params![session_id, draft],
1085 )?;
1086 ensure!(updated == 1, "unknown session {session_id}");
1087 tx.commit()?;
1088 Ok(())
1089}
1090
1091pub fn clear_session_draft_input_if_matches(session_id: &str, expected: &str) -> Result<()> {
1093 let session_id = session_id.to_owned();
1094 let expected = expected.to_owned();
1095 submit_database_write("clear_session_draft_input_if_matches", move |connection| {
1096 connection.execute(
1097 "UPDATE sessions SET draft_input = '' WHERE session_id = ?1 AND draft_input = ?2",
1098 params![session_id, expected],
1099 )?;
1100 Ok(())
1101 })
1102}
1103
1104pub fn record_recovery_success(
1105 session_id: &str,
1106 native_session_id: &str,
1107 checkpoint: &CheckpointMetadata,
1108) -> Result<()> {
1109 let session_id = session_id.to_owned();
1110 let native_session_id = native_session_id.to_owned();
1111 let checkpoint = checkpoint.clone();
1112 submit_database_write("record_recovery_success", move |_| {
1113 record_recovery_success_to(
1114 &database_path(),
1115 &session_id,
1116 &native_session_id,
1117 &checkpoint,
1118 )
1119 })
1120}
1121
1122pub(super) fn record_recovery_success_to(
1123 path: &Path,
1124 session_id: &str,
1125 native_session_id: &str,
1126 checkpoint: &CheckpointMetadata,
1127) -> Result<()> {
1128 let mut connection = open(path)?;
1129 let tx = connection.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
1130 let changed = tx.execute(
1131 "UPDATE sessions
1132 SET native_session_id = ?2, last_checkpoint_error = NULL
1133 WHERE session_id = ?1",
1134 params![session_id, native_session_id],
1135 )?;
1136 if changed != 1 {
1137 bail!("unknown session {session_id}");
1138 }
1139 tx.execute(
1140 "INSERT INTO session_checkpoints(
1141 session_id, archive_path, sha256, created_at, event_frontier
1142 ) VALUES (?1,?2,?3,?4,?5)
1143 ON CONFLICT(session_id) DO UPDATE SET
1144 archive_path = excluded.archive_path,
1145 sha256 = excluded.sha256,
1146 created_at = excluded.created_at,
1147 event_frontier = excluded.event_frontier",
1148 params![
1149 session_id,
1150 path_to_blob(&checkpoint.archive_path),
1151 checkpoint.sha256,
1152 checkpoint.created_at,
1153 checkpoint.event_frontier,
1154 ],
1155 )?;
1156 tx.commit()?;
1157 Ok(())
1158}
1159
1160pub fn record_recovery_failure(session_id: &str, detail: &str) -> Result<()> {
1161 let session_id = session_id.to_owned();
1162 let detail = detail.to_owned();
1163 submit_database_write("record_recovery_failure", move |_| {
1164 record_recovery_failure_to(&database_path(), &session_id, &detail)
1165 })
1166}
1167
1168pub(super) fn record_recovery_failure_to(
1169 path: &Path,
1170 session_id: &str,
1171 detail: &str,
1172) -> Result<()> {
1173 let connection = open(path)?;
1174 let changed = connection.execute(
1175 "UPDATE sessions SET last_checkpoint_error = ?2 WHERE session_id = ?1",
1176 params![session_id, detail],
1177 )?;
1178 if changed != 1 {
1179 bail!("unknown session {session_id}");
1180 }
1181 Ok(())
1182}
1183
1184pub fn rebind_session_bundle(session_id: &str, bundle_id: &str) -> Result<()> {
1191 let session_id = session_id.to_owned();
1192 let bundle_id = bundle_id.to_owned();
1193 submit_database_write("rebind_session_bundle", move |_| {
1194 rebind_session_bundle_to(&database_path(), &session_id, &bundle_id)
1195 })
1196}
1197
1198pub(super) fn rebind_session_bundle_to(
1199 path: &Path,
1200 session_id: &str,
1201 bundle_id: &str,
1202) -> Result<()> {
1203 let mut connection = open(path)?;
1204 let tx = connection.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
1205 let changed = tx.execute(
1206 "UPDATE session_contexts SET bundle_id = ?2 WHERE session_id = ?1",
1207 params![session_id, bundle_id],
1208 )?;
1209 if changed == 0 {
1210 tx.execute(
1211 "INSERT INTO session_contexts(session_id, bundle_id, created_at) VALUES (?1, ?2, ?3)",
1212 params![session_id, bundle_id, Utc::now().to_rfc3339()],
1213 )?;
1214 }
1215 tx.commit()?;
1216 Ok(())
1217}