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