mj_controller/daemon/
create.rs1use super::*;
2
3impl RuntimeState {
4 pub(super) async fn start_create_session(
5 self: &Arc<Self>,
6 request: CreateSessionRequest,
7 ) -> Result<RegisteredSession> {
8 self.start_create_session_inner(request, CreateSessionControl::default(), None)
9 .await
10 }
11
12 pub async fn start_subagent_session(
14 self: &Arc<Self>,
15 request: crate::controller::RegisterSubagentRequest,
16 ) -> Result<mj_core::subagent::SubagentRecord> {
17 let relation = blocking(move || {
18 let mut controller = Controller::load()?;
19 controller.register_subagent(request)
20 })
21 .await?;
22 let session_id = relation.child_session_id.clone();
23 self.start_or_join_lifecycle_controlled(
24 session_id.clone(),
25 LifecycleKind::Create,
26 None,
27 Some(relation.request_key.clone()),
28 None,
29 move |state, session_id, cancelled| async move {
30 let mut controller = tokio::task::spawn_blocking(Controller::load)
31 .await
32 .context("load controller for sub-agent startup")??;
33 let executor = DaemonStageReportingExecutor::new(
34 CancellableProcessExecutor::new(cancelled),
35 state,
36 session_id.clone(),
37 );
38 controller
39 .provision_subagent_session_controlled(&session_id, &executor)
40 .await?;
41 Ok(DaemonLifecycleResult::Done)
42 },
43 )?;
44 self.reload_controller().await?;
45 Ok(relation)
46 }
47
48 pub async fn start_create_session_controlled(
49 self: &Arc<Self>,
50 request: CreateSessionRequest,
51 control: CreateSessionControl,
52 publication: tokio::sync::oneshot::Receiver<std::result::Result<(), String>>,
53 ) -> Result<RegisteredSession> {
54 self.start_create_session_inner(request, control, Some(publication))
55 .await
56 }
57
58 pub(super) async fn start_create_session_inner(
59 self: &Arc<Self>,
60 request: CreateSessionRequest,
61 control: CreateSessionControl,
62 publication: Option<tokio::sync::oneshot::Receiver<std::result::Result<(), String>>>,
63 ) -> Result<RegisteredSession> {
64 let path_cancelled = control.cancelled.clone();
65 let registered = blocking(move || {
66 let mut controller = Controller::load()?;
67 let path_executor = crate::targets::CancellableProcessExecutor::new(path_cancelled)
68 .with_deadline(Duration::from_secs(30));
69 let project_directory = request
70 .project_directory
71 .as_deref()
72 .map(|path| {
73 controller.resolve_project_directory(
74 &request.target_template_id,
75 path,
76 &path_executor,
77 )
78 })
79 .transpose()?;
80 let session_id = controller.register_session_with_resources(
81 &request.profile_id,
82 &request.bundle_id,
83 &request.target_template_id,
84 request.title,
85 SessionLaunchOptions {
86 create_managed_worktree: request.create_managed_worktree,
87 mjolnir_subagents: request.mjolnir_subagents,
88 initial_prompt: request.initial_prompt,
89 workspace_id: request.workspace_id,
90 additional_mounts: request.additional_mounts,
91 resource_allocation: request.resource_allocation,
92 project_directory,
93 session_title_override: request.session_title_override,
94 },
95 )?;
96 let session = controller
97 .state
98 .sessions
99 .get(&session_id)
100 .expect("newly registered session exists")
101 .clone();
102 let remembered_container_size = controller
103 .config
104 .targets
105 .get(&request.target_template_id)
106 .and_then(mj_core::config::container_size_host)
107 .and_then(|host| {
108 controller
109 .state
110 .container_sizes
111 .get(host)
112 .copied()
113 .map(|size| (host.to_owned(), size))
114 });
115 Ok(RegisteredSession {
116 session,
117 remembered_container_size,
118 })
119 })
120 .await?;
121 let session_id = registered.session.id.clone();
122 self.start_or_join_lifecycle_controlled(
123 session_id,
124 LifecycleKind::Create,
125 None,
126 None,
127 Some(control.clone()),
128 move |state, session_id, cancelled| async move {
129 let mut controller = tokio::task::spawn_blocking(Controller::load)
130 .await
131 .context("load controller for daemon create task")??;
132 let publication_error = if let Some(publication) = publication {
133 let published = tokio::select! {
134 result = publication => result.context("session publication owner stopped")
135 .and_then(|result| result.map_err(anyhow::Error::msg)),
136 () = async {
137 while !cancelled.load(Ordering::Acquire) {
138 tokio::time::sleep(Duration::from_millis(25)).await;
139 }
140 } => Err(anyhow!("session creation cancelled before publication")),
141 };
142 published.err()
143 } else {
144 None
145 };
146 if publication_error.is_some() {
147 control.request_cancel();
148 }
149 let executor = DaemonStageReportingExecutor::new(
150 CancellableProcessExecutor::new(cancelled),
151 state,
152 session_id.clone(),
153 );
154 let provision = controller
155 .provision_session_controlled_with_commit(&session_id, &executor, || {
156 ensure!(
157 control.grant_commit(),
158 "session creation cancelled before commit"
159 );
160 Ok(())
161 })
162 .await;
163 if let Some(error) = publication_error {
164 return match provision {
165 Ok(()) => Err(error),
166 Err(rollback) => {
167 Err(error.context(format!("discard unpublished session: {rollback:#}")))
168 }
169 };
170 }
171 provision?;
172 Ok(DaemonLifecycleResult::Done)
173 },
174 )?;
175 self.reload_controller().await?;
176 Ok(registered)
177 }
178
179 pub async fn wait_create_session(&self, session_id: &str) -> Result<()> {
180 let result = {
181 let lifecycle = self
182 .lifecycle
183 .lock()
184 .unwrap_or_else(PoisonError::into_inner);
185 let active = lifecycle
186 .get(session_id)
187 .with_context(|| format!("no create operation exists for session {session_id}"))?;
188 ensure!(
189 active.kind == LifecycleKind::Create,
190 "session {session_id} is no longer being created"
191 );
192 active.result.clone()
193 };
194 let channel = result.clone();
195 let outcome = Self::wait_lifecycle_result(result).await;
196 self.remove_completed_lifecycle(&channel);
197 match outcome? {
198 DaemonLifecycleResult::Done => Ok(()),
199 DaemonLifecycleResult::Move(_) => unreachable!("cleanup cannot return a move outcome"),
200 DaemonLifecycleResult::DeferredCleanup => {
201 unreachable!("session creation cannot schedule target cleanup")
202 }
203 }
204 }
205}