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 mjolnir_subagents: request.mjolnir_subagents,
105 initial_prompt: request.initial_prompt,
106 workspace_id: request.workspace_id,
107 additional_mounts: request.additional_mounts,
108 resource_allocation: request.resource_allocation,
109 project_directory,
110 session_title_override: request.session_title_override,
111 },
112 )?;
113 if let Err(error) = mj_core::go::GoPreferences::remember_first_pair(
118 &mj_core::go::GoPreferences::path(),
119 &request.profile_id,
120 &request.target_template_id,
121 ) {
122 tracing::warn!(%error, "could not save the default profile and target");
123 }
124 let session = controller
125 .state
126 .sessions
127 .get(&session_id)
128 .expect("newly registered session exists")
129 .clone();
130 let remembered_container_size = controller
131 .config
132 .targets
133 .get(&request.target_template_id)
134 .and_then(mj_core::config::container_size_host)
135 .and_then(|host| {
136 controller
137 .state
138 .container_sizes
139 .get(host)
140 .copied()
141 .map(|size| (host.to_owned(), size))
142 });
143 Ok(RegisteredSession {
144 session,
145 remembered_container_size,
146 })
147 })
148 .await?;
149 let session_id = registered.session.id.clone();
150 self.start_or_join_lifecycle_controlled(
151 session_id,
152 LifecycleKind::Create,
153 None,
154 None,
155 Some(control.clone()),
156 move |state, session_id, cancelled| async move {
157 let mut controller = tokio::task::spawn_blocking(Controller::load)
158 .await
159 .context("load controller for daemon create task")??;
160 let publication_error = if let Some(publication) = publication {
161 let published = tokio::select! {
162 result = publication => result.context("session publication owner stopped")
163 .and_then(|result| result.map_err(anyhow::Error::msg)),
164 () = async {
165 while !cancelled.load(Ordering::Acquire) {
166 tokio::time::sleep(Duration::from_millis(25)).await;
167 }
168 } => Err(anyhow!("session creation cancelled before publication")),
169 };
170 published.err()
171 } else {
172 None
173 };
174 if publication_error.is_some() {
175 control.request_cancel();
176 }
177 let executor = DaemonStageReportingExecutor::new(
178 CancellableProcessExecutor::new(cancelled),
179 state,
180 session_id.clone(),
181 );
182 let provision = controller
183 .provision_session_controlled_with_commit(&session_id, &executor, || {
184 ensure!(
185 control.grant_commit(),
186 "session creation cancelled before commit"
187 );
188 Ok(())
189 })
190 .await;
191 if let Some(error) = publication_error {
192 return match provision {
193 Ok(()) => Err(error),
194 Err(rollback) => {
195 Err(error.context(format!("discard unpublished session: {rollback:#}")))
196 }
197 };
198 }
199 provision?;
200 Ok(DaemonLifecycleResult::Done)
201 },
202 )?;
203 self.reload_controller().await?;
204 Ok(registered)
205 }
206
207 pub async fn wait_create_session(&self, session_id: &str) -> Result<()> {
208 let result = {
209 let lifecycle = self
210 .lifecycle
211 .lock()
212 .unwrap_or_else(PoisonError::into_inner);
213 let active = lifecycle
214 .get(session_id)
215 .with_context(|| format!("no create operation exists for session {session_id}"))?;
216 ensure!(
217 active.kind == LifecycleKind::Create,
218 "session {session_id} is no longer being created"
219 );
220 active.result.clone()
221 };
222 let channel = result.clone();
223 let outcome = Self::wait_lifecycle_result(result).await;
224 self.remove_completed_lifecycle(&channel);
225 match outcome? {
226 DaemonLifecycleResult::Done => Ok(()),
227 DaemonLifecycleResult::Move(_) => unreachable!("cleanup cannot return a move outcome"),
228 DaemonLifecycleResult::DeferredCleanup => {
229 unreachable!("session creation cannot schedule target cleanup")
230 }
231 }
232 }
233}