use std::collections::BTreeMap;
use std::path::{Component, Path, PathBuf};
use std::sync::{Arc, Mutex};
use std::time::Duration;
use agent_client_protocol::schema::v1::{ContentBlock, TextContent};
use anyhow::{Context, Result, anyhow, bail, ensure};
use mj_core::state::{MaterializedExecutionState, SessionState};
use mj_core::subagent::{DEFAULT_WAIT_SECONDS, ReportState, bounded_report};
use crate::quota::ProfileQuota;
use crate::controller::{Controller, SessionExportLayout};
use crate::server::api::{
BundleExport, ExportError, PushedBranch, StartFollowup, StartStatus, SubagentBackend,
TranscriptPage, TurnSpan, TurnState, TurnSummary,
};
use crate::targets::{self, CancellableProcessExecutor, CommandExecutor, CommandOutput};
use mj_client::session::{BoxFuture, SessionControl, SessionHandle, new_command_id};
use mj_core::relay::RelayCommand;
mod subagent_input;
use crate::daemon::RuntimeState;
use mj_client::daemon::{WikiHitTranscript, WikiRestoreRequest, WikiSearchPage, WikiSessionInfo};
pub type SessionStateSource = Arc<dyn Fn(&str) -> Option<SessionState> + Send + Sync>;
pub trait ExportRuntime: Send + Sync {
fn queue_startup(
self: Arc<Self>,
_session_id: String,
_group_id: String,
_followup: StartFollowup,
) -> BoxFuture<'static, Result<()>> {
Box::pin(async { bail!("durable startup delivery is unavailable") })
}
fn startup_status(&self, _session_id: String) -> BoxFuture<'_, Result<Option<StartStatus>>> {
Box::pin(async { Ok(None) })
}
fn startup_context(
&self,
session_id: String,
) -> BoxFuture<'_, Result<(Option<String>, Option<StartStatus>)>> {
Box::pin(async move { Ok((None, self.startup_status(session_id).await?)) })
}
fn dismiss_startup_status(
&self,
_session_id: String,
_group_id: Option<String>,
) -> BoxFuture<'_, Result<()>> {
Box::pin(async { Ok(()) })
}
fn cancel_startup(self: Arc<Self>, _session_id: String) -> BoxFuture<'static, Result<()>> {
Box::pin(async { Ok(()) })
}
fn session_record(&self, session_id: &str) -> Option<mj_core::state::SessionRecord>;
fn close_is_requested(&self, _session_id: &str) -> bool {
false
}
fn workspace_session(
&self,
_session_id: String,
) -> BoxFuture<'_, Result<crate::session_manager::ManagedSessionHandle>> {
Box::pin(async { anyhow::bail!("workspace file injection is unavailable") })
}
fn checkpoint_now(
&self,
session_id: String,
) -> BoxFuture<'_, Result<mj_core::state::CheckpointMetadata>>;
fn spawn_subagent(
self: Arc<Self>,
_request: crate::controller::RegisterSubagentRequest,
) -> BoxFuture<'static, Result<mj_core::subagent::SubagentRecord>> {
Box::pin(async { anyhow::bail!("sub-agent creation is unavailable") })
}
fn close_subagent_request(
self: Arc<Self>,
_session_id: String,
_parent: String,
_request: String,
_expected_incarnation: String,
) -> BoxFuture<'static, Result<()>> {
Box::pin(async { anyhow::bail!("durable sub-agent close admission is unavailable") })
}
fn close_subagent(self: Arc<Self>, _session_id: String) -> BoxFuture<'static, Result<()>> {
Box::pin(async { anyhow::bail!("sub-agent close is unavailable") })
}
fn subagent_park_running(&self, _session_id: &str) -> bool {
false
}
fn park_subagent(
self: Arc<Self>,
_session_id: String,
) -> BoxFuture<'static, Result<crate::controller::ParkOutcome>> {
Box::pin(async { anyhow::bail!("parking a sub-agent is unavailable") })
}
fn unpark_subagent(self: Arc<Self>, _session_id: String) -> BoxFuture<'static, Result<()>> {
Box::pin(async { anyhow::bail!("restarting a parked sub-agent is unavailable") })
}
fn refresh_workspaces(&self) -> BoxFuture<'_, Result<()>> {
Box::pin(async { Ok(()) })
}
fn wiki_sync_is_stale(&self) -> bool {
false
}
fn wiki_request_sync(&self) {}
fn wiki_search(&self, _query: String, _limit: usize) -> BoxFuture<'_, Result<WikiSearchPage>> {
Box::pin(async { anyhow::bail!("SessionWiki search is unavailable") })
}
fn wiki_brief(
&self,
_wiki_id: String,
_max_chars: usize,
) -> BoxFuture<'_, Result<Option<String>>> {
Box::pin(async { anyhow::bail!("SessionWiki briefings are unavailable") })
}
fn wiki_hits(
&self,
_wiki_id: String,
_query: String,
_context_messages: usize,
_per_message_chars: usize,
) -> BoxFuture<'_, Result<Option<WikiHitTranscript>>> {
Box::pin(async { anyhow::bail!("SessionWiki transcript hits are unavailable") })
}
fn wiki_session(&self, _wiki_id: String) -> BoxFuture<'_, Result<Option<WikiSessionInfo>>> {
Box::pin(async { anyhow::bail!("SessionWiki lookups are unavailable") })
}
fn wiki_restore(
self: Arc<Self>,
_request: WikiRestoreRequest,
) -> BoxFuture<'static, Result<Option<String>>> {
Box::pin(async { anyhow::bail!("SessionWiki restore is unavailable") })
}
}
impl ExportRuntime for RuntimeState {
fn queue_startup(
self: Arc<Self>,
session_id: String,
group_id: String,
followup: StartFollowup,
) -> BoxFuture<'static, Result<()>> {
Box::pin(async move {
self.queue_api_followup(&session_id, group_id, followup)
.await
})
}
fn startup_status(&self, session_id: String) -> BoxFuture<'_, Result<Option<StartStatus>>> {
Box::pin(load_startup_status(session_id))
}
fn startup_context(
&self,
session_id: String,
) -> BoxFuture<'_, Result<(Option<String>, Option<StartStatus>)>> {
Box::pin(load_startup_context(session_id))
}
fn dismiss_startup_status(
&self,
session_id: String,
group_id: Option<String>,
) -> BoxFuture<'_, Result<()>> {
Box::pin(async move {
if let Some(group_id) = group_id {
blocking("dismiss completed startup status", move || {
crate::database::dismiss_startup_group(&session_id, &group_id)
})
.await?;
}
Ok(())
})
}
fn cancel_startup(self: Arc<Self>, session_id: String) -> BoxFuture<'static, Result<()>> {
Box::pin(async move { self.cancel_api_followup(&session_id).await })
}
fn session_record(&self, session_id: &str) -> Option<mj_core::state::SessionRecord> {
RuntimeState::session_record(self, session_id)
}
fn close_is_requested(&self, session_id: &str) -> bool {
RuntimeState::close_is_requested(self, session_id)
}
fn workspace_session(
&self,
session_id: String,
) -> BoxFuture<'_, Result<crate::session_manager::ManagedSessionHandle>> {
Box::pin(async move { self.workspace_session_handle(&session_id).await })
}
fn checkpoint_now(
&self,
session_id: String,
) -> BoxFuture<'_, Result<mj_core::state::CheckpointMetadata>> {
Box::pin(async move { self.checkpoint_session_now(&session_id).await })
}
fn spawn_subagent(
self: Arc<Self>,
request: crate::controller::RegisterSubagentRequest,
) -> BoxFuture<'static, Result<mj_core::subagent::SubagentRecord>> {
Box::pin(async move { self.start_subagent_session(request).await })
}
fn close_subagent_request(
self: Arc<Self>,
session_id: String,
parent: String,
request: String,
expected_incarnation: String,
) -> BoxFuture<'static, Result<()>> {
Box::pin(async move {
RuntimeState::close_subagent_request(
&self,
session_id,
parent,
request,
expected_incarnation,
)
.await
})
}
fn close_subagent(self: Arc<Self>, session_id: String) -> BoxFuture<'static, Result<()>> {
Box::pin(async move { self.suspend_session(session_id).await })
}
fn subagent_park_running(&self, session_id: &str) -> bool {
RuntimeState::subagent_park_running(self, session_id)
}
fn park_subagent(
self: Arc<Self>,
session_id: String,
) -> BoxFuture<'static, Result<crate::controller::ParkOutcome>> {
Box::pin(async move { RuntimeState::park_subagent(&self, session_id).await })
}
fn unpark_subagent(self: Arc<Self>, session_id: String) -> BoxFuture<'static, Result<()>> {
Box::pin(async move { RuntimeState::unpark_subagent(&self, session_id).await })
}
fn refresh_workspaces(&self) -> BoxFuture<'_, Result<()>> {
Box::pin(RuntimeState::refresh_workspaces(self))
}
fn wiki_sync_is_stale(&self) -> bool {
crate::sessionwiki::sync_is_stale(self.wiki().last_success())
}
fn wiki_request_sync(&self) {
self.wiki().request_sync(false);
}
fn wiki_search(&self, query: String, limit: usize) -> BoxFuture<'_, Result<WikiSearchPage>> {
Box::pin(async move { RuntimeState::wiki_search(self, query, limit).await })
}
fn wiki_brief(
&self,
wiki_id: String,
max_chars: usize,
) -> BoxFuture<'_, Result<Option<String>>> {
Box::pin(async move { RuntimeState::wiki_brief(self, wiki_id, max_chars).await })
}
fn wiki_hits(
&self,
wiki_id: String,
query: String,
context_messages: usize,
per_message_chars: usize,
) -> BoxFuture<'_, Result<Option<WikiHitTranscript>>> {
Box::pin(async move {
RuntimeState::wiki_hits(self, wiki_id, query, context_messages, per_message_chars).await
})
}
fn wiki_session(&self, wiki_id: String) -> BoxFuture<'_, Result<Option<WikiSessionInfo>>> {
Box::pin(async move { RuntimeState::wiki_session(self, wiki_id).await })
}
fn wiki_restore(
self: Arc<Self>,
request: WikiRestoreRequest,
) -> BoxFuture<'static, Result<Option<String>>> {
Box::pin(async move {
Ok(self
.restore_wiki_session(request, &tokio_util::sync::CancellationToken::new())
.await?
.map(|registered| registered.session.id))
})
}
}
const START_POLL: Duration = Duration::from_secs(5);
const EXPORT_TIMEOUT: Duration = Duration::from_secs(5 * 60);
const UNPARK_ATTACH_TIMEOUT: Duration = Duration::from_secs(60);
const CLAP_USAGE_EXIT_CODE: i32 = 2;
pub struct ApiBackend {
sessions: SessionControl,
exports: Arc<dyn ExportRuntime>,
quota_reports: Arc<Mutex<BTreeMap<String, ProfileQuota>>>,
rejected_logins: Arc<Mutex<mj_core::credentials::RejectedLogins>>,
profile_catalog: Arc<super::profile_catalog::ProfileCatalog>,
}
impl ApiBackend {
pub(crate) async fn start_followup_with_id(
&self,
session_id: String,
followup: StartFollowup,
group_id: String,
) -> Result<()> {
if followup == StartFollowup::default() {
return Ok(());
}
Arc::clone(&self.exports)
.queue_startup(session_id, group_id, followup)
.await
}
pub fn new(
sessions: SessionControl,
_session_states: SessionStateSource,
exports: Arc<dyn ExportRuntime>,
) -> Self {
Self {
sessions,
exports,
quota_reports: Arc::new(Mutex::new(BTreeMap::new())),
rejected_logins: Arc::default(),
profile_catalog: super::profile_catalog::ProfileCatalog::new(
tokio_util::sync::CancellationToken::new(),
),
}
}
pub fn with_quota_reports(
mut self,
quota_reports: Arc<Mutex<BTreeMap<String, ProfileQuota>>>,
) -> Self {
self.quota_reports = quota_reports;
self
}
pub fn with_rejected_logins(
mut self,
rejected_logins: Arc<Mutex<mj_core::credentials::RejectedLogins>>,
) -> Self {
self.rejected_logins = rejected_logins;
self
}
pub fn with_profile_catalog(
mut self,
profile_catalog: Arc<super::profile_catalog::ProfileCatalog>,
) -> Self {
self.profile_catalog = profile_catalog;
self
}
pub(crate) async fn execute_subagent_tool_durable(
self: &Arc<Self>,
parent: String,
request: mj_core::subagent::SubagentToolRequest,
) -> Result<mj_core::subagent::SubagentToolResult> {
let stored = blocking("load delegation effect", {
let parent = parent.clone();
let id = request.request_id.clone();
move || crate::database::load_delegation(&parent, &id)
})
.await?;
let prepared = if let Some((prepared, result)) = stored {
if let Some(result) = result {
return Ok(result);
}
prepared
} else {
let turn_session = match &request.action {
mj_core::subagent::SubagentToolAction::InterruptAgent { child_session_id } => {
Some(child_session_id.clone())
}
mj_core::subagent::SubagentToolAction::Handback { .. } => None,
_ => None,
};
let turn_target = if matches!(
request.action,
mj_core::subagent::SubagentToolAction::Handback { .. }
) {
request.originating_command_id.clone()
} else if let Some(session) = turn_session {
self.session_handle(session)
.await?
.and_then(|handle| handle.view().snapshot)
.and_then(|snapshot| snapshot.materialized.active_turn)
.map(|turn| turn.command_id)
} else {
None
};
blocking("prepare delegation effect", {
let parent = parent.clone();
move || {
crate::database::prepare_delegation(
parent,
crate::database::PreparedDelegation {
request,
turn_target,
spawn: None,
close_incarnation: None,
},
)
}
})
.await?
};
blocking("mark delegation delivering", {
let parent = parent.clone();
let id = prepared.request.request_id.clone();
move || crate::database::delegation_delivering(parent, id)
})
.await?;
let result = self
.execute_subagent_tool_prepared(
parent.clone(),
prepared.request.clone(),
Some(&prepared),
)
.await;
blocking("persist delegation result", {
let result = result.clone();
move || crate::database::record_delegation_result(parent, result)
})
.await?;
Ok(result)
}
pub(crate) async fn release_delegation_receipt(
&self,
request: &mj_core::subagent::SubagentToolRequest,
) -> Result<()> {
use mj_core::subagent::SubagentToolAction;
let (child, prefix) = match &request.action {
SubagentToolAction::SendInput {
child_session_id, ..
} => (child_session_id, "subagent-input"),
SubagentToolAction::InterruptAgent { child_session_id } => {
(child_session_id, "subagent-interrupt")
}
_ => return Ok(()),
};
if self.exports.session_record(child).is_none() {
return Ok(());
}
let result = async {
let handle = self
.session_handle(child.clone())
.await?
.context("waiting for delegation receipt owner")?;
handle
.release_command_receipt(format!("{prefix}-{}", request.request_id))
.await?;
Ok::<_, anyhow::Error>(())
}
.await;
result.context(
"delegation result is durable; retrying receipt release without repeating the effect",
)
}
#[cfg(test)]
pub async fn execute_subagent_tool(
self: &Arc<Self>,
parent_session_id: String,
request: mj_core::subagent::SubagentToolRequest,
) -> mj_core::subagent::SubagentToolResult {
self.execute_subagent_tool_prepared(parent_session_id, request, None)
.await
}
async fn execute_subagent_tool_prepared(
self: &Arc<Self>,
parent_session_id: String,
request: mj_core::subagent::SubagentToolRequest,
prepared: Option<&crate::database::PreparedDelegation>,
) -> mj_core::subagent::SubagentToolResult {
let outcome = self
.execute_subagent_tool_inner(&parent_session_id, &request, prepared)
.await;
if let mj_core::subagent::SubagentToolAction::SendInput {
child_session_id, ..
} = &request.action
{
let (mut value, is_error) = match outcome {
Ok(value) => (value, false),
Err(error) => (
serde_json::json!({"status":"failed", "error":format!("{error:#}")}),
true,
),
};
value["child_session_id"] = child_session_id.clone().into();
value["request_id"] = request.request_id.clone().into();
value["created_at_ms"] = request.created_at_ms.into();
return mj_core::subagent::SubagentToolResult {
request_id: request.request_id,
completed_at_ms: mj_core::clock::epoch_millis(),
is_error,
message: value.to_string(),
};
}
let (is_error, message) = match outcome {
Ok(value) => (
false,
serde_json::to_string_pretty(&value).unwrap_or_else(|_| value.to_string()),
),
Err(error) => (true, format!("{error:#}")),
};
mj_core::subagent::SubagentToolResult {
request_id: request.request_id,
completed_at_ms: mj_core::clock::epoch_millis(),
is_error,
message,
}
}
async fn execute_subagent_tool_inner(
self: &Arc<Self>,
parent_session_id: &str,
request: &mj_core::subagent::SubagentToolRequest,
prepared: Option<&crate::database::PreparedDelegation>,
) -> Result<serde_json::Value> {
use mj_core::subagent::SubagentToolAction;
let request_created_at_ms = request.created_at_ms;
match &request.action {
SubagentToolAction::ListProfiles => {
let parent = self
.exports
.session_record(parent_session_id)
.context("parent session disappeared")?;
anyhow::ensure!(
parent.subagents == Some(mj_core::subagent::SubagentPolicy::AllModels),
"list_profiles is unavailable for this session's subagent policy"
);
let candidates = self
.subagent_candidates(parent.last_profile.clone())
.await?;
let profiles =
crate::server::api::merge_same_models(candidates.offered, &parent.last_profile)
.into_iter()
.map(|candidate| {
serde_json::json!({
"profile_id":candidate.profile_id,
"harness":candidate.harness.id(),
"default_model":candidate.choices.model,
"models":candidate.choices.models,
"efforts":candidate.choices.efforts,
})
})
.collect::<Vec<_>>();
if !candidates.unavailable.is_empty() {
let unavailable = candidates
.unavailable
.into_iter()
.map(|(profile_id, reason)| {
serde_json::json!({"profile_id":profile_id,"reason":reason})
})
.collect::<Vec<_>>();
return Ok(serde_json::json!({"profiles":profiles,"unavailable":unavailable}));
}
Ok(serde_json::json!({"profiles":profiles}))
}
SubagentToolAction::Spawn {
task_name,
instructions,
profile_id,
model,
effort,
working_directory,
context,
files,
} => {
let selection = if let Some(selection) = prepared.and_then(|p| p.spawn.clone()) {
selection
} else {
let parent = self
.exports
.session_record(parent_session_id)
.context("parent session disappeared")?;
let backend: Arc<dyn crate::server::api::SubagentBackend> = self.clone();
let selection = crate::server::api::resolve_subagent_policy_selection(
&backend,
parent_session_id,
&parent.last_profile,
&parent.subagents.clone().unwrap_or_default(),
profile_id.as_deref(),
model.as_deref(),
effort.as_deref(),
)
.await
.map_err(|failure| anyhow::anyhow!(failure.message))?;
let ranges = files
.iter()
.flat_map(|entry| {
let file = entry.file.clone();
entry.ranges.iter().map(move |range| {
crate::server::api::SubagentSourceRange {
file: file.clone(),
start: range.start,
end: range.end,
}
})
})
.collect::<Vec<_>>();
let prompt = crate::server::api::build_subagent_prompt(
&backend,
parent_session_id,
instructions,
context.as_deref(),
&ranges,
)
.await
.map_err(|failure| anyhow::anyhow!(failure.message))?;
let selected = crate::database::PreparedSpawn {
profile_id: selection.profile_id,
model: selection.model,
effort: selection.effort,
fast_mode: selection.fast_mode,
prompt,
};
if prepared.is_some() {
blocking("persist selected spawn", {
let parent = parent_session_id.to_owned();
let request = request.request_id.clone();
move || {
crate::database::prepare_delegation_spawn(parent, request, selected)
}
})
.await?
} else {
selected
}
};
let relation = self
.start_subagent(crate::controller::RegisterSubagentRequest {
parent_session_id: parent_session_id.to_owned(),
task_name: task_name.clone(),
profile_id: selection.profile_id,
model: Some(selection.model.clone()),
effort: selection.effort.clone(),
working_directory: working_directory.clone(),
initial_prompt: selection.prompt,
request_key: request.request_id.clone(),
report_root: None,
})
.await?;
self.start_followup_with_id(
relation.child_session_id.clone(),
crate::server::api::StartFollowup {
model: Some(selection.model),
effort: selection.effort,
prompt: Some(relation.initial_prompt.clone()),
fast_mode: selection.fast_mode,
},
format!("subagent-spawn-{}", request.request_id),
)
.await?;
let report_dir = blocking("load sub-agent report directory", {
let child_id = relation.child_session_id.clone();
move || crate::database::load_subagent_report(&child_id)
})
.await?
.report_dir;
Ok(serde_json::json!({
"child_session_id":relation.child_session_id,
"task_name":relation.task_name,
"profile_id":relation.profile_id,
"model":relation.model,
"report_dir":report_dir,
}))
}
SubagentToolAction::ListAgents => {
let inputs = self.subagent_input_progress(parent_session_id).await?;
let relations = self.list_subagents(parent_session_id.to_owned()).await?;
let child_ids = relations
.iter()
.map(|relation| relation.child_session_id.clone())
.collect::<Vec<_>>();
let summaries = tokio::task::spawn_blocking(move || {
child_ids
.into_iter()
.map(|id| {
let summary = crate::database::load_materialized_session_summary(&id)?;
let progress = load_child_progress(&id)?;
Ok((id, (summary, progress)))
})
.collect::<Result<std::collections::BTreeMap<_, _>>>()
})
.await??;
let mut starts = std::collections::BTreeMap::new();
for relation in &relations {
if let Some(status) =
self.start_status(relation.child_session_id.clone()).await?
{
starts.insert(relation.child_session_id.clone(), status);
}
}
let agents = relations
.into_iter()
.map(|relation| {
let record = self.exports.session_record(&relation.child_session_id);
let unknown = ChildProgress {
report: ReportState::Fallback,
awaited_ordinal: None,
answered_ordinal: None,
failed_turn: None,
login_failure: None,
report_dir: None,
};
let (summary, progress) = summaries
.get(&relation.child_session_id)
.map_or((None, &unknown), |(summary, progress)| {
(summary.as_ref(), progress)
});
let (state, _, _) = inputs.status(
&relation.child_session_id,
subagent_status(
record.as_ref(),
summary,
starts.get(&relation.child_session_id),
None,
self.exports.close_is_requested(&relation.child_session_id),
progress,
),
);
let mut entry = serde_json::json!({
"child_session_id":relation.child_session_id,
"task_name":relation.task_name,
"profile_id":relation.profile_id,
"state":state,
});
inputs.annotate(&relation.child_session_id, &mut entry);
mark_parked(&mut entry, record.as_ref());
entry
})
.collect::<Vec<_>>();
Ok(serde_json::json!({"agents":agents}))
}
SubagentToolAction::SendInput {
child_session_id,
message,
} => {
self.require_owned_child(parent_session_id, child_session_id)
.await?;
let turn_id = self
.deliver_subagent_input(parent_session_id, child_session_id, message, request)
.await?;
blocking("record sub-agent prompt", {
let child_id = child_session_id.clone();
move || crate::database::record_subagent_prompt(&child_id, turn_id)
})
.await?;
Ok(
serde_json::json!({"child_session_id":child_session_id,"turn_id":turn_id,"status":"submitted"}),
)
}
SubagentToolAction::WaitAgents {
child_session_ids,
timeout_seconds,
return_when,
} => {
for child_id in child_session_ids {
self.require_owned_child(parent_session_id, child_id)
.await?;
}
let started = tokio::time::Instant::now();
let remaining = mj_core::subagent::remaining_subagent_wait(
request_created_at_ms,
*timeout_seconds,
mj_core::clock::epoch_millis(),
);
tracing::info!(
parent_session_id,
children = child_session_ids.len(),
requested_seconds = timeout_seconds.unwrap_or(DEFAULT_WAIT_SECONDS),
remaining_seconds = remaining.as_secs(),
"starting a sub-agent wait"
);
let deadline = started + remaining;
loop {
let inputs = self.subagent_input_progress(parent_session_id).await?;
let ids = child_session_ids.clone();
let summaries = tokio::task::spawn_blocking(move || {
ids.into_iter()
.map(|id| {
let summary =
crate::database::load_materialized_session_summary(&id)?;
let progress = load_child_progress(&id)?;
Ok((id, summary, progress))
})
.collect::<Result<Vec<_>>>()
})
.await??;
let mut starts = std::collections::BTreeMap::new();
for (id, _, _) in &summaries {
if let Some(status) = self.start_status(id.clone()).await? {
starts.insert(id.clone(), status);
}
}
let finished = summaries
.iter()
.map(|(id, summary, progress)| {
let record = self.exports.session_record(id);
inputs
.status(
id,
subagent_status(
record.as_ref(),
summary.as_ref(),
starts.get(id),
None,
self.exports.close_is_requested(id),
progress,
),
)
.2
})
.collect::<Vec<_>>();
let complete = finished.iter().all(|done| *done);
if return_when.satisfied(&finished) || tokio::time::Instant::now() >= deadline {
let ids: Vec<String> =
summaries.iter().map(|(id, _, _)| id.clone()).collect();
let reports = tokio::task::spawn_blocking(move || {
ids.into_iter()
.map(|id| {
crate::database::load_materialized_finished_turn_message(&id)
.map(|message| (id, message))
})
.collect::<Result<std::collections::BTreeMap<_, _>>>()
})
.await??;
let agents = summaries
.into_iter()
.map(|(id, summary, progress)| {
let record = self.exports.session_record(&id);
let (state, output, finished) = inputs.status(
&id,
subagent_status(
record.as_ref(),
summary.as_ref(),
starts.get(&id),
reports.get(&id).and_then(Option::as_deref),
self.exports.close_is_requested(&id),
&progress,
),
);
let mut entry =
wait_agent_entry(&id, &state, output, finished, &progress);
inputs.annotate(&id, &mut entry);
mark_parked(&mut entry, record.as_ref());
entry
})
.collect::<Vec<_>>();
let unfinished = agents
.iter()
.filter(|agent| agent["finished"] != serde_json::Value::Bool(true))
.filter_map(|agent| agent["child_session_id"].as_str())
.map(str::to_owned)
.collect::<Vec<_>>();
let waited_seconds =
(mj_core::subagent::subagent_wait_timeout(*timeout_seconds)
- remaining
+ started.elapsed())
.as_secs();
tracing::info!(
parent_session_id,
complete,
waited_seconds,
"answering a sub-agent wait"
);
return Ok(serde_json::json!({
"status": if complete {
mj_core::subagent::WAIT_STATUS_COMPLETE
} else {
mj_core::subagent::WAIT_STATUS_STILL_RUNNING
},
"waited_seconds": waited_seconds,
"agents": agents,
"next_action": mj_core::subagent::next_action(
child_session_ids,
&unfinished,
),
}));
}
tokio::time::sleep(Duration::from_millis(250)).await;
}
}
SubagentToolAction::InterruptAgent { child_session_id } => {
self.require_owned_child(parent_session_id, child_session_id)
.await?;
let handle = self.session_handle(child_session_id.clone()).await?;
let target = match prepared {
Some(prepared) => prepared.turn_target.clone(),
None => handle
.as_ref()
.and_then(|handle| handle.view().snapshot)
.and_then(|snapshot| snapshot.materialized.active_turn)
.map(|turn| turn.command_id),
};
let Some(target) = target else {
return Ok(
serde_json::json!({"child_session_id":child_session_id,"interrupted":false}),
);
};
let handle = handle.context("interrupt target worker is unavailable")?;
let result = handle
.submit_durable(
format!("subagent-interrupt-{}", request.request_id),
RelayCommand::CancelTurnFor {
active_prompt_id: target.clone(),
},
)
.await;
if let Err(error) = result {
if error
.downcast_ref::<mj_client::session::DeliveryUnconfirmed>()
.is_some()
{
return Err(error);
}
handle.sync_now().await?;
let still_active = handle
.view()
.snapshot
.and_then(|s| s.materialized.active_turn)
.is_some_and(|turn| turn.command_id == target);
if still_active {
return Err(error);
}
return Ok(
serde_json::json!({"child_session_id":child_session_id,"interrupted":false}),
);
}
Ok(serde_json::json!({"child_session_id":child_session_id,"interrupted":true}))
}
SubagentToolAction::CloseAgent { child_session_id } => {
self.require_owned_child(parent_session_id, child_session_id)
.await?;
if let Some(prepared) = prepared {
Arc::clone(&self.exports)
.close_subagent_request(
child_session_id.clone(),
parent_session_id.to_owned(),
request.request_id.clone(),
prepared.close_incarnation.clone().context(
"child incarnation was unavailable at request preparation",
)?,
)
.await?;
} else {
Arc::clone(&self.exports)
.close_subagent(child_session_id.clone())
.await?;
}
Ok(serde_json::json!({"child_session_id":child_session_id,"closed":true}))
}
SubagentToolAction::Handback { message } => {
let child_id = parent_session_id.to_owned();
ensure!(
!message.trim().is_empty(),
"a report cannot be empty; call handback with your full report"
);
ensure!(
message.chars().count() <= mj_core::subagent::MAX_HANDBACK_CHARS,
"a report can be at most {} characters and this one has {}. Write the details \
to files in the report directory named in your first prompt, then call \
handback again with a short report that lists their paths",
mj_core::subagent::MAX_HANDBACK_CHARS,
message.chars().count()
);
let record = blocking("load sub-agent record", {
let child_id = child_id.clone();
move || crate::database::load_subagent(&child_id)
})
.await?;
ensure!(
record.is_some_and(|record| record.handback_tool),
"handback is only for a Mjolnir sub-agent that was given the tool"
);
let command_id = if let Some(prepared) = prepared {
prepared
.turn_target
.clone()
.context("no turn was running when this report was prepared")?
} else {
let live_turn = self
.session_handle(child_id.clone())
.await?
.and_then(|handle| handle.view().snapshot)
.and_then(|snapshot| snapshot.materialized.active_turn);
let active_turn = match live_turn {
Some(turn) => Some(turn),
None => blocking("load child turn", {
let child_id = child_id.clone();
move || crate::database::load_materialized_turn_outcome(&child_id)
})
.await?
.and_then(|(_, active, _)| active),
};
let turn = active_turn.context(
"no turn is running, so there is nothing to report on; hand back your report \
during the turn that did the work",
)?;
turn.command_id
};
let recorded = blocking("record sub-agent report", {
let child_id = child_id.clone();
let handback = mj_core::subagent::SubagentHandback {
command_id,
message: message.clone(),
recorded_at_ms: mj_core::clock::epoch_millis(),
};
move || crate::database::record_subagent_handback(&child_id, &handback)
})
.await?;
Ok(if recorded {
serde_json::json!({
"delivered": true,
"message": "Report delivered to the session that started you.",
})
} else {
serde_json::json!({
"delivered": false,
"message": "Nothing was sent: your report for this turn was already delivered. Stop now.",
})
})
}
}
}
async fn unpark_child(&self, child_id: &str) -> Result<bool> {
let parked = self
.exports
.session_record(child_id)
.is_some_and(|record| record.state == SessionState::Parked)
|| self.exports.subagent_park_running(child_id);
if !parked {
return Ok(false);
}
Arc::clone(&self.exports)
.unpark_subagent(child_id.to_owned())
.await
.context(
"could not start the parked sub-agent again; it is still parked, so you can retry",
)?;
if self
.exports
.session_record(child_id)
.is_none_or(|record| record.state != SessionState::Running)
{
return Ok(false);
}
let deadline = tokio::time::Instant::now() + UNPARK_ATTACH_TIMEOUT;
let mut handle = self
.sessions
.wait_for_session(child_id, UNPARK_ATTACH_TIMEOUT)
.await
.context("the restarted sub-agent did not reattach")?;
loop {
let view = handle.view();
if view.connected
&& view
.snapshot
.as_ref()
.is_some_and(|snapshot| snapshot.operational.native_session_is_ready())
{
return Ok(true);
}
ensure!(
tokio::time::Instant::now() < deadline,
"the restarted sub-agent was not ready for a prompt within {} seconds",
UNPARK_ATTACH_TIMEOUT.as_secs()
);
let _ = tokio::time::timeout(START_POLL, handle.changed()).await;
}
}
async fn require_owned_child(&self, parent_id: &str, child_id: &str) -> Result<()> {
anyhow::ensure!(
self.list_subagents(parent_id.to_owned())
.await?
.iter()
.any(|child| child.child_session_id == child_id),
"session {child_id} does not belong to parent {parent_id}"
);
Ok(())
}
pub async fn record_subagent_completion_notice(
&self,
parent_session_id: String,
child_session_id: &str,
child_title: &str,
outcome: &mj_core::state::MaterializedTurnOutcome,
) -> Result<()> {
ensure!(
!matches!(
self.start_status(parent_session_id.clone()).await?,
Some(StartStatus::Pending)
),
"session initialization is still running"
);
let handle = self
.sessions
.session(parent_session_id.clone())
.await
.with_context(|| format!("session {parent_session_id} is not running"))?;
let turn = match outcome.accepted_ordinal {
Some(turn) => format!("turn {turn}"),
None => "a turn".to_owned(),
};
submit_notice(
&handle,
format!(
"Subagent \"{child_title}\" ({}) finished {turn} ({}).",
mj_core::state::short_id(child_session_id),
outcome.outcome
),
)
.await?;
Ok(())
}
pub async fn remind_subagent_to_hand_back(
&self,
child_session_id: &str,
handback_tool: bool,
last_turn: &mj_core::state::MaterializedTurnOutcome,
in_flight: &[String],
) -> Result<bool> {
if !handback_tool || !in_flight.is_empty() {
return Ok(false);
}
let report = blocking("load sub-agent report", {
let child_id = child_session_id.to_owned();
move || crate::database::load_subagent_report(&child_id)
})
.await?;
if let Some(reminder) = &report.reminder {
self.sessions
.session(child_session_id.to_owned())
.await?
.release_command_receipt(reminder.command_id.clone())
.await?;
}
let state = mj_core::subagent::report_state(
handback_tool,
&report,
Some(last_turn),
&[],
mj_core::clock::epoch_millis(),
);
if state != (ReportState::Pending { remind: true }) {
return Ok(state == (ReportState::Pending { remind: false }));
}
let command_id = format!(
"{}-{}",
mj_core::subagent::HANDBACK_REMINDER_PREFIX,
last_turn.completed_ordinal
);
let submitted = async {
let handle = self
.sessions
.session(child_session_id.to_owned())
.await
.with_context(|| format!("session {child_session_id} is not running"))?;
handle
.submit_durable(
command_id.clone(),
RelayCommand::HandbackReminder {
completed_command_id: last_turn.command_id.clone(),
completed_ordinal: last_turn.completed_ordinal,
},
)
.await
}
.await;
let child_id = child_session_id.to_owned();
let for_command_id = last_turn.command_id.clone();
match submitted {
Ok(_) => {
let reminder = mj_core::subagent::HandbackReminder {
command_id: command_id.clone(),
for_command_id,
sent_at_ms: mj_core::clock::epoch_millis(),
};
blocking("record handback reminder", move || {
crate::database::record_handback_reminder(&child_id, &reminder)
})
.await?;
self.sessions
.session(child_session_id.to_owned())
.await?
.release_command_receipt(command_id)
.await?;
Ok(true)
}
Err(error) => {
if error
.downcast_ref::<mj_client::session::DeliveryUnconfirmed>()
.is_some()
{
return Err(error);
}
tracing::warn!(
child_session_id,
error = format!("{error:#}"),
"could not remind a sub-agent to hand back its report"
);
blocking("record failed handback reminder", move || {
crate::database::record_handback_reminder_failed(&child_id, &for_command_id)
})
.await?;
Ok(false)
}
}
}
pub async fn park_subagent(
&self,
child_session_id: &str,
) -> Result<crate::controller::ParkOutcome> {
Arc::clone(&self.exports)
.park_subagent(child_session_id.to_owned())
.await
}
fn require_live_target(&self, session_id: &str) -> Result<(), ExportError> {
let record = self
.exports
.session_record(session_id)
.ok_or_else(|| ExportError::Refused(format!("unknown session {session_id}")))?;
match record.target {
Some(_) => Ok(()),
None => Err(ExportError::Refused(format!(
"session {session_id} has no live target; export its checkpoint bundle instead"
))),
}
}
async fn require_idle_turn(&self, session_id: &str) -> Result<(), ExportError> {
let Ok(handle) = self.sessions.session(session_id.to_owned()).await else {
return Ok(());
};
let Some(snapshot) = handle.view().snapshot else {
return Ok(());
};
let running = matches!(
snapshot.materialized.execution,
MaterializedExecutionState::Running { .. }
) || snapshot.materialized.active_turn.is_some();
if running {
return Err(ExportError::Refused(
"this session is running a turn; cancel or wait for it before pushing".into(),
));
}
Ok(())
}
}
fn profile_remaining_percent(report: Option<&ProfileQuota>) -> Option<u8> {
let report = report.filter(|report| report.error.is_none())?;
if report.is_usage_priced() {
return Some(100);
}
report
.windows
.iter()
.filter_map(|window| window.remaining_percent)
.min()
}
fn subagent_status(
record: Option<&mj_core::state::SessionRecord>,
summary: Option<&mj_core::state::MaterializedSessionSummary>,
start: Option<&StartStatus>,
finished_turn_message: Option<&str>,
closing: bool,
progress: &ChildProgress,
) -> (String, Option<String>, bool) {
if closing && record.is_some_and(|record| record.state != SessionState::Stopped) {
return ("stopping".into(), None, false);
}
let recorded_cause = record
.filter(|record| record.state == SessionState::Error)
.and_then(|record| record.last_error.clone());
if let Some(StartStatus::Failed { message }) = start {
return (
"error".into(),
Some(recorded_cause.unwrap_or_else(|| message.clone())),
true,
);
}
let start_pending = matches!(start, Some(StartStatus::Pending));
let lifecycle = record.map(|record| record.state);
let idle = |summary: &mj_core::state::MaterializedSessionSummary| {
lifecycle == Some(SessionState::Parked)
|| matches!(summary.execution, MaterializedExecutionState::Idle)
};
match lifecycle {
Some(SessionState::Error) => ("error".into(), recorded_cause, true),
Some(SessionState::Lost) => ("lost".into(), None, true),
Some(SessionState::Stopped | SessionState::DestroyedWithDataLoss) | None => {
("stopped".into(), None, true)
}
Some(SessionState::Closing | SessionState::Destroying) => ("stopping".into(), None, false),
Some(SessionState::Provisioning) if summary.is_none() || start_pending => {
("preparing".into(), None, false)
}
_ if start_pending => ("running".into(), None, false),
_ => match summary {
Some(summary)
if idle(summary)
&& (progress.awaiting_prompt(match start {
Some(StartStatus::Submitted { turn_id }) => Some(*turn_id),
_ => None,
}) || matches!(progress.report, ReportState::Pending { .. })) =>
{
("running".into(), summary.last_agent_message.clone(), false)
}
Some(summary) if idle(summary) => {
match (
&progress.report,
&progress.login_failure,
&progress.failed_turn,
) {
(ReportState::Delivered(message), _, _) => {
("completed".into(), Some(message.clone()), true)
}
(_, Some(profile_id), _) => (
"failed".into(),
Some(format!(
"profile {profile_id}: {}",
mj_core::subagent::login_invalid_reason(profile_id)
)),
true,
),
(_, None, Some((state, reason))) => {
((*state).to_owned(), Some(reason.clone()), true)
}
_ => (
"completed".into(),
finished_turn_message
.map(str::to_owned)
.or_else(|| summary.last_agent_message.clone()),
true,
),
}
}
Some(summary) => ("running".into(), summary.last_agent_message.clone(), false),
None => ("preparing".into(), None, false),
},
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub(crate) struct ChildProgress {
pub report: ReportState,
pub awaited_ordinal: Option<u64>,
pub answered_ordinal: Option<u64>,
pub failed_turn: Option<(&'static str, String)>,
pub login_failure: Option<String>,
pub report_dir: Option<String>,
}
impl ChildProgress {
#[cfg(test)]
pub(crate) fn settled(report: ReportState) -> Self {
Self {
report,
awaited_ordinal: None,
answered_ordinal: None,
failed_turn: None,
login_failure: None,
report_dir: None,
}
}
fn awaiting_prompt(&self, submitted: Option<u64>) -> bool {
self.awaited_ordinal.max(submitted).is_some_and(|awaited| {
self.answered_ordinal
.is_none_or(|answered| answered < awaited)
})
}
}
pub(crate) fn load_child_progress(child_id: &str) -> Result<ChildProgress> {
let subagent = crate::database::load_subagent(child_id)?;
let handback_tool = subagent.as_ref().is_some_and(|record| record.handback_tool);
let recorded = crate::database::load_subagent_report(child_id)?;
let (active, last) = crate::database::load_materialized_turn_outcome(child_id)?
.map(|(_, active, last)| (active, last))
.unwrap_or_default();
let in_flight = active
.iter()
.map(|turn| turn.command_id.as_str())
.collect::<Vec<_>>();
let report = mj_core::subagent::report_state(
handback_tool,
&recorded,
last.as_ref(),
&in_flight,
mj_core::clock::epoch_millis(),
);
let failed_turn = match last.as_ref() {
Some(turn) if mj_core::subagent::failed_turn(turn, None).is_some() => {
let message = crate::database::load_materialized_finished_turn_message(child_id)?;
mj_core::subagent::failed_turn(turn, message.as_deref())
}
_ => None,
};
let login_failure = last
.as_ref()
.filter(|turn| mj_core::subagent::turn_failed_on_login(turn))
.and(subagent.map(|record| record.profile_id));
Ok(ChildProgress {
report,
awaited_ordinal: recorded.awaited_ordinal,
answered_ordinal: last.as_ref().and_then(|turn| turn.accepted_ordinal),
failed_turn,
login_failure,
report_dir: recorded.report_dir,
})
}
fn wait_agent_entry(
id: &str,
state: &str,
output: Option<String>,
finished: bool,
progress: &ChildProgress,
) -> serde_json::Value {
let output = output.map(|output| bounded_report(&output));
let mut agent = serde_json::json!({
"child_session_id":id,
"report_source":report_source(state, &progress.report),
"state":state,
"finished":finished,
"output":output.as_ref().map(|(output, _)| output),
"report_dir":progress.report_dir,
});
if output.is_some_and(|(_, truncated)| truncated) {
agent["truncated"] = serde_json::Value::Bool(true);
}
if state == "failed"
&& let Some(profile_id) = &progress.login_failure
{
agent["failure"] = serde_json::json!({
"kind": "login_invalid",
"profile_id": profile_id,
});
}
agent
}
fn mark_parked(entry: &mut serde_json::Value, record: Option<&mj_core::state::SessionRecord>) {
if record.is_some_and(|record| record.state == SessionState::Parked) {
entry["parked"] = serde_json::Value::Bool(true);
}
}
fn report_source(state: &str, report: &ReportState) -> Option<&'static str> {
(state == "completed").then_some(match report {
ReportState::Delivered(_) => "handback",
_ => "last_message",
})
}
async fn load_startup_status(session_id: String) -> Result<Option<StartStatus>> {
Ok(load_startup_context(session_id).await?.1)
}
async fn load_startup_context(session_id: String) -> Result<(Option<String>, Option<StartStatus>)> {
let steps = blocking("read durable startup status", move || {
crate::database::load_latest_startup_group(&session_id)
})
.await?;
let group_id = steps.last().and_then(|step| step.group_id.clone());
if let Some(step) = steps
.iter()
.find(|step| matches!(step.phase.as_str(), "failed" | "rejecting"))
{
return Ok((
group_id,
Some(StartStatus::Failed {
message: step
.error
.clone()
.unwrap_or_else(|| "session startup was rejected".to_owned()),
}),
));
}
if steps
.iter()
.any(|step| matches!(step.phase.as_str(), "pending" | "delivering" | "accepted"))
{
return Ok((group_id, Some(StartStatus::Pending)));
}
let Some(last) = steps.last() else {
return Ok((group_id, None));
};
if last.phase == "done"
&& matches!(
serde_json::from_str::<crate::daemon::StartupStep>(&last.step_json)?,
crate::daemon::StartupStep::ApiPrompt { .. }
)
{
return Ok((
group_id,
Some(StartStatus::Submitted {
turn_id: last
.accepted_ordinal
.context("completed startup prompt has no acceptance ordinal")?,
}),
));
}
Ok((group_id, None))
}
async fn submit_prompt(handle: &SessionHandle, text: String) -> Result<u64> {
handle
.submit(
new_command_id("api")?,
RelayCommand::Prompt {
prompt: vec![ContentBlock::Text(TextContent::new(text))],
},
)
.await
}
async fn submit_notice(handle: &SessionHandle, text: String) -> Result<u64> {
handle
.submit(
new_command_id("subagent")?,
RelayCommand::RecordNotice { text },
)
.await
}
async fn blocking<T, F>(label: &'static str, job: F) -> Result<T>
where
T: Send + 'static,
F: FnOnce() -> Result<T> + Send + 'static,
{
tokio::task::spawn_blocking(job)
.await
.with_context(|| format!("{label} task panicked"))?
}
fn checkpoint_export_error(error: anyhow::Error) -> ExportError {
if crate::controller::checkpoint_was_deferred(&error) {
ExportError::Refused(format!("{error:#}"))
} else {
ExportError::Failed(error)
}
}
async fn export_blocking<T, F>(label: &'static str, job: F) -> Result<T, ExportError>
where
T: Send + 'static,
F: FnOnce() -> Result<T> + Send + 'static,
{
match tokio::task::spawn_blocking(job).await {
Ok(result) => result.map_err(ExportError::Failed),
Err(error) => Err(ExportError::Failed(anyhow!(
"{label} task panicked: {error}"
))),
}
}
async fn export_layout(session_id: String) -> Result<SessionExportLayout, ExportError> {
export_blocking("resolve the session export layout", move || {
let executor = CancellableProcessExecutor::with_timeout(EXPORT_TIMEOUT);
Controller::load()?.session_export_layout(&session_id, &executor)
})
.await
}
async fn supervised_checkpoint(
exports: Arc<dyn ExportRuntime>,
session_id: String,
) -> Result<mj_core::state::CheckpointMetadata> {
let upgrade_work = crate::upgrade::activity("API checkpoint")?;
let checkpoint = tokio::spawn(async move {
let _upgrade_work = upgrade_work;
exports.checkpoint_now(session_id).await
});
tokio::spawn(async move {
let result = match checkpoint.await {
Ok(result) => result,
Err(error) => Err(anyhow!("the session checkpoint task failed: {error}")),
};
if let Err(error) = &result {
tracing::warn!(?error, "API bundle export checkpoint failed");
}
result
})
.await
.unwrap_or_else(|error| Err(anyhow!("the session checkpoint task failed: {error}")))
}
fn agent_working_directory(layout: &SessionExportLayout) -> Result<String, ExportError> {
Ok(target_join(
&layout.workspace_root,
&primary_repository(layout)?.relative_destination,
))
}
fn primary_repository(
layout: &SessionExportLayout,
) -> Result<&mj_checkpoint::checkpoint::CheckpointRepositorySpec, ExportError> {
layout
.repositories
.iter()
.find(|repository| repository.id == layout.primary_repository)
.ok_or_else(|| {
ExportError::Failed(anyhow!(
"session workspace has no repository {:?}",
layout.primary_repository
))
})
}
fn export_root_and_path(
layout: &SessionExportLayout,
relative: &Path,
) -> Result<(String, String), ExportError> {
let primary = primary_repository(layout)?;
let agent_directory = target_join(&layout.workspace_root, &primary.relative_destination);
let (root, mut resolved) = if layout.repositories.len() > 1 {
let prefix = primary
.relative_destination
.components()
.filter_map(|component| match component {
Component::Normal(part) => Some(part.to_string_lossy().into_owned()),
_ => None,
})
.collect();
(layout.workspace_root.clone(), prefix)
} else {
(agent_directory.clone(), Vec::new())
};
let climbed = || {
if root == agent_directory {
format!(
"{} climbs above {agent_directory}, the directory the agent runs in",
relative.display()
)
} else {
format!(
"{} climbs above the session workspace {root}; it was resolved in {agent_directory}",
relative.display()
)
}
};
for component in relative.components() {
match component {
Component::Normal(part) => resolved.push(part.to_string_lossy().into_owned()),
Component::CurDir => {}
Component::ParentDir => {
if resolved.pop().is_none() {
return Err(ExportError::Refused(climbed()));
}
}
Component::RootDir | Component::Prefix(_) => {
return Err(ExportError::Refused(format!(
"{} must be relative to {agent_directory}",
relative.display()
)));
}
}
}
if resolved.is_empty() {
return Err(ExportError::Refused(format!(
"{} names {root} itself, not a file in it",
relative.display()
)));
}
Ok((root, resolved.join("/")))
}
fn target_join(root: &str, relative: &Path) -> String {
let mut path = root.trim_end_matches('/').to_owned();
for component in relative.components() {
if let Component::Normal(part) = component {
path.push('/');
path.push_str(&part.to_string_lossy());
}
}
path
}
async fn write_workspace_file(
exports: Arc<dyn ExportRuntime>,
session_id: String,
path: PathBuf,
bytes: Vec<u8>,
overwrite: bool,
cancelled: Arc<std::sync::atomic::AtomicBool>,
) -> Result<(), ExportError> {
use std::sync::atomic::Ordering;
let record = exports
.session_record(&session_id)
.ok_or_else(|| ExportError::Refused("unknown session".into()))?;
let handle = exports
.workspace_session(session_id.clone())
.await
.map_err(|e| ExportError::Refused(format!("{e:#}")))?;
let layout = export_layout(session_id.clone()).await?;
let (root, relative) = export_root_and_path(&layout, &path)?;
if cancelled.load(Ordering::Acquire) {
return Err(ExportError::Refused("file upload cancelled".into()));
}
let mut lease = crate::controller::IdleWorkspaceLease::acquire(&handle, record.harness_kind)
.await
.map_err(|e| ExportError::Refused(format!("{e:#}")))?;
if cancelled.load(Ordering::Acquire) {
return Err(ExportError::Refused("file upload cancelled".into()));
}
let worker_cancelled = cancelled.clone();
let mut transfer = tokio::task::spawn_blocking(move || {
let binary = format!(
"{}/hel",
targets::worker_root(&layout.backend, &session_id)?
);
let mut argv = vec![
binary,
"worker".into(),
"write-file".into(),
"--length".into(),
bytes.len().to_string(),
"--root".into(),
root,
"--path".into(),
relative,
];
if overwrite {
argv.push("--overwrite".into());
}
let command =
targets::command_on_locator(&layout.backend, &session_id, argv, "session file write")?
.with_sensitive_stdin(bytes);
CancellableProcessExecutor::new(worker_cancelled)
.with_deadline(EXPORT_TIMEOUT)
.execute(&command)
});
let output = loop {
tokio::select! {
output = &mut transfer => break output.map_err(|e| ExportError::Failed(e.into()))?.map_err(ExportError::Failed)?,
() = tokio::time::sleep(Duration::from_millis(250)) => {
if let Err(error) = tokio::time::timeout(Duration::from_secs(10), lease.verify()).await.context("checking file write barrier timed out").and_then(|r| r) {
cancelled.store(true, Ordering::Release);
match transfer.await {
Ok(Ok(_)) => {},
Ok(Err(failure)) => tracing::warn!("cancelled file transfer: {failure:#}"),
Err(failure) => tracing::warn!("file transfer task failed: {failure}"),
}
return Err(ExportError::Failed(error));
}
}
}
};
let result = worker_output(output, "session file write").map(|_| ());
lease.release().await.map_err(|e| {
ExportError::Failed(
e.context("file transfer ended but the workspace barrier could not be released"),
)
})?;
result
}
async fn worker_command(
layout: SessionExportLayout,
session_id: String,
arguments: Vec<String>,
purpose: &'static str,
) -> Result<Vec<u8>, ExportError> {
let output: CommandOutput = export_blocking(purpose, move || {
let binary = format!(
"{}/hel",
targets::worker_root(&layout.backend, &session_id)?
);
let mut argv = vec![binary, "worker".to_owned()];
argv.extend(arguments);
let command = targets::command_on_locator(&layout.backend, &session_id, argv, purpose)?;
CancellableProcessExecutor::with_timeout(EXPORT_TIMEOUT).execute(&command)
})
.await?;
worker_output(output, purpose)
}
fn worker_output(output: CommandOutput, purpose: &str) -> Result<Vec<u8>, ExportError> {
match output.status {
0 => Ok(output.stdout),
CLAP_USAGE_EXIT_CODE => Err(ExportError::Refused(format!(
"worker on this target does not support {purpose}; resume the session to upgrade"
))),
mj_checkpoint::archive::EXPORT_REFUSED_EXIT_CODE => Err(ExportError::Refused(
refusal_reason(&output.stdout, &output.stderr, purpose),
)),
status => Err(ExportError::Failed(anyhow!(
"{purpose} failed with status {status}: {}",
String::from_utf8_lossy(&output.stderr).trim()
))),
}
}
fn refusal_reason(stdout: &[u8], stderr: &[u8], purpose: &str) -> String {
[stdout, stderr]
.into_iter()
.map(|stream| String::from_utf8_lossy(stream).trim().to_owned())
.find(|reason| !reason.is_empty())
.unwrap_or_else(|| format!("{purpose} was refused by the target"))
}
impl SubagentBackend for ApiBackend {
fn subagent_report(
&self,
session_id: String,
) -> BoxFuture<'_, Result<Option<(bool, mj_core::subagent::SubagentReport)>>> {
Box::pin(async move {
blocking("load sub-agent report", move || {
let Some(record) = crate::database::load_subagent(&session_id)? else {
return Ok(None);
};
let report = crate::database::load_subagent_report(&session_id)?;
Ok(Some((record.handback_tool, report)))
})
.await
})
}
fn create_workspace(
&self,
name: String,
) -> BoxFuture<'_, Result<mj_core::workspace::WorkspaceRecord>> {
Box::pin(async move {
let workspace = tokio::task::spawn_blocking(move || {
crate::database::create_or_get_workspace(&name)
})
.await??;
self.exports.refresh_workspaces().await?;
Ok(workspace)
})
}
fn subagent_candidates(
&self,
parent_profile: String,
) -> BoxFuture<'_, Result<crate::server::api::SubagentCandidates>> {
Box::pin(async move {
let mut candidates = crate::server::api::SubagentCandidates::default();
let mut ids = self.profile_catalog.candidates(&parent_profile)?;
{
let rejected = self
.rejected_logins
.lock()
.map_err(|_| anyhow!("refused logins lock poisoned"))?;
ids.retain(|(profile_id, _)| match rejected.refusal(profile_id) {
Some(reason) => {
candidates.unavailable.push((profile_id.clone(), reason));
false
}
None => true,
});
}
let discoveries = futures::future::join_all(
ids.iter()
.map(|(id, _)| self.profile_catalog.capabilities(std::slice::from_ref(id))),
)
.await;
let quota_reports = self
.quota_reports
.lock()
.map_err(|_| anyhow!("sub-agent quota reports lock poisoned"))?;
for ((profile_id, harness), discovery) in ids.into_iter().zip(discoveries) {
match discovery.map(|mut choices| choices.pop()) {
Ok(Some(choices)) => {
let remaining_percent =
profile_remaining_percent(quota_reports.get(&profile_id));
candidates
.offered
.push(crate::server::api::SubagentCandidate {
profile_id,
harness,
choices,
remaining_percent,
});
}
Ok(None) => candidates
.unavailable
.push((profile_id, "its discovery returned nothing".to_owned())),
Err(error) => {
tracing::warn!(profile_id, %error, "sub-agent profile left out");
candidates
.unavailable
.push((profile_id, format!("{error:#}")));
}
}
}
if candidates.offered.is_empty() && !candidates.unavailable.is_empty() {
bail!(
"no sub-agent profile is available: {}",
candidates
.unavailable
.iter()
.map(|(id, reason)| format!("{id} ({reason})"))
.collect::<Vec<_>>()
.join("; ")
);
}
Ok(candidates)
})
}
fn wiki_sync_is_stale(&self) -> bool {
self.exports.wiki_sync_is_stale()
}
fn wiki_request_sync(&self) {
self.exports.wiki_request_sync();
}
fn wiki_search(&self, query: String, limit: usize) -> BoxFuture<'_, Result<WikiSearchPage>> {
let runtime = Arc::clone(&self.exports);
Box::pin(async move { runtime.wiki_search(query, limit).await })
}
fn wiki_brief(
&self,
wiki_id: String,
max_chars: usize,
) -> BoxFuture<'_, Result<Option<String>>> {
let runtime = Arc::clone(&self.exports);
Box::pin(async move { runtime.wiki_brief(wiki_id, max_chars).await })
}
fn wiki_hits(
&self,
wiki_id: String,
query: String,
context_messages: usize,
per_message_chars: usize,
) -> BoxFuture<'_, Result<Option<WikiHitTranscript>>> {
let runtime = Arc::clone(&self.exports);
Box::pin(async move {
runtime
.wiki_hits(wiki_id, query, context_messages, per_message_chars)
.await
})
}
fn wiki_session(&self, wiki_id: String) -> BoxFuture<'_, Result<Option<WikiSessionInfo>>> {
let runtime = Arc::clone(&self.exports);
Box::pin(async move { runtime.wiki_session(wiki_id).await })
}
fn wiki_restore(&self, request: WikiRestoreRequest) -> BoxFuture<'_, Result<Option<String>>> {
let runtime = Arc::clone(&self.exports);
Box::pin(async move { runtime.wiki_restore(request).await })
}
fn start_subagent(
&self,
request: crate::controller::RegisterSubagentRequest,
) -> BoxFuture<'_, Result<mj_core::subagent::SubagentRecord>> {
let runtime = Arc::clone(&self.exports);
Box::pin(async move { runtime.spawn_subagent(request).await })
}
fn cancel_start(&self, session_id: String) -> BoxFuture<'_, Result<()>> {
Arc::clone(&self.exports).cancel_startup(session_id)
}
fn set_config(
&self,
session_id: String,
key: String,
value: String,
) -> BoxFuture<'_, Result<()>> {
Box::pin(async move {
let (group_id, status) = self.exports.startup_context(session_id.clone()).await?;
ensure!(
!matches!(status, Some(StartStatus::Pending)),
"session initialization is still running"
);
let handle = self.sessions.session(session_id.clone()).await?;
handle.set_config(key, value).await?;
self.exports
.dismiss_startup_status(session_id, group_id)
.await?;
Ok(())
})
}
fn session_handle(&self, session_id: String) -> BoxFuture<'_, Result<Option<SessionHandle>>> {
Box::pin(async move { Ok(self.sessions.session(session_id).await.ok()) })
}
fn prompt(&self, session_id: String, text: String) -> BoxFuture<'_, Result<u64>> {
Box::pin(async move {
let (group_id, status) = self.exports.startup_context(session_id.clone()).await?;
ensure!(
!matches!(status, Some(StartStatus::Pending)),
"session initialization is still running"
);
let handle = self
.sessions
.session(session_id.clone())
.await
.with_context(|| format!("session {session_id} is not running"))?;
let turn = submit_prompt(&handle, text).await?;
self.exports
.dismiss_startup_status(session_id, group_id)
.await?;
Ok(turn)
})
}
fn runtime_receipt(
&self,
session_id: String,
) -> BoxFuture<'_, Result<Option<mj_core::harness_runtime::RuntimeReceipt>>> {
Box::pin(async move {
blocking("load runtime receipt", move || {
crate::database::load_runtime_receipt(&session_id)
})
.await
})
}
fn turn_state(&self, session_id: String) -> BoxFuture<'_, Result<Option<TurnState>>> {
Box::pin(async move {
blocking("load turn outcome", move || {
Ok(
crate::database::load_materialized_turn_outcome(&session_id)?.map(
|(execution, active_turn, last_turn_outcome)| TurnState {
execution,
active_turn,
last_turn_outcome,
},
),
)
})
.await
})
}
fn turn_summary(
&self,
session_id: String,
turn: TurnSpan,
) -> BoxFuture<'_, Result<TurnSummary>> {
Box::pin(async move {
blocking("load turn summary", move || {
crate::database::load_materialized_turn_summary(
&session_id,
turn.start_position,
turn.completed_position,
)
})
.await
})
}
fn start_followup(
&self,
session_id: String,
followup: StartFollowup,
) -> BoxFuture<'_, Result<()>> {
Box::pin(async move {
self.start_followup_with_id(session_id, followup, new_command_id("api-startup")?)
.await
})
}
fn start_status(&self, session_id: String) -> BoxFuture<'_, Result<Option<StartStatus>>> {
self.exports.startup_status(session_id)
}
fn transcript(
&self,
session_id: String,
after_seq: u64,
limit: usize,
role: Option<mj_core::transcript::TranscriptRole>,
) -> BoxFuture<'_, Result<Option<TranscriptPage>>> {
Box::pin(async move {
blocking("load transcript page", move || {
crate::database::load_materialized_transcript_filtered(
&session_id,
after_seq,
limit,
role,
)
})
.await
})
}
fn diff(
&self,
session_id: String,
options: crate::server::api::DiffOptions,
) -> BoxFuture<'_, std::result::Result<String, ExportError>> {
Box::pin(async move {
self.require_live_target(&session_id)?;
let layout = export_layout(session_id.clone()).await?;
let repository = agent_working_directory(&layout)?;
let mut arguments = vec!["diff".to_owned(), "--repository".to_owned(), repository];
if let Some(base) = options.base {
arguments.push(format!("--base={base}"));
} else if let Some(worktree) = &layout.managed_worktree {
match &worktree.base_commit {
Some(base) => {
arguments.push("--base".to_owned());
arguments.push(base.clone());
}
None => {
arguments.push("--branch".to_owned());
arguments.push(worktree.branch.clone());
}
}
}
if options.json {
arguments.push("--json".to_owned());
}
let stdout = worker_command(layout, session_id, arguments, "session diff").await?;
String::from_utf8(stdout).map_err(|error| {
ExportError::Failed(anyhow!("the session diff was not UTF-8: {error}"))
})
})
}
fn read_file(
&self,
session_id: String,
path: PathBuf,
) -> BoxFuture<'_, std::result::Result<Vec<u8>, ExportError>> {
Box::pin(async move {
self.require_live_target(&session_id)?;
let layout = export_layout(session_id.clone()).await?;
let (root, relative) = export_root_and_path(&layout, &path)?;
let arguments = vec![
"read-file".to_owned(),
"--root".to_owned(),
root,
"--path".to_owned(),
relative,
];
worker_command(layout, session_id, arguments, "session file read").await
})
}
fn read_context_file(
&self,
session_id: String,
path: PathBuf,
) -> BoxFuture<'_, std::result::Result<Vec<u8>, ExportError>> {
Box::pin(async move {
self.require_live_target(&session_id)?;
let layout = export_layout(session_id.clone()).await?;
let root = agent_working_directory(&layout)?;
let arguments = vec![
"read-file".to_owned(),
"--root".to_owned(),
root,
"--path".to_owned(),
target_join("", &path).trim_start_matches('/').to_owned(),
];
worker_command(layout, session_id, arguments, "sub-agent context file read").await
})
}
fn write_file(
&self,
session_id: String,
path: PathBuf,
bytes: Vec<u8>,
overwrite: bool,
) -> BoxFuture<'_, Result<(), ExportError>> {
let exports = self.exports.clone();
Box::pin(async move {
if matches!(
self.start_status(session_id.clone())
.await
.map_err(ExportError::Failed)?,
Some(StartStatus::Pending)
) {
return Err(ExportError::Refused(
"session initialization is still running".into(),
));
}
let cancelled = Arc::new(std::sync::atomic::AtomicBool::new(false));
let _cancel_on_drop = super::ProcessCancellationGuard(cancelled.clone());
let task = tokio::spawn(write_workspace_file(
exports, session_id, path, bytes, overwrite, cancelled,
));
tokio::spawn(async move {
let result = task
.await
.map_err(|error| ExportError::Failed(error.into()))
.and_then(|result| result);
if let Err(error) = &result {
tracing::warn!(?error, "API file injection failed");
}
result
})
.await
.map_err(|error| ExportError::Failed(error.into()))?
})
}
fn push_branch(
&self,
session_id: String,
branch: String,
) -> BoxFuture<'_, std::result::Result<PushedBranch, ExportError>> {
Box::pin(async move {
self.require_live_target(&session_id)?;
self.require_idle_turn(&session_id).await?;
let layout = export_layout(session_id.clone()).await?;
let repository = agent_working_directory(&layout)?;
let arguments = vec![
"push-branch".to_owned(),
"--root".to_owned(),
targets::worker_root(&layout.backend, &session_id).map_err(ExportError::Failed)?,
"--repository".to_owned(),
repository,
"--branch".to_owned(),
branch,
];
let stdout =
worker_command(layout, session_id, arguments, "session branch push").await?;
let pushed: mj_checkpoint::archive::PushedBranch = serde_json::from_slice(&stdout)
.map_err(|error| {
ExportError::Failed(anyhow!("the worker's push result was unreadable: {error}"))
})?;
Ok(PushedBranch {
branch: pushed.branch,
remote: pushed.remote,
})
})
}
fn bundle(
&self,
session_id: String,
) -> BoxFuture<'_, std::result::Result<BundleExport, ExportError>> {
Box::pin(async move {
let record = self
.exports
.session_record(&session_id)
.ok_or_else(|| ExportError::Refused(format!("unknown session {session_id}")))?;
let archive_path = match record.target {
Some(_) => match supervised_checkpoint(self.exports.clone(), session_id.clone())
.await
{
Ok(checkpoint) => checkpoint.archive_path,
Err(error)
if error
.downcast_ref::<crate::daemon::SessionLifecycleBusy>()
.is_some() =>
{
record
.checkpoint
.ok_or_else(|| {
ExportError::Refused(format!(
"{error}, and has no earlier checkpoint to export a bundle from"
))
})?
.archive_path
}
Err(error) => return Err(checkpoint_export_error(error)),
},
None => {
record
.checkpoint
.ok_or_else(|| {
ExportError::Refused(
"this session has no checkpoint to export a bundle from".into(),
)
})?
.archive_path
}
};
let bundles = export_blocking("verify the checkpoint bundles", move || {
mj_checkpoint::archive::verify_repository_bundles_streaming(&archive_path)
})
.await?;
let repository = bundles
.repositories
.iter()
.find(|repository| repository.metadata.id == bundles.primary_repository)
.ok_or_else(|| {
ExportError::Failed(anyhow!(
"the checkpoint has no repository {:?}",
bundles.primary_repository
))
})?;
if repository.committed_bundle.is_empty() {
return Err(ExportError::Refused(
"no commits beyond the session base".into(),
));
}
Ok(BundleExport {
repository: repository.metadata.id.clone(),
bytes: repository.committed_bundle.clone(),
})
})
}
}
#[cfg(test)]
mod tests;