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