1use super::*;
2
3impl RuntimeState {
4 pub(super) fn cancel_lifecycle(&self, session_id: &str) -> Result<()> {
5 let mut controller_owner = self.owner();
6 let durable = durable_session_state(controller_owner.controller(), session_id);
7 let active = controller_owner
8 .lifecycle
9 .get_mut(session_id)
10 .with_context(|| {
11 format!("no lifecycle operation is running for session {session_id}")
12 })?;
13 ensure!(
14 lifecycle_cancellable(active.kind, durable),
15 "stop of {session_id} has passed its verified checkpoint and is removing the target; \
16 it cannot be cancelled"
17 );
18 ensure!(
19 active.request_cancel(),
20 "lifecycle operation is no longer cancellable"
21 );
22 drop(controller_owner);
23 self.publish_revision();
24 Ok(())
25 }
26
27 pub(super) async fn cancel_and_wait_lifecycles(&self) -> Result<()> {
32 let mut pending = {
33 let mut lifecycle_owner = self.owner();
34 let lifecycle = &mut lifecycle_owner.lifecycle;
35 lifecycle
36 .iter_mut()
37 .filter(|(_, active)| active.is_running())
38 .map(|(session_id, active)| {
39 if active.kind != LifecycleKind::Cleanup {
40 active.request_cancel();
41 }
42 let stage = active
43 .active_stages
44 .keys()
45 .next_back()
46 .map(|stage| stage.label())
47 .unwrap_or_else(|| "container cleanup".to_owned());
48 (
49 session_id.clone(),
50 active.operation_id.clone(),
51 active.kind,
52 stage,
53 active.result.clone(),
54 )
55 })
56 .collect::<Vec<_>>()
57 };
58 let cleanup_deadline = tokio::time::Instant::now() + Duration::from_secs(8);
59 for (session_id, operation_id, kind, stage, result) in &mut pending {
60 if *kind != LifecycleKind::Cleanup || result.borrow().is_some() {
61 continue;
62 }
63 tracing::info!(%session_id, %stage, "daemon shutdown is waiting for deferred cleanup");
64 self.set_lifecycle_notice(
65 session_id,
66 operation_id,
67 &format!("Daemon shutdown is waiting for {stage}"),
68 );
69 let finished = tokio::time::timeout_at(cleanup_deadline, async {
70 while result.borrow_and_update().is_none() {
71 result.changed().await.with_context(|| {
72 format!("cleanup owner stopped without a result for session {session_id}")
73 })?;
74 }
75 Ok::<_, anyhow::Error>(())
76 })
77 .await;
78 match finished {
79 Ok(result) => result?,
80 Err(_) => {
81 tracing::warn!(%session_id, %stage, "deferred cleanup exceeded the daemon shutdown drain deadline");
82 self.cancel_operation(session_id, operation_id);
83 }
84 }
85 }
86 let join_started = tokio::time::Instant::now();
87 let join_deadline = join_started + Duration::from_secs(1);
88 let startup_cleanup_deadline = join_started
91 + crate::controller::FAILED_STARTUP_CLEANUP_TIMEOUT
92 + Duration::from_secs(1);
93 for (session_id, operation_id, kind, stage, mut result) in pending {
94 let join_deadline = if matches!(
95 kind,
96 LifecycleKind::Create
97 | LifecycleKind::Resume
98 | LifecycleKind::Unpark
99 | LifecycleKind::StartupCleanup
100 ) {
101 startup_cleanup_deadline
102 } else {
103 join_deadline
104 };
105 self.cancel_operation(&session_id, &operation_id);
106 let joined = tokio::time::timeout_at(join_deadline, async {
107 while result.borrow_and_update().is_none() {
108 result.changed().await.with_context(|| {
109 format!("lifecycle owner stopped without a result for session {session_id}")
110 })?;
111 }
112 Ok::<_, anyhow::Error>(())
113 })
114 .await;
115 if joined.is_err() {
116 bail!(
117 "timed out cancelling lifecycle owner for session {session_id} while {stage}"
118 );
119 }
120 joined.expect("checked timeout")?;
121 }
122 Ok(())
123 }
124
125 pub fn active_lifecycles(&self) -> Vec<RuntimeLifecycleView> {
134 let controller_owner = self.owner();
135 Self::active_lifecycles_with(&controller_owner)
136 }
137
138 pub(super) fn active_lifecycles_with(owner: &RuntimeStateOwner) -> Vec<RuntimeLifecycleView> {
140 let controller = owner.controller();
141 owner
142 .lifecycle
143 .iter()
144 .filter(|(_, active)| active.is_visible())
145 .map(|(session_id, active)| RuntimeLifecycleView {
146 operation_id: active.operation_id.clone(),
147 cancellable: active.is_cancellable()
148 && lifecycle_cancellable(
149 active.kind,
150 durable_session_state(controller, session_id),
151 ),
152 session_id: session_id.clone(),
153 kind: active.kind.into(),
154 started_at_epoch_seconds: active.started_at_epoch_seconds,
155 active_stages: active
156 .active_stages
157 .iter()
158 .map(|(stage, (_, started_at))| (*stage, *started_at))
159 .collect(),
160 resume_destination: active.resume_destination.clone(),
161 notice: active.notice.clone(),
162 })
163 .collect()
164 }
165
166 pub fn session_state(&self, session_id: &str) -> Option<mj_core::state::SessionState> {
170 let owner = self.owner();
171 if owner.close_requested.contains(session_id) {
172 return Some(SessionState::Closing);
173 }
174 owner
175 .controller()
176 .state
177 .sessions
178 .get(session_id)
179 .map(|record| record.state)
180 }
181
182 pub fn session_record(&self, session_id: &str) -> Option<SessionRecord> {
184 self.owner()
185 .controller()
186 .state
187 .sessions
188 .get(session_id)
189 .cloned()
190 }
191
192 pub async fn workspace_session_handle(
193 &self,
194 session_id: &str,
195 ) -> Result<crate::session_manager::ManagedSessionHandle> {
196 let record = self.session_record(session_id).context("unknown session")?;
197 ensure!(
198 record.target.is_some()
199 && record.state == SessionState::Running
200 && !self.close_is_requested(session_id),
201 "session must have a live running target for file injection"
202 );
203 self.session_manager.session(session_id.to_owned()).await
204 }
205
206 pub async fn checkpoint_session_now(
216 &self,
217 session_id: &str,
218 ) -> Result<mj_core::state::CheckpointMetadata> {
219 let _upgrade_work = crate::upgrade::activity("requested checkpoint")?;
220 if let Some(busy) = self.session_lifecycle_busy(session_id) {
221 return Err(anyhow::Error::new(busy));
222 }
223 let session_id = session_id.to_owned();
224 let checkpoint = blocking(move || {
228 let mut controller = Controller::load()?;
229 mj_core::runtime::block_on(controller.checkpoint_session(&session_id))?
230 })
231 .await?;
232 refresh_runtime_controller(self).await;
233 Ok(checkpoint)
234 }
235
236 pub(super) fn session_lifecycle_busy(&self, session_id: &str) -> Option<SessionLifecycleBusy> {
240 let lifecycle_owner = self.owner();
241 let lifecycle = &lifecycle_owner.lifecycle;
242 let active = lifecycle.get(session_id)?;
243 active
244 .is_running()
245 .then(|| describe_lifecycle_busy(session_id, active))
246 }
247
248 pub(super) fn any_lifecycle_busy(&self) -> Option<SessionLifecycleBusy> {
252 self.owner()
253 .lifecycle
254 .iter()
255 .find(|(_, active)| active.is_running())
256 .map(|(session_id, active)| describe_lifecycle_busy(session_id, active))
257 }
258
259 pub fn session_projection(
263 &self,
264 ) -> (
265 mj_core::snapshot_map::SnapshotMap<String, SessionRecord>,
266 Vec<RuntimeLifecycleView>,
267 ) {
268 let controller_owner = self.owner();
269 let operations = Self::active_lifecycles_with(&controller_owner);
270 (controller_owner.projected_records(), operations)
271 }
272
273 pub(super) fn resume_candidates(&self) -> mj_client::daemon::ResumeCandidates {
277 let (records, subagents, moves, config) = {
278 let owner = self.owner();
279 (
280 owner.projected_records(),
281 owner.controller().state.subagents.clone(),
282 owner
283 .committed()
284 .map(|committed| committed.moves.clone())
285 .unwrap_or_default(),
286 owner.controller().config.clone(),
287 )
288 };
289 let mut candidates = mj_client::daemon::ResumeCandidates::default();
290 for (id, record) in &records {
291 if let Some(native_session_id) = &record.native_session_id {
292 candidates
293 .adopted_native_sessions
294 .push((record.harness_kind, native_session_id.clone()));
295 }
296 if let Some(checkout) = &record.managed_worktree
297 && checkout.target == mj_core::state::ManagedWorktreeTarget::Local
298 {
299 candidates
300 .local_checkout_roots
301 .push(checkout.worktree_root.clone());
302 }
303 if record.state.is_active()
304 || subagents.contains_key(id)
305 || mj_core::native_agent::is_view_id(id)
306 {
307 continue;
308 }
309 if let Some(operation) = moves.get(id) {
310 candidates.moves.push(operation.clone());
311 }
312 candidates
313 .candidates
314 .push(mj_client::daemon::ResumeCandidate::of(record, &config));
315 }
316 candidates
317 }
318
319 pub(super) fn go_startup_session(
322 &self,
323 workspace_id: &str,
324 last_session_id: Option<&str>,
325 ) -> Option<SessionRecord> {
326 let (records, subagents) = {
327 let owner = self.owner();
328 (
329 owner.projected_records(),
330 owner.controller().state.subagents.clone(),
331 )
332 };
333 let eligible = |session: &&SessionRecord| {
334 session.workspace_id == workspace_id
335 && !session.archived
336 && !subagents.contains_key(&session.id)
337 && !mj_core::native_agent::is_view_id(&session.id)
338 && session.state != SessionState::DestroyedWithDataLoss
339 };
340 last_session_id
341 .and_then(|id| records.get(id))
342 .filter(eligible)
343 .or_else(|| {
344 records
345 .values()
346 .filter(eligible)
347 .max_by_key(|session| &session.updated_at)
348 })
349 .cloned()
350 }
351
352 pub(crate) fn worker_controller_projection(&self) -> Controller {
353 self.owner().pollable_worker_inputs().controller()
354 }
355
356 pub(crate) fn controller_projection(&self) -> Controller {
357 let owner = self.owner();
358 let mut state = owner.controller().state.clone();
359 state.sessions = owner.projected_records();
360 Controller {
361 config: owner.controller().config.clone(),
362 state,
363 }
364 }
365
366 pub(crate) fn active_controller_projection(&self) -> Controller {
367 let owner = self.owner();
368 Controller {
369 config: owner.controller().config.clone(),
370 state: mj_core::state::State {
371 sessions: owner
372 .indexes
373 .active
374 .keys()
375 .filter_map(|id| {
376 owner
377 .controller()
378 .state
379 .sessions
380 .get(id)
381 .map(|record| (id.clone(), record.clone()))
382 })
383 .collect(),
384 ..Default::default()
385 },
386 }
387 }
388
389 pub fn cancel_lifecycle_if_active(&self, session_id: &str) {
390 if let Some(active) = self.owner().lifecycle.get_mut(session_id) {
391 active.request_cancel();
392 self.publish_revision();
393 }
394 }
395
396 pub(super) fn set_lifecycle_resume_destination(
397 &self,
398 session_id: &str,
399 profile_id: String,
400 target_id: String,
401 ) {
402 if let Some(active) = self.owner().lifecycle.get_mut(session_id) {
403 active.resume_destination = Some((profile_id, target_id));
404 self.publish_revision();
405 }
406 }
407
408 pub(super) fn change_lifecycle_stage(
409 &self,
410 session_id: &str,
411 operation_id: &str,
412 stage: ProvisionStage,
413 active: bool,
414 ) {
415 let changed = {
416 let mut lifecycle_owner = self.owner();
417 let lifecycle = &mut lifecycle_owner.lifecycle;
418 let Some(operation) = lifecycle.get_mut(session_id) else {
419 return;
420 };
421 if operation.operation_id != operation_id || !operation.is_running() {
422 return;
423 }
424 if active {
425 let entry = operation
426 .active_stages
427 .entry(stage)
428 .or_insert_with(|| (0, epoch_seconds()));
429 entry.0 += 1;
430 entry.0 == 1
431 } else {
432 let Some((count, _)) = operation.active_stages.get_mut(&stage) else {
433 return;
434 };
435 *count -= 1;
436 if *count == 0 {
437 operation.active_stages.remove(&stage);
438 true
439 } else {
440 false
441 }
442 }
443 };
444 if changed {
445 self.publish_revision();
446 }
447 }
448
449 pub(crate) fn push_notice(&self, session_id: &str, text: impl Into<String>) {
452 const RETAINED_NOTICES: usize = 32;
453
454 let notice = RuntimeNotice {
455 id: self.next_notice_id.fetch_add(1, Ordering::AcqRel),
456 session_id: session_id.to_owned(),
457 text: text.into(),
458 };
459 {
460 let mut notices = self.notices.lock().unwrap_or_else(PoisonError::into_inner);
461 notices.push_back(notice);
462 while notices.len() > RETAINED_NOTICES {
463 notices.pop_front();
464 }
465 }
466 self.publish_revision();
467 }
468
469 pub(crate) fn attach_quota_refresh(&self, refresh: tokio::sync::mpsc::Sender<()>) {
473 self.quota
474 .lock()
475 .unwrap_or_else(PoisonError::into_inner)
476 .refresh = Some(refresh);
477 }
478
479 pub(crate) fn request_quota_refresh(&self) -> Result<()> {
482 let quota = self.quota.lock().unwrap_or_else(PoisonError::into_inner);
483 let refresh = quota
484 .refresh
485 .as_ref()
486 .context("the daemon's quota service is not running")?;
487 let _ = refresh.try_send(());
489 Ok(())
490 }
491
492 pub(crate) fn publish_quotas(&self, snapshot: mj_client::quota::QuotaSnapshot) {
495 {
496 let mut quota = self.quota.lock().unwrap_or_else(PoisonError::into_inner);
497 if quota.snapshot == snapshot {
498 return;
499 }
500 quota.snapshot = snapshot;
501 }
502 self.publish_revision();
503 }
504
505 pub(super) fn reserve_move_destination(&self, session_id: &str, operation_id: &str) {
506 let mut owner = self.owner();
507 if let Some(active) = owner.lifecycle.get_mut(session_id)
508 && active.operation_id == operation_id
509 && active.kind == LifecycleKind::Move
510 {
511 active.phase = match active.phase {
512 LifecyclePhase::Executing => LifecyclePhase::MovingDestination,
513 LifecyclePhase::Cancelling => LifecyclePhase::CancellingMoveDestination,
514 _ => return,
515 };
516 self.publish_revision();
517 }
518 }
519
520 pub(super) fn set_lifecycle_notice(&self, session_id: &str, operation_id: &str, notice: &str) {
521 if let Some(active) = self.owner().lifecycle.get_mut(session_id)
522 && active.operation_id == operation_id
523 && active.is_running()
524 {
525 active.notice = Some(notice.to_owned());
526 self.publish_revision();
527 }
528 }
529
530 fn cancel_operation(&self, session_id: &str, operation_id: &str) {
531 if let Some(active) = self.owner().lifecycle.get_mut(session_id)
532 && active.operation_id == operation_id
533 && active.request_cancel()
534 {
535 self.publish_revision();
536 }
537 }
538}