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