Skip to main content

mj_controller/database/
state_io.rs

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