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