Skip to main content

mj_controller/database/
state_io.rs

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 connection = open_reader(path)?;
9    let mut state = State::default();
10    let mut statement = connection.prepare(
11        "SELECT s.session_id, s.title, s.harness_kind, s.last_profile, c.bundle_id,
12                s.target_template_id, s.state, s.native_session_id, s.acp_session_title,
13                s.session_title_override, c.created_at, s.updated_at,
14                s.viewed_through_event_ordinal, s.last_error, s.resource_allocation,
15                s.last_checkpoint_error, s.project_directory, s.managed_worktree,
16                s.draft_input, s.container_cpus, s.container_memory, s.archived
17                , c.workspace_id, s.create_managed_worktree, s.mjolnir_subagents,
18                s.container_workspace, s.build_cache_json
19         FROM sessions s JOIN session_contexts c USING(session_id)
20         ORDER BY s.session_id",
21    )?;
22    let rows = statement.query_map([], |row| {
23        // A harness Mjolnir no longer supports can still own rows an earlier
24        // release wrote. Skip such a session with a warning rather than
25        // failing the whole listing and hiding every other session with it.
26        let harness_text: String = row.get(2)?;
27        let Ok(harness_kind) = harness_text.parse() else {
28            let session_id: String = row.get(0)?;
29            tracing::warn!(
30                session_id,
31                harness = %harness_text,
32                "session harness is no longer supported; the session is not listed"
33            );
34            return Ok(None);
35        };
36        Ok(Some(SessionRecord {
37            harness_kind,
38            create_managed_worktree: row.get(23)?,
39            mjolnir_subagents: row.get(24)?,
40            container_workspace: row.get::<_, Option<String>>(25)?.map(PathBuf::from),
41            build_cache: row
42                .get::<_, Option<String>>(26)?
43                .as_deref()
44                .and_then(|text| match serde_json::from_str(text) {
45                    Ok(build_cache) => Some(build_cache),
46                    Err(error) => {
47                        tracing::warn!(%error, "session build cache record is unreadable");
48                        None
49                    }
50                }),
51            workspace_id: row.get(22)?,
52            archived: row.get(21)?,
53            container_cpus: row.get(19)?,
54            container_memory: row.get(20)?,
55            id: row.get(0)?,
56            title: row.get(1)?,
57            last_profile: row.get(3)?,
58            bundle_id: row.get(4)?,
59            project_directory: row.get_ref(16)?.blob_or_null()?.map(blob_to_path),
60            managed_worktree: row
61                .get::<_, Option<String>>(17)?
62                .map(|json| serde_json::from_str::<ManagedWorktree>(&json))
63                .transpose()
64                .map_err(|error| {
65                    rusqlite::Error::FromSqlConversionFailure(
66                        17,
67                        rusqlite::types::Type::Text,
68                        Box::new(error),
69                    )
70                })?,
71            target_template_id: row.get(5)?,
72            resource_allocation: row
73                .get::<_, Option<String>>(14)?
74                .map(|json| serde_json::from_str::<SessionResourceAllocation>(&json))
75                .transpose()
76                .map_err(|error| {
77                    rusqlite::Error::FromSqlConversionFailure(
78                        14,
79                        rusqlite::types::Type::Text,
80                        Box::new(error),
81                    )
82                })?,
83            additional_mounts: Vec::new(),
84            state: stored_session_state(&row.get::<_, String>(6)?),
85            target: None,
86            native_session_id: row.get(7)?,
87            acp_session_title: row
88                .get::<_, Option<String>>(8)?
89                .as_deref()
90                .and_then(mj_core::state::normalize_session_title),
91            session_title_override: row.get(9)?,
92            created_at: row.get(10)?,
93            updated_at: row.get(11)?,
94            viewed_through_event_ordinal: row.get::<_, u64>(12)?,
95            draft_input: row.get(18)?,
96            last_error: row.get(13)?,
97            last_checkpoint_error: row.get(15)?,
98            checkpoint: None,
99        }))
100    })?;
101    for row in rows {
102        if let Some(session) = row? {
103            state.sessions.insert(session.id.clone(), session);
104        }
105    }
106    let mut statement = connection.prepare(
107        "SELECT child_session_id, record_json FROM subagent_sessions ORDER BY child_session_id",
108    )?;
109    let rows = statement.query_map([], |row| {
110        let child_id = row.get::<_, String>(0)?;
111        let json = row.get::<_, String>(1)?;
112        let record = serde_json::from_str::<SubagentRecord>(&json).map_err(|error| {
113            rusqlite::Error::FromSqlConversionFailure(1, Type::Text, Box::new(error))
114        })?;
115        Ok((child_id, record))
116    })?;
117    for row in rows {
118        let (child_id, record) = row?;
119        // A relation whose child or parent is not among the sessions this load
120        // returned describes nothing. Keeping it would fail the state check
121        // below and make every later operation fail with it, which is how a
122        // sub-agent spawn came to be refused with "sub-agent ... has no child
123        // session" long after the child in question was gone (#1065). The load
124        // already skips a session whose harness it cannot parse; a relation
125        // that pointed at such a session is the same kind of residue, and
126        // `save_state_to` deletes these rows on the next save.
127        let missing = if !state.sessions.contains_key(&child_id) {
128            Some("child")
129        } else if !state.sessions.contains_key(&record.parent_session_id) {
130            Some("parent")
131        } else {
132            None
133        };
134        if let Some(missing) = missing {
135            tracing::warn!(
136                child_session_id = child_id,
137                parent_session_id = record.parent_session_id,
138                missing,
139                "dropping a sub-agent relation whose session is not in this state"
140            );
141            continue;
142        }
143        state.subagents.insert(child_id, record);
144    }
145    load_targets(&connection, &mut state)?;
146    load_mounts(&connection, &mut state)?;
147    load_checkpoints(&connection, &mut state)?;
148    let mut statement =
149        connection.prepare("SELECT host, source FROM mount_history ORDER BY host, ordinal")?;
150    let rows = statement.query_map([], |row| {
151        Ok((
152            row.get::<_, String>(0)?,
153            blob_to_path(row.get_ref(1)?.as_blob()?),
154        ))
155    })?;
156    for row in rows {
157        let (host, source) = row?;
158        state.mount_history.entry(host).or_default().push(source);
159    }
160    let mut statement = connection
161        .prepare("SELECT host, cpus, memory_bytes FROM host_container_sizes ORDER BY host")?;
162    let rows = statement.query_map([], |row| {
163        Ok((
164            row.get::<_, String>(0)?,
165            HostContainerSize {
166                cpus: row.get::<_, i64>(1)? as u64,
167                memory_bytes: row.get::<_, i64>(2)? as u64,
168            },
169        ))
170    })?;
171    for row in rows {
172        let (host, size) = row?;
173        state.container_sizes.insert(host, size);
174    }
175    state.validate()?;
176    Ok(state)
177}
178
179pub fn save_state(state: &State) -> Result<()> {
180    let state = state.clone();
181    submit_database_write("save_state", move |_| {
182        save_state_to(&database_path(), &state)
183    })
184}
185
186pub fn save_state_to(path: &Path, state: &State) -> Result<()> {
187    state.validate()?;
188    let mut connection = open(path)?;
189    let tx = connection.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
190    let existing_contexts = existing_contexts(&tx)?;
191    let existing_sessions = {
192        let mut statement = tx.prepare("SELECT session_id FROM sessions")?;
193        statement
194            .query_map([], |row| row.get::<_, String>(0))?
195            .collect::<rusqlite::Result<Vec<_>>>()?
196    };
197    tx.execute(
198        "DELETE FROM subagent_sessions
199         WHERE child_session_id NOT IN (SELECT session_id FROM sessions)
200            OR parent_session_id NOT IN (SELECT session_id FROM sessions)",
201        [],
202    )?;
203    let existing_subagents = {
204        let mut statement =
205            tx.prepare("SELECT child_session_id, parent_session_id FROM subagent_sessions")?;
206        statement
207            .query_map([], |row| {
208                Ok((row.get::<_, String>(0)?, row.get::<_, String>(1)?))
209            })?
210            .collect::<rusqlite::Result<Vec<_>>>()?
211    };
212    for (child_id, parent_id) in existing_subagents {
213        if !state.subagents.contains_key(&child_id)
214            || !state.sessions.contains_key(&child_id)
215            || !state.sessions.contains_key(&parent_id)
216        {
217            tx.execute(
218                "DELETE FROM subagent_sessions WHERE child_session_id = ?1",
219                [child_id],
220            )?;
221        }
222    }
223    for session_id in existing_sessions {
224        if !state.sessions.contains_key(&session_id) {
225            tx.execute("DELETE FROM sessions WHERE session_id = ?1", [session_id])?;
226        }
227    }
228    tx.execute("DELETE FROM mount_history", [])?;
229    tx.execute("DELETE FROM host_container_sizes", [])?;
230    for session in state.sessions.values() {
231        if let Some((existing_bundle, existing_workspace)) = existing_contexts.get(&session.id) {
232            ensure!(
233                existing_bundle == &session.bundle_id,
234                "session {} was already associated with bundle {}, not {}",
235                session.id,
236                existing_bundle,
237                session.bundle_id
238            );
239            ensure!(
240                existing_workspace == &session.workspace_id,
241                "session {} was already associated with workspace {}, not {}",
242                session.id,
243                existing_workspace,
244                session.workspace_id
245            );
246        }
247        insert_session(&tx, session)?;
248    }
249    for subagent in state.subagents.values() {
250        let record_json = serde_json::to_string(subagent)?;
251        tx.execute(
252            "INSERT INTO subagent_sessions(
253                 child_session_id, parent_session_id, request_key, record_json
254             ) VALUES (?1, ?2, ?3, ?4)
255             ON CONFLICT(child_session_id) DO UPDATE SET
256                 parent_session_id = excluded.parent_session_id,
257                 request_key = excluded.request_key,
258                 record_json = excluded.record_json",
259            params![
260                subagent.child_session_id,
261                subagent.parent_session_id,
262                subagent.request_key,
263                record_json
264            ],
265        )?;
266    }
267    for (host, sources) in &state.mount_history {
268        for (ordinal, source) in sources.iter().enumerate() {
269            tx.execute(
270                "INSERT INTO mount_history(host, source, ordinal) VALUES (?1, ?2, ?3)",
271                params![host, path_to_blob(source), ordinal as i64],
272            )?;
273        }
274    }
275    for (host, size) in &state.container_sizes {
276        write_host_container_size(&tx, host, *size)?;
277    }
278    tx.commit()?;
279    Ok(())
280}
281
282pub(super) fn existing_contexts(
283    tx: &Transaction<'_>,
284) -> Result<BTreeMap<String, (String, String)>> {
285    let mut statement =
286        tx.prepare("SELECT session_id, bundle_id, workspace_id FROM session_contexts")?;
287    let rows = statement.query_map([], |row| Ok((row.get(0)?, (row.get(1)?, row.get(2)?))))?;
288    rows.collect::<rusqlite::Result<_>>().map_err(Into::into)
289}
290
291pub(super) fn session_exists(tx: &Transaction<'_>, session_id: &str) -> Result<bool> {
292    Ok(tx
293        .query_row(
294            "SELECT 1 FROM sessions WHERE session_id = ?1",
295            [session_id],
296            |_| Ok(()),
297        )
298        .optional()?
299        .is_some())
300}
301
302pub(super) fn write_materialized_session(
303    tx: &Transaction<'_>,
304    materialized: &MaterializedSession,
305) -> Result<()> {
306    let (execution, running_started_at_ms) = materialized_execution_columns(materialized.execution);
307    tx.execute(
308        "INSERT INTO materialized_sessions(
309             session_id, applied_event_ordinal, applied_event_digest, execution_state,
310             running_started_at_ms, session_title, configuration_json, last_activity_at_ms,
311             pending_elicitations_json, active_turn_json, last_turn_outcome_json
312         ) VALUES (?1,?2,?3,?4,?5,?6,?7,?8,?9,?10,?11)
313         ON CONFLICT(session_id) DO UPDATE SET
314             applied_event_ordinal = excluded.applied_event_ordinal,
315             applied_event_digest = excluded.applied_event_digest,
316             execution_state = excluded.execution_state,
317             running_started_at_ms = excluded.running_started_at_ms,
318             session_title = excluded.session_title,
319             configuration_json = excluded.configuration_json,
320             last_activity_at_ms = excluded.last_activity_at_ms,
321             pending_elicitations_json = excluded.pending_elicitations_json,
322             active_turn_json = excluded.active_turn_json,
323             last_turn_outcome_json = excluded.last_turn_outcome_json",
324        params![
325            materialized.session_id,
326            materialized.applied_event_ordinal,
327            materialized.applied_event_digest,
328            execution,
329            running_started_at_ms,
330            materialized.session_title,
331            serde_json::to_string(&materialized.configuration)?,
332            materialized.last_activity_at_ms,
333            serde_json::to_string(&materialized.pending_elicitations)?,
334            materialized
335                .active_turn
336                .as_ref()
337                .map(serde_json::to_string)
338                .transpose()?,
339            materialized
340                .last_turn_outcome
341                .as_ref()
342                .map(serde_json::to_string)
343                .transpose()?,
344        ],
345    )?;
346    tx.execute(
347        "DELETE FROM materialized_transcript_items WHERE session_id = ?1",
348        [materialized.session_id.as_str()],
349    )?;
350    for item in &materialized.transcript {
351        upsert_transcript_item(tx, &materialized.session_id, item)?;
352    }
353    replace_materialized_queue(tx, &materialized.session_id, &materialized.queued_prompts)?;
354    Ok(())
355}
356
357pub(super) fn upsert_transcript_item(
358    tx: &Transaction<'_>,
359    session_id: &str,
360    item: &TranscriptItem,
361) -> Result<()> {
362    let existing = tx
363        .query_row(
364            "SELECT position, latest_content_event_ordinal, created_at_ms, last_changed_at_ms
365             FROM materialized_transcript_items
366             WHERE session_id = ?1 AND stable_id = ?2",
367            params![session_id, item.stable_id],
368            |row| {
369                Ok((
370                    row.get::<_, u64>(0)?,
371                    row.get::<_, Option<u64>>(1)?,
372                    row.get::<_, i64>(2)?,
373                    row.get::<_, i64>(3)?,
374                ))
375            },
376        )
377        .optional()?;
378    if let Some((position, latest_content_event_ordinal, created_at_ms, last_changed_at_ms)) =
379        existing
380    {
381        if position != item.position || created_at_ms != item.created_at_ms {
382            return Err(ProjectionIntegrityError(format!(
383                "transcript item {:?} changed immutable identity fields",
384                item.stable_id
385            ))
386            .into());
387        }
388        if item.last_changed_at_ms < last_changed_at_ms {
389            return Err(ProjectionIntegrityError(format!(
390                "transcript item {:?} moved its changed timestamp backwards",
391                item.stable_id
392            ))
393            .into());
394        }
395        if latest_content_event_ordinal.is_some_and(|existing| {
396            item.latest_content_event_ordinal
397                .is_none_or(|next| next < existing)
398        }) {
399            return Err(ProjectionIntegrityError(format!(
400                "transcript item {:?} moved its latest content ordinal backwards",
401                item.stable_id
402            ))
403            .into());
404        }
405        tx.execute(
406            "UPDATE materialized_transcript_items
407             SET latest_content_event_ordinal = ?3, last_changed_at_ms = ?4, body_json = ?5
408             WHERE session_id = ?1 AND stable_id = ?2",
409            params![
410                session_id,
411                item.stable_id,
412                item.latest_content_event_ordinal,
413                item.last_changed_at_ms,
414                serde_json::to_string(&item.body)?,
415            ],
416        )?;
417    } else {
418        tx.execute(
419            "INSERT INTO materialized_transcript_items(
420                 session_id, stable_id, position, latest_content_event_ordinal,
421                 created_at_ms, last_changed_at_ms, body_json
422             ) VALUES (?1,?2,?3,?4,?5,?6,?7)",
423            params![
424                session_id,
425                item.stable_id,
426                item.position,
427                item.latest_content_event_ordinal,
428                item.created_at_ms,
429                item.last_changed_at_ms,
430                serde_json::to_string(&item.body)?,
431            ],
432        )?;
433    }
434    Ok(())
435}
436
437pub(super) fn replace_materialized_queue(
438    tx: &Transaction<'_>,
439    session_id: &str,
440    queued_prompts: &[MaterializedQueuedPrompt],
441) -> Result<()> {
442    let mut command_ids = BTreeSet::new();
443    for prompt in queued_prompts {
444        if prompt.command_id.trim().is_empty() {
445            bail!("materialized prompt queue has an empty command id");
446        }
447        if !command_ids.insert(prompt.command_id.as_str()) {
448            bail!(
449                "materialized prompt queue contains duplicate command {:?}",
450                prompt.command_id
451            );
452        }
453    }
454    tx.execute(
455        "DELETE FROM materialized_queued_prompts WHERE session_id = ?1",
456        [session_id],
457    )?;
458    for (ordinal, prompt) in queued_prompts.iter().enumerate() {
459        tx.execute(
460            "INSERT INTO materialized_queued_prompts(
461                 session_id, ordinal, command_id, kind_json, content_json, queued_at_ms,
462                 accepted_ordinal
463             ) VALUES (?1,?2,?3,?4,?5,?6,?7)",
464            params![
465                session_id,
466                ordinal as i64,
467                prompt.command_id,
468                serde_json::to_string(&prompt.kind)?,
469                serde_json::to_string(&prompt.content)?,
470                prompt.queued_at_ms,
471                prompt.accepted_ordinal,
472            ],
473        )?;
474    }
475    Ok(())
476}
477
478pub(super) fn materialized_execution_columns(
479    execution: MaterializedExecutionState,
480) -> (&'static str, Option<i64>) {
481    match execution {
482        MaterializedExecutionState::Idle => ("idle", None),
483        MaterializedExecutionState::Running { started_at_ms } => ("running", Some(started_at_ms)),
484        MaterializedExecutionState::Closing => ("closing", None),
485        MaterializedExecutionState::Closed => ("closed", None),
486    }
487}
488
489pub(super) fn parse_materialized_execution(
490    execution: &str,
491    running_started_at_ms: Option<i64>,
492) -> Result<MaterializedExecutionState> {
493    match (execution, running_started_at_ms) {
494        ("idle", None) => Ok(MaterializedExecutionState::Idle),
495        ("running", Some(started_at_ms)) => {
496            Ok(MaterializedExecutionState::Running { started_at_ms })
497        }
498        ("closing", None) => Ok(MaterializedExecutionState::Closing),
499        ("closed", None) => Ok(MaterializedExecutionState::Closed),
500        _ => bail!("invalid materialized execution state {execution:?}"),
501    }
502}
503
504/// Write every field of a session, including the ones other writers own.
505/// Only a flow that authors the whole record — creation, import, resume, or
506/// orphan adoption — may use this.
507pub(super) fn insert_session(tx: &Transaction<'_>, session: &SessionRecord) -> Result<()> {
508    tx.execute(
509        "INSERT INTO session_contexts(session_id, bundle_id, created_at, workspace_id)
510         VALUES (?1, ?2, ?3, ?4)
511         ON CONFLICT(session_id) DO NOTHING",
512        params![
513            session.id,
514            session.bundle_id,
515            session.created_at,
516            session.workspace_id
517        ],
518    )?;
519    let (stored_bundle, stored_workspace): (String, String) = tx.query_row(
520        "SELECT bundle_id, workspace_id FROM session_contexts WHERE session_id = ?1",
521        [session.id.as_str()],
522        |row| Ok((row.get(0)?, row.get(1)?)),
523    )?;
524    ensure!(
525        stored_bundle == session.bundle_id,
526        "session {} belongs to bundle {}, not {}",
527        session.id,
528        stored_bundle,
529        session.bundle_id
530    );
531    ensure!(
532        stored_workspace == session.workspace_id,
533        "session {} belongs to workspace {}, not {}",
534        session.id,
535        stored_workspace,
536        session.workspace_id
537    );
538    tx.execute(
539        "INSERT INTO sessions(
540             session_id, title, harness_kind, last_profile, target_template_id, state,
541             native_session_id, acp_session_title, session_title_override, updated_at,
542             viewed_through_event_ordinal, last_error, resource_allocation,
543             last_checkpoint_error, project_directory, managed_worktree,
544             container_cpus, container_memory, archived, draft_input, create_managed_worktree,
545             mjolnir_subagents, container_workspace, build_cache_json
546         ) 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)
547         ON CONFLICT(session_id) DO UPDATE SET
548             title = excluded.title,
549             harness_kind = excluded.harness_kind,
550             last_profile = excluded.last_profile,
551             target_template_id = excluded.target_template_id,
552             state = excluded.state,
553             native_session_id = excluded.native_session_id,
554             acp_session_title = excluded.acp_session_title,
555             session_title_override = excluded.session_title_override,
556             updated_at = excluded.updated_at,
557             viewed_through_event_ordinal = max(
558                 sessions.viewed_through_event_ordinal,
559                 excluded.viewed_through_event_ordinal
560             ),
561             last_error = excluded.last_error,
562             resource_allocation = excluded.resource_allocation,
563             last_checkpoint_error = excluded.last_checkpoint_error,
564             project_directory = excluded.project_directory,
565             managed_worktree = excluded.managed_worktree,
566             container_cpus = excluded.container_cpus,
567             container_memory = excluded.container_memory,
568             archived = excluded.archived,
569             create_managed_worktree = excluded.create_managed_worktree,
570             mjolnir_subagents = excluded.mjolnir_subagents,
571             container_workspace = excluded.container_workspace,
572             build_cache_json = excluded.build_cache_json",
573        params![
574            session.id,
575            session.title,
576            session.harness_kind.id(),
577            session.last_profile,
578            session.target_template_id,
579            session.state.as_str(),
580            session.native_session_id,
581            session.acp_session_title,
582            session.session_title_override,
583            session.updated_at,
584            session.viewed_through_event_ordinal,
585            session.last_error,
586            session
587                .resource_allocation
588                .as_ref()
589                .map(serde_json::to_string)
590                .transpose()?,
591            session.last_checkpoint_error,
592            session
593                .project_directory
594                .as_ref()
595                .map(|path| path_to_blob(path)),
596            session
597                .managed_worktree
598                .as_ref()
599                .map(serde_json::to_string)
600                .transpose()?,
601            session.container_cpus,
602            session.container_memory,
603            session.archived,
604            session.draft_input,
605            session.create_managed_worktree,
606            session.mjolnir_subagents,
607            session
608                .container_workspace
609                .as_ref()
610                .map(|path| path.to_string_lossy().into_owned()),
611            session
612                .build_cache
613                .as_ref()
614                .map(serde_json::to_string)
615                .transpose()?,
616        ],
617    )?;
618    tx.execute(
619        "INSERT INTO materialized_sessions(session_id) VALUES (?1)
620         ON CONFLICT(session_id) DO NOTHING",
621        [session.id.as_str()],
622    )?;
623    replace_targets(tx, session)?;
624    replace_mounts(tx, &session.id, &session.additional_mounts)?;
625    replace_checkpoint(tx, session)?;
626    Ok(())
627}
628
629/// Update the columns a lifecycle transition owns, plus the target locator
630/// that provisioning and teardown maintain with them. The row must exist:
631/// a transition never resurrects a session another writer deleted.
632pub(super) fn update_lifecycle_fields(tx: &Transaction<'_>, session: &SessionRecord) -> Result<()> {
633    // Destructured exhaustively and without `..` on purpose. This statement is
634    // the only thing standing between a new `SessionRecord` field and a value
635    // that is set in memory, read back as its default, and never missed until
636    // somebody inspects the database. A new field breaks this binding, and
637    // whoever adds it decides then whether a lifecycle transition owns it.
638    // Everything bound to `_` is owned by `upsert_session` instead.
639    let SessionRecord {
640        id,
641        title,
642        harness_kind,
643        last_profile,
644        target_template_id,
645        state,
646        updated_at,
647        viewed_through_event_ordinal,
648        last_error,
649        resource_allocation,
650        last_checkpoint_error,
651        project_directory,
652        managed_worktree,
653        build_cache,
654        workspace_id: _,
655        bundle_id: _,
656        create_managed_worktree: _,
657        mjolnir_subagents: _,
658        additional_mounts: _,
659        container_cpus: _,
660        container_memory: _,
661        container_workspace: _,
662        archived: _,
663        // Written by `replace_targets` below rather than by this statement.
664        target: _,
665        native_session_id: _,
666        acp_session_title: _,
667        session_title_override: _,
668        created_at: _,
669        draft_input: _,
670        // Written by `replace_checkpoint`.
671        checkpoint: _,
672    } = session;
673    let changed = tx.execute(
674        // The detach ordinal only ever moves forward, so a transition that
675        // started before a detach receipt cannot rewind it.
676        "UPDATE sessions
677         SET title = ?2,
678             harness_kind = ?3,
679             last_profile = ?4,
680             target_template_id = ?5,
681             state = ?6,
682             updated_at = ?7,
683             viewed_through_event_ordinal = max(viewed_through_event_ordinal, ?8),
684             last_error = ?9,
685             resource_allocation = ?10,
686             last_checkpoint_error = ?11,
687             project_directory = ?12,
688             managed_worktree = ?13,
689             build_cache_json = ?14
690         WHERE session_id = ?1",
691        params![
692            id,
693            title,
694            harness_kind.id(),
695            last_profile,
696            target_template_id,
697            state.as_str(),
698            updated_at,
699            viewed_through_event_ordinal,
700            last_error,
701            resource_allocation
702                .as_ref()
703                .map(serde_json::to_string)
704                .transpose()?,
705            last_checkpoint_error,
706            project_directory.as_ref().map(|path| path_to_blob(path)),
707            managed_worktree
708                .as_ref()
709                .map(serde_json::to_string)
710                .transpose()?,
711            // Resolved while a session is provisioned and assigned to the
712            // record right before this write, so the lifecycle path owns it.
713            build_cache
714                .as_ref()
715                .map(serde_json::to_string)
716                .transpose()?,
717        ],
718    )?;
719    if changed != 1 {
720        bail!("unknown session {id}");
721    }
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    let mut statement = connection.prepare(
891        "SELECT session_id, kind, host, resource_id, address, workspace, worker_id, workspace_storage, borrowed_from
892         FROM session_targets",
893    )?;
894    let rows = statement.query_map([], |row| {
895        let session_id: String = row.get(0)?;
896        let kind: String = row.get(1)?;
897        let host: Option<String> = row.get(2)?;
898        let resource: Option<String> = row.get(3)?;
899        let address: Option<String> = row.get(4)?;
900        let workspace = row.get_ref(5)?.blob_or_null()?.map(blob_to_path);
901        let worker_id: Option<String> = row.get(6)?;
902        let workspace_storage = row
903            .get::<_, Option<String>>(7)?
904            .map(|serialized| {
905                serde_json::from_str(&serialized).map_err(|error| {
906                    rusqlite::Error::FromSqlConversionFailure(7, Type::Text, Box::new(error))
907                })
908            })
909            .transpose()?
910            .unwrap_or_default();
911        let borrowed_from: Option<String> = row.get(8)?;
912        let target = match kind.as_str() {
913            "local-bare" => TargetLocator::LocalBare {
914                worker_root: workspace.unwrap(),
915            },
916            "local-podman" => TargetLocator::LocalPodman {
917                borrowed_from,
918                container_id: resource.unwrap(),
919                workspace_storage,
920            },
921            "local-docker" => TargetLocator::LocalDocker {
922                borrowed_from,
923                container_id: resource.unwrap(),
924            },
925            "apple-container" => TargetLocator::AppleContainer {
926                borrowed_from,
927                container_id: resource.unwrap(),
928            },
929            "aws-ec2" => TargetLocator::AwsEc2 {
930                instance_id: resource.unwrap(),
931                address,
932            },
933            "ssh-bare" => TargetLocator::SshBare {
934                host: host.unwrap(),
935                workspace: workspace.unwrap(),
936                worker_id,
937            },
938            "ssh-docker" => TargetLocator::SshDocker {
939                borrowed_from,
940                host: host.unwrap(),
941                container_id: resource.unwrap(),
942            },
943            "ssh-podman" => TargetLocator::SshPodman {
944                borrowed_from,
945                host: host.unwrap(),
946                container_id: resource.unwrap(),
947                workspace_storage,
948            },
949            _ => unreachable!("target kind constrained by schema"),
950        };
951        Ok((session_id, target))
952    })?;
953    for row in rows {
954        let (session_id, target) = row?;
955        // A session skipped for an unsupported harness has no entry to attach
956        // its target, mounts, or checkpoint to.
957        if let Some(session) = state.sessions.get_mut(&session_id) {
958            session.target = Some(target);
959        }
960    }
961    Ok(())
962}
963
964/// Rewrite a session's attached directories.
965///
966/// `session_mounts.read_only` keeps the meaning older builds understand, so
967/// read-write mounts are stored there as not read-only and recorded again in
968/// `session_mount_access`. Older builds rewrite `session_mounts` without
969/// touching that table, which is what keeps the read-write choice.
970pub(super) fn replace_mounts(
971    tx: &rusqlite::Transaction<'_>,
972    session_id: &str,
973    mounts: &[AdditionalMount],
974) -> Result<()> {
975    tx.execute(
976        "DELETE FROM session_mounts WHERE session_id = ?1",
977        [session_id],
978    )?;
979    tx.execute(
980        "DELETE FROM session_mount_access WHERE session_id = ?1",
981        [session_id],
982    )?;
983    for (ordinal, mount) in mounts.iter().enumerate() {
984        tx.execute(
985            "INSERT INTO session_mounts(session_id, ordinal, source, destination, read_only)
986             VALUES (?1, ?2, ?3, ?4, ?5)",
987            params![
988                session_id,
989                ordinal as i64,
990                path_to_blob(&mount.source),
991                path_to_blob(&mount.destination),
992                mount.access == MountAccess::Ro
993            ],
994        )?;
995        if mount.access == MountAccess::Rw {
996            tx.execute(
997                "INSERT INTO session_mount_access(session_id, source, destination, access)
998                 VALUES (?1, ?2, ?3, 'rw')",
999                params![
1000                    session_id,
1001                    path_to_blob(&mount.source),
1002                    path_to_blob(&mount.destination)
1003                ],
1004            )?;
1005        }
1006    }
1007    Ok(())
1008}
1009
1010pub(super) fn load_mounts(connection: &Connection, state: &mut State) -> Result<()> {
1011    let mut statement = connection.prepare(
1012        "SELECT m.session_id, m.source, m.destination, m.read_only, a.access IS NOT NULL
1013         FROM session_mounts m
1014         LEFT JOIN session_mount_access a
1015             ON a.session_id = m.session_id
1016             AND a.source = m.source
1017             AND a.destination = m.destination
1018         ORDER BY m.session_id, m.ordinal",
1019    )?;
1020    let rows = statement.query_map([], |row| {
1021        // An older build that made the mount read-only left the access row
1022        // behind; its later choice wins.
1023        let access = match (row.get::<_, bool>(3)?, row.get::<_, bool>(4)?) {
1024            (true, _) => MountAccess::Ro,
1025            (false, true) => MountAccess::Rw,
1026            (false, false) => MountAccess::Cow,
1027        };
1028        Ok((
1029            row.get::<_, String>(0)?,
1030            AdditionalMount {
1031                source: blob_to_path(row.get_ref(1)?.as_blob()?),
1032                destination: blob_to_path(row.get_ref(2)?.as_blob()?),
1033                access,
1034            },
1035        ))
1036    })?;
1037    for row in rows {
1038        let (session_id, mount) = row?;
1039        if let Some(session) = state.sessions.get_mut(&session_id) {
1040            session.additional_mounts.push(mount);
1041        }
1042    }
1043    Ok(())
1044}
1045
1046pub(super) fn load_checkpoints(connection: &Connection, state: &mut State) -> Result<()> {
1047    let mut statement = connection.prepare(
1048        "SELECT session_id, archive_path, sha256, created_at, event_frontier FROM session_checkpoints",
1049    )?;
1050    let rows = statement.query_map([], |row| {
1051        Ok((
1052            row.get::<_, String>(0)?,
1053            CheckpointMetadata {
1054                archive_path: blob_to_path(row.get_ref(1)?.as_blob()?),
1055                sha256: row.get(2)?,
1056                created_at: row.get(3)?,
1057                event_frontier: row.get(4)?,
1058            },
1059        ))
1060    })?;
1061    for row in rows {
1062        let (session_id, checkpoint) = row?;
1063        if let Some(session) = state.sessions.get_mut(&session_id) {
1064            session.checkpoint = Some(checkpoint);
1065        }
1066    }
1067    Ok(())
1068}