1use super::*;
2
3pub fn load_state() -> Result<State> {
4 load_state_from(&database_path())
5}
6
7pub fn load_state_from(path: &Path) -> Result<State> {
8 let mut reader = open_reader(path)?;
9 let connection = reader.transaction()?;
12 let mut state = State::default();
13 let mut statement = connection.prepare(
14 "SELECT s.session_id, s.title, s.harness_kind, s.last_profile, c.bundle_id,
15 s.target_template_id, s.state, s.native_session_id, s.acp_session_title,
16 s.session_title_override, c.created_at, s.updated_at,
17 s.viewed_through_event_ordinal, s.last_error, s.resource_allocation,
18 s.last_checkpoint_error, s.project_directory, s.managed_worktree,
19 s.draft_input, s.container_cpus, s.container_memory, s.archived
20 , c.workspace_id, s.create_managed_worktree, s.mjolnir_subagents,
21 s.container_workspace, s.build_cache_json
22 FROM sessions s JOIN session_contexts c USING(session_id)
23 ORDER BY s.session_id",
24 )?;
25 let rows = statement.query_map([], |row| {
26 let harness_text: String = row.get(2)?;
30 let Ok(harness_kind) = harness_text.parse() else {
31 let session_id: String = row.get(0)?;
32 tracing::warn!(
33 session_id,
34 harness = %harness_text,
35 "session harness is no longer supported; the session is not listed"
36 );
37 return Ok(None);
38 };
39 Ok(Some(SessionRecord {
40 harness_kind,
41 create_managed_worktree: row.get(23)?,
42 mjolnir_subagents: row.get(24)?,
43 container_workspace: row.get::<_, Option<String>>(25)?.map(PathBuf::from),
44 build_cache: row
45 .get::<_, Option<String>>(26)?
46 .as_deref()
47 .and_then(|text| match serde_json::from_str(text) {
48 Ok(build_cache) => Some(build_cache),
49 Err(error) => {
50 tracing::warn!(%error, "session build cache record is unreadable");
51 None
52 }
53 }),
54 workspace_id: row.get(22)?,
55 archived: row.get(21)?,
56 container_cpus: row.get(19)?,
57 container_memory: row.get(20)?,
58 id: row.get(0)?,
59 title: row.get(1)?,
60 last_profile: row.get(3)?,
61 bundle_id: row.get(4)?,
62 project_directory: row.get_ref(16)?.blob_or_null()?.map(blob_to_path),
63 managed_worktree: row
64 .get::<_, Option<String>>(17)?
65 .map(|json| serde_json::from_str::<ManagedWorktree>(&json))
66 .transpose()
67 .map_err(|error| {
68 rusqlite::Error::FromSqlConversionFailure(
69 17,
70 rusqlite::types::Type::Text,
71 Box::new(error),
72 )
73 })?,
74 target_template_id: row.get(5)?,
75 resource_allocation: row
76 .get::<_, Option<String>>(14)?
77 .map(|json| serde_json::from_str::<SessionResourceAllocation>(&json))
78 .transpose()
79 .map_err(|error| {
80 rusqlite::Error::FromSqlConversionFailure(
81 14,
82 rusqlite::types::Type::Text,
83 Box::new(error),
84 )
85 })?,
86 additional_mounts: Vec::new(),
87 state: stored_session_state(&row.get::<_, String>(6)?),
88 target: None,
89 native_session_id: row.get(7)?,
90 acp_session_title: row
91 .get::<_, Option<String>>(8)?
92 .as_deref()
93 .and_then(mj_core::state::normalize_session_title),
94 session_title_override: row.get(9)?,
95 created_at: row.get(10)?,
96 updated_at: row.get(11)?,
97 viewed_through_event_ordinal: row.get::<_, u64>(12)?,
98 draft_input: row.get(18)?,
99 last_error: row.get(13)?,
100 last_checkpoint_error: row.get(15)?,
101 checkpoint: None,
102 }))
103 })?;
104 for row in rows {
105 if let Some(session) = row? {
106 state.sessions.insert(session.id.clone(), session);
107 }
108 }
109 #[cfg(test)]
110 super::tests::after_state_sessions_read();
111 let mut statement = connection.prepare(
112 "SELECT child_session_id, record_json FROM subagent_sessions ORDER BY child_session_id",
113 )?;
114 let rows = statement.query_map([], |row| {
115 let child_id = row.get::<_, String>(0)?;
116 let json = row.get::<_, String>(1)?;
117 let record = serde_json::from_str::<SubagentRecord>(&json).map_err(|error| {
118 rusqlite::Error::FromSqlConversionFailure(1, Type::Text, Box::new(error))
119 })?;
120 Ok((child_id, record))
121 })?;
122 for row in rows {
123 let (child_id, record) = row?;
124 let missing = if !state.sessions.contains_key(&child_id) {
133 Some("child")
134 } else if !state.sessions.contains_key(&record.parent_session_id) {
135 Some("parent")
136 } else {
137 None
138 };
139 if let Some(missing) = missing {
140 tracing::warn!(
141 child_session_id = child_id,
142 parent_session_id = record.parent_session_id,
143 missing,
144 "dropping a sub-agent relation whose session is not in this state"
145 );
146 continue;
147 }
148 state.subagents.insert(child_id, record);
149 }
150 load_targets(&connection, &mut state)?;
151 load_mounts(&connection, &mut state)?;
152 load_checkpoints(&connection, &mut state)?;
153 let mut statement =
154 connection.prepare("SELECT host, source FROM mount_history ORDER BY host, ordinal")?;
155 let rows = statement.query_map([], |row| {
156 Ok((
157 row.get::<_, String>(0)?,
158 blob_to_path(row.get_ref(1)?.as_blob()?),
159 ))
160 })?;
161 for row in rows {
162 let (host, source) = row?;
163 state.mount_history.entry(host).or_default().push(source);
164 }
165 let mut statement = connection
166 .prepare("SELECT host, cpus, memory_bytes FROM host_container_sizes ORDER BY host")?;
167 let rows = statement.query_map([], |row| {
168 Ok((
169 row.get::<_, String>(0)?,
170 HostContainerSize {
171 cpus: row.get::<_, i64>(1)? as u64,
172 memory_bytes: row.get::<_, i64>(2)? as u64,
173 },
174 ))
175 })?;
176 for row in rows {
177 let (host, size) = row?;
178 state.container_sizes.insert(host, size);
179 }
180 state.validate()?;
181 Ok(state)
182}
183
184pub fn save_state(state: &State) -> Result<()> {
185 let state = state.clone();
186 submit_database_write("save_state", move |_| {
187 save_state_to(&database_path(), &state)
188 })
189}
190
191pub fn save_state_to(path: &Path, state: &State) -> Result<()> {
192 state.validate()?;
193 let mut connection = open(path)?;
194 let tx = connection.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
195 let existing_contexts = existing_contexts(&tx)?;
196 let existing_sessions = {
197 let mut statement = tx.prepare("SELECT session_id FROM sessions")?;
198 statement
199 .query_map([], |row| row.get::<_, String>(0))?
200 .collect::<rusqlite::Result<Vec<_>>>()?
201 };
202 tx.execute(
203 "DELETE FROM subagent_sessions
204 WHERE child_session_id NOT IN (SELECT session_id FROM sessions)
205 OR parent_session_id NOT IN (SELECT session_id FROM sessions)",
206 [],
207 )?;
208 let existing_subagents = {
209 let mut statement =
210 tx.prepare("SELECT child_session_id, parent_session_id FROM subagent_sessions")?;
211 statement
212 .query_map([], |row| {
213 Ok((row.get::<_, String>(0)?, row.get::<_, String>(1)?))
214 })?
215 .collect::<rusqlite::Result<Vec<_>>>()?
216 };
217 for (child_id, parent_id) in existing_subagents {
218 if !state.subagents.contains_key(&child_id)
219 || !state.sessions.contains_key(&child_id)
220 || !state.sessions.contains_key(&parent_id)
221 {
222 tx.execute(
223 "DELETE FROM subagent_sessions WHERE child_session_id = ?1",
224 [child_id],
225 )?;
226 }
227 }
228 for session_id in existing_sessions {
229 if !state.sessions.contains_key(&session_id) {
230 tx.execute("DELETE FROM sessions WHERE session_id = ?1", [session_id])?;
231 }
232 }
233 tx.execute("DELETE FROM mount_history", [])?;
234 tx.execute("DELETE FROM host_container_sizes", [])?;
235 for session in state.sessions.values() {
236 if let Some((existing_bundle, existing_workspace)) = existing_contexts.get(&session.id) {
237 ensure!(
238 existing_bundle == &session.bundle_id,
239 "session {} was already associated with bundle {}, not {}",
240 session.id,
241 existing_bundle,
242 session.bundle_id
243 );
244 ensure!(
245 existing_workspace == &session.workspace_id,
246 "session {} was already associated with workspace {}, not {}",
247 session.id,
248 existing_workspace,
249 session.workspace_id
250 );
251 }
252 insert_session(&tx, session)?;
253 }
254 for subagent in state.subagents.values() {
255 let record_json = serde_json::to_string(subagent)?;
256 tx.execute(
257 "INSERT INTO subagent_sessions(
258 child_session_id, parent_session_id, request_key, record_json
259 ) VALUES (?1, ?2, ?3, ?4)
260 ON CONFLICT(child_session_id) DO UPDATE SET
261 parent_session_id = excluded.parent_session_id,
262 request_key = excluded.request_key,
263 record_json = excluded.record_json",
264 params![
265 subagent.child_session_id,
266 subagent.parent_session_id,
267 subagent.request_key,
268 record_json
269 ],
270 )?;
271 }
272 for (host, sources) in &state.mount_history {
273 for (ordinal, source) in sources.iter().enumerate() {
274 tx.execute(
275 "INSERT INTO mount_history(host, source, ordinal) VALUES (?1, ?2, ?3)",
276 params![host, path_to_blob(source), ordinal as i64],
277 )?;
278 }
279 }
280 for (host, size) in &state.container_sizes {
281 write_host_container_size(&tx, host, *size)?;
282 }
283 tx.commit()?;
284 Ok(())
285}
286
287pub(super) fn existing_contexts(
288 tx: &Transaction<'_>,
289) -> Result<BTreeMap<String, (String, String)>> {
290 let mut statement =
291 tx.prepare("SELECT session_id, bundle_id, workspace_id FROM session_contexts")?;
292 let rows = statement.query_map([], |row| Ok((row.get(0)?, (row.get(1)?, row.get(2)?))))?;
293 rows.collect::<rusqlite::Result<_>>().map_err(Into::into)
294}
295
296pub(super) fn session_exists(tx: &Transaction<'_>, session_id: &str) -> Result<bool> {
297 Ok(tx
298 .query_row(
299 "SELECT 1 FROM sessions WHERE session_id = ?1",
300 [session_id],
301 |_| Ok(()),
302 )
303 .optional()?
304 .is_some())
305}
306
307pub(super) fn write_materialized_session(
308 tx: &Transaction<'_>,
309 materialized: &MaterializedSession,
310) -> Result<()> {
311 let (execution, running_started_at_ms) = materialized_execution_columns(materialized.execution);
312 tx.execute(
313 "INSERT INTO materialized_sessions(
314 session_id, applied_event_ordinal, applied_event_digest, execution_state,
315 running_started_at_ms, session_title, configuration_json, last_activity_at_ms,
316 pending_elicitations_json, active_turn_json, last_turn_outcome_json
317 ) VALUES (?1,?2,?3,?4,?5,?6,?7,?8,?9,?10,?11)
318 ON CONFLICT(session_id) DO UPDATE SET
319 applied_event_ordinal = excluded.applied_event_ordinal,
320 applied_event_digest = excluded.applied_event_digest,
321 execution_state = excluded.execution_state,
322 running_started_at_ms = excluded.running_started_at_ms,
323 session_title = excluded.session_title,
324 configuration_json = excluded.configuration_json,
325 last_activity_at_ms = excluded.last_activity_at_ms,
326 pending_elicitations_json = excluded.pending_elicitations_json,
327 active_turn_json = excluded.active_turn_json,
328 last_turn_outcome_json = excluded.last_turn_outcome_json",
329 params![
330 materialized.session_id,
331 materialized.applied_event_ordinal,
332 materialized.applied_event_digest,
333 execution,
334 running_started_at_ms,
335 materialized.session_title,
336 serde_json::to_string(&materialized.configuration)?,
337 materialized.last_activity_at_ms,
338 serde_json::to_string(&materialized.pending_elicitations)?,
339 materialized
340 .active_turn
341 .as_ref()
342 .map(serde_json::to_string)
343 .transpose()?,
344 materialized
345 .last_turn_outcome
346 .as_ref()
347 .map(serde_json::to_string)
348 .transpose()?,
349 ],
350 )?;
351 tx.execute(
352 "DELETE FROM materialized_transcript_items WHERE session_id = ?1",
353 [materialized.session_id.as_str()],
354 )?;
355 for item in &materialized.transcript {
356 upsert_transcript_item(tx, &materialized.session_id, item)?;
357 }
358 replace_materialized_queue(tx, &materialized.session_id, &materialized.queued_prompts)?;
359 Ok(())
360}
361
362pub(super) fn upsert_transcript_item(
363 tx: &Transaction<'_>,
364 session_id: &str,
365 item: &TranscriptItem,
366) -> Result<()> {
367 let existing = tx
368 .query_row(
369 "SELECT position, latest_content_event_ordinal, created_at_ms, last_changed_at_ms
370 FROM materialized_transcript_items
371 WHERE session_id = ?1 AND stable_id = ?2",
372 params![session_id, item.stable_id],
373 |row| {
374 Ok((
375 row.get::<_, u64>(0)?,
376 row.get::<_, Option<u64>>(1)?,
377 row.get::<_, i64>(2)?,
378 row.get::<_, i64>(3)?,
379 ))
380 },
381 )
382 .optional()?;
383 if let Some((position, latest_content_event_ordinal, created_at_ms, last_changed_at_ms)) =
384 existing
385 {
386 if position != item.position || created_at_ms != item.created_at_ms {
387 return Err(ProjectionIntegrityError(format!(
388 "transcript item {:?} changed immutable identity fields",
389 item.stable_id
390 ))
391 .into());
392 }
393 if item.last_changed_at_ms < last_changed_at_ms {
394 return Err(ProjectionIntegrityError(format!(
395 "transcript item {:?} moved its changed timestamp backwards",
396 item.stable_id
397 ))
398 .into());
399 }
400 if latest_content_event_ordinal.is_some_and(|existing| {
401 item.latest_content_event_ordinal
402 .is_none_or(|next| next < existing)
403 }) {
404 return Err(ProjectionIntegrityError(format!(
405 "transcript item {:?} moved its latest content ordinal backwards",
406 item.stable_id
407 ))
408 .into());
409 }
410 tx.execute(
411 "UPDATE materialized_transcript_items
412 SET latest_content_event_ordinal = ?3, last_changed_at_ms = ?4, body_json = ?5
413 WHERE session_id = ?1 AND stable_id = ?2",
414 params![
415 session_id,
416 item.stable_id,
417 item.latest_content_event_ordinal,
418 item.last_changed_at_ms,
419 serde_json::to_string(&item.body)?,
420 ],
421 )?;
422 } else {
423 tx.execute(
424 "INSERT INTO materialized_transcript_items(
425 session_id, stable_id, position, latest_content_event_ordinal,
426 created_at_ms, last_changed_at_ms, body_json
427 ) VALUES (?1,?2,?3,?4,?5,?6,?7)",
428 params![
429 session_id,
430 item.stable_id,
431 item.position,
432 item.latest_content_event_ordinal,
433 item.created_at_ms,
434 item.last_changed_at_ms,
435 serde_json::to_string(&item.body)?,
436 ],
437 )?;
438 }
439 Ok(())
440}
441
442pub(super) fn replace_materialized_queue(
443 tx: &Transaction<'_>,
444 session_id: &str,
445 queued_prompts: &[MaterializedQueuedPrompt],
446) -> Result<()> {
447 let mut command_ids = BTreeSet::new();
448 for prompt in queued_prompts {
449 if prompt.command_id.trim().is_empty() {
450 bail!("materialized prompt queue has an empty command id");
451 }
452 if !command_ids.insert(prompt.command_id.as_str()) {
453 bail!(
454 "materialized prompt queue contains duplicate command {:?}",
455 prompt.command_id
456 );
457 }
458 }
459 tx.execute(
460 "DELETE FROM materialized_queued_prompts WHERE session_id = ?1",
461 [session_id],
462 )?;
463 for (ordinal, prompt) in queued_prompts.iter().enumerate() {
464 tx.execute(
465 "INSERT INTO materialized_queued_prompts(
466 session_id, ordinal, command_id, kind_json, content_json, queued_at_ms,
467 accepted_ordinal
468 ) VALUES (?1,?2,?3,?4,?5,?6,?7)",
469 params![
470 session_id,
471 ordinal as i64,
472 prompt.command_id,
473 serde_json::to_string(&prompt.kind)?,
474 serde_json::to_string(&prompt.content)?,
475 prompt.queued_at_ms,
476 prompt.accepted_ordinal,
477 ],
478 )?;
479 }
480 Ok(())
481}
482
483pub(super) fn materialized_execution_columns(
484 execution: MaterializedExecutionState,
485) -> (&'static str, Option<i64>) {
486 match execution {
487 MaterializedExecutionState::Idle => ("idle", None),
488 MaterializedExecutionState::Running { started_at_ms } => ("running", Some(started_at_ms)),
489 MaterializedExecutionState::Closing => ("closing", None),
490 MaterializedExecutionState::Closed => ("closed", None),
491 }
492}
493
494pub(super) fn parse_materialized_execution(
495 execution: &str,
496 running_started_at_ms: Option<i64>,
497) -> Result<MaterializedExecutionState> {
498 match (execution, running_started_at_ms) {
499 ("idle", None) => Ok(MaterializedExecutionState::Idle),
500 ("running", Some(started_at_ms)) => {
501 Ok(MaterializedExecutionState::Running { started_at_ms })
502 }
503 ("closing", None) => Ok(MaterializedExecutionState::Closing),
504 ("closed", None) => Ok(MaterializedExecutionState::Closed),
505 _ => bail!("invalid materialized execution state {execution:?}"),
506 }
507}
508
509pub(super) fn insert_session(tx: &Transaction<'_>, session: &SessionRecord) -> Result<()> {
513 tx.execute(
514 "INSERT INTO session_contexts(session_id, bundle_id, created_at, workspace_id)
515 VALUES (?1, ?2, ?3, ?4)
516 ON CONFLICT(session_id) DO NOTHING",
517 params![
518 session.id,
519 session.bundle_id,
520 session.created_at,
521 session.workspace_id
522 ],
523 )?;
524 let (stored_bundle, stored_workspace): (String, String) = tx.query_row(
525 "SELECT bundle_id, workspace_id FROM session_contexts WHERE session_id = ?1",
526 [session.id.as_str()],
527 |row| Ok((row.get(0)?, row.get(1)?)),
528 )?;
529 ensure!(
530 stored_bundle == session.bundle_id,
531 "session {} belongs to bundle {}, not {}",
532 session.id,
533 stored_bundle,
534 session.bundle_id
535 );
536 ensure!(
537 stored_workspace == session.workspace_id,
538 "session {} belongs to workspace {}, not {}",
539 session.id,
540 stored_workspace,
541 session.workspace_id
542 );
543 tx.execute(
544 "INSERT INTO sessions(
545 session_id, title, harness_kind, last_profile, target_template_id, state,
546 native_session_id, acp_session_title, session_title_override, updated_at,
547 viewed_through_event_ordinal, last_error, resource_allocation,
548 last_checkpoint_error, project_directory, managed_worktree,
549 container_cpus, container_memory, archived, draft_input, create_managed_worktree,
550 mjolnir_subagents, container_workspace, build_cache_json
551 ) VALUES (?1,?2,?3,?4,?5,?6,?7,?8,?9,?10,?11,?12,?13,?14,?15,?16,?17,?18,?19,?20,?21,?22,?23,?24)
552 ON CONFLICT(session_id) DO UPDATE SET
553 title = excluded.title,
554 harness_kind = excluded.harness_kind,
555 last_profile = excluded.last_profile,
556 target_template_id = excluded.target_template_id,
557 state = excluded.state,
558 native_session_id = excluded.native_session_id,
559 acp_session_title = excluded.acp_session_title,
560 session_title_override = excluded.session_title_override,
561 updated_at = excluded.updated_at,
562 viewed_through_event_ordinal = max(
563 sessions.viewed_through_event_ordinal,
564 excluded.viewed_through_event_ordinal
565 ),
566 last_error = excluded.last_error,
567 resource_allocation = excluded.resource_allocation,
568 last_checkpoint_error = excluded.last_checkpoint_error,
569 project_directory = excluded.project_directory,
570 managed_worktree = excluded.managed_worktree,
571 container_cpus = excluded.container_cpus,
572 container_memory = excluded.container_memory,
573 archived = excluded.archived,
574 create_managed_worktree = excluded.create_managed_worktree,
575 mjolnir_subagents = excluded.mjolnir_subagents,
576 container_workspace = excluded.container_workspace,
577 build_cache_json = excluded.build_cache_json",
578 params![
579 session.id,
580 session.title,
581 session.harness_kind.id(),
582 session.last_profile,
583 session.target_template_id,
584 session.state.as_str(),
585 session.native_session_id,
586 session.acp_session_title,
587 session.session_title_override,
588 session.updated_at,
589 session.viewed_through_event_ordinal,
590 session.last_error,
591 session
592 .resource_allocation
593 .as_ref()
594 .map(serde_json::to_string)
595 .transpose()?,
596 session.last_checkpoint_error,
597 session
598 .project_directory
599 .as_ref()
600 .map(|path| path_to_blob(path)),
601 session
602 .managed_worktree
603 .as_ref()
604 .map(serde_json::to_string)
605 .transpose()?,
606 session.container_cpus,
607 session.container_memory,
608 session.archived,
609 session.draft_input,
610 session.create_managed_worktree,
611 session.mjolnir_subagents,
612 session
613 .container_workspace
614 .as_ref()
615 .map(|path| path.to_string_lossy().into_owned()),
616 session
617 .build_cache
618 .as_ref()
619 .map(serde_json::to_string)
620 .transpose()?,
621 ],
622 )?;
623 tx.execute(
624 "INSERT INTO materialized_sessions(session_id) VALUES (?1)
625 ON CONFLICT(session_id) DO NOTHING",
626 [session.id.as_str()],
627 )?;
628 replace_targets(tx, session)?;
629 replace_mounts(tx, &session.id, &session.additional_mounts)?;
630 replace_checkpoint(tx, session)?;
631 Ok(())
632}
633
634pub(super) fn update_lifecycle_fields(tx: &Transaction<'_>, session: &SessionRecord) -> Result<()> {
638 let SessionRecord {
645 id,
646 title,
647 harness_kind,
648 last_profile,
649 target_template_id,
650 state,
651 updated_at,
652 viewed_through_event_ordinal,
653 last_error,
654 resource_allocation,
655 last_checkpoint_error,
656 project_directory,
657 managed_worktree,
658 build_cache,
659 workspace_id: _,
660 bundle_id: _,
661 create_managed_worktree: _,
662 mjolnir_subagents: _,
663 additional_mounts: _,
664 container_cpus: _,
665 container_memory: _,
666 container_workspace: _,
667 archived: _,
668 target: _,
670 native_session_id: _,
671 acp_session_title: _,
672 session_title_override: _,
673 created_at: _,
674 draft_input: _,
675 checkpoint: _,
677 } = session;
678 let changed = tx.execute(
679 "UPDATE sessions
682 SET title = ?2,
683 harness_kind = ?3,
684 last_profile = ?4,
685 target_template_id = ?5,
686 state = ?6,
687 updated_at = ?7,
688 viewed_through_event_ordinal = max(viewed_through_event_ordinal, ?8),
689 last_error = ?9,
690 resource_allocation = ?10,
691 last_checkpoint_error = ?11,
692 project_directory = ?12,
693 managed_worktree = ?13,
694 build_cache_json = ?14
695 WHERE session_id = ?1",
696 params![
697 id,
698 title,
699 harness_kind.id(),
700 last_profile,
701 target_template_id,
702 state.as_str(),
703 updated_at,
704 viewed_through_event_ordinal,
705 last_error,
706 resource_allocation
707 .as_ref()
708 .map(serde_json::to_string)
709 .transpose()?,
710 last_checkpoint_error,
711 project_directory.as_ref().map(|path| path_to_blob(path)),
712 managed_worktree
713 .as_ref()
714 .map(serde_json::to_string)
715 .transpose()?,
716 build_cache
719 .as_ref()
720 .map(serde_json::to_string)
721 .transpose()?,
722 ],
723 )?;
724 if changed != 1 {
725 bail!("unknown session {id}");
726 }
727 replace_targets(tx, session)
728}
729
730pub(super) fn replace_targets(tx: &Transaction<'_>, session: &SessionRecord) -> Result<()> {
731 tx.execute(
732 "DELETE FROM session_targets WHERE session_id = ?1",
733 [session.id.as_str()],
734 )?;
735 if let Some(target) = &session.target {
736 insert_target(tx, &session.id, target)?;
737 }
738 Ok(())
739}
740
741pub(super) fn replace_checkpoint(tx: &Transaction<'_>, session: &SessionRecord) -> Result<()> {
742 tx.execute(
743 "DELETE FROM session_checkpoints WHERE session_id = ?1",
744 [session.id.as_str()],
745 )?;
746 if let Some(checkpoint) = &session.checkpoint {
747 tx.execute(
748 "INSERT INTO session_checkpoints(session_id, archive_path, sha256, created_at, event_frontier)
749 VALUES (?1, ?2, ?3, ?4, ?5)",
750 params![
751 session.id,
752 path_to_blob(&checkpoint.archive_path),
753 checkpoint.sha256,
754 checkpoint.created_at,
755 checkpoint.event_frontier,
756 ],
757 )?;
758 }
759 Ok(())
760}
761
762pub(super) fn insert_target(
763 tx: &Transaction<'_>,
764 session_id: &str,
765 target: &TargetLocator,
766) -> Result<()> {
767 let (kind, host, resource, address, workspace, worker_id, workspace_storage, borrowed_from) =
768 match target {
769 TargetLocator::LocalBare { worker_root } => (
770 "local-bare",
771 None,
772 None,
773 None,
774 Some(path_to_blob(worker_root)),
775 None,
776 None,
777 None,
778 ),
779 TargetLocator::LocalPodman {
780 container_id,
781 workspace_storage,
782 borrowed_from,
783 } => (
784 "local-podman",
785 None,
786 Some(container_id.as_str()),
787 None,
788 None,
789 None,
790 Some(serde_json::to_string(workspace_storage)?),
791 borrowed_from.as_deref(),
792 ),
793 TargetLocator::LocalDocker {
794 container_id,
795 borrowed_from,
796 } => (
797 "local-docker",
798 None,
799 Some(container_id.as_str()),
800 None,
801 None,
802 None,
803 None,
804 borrowed_from.as_deref(),
805 ),
806 TargetLocator::SshDocker {
807 host,
808 container_id,
809 borrowed_from,
810 } => (
811 "ssh-docker",
812 Some(host.as_str()),
813 Some(container_id.as_str()),
814 None,
815 None,
816 None,
817 None,
818 borrowed_from.as_deref(),
819 ),
820 TargetLocator::AppleContainer {
821 container_id,
822 borrowed_from,
823 } => (
824 "apple-container",
825 None,
826 Some(container_id.as_str()),
827 None,
828 None,
829 None,
830 None,
831 borrowed_from.as_deref(),
832 ),
833 TargetLocator::AwsEc2 {
834 instance_id,
835 address,
836 } => (
837 "aws-ec2",
838 None,
839 Some(instance_id.as_str()),
840 address.as_deref(),
841 None,
842 None,
843 None,
844 None,
845 ),
846 TargetLocator::SshBare {
847 host,
848 workspace,
849 worker_id,
850 } => (
851 "ssh-bare",
852 Some(host.as_str()),
853 None,
854 None,
855 Some(path_to_blob(workspace)),
856 worker_id.as_deref(),
857 None,
858 None,
859 ),
860 TargetLocator::SshPodman {
861 host,
862 container_id,
863 workspace_storage,
864 borrowed_from,
865 } => (
866 "ssh-podman",
867 Some(host.as_str()),
868 Some(container_id.as_str()),
869 None,
870 None,
871 None,
872 Some(serde_json::to_string(workspace_storage)?),
873 borrowed_from.as_deref(),
874 ),
875 };
876 tx.execute(
877 "INSERT INTO session_targets(session_id, kind, host, resource_id, address, workspace, worker_id, workspace_storage, borrowed_from)
878 VALUES (?1,?2,?3,?4,?5,?6,?7,?8,?9)",
879 params![
880 session_id,
881 kind,
882 host,
883 resource,
884 address,
885 workspace,
886 worker_id,
887 workspace_storage,
888 borrowed_from
889 ],
890 )?;
891 Ok(())
892}
893
894pub(super) fn load_targets(connection: &Connection, state: &mut State) -> Result<()> {
895 let mut statement = connection.prepare(
896 "SELECT session_id, kind, host, resource_id, address, workspace, worker_id, workspace_storage, borrowed_from
897 FROM session_targets",
898 )?;
899 let rows = statement.query_map([], |row| {
900 let session_id: String = row.get(0)?;
901 let kind: String = row.get(1)?;
902 let host: Option<String> = row.get(2)?;
903 let resource: Option<String> = row.get(3)?;
904 let address: Option<String> = row.get(4)?;
905 let workspace = row.get_ref(5)?.blob_or_null()?.map(blob_to_path);
906 let worker_id: Option<String> = row.get(6)?;
907 let workspace_storage = row
908 .get::<_, Option<String>>(7)?
909 .map(|serialized| {
910 serde_json::from_str(&serialized).map_err(|error| {
911 rusqlite::Error::FromSqlConversionFailure(7, Type::Text, Box::new(error))
912 })
913 })
914 .transpose()?
915 .unwrap_or_default();
916 let borrowed_from: Option<String> = row.get(8)?;
917 let target = match kind.as_str() {
918 "local-bare" => TargetLocator::LocalBare {
919 worker_root: workspace.unwrap(),
920 },
921 "local-podman" => TargetLocator::LocalPodman {
922 borrowed_from,
923 container_id: resource.unwrap(),
924 workspace_storage,
925 },
926 "local-docker" => TargetLocator::LocalDocker {
927 borrowed_from,
928 container_id: resource.unwrap(),
929 },
930 "apple-container" => TargetLocator::AppleContainer {
931 borrowed_from,
932 container_id: resource.unwrap(),
933 },
934 "aws-ec2" => TargetLocator::AwsEc2 {
935 instance_id: resource.unwrap(),
936 address,
937 },
938 "ssh-bare" => TargetLocator::SshBare {
939 host: host.unwrap(),
940 workspace: workspace.unwrap(),
941 worker_id,
942 },
943 "ssh-docker" => TargetLocator::SshDocker {
944 borrowed_from,
945 host: host.unwrap(),
946 container_id: resource.unwrap(),
947 },
948 "ssh-podman" => TargetLocator::SshPodman {
949 borrowed_from,
950 host: host.unwrap(),
951 container_id: resource.unwrap(),
952 workspace_storage,
953 },
954 _ => unreachable!("target kind constrained by schema"),
955 };
956 Ok((session_id, target))
957 })?;
958 for row in rows {
959 let (session_id, target) = row?;
960 if let Some(session) = state.sessions.get_mut(&session_id) {
963 session.target = Some(target);
964 }
965 }
966 Ok(())
967}
968
969pub(super) fn replace_mounts(
976 tx: &rusqlite::Transaction<'_>,
977 session_id: &str,
978 mounts: &[AdditionalMount],
979) -> Result<()> {
980 tx.execute(
981 "DELETE FROM session_mounts WHERE session_id = ?1",
982 [session_id],
983 )?;
984 tx.execute(
985 "DELETE FROM session_mount_access WHERE session_id = ?1",
986 [session_id],
987 )?;
988 for (ordinal, mount) in mounts.iter().enumerate() {
989 tx.execute(
990 "INSERT INTO session_mounts(session_id, ordinal, source, destination, read_only)
991 VALUES (?1, ?2, ?3, ?4, ?5)",
992 params![
993 session_id,
994 ordinal as i64,
995 path_to_blob(&mount.source),
996 path_to_blob(&mount.destination),
997 mount.access == MountAccess::Ro
998 ],
999 )?;
1000 if mount.access == MountAccess::Rw {
1001 tx.execute(
1002 "INSERT INTO session_mount_access(session_id, source, destination, access)
1003 VALUES (?1, ?2, ?3, 'rw')",
1004 params![
1005 session_id,
1006 path_to_blob(&mount.source),
1007 path_to_blob(&mount.destination)
1008 ],
1009 )?;
1010 }
1011 }
1012 Ok(())
1013}
1014
1015pub(super) fn load_mounts(connection: &Connection, state: &mut State) -> Result<()> {
1016 let mut statement = connection.prepare(
1017 "SELECT m.session_id, m.source, m.destination, m.read_only, a.access IS NOT NULL
1018 FROM session_mounts m
1019 LEFT JOIN session_mount_access a
1020 ON a.session_id = m.session_id
1021 AND a.source = m.source
1022 AND a.destination = m.destination
1023 ORDER BY m.session_id, m.ordinal",
1024 )?;
1025 let rows = statement.query_map([], |row| {
1026 let access = match (row.get::<_, bool>(3)?, row.get::<_, bool>(4)?) {
1029 (true, _) => MountAccess::Ro,
1030 (false, true) => MountAccess::Rw,
1031 (false, false) => MountAccess::Cow,
1032 };
1033 Ok((
1034 row.get::<_, String>(0)?,
1035 AdditionalMount {
1036 source: blob_to_path(row.get_ref(1)?.as_blob()?),
1037 destination: blob_to_path(row.get_ref(2)?.as_blob()?),
1038 access,
1039 },
1040 ))
1041 })?;
1042 for row in rows {
1043 let (session_id, mount) = row?;
1044 if let Some(session) = state.sessions.get_mut(&session_id) {
1045 session.additional_mounts.push(mount);
1046 }
1047 }
1048 Ok(())
1049}
1050
1051pub(super) fn load_checkpoints(connection: &Connection, state: &mut State) -> Result<()> {
1052 let mut statement = connection.prepare(
1053 "SELECT session_id, archive_path, sha256, created_at, event_frontier FROM session_checkpoints",
1054 )?;
1055 let rows = statement.query_map([], |row| {
1056 Ok((
1057 row.get::<_, String>(0)?,
1058 CheckpointMetadata {
1059 archive_path: blob_to_path(row.get_ref(1)?.as_blob()?),
1060 sha256: row.get(2)?,
1061 created_at: row.get(3)?,
1062 event_frontier: row.get(4)?,
1063 },
1064 ))
1065 })?;
1066 for row in rows {
1067 let (session_id, checkpoint) = row?;
1068 if let Some(session) = state.sessions.get_mut(&session_id) {
1069 session.checkpoint = Some(checkpoint);
1070 }
1071 }
1072 Ok(())
1073}