1use super::*;
2
3pub fn load_state() -> Result<State> {
4 if let Some(committed) = committed_state()? {
5 return Ok(committed.state.clone());
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)?.into_iter().collect();
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 review_json
516 ) 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,?31)
517 ON CONFLICT(session_id) DO UPDATE SET
518 title = excluded.title,
519 harness_kind = excluded.harness_kind,
520 last_profile = excluded.last_profile,
521 target_template_id = excluded.target_template_id,
522 state = excluded.state,
523 native_session_id = excluded.native_session_id,
524 acp_session_title = excluded.acp_session_title,
525 session_title_override = excluded.session_title_override,
526 updated_at = excluded.updated_at,
527 viewed_through_event_ordinal = max(
528 sessions.viewed_through_event_ordinal,
529 excluded.viewed_through_event_ordinal
530 ),
531 last_error = excluded.last_error,
532 resource_allocation = excluded.resource_allocation,
533 last_checkpoint_error = excluded.last_checkpoint_error,
534 project_directory = excluded.project_directory,
535 managed_worktree = excluded.managed_worktree,
536 container_cpus = excluded.container_cpus,
537 container_memory = excluded.container_memory,
538 archived = excluded.archived,
539 create_managed_worktree = excluded.create_managed_worktree,
540 subagents = excluded.subagents,
541 container_workspace = excluded.container_workspace,
542 build_cache_json = excluded.build_cache_json,
543 launch_base = excluded.launch_base,
544 target_runtime_json = excluded.target_runtime_json,
545 launch_branch = excluded.launch_branch,
546 publication_json = excluded.publication_json,
547 checkout_json = excluded.checkout_json,
548 project_json = coalesce(excluded.project_json, sessions.project_json),
549 review_json = excluded.review_json",
550 params![
551 session.id,
552 session.title,
553 session.harness_kind.id(),
554 session.last_profile,
555 session.target_template_id,
556 session.state.as_str(),
557 session.native_session_id,
558 session.acp_session_title,
559 session.session_title_override,
560 session.updated_at,
561 session.viewed_through_event_ordinal,
562 session.last_error,
563 session
564 .resource_allocation
565 .as_ref()
566 .map(serde_json::to_string)
567 .transpose()?,
568 session.last_checkpoint_error,
569 session
570 .project_directory
571 .as_ref()
572 .map(|path| path_to_blob(path)),
573 session
574 .managed_worktree
575 .as_ref()
576 .map(serde_json::to_string)
577 .transpose()?,
578 session.container_cpus,
579 session.container_memory,
580 session.archived,
581 session.draft_input,
582 session.create_managed_worktree,
583 session.subagents.as_ref().map(serde_json::to_string).transpose()?,
584 session
585 .container_workspace
586 .as_ref()
587 .map(|path| path.to_string_lossy().into_owned()),
588 session
589 .build_cache
590 .as_ref()
591 .map(serde_json::to_string)
592 .transpose()?,
593 session.launch_base,
594 session.target_runtime.as_ref().map(serde_json::to_string).transpose()?,
595 session.launch_branch,
596 session.publication.as_ref().map(serde_json::to_string).transpose()?,
597 session.checkout.as_ref().map(serde_json::to_string).transpose()?,
598 session.project.as_ref().map(serde_json::to_string).transpose()?,
599 session.review.as_ref().map(serde_json::to_string).transpose()?,
600 ],
601 )?;
602 tx.execute(
603 "INSERT INTO materialized_sessions(session_id) VALUES (?1)
604 ON CONFLICT(session_id) DO NOTHING",
605 [session.id.as_str()],
606 )?;
607 replace_targets(tx, session)?;
608 replace_mounts(tx, &session.id, &session.additional_mounts)?;
609 replace_checkpoint(tx, session)?;
610 super::events::record_session_fault_transition(tx, session, previous_error)?;
611 Ok(())
612}
613
614pub(super) fn update_lifecycle_fields(tx: &Transaction<'_>, session: &SessionRecord) -> Result<()> {
618 let previous_error = super::events::previous_session_error(tx, &session.id)?;
619 let SessionRecord {
626 target_runtime,
627 id,
628 title,
629 harness_kind,
630 last_profile,
631 target_template_id,
632 state,
633 updated_at,
634 viewed_through_event_ordinal,
635 last_error,
636 resource_allocation,
637 last_checkpoint_error,
638 project_directory,
639 managed_worktree,
640 build_cache,
641 workspace_id: _,
642 bundle_id: _,
643 project,
644 create_managed_worktree: _,
645 launch_base: _,
646 launch_branch: _,
647 checkout: _,
648 publication: _,
649 review: _,
651 subagents,
653 additional_mounts: _,
654 container_cpus: _,
655 container_memory: _,
656 container_workspace: _,
657 archived: _,
658 target: _,
660 native_session_id: _,
661 acp_session_title: _,
662 session_title_override: _,
663 created_at: _,
664 draft_input: _,
665 checkpoint: _,
667 } = session;
668 let changed = tx.execute(
669 "UPDATE sessions
672 SET title = ?2,
673 harness_kind = ?3,
674 last_profile = ?4,
675 target_template_id = ?5,
676 state = ?6,
677 updated_at = ?7,
678 viewed_through_event_ordinal = max(viewed_through_event_ordinal, ?8),
679 last_error = ?9,
680 resource_allocation = ?10,
681 last_checkpoint_error = ?11,
682 project_directory = ?12,
683 managed_worktree = ?13,
684 build_cache_json = ?14,
685 target_runtime_json = ?15,
686 subagents = ?16,
687 project_json = coalesce(?17,project_json)
688 WHERE session_id = ?1",
689 params![
690 id,
691 title,
692 harness_kind.id(),
693 last_profile,
694 target_template_id,
695 state.as_str(),
696 updated_at,
697 viewed_through_event_ordinal,
698 last_error,
699 resource_allocation
700 .as_ref()
701 .map(serde_json::to_string)
702 .transpose()?,
703 last_checkpoint_error,
704 project_directory.as_ref().map(|path| path_to_blob(path)),
705 managed_worktree
706 .as_ref()
707 .map(serde_json::to_string)
708 .transpose()?,
709 build_cache
712 .as_ref()
713 .map(serde_json::to_string)
714 .transpose()?,
715 target_runtime
716 .as_ref()
717 .map(serde_json::to_string)
718 .transpose()?,
719 subagents.as_ref().map(serde_json::to_string).transpose()?,
720 project.as_ref().map(serde_json::to_string).transpose()?,
721 ],
722 )?;
723 if changed != 1 {
724 bail!("unknown session {id}");
725 }
726 super::events::record_session_fault_transition(tx, session, previous_error)?;
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 load_targets_selected(connection, state, None)
896}
897
898fn load_targets_selected(
899 connection: &Connection,
900 state: &mut State,
901 id: Option<&str>,
902) -> Result<()> {
903 let predicate = if id.is_some() {
904 " WHERE session_id=?1"
905 } else {
906 ""
907 };
908 let mut statement = connection.prepare(&format!("SELECT session_id, kind, host, resource_id, address, workspace, worker_id, workspace_storage, borrowed_from
909 FROM session_targets{predicate}"))?;
910 let arguments: Vec<&dyn rusqlite::ToSql> = id
911 .as_ref()
912 .map(|id| vec![id as &dyn rusqlite::ToSql])
913 .unwrap_or_default();
914 let rows = statement.query_map(arguments.as_slice(), |row| {
915 let session_id: String = row.get(0)?;
916 let kind: String = row.get(1)?;
917 let host: Option<String> = row.get(2)?;
918 let resource: Option<String> = row.get(3)?;
919 let address: Option<String> = row.get(4)?;
920 let workspace = row.get_ref(5)?.blob_or_null()?.map(blob_to_path);
921 let worker_id: Option<String> = row.get(6)?;
922 let workspace_storage = row
923 .get::<_, Option<String>>(7)?
924 .map(|serialized| {
925 serde_json::from_str(&serialized).map_err(|error| {
926 rusqlite::Error::FromSqlConversionFailure(7, Type::Text, Box::new(error))
927 })
928 })
929 .transpose()?
930 .unwrap_or_default();
931 let borrowed_from: Option<String> = row.get(8)?;
932 let target = match kind.as_str() {
933 "local-bare" => TargetLocator::LocalBare {
934 worker_root: workspace.unwrap(),
935 },
936 "local-podman" => TargetLocator::LocalPodman {
937 borrowed_from,
938 container_id: resource.unwrap(),
939 workspace_storage,
940 },
941 "local-docker" => TargetLocator::LocalDocker {
942 borrowed_from,
943 container_id: resource.unwrap(),
944 },
945 "apple-container" => TargetLocator::AppleContainer {
946 borrowed_from,
947 container_id: resource.unwrap(),
948 },
949 "aws-ec2" => TargetLocator::AwsEc2 {
950 instance_id: resource.unwrap(),
951 address,
952 },
953 "ssh-bare" => TargetLocator::SshBare {
954 host: host.unwrap(),
955 workspace: workspace.unwrap(),
956 worker_id,
957 },
958 "ssh-docker" => TargetLocator::SshDocker {
959 borrowed_from,
960 host: host.unwrap(),
961 container_id: resource.unwrap(),
962 },
963 "ssh-podman" => TargetLocator::SshPodman {
964 borrowed_from,
965 host: host.unwrap(),
966 container_id: resource.unwrap(),
967 workspace_storage,
968 },
969 _ => unreachable!("target kind constrained by schema"),
970 };
971 Ok((session_id, target))
972 })?;
973 for row in rows {
974 let (session_id, target) = row?;
975 if let Some(session) = state.sessions.get_mut(&session_id) {
978 session.target = Some(target);
979 }
980 }
981 Ok(())
982}
983
984pub(super) fn replace_mounts(
991 tx: &rusqlite::Transaction<'_>,
992 session_id: &str,
993 mounts: &[AdditionalMount],
994) -> Result<()> {
995 tx.execute(
996 "DELETE FROM session_mounts WHERE session_id = ?1",
997 [session_id],
998 )?;
999 tx.execute(
1000 "DELETE FROM session_mount_access WHERE session_id = ?1",
1001 [session_id],
1002 )?;
1003 for (ordinal, mount) in mounts.iter().enumerate() {
1004 tx.execute(
1005 "INSERT INTO session_mounts(session_id, ordinal, source, destination, read_only)
1006 VALUES (?1, ?2, ?3, ?4, ?5)",
1007 params![
1008 session_id,
1009 ordinal as i64,
1010 path_to_blob(&mount.source),
1011 path_to_blob(&mount.destination),
1012 mount.access == MountAccess::Ro
1013 ],
1014 )?;
1015 if mount.access == MountAccess::Rw {
1016 tx.execute(
1017 "INSERT INTO session_mount_access(session_id, source, destination, access)
1018 VALUES (?1, ?2, ?3, 'rw')",
1019 params![
1020 session_id,
1021 path_to_blob(&mount.source),
1022 path_to_blob(&mount.destination)
1023 ],
1024 )?;
1025 }
1026 }
1027 Ok(())
1028}
1029
1030pub(super) fn load_mounts(connection: &Connection, state: &mut State) -> Result<()> {
1031 load_mounts_selected(connection, state, None)
1032}
1033
1034fn load_mounts_selected(
1035 connection: &Connection,
1036 state: &mut State,
1037 id: Option<&str>,
1038) -> Result<()> {
1039 let predicate = if id.is_some() {
1040 " WHERE m.session_id=?1"
1041 } else {
1042 ""
1043 };
1044 let mut statement = connection.prepare(&format!(
1045 "SELECT m.session_id, m.source, m.destination, m.read_only, a.access IS NOT NULL
1046 FROM session_mounts m
1047 LEFT JOIN session_mount_access a
1048 ON a.session_id = m.session_id
1049 AND a.source = m.source
1050 AND a.destination = m.destination
1051 {predicate} ORDER BY m.session_id, m.ordinal"
1052 ))?;
1053 let arguments: Vec<&dyn rusqlite::ToSql> = id
1054 .as_ref()
1055 .map(|id| vec![id as &dyn rusqlite::ToSql])
1056 .unwrap_or_default();
1057 let rows = statement.query_map(arguments.as_slice(), |row| {
1058 let access = match (row.get::<_, bool>(3)?, row.get::<_, bool>(4)?) {
1061 (true, _) => MountAccess::Ro,
1062 (false, true) => MountAccess::Rw,
1063 (false, false) => MountAccess::Cow,
1064 };
1065 Ok((
1066 row.get::<_, String>(0)?,
1067 AdditionalMount {
1068 source: blob_to_path(row.get_ref(1)?.as_blob()?),
1069 destination: blob_to_path(row.get_ref(2)?.as_blob()?),
1070 access,
1071 },
1072 ))
1073 })?;
1074 for row in rows {
1075 let (session_id, mount) = row?;
1076 if let Some(session) = state.sessions.get_mut(&session_id) {
1077 session.additional_mounts.push(mount);
1078 }
1079 }
1080 Ok(())
1081}
1082
1083pub(super) fn load_checkpoints(connection: &Connection, state: &mut State) -> Result<()> {
1084 load_checkpoints_selected(connection, state, None)
1085}
1086
1087fn load_checkpoints_selected(
1088 connection: &Connection,
1089 state: &mut State,
1090 id: Option<&str>,
1091) -> Result<()> {
1092 let predicate = if id.is_some() {
1093 " WHERE session_id=?1"
1094 } else {
1095 ""
1096 };
1097 let mut statement = connection.prepare(&format!("SELECT session_id, archive_path, sha256, created_at, event_frontier FROM session_checkpoints{predicate}"))?;
1098 let arguments: Vec<&dyn rusqlite::ToSql> = id
1099 .as_ref()
1100 .map(|id| vec![id as &dyn rusqlite::ToSql])
1101 .unwrap_or_default();
1102 let rows = statement.query_map(arguments.as_slice(), |row| {
1103 Ok((
1104 row.get::<_, String>(0)?,
1105 CheckpointMetadata {
1106 archive_path: blob_to_path(row.get_ref(1)?.as_blob()?),
1107 sha256: row.get(2)?,
1108 created_at: row.get(3)?,
1109 event_frontier: row.get(4)?,
1110 },
1111 ))
1112 })?;
1113 for row in rows {
1114 let (session_id, checkpoint) = row?;
1115 if let Some(session) = state.sessions.get_mut(&session_id) {
1116 session.checkpoint = Some(checkpoint);
1117 }
1118 }
1119 Ok(())
1120}
1121
1122pub fn backfill_target_runtime(
1126 session_id: &str,
1127 target_template_id: &str,
1128 runtime: &mj_core::state::TargetRuntimeSettings,
1129) -> Result<Option<mj_core::state::TargetRuntimeSettings>> {
1130 let session_id = session_id.to_owned();
1131 let target_template_id = target_template_id.to_owned();
1132 let runtime = serde_json::to_string(runtime)?;
1133 submit_database_write("backfill_target_runtime", move |_| {
1134 let mut connection = open(&database_path())?;
1135 let tx = connection.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
1136 tx.execute("UPDATE sessions SET target_runtime_json = ?2 WHERE session_id = ?1 AND target_template_id = ?3 AND target_runtime_json IS NULL",
1137 params![session_id, runtime, target_template_id])?;
1138 let stored: Option<String> = tx.query_row("SELECT target_runtime_json FROM sessions WHERE session_id = ?1 AND target_template_id = ?2",
1140 params![session_id, target_template_id], |row| row.get(0)).optional()?;
1141 let runtime = stored.as_deref().map(serde_json::from_str).transpose()?;
1142 tx.commit()?;
1143 Ok(runtime)
1144 })
1145}
1146
1147const SESSION_QUERY: &str =
1148 "SELECT s.session_id, s.title, s.harness_kind, s.last_profile, c.bundle_id,
1149 s.target_template_id, s.state, s.native_session_id, s.acp_session_title,
1150 s.session_title_override, c.created_at, s.updated_at,
1151 s.viewed_through_event_ordinal, s.last_error, s.resource_allocation,
1152 s.last_checkpoint_error, s.project_directory, s.managed_worktree,
1153 s.draft_input, s.container_cpus, s.container_memory, s.archived
1154 , c.workspace_id, s.create_managed_worktree, s.subagents,
1155 s.container_workspace, s.build_cache_json, s.launch_base, s.target_runtime_json,
1156 s.launch_branch, s.publication_json, s.checkout_json, s.project_json,
1157 s.review_json
1158 FROM sessions s JOIN session_contexts c USING(session_id)";
1159
1160fn decode_session(row: &rusqlite::Row<'_>) -> rusqlite::Result<Option<SessionRecord>> {
1161 let harness_text: String = row.get(2)?;
1165 let Ok(harness_kind) = harness_text.parse() else {
1166 let session_id: String = row.get(0)?;
1167 tracing::warn!(
1168 session_id,
1169 harness = %harness_text,
1170 "session harness is no longer supported; the session is not listed"
1171 );
1172 return Ok(None);
1173 };
1174 Ok(Some(SessionRecord {
1175 project: row
1176 .get::<_, Option<String>>(32)?
1177 .map(|json| {
1178 serde_json::from_str(&json).map_err(|error| {
1179 rusqlite::Error::FromSqlConversionFailure(32, Type::Text, Box::new(error))
1180 })
1181 })
1182 .transpose()?,
1183 target_runtime: row
1184 .get::<_, Option<String>>(28)?
1185 .map(|json| {
1186 serde_json::from_str(&json).map_err(|error| {
1187 rusqlite::Error::FromSqlConversionFailure(28, Type::Text, Box::new(error))
1188 })
1189 })
1190 .transpose()?,
1191 harness_kind,
1192 create_managed_worktree: row.get(23)?,
1193 launch_base: row.get(27)?,
1194 launch_branch: row.get(29)?,
1195 checkout: row
1196 .get::<_, Option<String>>(31)?
1197 .map(|json| {
1198 serde_json::from_str(&json).map_err(|error| {
1199 rusqlite::Error::FromSqlConversionFailure(31, Type::Text, Box::new(error))
1200 })
1201 })
1202 .transpose()?,
1203 publication: row
1204 .get::<_, Option<String>>(30)?
1205 .map(|json| {
1206 serde_json::from_str(&json).map_err(|error| {
1207 rusqlite::Error::FromSqlConversionFailure(30, Type::Text, Box::new(error))
1208 })
1209 })
1210 .transpose()?,
1211 subagents: row
1212 .get::<_, Option<String>>(24)?
1213 .map(|json| {
1214 serde_json::from_str(&json).map_err(|error| {
1215 rusqlite::Error::FromSqlConversionFailure(24, Type::Text, Box::new(error))
1216 })
1217 })
1218 .transpose()?,
1219 container_workspace: row.get::<_, Option<String>>(25)?.map(PathBuf::from),
1220 build_cache: row
1221 .get::<_, Option<String>>(26)?
1222 .as_deref()
1223 .and_then(|text| match serde_json::from_str(text) {
1224 Ok(build_cache) => Some(build_cache),
1225 Err(error) => {
1226 tracing::warn!(%error, "session build cache record is unreadable");
1227 None
1228 }
1229 }),
1230 workspace_id: row.get(22)?,
1231 archived: row.get(21)?,
1232 container_cpus: row.get(19)?,
1233 container_memory: row.get(20)?,
1234 id: row.get(0)?,
1235 title: row.get(1)?,
1236 last_profile: row.get(3)?,
1237 bundle_id: row.get(4)?,
1238 project_directory: row.get_ref(16)?.blob_or_null()?.map(blob_to_path),
1239 managed_worktree: row
1240 .get::<_, Option<String>>(17)?
1241 .map(|json| serde_json::from_str::<ManagedWorktree>(&json))
1242 .transpose()
1243 .map_err(|error| {
1244 rusqlite::Error::FromSqlConversionFailure(
1245 17,
1246 rusqlite::types::Type::Text,
1247 Box::new(error),
1248 )
1249 })?,
1250 review: row
1253 .get::<_, Option<String>>(33)?
1254 .as_deref()
1255 .and_then(|text| match serde_json::from_str(text) {
1256 Ok(review) => Some(review),
1257 Err(error) => {
1258 tracing::warn!(%error, "session review choice is unreadable");
1259 None
1260 }
1261 }),
1262 target_template_id: row.get(5)?,
1263 resource_allocation: row
1264 .get::<_, Option<String>>(14)?
1265 .map(|json| serde_json::from_str::<SessionResourceAllocation>(&json))
1266 .transpose()
1267 .map_err(|error| {
1268 rusqlite::Error::FromSqlConversionFailure(
1269 14,
1270 rusqlite::types::Type::Text,
1271 Box::new(error),
1272 )
1273 })?,
1274 additional_mounts: Vec::new(),
1275 state: stored_session_state(&row.get::<_, String>(6)?),
1276 target: None,
1277 native_session_id: row.get(7)?,
1278 acp_session_title: row
1279 .get::<_, Option<String>>(8)?
1280 .as_deref()
1281 .and_then(mj_core::state::normalize_session_title),
1282 session_title_override: row.get(9)?,
1283 created_at: row.get(10)?,
1284 updated_at: row.get(11)?,
1285 viewed_through_event_ordinal: row.get::<_, u64>(12)?,
1286 draft_input: row.get(18)?,
1287 last_error: row.get(13)?,
1288 last_checkpoint_error: row.get(15)?,
1289 checkpoint: None,
1290 }))
1291}
1292
1293pub fn load_session_record(id: &str) -> Result<Option<SessionRecord>> {
1295 if let Some(committed) = committed_state()? {
1296 return Ok(committed.state.sessions.get(id).cloned());
1297 }
1298 read_durable_session_record(id)
1299}
1300
1301pub(crate) fn read_durable_session_record(id: &str) -> Result<Option<SessionRecord>> {
1304 let mut connection = open_reader(&database_path())?;
1305 let snapshot = connection.transaction()?;
1306 load_session_with(&snapshot, id)
1307}
1308
1309pub(super) fn load_session_with(
1310 connection: &Connection,
1311 id: &str,
1312) -> Result<Option<SessionRecord>> {
1313 let record = connection
1314 .query_row(
1315 &format!("{SESSION_QUERY} WHERE s.session_id=?1"),
1316 [id],
1317 decode_session,
1318 )
1319 .optional()?
1320 .flatten();
1321 let Some(record) = record else {
1322 return Ok(None);
1323 };
1324 let mut state = State::default();
1325 state.sessions.insert(id.to_owned(), record);
1326 load_targets_selected(connection, &mut state, Some(id))?;
1327 load_mounts_selected(connection, &mut state, Some(id))?;
1328 load_checkpoints_selected(connection, &mut state, Some(id))?;
1329 state.validate()?;
1330 Ok(state.sessions.remove(id))
1331}
1332
1333#[cfg(test)]
1334mod targeted_read_tests {
1335 use super::*;
1336
1337 #[test]
1338 fn single_session_read_includes_related_rows_and_ignores_unrelated_invalid_records() {
1339 let directory = tempfile::tempdir().unwrap();
1340 let path = directory.path().join("controller.sqlite");
1341 let expected = super::super::tests::session("selected", "project");
1342 save_session_to(&path, &expected).unwrap();
1343 save_session_to(&path, &super::super::tests::session("unrelated", "project")).unwrap();
1344 let connection = open(&path).unwrap();
1345 connection
1348 .execute(
1349 "UPDATE sessions SET resource_allocation = '[]' WHERE session_id = 'unrelated'",
1350 [],
1351 )
1352 .unwrap();
1353 assert!(load_state_from(&path).is_err());
1354 let mut reader = open_reader(&path).unwrap();
1355 let snapshot = reader.transaction().unwrap();
1356 assert_eq!(
1357 load_session_with(&snapshot, "selected").unwrap(),
1358 Some(expected)
1359 );
1360 assert_eq!(load_session_with(&snapshot, "missing").unwrap(), None);
1361 }
1362}