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