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
13/// Update publication evidence only while the exact stopped checkpoint is
14/// still current. An age scan may race a user's Resume or a new checkpoint.
15pub fn set_publication_assessment_if_current(
16    session_id: &str,
17    assessment: &mj_core::state::PublicationAssessment,
18) -> Result<bool> {
19    let session_id = session_id.to_owned();
20    let assessment = assessment.clone();
21    submit_database_write("set_publication_assessment_if_current", move |_| {
22        let connection = open(&database_path())?;
23        let updated = connection.execute(
24            "UPDATE sessions SET publication_json = ?2
25             WHERE session_id = ?1 AND state = 'stopped'
26               AND EXISTS (SELECT 1 FROM session_checkpoints c
27                           WHERE c.session_id = ?1 AND c.sha256 = ?3)",
28            params![
29                session_id,
30                serde_json::to_string(&assessment)?,
31                assessment.checkpoint_sha256
32            ],
33        )?;
34        Ok(updated == 1)
35    })
36}
37
38/// Persist a borrowed-target child and its parent relationship atomically.
39pub fn save_subagent_session(
40    session: &SessionRecord,
41    subagent: &mj_core::subagent::SubagentRecord,
42) -> Result<()> {
43    let session = session.clone();
44    let subagent = subagent.clone();
45    submit_database_write("save_subagent_session", move |_| {
46        let mut connection = open(&database_path())?;
47        let tx = connection.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
48        insert_session(&tx, &session)?;
49        tx.execute(
50            "INSERT INTO subagent_sessions(
51                 child_session_id, parent_session_id, request_key, record_json
52             ) VALUES (?1, ?2, ?3, ?4)",
53            params![
54                subagent.child_session_id,
55                subagent.parent_session_id,
56                subagent.request_key,
57                serde_json::to_string(&subagent)?,
58            ],
59        )?;
60        tx.commit()?;
61        Ok(())
62    })
63}
64
65/// Record the child turn whose completion notice the parent already has.
66pub fn mark_subagent_turn_noticed(child_session_id: &str, turn: u64) -> Result<()> {
67    let child_session_id = child_session_id.to_owned();
68    submit_database_write("mark_subagent_turn_noticed", move |_| {
69        let mut relation = load_subagent(&child_session_id)?
70            .with_context(|| format!("unknown sub-agent session {child_session_id}"))?;
71        relation.noticed_turn = Some(turn);
72        let json = serde_json::to_string(&relation)?;
73        let connection = open(&database_path())?;
74        connection.execute(
75            "UPDATE subagent_sessions SET record_json = ?2 WHERE child_session_id = ?1",
76            params![child_session_id, json],
77        )?;
78        Ok(())
79    })
80}
81
82pub fn load_subagent(child_session_id: &str) -> Result<Option<mj_core::subagent::SubagentRecord>> {
83    let connection = open_reader(&database_path())?;
84    connection
85        .query_row(
86            "SELECT record_json FROM subagent_sessions WHERE child_session_id = ?1",
87            [child_session_id],
88            |row| row.get::<_, String>(0),
89        )
90        .optional()?
91        .map(|json| serde_json::from_str(&json).context("decode sub-agent record"))
92        .transpose()
93}
94
95pub fn list_subagents(parent_session_id: &str) -> Result<Vec<mj_core::subagent::SubagentRecord>> {
96    let connection = open_reader(&database_path())?;
97    let mut statement = connection.prepare(
98        "SELECT record_json FROM subagent_sessions
99         WHERE parent_session_id = ?1 ORDER BY rowid",
100    )?;
101    statement
102        .query_map([parent_session_id], |row| row.get::<_, String>(0))?
103        .map(|row| serde_json::from_str(&row?).context("decode sub-agent record"))
104        .collect()
105}
106
107/// Everything recorded about one child's report. A child with nothing
108/// recorded yet reads as the empty report.
109pub fn load_subagent_report(child_session_id: &str) -> Result<mj_core::subagent::SubagentReport> {
110    load_subagent_report_from(&database_path(), child_session_id)
111}
112
113pub(super) fn load_subagent_report_from(
114    path: &Path,
115    child_session_id: &str,
116) -> Result<mj_core::subagent::SubagentReport> {
117    let connection = open_reader(path)?;
118    let row = connection
119        .query_row(
120            "SELECT handback_command_id, handback_message, handback_recorded_at_ms,
121                    reminder_command_id, reminder_for_command_id, reminder_sent_at_ms,
122                    reminder_failed_for_command_id, awaited_ordinal
123             FROM subagent_handbacks WHERE child_session_id = ?1",
124            [child_session_id],
125            |row| {
126                Ok((
127                    row.get::<_, Option<String>>(0)?,
128                    row.get::<_, Option<String>>(1)?,
129                    row.get::<_, Option<i64>>(2)?,
130                    row.get::<_, Option<String>>(3)?,
131                    row.get::<_, Option<String>>(4)?,
132                    row.get::<_, Option<i64>>(5)?,
133                    row.get::<_, Option<String>>(6)?,
134                    row.get::<_, Option<i64>>(7)?,
135                ))
136            },
137        )
138        .optional()?;
139    let Some((
140        handback_command,
141        handback_message,
142        handback_at,
143        reminder_command,
144        reminder_for,
145        reminder_at,
146        reminder_failed_for,
147        awaited_ordinal,
148    )) = row
149    else {
150        return Ok(mj_core::subagent::SubagentReport::default());
151    };
152    Ok(mj_core::subagent::SubagentReport {
153        handback: match (handback_command, handback_message, handback_at) {
154            (Some(command_id), Some(message), Some(recorded_at_ms)) => {
155                Some(mj_core::subagent::SubagentHandback {
156                    command_id,
157                    message,
158                    recorded_at_ms,
159                })
160            }
161            _ => None,
162        },
163        reminder: match (reminder_command, reminder_for, reminder_at) {
164            (Some(command_id), Some(for_command_id), Some(sent_at_ms)) => {
165                Some(mj_core::subagent::HandbackReminder {
166                    command_id,
167                    for_command_id,
168                    sent_at_ms,
169                })
170            }
171            _ => None,
172        },
173        reminder_failed_for,
174        awaited_ordinal: awaited_ordinal.and_then(|ordinal| u64::try_from(ordinal).ok()),
175    })
176}
177
178/// Record that the parent gave a child a prompt, accepted at `ordinal`. It is
179/// recorded only for a sub-agent child, and never moves backwards.
180pub fn record_subagent_prompt(child_session_id: &str, ordinal: u64) -> Result<()> {
181    let child_session_id = child_session_id.to_owned();
182    submit_database_write("record_subagent_prompt", move |_| {
183        record_subagent_prompt_to(&database_path(), &child_session_id, ordinal)
184    })
185}
186
187pub(super) fn record_subagent_prompt_to(
188    path: &Path,
189    child_session_id: &str,
190    ordinal: u64,
191) -> Result<()> {
192    let ordinal = i64::try_from(ordinal).context("prompt ordinal exceeds the store's range")?;
193    open(path)?.execute(
194        "INSERT INTO subagent_handbacks(child_session_id, awaited_ordinal)
195         SELECT ?1, ?2 WHERE EXISTS (
196             SELECT 1 FROM subagent_sessions WHERE child_session_id = ?1
197         )
198         ON CONFLICT(child_session_id) DO UPDATE SET
199             awaited_ordinal = max(coalesce(awaited_ordinal, 0), excluded.awaited_ordinal)",
200        params![child_session_id, ordinal],
201    )?;
202    Ok(())
203}
204
205/// Record a child's report for the turn it names, unless a report for that
206/// turn is already recorded. Returns whether this one was recorded: a child
207/// delivers one report per turn, and the check and the write are one
208/// transaction so two racing calls cannot both be accepted.
209pub fn record_subagent_handback(
210    child_session_id: &str,
211    handback: &mj_core::subagent::SubagentHandback,
212) -> Result<bool> {
213    let child_session_id = child_session_id.to_owned();
214    let handback = handback.clone();
215    submit_database_write("record_subagent_handback", move |_| {
216        record_subagent_handback_to(&database_path(), &child_session_id, &handback)
217    })
218}
219
220pub(super) fn record_subagent_handback_to(
221    path: &Path,
222    child_session_id: &str,
223    handback: &mj_core::subagent::SubagentHandback,
224) -> Result<bool> {
225    let mut connection = open(path)?;
226    let tx = connection.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
227    let recorded_for: Option<String> = tx
228        .query_row(
229            "SELECT handback_command_id FROM subagent_handbacks WHERE child_session_id = ?1",
230            [child_session_id],
231            |row| row.get(0),
232        )
233        .optional()?
234        .flatten();
235    if recorded_for.as_deref() == Some(handback.command_id.as_str()) {
236        return Ok(false);
237    }
238    tx.execute(
239        "INSERT INTO subagent_handbacks(
240             child_session_id, handback_command_id, handback_message, handback_recorded_at_ms
241         ) VALUES (?1, ?2, ?3, ?4)
242         ON CONFLICT(child_session_id) DO UPDATE SET
243             handback_command_id = excluded.handback_command_id,
244             handback_message = excluded.handback_message,
245             handback_recorded_at_ms = excluded.handback_recorded_at_ms",
246        params![
247            child_session_id,
248            handback.command_id,
249            handback.message,
250            handback.recorded_at_ms
251        ],
252    )?;
253    tx.commit()?;
254    Ok(true)
255}
256
257/// Record the reminder sent after a child's turn ended without a report.
258pub fn record_handback_reminder(
259    child_session_id: &str,
260    reminder: &mj_core::subagent::HandbackReminder,
261) -> Result<()> {
262    let child_session_id = child_session_id.to_owned();
263    let reminder = reminder.clone();
264    submit_database_write("record_handback_reminder", move |_| {
265        open(&database_path())?.execute(
266            "INSERT INTO subagent_handbacks(
267                 child_session_id, reminder_command_id, reminder_for_command_id, reminder_sent_at_ms
268             ) VALUES (?1, ?2, ?3, ?4)
269             ON CONFLICT(child_session_id) DO UPDATE SET
270                 reminder_command_id = excluded.reminder_command_id,
271                 reminder_for_command_id = excluded.reminder_for_command_id,
272                 reminder_sent_at_ms = excluded.reminder_sent_at_ms",
273            params![
274                child_session_id,
275                reminder.command_id,
276                reminder.for_command_id,
277                reminder.sent_at_ms
278            ],
279        )?;
280        Ok(())
281    })
282}
283
284/// Record that the reminder for a turn could not be sent, so the child's
285/// last message stands as its report.
286pub fn record_handback_reminder_failed(child_session_id: &str, for_command_id: &str) -> Result<()> {
287    let child_session_id = child_session_id.to_owned();
288    let for_command_id = for_command_id.to_owned();
289    submit_database_write("record_handback_reminder_failed", move |_| {
290        open(&database_path())?.execute(
291            "INSERT INTO subagent_handbacks(child_session_id, reminder_failed_for_command_id)
292             VALUES (?1, ?2)
293             ON CONFLICT(child_session_id) DO UPDATE SET
294                 reminder_failed_for_command_id = excluded.reminder_failed_for_command_id",
295            params![child_session_id, for_command_id],
296        )?;
297        Ok(())
298    })
299}
300
301pub fn lookup_subagent_request(
302    parent_session_id: &str,
303    request_key: &str,
304) -> Result<Option<mj_core::subagent::SubagentRecord>> {
305    let connection = open_reader(&database_path())?;
306    connection
307        .query_row(
308            "SELECT record_json FROM subagent_sessions
309             WHERE parent_session_id = ?1 AND request_key = ?2",
310            params![parent_session_id, request_key],
311            |row| row.get::<_, String>(0),
312        )
313        .optional()?
314        .map(|json| serde_json::from_str(&json).context("decode sub-agent record"))
315        .transpose()
316}
317
318/// Persist a session and the container size it most recently launched on its
319/// host in one transaction.
320pub fn save_session_with_container_size(
321    session: &SessionRecord,
322    host: &str,
323    size: HostContainerSize,
324) -> Result<()> {
325    let session = session.clone();
326    let host = host.to_owned();
327    submit_database_write("save_session_with_container_size", move |_| {
328        save_session_with_container_size_to(&database_path(), &session, Some((&host, size)))
329    })
330}
331
332/// Update only the fields a lifecycle transition owns on a session that
333/// already exists. Everything else — display titles, checkpoints, container
334/// settings, and attached directories — stays with its own writer.
335pub fn save_lifecycle_session(session: &SessionRecord) -> Result<()> {
336    let session = session.clone();
337    submit_database_write("save_lifecycle_session", move |_| {
338        save_lifecycle_session_to(&database_path(), &session)
339    })
340}
341
342/// Install a lifecycle transition together with the checkpoint it just
343/// verified and the harness session id that produced it.
344pub fn save_checkpointed_session(session: &SessionRecord) -> Result<()> {
345    let session = session.clone();
346    submit_database_write("save_checkpointed_session", move |_| {
347        save_checkpointed_session_to(&database_path(), &session)
348    })
349}
350
351/// Recover lifecycle rows stranded by a process exit during checkpoint
352/// creation. This must be called once by the top-level controller process
353/// while it owns the controller-store guard, not by per-operation reloads.
354pub fn recover_interrupted_checkpointing_sessions(updated_at: &str) -> Result<usize> {
355    let updated_at = updated_at.to_owned();
356    submit_database_write("recover_interrupted_checkpointing_sessions", move |_| {
357        recover_interrupted_checkpointing_sessions_to(&database_path(), &updated_at)
358    })
359}
360
361/// Change only the user-owned display name. This avoids writing a stale
362/// SessionRecord over independently committed checkpoint or relay metadata.
363pub fn set_session_title_override(session_id: &str, title: &str, updated_at: &str) -> Result<()> {
364    let session_id = session_id.to_owned();
365    let title = title.to_owned();
366    let updated_at = updated_at.to_owned();
367    submit_database_write("set_session_title_override", move |_| {
368        set_session_title_override_to(&database_path(), &session_id, &title, &updated_at)
369    })
370}
371
372/// Rewrite a configured profile id in every persisted session in one SQLite
373/// transaction. Configuration is stored separately, so the controller owns
374/// coordinating this update with the matching config-map rename.
375pub fn rename_profile_references(old_id: &str, new_id: &str) -> Result<usize> {
376    rename_session_reference("last_profile", old_id, new_id)
377}
378
379/// Rewrite a configured target id in every persisted session in one SQLite
380/// transaction.
381pub fn rename_target_references(old_id: &str, new_id: &str) -> Result<usize> {
382    rename_session_reference("target_template_id", old_id, new_id)
383}
384
385pub(super) fn rename_session_reference(
386    column: &'static str,
387    old_id: &str,
388    new_id: &str,
389) -> Result<usize> {
390    ensure!(
391        matches!(column, "last_profile" | "target_template_id"),
392        "unsupported session reference column"
393    );
394    let old_id = old_id.to_owned();
395    let new_id = new_id.to_owned();
396    submit_database_write("rename_session_reference", move |_| {
397        rename_session_reference_at(&database_path(), column, &old_id, &new_id)
398    })
399}
400
401pub(super) fn rename_session_reference_at(
402    path: &Path,
403    column: &str,
404    old_id: &str,
405    new_id: &str,
406) -> Result<usize> {
407    let mut connection = open(path)?;
408    let tx = connection.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
409    let changed = tx.execute(
410        &format!("UPDATE sessions SET {column} = ?2 WHERE {column} = ?1"),
411        params![old_id, new_id],
412    )?;
413    tx.commit()?;
414    Ok(changed)
415}
416
417/// Change only whether the resume dialog hides this session. Archiving is a
418/// display choice, so it has its own writer and never rewrites lifecycle,
419/// checkpoint, or title columns another task owns.
420pub fn set_session_archived(session_id: &str, archived: bool) -> Result<()> {
421    let session_id = session_id.to_owned();
422    submit_database_write("set_session_archived", move |_| {
423        set_session_archived_to(&database_path(), &session_id, archived)
424    })
425}
426
427/// Record that the managed target of an otherwise live session is definitively
428/// gone. A verified checkpoint keeps the session recoverable as an error on the
429/// dashboard; without one, the session is lost. The state predicate keeps a
430/// late poll result from overwriting a concurrent lifecycle transition.
431pub fn mark_session_target_missing(
432    session_id: &str,
433    detail: &str,
434    updated_at: &str,
435) -> Result<Option<SessionState>> {
436    let session_id = session_id.to_owned();
437    let detail = detail.to_owned();
438    let updated_at = updated_at.to_owned();
439    submit_database_write("mark_session_target_missing", move |_| {
440        mark_session_target_missing_to(&database_path(), &session_id, &detail, &updated_at)
441    })
442}
443
444pub(super) fn mark_session_target_missing_to(
445    path: &Path,
446    session_id: &str,
447    detail: &str,
448    updated_at: &str,
449) -> Result<Option<SessionState>> {
450    mark_session_target_missing_if_current_to(path, session_id, detail, updated_at, None)
451}
452
453/// Record a definitive worker failure only while the observed session record
454/// is still current. A delayed background write must not invalidate a resume.
455pub fn mark_session_target_missing_if_current(
456    session_id: &str,
457    detail: &str,
458    updated_at: &str,
459    observed_updated_at: &str,
460) -> Result<Option<SessionState>> {
461    let session_id = session_id.to_owned();
462    let detail = detail.to_owned();
463    let updated_at = updated_at.to_owned();
464    let observed_updated_at = observed_updated_at.to_owned();
465    submit_database_write("mark_session_target_missing_if_current", move |_| {
466        mark_session_target_missing_if_current_to(
467            &database_path(),
468            &session_id,
469            &detail,
470            &updated_at,
471            Some(&observed_updated_at),
472        )
473    })
474}
475
476pub(super) fn mark_session_target_missing_if_current_to(
477    path: &Path,
478    session_id: &str,
479    detail: &str,
480    updated_at: &str,
481    observed_updated_at: Option<&str>,
482) -> Result<Option<SessionState>> {
483    let mut connection = open(path)?;
484    let tx = connection.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
485    let changed = tx.execute(
486        "UPDATE sessions
487         SET state = CASE
488                 WHEN EXISTS(
489                     SELECT 1 FROM session_checkpoints
490                     WHERE session_checkpoints.session_id = sessions.session_id
491                 ) THEN 'error'
492                 ELSE 'lost'
493             END,
494             last_error = ?2,
495             updated_at = ?3
496         WHERE session_id = ?1
497           AND (?4 IS NULL OR updated_at = ?4)
498           AND state IN ('provisioning', 'running', 'disconnected', 'error')",
499        params![session_id, detail, updated_at, observed_updated_at],
500    )?;
501    ensure!(changed <= 1, "updated {changed} sessions for {session_id}");
502    let state = if changed == 1 {
503        let stored: String = tx.query_row(
504            "SELECT state FROM sessions WHERE session_id = ?1",
505            [session_id],
506            |row| row.get(0),
507        )?;
508        Some(stored_session_state(&stored))
509    } else {
510        None
511    };
512    tx.commit()?;
513    Ok(state)
514}
515
516pub(super) fn set_session_archived_to(path: &Path, session_id: &str, archived: bool) -> Result<()> {
517    let connection = open(path)?;
518    let changed = connection.execute(
519        "UPDATE sessions SET archived = ?2 WHERE session_id = ?1",
520        params![session_id, archived],
521    )?;
522    if changed != 1 {
523        bail!("unknown session {session_id}");
524    }
525    Ok(())
526}
527
528/// Native sessions the resume dialog hides. Hel never writes into a harness
529/// home, so the hidden set lives here instead of in the harness's own store.
530pub fn hidden_native_sessions() -> Result<BTreeSet<(mj_core::config::HarnessKind, String)>> {
531    hidden_native_sessions_from(&database_path())
532}
533
534pub(super) fn hidden_native_sessions_from(
535    path: &Path,
536) -> Result<BTreeSet<(mj_core::config::HarnessKind, String)>> {
537    let connection = open_reader(path)?;
538    let mut statement =
539        connection.prepare("SELECT harness_kind, native_session_id FROM hidden_native_sessions")?;
540    let rows = statement.query_map([], |row| {
541        Ok((row.get::<_, String>(0)?, row.get::<_, String>(1)?))
542    })?;
543    let mut hidden = BTreeSet::new();
544    for row in rows {
545        let (harness, native_session_id) = row?;
546        // Rows for a harness this release no longer supports are ignored, not
547        // fatal; they simply hide nothing.
548        match harness.parse::<mj_core::config::HarnessKind>() {
549            Ok(harness) => {
550                hidden.insert((harness, native_session_id));
551            }
552            Err(_) => tracing::warn!(
553                harness = %harness,
554                "ignoring a hidden native session for a harness that is no longer supported"
555            ),
556        }
557    }
558    Ok(hidden)
559}
560
561/// Hide or reveal one native session in the resume dialog.
562pub fn set_native_session_hidden(
563    harness: mj_core::config::HarnessKind,
564    native_session_id: &str,
565    hidden: bool,
566) -> Result<()> {
567    let native_session_id = native_session_id.to_owned();
568    submit_database_write("set_native_session_hidden", move |_| {
569        set_native_session_hidden_to(&database_path(), harness, &native_session_id, hidden)
570    })
571}
572
573pub(super) fn set_native_session_hidden_to(
574    path: &Path,
575    harness: mj_core::config::HarnessKind,
576    native_session_id: &str,
577    hidden: bool,
578) -> Result<()> {
579    if native_session_id.trim().is_empty() {
580        bail!("native session id is empty");
581    }
582    let connection = open(path)?;
583    if hidden {
584        connection.execute(
585            "INSERT INTO hidden_native_sessions(harness_kind, native_session_id, hidden_at)
586             VALUES (?1, ?2, ?3)
587             ON CONFLICT(harness_kind, native_session_id) DO NOTHING",
588            params![harness.id(), native_session_id, Utc::now().to_rfc3339()],
589        )?;
590    } else {
591        connection.execute(
592            "DELETE FROM hidden_native_sessions
593             WHERE harness_kind = ?1 AND native_session_id = ?2",
594            params![harness.id(), native_session_id],
595        )?;
596    }
597    Ok(())
598}
599
600/// Change only the per-session container provisioning inputs: the size
601/// overrides and the attached directories. Everything else the session row
602/// owns is left to its own writer.
603pub fn set_session_container_settings(
604    session_id: &str,
605    cpus: Option<&str>,
606    memory: Option<&str>,
607    mounts: &[AdditionalMount],
608    updated_at: &str,
609) -> Result<()> {
610    let session_id = session_id.to_owned();
611    let cpus = cpus.map(str::to_owned);
612    let memory = memory.map(str::to_owned);
613    let mounts = mounts.to_vec();
614    let updated_at = updated_at.to_owned();
615    submit_database_write("set_session_container_settings", move |_| {
616        set_session_container_settings_to(
617            &database_path(),
618            &session_id,
619            cpus.as_deref(),
620            memory.as_deref(),
621            &mounts,
622            &updated_at,
623        )
624    })
625}
626
627pub(super) fn set_session_container_settings_to(
628    path: &Path,
629    session_id: &str,
630    cpus: Option<&str>,
631    memory: Option<&str>,
632    mounts: &[AdditionalMount],
633    updated_at: &str,
634) -> Result<()> {
635    if updated_at.trim().is_empty() {
636        bail!("session update timestamp is empty");
637    }
638    crate::targets::validate_additional_mounts(mounts)?;
639    let mut connection = open(path)?;
640    let tx = connection.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
641    let changed = tx.execute(
642        "UPDATE sessions
643         SET container_cpus = ?2, container_memory = ?3, updated_at = ?4
644         WHERE session_id = ?1",
645        params![session_id, cpus, memory, updated_at],
646    )?;
647    if changed != 1 {
648        bail!("unknown session {session_id}");
649    }
650    replace_mounts(&tx, session_id, mounts)?;
651    tx.commit()?;
652    Ok(())
653}
654
655pub(super) fn set_session_title_override_to(
656    path: &Path,
657    session_id: &str,
658    title: &str,
659    updated_at: &str,
660) -> Result<()> {
661    if title.trim().is_empty() {
662        bail!("session title is empty");
663    }
664    if updated_at.trim().is_empty() {
665        bail!("session update timestamp is empty");
666    }
667    let connection = open(path)?;
668    let changed = connection.execute(
669        "UPDATE sessions
670         SET session_title_override = ?2, updated_at = ?3
671         WHERE session_id = ?1",
672        params![session_id, title, updated_at],
673    )?;
674    if changed != 1 {
675        bail!("unknown session {session_id}");
676    }
677    Ok(())
678}
679
680/// Persist the latest ACP-provided title without replacing unrelated session
681/// fields that may have changed in another supervised controller task.
682pub fn set_session_acp_title(session_id: &str, title: Option<&str>) -> Result<()> {
683    let session_id = session_id.to_owned();
684    let title = title.map(str::to_owned);
685    submit_database_write("set_session_acp_title", move |_| {
686        set_session_acp_title_to(&database_path(), &session_id, title.as_deref())
687    })
688}
689
690pub(super) fn set_session_acp_title_to(
691    path: &Path,
692    session_id: &str,
693    title: Option<&str>,
694) -> Result<()> {
695    if title.is_some_and(|title| title.trim().is_empty()) {
696        bail!("ACP session title is empty");
697    }
698    let title = title.and_then(mj_core::state::normalize_session_title);
699    let connection = open(path)?;
700    let changed = connection.execute(
701        "UPDATE sessions SET acp_session_title = ?2 WHERE session_id = ?1",
702        params![session_id, title],
703    )?;
704    if changed != 1 {
705        bail!("unknown session {session_id}");
706    }
707    Ok(())
708}
709
710/// Commit the successful handshake for a newly provisioned worker without
711/// replacing checkpoint or display metadata owned by other controller tasks.
712pub fn mark_session_worker_connected(
713    session_id: &str,
714    native_session_id: Option<&str>,
715    updated_at: &str,
716) -> Result<()> {
717    let session_id = session_id.to_owned();
718    let native_session_id = native_session_id.map(str::to_owned);
719    let updated_at = updated_at.to_owned();
720    submit_database_write("mark_session_worker_connected", move |_| {
721        mark_session_worker_connected_to(
722            &database_path(),
723            &session_id,
724            native_session_id.as_deref(),
725            &updated_at,
726        )
727    })
728}
729
730/// Point a session at a native session its worker opened on its own. Only that
731/// column moves: the session's lifecycle state belongs to whatever operation is
732/// running.
733pub fn adopt_native_session_id(session_id: &str, native_session_id: &str) -> Result<()> {
734    let session_id = session_id.to_owned();
735    let native_session_id = native_session_id.to_owned();
736    submit_database_write("adopt_native_session_id", move |_| {
737        adopt_native_session_id_to(&database_path(), &session_id, &native_session_id)
738    })
739}
740
741pub(super) fn adopt_native_session_id_to(
742    path: &Path,
743    session_id: &str,
744    native_session_id: &str,
745) -> Result<()> {
746    let connection = open(path)?;
747    let changed = connection.execute(
748        "UPDATE sessions SET native_session_id = ?2 WHERE session_id = ?1",
749        params![session_id, native_session_id],
750    )?;
751    if changed != 1 {
752        bail!("unknown session {session_id}");
753    }
754    Ok(())
755}
756
757pub(super) fn mark_session_worker_connected_to(
758    path: &Path,
759    session_id: &str,
760    native_session_id: Option<&str>,
761    updated_at: &str,
762) -> Result<()> {
763    if updated_at.trim().is_empty() {
764        bail!("worker connection timestamp is empty");
765    }
766    let connection = open(path)?;
767    let changed = connection.execute(
768        "UPDATE sessions
769         SET state = 'running',
770             native_session_id = coalesce(?2, native_session_id),
771             updated_at = ?3,
772             last_error = NULL
773         WHERE session_id = ?1",
774        params![session_id, native_session_id, updated_at],
775    )?;
776    if changed != 1 {
777        bail!("unknown session {session_id}");
778    }
779    Ok(())
780}
781
782pub(super) fn recover_interrupted_checkpointing_sessions_to(
783    path: &Path,
784    updated_at: &str,
785) -> Result<usize> {
786    if updated_at.trim().is_empty() {
787        bail!("checkpoint recovery timestamp is empty");
788    }
789    let connection = open(path)?;
790    connection
791        .execute(
792            "UPDATE sessions
793             SET state = 'running', updated_at = ?1, last_checkpoint_error = ?2
794             WHERE state = 'checkpointing'",
795            params![
796                updated_at,
797                "checkpointing was interrupted by a controller restart; the target was left running"
798            ],
799        )
800        .context("recover interrupted checkpointing sessions")
801}
802
803pub(super) fn save_session_to(path: &Path, session: &SessionRecord) -> Result<()> {
804    save_session_with_container_size_to(path, session, None)
805}
806
807pub(super) fn save_session_with_container_size_to(
808    path: &Path,
809    session: &SessionRecord,
810    container_size: Option<(&str, HostContainerSize)>,
811) -> Result<()> {
812    validate_session_record(session)?;
813
814    let mut connection = open(path)?;
815    let tx = connection.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
816    if let Some(existing_bundle) = tx
817        .query_row(
818            "SELECT bundle_id FROM session_contexts WHERE session_id = ?1",
819            [session.id.as_str()],
820            |row| row.get::<_, String>(0),
821        )
822        .optional()?
823        && existing_bundle != session.bundle_id
824    {
825        bail!(
826            "session {} was already associated with bundle {}, not {}",
827            session.id,
828            existing_bundle,
829            session.bundle_id
830        );
831    }
832    let mut session = session.clone();
833    let moving: bool = tx.query_row(
834        "SELECT EXISTS(SELECT 1 FROM session_moves WHERE session_id=?1
835         AND json_extract(operation_json, '$.phase') IN ('preparing','closing_source','resuming_destination','starting_queue'))",
836        [&session.id], |row| row.get(0),
837    )?;
838    if moving {
839        // A Move may provision for minutes while clients keep editing drafts
840        // and titles. Merge these independently owned fields in this same
841        // transaction rather than restoring the lifecycle's earlier copy.
842        let (draft, title, acp_title, viewed, archived) = tx.query_row(
843            "SELECT draft_input, session_title_override, acp_session_title, viewed_through_event_ordinal, archived
844             FROM sessions WHERE session_id=?1", [&session.id], |row| Ok((
845                row.get::<_, String>(0)?, row.get::<_, Option<String>>(1)?, row.get::<_, Option<String>>(2)?,
846                row.get::<_, u64>(3)?, row.get::<_, bool>(4)?,
847            )),
848        )?;
849        session.draft_input = draft;
850        session.session_title_override = title;
851        session.acp_session_title = acp_title;
852        session.viewed_through_event_ordinal = viewed;
853        session.archived = archived;
854    }
855    insert_session(&tx, &session)?;
856    if let Some((host, size)) = container_size {
857        write_host_container_size(&tx, host, size)?;
858    }
859    tx.commit()?;
860    Ok(())
861}
862
863pub(super) fn save_lifecycle_session_to(path: &Path, session: &SessionRecord) -> Result<()> {
864    validate_session_record(session)?;
865
866    let mut connection = open(path)?;
867    let tx = connection.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
868    update_lifecycle_fields(&tx, session)?;
869    tx.commit()?;
870    Ok(())
871}
872
873pub(super) fn save_checkpointed_session_to(path: &Path, session: &SessionRecord) -> Result<()> {
874    validate_session_record(session)?;
875
876    let mut connection = open(path)?;
877    let tx = connection.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
878    update_lifecycle_fields(&tx, session)?;
879    tx.execute(
880        "UPDATE sessions SET native_session_id = ?2 WHERE session_id = ?1",
881        params![session.id, session.native_session_id],
882    )?;
883    replace_checkpoint(&tx, session)?;
884    tx.commit()?;
885    Ok(())
886}
887
888pub(super) fn validate_session_record(session: &SessionRecord) -> Result<()> {
889    let mut validation = State::default();
890    validation
891        .sessions
892        .insert(session.id.clone(), session.clone());
893    validation.validate()
894}
895
896/// Remove one operational session while retaining its relational history
897/// context and prompt history.
898pub fn delete_session(session_id: &str) -> Result<()> {
899    let session_id = session_id.to_owned();
900    submit_database_write("delete_session", move |_| {
901        delete_session_from(&database_path(), &session_id)
902    })
903}
904
905pub(super) fn delete_session_from(path: &Path, session_id: &str) -> Result<()> {
906    let connection = open(path)?;
907    connection.execute("DELETE FROM sessions WHERE session_id = ?1", [session_id])?;
908    Ok(())
909}
910
911/// Overwrite the unsent chat input carried across a detach. Unlike the read
912/// receipt this is not monotonic: a draft can shrink, and an empty string
913/// clears it.
914pub fn set_session_draft_input(session_id: &str, draft: &str) -> Result<()> {
915    let session_id = session_id.to_owned();
916    let draft = draft.to_owned();
917    submit_database_write("set_session_draft_input", move |_| {
918        set_session_draft_input_at(&database_path(), &session_id, &draft)
919    })
920}
921
922pub(super) fn set_session_draft_input_at(path: &Path, session_id: &str, draft: &str) -> Result<()> {
923    let mut connection = open(path)?;
924    let tx = connection.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
925    let updated = tx.execute(
926        "UPDATE sessions SET draft_input = ?2 WHERE session_id = ?1",
927        params![session_id, draft],
928    )?;
929    ensure!(updated == 1, "unknown session {session_id}");
930    tx.commit()?;
931    Ok(())
932}
933
934/// Retire a submitted shared draft without erasing a newer client's edit.
935pub fn clear_session_draft_input_if_matches(session_id: &str, expected: &str) -> Result<()> {
936    let session_id = session_id.to_owned();
937    let expected = expected.to_owned();
938    submit_database_write("clear_session_draft_input_if_matches", move |connection| {
939        connection.execute(
940            "UPDATE sessions SET draft_input = '' WHERE session_id = ?1 AND draft_input = ?2",
941            params![session_id, expected],
942        )?;
943        Ok(())
944    })
945}
946
947pub fn record_recovery_success(
948    session_id: &str,
949    native_session_id: &str,
950    checkpoint: &CheckpointMetadata,
951) -> Result<()> {
952    let session_id = session_id.to_owned();
953    let native_session_id = native_session_id.to_owned();
954    let checkpoint = checkpoint.clone();
955    submit_database_write("record_recovery_success", move |_| {
956        record_recovery_success_to(
957            &database_path(),
958            &session_id,
959            &native_session_id,
960            &checkpoint,
961        )
962    })
963}
964
965pub(super) fn record_recovery_success_to(
966    path: &Path,
967    session_id: &str,
968    native_session_id: &str,
969    checkpoint: &CheckpointMetadata,
970) -> Result<()> {
971    let mut connection = open(path)?;
972    let tx = connection.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
973    let changed = tx.execute(
974        "UPDATE sessions
975         SET native_session_id = ?2, last_checkpoint_error = NULL
976         WHERE session_id = ?1",
977        params![session_id, native_session_id],
978    )?;
979    if changed != 1 {
980        bail!("unknown session {session_id}");
981    }
982    tx.execute(
983        "INSERT INTO session_checkpoints(
984             session_id, archive_path, sha256, created_at, event_frontier
985         ) VALUES (?1,?2,?3,?4,?5)
986         ON CONFLICT(session_id) DO UPDATE SET
987             archive_path = excluded.archive_path,
988             sha256 = excluded.sha256,
989             created_at = excluded.created_at,
990             event_frontier = excluded.event_frontier",
991        params![
992            session_id,
993            path_to_blob(&checkpoint.archive_path),
994            checkpoint.sha256,
995            checkpoint.created_at,
996            checkpoint.event_frontier,
997        ],
998    )?;
999    tx.commit()?;
1000    Ok(())
1001}
1002
1003pub fn record_recovery_failure(session_id: &str, detail: &str) -> Result<()> {
1004    let session_id = session_id.to_owned();
1005    let detail = detail.to_owned();
1006    submit_database_write("record_recovery_failure", move |_| {
1007        record_recovery_failure_to(&database_path(), &session_id, &detail)
1008    })
1009}
1010
1011pub(super) fn record_recovery_failure_to(
1012    path: &Path,
1013    session_id: &str,
1014    detail: &str,
1015) -> Result<()> {
1016    let connection = open(path)?;
1017    let changed = connection.execute(
1018        "UPDATE sessions SET last_checkpoint_error = ?2 WHERE session_id = ?1",
1019        params![session_id, detail],
1020    )?;
1021    if changed != 1 {
1022        bail!("unknown session {session_id}");
1023    }
1024    Ok(())
1025}
1026
1027/// Re-associate a session with another project bundle.
1028///
1029/// A session's bundle is otherwise fixed, because prompt history is grouped by
1030/// it. Resume calls this when it converts a session between its raw and bundle
1031/// representations: the project is the same, so its history follows it, and
1032/// only the name Hel files it under changes.
1033pub fn rebind_session_bundle(session_id: &str, bundle_id: &str) -> Result<()> {
1034    let session_id = session_id.to_owned();
1035    let bundle_id = bundle_id.to_owned();
1036    submit_database_write("rebind_session_bundle", move |_| {
1037        rebind_session_bundle_to(&database_path(), &session_id, &bundle_id)
1038    })
1039}
1040
1041pub(super) fn rebind_session_bundle_to(
1042    path: &Path,
1043    session_id: &str,
1044    bundle_id: &str,
1045) -> Result<()> {
1046    let mut connection = open(path)?;
1047    let tx = connection.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
1048    let changed = tx.execute(
1049        "UPDATE session_contexts SET bundle_id = ?2 WHERE session_id = ?1",
1050        params![session_id, bundle_id],
1051    )?;
1052    if changed == 0 {
1053        tx.execute(
1054            "INSERT INTO session_contexts(session_id, bundle_id, created_at) VALUES (?1, ?2, ?3)",
1055            params![session_id, bundle_id, Utc::now().to_rfc3339()],
1056        )?;
1057    }
1058    tx.commit()?;
1059    Ok(())
1060}