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