use super::*;
const REPORT_ROOT_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(60);
impl RuntimeState {
pub(super) async fn start_create_session(
self: &Arc<Self>,
request: CreateSessionRequest,
) -> Result<RegisteredSession> {
self.start_create_session_inner(request, CreateSessionControl::default(), None)
.await
}
pub async fn start_subagent_session(
self: &Arc<Self>,
request: crate::controller::RegisterSubagentRequest,
) -> Result<mj_core::subagent::SubagentRecord> {
let _upgrade_work = crate::upgrade::activity("subagent admission")?;
let (relation, needs_provisioning) = blocking(move || {
let mut controller = Controller::load()?;
if let Some(existing) = crate::database::lookup_subagent_request(
&request.parent_session_id,
&request.request_key,
)? {
let needs_provisioning = controller
.state
.sessions
.get(&existing.child_session_id)
.is_some_and(|record| {
record.state == mj_core::state::SessionState::Provisioning
});
return Ok((existing, needs_provisioning));
}
let mut request = request;
let executor = CancellableProcessExecutor::with_timeout(REPORT_ROOT_TIMEOUT);
request.report_root = Some(
controller
.prepare_subagent_report_root(&request.parent_session_id, &executor)
.context("create the sub-agent report directory")?,
);
controller
.register_subagent(request)
.map(|relation| (relation, true))
})
.await?;
if !needs_provisioning {
return Ok(relation);
}
let session_id = relation.child_session_id.clone();
self.start_or_join_lifecycle_controlled(
session_id.clone(),
LifecycleKind::Create,
None,
Some(relation.request_key.clone()),
None,
move |state, session_id, cancelled| async move {
let mut controller = tokio::task::spawn_blocking(Controller::load)
.await
.context("load controller for sub-agent startup")??;
if controller
.state
.sessions
.get(&session_id)
.is_none_or(|record| record.state != mj_core::state::SessionState::Provisioning)
{
return Ok(DaemonLifecycleResult::Done);
}
let executor = DaemonStageReportingExecutor::new(
CancellableProcessExecutor::new(cancelled),
state,
session_id.clone(),
);
controller
.provision_subagent_session_controlled(&session_id, &executor)
.await?;
Ok(DaemonLifecycleResult::Done)
},
)?;
self.reload_controller().await?;
Ok(relation)
}
pub async fn start_create_session_controlled(
self: &Arc<Self>,
request: CreateSessionRequest,
control: CreateSessionControl,
publication: tokio::sync::oneshot::Receiver<std::result::Result<(), String>>,
) -> Result<RegisteredSession> {
self.start_create_session_inner(request, control, Some(publication))
.await
}
pub(super) async fn start_create_session_inner(
self: &Arc<Self>,
request: CreateSessionRequest,
control: CreateSessionControl,
publication: Option<tokio::sync::oneshot::Receiver<std::result::Result<(), String>>>,
) -> Result<RegisteredSession> {
let mut request = request;
let parent = request.profile_id.clone();
let supplied = request.subagents.clone();
let policy = blocking(move || {
let controller = Controller::load()?;
let kind = controller
.config
.enabled_profile(&parent)
.context("parent profile unavailable")?
.kind;
let policy = supplied.unwrap_or_else(|| {
if kind.supports_delegation_tools() {
controller.state.last_subagent_policy.clone()
} else {
Default::default()
}
});
anyhow::ensure!(
kind.supports_delegation_tools()
|| policy == mj_core::subagent::SubagentPolicy::Native,
"subagent policies are supported only by Claude and Codex"
);
Ok(policy)
})
.await?;
if let mj_core::subagent::SubagentPolicy::SingleModel { model, .. } = &policy {
let options = crate::controller::profile_config::subagent_options(
request.profile_id.clone(),
Some(model.clone()),
)
.await?;
options.validate(&policy).map_err(anyhow::Error::msg)?;
}
let _upgrade_work = crate::upgrade::activity("session admission")?;
request.subagents = Some(policy);
let path_cancelled = control.cancelled.clone();
let registered = blocking(move || {
let mut controller = Controller::load()?;
let path_executor = crate::targets::CancellableProcessExecutor::new(path_cancelled)
.with_deadline(Duration::from_secs(30));
let project_directory = request
.project_directory
.as_deref()
.map(|path| {
controller.resolve_project_directory(
&request.target_template_id,
path,
&path_executor,
)
})
.transpose()?;
let session_id = controller.register_session_with_resources(
&request.profile_id,
&request.bundle_id,
&request.target_template_id,
request.title,
SessionLaunchOptions {
create_managed_worktree: request.create_managed_worktree,
launch_base: request.launch_base,
launch_branch: request.launch_branch,
checkout: request.checkout,
expected_runtime_identity: request.expected_runtime_identity,
subagents: request.subagents,
initial_prompt: request.initial_prompt,
workspace_id: request.workspace_id,
additional_mounts: request.additional_mounts,
resource_allocation: request.resource_allocation,
project_directory,
session_title_override: request.session_title_override,
},
)?;
if let Err(error) = mj_core::go::GoPreferences::remember_first_pair(
&mj_core::go::GoPreferences::path(),
&request.profile_id,
&request.target_template_id,
) {
tracing::warn!(%error, "could not save the default profile and target");
}
let session = controller
.state
.sessions
.get(&session_id)
.expect("newly registered session exists")
.clone();
let remembered_container_size = controller
.config
.targets
.get(&request.target_template_id)
.and_then(mj_core::config::container_size_host)
.and_then(|host| {
controller
.state
.container_sizes
.get(host)
.copied()
.map(|size| (host.to_owned(), size))
});
Ok(RegisteredSession {
session,
remembered_container_size,
})
})
.await?;
let session_id = registered.session.id.clone();
self.start_or_join_lifecycle_controlled(
session_id,
LifecycleKind::Create,
None,
None,
Some(control.clone()),
move |state, session_id, cancelled| async move {
let mut controller = tokio::task::spawn_blocking(Controller::load)
.await
.context("load controller for daemon create task")??;
let publication_error = if let Some(publication) = publication {
let published = tokio::select! {
result = publication => result.context("session publication owner stopped")
.and_then(|result| result.map_err(anyhow::Error::msg)),
() = async {
while !cancelled.load(Ordering::Acquire) {
tokio::time::sleep(Duration::from_millis(25)).await;
}
} => Err(anyhow!("session creation cancelled before publication")),
};
published.err()
} else {
None
};
if publication_error.is_some() {
control.request_cancel();
}
let executor = DaemonStageReportingExecutor::new(
CancellableProcessExecutor::new(cancelled),
state,
session_id.clone(),
);
let provision = controller
.provision_session_controlled_with_commit(&session_id, &executor, || {
ensure!(
control.grant_commit(),
"session creation cancelled before commit"
);
Ok(())
})
.await;
if let Some(error) = publication_error {
return match provision {
Ok(()) => Err(error),
Err(rollback) => {
Err(error.context(format!("discard unpublished session: {rollback:#}")))
}
};
}
provision?;
Ok(DaemonLifecycleResult::Done)
},
)?;
self.reload_controller().await?;
Ok(registered)
}
pub async fn wait_create_session(&self, session_id: &str) -> Result<()> {
let result = {
let lifecycle_owner = self.owner();
let lifecycle = &lifecycle_owner.lifecycle;
let active = lifecycle
.get(session_id)
.with_context(|| format!("no create operation exists for session {session_id}"))?;
ensure!(
active.kind == LifecycleKind::Create,
"session {session_id} is no longer being created"
);
active.result.clone()
};
let channel = result.clone();
let outcome = Self::wait_lifecycle_result(result).await;
self.remove_completed_lifecycle(&channel);
match outcome? {
DaemonLifecycleResult::Done => Ok(()),
DaemonLifecycleResult::Move(_)
| DaemonLifecycleResult::Park(_)
| DaemonLifecycleResult::Superseded => {
unreachable!("cleanup cannot return a move outcome")
}
DaemonLifecycleResult::DeferredCleanup => {
unreachable!("session creation cannot schedule target cleanup")
}
}
}
}
#[cfg(test)]
mod delegation_replay_tests {
use crate::controller::test_support::{IsolatedTest, test_name};
#[tokio::test]
async fn replayed_spawn_does_not_reprovision_an_existing_child() {
const CHILD: &str = "MJ_TEST_SPAWN_REPLAY";
if std::env::var_os(CHILD).is_none() {
let root = tempfile::tempdir().unwrap();
IsolatedTest::new(test_name(
module_path!(),
"replayed_spawn_does_not_reprovision_an_existing_child",
))
.env(CHILD, "1")
.env("MJ_INSTANCE", "concurrency-sweep-spawn")
.isolated_store(root.path())
.run();
return;
}
let _writer = crate::database::install_isolated_test_writer();
let workspace = crate::database::create_workspace("spawn replay").unwrap();
let parent = crate::daemon::tests::runtime_test_session(
"parent",
&workspace.id,
mj_core::state::SessionState::Running,
);
crate::database::save_session(&parent).unwrap();
for (index, status) in [
mj_core::state::SessionState::Running,
mj_core::state::SessionState::Parked,
mj_core::state::SessionState::Stopped,
]
.into_iter()
.enumerate()
{
let child_id = format!("child-{index}");
let child =
crate::daemon::tests::runtime_test_session(&child_id, &workspace.id, status);
let relation = crate::daemon::tests::runtime_test_subagent(&child_id, "parent");
crate::database::save_subagent_session(&child, &relation).unwrap();
let runtime = crate::daemon::tests::test_runtime_state();
let replayed = runtime
.start_subagent_session(crate::controller::RegisterSubagentRequest {
parent_session_id: "parent".into(),
task_name: "must reuse original".into(),
profile_id: "missing-profile".into(),
model: None,
effort: None,
working_directory: Default::default(),
initial_prompt: "must not resend".into(),
request_key: relation.request_key.clone(),
report_root: None,
})
.await
.unwrap();
assert_eq!(replayed, relation);
assert!(runtime.active_lifecycles().is_empty());
assert_eq!(
crate::database::load_session_record(&child_id)
.unwrap()
.unwrap()
.state,
status
);
}
}
}