use super::*;
use mj_core::subagent::{SubagentToolRequest, SubagentToolResult};
#[derive(Clone, Debug, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
pub(crate) struct StoredDelegationResult {
pub result: SubagentToolResult,
#[serde(default)]
pub reported_finishes: Vec<(String, mj_core::subagent::SubagentFinishIdentity)>,
}
fn decode_stored_result(body: &str) -> Result<StoredDelegationResult> {
let value: serde_json::Value = serde_json::from_str(body)?;
if value.get("result").is_some() {
Ok(serde_json::from_value(value)?)
} else {
Ok(StoredDelegationResult {
result: serde_json::from_value(value)?,
reported_finishes: Vec::new(),
})
}
}
pub(crate) fn list_subagent_parent_ids() -> Result<Vec<String>> {
let connection = open_reader(&database_path())?;
let mut statement = connection.prepare(
"SELECT DISTINCT parent_session_id FROM subagent_sessions ORDER BY parent_session_id",
)?;
Ok(statement
.query_map([], |row| row.get(0))?
.collect::<std::result::Result<Vec<_>, _>>()?)
}
#[derive(Clone, Debug, serde::Serialize, serde::Deserialize)]
pub(crate) struct PreparedDelegation {
pub request: SubagentToolRequest,
pub turn_target: Option<String>,
#[serde(default)]
pub spawn: Option<PreparedSpawn>,
}
#[derive(Clone, Debug, serde::Serialize, serde::Deserialize)]
pub(crate) struct PreparedSpawn {
pub profile_id: String,
pub model: String,
pub effort: Option<String>,
pub fast_mode: bool,
pub prompt: String,
}
pub(crate) fn prepare_delegation_spawn(
parent: String,
request: String,
spawn: PreparedSpawn,
) -> Result<PreparedSpawn> {
submit_database_write("prepare delegation spawn", move |connection| {
let stored: String = connection.query_row("SELECT prepared_json FROM delegation_effects WHERE parent_session_id=?1 AND request_id=?2", params![parent, request], |row| row.get(0))?;
let mut prepared: PreparedDelegation = serde_json::from_str(&stored)?;
if let Some(existing) = prepared.spawn {
return Ok(existing);
}
prepared.spawn = Some(spawn.clone());
connection.execute("UPDATE delegation_effects SET prepared_json=?3 WHERE parent_session_id=?1 AND request_id=?2", params![parent, request, serde_json::to_string(&prepared)?])?;
Ok(spawn)
})
}
pub(crate) fn load_delegation(
parent: &str,
request: &str,
) -> Result<Option<(PreparedDelegation, Option<StoredDelegationResult>)>> {
let connection = open_reader(&database_path())?;
let row = connection.query_row("SELECT prepared_json, result_json FROM delegation_effects WHERE parent_session_id=?1 AND request_id=?2", params![parent, request], |row| Ok((row.get::<_, String>(0)?, row.get::<_, Option<String>>(1)?))).optional()?;
row.map(|(prepared, result)| {
Ok((
serde_json::from_str(&prepared)?,
result.map(|s| decode_stored_result(&s)).transpose()?,
))
})
.transpose()
}
pub(crate) fn has_pending_mutating_delegations(parent: &str) -> Result<bool> {
let connection = open_reader(&database_path())?;
let mut statement = connection.prepare(
"SELECT prepared_json FROM delegation_effects
WHERE parent_session_id=?1 AND result_json IS NULL",
)?;
let rows = statement.query_map([parent], |row| row.get::<_, String>(0))?;
for row in rows {
let prepared: PreparedDelegation = serde_json::from_str(&row?)?;
if prepared.request.action.mutates_child_state() {
return Ok(true);
}
}
Ok(false)
}
pub(crate) fn prepare_delegation(
parent: String,
prepared: PreparedDelegation,
) -> Result<PreparedDelegation> {
submit_database_write("prepare delegation", move |connection| {
let transaction = connection.transaction()?;
transaction.execute("INSERT INTO delegation_effects(parent_session_id,request_id,phase,prepared_json,result_json) VALUES (?1,?2,'prepared',?3,NULL) ON CONFLICT(parent_session_id,request_id) DO NOTHING", params![parent, prepared.request.request_id, serde_json::to_string(&prepared)?])?;
let stored: String = transaction.query_row("SELECT prepared_json FROM delegation_effects WHERE parent_session_id=?1 AND request_id=?2", params![parent, prepared.request.request_id], |row| row.get(0))?;
transaction.commit()?;
Ok(serde_json::from_str(&stored)?)
})
}
pub(crate) fn delegation_delivering(parent: String, request: String) -> Result<()> {
submit_database_write("deliver delegation effect", move |connection| {
connection.execute("UPDATE delegation_effects SET phase='delivering' WHERE parent_session_id=?1 AND request_id=?2 AND result_json IS NULL", params![parent, request])?;
Ok(())
})
}
#[cfg(test)]
pub(crate) fn record_delegation_result(parent: String, result: SubagentToolResult) -> Result<()> {
record_delegation_result_with_reports(
parent,
StoredDelegationResult {
result,
reported_finishes: Vec::new(),
},
)
}
pub(crate) fn record_delegation_result_with_reports(
parent: String,
stored_result: StoredDelegationResult,
) -> Result<()> {
submit_database_write("record delegation result", move |connection| {
let result = stored_result.result;
let transaction =
connection.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
let body = serde_json::to_string(&StoredDelegationResult {
result: result.clone(),
reported_finishes: stored_result.reported_finishes.clone(),
})?;
let changed = transaction.execute("UPDATE delegation_effects SET phase='result', result_json=?3 WHERE parent_session_id=?1 AND request_id=?2 AND result_json IS NULL", params![parent, result.request_id, body])?;
ensure!(
changed == 1,
"delegation result has no pending durable effect"
);
for (child_id, identity) in stored_result.reported_finishes {
let stored: String = transaction.query_row(
"SELECT record_json FROM subagent_sessions WHERE child_session_id=?1",
[&child_id],
|row| row.get(0),
)?;
let mut relation: mj_core::subagent::SubagentRecord = serde_json::from_str(&stored)?;
ensure!(
relation.parent_session_id == parent,
"reported child {child_id} does not belong to parent {parent}"
);
relation.reported_finish = Some(identity);
transaction.execute(
"UPDATE subagent_sessions SET record_json=?2 WHERE child_session_id=?1",
params![child_id, serde_json::to_string(&relation)?],
)?;
}
transaction.commit()?;
Ok(())
})
}
pub(crate) fn unreport_delegation_finishes(
parent: String,
reported_finishes: Vec<(String, mj_core::subagent::SubagentFinishIdentity)>,
) -> Result<()> {
submit_database_write(
"restore undelivered delegation reports",
move |connection| {
let transaction =
connection.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
for (child_id, identity) in reported_finishes {
let stored: Option<String> = transaction
.query_row(
"SELECT record_json FROM subagent_sessions WHERE child_session_id=?1 AND parent_session_id=?2",
params![child_id, parent],
|row| row.get(0),
)
.optional()?;
let Some(stored) = stored else { continue };
let mut relation: mj_core::subagent::SubagentRecord =
serde_json::from_str(&stored)?;
if relation.reported_finish.as_ref() != Some(&identity) {
continue;
}
relation.reported_finish = None;
transaction.execute(
"UPDATE subagent_sessions SET record_json=?2 WHERE child_session_id=?1 AND parent_session_id=?3",
params![child_id, serde_json::to_string(&relation)?, parent],
)?;
}
transaction.commit()?;
Ok(())
},
)
}
pub(crate) fn acknowledge_delegation(parent: String, request: String) -> Result<()> {
submit_database_write("acknowledge delegation", move |connection| {
connection.execute("DELETE FROM delegation_effects WHERE parent_session_id=?1 AND request_id=?2 AND result_json IS NOT NULL", params![parent, request])?;
Ok(())
})
}
pub(crate) fn prune_acknowledged_delegations() -> Result<()> {
submit_database_write("prune acknowledged delegations", move |connection| {
connection.execute(
"DELETE FROM delegation_effects WHERE phase='acknowledged'",
[],
)?;
Ok(())
})
}
#[cfg(test)]
mod tests {
use super::*;
use crate::controller::test_support::{IsolatedTest, test_name};
#[test]
fn restart_preserves_selected_turn_and_result_until_acknowledged() {
const CHILD: &str = "MJ_TEST_DURABLE_DELEGATION";
if std::env::var_os(CHILD).is_none() {
let root = tempfile::tempdir().unwrap();
IsolatedTest::new(test_name(
module_path!(),
"restart_preserves_selected_turn_and_result_until_acknowledged",
))
.env(CHILD, "1")
.env("MJ_INSTANCE", "concurrency-sweep-delegation")
.isolated_store(root.path())
.run();
return;
}
let writer = install_isolated_test_writer();
let prepared = PreparedDelegation {
request: SubagentToolRequest {
originating_command_id: None,
request_id: "request".into(),
created_at_ms: 1,
action: mj_core::subagent::SubagentToolAction::InterruptAgent {
child_session_id: "child".into(),
},
},
turn_target: Some("original-turn".into()),
spawn: None,
};
prepare_delegation("parent".into(), prepared.clone()).unwrap();
delegation_delivering("parent".into(), "request".into()).unwrap();
assert!(has_pending_mutating_delegations("parent").unwrap());
drop(writer);
let _writer = install_isolated_test_writer();
let mut replacement = prepared;
replacement.turn_target = Some("later-turn".into());
let restored = prepare_delegation("parent".into(), replacement).unwrap();
assert_eq!(restored.turn_target.as_deref(), Some("original-turn"));
let result = SubagentToolResult {
request_id: "request".into(),
completed_at_ms: 2,
is_error: false,
message: "original result".into(),
};
record_delegation_result("parent".into(), result).unwrap();
assert!(!has_pending_mutating_delegations("parent").unwrap());
let (restored, result) = load_delegation("parent", "request").unwrap().unwrap();
assert_eq!(restored.turn_target.as_deref(), Some("original-turn"));
assert_eq!(result.unwrap().result.message, "original result");
acknowledge_delegation("parent".into(), "request".into()).unwrap();
assert!(load_delegation("parent", "request").unwrap().is_none());
prepare_delegation("parent".into(), restored.clone()).unwrap();
submit_database_write("mark legacy acknowledged", |connection| {
connection.execute("UPDATE delegation_effects SET phase='acknowledged'", [])?;
Ok(())
})
.unwrap();
prune_acknowledged_delegations().unwrap();
assert!(load_delegation("parent", "request").unwrap().is_none());
}
#[test]
fn legacy_raw_delegation_results_load_without_report_markers() {
let result = SubagentToolResult {
request_id: "legacy".into(),
completed_at_ms: 2,
is_error: false,
message: "old result shape".into(),
};
let restored = decode_stored_result(&serde_json::to_string(&result).unwrap()).unwrap();
assert_eq!(restored.result, result);
assert!(restored.reported_finishes.is_empty());
}
}