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