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