use super::*;
pub fn session_incarnation(session_id: &str) -> Result<Option<String>> {
let connection = open_reader(&database_path())?;
connection
.query_row(
"SELECT identity FROM session_incarnations WHERE session_id = ?1",
[session_id],
|row| row.get(0),
)
.optional()
.map_err(Into::into)
}
pub fn save_session(session: &SessionRecord) -> Result<()> {
let session = session.clone();
submit_database_write("save_session", move |_| {
save_session_to(&database_path(), &session)
})
}
pub fn save_new_session(
session: &SessionRecord,
container_size: Option<(String, HostContainerSize)>,
) -> Result<()> {
let session = session.clone();
submit_database_write("save_new_session", move |_| {
save_new_session_to(&database_path(), &session, container_size)
})
}
pub(super) fn save_new_session_to(
path: &Path,
session: &SessionRecord,
container_size: Option<(String, HostContainerSize)>,
) -> Result<()> {
let mut connection = open(path)?;
let tx = connection.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
validate_session_record(session)?;
insert_session(&tx, session)?;
if let Some((host, size)) = container_size {
write_host_container_size(&tx, &host, size)?;
}
if session.harness_kind.supports_delegation_tools() {
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())?])?;
}
tx.commit()?;
Ok(())
}
pub fn set_publication_assessment_if_current(
session_id: &str,
assessment: &mj_core::state::PublicationAssessment,
) -> Result<bool> {
let session_id = session_id.to_owned();
let assessment = assessment.clone();
submit_database_write("set_publication_assessment_if_current", move |_| {
let connection = open(&database_path())?;
let updated = connection.execute(
"UPDATE sessions SET publication_json = ?2
WHERE session_id = ?1 AND state = 'stopped'
AND EXISTS (SELECT 1 FROM session_checkpoints c
WHERE c.session_id = ?1 AND c.sha256 = ?3)",
params![
session_id,
serde_json::to_string(&assessment)?,
assessment.checkpoint_sha256
],
)?;
Ok(updated == 1)
})
}
pub fn save_subagent_session(
session: &SessionRecord,
subagent: &mj_core::subagent::SubagentRecord,
) -> Result<()> {
let session = session.clone();
let subagent = subagent.clone();
submit_database_write("save_subagent_session", move |_| {
let mut connection = open(&database_path())?;
let tx = connection.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
insert_session(&tx, &session)?;
tx.execute(
"INSERT INTO subagent_sessions(
child_session_id, parent_session_id, request_key, record_json
) VALUES (?1, ?2, ?3, ?4)",
params![
subagent.child_session_id,
subagent.parent_session_id,
subagent.request_key,
serde_json::to_string(&subagent)?,
],
)?;
tx.commit()?;
Ok(())
})
}
pub fn mark_subagent_turn_noticed(child_session_id: &str, turn: u64) -> Result<()> {
let child_session_id = child_session_id.to_owned();
submit_database_write("mark_subagent_turn_noticed", move |_| {
let mut relation = load_subagent(&child_session_id)?
.with_context(|| format!("unknown sub-agent session {child_session_id}"))?;
relation.noticed_turn = Some(turn);
let json = serde_json::to_string(&relation)?;
let connection = open(&database_path())?;
connection.execute(
"UPDATE subagent_sessions SET record_json = ?2 WHERE child_session_id = ?1",
params![child_session_id, json],
)?;
Ok(())
})
}
pub fn load_subagent(child_session_id: &str) -> Result<Option<mj_core::subagent::SubagentRecord>> {
let connection = open_reader(&database_path())?;
connection
.query_row(
"SELECT record_json FROM subagent_sessions WHERE child_session_id = ?1",
[child_session_id],
|row| row.get::<_, String>(0),
)
.optional()?
.map(|json| serde_json::from_str(&json).context("decode sub-agent record"))
.transpose()
}
pub fn list_subagents(parent_session_id: &str) -> Result<Vec<mj_core::subagent::SubagentRecord>> {
let connection = open_reader(&database_path())?;
let mut statement = connection.prepare(
"SELECT record_json FROM subagent_sessions
WHERE parent_session_id = ?1 ORDER BY rowid",
)?;
statement
.query_map([parent_session_id], |row| row.get::<_, String>(0))?
.map(|row| serde_json::from_str(&row?).context("decode sub-agent record"))
.collect()
}
pub fn record_stopped_subagents(
parent_session_id: &str,
stopped: &[mj_core::subagent::StoppedSubagent],
) -> Result<()> {
let parent_session_id = parent_session_id.to_owned();
let stopped = stopped.to_vec();
submit_database_write("record_stopped_subagents", move |_| {
record_stopped_subagents_to(&database_path(), &parent_session_id, &stopped)
})
}
pub(super) fn record_stopped_subagents_to(
path: &Path,
parent_session_id: &str,
stopped: &[mj_core::subagent::StoppedSubagent],
) -> Result<()> {
let mut connection = open(path)?;
let tx = connection.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
for child in stopped {
tx.execute(
"INSERT INTO stopped_subagents(parent_session_id, child_session_id, record_json)
VALUES (?1, ?2, ?3)
ON CONFLICT(parent_session_id, child_session_id) DO UPDATE SET
record_json = excluded.record_json",
params![
parent_session_id,
child.child_session_id,
serde_json::to_string(child)?
],
)?;
}
tx.commit()?;
Ok(())
}
pub fn load_stopped_subagents(
parent_session_id: &str,
) -> Result<Vec<mj_core::subagent::StoppedSubagent>> {
load_stopped_subagents_from(&database_path(), parent_session_id)
}
pub(super) fn load_stopped_subagents_from(
path: &Path,
parent_session_id: &str,
) -> Result<Vec<mj_core::subagent::StoppedSubagent>> {
let connection = open_reader(path)?;
let mut statement = connection.prepare(
"SELECT record_json FROM stopped_subagents
WHERE parent_session_id = ?1 ORDER BY rowid",
)?;
statement
.query_map([parent_session_id], |row| row.get::<_, String>(0))?
.map(|row| serde_json::from_str(&row?).context("decode stopped sub-agent record"))
.collect()
}
pub fn clear_stopped_subagents(
parent_session_id: &str,
child_session_ids: &[String],
) -> Result<()> {
let parent_session_id = parent_session_id.to_owned();
let child_session_ids = child_session_ids.to_vec();
submit_database_write("clear_stopped_subagents", move |_| {
clear_stopped_subagents_from(&database_path(), &parent_session_id, &child_session_ids)
})
}
pub(super) fn clear_stopped_subagents_from(
path: &Path,
parent_session_id: &str,
child_session_ids: &[String],
) -> Result<()> {
let mut connection = open(path)?;
let tx = connection.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
for child_session_id in child_session_ids {
tx.execute(
"DELETE FROM stopped_subagents
WHERE parent_session_id = ?1 AND child_session_id = ?2",
params![parent_session_id, child_session_id],
)?;
}
tx.commit()?;
Ok(())
}
pub fn load_session_state(session_id: &str) -> Result<Option<SessionState>> {
let connection = open_reader(&database_path())?;
let stored = connection
.query_row(
"SELECT state FROM sessions WHERE session_id = ?1",
[session_id],
|row| row.get::<_, String>(0),
)
.optional()?;
Ok(stored.as_deref().map(stored_session_state))
}
pub fn load_subagent_report(child_session_id: &str) -> Result<mj_core::subagent::SubagentReport> {
load_subagent_report_from(&database_path(), child_session_id)
}
pub(super) fn load_subagent_report_from(
path: &Path,
child_session_id: &str,
) -> Result<mj_core::subagent::SubagentReport> {
let connection = open_reader(path)?;
let row = connection
.query_row(
"SELECT handback_command_id, handback_message, handback_recorded_at_ms,
reminder_command_id, reminder_for_command_id, reminder_sent_at_ms,
reminder_failed_for_command_id, awaited_ordinal, report_dir
FROM subagent_handbacks WHERE child_session_id = ?1",
[child_session_id],
|row| {
Ok((
row.get::<_, Option<String>>(0)?,
row.get::<_, Option<String>>(1)?,
row.get::<_, Option<i64>>(2)?,
row.get::<_, Option<String>>(3)?,
row.get::<_, Option<String>>(4)?,
row.get::<_, Option<i64>>(5)?,
row.get::<_, Option<String>>(6)?,
row.get::<_, Option<i64>>(7)?,
row.get::<_, Option<String>>(8)?,
))
},
)
.optional()?;
let Some((
handback_command,
handback_message,
handback_at,
reminder_command,
reminder_for,
reminder_at,
reminder_failed_for,
awaited_ordinal,
report_dir,
)) = row
else {
return Ok(mj_core::subagent::SubagentReport::default());
};
Ok(mj_core::subagent::SubagentReport {
handback: match (handback_command, handback_message, handback_at) {
(Some(command_id), Some(message), Some(recorded_at_ms)) => {
Some(mj_core::subagent::SubagentHandback {
command_id,
message,
recorded_at_ms,
})
}
_ => None,
},
reminder: match (reminder_command, reminder_for, reminder_at) {
(Some(command_id), Some(for_command_id), Some(sent_at_ms)) => {
Some(mj_core::subagent::HandbackReminder {
command_id,
for_command_id,
sent_at_ms,
})
}
_ => None,
},
reminder_failed_for,
awaited_ordinal: awaited_ordinal.and_then(|ordinal| u64::try_from(ordinal).ok()),
report_dir,
})
}
pub fn record_subagent_report_dir(child_session_id: &str, report_dir: &str) -> Result<()> {
let child_session_id = child_session_id.to_owned();
let report_dir = report_dir.to_owned();
submit_database_write("record_subagent_report_dir", move |_| {
record_subagent_report_dir_to(&database_path(), &child_session_id, &report_dir)
})
}
pub(super) fn record_subagent_report_dir_to(
path: &Path,
child_session_id: &str,
report_dir: &str,
) -> Result<()> {
open(path)?.execute(
"INSERT INTO subagent_handbacks(child_session_id, report_dir)
SELECT ?1, ?2 WHERE EXISTS (
SELECT 1 FROM subagent_sessions WHERE child_session_id = ?1
)
ON CONFLICT(child_session_id) DO UPDATE SET report_dir = excluded.report_dir",
params![child_session_id, report_dir],
)?;
Ok(())
}
pub fn record_subagent_prompt(child_session_id: &str, ordinal: u64) -> Result<()> {
let child_session_id = child_session_id.to_owned();
submit_database_write("record_subagent_prompt", move |_| {
record_subagent_prompt_to(&database_path(), &child_session_id, ordinal)
})
}
pub(super) fn record_subagent_prompt_to(
path: &Path,
child_session_id: &str,
ordinal: u64,
) -> Result<()> {
let ordinal = i64::try_from(ordinal).context("prompt ordinal exceeds the store's range")?;
open(path)?.execute(
"INSERT INTO subagent_handbacks(child_session_id, awaited_ordinal)
SELECT ?1, ?2 WHERE EXISTS (
SELECT 1 FROM subagent_sessions WHERE child_session_id = ?1
)
ON CONFLICT(child_session_id) DO UPDATE SET
awaited_ordinal = max(coalesce(awaited_ordinal, 0), excluded.awaited_ordinal)",
params![child_session_id, ordinal],
)?;
Ok(())
}
pub fn record_subagent_handback(
child_session_id: &str,
handback: &mj_core::subagent::SubagentHandback,
) -> Result<bool> {
let child_session_id = child_session_id.to_owned();
let handback = handback.clone();
submit_database_write("record_subagent_handback", move |_| {
record_subagent_handback_to(&database_path(), &child_session_id, &handback)
})
}
pub(super) fn record_subagent_handback_to(
path: &Path,
child_session_id: &str,
handback: &mj_core::subagent::SubagentHandback,
) -> Result<bool> {
let mut connection = open(path)?;
let tx = connection.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
let recorded_for: Option<String> = tx
.query_row(
"SELECT handback_command_id FROM subagent_handbacks WHERE child_session_id = ?1",
[child_session_id],
|row| row.get(0),
)
.optional()?
.flatten();
if recorded_for.as_deref() == Some(handback.command_id.as_str()) {
return Ok(false);
}
tx.execute(
"INSERT INTO subagent_handbacks(
child_session_id, handback_command_id, handback_message, handback_recorded_at_ms
) VALUES (?1, ?2, ?3, ?4)
ON CONFLICT(child_session_id) DO UPDATE SET
handback_command_id = excluded.handback_command_id,
handback_message = excluded.handback_message,
handback_recorded_at_ms = excluded.handback_recorded_at_ms",
params![
child_session_id,
handback.command_id,
handback.message,
handback.recorded_at_ms
],
)?;
tx.commit()?;
Ok(true)
}
pub fn record_handback_reminder(
child_session_id: &str,
reminder: &mj_core::subagent::HandbackReminder,
) -> Result<()> {
let child_session_id = child_session_id.to_owned();
let reminder = reminder.clone();
submit_database_write("record_handback_reminder", move |_| {
open(&database_path())?.execute(
"INSERT INTO subagent_handbacks(
child_session_id, reminder_command_id, reminder_for_command_id, reminder_sent_at_ms
) VALUES (?1, ?2, ?3, ?4)
ON CONFLICT(child_session_id) DO UPDATE SET
reminder_command_id = excluded.reminder_command_id,
reminder_for_command_id = excluded.reminder_for_command_id,
reminder_sent_at_ms = excluded.reminder_sent_at_ms",
params![
child_session_id,
reminder.command_id,
reminder.for_command_id,
reminder.sent_at_ms
],
)?;
Ok(())
})
}
pub fn record_handback_reminder_failed(child_session_id: &str, for_command_id: &str) -> Result<()> {
let child_session_id = child_session_id.to_owned();
let for_command_id = for_command_id.to_owned();
submit_database_write("record_handback_reminder_failed", move |_| {
open(&database_path())?.execute(
"INSERT INTO subagent_handbacks(child_session_id, reminder_failed_for_command_id)
VALUES (?1, ?2)
ON CONFLICT(child_session_id) DO UPDATE SET
reminder_failed_for_command_id = excluded.reminder_failed_for_command_id",
params![child_session_id, for_command_id],
)?;
Ok(())
})
}
pub fn lookup_subagent_request(
parent_session_id: &str,
request_key: &str,
) -> Result<Option<mj_core::subagent::SubagentRecord>> {
let connection = open_reader(&database_path())?;
connection
.query_row(
"SELECT record_json FROM subagent_sessions
WHERE parent_session_id = ?1 AND request_key = ?2",
params![parent_session_id, request_key],
|row| row.get::<_, String>(0),
)
.optional()?
.map(|json| serde_json::from_str(&json).context("decode sub-agent record"))
.transpose()
}
pub fn save_session_with_container_size(
session: &SessionRecord,
host: &str,
size: HostContainerSize,
) -> Result<()> {
let session = session.clone();
let host = host.to_owned();
submit_database_write("save_session_with_container_size", move |_| {
save_session_with_container_size_to(&database_path(), &session, Some((&host, size)))
})
}
pub fn save_resumed_session(
session: &SessionRecord,
container_size: Option<(&str, HostContainerSize)>,
) -> Result<()> {
let session = session.clone();
let container_size = container_size.map(|(host, size)| (host.to_owned(), size));
submit_database_write("save_resumed_session", move |_| {
save_resumed_session_to(
&database_path(),
&session,
container_size
.as_ref()
.map(|(host, size)| (host.as_str(), *size)),
)
})
}
pub(super) fn save_resumed_session_to(
path: &Path,
session: &SessionRecord,
container_size: Option<(&str, HostContainerSize)>,
) -> Result<()> {
validate_session_record(session)?;
let mut connection = open(path)?;
let tx = connection.transaction()?;
let (bundle, workspace): (String, String) = tx.query_row(
"SELECT bundle_id, workspace_id FROM session_contexts WHERE session_id = ?1",
[&session.id],
|row| Ok((row.get(0)?, row.get(1)?)),
)?;
ensure!(
bundle == session.bundle_id && workspace == session.workspace_id,
"session {} context changed before resume publication",
session.id
);
update_lifecycle_fields(&tx, session)?;
tx.execute(
"UPDATE sessions SET native_session_id = ?2, container_cpus = ?3,
container_memory = ?4, container_workspace = ?5,
create_managed_worktree = ?6, launch_base = ?7, launch_branch = ?8,
checkout_json = ?9, expected_runtime_identity = ?10
WHERE session_id = ?1",
params![
session.id,
session.native_session_id,
session.container_cpus,
session.container_memory,
session
.container_workspace
.as_ref()
.map(|path| path.to_string_lossy().into_owned()),
session.create_managed_worktree,
session.launch_base,
session.launch_branch,
session
.checkout
.as_ref()
.map(serde_json::to_string)
.transpose()?,
session.expected_runtime_identity,
],
)?;
replace_mounts(&tx, &session.id, &session.additional_mounts)?;
replace_checkpoint(&tx, session)?;
if let Some((host, size)) = container_size {
write_host_container_size(&tx, host, size)?;
}
tx.commit()?;
Ok(())
}
pub fn save_lifecycle_session(session: &SessionRecord) -> Result<()> {
let session = session.clone();
submit_database_write("save_lifecycle_session", move |_| {
save_lifecycle_session_to(&database_path(), &session)
})
}
pub fn save_checkpointed_session(session: &SessionRecord) -> Result<()> {
let session = session.clone();
submit_database_write("save_checkpointed_session", move |_| {
save_checkpointed_session_to(&database_path(), &session)
})
}
pub fn recover_interrupted_checkpointing_sessions(updated_at: &str) -> Result<usize> {
let updated_at = updated_at.to_owned();
submit_database_write("recover_interrupted_checkpointing_sessions", move |_| {
recover_interrupted_checkpointing_sessions_to(&database_path(), &updated_at)
})
}
pub fn set_session_title_override(session_id: &str, title: &str, updated_at: &str) -> Result<()> {
let session_id = session_id.to_owned();
let title = title.to_owned();
let updated_at = updated_at.to_owned();
submit_database_write("set_session_title_override", move |_| {
set_session_title_override_to(&database_path(), &session_id, &title, &updated_at)
})
}
pub fn rename_profile_references(old_id: &str, new_id: &str) -> Result<usize> {
rename_session_reference("last_profile", old_id, new_id)
}
pub fn rename_target_references(old_id: &str, new_id: &str) -> Result<usize> {
rename_session_reference("target_template_id", old_id, new_id)
}
pub(super) fn rename_session_reference(
column: &'static str,
old_id: &str,
new_id: &str,
) -> Result<usize> {
ensure!(
matches!(column, "last_profile" | "target_template_id"),
"unsupported session reference column"
);
let old_id = old_id.to_owned();
let new_id = new_id.to_owned();
submit_database_write("rename_session_reference", move |_| {
rename_session_reference_at(&database_path(), column, &old_id, &new_id)
})
}
pub(super) fn rename_session_reference_at(
path: &Path,
column: &str,
old_id: &str,
new_id: &str,
) -> Result<usize> {
let mut connection = open(path)?;
let tx = connection.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
let changed = tx.execute(
&format!("UPDATE sessions SET {column} = ?2 WHERE {column} = ?1"),
params![old_id, new_id],
)?;
tx.commit()?;
Ok(changed)
}
pub fn set_session_archived(session_id: &str, archived: bool) -> Result<()> {
let session_id = session_id.to_owned();
submit_database_write("set_session_archived", move |_| {
set_session_archived_to(&database_path(), &session_id, archived)
})
}
pub fn mark_session_target_missing(
session_id: &str,
detail: &str,
updated_at: &str,
) -> Result<Option<SessionState>> {
let session_id = session_id.to_owned();
let detail = detail.to_owned();
let updated_at = updated_at.to_owned();
submit_database_write("mark_session_target_missing", move |_| {
mark_session_target_missing_to(&database_path(), &session_id, &detail, &updated_at)
})
}
pub(super) fn mark_session_target_missing_to(
path: &Path,
session_id: &str,
detail: &str,
updated_at: &str,
) -> Result<Option<SessionState>> {
mark_session_target_missing_if_current_to(path, session_id, detail, updated_at, None)
}
pub fn mark_session_target_missing_if_current(
session_id: &str,
detail: &str,
updated_at: &str,
observed_updated_at: &str,
) -> Result<Option<SessionState>> {
let session_id = session_id.to_owned();
let detail = detail.to_owned();
let updated_at = updated_at.to_owned();
let observed_updated_at = observed_updated_at.to_owned();
submit_database_write("mark_session_target_missing_if_current", move |_| {
mark_session_target_missing_if_current_to(
&database_path(),
&session_id,
&detail,
&updated_at,
Some(&observed_updated_at),
)
})
}
pub(super) fn mark_session_target_missing_if_current_to(
path: &Path,
session_id: &str,
detail: &str,
updated_at: &str,
observed_updated_at: Option<&str>,
) -> Result<Option<SessionState>> {
let mut connection = open(path)?;
let tx = connection.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
let previous_error = super::events::previous_session_error(&tx, session_id)?;
let changed = tx.execute(
"UPDATE sessions
SET state = CASE
WHEN EXISTS(
SELECT 1 FROM session_checkpoints
WHERE session_checkpoints.session_id = sessions.session_id
) THEN 'error'
ELSE 'lost'
END,
last_error = ?2,
updated_at = ?3
WHERE session_id = ?1
AND (?4 IS NULL OR updated_at = ?4)
AND state IN ('provisioning', 'running', 'disconnected', 'error')",
params![session_id, detail, updated_at, observed_updated_at],
)?;
ensure!(changed <= 1, "updated {changed} sessions for {session_id}");
let state = if changed == 1 {
let stored: String = tx.query_row(
"SELECT state FROM sessions WHERE session_id = ?1",
[session_id],
|row| row.get(0),
)?;
Some(stored_session_state(&stored))
} else {
None
};
if changed == 1 && previous_error.as_deref() != Some(detail) {
super::events::insert_api_event(
&tx,
session_id,
Utc::now().timestamp_millis(),
&ApiEventData::SessionFault {
reason: mj_core::event_outcome::OutcomeReason::RuntimeUnavailable,
message: detail.into(),
command_id: None,
},
)?;
}
tx.commit()?;
Ok(state)
}
pub(super) fn set_session_archived_to(path: &Path, session_id: &str, archived: bool) -> Result<()> {
let connection = open(path)?;
let changed = connection.execute(
"UPDATE sessions SET archived = ?2 WHERE session_id = ?1",
params![session_id, archived],
)?;
if changed != 1 {
bail!("unknown session {session_id}");
}
Ok(())
}
pub fn hidden_native_sessions() -> Result<BTreeSet<(mj_core::config::HarnessKind, String)>> {
hidden_native_sessions_from(&database_path())
}
pub(super) fn hidden_native_sessions_from(
path: &Path,
) -> Result<BTreeSet<(mj_core::config::HarnessKind, String)>> {
let connection = open_reader(path)?;
let mut statement =
connection.prepare("SELECT harness_kind, native_session_id FROM hidden_native_sessions")?;
let rows = statement.query_map([], |row| {
Ok((row.get::<_, String>(0)?, row.get::<_, String>(1)?))
})?;
let mut hidden = BTreeSet::new();
for row in rows {
let (harness, native_session_id) = row?;
match harness.parse::<mj_core::config::HarnessKind>() {
Ok(harness) => {
hidden.insert((harness, native_session_id));
}
Err(_) => tracing::warn!(
harness = %harness,
"ignoring a hidden native session for a harness that is no longer supported"
),
}
}
Ok(hidden)
}
pub fn set_native_session_hidden(
harness: mj_core::config::HarnessKind,
native_session_id: &str,
hidden: bool,
) -> Result<()> {
let native_session_id = native_session_id.to_owned();
submit_database_write("set_native_session_hidden", move |_| {
set_native_session_hidden_to(&database_path(), harness, &native_session_id, hidden)
})
}
pub(super) fn set_native_session_hidden_to(
path: &Path,
harness: mj_core::config::HarnessKind,
native_session_id: &str,
hidden: bool,
) -> Result<()> {
if native_session_id.trim().is_empty() {
bail!("native session id is empty");
}
let connection = open(path)?;
if hidden {
connection.execute(
"INSERT INTO hidden_native_sessions(harness_kind, native_session_id, hidden_at)
VALUES (?1, ?2, ?3)
ON CONFLICT(harness_kind, native_session_id) DO NOTHING",
params![harness.id(), native_session_id, Utc::now().to_rfc3339()],
)?;
} else {
connection.execute(
"DELETE FROM hidden_native_sessions
WHERE harness_kind = ?1 AND native_session_id = ?2",
params![harness.id(), native_session_id],
)?;
}
Ok(())
}
pub fn set_session_container_settings(
session_id: &str,
cpus: Option<&str>,
memory: Option<&str>,
mounts: &[AdditionalMount],
updated_at: &str,
) -> Result<()> {
let session_id = session_id.to_owned();
let cpus = cpus.map(str::to_owned);
let memory = memory.map(str::to_owned);
let mounts = mounts.to_vec();
let updated_at = updated_at.to_owned();
submit_database_write("set_session_container_settings", move |_| {
set_session_container_settings_to(
&database_path(),
&session_id,
cpus.as_deref(),
memory.as_deref(),
&mounts,
&updated_at,
)
})
}
pub(super) fn set_session_container_settings_to(
path: &Path,
session_id: &str,
cpus: Option<&str>,
memory: Option<&str>,
mounts: &[AdditionalMount],
updated_at: &str,
) -> Result<()> {
if updated_at.trim().is_empty() {
bail!("session update timestamp is empty");
}
crate::targets::validate_additional_mounts(mounts)?;
let mut connection = open(path)?;
let tx = connection.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
let changed = tx.execute(
"UPDATE sessions
SET container_cpus = ?2, container_memory = ?3, updated_at = ?4
WHERE session_id = ?1",
params![session_id, cpus, memory, updated_at],
)?;
if changed != 1 {
bail!("unknown session {session_id}");
}
replace_mounts(&tx, session_id, mounts)?;
tx.commit()?;
Ok(())
}
pub(super) fn set_session_title_override_to(
path: &Path,
session_id: &str,
title: &str,
updated_at: &str,
) -> Result<()> {
if title.trim().is_empty() {
bail!("session title is empty");
}
if updated_at.trim().is_empty() {
bail!("session update timestamp is empty");
}
let connection = open(path)?;
let changed = connection.execute(
"UPDATE sessions
SET session_title_override = ?2, updated_at = ?3
WHERE session_id = ?1",
params![session_id, title, updated_at],
)?;
if changed != 1 {
bail!("unknown session {session_id}");
}
Ok(())
}
pub fn set_session_acp_title(session_id: &str, title: Option<&str>) -> Result<()> {
let session_id = session_id.to_owned();
let title = title.map(str::to_owned);
submit_database_write("set_session_acp_title", move |_| {
set_session_acp_title_to(&database_path(), &session_id, title.as_deref())
})
}
pub(super) fn set_session_acp_title_to(
path: &Path,
session_id: &str,
title: Option<&str>,
) -> Result<()> {
if title.is_some_and(|title| title.trim().is_empty()) {
bail!("ACP session title is empty");
}
let title = title.and_then(mj_core::state::normalize_session_title);
let connection = open(path)?;
let changed = connection.execute(
"UPDATE sessions SET acp_session_title = ?2 WHERE session_id = ?1",
params![session_id, title],
)?;
if changed != 1 {
bail!("unknown session {session_id}");
}
Ok(())
}
pub fn mark_session_worker_connected(
session_id: &str,
native_session_id: Option<&str>,
updated_at: &str,
) -> Result<()> {
let session_id = session_id.to_owned();
let native_session_id = native_session_id.map(str::to_owned);
let updated_at = updated_at.to_owned();
submit_database_write("mark_session_worker_connected", move |_| {
mark_session_worker_connected_to(
&database_path(),
&session_id,
native_session_id.as_deref(),
&updated_at,
)
})
}
pub fn adopt_native_session_id(session_id: &str, native_session_id: &str) -> Result<()> {
let session_id = session_id.to_owned();
let native_session_id = native_session_id.to_owned();
submit_database_write("adopt_native_session_id", move |_| {
adopt_native_session_id_to(&database_path(), &session_id, &native_session_id)
})
}
pub(super) fn adopt_native_session_id_to(
path: &Path,
session_id: &str,
native_session_id: &str,
) -> Result<()> {
let connection = open(path)?;
let changed = connection.execute(
"UPDATE sessions SET native_session_id = ?2 WHERE session_id = ?1",
params![session_id, native_session_id],
)?;
if changed != 1 {
bail!("unknown session {session_id}");
}
Ok(())
}
pub(super) fn mark_session_worker_connected_to(
path: &Path,
session_id: &str,
native_session_id: Option<&str>,
updated_at: &str,
) -> Result<()> {
if updated_at.trim().is_empty() {
bail!("worker connection timestamp is empty");
}
let connection = open(path)?;
let changed = connection.execute(
"UPDATE sessions
SET state = 'running',
native_session_id = coalesce(?2, native_session_id),
updated_at = ?3,
last_error = NULL
WHERE session_id = ?1",
params![session_id, native_session_id, updated_at],
)?;
if changed != 1 {
bail!("unknown session {session_id}");
}
Ok(())
}
pub(super) fn recover_interrupted_checkpointing_sessions_to(
path: &Path,
updated_at: &str,
) -> Result<usize> {
ensure!(
!updated_at.trim().is_empty(),
"checkpoint recovery timestamp is empty"
);
let mut connection = open(path)?;
let tx = connection.transaction()?;
let ids = {
let mut query = tx.prepare("SELECT session_id FROM checkpoint_operations")?;
query
.query_map([], |row| row.get::<_, String>(0))?
.collect::<rusqlite::Result<Vec<_>>>()?
};
for id in ids {
super::events::finish_checkpoint_operation(
&tx,
&id,
mj_core::event_outcome::CommandResultKind::Failed,
Some(mj_core::event_outcome::OutcomeReason::ControllerRestarted),
Some("Checkpoint operation interrupted by controller restart".into()),
)?;
}
let changed = tx.execute("UPDATE sessions SET state = 'running', updated_at = ?1, last_checkpoint_error = ?2 WHERE state = 'checkpointing'",
params![updated_at, "checkpointing was interrupted by a controller restart; the target was left running"])?;
tx.commit()?;
Ok(changed)
}
pub(super) fn save_session_to(path: &Path, session: &SessionRecord) -> Result<()> {
save_session_with_container_size_to(path, session, None)
}
pub(super) fn save_session_with_container_size_to(
path: &Path,
session: &SessionRecord,
container_size: Option<(&str, HostContainerSize)>,
) -> Result<()> {
validate_session_record(session)?;
let mut connection = open(path)?;
let tx = connection.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
if let Some(existing_bundle) = tx
.query_row(
"SELECT bundle_id FROM session_contexts WHERE session_id = ?1",
[session.id.as_str()],
|row| row.get::<_, String>(0),
)
.optional()?
&& existing_bundle != session.bundle_id
{
bail!(
"session {} was already associated with bundle {}, not {}",
session.id,
existing_bundle,
session.bundle_id
);
}
let mut session = session.clone();
let moving: bool = tx.query_row(
"SELECT EXISTS(SELECT 1 FROM session_moves WHERE session_id=?1
AND json_extract(operation_json, '$.phase') IN ('preparing','closing_source','resuming_destination','starting_queue'))",
[&session.id], |row| row.get(0),
)?;
if moving {
let (draft, title, acp_title, viewed, archived) = tx.query_row(
"SELECT draft_input, session_title_override, acp_session_title, viewed_through_event_ordinal, archived
FROM sessions WHERE session_id=?1", [&session.id], |row| Ok((
row.get::<_, String>(0)?, row.get::<_, Option<String>>(1)?, row.get::<_, Option<String>>(2)?,
row.get::<_, u64>(3)?, row.get::<_, bool>(4)?,
)),
)?;
session.draft_input = draft;
session.session_title_override = title;
session.acp_session_title = acp_title;
session.viewed_through_event_ordinal = viewed;
session.archived = archived;
}
insert_session(&tx, &session)?;
if let Some((host, size)) = container_size {
write_host_container_size(&tx, host, size)?;
}
tx.commit()?;
Ok(())
}
pub(super) fn save_lifecycle_session_to(path: &Path, session: &SessionRecord) -> Result<()> {
validate_session_record(session)?;
let mut connection = open(path)?;
let tx = connection.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
update_lifecycle_fields(&tx, session)?;
tx.commit()?;
Ok(())
}
pub(super) fn save_checkpointed_session_to(path: &Path, session: &SessionRecord) -> Result<()> {
validate_session_record(session)?;
let mut connection = open(path)?;
let tx = connection.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
update_lifecycle_fields(&tx, session)?;
tx.execute(
"UPDATE sessions SET native_session_id = ?2 WHERE session_id = ?1",
params![session.id, session.native_session_id],
)?;
replace_checkpoint(&tx, session)?;
tx.commit()?;
Ok(())
}
pub(super) fn validate_session_record(session: &SessionRecord) -> Result<()> {
let mut validation = State::default();
validation
.sessions
.insert(session.id.clone(), session.clone());
validation.validate()
}
pub fn delete_session(session_id: &str) -> Result<()> {
let session_id = session_id.to_owned();
submit_database_write("delete_session", move |_| {
delete_session_from(&database_path(), &session_id)
})
}
pub(super) fn delete_session_from(path: &Path, session_id: &str) -> Result<()> {
let connection = open(path)?;
connection.execute("DELETE FROM sessions WHERE session_id = ?1", [session_id])?;
Ok(())
}
pub fn set_session_draft_input(session_id: &str, draft: &str) -> Result<()> {
let session_id = session_id.to_owned();
let draft = draft.to_owned();
submit_database_write("set_session_draft_input", move |_| {
set_session_draft_input_at(&database_path(), &session_id, &draft)
})
}
pub(super) fn set_session_draft_input_at(path: &Path, session_id: &str, draft: &str) -> Result<()> {
let mut connection = open(path)?;
let tx = connection.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
let updated = tx.execute(
"UPDATE sessions SET draft_input = ?2 WHERE session_id = ?1",
params![session_id, draft],
)?;
ensure!(updated == 1, "unknown session {session_id}");
tx.commit()?;
Ok(())
}
pub fn append_session_draft_input(session_id: &str, text: &str) -> Result<()> {
let session_id = session_id.to_owned();
let text = text.to_owned();
submit_database_write("append_session_draft_input", move |connection| {
let updated = connection.execute(
"UPDATE sessions SET draft_input = CASE WHEN ?2 = '' THEN draft_input
WHEN draft_input = '' THEN ?2 ELSE draft_input || char(10) || char(10) || ?2 END
WHERE session_id = ?1",
params![session_id, text],
)?;
ensure!(updated == 1, "unknown session {session_id}");
Ok(())
})
}
pub fn clear_session_draft_input_if_matches(session_id: &str, expected: &str) -> Result<()> {
let session_id = session_id.to_owned();
let expected = expected.to_owned();
submit_database_write("clear_session_draft_input_if_matches", move |connection| {
connection.execute(
"UPDATE sessions SET draft_input = '' WHERE session_id = ?1 AND draft_input = ?2",
params![session_id, expected],
)?;
Ok(())
})
}
pub fn record_recovery_success(
session_id: &str,
native_session_id: &str,
checkpoint: &CheckpointMetadata,
) -> Result<()> {
let session_id = session_id.to_owned();
let native_session_id = native_session_id.to_owned();
let checkpoint = checkpoint.clone();
submit_database_write("record_recovery_success", move |_| {
record_recovery_success_to(
&database_path(),
&session_id,
&native_session_id,
&checkpoint,
)
})
}
pub(super) fn record_recovery_success_to(
path: &Path,
session_id: &str,
native_session_id: &str,
checkpoint: &CheckpointMetadata,
) -> Result<()> {
let mut connection = open(path)?;
let tx = connection.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
let changed = tx.execute(
"UPDATE sessions
SET native_session_id = ?2, last_checkpoint_error = NULL
WHERE session_id = ?1",
params![session_id, native_session_id],
)?;
if changed != 1 {
bail!("unknown session {session_id}");
}
tx.execute(
"INSERT INTO session_checkpoints(
session_id, archive_path, sha256, created_at, event_frontier
) VALUES (?1,?2,?3,?4,?5)
ON CONFLICT(session_id) DO UPDATE SET
archive_path = excluded.archive_path,
sha256 = excluded.sha256,
created_at = excluded.created_at,
event_frontier = excluded.event_frontier",
params![
session_id,
path_to_blob(&checkpoint.archive_path),
checkpoint.sha256,
checkpoint.created_at,
checkpoint.event_frontier,
],
)?;
tx.commit()?;
Ok(())
}
pub fn record_recovery_success_if_current(
session_id: &str,
expected_target: &TargetLocator,
expected_checkpoint: Option<&CheckpointMetadata>,
native_session_id: &str,
checkpoint: &CheckpointMetadata,
) -> Result<bool> {
let session_id = session_id.to_owned();
let expected_target = expected_target.clone();
let expected_checkpoint = expected_checkpoint.cloned();
let native_session_id = native_session_id.to_owned();
let checkpoint = checkpoint.clone();
submit_database_write("record_recovery_success_if_current", move |connection| {
record_recovery_success_if_current_with(
connection,
&session_id,
&expected_target,
expected_checkpoint.as_ref(),
&native_session_id,
&checkpoint,
)
})
}
fn record_recovery_success_if_current_with(
connection: &mut Connection,
session_id: &str,
expected_target: &TargetLocator,
expected_checkpoint: Option<&CheckpointMetadata>,
native_session_id: &str,
checkpoint: &CheckpointMetadata,
) -> Result<bool> {
let tx = connection.transaction()?;
let Some(mut current) = load_session_with(&tx, session_id)? else {
return Ok(false);
};
if current.target.as_ref() != Some(expected_target)
|| current.checkpoint.as_ref() != expected_checkpoint
{
return Ok(false);
}
tx.execute(
"UPDATE sessions SET native_session_id = ?2, last_checkpoint_error = NULL
WHERE session_id = ?1",
params![session_id, native_session_id],
)?;
current.checkpoint = Some(checkpoint.clone());
replace_checkpoint(&tx, ¤t)?;
tx.commit()?;
Ok(true)
}
pub fn record_recovery_failure(session_id: &str, detail: &str) -> Result<()> {
let session_id = session_id.to_owned();
let detail = detail.to_owned();
submit_database_write("record_recovery_failure", move |_| {
record_recovery_failure_to(&database_path(), &session_id, &detail)
})
}
pub fn record_recovery_failure_if_current(
session_id: &str,
expected_target: &TargetLocator,
expected_checkpoint: Option<&CheckpointMetadata>,
detail: &str,
) -> Result<bool> {
let session_id = session_id.to_owned();
let expected_target = expected_target.clone();
let expected_checkpoint = expected_checkpoint.cloned();
let detail = detail.to_owned();
submit_database_write("record_recovery_failure_if_current", move |connection| {
record_recovery_failure_if_current_with(
connection,
&session_id,
&expected_target,
expected_checkpoint.as_ref(),
&detail,
)
})
}
fn record_recovery_failure_if_current_with(
connection: &mut Connection,
session_id: &str,
expected_target: &TargetLocator,
expected_checkpoint: Option<&CheckpointMetadata>,
detail: &str,
) -> Result<bool> {
let tx = connection.transaction()?;
let Some(current) = load_session_with(&tx, session_id)? else {
return Ok(false);
};
if current.target.as_ref() != Some(expected_target)
|| current.checkpoint.as_ref() != expected_checkpoint
{
return Ok(false);
}
tx.execute(
"UPDATE sessions SET last_checkpoint_error = ?2 WHERE session_id = ?1",
params![session_id, detail],
)?;
tx.commit()?;
Ok(true)
}
pub(super) fn record_recovery_failure_to(
path: &Path,
session_id: &str,
detail: &str,
) -> Result<()> {
let connection = open(path)?;
let changed = connection.execute(
"UPDATE sessions SET last_checkpoint_error = ?2 WHERE session_id = ?1",
params![session_id, detail],
)?;
if changed != 1 {
bail!("unknown session {session_id}");
}
Ok(())
}
pub fn rebind_session_bundle(session_id: &str, bundle_id: &str) -> Result<()> {
let session_id = session_id.to_owned();
let bundle_id = bundle_id.to_owned();
submit_database_write("rebind_session_bundle", move |_| {
rebind_session_bundle_to(&database_path(), &session_id, &bundle_id)
})
}
pub(super) fn rebind_session_bundle_to(
path: &Path,
session_id: &str,
bundle_id: &str,
) -> Result<()> {
let mut connection = open(path)?;
let tx = connection.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
let changed = tx.execute(
"UPDATE session_contexts SET bundle_id = ?2 WHERE session_id = ?1",
params![session_id, bundle_id],
)?;
if changed == 0 {
tx.execute(
"INSERT INTO session_contexts(session_id, bundle_id, created_at) VALUES (?1, ?2, ?3)",
params![session_id, bundle_id, Utc::now().to_rfc3339()],
)?;
}
tx.commit()?;
Ok(())
}
#[cfg(test)]
mod ownership_tests {
use super::*;
#[test]
fn lifecycle_identity_survives_close_but_rotates_before_resume() {
let directory = tempfile::tempdir().unwrap();
let path = directory.path().join("incarnation.sqlite3");
let mut session = super::super::tests::session("incarnation", "project");
session.state = SessionState::Running;
save_session_to(&path, &session).unwrap();
let connection = open(&path).unwrap();
let identity = || {
connection
.query_row(
"SELECT identity FROM session_incarnations WHERE session_id = ?1",
[&session.id],
|row| row.get::<_, String>(0),
)
.unwrap()
};
let original = identity();
for state in ["disconnected", "running", "closing", "stopped"] {
connection
.execute(
"UPDATE sessions SET state = ?2 WHERE session_id = ?1",
params![session.id, state],
)
.unwrap();
assert_eq!(
identity(),
original,
"close must keep its execution identity"
);
}
connection
.execute(
"UPDATE sessions SET state = 'provisioning' WHERE session_id = ?1",
[&session.id],
)
.unwrap();
let resumed = identity();
assert_ne!(resumed, original);
connection
.execute(
"UPDATE sessions SET state = 'running' WHERE session_id = ?1",
[&session.id],
)
.unwrap();
assert_eq!(
identity(),
resumed,
"publishing a resume keeps its admitted identity"
);
}
#[test]
fn resume_and_rollback_preserve_edits_committed_after_their_snapshot() {
let directory = tempfile::tempdir().unwrap();
let path = directory.path().join("resume.sqlite3");
let original = super::super::tests::session("resume-owner", "project");
save_session_to(&path, &original).unwrap();
let connection = open(&path).unwrap();
connection
.execute(
"UPDATE sessions SET session_title_override = 'renamed during resume',
acp_session_title = 'new worker title', draft_input = 'new draft',
archived = 1, viewed_through_event_ordinal = 99 WHERE session_id = ?1",
[&original.id],
)
.unwrap();
let mut provisioning = original.clone();
provisioning.state = SessionState::Provisioning;
provisioning.target = None;
provisioning.native_session_id = Some("resumed-native".into());
provisioning.additional_mounts.clear();
save_resumed_session_to(&path, &provisioning, None).unwrap();
let current = load_session_with(&connection, &original.id)
.unwrap()
.unwrap();
assert_eq!(current.state, SessionState::Provisioning);
assert_eq!(current.native_session_id, provisioning.native_session_id);
assert!(current.additional_mounts.is_empty());
assert_eq!(
current.session_title_override.as_deref(),
Some("renamed during resume")
);
assert_eq!(
current.acp_session_title.as_deref(),
Some("new worker title")
);
assert_eq!(current.draft_input, "new draft");
assert!(current.archived);
assert_eq!(current.viewed_through_event_ordinal, 99);
save_resumed_session_to(&path, &original, None).unwrap();
let rolled_back = load_session_with(&connection, &original.id)
.unwrap()
.unwrap();
assert_eq!(rolled_back.state, original.state);
assert_eq!(rolled_back.additional_mounts, original.additional_mounts);
assert_eq!(
rolled_back.session_title_override,
current.session_title_override
);
assert_eq!(rolled_back.acp_session_title, current.acp_session_title);
assert_eq!(rolled_back.draft_input, current.draft_input);
assert_eq!(rolled_back.archived, current.archived);
assert_eq!(rolled_back.viewed_through_event_ordinal, 99);
}
#[test]
fn recovery_completion_cannot_replace_a_newer_checkpoint_or_target() {
let directory = tempfile::tempdir().unwrap();
let path = directory.path().join("recovery.sqlite3");
let original = super::super::tests::session("recovery-owner", "project");
save_session_to(&path, &original).unwrap();
let target = original.target.as_ref().unwrap();
let previous = original.checkpoint.as_ref().unwrap();
let mut next = previous.clone();
next.sha256 = "b".repeat(64);
next.event_frontier += 1;
let mut connection = open(&path).unwrap();
assert!(
record_recovery_success_if_current_with(
&mut connection,
&original.id,
target,
Some(previous),
"new-native",
&next,
)
.unwrap()
);
assert!(
!record_recovery_failure_if_current_with(
&mut connection,
&original.id,
target,
Some(previous),
"old failure",
)
.unwrap()
);
assert!(
!record_recovery_success_if_current_with(
&mut connection,
&original.id,
target,
Some(previous),
"old-native",
previous,
)
.unwrap()
);
let current = load_session_with(&connection, &original.id)
.unwrap()
.unwrap();
assert_eq!(current.checkpoint.as_ref(), Some(&next));
assert_eq!(current.native_session_id.as_deref(), Some("new-native"));
assert_eq!(current.last_checkpoint_error, None);
let mut replacement = current.clone();
replacement.target = None;
replacement.state = SessionState::Stopped;
save_resumed_session_to(&path, &replacement, None).unwrap();
assert!(
!record_recovery_failure_if_current_with(
&mut connection,
&original.id,
target,
Some(&next),
"old target failure",
)
.unwrap()
);
assert!(
!record_recovery_success_if_current_with(
&mut connection,
&original.id,
target,
Some(&next),
"old-native",
previous,
)
.unwrap()
);
}
#[test]
fn current_recovery_failure_settles_without_resurrecting_a_removed_session() {
let directory = tempfile::tempdir().unwrap();
let path = directory.path().join("recovery.sqlite3");
let original = super::super::tests::session("recovery-owner", "project");
save_session_to(&path, &original).unwrap();
let mut connection = open(&path).unwrap();
assert!(
record_recovery_failure_if_current_with(
&mut connection,
&original.id,
original.target.as_ref().unwrap(),
original.checkpoint.as_ref(),
"current failure",
)
.unwrap()
);
assert_eq!(
load_session_with(&connection, &original.id)
.unwrap()
.unwrap()
.last_checkpoint_error
.as_deref(),
Some("current failure")
);
connection
.execute("DELETE FROM sessions WHERE session_id = ?1", [&original.id])
.unwrap();
assert!(
!record_recovery_failure_if_current_with(
&mut connection,
&original.id,
original.target.as_ref().unwrap(),
original.checkpoint.as_ref(),
"late failure",
)
.unwrap()
);
assert!(
load_session_with(&connection, &original.id)
.unwrap()
.is_none()
);
assert!(save_resumed_session_to(&path, &original, None).is_err());
}
}