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