Skip to main content

mj_controller/database/
sessions.rs

1use super::*;
2
3/// Persist one operational session without rewriting unrelated controller
4/// state. Dashboard lifecycle jobs use this path so independent jobs can
5/// commit concurrently without restoring stale copies of other sessions.
6pub fn save_session(session: &SessionRecord) -> Result<()> {
7    let session = session.clone();
8    submit_database_write("save_session", move |_| {
9        save_session_to(&database_path(), &session)
10    })
11}
12
13pub fn save_new_session(
14    session: &SessionRecord,
15    container_size: Option<(String, HostContainerSize)>,
16) -> Result<()> {
17    let session = session.clone();
18    submit_database_write("save_new_session", move |_| {
19        save_new_session_to(&database_path(), &session, container_size)
20    })
21}
22
23pub(super) fn save_new_session_to(
24    path: &Path,
25    session: &SessionRecord,
26    container_size: Option<(String, HostContainerSize)>,
27) -> Result<()> {
28    let mut connection = open(path)?;
29    let tx = connection.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
30    validate_session_record(session)?;
31    insert_session(&tx, session)?;
32    if let Some((host, size)) = container_size {
33        write_host_container_size(&tx, &host, size)?;
34    }
35    if session.harness_kind.supports_delegation_tools() {
36        tx.execute("INSERT INTO subagent_preference(singleton, policy) VALUES(1, ?1) ON CONFLICT(singleton) DO UPDATE SET policy = excluded.policy", [serde_json::to_string(&session.subagents.clone().unwrap_or_default())?])?;
37    }
38    tx.commit()?;
39    Ok(())
40}
41
42/// Update publication evidence only while the exact stopped checkpoint is
43/// still current. An age scan may race a user's Resume or a new checkpoint.
44pub fn set_publication_assessment_if_current(
45    session_id: &str,
46    assessment: &mj_core::state::PublicationAssessment,
47) -> Result<bool> {
48    let session_id = session_id.to_owned();
49    let assessment = assessment.clone();
50    submit_database_write("set_publication_assessment_if_current", move |_| {
51        let connection = open(&database_path())?;
52        let updated = connection.execute(
53            "UPDATE sessions SET publication_json = ?2
54             WHERE session_id = ?1 AND state = 'stopped'
55               AND EXISTS (SELECT 1 FROM session_checkpoints c
56                           WHERE c.session_id = ?1 AND c.sha256 = ?3)",
57            params![
58                session_id,
59                serde_json::to_string(&assessment)?,
60                assessment.checkpoint_sha256
61            ],
62        )?;
63        Ok(updated == 1)
64    })
65}
66
67/// Persist a borrowed-target child and its parent relationship atomically.
68pub fn save_subagent_session(
69    session: &SessionRecord,
70    subagent: &mj_core::subagent::SubagentRecord,
71) -> Result<()> {
72    let session = session.clone();
73    let subagent = subagent.clone();
74    submit_database_write("save_subagent_session", move |_| {
75        let mut connection = open(&database_path())?;
76        let tx = connection.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
77        insert_session(&tx, &session)?;
78        super::usage::record_child_accounting(&tx, &subagent)?;
79        tx.execute(
80            "INSERT INTO subagent_sessions(
81                 child_session_id, parent_session_id, request_key, record_json
82             ) VALUES (?1, ?2, ?3, ?4)",
83            params![
84                subagent.child_session_id,
85                subagent.parent_session_id,
86                subagent.request_key,
87                serde_json::to_string(&subagent)?,
88            ],
89        )?;
90        tx.commit()?;
91        Ok(())
92    })
93}
94
95/// Record the child turn whose completion notice the parent already has.
96pub fn mark_subagent_turn_noticed(child_session_id: &str, turn: u64) -> Result<()> {
97    let child_session_id = child_session_id.to_owned();
98    submit_database_write("mark_subagent_turn_noticed", move |_| {
99        let mut relation = load_subagent(&child_session_id)?
100            .with_context(|| format!("unknown sub-agent session {child_session_id}"))?;
101        relation.noticed_turn = Some(turn);
102        let json = serde_json::to_string(&relation)?;
103        let connection = open(&database_path())?;
104        connection.execute(
105            "UPDATE subagent_sessions SET record_json = ?2 WHERE child_session_id = ?1",
106            params![child_session_id, json],
107        )?;
108        Ok(())
109    })
110}
111
112pub fn load_subagent(child_session_id: &str) -> Result<Option<mj_core::subagent::SubagentRecord>> {
113    let connection = open_reader(&database_path())?;
114    connection
115        .query_row(
116            "SELECT record_json FROM subagent_sessions WHERE child_session_id = ?1",
117            [child_session_id],
118            |row| row.get::<_, String>(0),
119        )
120        .optional()?
121        .map(|json| serde_json::from_str(&json).context("decode sub-agent record"))
122        .transpose()
123}
124
125pub fn list_subagents(parent_session_id: &str) -> Result<Vec<mj_core::subagent::SubagentRecord>> {
126    let connection = open_reader(&database_path())?;
127    let mut statement = connection.prepare(
128        "SELECT record_json FROM subagent_sessions
129         WHERE parent_session_id = ?1 ORDER BY rowid",
130    )?;
131    statement
132        .query_map([parent_session_id], |row| row.get::<_, String>(0))?
133        .map(|row| serde_json::from_str(&row?).context("decode sub-agent record"))
134        .collect()
135}
136
137/// Record, on the parent's session, the sub-agents its suspend is about to
138/// stop. A child already listed keeps its place and takes the newer details,
139/// so a suspend that runs again after a restart lists each child once.
140pub fn record_stopped_subagents(
141    parent_session_id: &str,
142    stopped: &[mj_core::subagent::StoppedSubagent],
143) -> Result<()> {
144    let parent_session_id = parent_session_id.to_owned();
145    let stopped = stopped.to_vec();
146    submit_database_write("record_stopped_subagents", move |_| {
147        record_stopped_subagents_to(&database_path(), &parent_session_id, &stopped)
148    })
149}
150
151pub(super) fn record_stopped_subagents_to(
152    path: &Path,
153    parent_session_id: &str,
154    stopped: &[mj_core::subagent::StoppedSubagent],
155) -> Result<()> {
156    let mut connection = open(path)?;
157    let tx = connection.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
158    for child in stopped {
159        tx.execute(
160            "INSERT INTO stopped_subagents(parent_session_id, child_session_id, record_json)
161             VALUES (?1, ?2, ?3)
162             ON CONFLICT(parent_session_id, child_session_id) DO UPDATE SET
163                 record_json = excluded.record_json",
164            params![
165                parent_session_id,
166                child.child_session_id,
167                serde_json::to_string(child)?
168            ],
169        )?;
170    }
171    tx.commit()?;
172    Ok(())
173}
174
175/// The sub-agents a suspend of this session stopped and its model has not
176/// been told about yet, in the order they were recorded.
177pub fn load_stopped_subagents(
178    parent_session_id: &str,
179) -> Result<Vec<mj_core::subagent::StoppedSubagent>> {
180    load_stopped_subagents_from(&database_path(), parent_session_id)
181}
182
183pub(super) fn load_stopped_subagents_from(
184    path: &Path,
185    parent_session_id: &str,
186) -> Result<Vec<mj_core::subagent::StoppedSubagent>> {
187    let connection = open_reader(path)?;
188    let mut statement = connection.prepare(
189        "SELECT record_json FROM stopped_subagents
190         WHERE parent_session_id = ?1 ORDER BY rowid",
191    )?;
192    statement
193        .query_map([parent_session_id], |row| row.get::<_, String>(0))?
194        .map(|row| serde_json::from_str(&row?).context("decode stopped sub-agent record"))
195        .collect()
196}
197
198/// Forget the stopped sub-agents whose note the parent's relay has taken.
199/// Only the named children go, so a child listed after the note was built
200/// waits for the next one.
201pub fn clear_stopped_subagents(
202    parent_session_id: &str,
203    child_session_ids: &[String],
204) -> Result<()> {
205    let parent_session_id = parent_session_id.to_owned();
206    let child_session_ids = child_session_ids.to_vec();
207    submit_database_write("clear_stopped_subagents", move |_| {
208        clear_stopped_subagents_from(&database_path(), &parent_session_id, &child_session_ids)
209    })
210}
211
212pub(super) fn clear_stopped_subagents_from(
213    path: &Path,
214    parent_session_id: &str,
215    child_session_ids: &[String],
216) -> Result<()> {
217    let mut connection = open(path)?;
218    let tx = connection.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
219    for child_session_id in child_session_ids {
220        tx.execute(
221            "DELETE FROM stopped_subagents
222             WHERE parent_session_id = ?1 AND child_session_id = ?2",
223            params![parent_session_id, child_session_id],
224        )?;
225    }
226    tx.commit()?;
227    Ok(())
228}
229
230/// Everything recorded about one child's report. A child with nothing
231/// recorded yet reads as the empty report.
232/// One session's stored lifecycle state, or `None` when the store holds no
233/// such session.
234pub fn load_session_state(session_id: &str) -> Result<Option<SessionState>> {
235    let connection = open_reader(&database_path())?;
236    let stored = connection
237        .query_row(
238            "SELECT state FROM sessions WHERE session_id = ?1",
239            [session_id],
240            |row| row.get::<_, String>(0),
241        )
242        .optional()?;
243    Ok(stored.as_deref().map(stored_session_state))
244}
245
246pub fn load_subagent_report(child_session_id: &str) -> Result<mj_core::subagent::SubagentReport> {
247    load_subagent_report_from(&database_path(), child_session_id)
248}
249
250pub(super) fn load_subagent_report_from(
251    path: &Path,
252    child_session_id: &str,
253) -> Result<mj_core::subagent::SubagentReport> {
254    Ok(load_subagent_report_with(&*open_reader(path)?, child_session_id)?.unwrap_or_default())
255}
256
257/// A child's recorded report, or `None` when nothing is recorded for it.
258pub(super) fn load_subagent_report_with(
259    connection: &Connection,
260    child_session_id: &str,
261) -> Result<Option<mj_core::subagent::SubagentReport>> {
262    let row = connection
263        .prepare_cached(
264            "SELECT handback_command_id, handback_message, handback_recorded_at_ms,
265                    reminder_command_id, reminder_for_command_id, reminder_sent_at_ms,
266                    reminder_failed_for_command_id, awaited_ordinal, report_dir
267             FROM subagent_handbacks WHERE child_session_id = ?1",
268        )?
269        .query_row([child_session_id], |row| {
270            Ok((
271                row.get::<_, Option<String>>(0)?,
272                row.get::<_, Option<String>>(1)?,
273                row.get::<_, Option<i64>>(2)?,
274                row.get::<_, Option<String>>(3)?,
275                row.get::<_, Option<String>>(4)?,
276                row.get::<_, Option<i64>>(5)?,
277                row.get::<_, Option<String>>(6)?,
278                row.get::<_, Option<i64>>(7)?,
279                row.get::<_, Option<String>>(8)?,
280            ))
281        })
282        .optional()?;
283    let Some((
284        handback_command,
285        handback_message,
286        handback_at,
287        reminder_command,
288        reminder_for,
289        reminder_at,
290        reminder_failed_for,
291        awaited_ordinal,
292        report_dir,
293    )) = row
294    else {
295        return Ok(None);
296    };
297    Ok(Some(mj_core::subagent::SubagentReport {
298        handback: match (handback_command, handback_message, handback_at) {
299            (Some(command_id), Some(message), Some(recorded_at_ms)) => {
300                Some(mj_core::subagent::SubagentHandback {
301                    command_id,
302                    message,
303                    recorded_at_ms,
304                })
305            }
306            _ => None,
307        },
308        reminder: match (reminder_command, reminder_for, reminder_at) {
309            (Some(command_id), Some(for_command_id), Some(sent_at_ms)) => {
310                Some(mj_core::subagent::HandbackReminder {
311                    command_id,
312                    for_command_id,
313                    sent_at_ms,
314                })
315            }
316            _ => None,
317        },
318        reminder_failed_for,
319        awaited_ordinal: awaited_ordinal.and_then(|ordinal| u64::try_from(ordinal).ok()),
320        report_dir,
321    }))
322}
323
324/// Record the directory Mjolnir created for a child's report files. It is
325/// recorded only for a sub-agent child.
326pub fn record_subagent_report_dir(child_session_id: &str, report_dir: &str) -> Result<()> {
327    let child_session_id = child_session_id.to_owned();
328    let report_dir = report_dir.to_owned();
329    submit_database_write("record_subagent_report_dir", move |_| {
330        record_subagent_report_dir_to(&database_path(), &child_session_id, &report_dir)
331    })
332}
333
334pub(super) fn record_subagent_report_dir_to(
335    path: &Path,
336    child_session_id: &str,
337    report_dir: &str,
338) -> Result<()> {
339    open(path)?.execute(
340        "INSERT INTO subagent_handbacks(child_session_id, report_dir)
341         SELECT ?1, ?2 WHERE EXISTS (
342             SELECT 1 FROM subagent_sessions WHERE child_session_id = ?1
343         )
344         ON CONFLICT(child_session_id) DO UPDATE SET report_dir = excluded.report_dir",
345        params![child_session_id, report_dir],
346    )?;
347    Ok(())
348}
349
350/// Record that the parent gave a child a prompt, accepted at `ordinal`. It is
351/// recorded only for a sub-agent child, and never moves backwards.
352pub fn record_subagent_prompt(child_session_id: &str, ordinal: u64) -> Result<()> {
353    let child_session_id = child_session_id.to_owned();
354    submit_database_write("record_subagent_prompt", move |_| {
355        record_subagent_prompt_to(&database_path(), &child_session_id, ordinal)
356    })
357}
358
359pub(super) fn record_subagent_prompt_to(
360    path: &Path,
361    child_session_id: &str,
362    ordinal: u64,
363) -> Result<()> {
364    let ordinal = i64::try_from(ordinal).context("prompt ordinal exceeds the store's range")?;
365    open(path)?.execute(
366        "INSERT INTO subagent_handbacks(child_session_id, awaited_ordinal)
367         SELECT ?1, ?2 WHERE EXISTS (
368             SELECT 1 FROM subagent_sessions WHERE child_session_id = ?1
369         )
370         ON CONFLICT(child_session_id) DO UPDATE SET
371             awaited_ordinal = max(coalesce(awaited_ordinal, 0), excluded.awaited_ordinal)",
372        params![child_session_id, ordinal],
373    )?;
374    Ok(())
375}
376
377/// Record a child's report for the turn it names, unless a report for that
378/// turn is already recorded. Returns whether this one was recorded: a child
379/// delivers one report per turn, and the check and the write are one
380/// transaction so two racing calls cannot both be accepted.
381pub fn record_subagent_handback(
382    child_session_id: &str,
383    handback: &mj_core::subagent::SubagentHandback,
384) -> Result<bool> {
385    let child_session_id = child_session_id.to_owned();
386    let handback = handback.clone();
387    submit_database_write("record_subagent_handback", move |_| {
388        record_subagent_handback_to(&database_path(), &child_session_id, &handback)
389    })
390}
391
392pub(super) fn record_subagent_handback_to(
393    path: &Path,
394    child_session_id: &str,
395    handback: &mj_core::subagent::SubagentHandback,
396) -> Result<bool> {
397    let mut connection = open(path)?;
398    let tx = connection.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
399    let recorded_for: Option<String> = tx
400        .query_row(
401            "SELECT handback_command_id FROM subagent_handbacks WHERE child_session_id = ?1",
402            [child_session_id],
403            |row| row.get(0),
404        )
405        .optional()?
406        .flatten();
407    if recorded_for.as_deref() == Some(handback.command_id.as_str()) {
408        return Ok(false);
409    }
410    tx.execute(
411        "INSERT INTO subagent_handbacks(
412             child_session_id, handback_command_id, handback_message, handback_recorded_at_ms
413         ) VALUES (?1, ?2, ?3, ?4)
414         ON CONFLICT(child_session_id) DO UPDATE SET
415             handback_command_id = excluded.handback_command_id,
416             handback_message = excluded.handback_message,
417             handback_recorded_at_ms = excluded.handback_recorded_at_ms",
418        params![
419            child_session_id,
420            handback.command_id,
421            handback.message,
422            handback.recorded_at_ms
423        ],
424    )?;
425    tx.commit()?;
426    Ok(true)
427}
428
429/// Record the reminder sent after a child's turn ended without a report.
430pub fn record_handback_reminder(
431    child_session_id: &str,
432    reminder: &mj_core::subagent::HandbackReminder,
433) -> Result<()> {
434    let child_session_id = child_session_id.to_owned();
435    let reminder = reminder.clone();
436    submit_database_write("record_handback_reminder", move |_| {
437        open(&database_path())?.execute(
438            "INSERT INTO subagent_handbacks(
439                 child_session_id, reminder_command_id, reminder_for_command_id, reminder_sent_at_ms
440             ) VALUES (?1, ?2, ?3, ?4)
441             ON CONFLICT(child_session_id) DO UPDATE SET
442                 reminder_command_id = excluded.reminder_command_id,
443                 reminder_for_command_id = excluded.reminder_for_command_id,
444                 reminder_sent_at_ms = excluded.reminder_sent_at_ms",
445            params![
446                child_session_id,
447                reminder.command_id,
448                reminder.for_command_id,
449                reminder.sent_at_ms
450            ],
451        )?;
452        Ok(())
453    })
454}
455
456/// Record that the reminder for a turn could not be sent, so the child's
457/// last message stands as its report.
458pub fn record_handback_reminder_failed(child_session_id: &str, for_command_id: &str) -> Result<()> {
459    let child_session_id = child_session_id.to_owned();
460    let for_command_id = for_command_id.to_owned();
461    submit_database_write("record_handback_reminder_failed", move |_| {
462        open(&database_path())?.execute(
463            "INSERT INTO subagent_handbacks(child_session_id, reminder_failed_for_command_id)
464             VALUES (?1, ?2)
465             ON CONFLICT(child_session_id) DO UPDATE SET
466                 reminder_failed_for_command_id = excluded.reminder_failed_for_command_id",
467            params![child_session_id, for_command_id],
468        )?;
469        Ok(())
470    })
471}
472
473pub fn lookup_subagent_request(
474    parent_session_id: &str,
475    request_key: &str,
476) -> Result<Option<mj_core::subagent::SubagentRecord>> {
477    let connection = open_reader(&database_path())?;
478    connection
479        .query_row(
480            "SELECT record_json FROM subagent_sessions
481             WHERE parent_session_id = ?1 AND request_key = ?2",
482            params![parent_session_id, request_key],
483            |row| row.get::<_, String>(0),
484        )
485        .optional()?
486        .map(|json| serde_json::from_str(&json).context("decode sub-agent record"))
487        .transpose()
488}
489
490/// Persist a session and the container size it most recently launched on its
491/// host in one transaction.
492pub fn save_session_with_container_size(
493    session: &SessionRecord,
494    host: &str,
495    size: HostContainerSize,
496) -> Result<()> {
497    let session = session.clone();
498    let host = host.to_owned();
499    submit_database_write("save_session_with_container_size", move |_| {
500        save_session_with_container_size_to(&database_path(), &session, Some((&host, size)))
501    })
502}
503
504/// Resume owns the target, checkout and attached resources, never a client's
505/// title, draft, archive visibility or the worker's current display title.
506pub fn save_resumed_session(
507    session: &SessionRecord,
508    container_size: Option<(&str, HostContainerSize)>,
509) -> Result<()> {
510    let session = session.clone();
511    let container_size = container_size.map(|(host, size)| (host.to_owned(), size));
512    submit_database_write("save_resumed_session", move |_| {
513        save_resumed_session_to(
514            &database_path(),
515            &session,
516            container_size
517                .as_ref()
518                .map(|(host, size)| (host.as_str(), *size)),
519        )
520    })
521}
522
523pub(super) fn save_resumed_session_to(
524    path: &Path,
525    session: &SessionRecord,
526    container_size: Option<(&str, HostContainerSize)>,
527) -> Result<()> {
528    validate_session_record(session)?;
529    let mut connection = open(path)?;
530    let tx = connection.transaction()?;
531    let (bundle, workspace): (String, String) = tx.query_row(
532        "SELECT bundle_id, workspace_id FROM session_contexts WHERE session_id = ?1",
533        [&session.id],
534        |row| Ok((row.get(0)?, row.get(1)?)),
535    )?;
536    ensure!(
537        bundle == session.bundle_id && workspace == session.workspace_id,
538        "session {} context changed before resume publication",
539        session.id
540    );
541    update_lifecycle_fields(&tx, session)?;
542    tx.execute(
543        "UPDATE sessions SET native_session_id = ?2, container_cpus = ?3,
544             container_memory = ?4, container_workspace = ?5,
545             create_managed_worktree = ?6, launch_base = ?7, launch_branch = ?8,
546             checkout_json = ?9
547         WHERE session_id = ?1",
548        params![
549            session.id,
550            session.native_session_id,
551            session.container_cpus,
552            session.container_memory,
553            session
554                .container_workspace
555                .as_ref()
556                .map(|path| path.to_string_lossy().into_owned()),
557            session.create_managed_worktree,
558            session.launch_base,
559            session.launch_branch,
560            session
561                .checkout
562                .as_ref()
563                .map(serde_json::to_string)
564                .transpose()?,
565        ],
566    )?;
567    replace_mounts(&tx, &session.id, &session.additional_mounts)?;
568    replace_checkpoint(&tx, session)?;
569    if let Some((host, size)) = container_size {
570        write_host_container_size(&tx, host, size)?;
571    }
572    tx.commit()?;
573    Ok(())
574}
575
576/// Update only the fields a lifecycle transition owns on a session that
577/// already exists. Everything else — display titles, checkpoints, container
578/// settings, and attached directories — stays with its own writer.
579pub fn save_lifecycle_session(session: &SessionRecord) -> Result<()> {
580    let session = session.clone();
581    submit_database_write("save_lifecycle_session", move |_| {
582        save_lifecycle_session_to(&database_path(), &session)
583    })
584}
585
586/// Install a lifecycle transition together with the checkpoint it just
587/// verified and the harness session id that produced it.
588pub fn save_checkpointed_session(session: &SessionRecord) -> Result<()> {
589    let session = session.clone();
590    submit_database_write("save_checkpointed_session", move |_| {
591        save_checkpointed_session_to(&database_path(), &session)
592    })
593}
594
595/// Recover lifecycle rows stranded by a process exit during checkpoint
596/// creation. This must be called once by the top-level controller process
597/// while it owns the controller-store guard, not by per-operation reloads.
598pub fn recover_interrupted_checkpointing_sessions(updated_at: &str) -> Result<usize> {
599    let updated_at = updated_at.to_owned();
600    submit_database_write("recover_interrupted_checkpointing_sessions", move |_| {
601        recover_interrupted_checkpointing_sessions_to(&database_path(), &updated_at)
602    })
603}
604
605/// Change only the user-owned display name. This avoids writing a stale
606/// SessionRecord over independently committed checkpoint or relay metadata.
607pub fn set_session_title_override(session_id: &str, title: &str, updated_at: &str) -> Result<()> {
608    let session_id = session_id.to_owned();
609    let title = title.to_owned();
610    let updated_at = updated_at.to_owned();
611    submit_database_write("set_session_title_override", move |_| {
612        set_session_title_override_to(&database_path(), &session_id, &title, &updated_at)
613    })
614}
615
616/// Rewrite a configured profile id in every persisted session in one SQLite
617/// transaction. Configuration is stored separately, so the controller owns
618/// coordinating this update with the matching config-map rename.
619pub fn rename_profile_references(old_id: &str, new_id: &str) -> Result<usize> {
620    rename_session_reference("last_profile", old_id, new_id)
621}
622
623/// Rewrite a configured target id in every persisted session in one SQLite
624/// transaction.
625pub fn rename_target_references(old_id: &str, new_id: &str) -> Result<usize> {
626    rename_session_reference("target_template_id", old_id, new_id)
627}
628
629pub(super) fn rename_session_reference(
630    column: &'static str,
631    old_id: &str,
632    new_id: &str,
633) -> Result<usize> {
634    ensure!(
635        matches!(column, "last_profile" | "target_template_id"),
636        "unsupported session reference column"
637    );
638    let old_id = old_id.to_owned();
639    let new_id = new_id.to_owned();
640    submit_database_write("rename_session_reference", move |_| {
641        rename_session_reference_at(&database_path(), column, &old_id, &new_id)
642    })
643}
644
645pub(super) fn rename_session_reference_at(
646    path: &Path,
647    column: &str,
648    old_id: &str,
649    new_id: &str,
650) -> Result<usize> {
651    let mut connection = open(path)?;
652    let tx = connection.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
653    let changed = tx.execute(
654        &format!("UPDATE sessions SET {column} = ?2 WHERE {column} = ?1"),
655        params![old_id, new_id],
656    )?;
657    tx.commit()?;
658    Ok(changed)
659}
660
661/// Change only whether the resume dialog hides this session. Archiving is a
662/// display choice, so it has its own writer and never rewrites lifecycle,
663/// checkpoint, or title columns another task owns.
664pub fn set_session_archived(session_id: &str, archived: bool) -> Result<()> {
665    let session_id = session_id.to_owned();
666    submit_database_write("set_session_archived", move |_| {
667        set_session_archived_to(&database_path(), &session_id, archived)
668    })
669}
670
671/// Record that the managed target of an otherwise live session is definitively
672/// gone. A verified checkpoint keeps the session recoverable as an error on the
673/// dashboard; without one, the session is lost. The state predicate keeps a
674/// late poll result from overwriting a concurrent lifecycle transition.
675pub fn mark_session_target_missing(
676    session_id: &str,
677    detail: &str,
678    updated_at: &str,
679) -> Result<Option<SessionState>> {
680    let session_id = session_id.to_owned();
681    let detail = detail.to_owned();
682    let updated_at = updated_at.to_owned();
683    submit_database_write("mark_session_target_missing", move |_| {
684        mark_session_target_missing_to(&database_path(), &session_id, &detail, &updated_at)
685    })
686}
687
688pub(super) fn mark_session_target_missing_to(
689    path: &Path,
690    session_id: &str,
691    detail: &str,
692    updated_at: &str,
693) -> Result<Option<SessionState>> {
694    mark_session_target_missing_if_current_to(path, session_id, detail, updated_at, None)
695}
696
697/// Record a definitive worker failure only while the observed session record
698/// is still current. A delayed background write must not invalidate a resume.
699pub fn mark_session_target_missing_if_current(
700    session_id: &str,
701    detail: &str,
702    updated_at: &str,
703    observed_updated_at: &str,
704) -> Result<Option<SessionState>> {
705    let session_id = session_id.to_owned();
706    let detail = detail.to_owned();
707    let updated_at = updated_at.to_owned();
708    let observed_updated_at = observed_updated_at.to_owned();
709    submit_database_write("mark_session_target_missing_if_current", move |_| {
710        mark_session_target_missing_if_current_to(
711            &database_path(),
712            &session_id,
713            &detail,
714            &updated_at,
715            Some(&observed_updated_at),
716        )
717    })
718}
719
720pub(super) fn mark_session_target_missing_if_current_to(
721    path: &Path,
722    session_id: &str,
723    detail: &str,
724    updated_at: &str,
725    observed_updated_at: Option<&str>,
726) -> Result<Option<SessionState>> {
727    let mut connection = open(path)?;
728    let tx = connection.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
729    let previous_error = super::events::previous_session_error(&tx, session_id)?;
730    let changed = tx.execute(
731        "UPDATE sessions
732         SET state = CASE
733                 WHEN EXISTS(
734                     SELECT 1 FROM session_checkpoints
735                     WHERE session_checkpoints.session_id = sessions.session_id
736                 ) THEN 'error'
737                 ELSE 'lost'
738             END,
739             last_error = ?2,
740             updated_at = ?3
741         WHERE session_id = ?1
742           AND (?4 IS NULL OR updated_at = ?4)
743           AND state IN ('provisioning', 'running', 'disconnected', 'error')",
744        params![session_id, detail, updated_at, observed_updated_at],
745    )?;
746    ensure!(changed <= 1, "updated {changed} sessions for {session_id}");
747    let state = if changed == 1 {
748        let stored: String = tx.query_row(
749            "SELECT state FROM sessions WHERE session_id = ?1",
750            [session_id],
751            |row| row.get(0),
752        )?;
753        Some(stored_session_state(&stored))
754    } else {
755        None
756    };
757    if changed == 1 && previous_error.as_deref() != Some(detail) {
758        super::events::insert_api_event(
759            &tx,
760            session_id,
761            Utc::now().timestamp_millis(),
762            &ApiEventData::SessionFault {
763                reason: mj_core::event_outcome::OutcomeReason::RuntimeUnavailable,
764                message: detail.into(),
765                command_id: None,
766            },
767        )?;
768    }
769    tx.commit()?;
770    Ok(state)
771}
772
773pub(super) fn set_session_archived_to(path: &Path, session_id: &str, archived: bool) -> Result<()> {
774    let connection = open(path)?;
775    let changed = connection.execute(
776        "UPDATE sessions SET archived = ?2 WHERE session_id = ?1",
777        params![session_id, archived],
778    )?;
779    if changed != 1 {
780        bail!("unknown session {session_id}");
781    }
782    Ok(())
783}
784
785/// Native sessions the resume dialog hides. Hel never writes into a harness
786/// home, so the hidden set lives here instead of in the harness's own store.
787pub fn hidden_native_sessions() -> Result<BTreeSet<(mj_core::config::HarnessKind, String)>> {
788    hidden_native_sessions_from(&database_path())
789}
790
791pub(super) fn hidden_native_sessions_from(
792    path: &Path,
793) -> Result<BTreeSet<(mj_core::config::HarnessKind, String)>> {
794    let connection = open_reader(path)?;
795    let mut statement =
796        connection.prepare("SELECT harness_kind, native_session_id FROM hidden_native_sessions")?;
797    let rows = statement.query_map([], |row| {
798        Ok((row.get::<_, String>(0)?, row.get::<_, String>(1)?))
799    })?;
800    let mut hidden = BTreeSet::new();
801    for row in rows {
802        let (harness, native_session_id) = row?;
803        // Rows for a harness this release no longer supports are ignored, not
804        // fatal; they simply hide nothing.
805        match harness.parse::<mj_core::config::HarnessKind>() {
806            Ok(harness) => {
807                hidden.insert((harness, native_session_id));
808            }
809            Err(_) => tracing::warn!(
810                harness = %harness,
811                "ignoring a hidden native session for a harness that is no longer supported"
812            ),
813        }
814    }
815    Ok(hidden)
816}
817
818/// Hide or reveal one native session in the resume dialog.
819pub fn set_native_session_hidden(
820    harness: mj_core::config::HarnessKind,
821    native_session_id: &str,
822    hidden: bool,
823) -> Result<()> {
824    let native_session_id = native_session_id.to_owned();
825    submit_database_write("set_native_session_hidden", move |_| {
826        set_native_session_hidden_to(&database_path(), harness, &native_session_id, hidden)
827    })
828}
829
830pub(super) fn set_native_session_hidden_to(
831    path: &Path,
832    harness: mj_core::config::HarnessKind,
833    native_session_id: &str,
834    hidden: bool,
835) -> Result<()> {
836    if native_session_id.trim().is_empty() {
837        bail!("native session id is empty");
838    }
839    let connection = open(path)?;
840    if hidden {
841        connection.execute(
842            "INSERT INTO hidden_native_sessions(harness_kind, native_session_id, hidden_at)
843             VALUES (?1, ?2, ?3)
844             ON CONFLICT(harness_kind, native_session_id) DO NOTHING",
845            params![harness.id(), native_session_id, Utc::now().to_rfc3339()],
846        )?;
847    } else {
848        connection.execute(
849            "DELETE FROM hidden_native_sessions
850             WHERE harness_kind = ?1 AND native_session_id = ?2",
851            params![harness.id(), native_session_id],
852        )?;
853    }
854    Ok(())
855}
856
857/// Change only the per-session container provisioning inputs: the size
858/// overrides and the attached directories. Everything else the session row
859/// owns is left to its own writer.
860pub fn set_session_container_settings(
861    session_id: &str,
862    cpus: Option<&str>,
863    memory: Option<&str>,
864    mounts: &[AdditionalMount],
865    updated_at: &str,
866) -> Result<()> {
867    let session_id = session_id.to_owned();
868    let cpus = cpus.map(str::to_owned);
869    let memory = memory.map(str::to_owned);
870    let mounts = mounts.to_vec();
871    let updated_at = updated_at.to_owned();
872    submit_database_write("set_session_container_settings", move |_| {
873        set_session_container_settings_to(
874            &database_path(),
875            &session_id,
876            cpus.as_deref(),
877            memory.as_deref(),
878            &mounts,
879            &updated_at,
880        )
881    })
882}
883
884pub(super) fn set_session_container_settings_to(
885    path: &Path,
886    session_id: &str,
887    cpus: Option<&str>,
888    memory: Option<&str>,
889    mounts: &[AdditionalMount],
890    updated_at: &str,
891) -> Result<()> {
892    if updated_at.trim().is_empty() {
893        bail!("session update timestamp is empty");
894    }
895    crate::targets::validate_additional_mounts(mounts)?;
896    let mut connection = open(path)?;
897    let tx = connection.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
898    let changed = tx.execute(
899        "UPDATE sessions
900         SET container_cpus = ?2, container_memory = ?3, updated_at = ?4
901         WHERE session_id = ?1",
902        params![session_id, cpus, memory, updated_at],
903    )?;
904    if changed != 1 {
905        bail!("unknown session {session_id}");
906    }
907    replace_mounts(&tx, session_id, mounts)?;
908    tx.commit()?;
909    Ok(())
910}
911
912pub(super) fn set_session_title_override_to(
913    path: &Path,
914    session_id: &str,
915    title: &str,
916    updated_at: &str,
917) -> Result<()> {
918    if title.trim().is_empty() {
919        bail!("session title is empty");
920    }
921    if updated_at.trim().is_empty() {
922        bail!("session update timestamp is empty");
923    }
924    let connection = open(path)?;
925    let changed = connection.execute(
926        "UPDATE sessions
927         SET session_title_override = ?2, updated_at = ?3
928         WHERE session_id = ?1",
929        params![session_id, title, updated_at],
930    )?;
931    if changed != 1 {
932        bail!("unknown session {session_id}");
933    }
934    Ok(())
935}
936
937/// Persist the latest ACP-provided title without replacing unrelated session
938/// fields that may have changed in another supervised controller task.
939pub fn set_session_acp_title(session_id: &str, title: Option<&str>) -> Result<()> {
940    let session_id = session_id.to_owned();
941    let title = title.map(str::to_owned);
942    submit_database_write("set_session_acp_title", move |_| {
943        set_session_acp_title_to(&database_path(), &session_id, title.as_deref())
944    })
945}
946
947pub(super) fn set_session_acp_title_to(
948    path: &Path,
949    session_id: &str,
950    title: Option<&str>,
951) -> Result<()> {
952    if title.is_some_and(|title| title.trim().is_empty()) {
953        bail!("ACP session title is empty");
954    }
955    let title = title.and_then(mj_core::state::normalize_session_title);
956    let connection = open(path)?;
957    let changed = connection.execute(
958        "UPDATE sessions SET acp_session_title = ?2 WHERE session_id = ?1",
959        params![session_id, title],
960    )?;
961    if changed != 1 {
962        bail!("unknown session {session_id}");
963    }
964    Ok(())
965}
966
967/// Commit the successful handshake for a newly provisioned worker without
968/// replacing checkpoint or display metadata owned by other controller tasks.
969pub fn mark_session_worker_connected(
970    session_id: &str,
971    native_session_id: Option<&str>,
972    updated_at: &str,
973) -> Result<()> {
974    let session_id = session_id.to_owned();
975    let native_session_id = native_session_id.map(str::to_owned);
976    let updated_at = updated_at.to_owned();
977    submit_database_write("mark_session_worker_connected", move |_| {
978        mark_session_worker_connected_to(
979            &database_path(),
980            &session_id,
981            native_session_id.as_deref(),
982            &updated_at,
983        )
984    })
985}
986
987/// Point a session at a native session its worker opened on its own. Only that
988/// column moves: the session's lifecycle state belongs to whatever operation is
989/// running.
990pub fn adopt_native_session_id(session_id: &str, native_session_id: &str) -> Result<()> {
991    let session_id = session_id.to_owned();
992    let native_session_id = native_session_id.to_owned();
993    submit_database_write("adopt_native_session_id", move |_| {
994        adopt_native_session_id_to(&database_path(), &session_id, &native_session_id)
995    })
996}
997
998pub(super) fn adopt_native_session_id_to(
999    path: &Path,
1000    session_id: &str,
1001    native_session_id: &str,
1002) -> Result<()> {
1003    let connection = open(path)?;
1004    let changed = connection.execute(
1005        "UPDATE sessions SET native_session_id = ?2 WHERE session_id = ?1",
1006        params![session_id, native_session_id],
1007    )?;
1008    if changed != 1 {
1009        bail!("unknown session {session_id}");
1010    }
1011    Ok(())
1012}
1013
1014pub(super) fn mark_session_worker_connected_to(
1015    path: &Path,
1016    session_id: &str,
1017    native_session_id: Option<&str>,
1018    updated_at: &str,
1019) -> Result<()> {
1020    if updated_at.trim().is_empty() {
1021        bail!("worker connection timestamp is empty");
1022    }
1023    let connection = open(path)?;
1024    let changed = connection.execute(
1025        "UPDATE sessions
1026         SET state = 'running',
1027             native_session_id = coalesce(?2, native_session_id),
1028             updated_at = ?3,
1029             last_error = NULL
1030         WHERE session_id = ?1",
1031        params![session_id, native_session_id, updated_at],
1032    )?;
1033    if changed != 1 {
1034        bail!("unknown session {session_id}");
1035    }
1036    Ok(())
1037}
1038
1039pub(super) fn recover_interrupted_checkpointing_sessions_to(
1040    path: &Path,
1041    updated_at: &str,
1042) -> Result<usize> {
1043    ensure!(
1044        !updated_at.trim().is_empty(),
1045        "checkpoint recovery timestamp is empty"
1046    );
1047    let mut connection = open(path)?;
1048    let tx = connection.transaction()?;
1049    let ids = {
1050        let mut query = tx.prepare("SELECT session_id FROM checkpoint_operations")?;
1051        query
1052            .query_map([], |row| row.get::<_, String>(0))?
1053            .collect::<rusqlite::Result<Vec<_>>>()?
1054    };
1055    for id in ids {
1056        super::events::finish_checkpoint_operation(
1057            &tx,
1058            &id,
1059            mj_core::event_outcome::CommandResultKind::Failed,
1060            Some(mj_core::event_outcome::OutcomeReason::ControllerRestarted),
1061            Some("Checkpoint operation interrupted by controller restart".into()),
1062        )?;
1063    }
1064    let changed = tx.execute("UPDATE sessions SET state = 'running', updated_at = ?1, last_checkpoint_error = ?2 WHERE state = 'checkpointing'",
1065        params![updated_at, "checkpointing was interrupted by a controller restart; the target was left running"])?;
1066    tx.commit()?;
1067    Ok(changed)
1068}
1069
1070pub(super) fn save_session_to(path: &Path, session: &SessionRecord) -> Result<()> {
1071    save_session_with_container_size_to(path, session, None)
1072}
1073
1074pub(super) fn save_session_with_container_size_to(
1075    path: &Path,
1076    session: &SessionRecord,
1077    container_size: Option<(&str, HostContainerSize)>,
1078) -> Result<()> {
1079    validate_session_record(session)?;
1080
1081    let mut connection = open(path)?;
1082    let tx = connection.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
1083    if let Some(existing_bundle) = tx
1084        .query_row(
1085            "SELECT bundle_id FROM session_contexts WHERE session_id = ?1",
1086            [session.id.as_str()],
1087            |row| row.get::<_, String>(0),
1088        )
1089        .optional()?
1090        && existing_bundle
1091            != super::projects::session_project_id(&tx, &session.id, &session.bundle_id)?
1092    {
1093        bail!(
1094            "session {} was already associated with bundle {}, not {}",
1095            session.id,
1096            existing_bundle,
1097            session.bundle_id
1098        );
1099    }
1100    let mut session = session.clone();
1101    let moving: bool = tx.query_row(
1102        "SELECT EXISTS(SELECT 1 FROM session_moves WHERE session_id=?1
1103         AND json_extract(operation_json, '$.phase') IN ('preparing','closing_source','resuming_destination','starting_queue'))",
1104        [&session.id], |row| row.get(0),
1105    )?;
1106    if moving {
1107        // A Move may provision for minutes while clients keep editing drafts
1108        // and titles. Merge these independently owned fields in this same
1109        // transaction rather than restoring the lifecycle's earlier copy.
1110        let (draft, title, acp_title, viewed, archived) = tx.query_row(
1111            "SELECT draft_input, session_title_override, acp_session_title, viewed_through_event_ordinal, archived
1112             FROM sessions WHERE session_id=?1", [&session.id], |row| Ok((
1113                row.get::<_, String>(0)?, row.get::<_, Option<String>>(1)?, row.get::<_, Option<String>>(2)?,
1114                row.get::<_, u64>(3)?, row.get::<_, bool>(4)?,
1115            )),
1116        )?;
1117        session.draft_input = draft;
1118        session.session_title_override = title;
1119        session.acp_session_title = acp_title;
1120        session.viewed_through_event_ordinal = viewed;
1121        session.archived = archived;
1122    }
1123    insert_session(&tx, &session)?;
1124    if let Some((host, size)) = container_size {
1125        write_host_container_size(&tx, host, size)?;
1126    }
1127    tx.commit()?;
1128    Ok(())
1129}
1130
1131pub(super) fn save_lifecycle_session_to(path: &Path, session: &SessionRecord) -> Result<()> {
1132    validate_session_record(session)?;
1133
1134    let mut connection = open(path)?;
1135    let tx = connection.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
1136    update_lifecycle_fields(&tx, session)?;
1137    tx.commit()?;
1138    Ok(())
1139}
1140
1141/// Settle teardown and archive failed children in the same owner transaction.
1142pub(crate) fn save_startup_cleanup_outcome(session: &SessionRecord, archive: bool) -> Result<()> {
1143    let session = session.clone();
1144    submit_database_write("save_startup_cleanup_outcome", move |_| {
1145        validate_session_record(&session)?;
1146        let mut connection = open(&database_path())?;
1147        let tx = connection.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
1148        update_lifecycle_fields(&tx, &session)?;
1149        if archive {
1150            tx.execute(
1151                "UPDATE sessions SET archived=1 WHERE session_id=?1",
1152                [&session.id],
1153            )?;
1154        }
1155        tx.commit()?;
1156        Ok(())
1157    })
1158}
1159
1160pub(super) fn save_checkpointed_session_to(path: &Path, session: &SessionRecord) -> Result<()> {
1161    validate_session_record(session)?;
1162
1163    let mut connection = open(path)?;
1164    let tx = connection.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
1165    update_lifecycle_fields(&tx, session)?;
1166    tx.execute(
1167        "UPDATE sessions SET native_session_id = ?2 WHERE session_id = ?1",
1168        params![session.id, session.native_session_id],
1169    )?;
1170    replace_checkpoint(&tx, session)?;
1171    tx.commit()?;
1172    Ok(())
1173}
1174
1175pub(super) fn validate_session_record(session: &SessionRecord) -> Result<()> {
1176    let mut validation = State::default();
1177    validation
1178        .sessions
1179        .insert(session.id.clone(), session.clone());
1180    validation.validate()
1181}
1182
1183/// Remove one operational session while retaining its relational history
1184/// context and prompt history.
1185pub fn delete_session(session_id: &str) -> Result<()> {
1186    let session_id = session_id.to_owned();
1187    submit_database_write("delete_session", move |_| {
1188        delete_session_from(&database_path(), &session_id)
1189    })
1190}
1191
1192pub(super) fn delete_session_from(path: &Path, session_id: &str) -> Result<()> {
1193    let connection = open(path)?;
1194    connection.execute("DELETE FROM sessions WHERE session_id = ?1", [session_id])?;
1195    Ok(())
1196}
1197
1198/// Overwrite the unsent chat input carried across a detach. Unlike the read
1199/// receipt this is not monotonic: a draft can shrink, and an empty string
1200/// clears it.
1201pub fn set_session_draft_input(session_id: &str, draft: &str) -> Result<()> {
1202    let session_id = session_id.to_owned();
1203    let draft = draft.to_owned();
1204    submit_database_write("set_session_draft_input", move |_| {
1205        set_session_draft_input_at(&database_path(), &session_id, &draft)
1206    })
1207}
1208
1209pub(super) fn set_session_draft_input_at(path: &Path, session_id: &str, draft: &str) -> Result<()> {
1210    let mut connection = open(path)?;
1211    let tx = connection.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
1212    let updated = tx.execute(
1213        "UPDATE sessions SET draft_input = ?2 WHERE session_id = ?1",
1214        params![session_id, draft],
1215    )?;
1216    ensure!(updated == 1, "unknown session {session_id}");
1217    tx.commit()?;
1218    Ok(())
1219}
1220
1221/// Append recovered input on the writer lane, preserving every preceding
1222/// durable edit instead of combining with a stale daemon snapshot.
1223pub fn append_session_draft_input(session_id: &str, text: &str) -> Result<()> {
1224    let session_id = session_id.to_owned();
1225    let text = text.to_owned();
1226    submit_database_write("append_session_draft_input", move |connection| {
1227        let updated = connection.execute(
1228            "UPDATE sessions SET draft_input = CASE WHEN ?2 = '' THEN draft_input
1229             WHEN draft_input = '' THEN ?2 ELSE draft_input || char(10) || char(10) || ?2 END
1230             WHERE session_id = ?1",
1231            params![session_id, text],
1232        )?;
1233        ensure!(updated == 1, "unknown session {session_id}");
1234        Ok(())
1235    })
1236}
1237
1238/// Retire a submitted shared draft without erasing a newer client's edit.
1239pub fn clear_session_draft_input_if_matches(session_id: &str, expected: &str) -> Result<()> {
1240    let session_id = session_id.to_owned();
1241    let expected = expected.to_owned();
1242    submit_database_write("clear_session_draft_input_if_matches", move |connection| {
1243        connection.execute(
1244            "UPDATE sessions SET draft_input = '' WHERE session_id = ?1 AND draft_input = ?2",
1245            params![session_id, expected],
1246        )?;
1247        Ok(())
1248    })
1249}
1250
1251pub fn record_recovery_success(
1252    session_id: &str,
1253    native_session_id: &str,
1254    checkpoint: &CheckpointMetadata,
1255) -> Result<()> {
1256    let session_id = session_id.to_owned();
1257    let native_session_id = native_session_id.to_owned();
1258    let checkpoint = checkpoint.clone();
1259    submit_database_write("record_recovery_success", move |_| {
1260        record_recovery_success_to(
1261            &database_path(),
1262            &session_id,
1263            &native_session_id,
1264            &checkpoint,
1265        )
1266    })
1267}
1268
1269pub(super) fn record_recovery_success_to(
1270    path: &Path,
1271    session_id: &str,
1272    native_session_id: &str,
1273    checkpoint: &CheckpointMetadata,
1274) -> Result<()> {
1275    let mut connection = open(path)?;
1276    let tx = connection.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
1277    let changed = tx.execute(
1278        "UPDATE sessions
1279         SET native_session_id = ?2, last_checkpoint_error = NULL
1280         WHERE session_id = ?1",
1281        params![session_id, native_session_id],
1282    )?;
1283    if changed != 1 {
1284        bail!("unknown session {session_id}");
1285    }
1286    tx.execute(
1287        "INSERT INTO session_checkpoints(
1288             session_id, archive_path, sha256, created_at, event_frontier
1289         ) VALUES (?1,?2,?3,?4,?5)
1290         ON CONFLICT(session_id) DO UPDATE SET
1291             archive_path = excluded.archive_path,
1292             sha256 = excluded.sha256,
1293             created_at = excluded.created_at,
1294             event_frontier = excluded.event_frontier",
1295        params![
1296            session_id,
1297            path_to_blob(&checkpoint.archive_path),
1298            checkpoint.sha256,
1299            checkpoint.created_at,
1300            checkpoint.event_frontier,
1301        ],
1302    )?;
1303    tx.commit()?;
1304    Ok(())
1305}
1306
1307pub fn record_recovery_success_if_current(
1308    session_id: &str,
1309    expected_target: &TargetLocator,
1310    expected_checkpoint: Option<&CheckpointMetadata>,
1311    native_session_id: &str,
1312    checkpoint: &CheckpointMetadata,
1313) -> Result<bool> {
1314    let session_id = session_id.to_owned();
1315    let expected_target = expected_target.clone();
1316    let expected_checkpoint = expected_checkpoint.cloned();
1317    let native_session_id = native_session_id.to_owned();
1318    let checkpoint = checkpoint.clone();
1319    submit_database_write("record_recovery_success_if_current", move |connection| {
1320        record_recovery_success_if_current_with(
1321            connection,
1322            &session_id,
1323            &expected_target,
1324            expected_checkpoint.as_ref(),
1325            &native_session_id,
1326            &checkpoint,
1327        )
1328    })
1329}
1330
1331fn record_recovery_success_if_current_with(
1332    connection: &mut Connection,
1333    session_id: &str,
1334    expected_target: &TargetLocator,
1335    expected_checkpoint: Option<&CheckpointMetadata>,
1336    native_session_id: &str,
1337    checkpoint: &CheckpointMetadata,
1338) -> Result<bool> {
1339    let tx = connection.transaction()?;
1340    let Some(mut current) = load_session_with(&tx, session_id)? else {
1341        return Ok(false);
1342    };
1343    if current.target.as_ref() != Some(expected_target)
1344        || current.checkpoint.as_ref() != expected_checkpoint
1345    {
1346        return Ok(false);
1347    }
1348    tx.execute(
1349        "UPDATE sessions SET native_session_id = ?2, last_checkpoint_error = NULL
1350         WHERE session_id = ?1",
1351        params![session_id, native_session_id],
1352    )?;
1353    current.checkpoint = Some(checkpoint.clone());
1354    replace_checkpoint(&tx, &current)?;
1355    tx.commit()?;
1356    Ok(true)
1357}
1358
1359pub fn record_recovery_failure(session_id: &str, detail: &str) -> Result<()> {
1360    let session_id = session_id.to_owned();
1361    let detail = detail.to_owned();
1362    submit_database_write("record_recovery_failure", move |_| {
1363        record_recovery_failure_to(&database_path(), &session_id, &detail)
1364    })
1365}
1366
1367/// A completed copy can report only against the target and checkpoint it
1368/// observed. A later successful checkpoint or target replacement wins.
1369pub fn record_recovery_failure_if_current(
1370    session_id: &str,
1371    expected_target: &TargetLocator,
1372    expected_checkpoint: Option<&CheckpointMetadata>,
1373    detail: &str,
1374) -> Result<bool> {
1375    let session_id = session_id.to_owned();
1376    let expected_target = expected_target.clone();
1377    let expected_checkpoint = expected_checkpoint.cloned();
1378    let detail = detail.to_owned();
1379    submit_database_write("record_recovery_failure_if_current", move |connection| {
1380        record_recovery_failure_if_current_with(
1381            connection,
1382            &session_id,
1383            &expected_target,
1384            expected_checkpoint.as_ref(),
1385            &detail,
1386        )
1387    })
1388}
1389
1390fn record_recovery_failure_if_current_with(
1391    connection: &mut Connection,
1392    session_id: &str,
1393    expected_target: &TargetLocator,
1394    expected_checkpoint: Option<&CheckpointMetadata>,
1395    detail: &str,
1396) -> Result<bool> {
1397    let tx = connection.transaction()?;
1398    let Some(current) = load_session_with(&tx, session_id)? else {
1399        return Ok(false);
1400    };
1401    if current.target.as_ref() != Some(expected_target)
1402        || current.checkpoint.as_ref() != expected_checkpoint
1403    {
1404        return Ok(false);
1405    }
1406    tx.execute(
1407        "UPDATE sessions SET last_checkpoint_error = ?2 WHERE session_id = ?1",
1408        params![session_id, detail],
1409    )?;
1410    tx.commit()?;
1411    Ok(true)
1412}
1413
1414pub(super) fn record_recovery_failure_to(
1415    path: &Path,
1416    session_id: &str,
1417    detail: &str,
1418) -> Result<()> {
1419    let connection = open(path)?;
1420    let changed = connection.execute(
1421        "UPDATE sessions SET last_checkpoint_error = ?2 WHERE session_id = ?1",
1422        params![session_id, detail],
1423    )?;
1424    if changed != 1 {
1425        bail!("unknown session {session_id}");
1426    }
1427    Ok(())
1428}
1429
1430/// Re-associate a session with another project bundle.
1431///
1432/// A session's bundle is otherwise fixed, because prompt history is grouped by
1433/// it. Resume calls this when it converts a session between its raw and bundle
1434/// representations: the project is the same, so its history follows it, and
1435/// only the name Hel files it under changes.
1436/// Files the session under `bundle_id`'s project and returns the bundle id the
1437/// session context now holds: the canonical project when the catalog knows
1438/// `bundle_id` as an alias of another bundle. A caller keeping the session
1439/// record must adopt the returned id, because resume publication checks the
1440/// record's bundle against the context.
1441pub fn rebind_session_bundle(session_id: &str, bundle_id: &str) -> Result<String> {
1442    let session_id = session_id.to_owned();
1443    let bundle_id = bundle_id.to_owned();
1444    submit_database_write("rebind_session_bundle", move |_| {
1445        rebind_session_bundle_to(&database_path(), &session_id, &bundle_id)
1446    })
1447}
1448
1449pub(super) fn rebind_session_bundle_to(
1450    path: &Path,
1451    session_id: &str,
1452    bundle_id: &str,
1453) -> Result<String> {
1454    let mut connection = open(path)?;
1455    let tx = connection.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
1456    let canonical_bundle = super::projects::canonical_project_id(&tx, bundle_id)?;
1457    tx.execute("INSERT OR IGNORE INTO project_session_aliases(session_id,bundle_id) SELECT session_id,bundle_id FROM session_contexts WHERE session_id=?1",[session_id])?;
1458    let changed = tx.execute(
1459        "UPDATE session_contexts SET bundle_id = ?2 WHERE session_id = ?1",
1460        params![session_id, canonical_bundle],
1461    )?;
1462    if changed == 0 {
1463        tx.execute(
1464            "INSERT INTO session_contexts(session_id, bundle_id, created_at) VALUES (?1, ?2, ?3)",
1465            params![session_id, canonical_bundle, Utc::now().to_rfc3339()],
1466        )?;
1467    }
1468    tx.commit()?;
1469    Ok(canonical_bundle)
1470}
1471
1472#[cfg(test)]
1473mod ownership_tests {
1474    use super::*;
1475
1476    #[test]
1477    fn resume_and_rollback_preserve_edits_committed_after_their_snapshot() {
1478        let directory = tempfile::tempdir().unwrap();
1479        let path = directory.path().join("resume.sqlite3");
1480        let original = super::super::tests::session("resume-owner", "project");
1481        save_session_to(&path, &original).unwrap();
1482        let connection = open(&path).unwrap();
1483        connection
1484            .execute(
1485                "UPDATE sessions SET session_title_override = 'renamed during resume',
1486                 acp_session_title = 'new worker title', draft_input = 'new draft',
1487                 archived = 1, viewed_through_event_ordinal = 99 WHERE session_id = ?1",
1488                [&original.id],
1489            )
1490            .unwrap();
1491        let mut provisioning = original.clone();
1492        provisioning.state = SessionState::Provisioning;
1493        provisioning.target = None;
1494        provisioning.native_session_id = Some("resumed-native".into());
1495        provisioning.additional_mounts.clear();
1496        save_resumed_session_to(&path, &provisioning, None).unwrap();
1497        let current = load_session_with(&connection, &original.id)
1498            .unwrap()
1499            .unwrap();
1500        assert_eq!(current.state, SessionState::Provisioning);
1501        assert_eq!(current.native_session_id, provisioning.native_session_id);
1502        assert!(current.additional_mounts.is_empty());
1503        assert_eq!(
1504            current.session_title_override.as_deref(),
1505            Some("renamed during resume")
1506        );
1507        assert_eq!(
1508            current.acp_session_title.as_deref(),
1509            Some("new worker title")
1510        );
1511        assert_eq!(current.draft_input, "new draft");
1512        assert!(current.archived);
1513        assert_eq!(current.viewed_through_event_ordinal, 99);
1514
1515        // Failure restores a snapshot captured before those independent edits.
1516        save_resumed_session_to(&path, &original, None).unwrap();
1517        let rolled_back = load_session_with(&connection, &original.id)
1518            .unwrap()
1519            .unwrap();
1520        assert_eq!(rolled_back.state, original.state);
1521        assert_eq!(rolled_back.additional_mounts, original.additional_mounts);
1522        assert_eq!(
1523            rolled_back.session_title_override,
1524            current.session_title_override
1525        );
1526        assert_eq!(rolled_back.acp_session_title, current.acp_session_title);
1527        assert_eq!(rolled_back.draft_input, current.draft_input);
1528        assert_eq!(rolled_back.archived, current.archived);
1529        assert_eq!(rolled_back.viewed_through_event_ordinal, 99);
1530    }
1531
1532    #[test]
1533    fn recovery_completion_cannot_replace_a_newer_checkpoint_or_target() {
1534        let directory = tempfile::tempdir().unwrap();
1535        let path = directory.path().join("recovery.sqlite3");
1536        let original = super::super::tests::session("recovery-owner", "project");
1537        save_session_to(&path, &original).unwrap();
1538        let target = original.target.as_ref().unwrap();
1539        let previous = original.checkpoint.as_ref().unwrap();
1540        let mut next = previous.clone();
1541        next.sha256 = "b".repeat(64);
1542        next.event_frontier += 1;
1543        let mut connection = open(&path).unwrap();
1544        assert!(
1545            record_recovery_success_if_current_with(
1546                &mut connection,
1547                &original.id,
1548                target,
1549                Some(previous),
1550                "new-native",
1551                &next,
1552            )
1553            .unwrap()
1554        );
1555        assert!(
1556            !record_recovery_failure_if_current_with(
1557                &mut connection,
1558                &original.id,
1559                target,
1560                Some(previous),
1561                "old failure",
1562            )
1563            .unwrap()
1564        );
1565        assert!(
1566            !record_recovery_success_if_current_with(
1567                &mut connection,
1568                &original.id,
1569                target,
1570                Some(previous),
1571                "old-native",
1572                previous,
1573            )
1574            .unwrap()
1575        );
1576        let current = load_session_with(&connection, &original.id)
1577            .unwrap()
1578            .unwrap();
1579        assert_eq!(current.checkpoint.as_ref(), Some(&next));
1580        assert_eq!(current.native_session_id.as_deref(), Some("new-native"));
1581        assert_eq!(current.last_checkpoint_error, None);
1582
1583        let mut replacement = current.clone();
1584        replacement.target = None;
1585        replacement.state = SessionState::Stopped;
1586        save_resumed_session_to(&path, &replacement, None).unwrap();
1587        assert!(
1588            !record_recovery_failure_if_current_with(
1589                &mut connection,
1590                &original.id,
1591                target,
1592                Some(&next),
1593                "old target failure",
1594            )
1595            .unwrap()
1596        );
1597        assert!(
1598            !record_recovery_success_if_current_with(
1599                &mut connection,
1600                &original.id,
1601                target,
1602                Some(&next),
1603                "old-native",
1604                previous,
1605            )
1606            .unwrap()
1607        );
1608    }
1609
1610    #[test]
1611    fn current_recovery_failure_settles_without_resurrecting_a_removed_session() {
1612        let directory = tempfile::tempdir().unwrap();
1613        let path = directory.path().join("recovery.sqlite3");
1614        let original = super::super::tests::session("recovery-owner", "project");
1615        save_session_to(&path, &original).unwrap();
1616        let mut connection = open(&path).unwrap();
1617        assert!(
1618            record_recovery_failure_if_current_with(
1619                &mut connection,
1620                &original.id,
1621                original.target.as_ref().unwrap(),
1622                original.checkpoint.as_ref(),
1623                "current failure",
1624            )
1625            .unwrap()
1626        );
1627        assert_eq!(
1628            load_session_with(&connection, &original.id)
1629                .unwrap()
1630                .unwrap()
1631                .last_checkpoint_error
1632                .as_deref(),
1633            Some("current failure")
1634        );
1635        connection
1636            .execute("DELETE FROM sessions WHERE session_id = ?1", [&original.id])
1637            .unwrap();
1638        assert!(
1639            !record_recovery_failure_if_current_with(
1640                &mut connection,
1641                &original.id,
1642                original.target.as_ref().unwrap(),
1643                original.checkpoint.as_ref(),
1644                "late failure",
1645            )
1646            .unwrap()
1647        );
1648        assert!(
1649            load_session_with(&connection, &original.id)
1650                .unwrap()
1651                .is_none()
1652        );
1653        assert!(save_resumed_session_to(&path, &original, None).is_err());
1654    }
1655}