1mod session_move;
4use crate::controller::move_session::{
5 MoveMutationGuard, MoveOutcome, MovePreparation, MoveSelection, MoveSessionRequest,
6};
7pub use mj_client::daemon::*;
8
9use std::collections::{BTreeMap, BTreeSet, VecDeque};
10use std::fs::{self, OpenOptions};
11use std::io::Write;
12use std::net::{IpAddr, Ipv4Addr, SocketAddr};
13use std::path::Path;
14use std::sync::atomic::{AtomicBool, AtomicU8, AtomicU64, Ordering};
15use std::sync::{Arc, Mutex, PoisonError};
16use std::time::{Duration, SystemTime, UNIX_EPOCH};
17
18use crate::database::StoreSchemaMismatch;
19use crate::recovery_gate::RecoveryObserver;
20use crate::targets::{
21 CancellableProcessExecutor, CommandExecutor, CommandOutput, CommandSpec, ProcessExecutor,
22 ProvisionStage, ProvisionStageGuard,
23};
24use anyhow::{Context, Result, anyhow, bail, ensure};
25use mj_core::config::Config;
26use mj_core::relay::RelayCommand;
27use mj_core::state::{RecoveryObservation, SessionRecord, SessionState};
28
29use crate::controller::{
30 Controller, ControllerStoreGuard, SessionLaunchOptions, SessionResumeOptions,
31};
32use crate::review_host::TurnReviewHost;
33use crate::session_manager::{
34 ManagedSessionView, RemoteSessionPublisher, RemoteSessionRequest, SessionManagerChannels,
35 SessionManagerControl, ViewError, new_command_id, spawn_remote_session_manager,
36 spawn_session_manager,
37};
38#[cfg(test)]
39use crate::session_manager::{RelaySessionTarget, RemoteSessionRequests, SessionManagerShutdown};
40use crate::worker_upgrade::{WorkerUpgradeObservation, WorkerUpgradeObserver};
41use mj_core::workspace::WorkspaceRecord;
42use tokio::net::{TcpListener, TcpStream};
43use tokio_util::sync::CancellationToken;
44
45use crate::pollers::{
46 dashboard_worker_targets, dashboard_worker_targets_excluding, interrupted_close_session_ids,
47 reserve_recovery_or_cancel, spawn_image_refresher, spawn_interrupted_close_recovery,
48};
49
50const SHUTDOWN_FORCE_EXIT_TIMEOUT: Duration = Duration::from_secs(10);
60
61const FORCE_DESTROY_PREEMPT_TIMEOUT: Duration = Duration::from_secs(8);
67
68#[derive(Clone, Default)]
70pub struct CreateSessionControl {
71 state: Arc<AtomicU8>,
72 pub cancelled: Arc<AtomicBool>,
73}
74
75impl CreateSessionControl {
76 pub fn request_cancel(&self) -> bool {
77 let accepted = self
78 .state
79 .compare_exchange(0, 1, Ordering::AcqRel, Ordering::Acquire)
80 .is_ok();
81 if accepted {
82 self.cancelled.store(true, Ordering::Release);
83 }
84 accepted
85 }
86
87 pub fn grant_commit(&self) -> bool {
88 self.state
89 .compare_exchange(0, 2, Ordering::AcqRel, Ordering::Acquire)
90 .is_ok()
91 }
92
93 fn is_cancellable(&self) -> bool {
94 self.state.load(Ordering::Acquire) == 0
95 }
96}
97
98#[derive(Debug, Clone)]
99struct Attachment {
100 pid: u32,
101}
102
103pub struct RuntimeState {
104 attachments: Mutex<BTreeMap<String, Attachment>>,
105 phone_status: Mutex<WebViewerStatus>,
106 pub web_viewer: crate::web_viewer::ViewerControl,
107 ever_attached: AtomicBool,
108 sessions: Mutex<BTreeMap<String, RuntimeSessionView>>,
109 revisions: RuntimeRevisions,
110 workspaces_tx: tokio::sync::watch::Sender<Vec<WorkspaceRecord>>,
111 session_manager: SessionManagerControl,
112 lifecycle: Mutex<BTreeMap<String, ActiveLifecycle>>,
113 close_requested: Mutex<BTreeSet<String>>,
114 controller: Mutex<Controller>,
115 controller_loader: fn() -> Result<Controller>,
116 config_mutation: tokio::sync::Mutex<()>,
117 recovery_observer: RecoveryObserver,
118 worker_upgrade_observer: WorkerUpgradeObserver,
119 notices: Mutex<VecDeque<RuntimeNotice>>,
122 next_notice_id: AtomicU64,
123 review_config: Arc<Mutex<mj_core::config::ReviewConfig>>,
125 review_host: TurnReviewHost,
128}
129
130#[derive(Clone)]
136struct RuntimeRevisions {
137 allocated: Arc<std::sync::atomic::AtomicU64>,
138 published: tokio::sync::watch::Sender<u64>,
139}
140
141impl RuntimeRevisions {
142 fn new(initial: u64) -> Self {
143 let (published, _) = tokio::sync::watch::channel(initial);
144 Self {
145 allocated: Arc::new(std::sync::atomic::AtomicU64::new(initial)),
146 published,
147 }
148 }
149
150 fn allocate(&self) -> u64 {
151 self.allocated.fetch_add(1, Ordering::AcqRel) + 1
152 }
153
154 fn publish(&self) -> u64 {
155 let revision = self.allocate();
156 self.publish_allocated(revision);
157 revision
158 }
159
160 fn publish_allocated(&self, revision: u64) {
161 self.published.send_if_modified(|visible| {
162 if revision > *visible {
163 *visible = revision;
164 true
165 } else {
166 false
167 }
168 });
169 }
170
171 fn notifier(&self) -> Arc<dyn Fn() + Send + Sync> {
172 let revisions = self.clone();
173 Arc::new(move || {
174 revisions.publish();
175 })
176 }
177
178 fn subscribe(&self) -> tokio::sync::watch::Receiver<u64> {
179 self.published.subscribe()
180 }
181
182 fn current(&self) -> u64 {
183 self.allocated.load(Ordering::Acquire)
184 }
185}
186
187#[derive(Debug, Clone, Copy, PartialEq, Eq)]
188enum LifecycleKind {
189 Create,
190 Close,
191 Resume,
192 Move,
193 ForceStop,
194 DestroyStopped,
195 ForceDestroy,
196 Cleanup,
197}
198
199fn lifecycle_owns_worker_target(kind: LifecycleKind, state: Option<SessionState>) -> bool {
204 match kind {
205 LifecycleKind::Close => state == Some(SessionState::Destroying),
206 LifecycleKind::Move => !matches!(
207 state,
208 Some(
209 SessionState::Running
210 | SessionState::Disconnected
211 | SessionState::Checkpointing
212 | SessionState::Closing
213 )
214 ),
215 _ => true,
216 }
217}
218
219struct ActiveLifecycle {
220 operation_id: String,
221 create_control: Option<CreateSessionControl>,
222 kind: LifecycleKind,
223 cancelled: Arc<AtomicBool>,
224 started_at_epoch_seconds: u64,
225 active_stages: BTreeMap<ProvisionStage, (usize, u64)>,
226 resume_workspace_id: Option<String>,
229 resume_destination: Option<(String, String)>,
230 notice: Option<String>,
231 request_key: Option<String>,
232 _move_guard: Option<MoveMutationGuard>,
233 move_source_closed: bool,
234 result:
235 tokio::sync::watch::Receiver<Option<std::result::Result<DaemonLifecycleResult, String>>>,
236}
237
238impl ActiveLifecycle {
239 fn is_visible(&self) -> bool {
240 let result = self.result.borrow();
241 result.is_none()
242 || matches!(
243 result.as_ref(),
244 Some(Ok(DaemonLifecycleResult::DeferredCleanup))
245 )
246 }
247
248 fn request_cancel(&self) -> bool {
249 if let Some(control) = &self.create_control {
250 control.request_cancel()
251 } else {
252 !self.cancelled.swap(true, Ordering::AcqRel)
253 }
254 }
255
256 fn is_cancellable(&self) -> bool {
257 self.result.borrow().is_none()
258 && self.create_control.as_ref().map_or_else(
259 || !self.cancelled.load(Ordering::Acquire),
260 CreateSessionControl::is_cancellable,
261 )
262 }
263}
264
265#[derive(Debug, Clone)]
266enum DaemonLifecycleResult {
267 Done,
268 DeferredCleanup,
269 Move(MoveOutcome),
270}
271
272impl From<LifecycleKind> for RuntimeLifecycleKind {
273 fn from(kind: LifecycleKind) -> Self {
274 match kind {
275 LifecycleKind::Create => Self::Create,
276 LifecycleKind::Close => Self::Close,
277 LifecycleKind::Resume => Self::Resume,
278 LifecycleKind::Move => Self::Move,
279 LifecycleKind::ForceStop => Self::ForceStop,
280 LifecycleKind::DestroyStopped => Self::DestroyStopped,
281 LifecycleKind::ForceDestroy => Self::ForceDestroy,
282 LifecycleKind::Cleanup => Self::Cleanup,
283 }
284 }
285}
286
287impl RuntimeState {
288 fn new(
289 session_manager: SessionManagerControl,
290 controller: Controller,
291 recovery_observer: RecoveryObserver,
292 worker_upgrade_observer: WorkerUpgradeObserver,
293 workspaces: Vec<WorkspaceRecord>,
294 ) -> Self {
295 Self::new_with_controller_loader(
296 session_manager,
297 controller,
298 recovery_observer,
299 worker_upgrade_observer,
300 workspaces,
301 Controller::load,
302 )
303 }
304
305 fn new_with_controller_loader(
306 session_manager: SessionManagerControl,
307 controller: Controller,
308 recovery_observer: RecoveryObserver,
309 worker_upgrade_observer: WorkerUpgradeObserver,
310 workspaces: Vec<WorkspaceRecord>,
311 controller_loader: fn() -> Result<Controller>,
312 ) -> Self {
313 let initial_revision = u64::try_from(chrono::Utc::now().timestamp_micros()).unwrap_or(1);
318 let revisions = RuntimeRevisions::new(initial_revision);
319 let (workspaces_tx, _) = tokio::sync::watch::channel(workspaces);
320 let review_config = Arc::new(Mutex::new(controller.config.review.clone()));
324 let review_host = TurnReviewHost::spawn_notifying(
325 session_manager.clone(),
326 {
327 let installed = review_config.clone();
328 Arc::new(move || {
329 installed
330 .lock()
331 .unwrap_or_else(PoisonError::into_inner)
332 .clone()
333 })
334 },
335 revisions.notifier(),
336 );
337 Self {
338 attachments: Mutex::new(BTreeMap::new()),
339 phone_status: Mutex::new(WebViewerStatus::Starting),
340 web_viewer: crate::web_viewer::ViewerControl::new(),
341 ever_attached: AtomicBool::new(false),
342 sessions: Mutex::new(BTreeMap::new()),
343 revisions,
344 workspaces_tx,
345 session_manager,
346 lifecycle: Mutex::new(BTreeMap::new()),
347 close_requested: Mutex::new(BTreeSet::new()),
348 controller: Mutex::new(controller),
349 controller_loader,
350 config_mutation: tokio::sync::Mutex::new(()),
351 recovery_observer,
352 worker_upgrade_observer,
353 notices: Mutex::new(VecDeque::new()),
354 next_notice_id: AtomicU64::new(1),
355 review_config,
356 review_host,
357 }
358 }
359
360 pub fn review_host(&self) -> &TurnReviewHost {
362 &self.review_host
363 }
364
365 pub fn allocate_revision(&self) -> u64 {
366 self.revisions.allocate()
367 }
368
369 fn publish_revision(&self) -> u64 {
370 self.revisions.publish()
371 }
372
373 fn attachments(&self) -> std::sync::MutexGuard<'_, BTreeMap<String, Attachment>> {
374 self.attachments
375 .lock()
376 .unwrap_or_else(PoisonError::into_inner)
377 }
378
379 fn prune_dead_clients(&self) {
380 self.attachments()
381 .retain(|_, attachment| process_is_alive(attachment.pid));
382 }
383
384 fn workspace_has_active_resume(&self, workspace_id: &str) -> bool {
385 self.lifecycle
386 .lock()
387 .unwrap_or_else(PoisonError::into_inner)
388 .values()
389 .any(|active| {
390 active.result.borrow().is_none()
391 && active.resume_workspace_id.as_deref() == Some(workspace_id)
392 })
393 }
394
395 pub fn publish_web_access(&self, access: crate::server::WebViewerAccess) {
396 use crate::server::WebViewerAccess;
397 let status = match &access {
398 WebViewerAccess::Starting => WebViewerStatus::Starting,
399 WebViewerAccess::Ready {
400 viewer_url,
401 viewer_code,
402 qr_login_url,
403 fallback_reason,
404 } => WebViewerStatus::Ready {
405 viewer_url: viewer_url.clone(),
406 viewer_code: viewer_code.clone(),
407 qr_login_url: qr_login_url.clone(),
408 fallback_reason: fallback_reason.clone(),
409 },
410 WebViewerAccess::Failed {
411 address, message, ..
412 } => WebViewerStatus::Error {
413 message: format!("{message} Address: {address}"),
414 },
415 WebViewerAccess::Unavailable(message) => WebViewerStatus::Error {
416 message: message.clone(),
417 },
418 };
419 self.web_viewer.publish(access);
420 self.set_phone_status(status);
421 }
422
423 fn set_phone_status(&self, status: WebViewerStatus) {
424 *self
425 .phone_status
426 .lock()
427 .unwrap_or_else(PoisonError::into_inner) = status;
428 }
429
430 fn phone_status(&self) -> WebViewerStatus {
431 self.phone_status
432 .lock()
433 .unwrap_or_else(PoisonError::into_inner)
434 .clone()
435 }
436
437 fn workspaces(&self) -> tokio::sync::watch::Receiver<Vec<WorkspaceRecord>> {
438 self.workspaces_tx.subscribe()
439 }
440
441 fn worker_poll_exclusion_session_ids(&self, controller: &Controller) -> BTreeSet<String> {
442 self.lifecycle
443 .lock()
444 .unwrap_or_else(PoisonError::into_inner)
445 .iter()
446 .filter(|(session_id, active)| {
447 active.result.borrow().is_none()
448 && (active.move_source_closed
449 || lifecycle_owns_worker_target(
450 active.kind,
451 controller
452 .state
453 .sessions
454 .get(*session_id)
455 .map(|session| session.state),
456 ))
457 })
458 .map(|(session_id, _)| session_id.clone())
459 .collect()
460 }
461
462 pub fn revisions(&self) -> tokio::sync::watch::Receiver<u64> {
463 self.revisions.subscribe()
464 }
465
466 pub fn with_config<T>(&self, read: impl FnOnce(&Config) -> T) -> T {
469 read(
470 &self
471 .controller
472 .lock()
473 .unwrap_or_else(PoisonError::into_inner)
474 .config,
475 )
476 }
477
478 pub async fn create_quick_bundle(
483 &self,
484 source: String,
485 ) -> std::result::Result<
486 crate::controller::QuickBundleCreation,
487 crate::controller::QuickBundleFailure,
488 > {
489 let _mutation = self.config_mutation.lock().await;
490 tokio::task::spawn_blocking(move || crate::controller::create_quick_bundle(&source))
491 .await
492 .map_err(|error| {
493 crate::controller::QuickBundleFailure::Persistence(anyhow!(
494 "bundle creation task panicked: {error}"
495 ))
496 })?
497 }
498
499 fn publish_workspaces(&self, workspaces: Vec<WorkspaceRecord>) {
500 self.workspaces_tx.send_replace(workspaces);
501 self.publish_revision();
502 }
503
504 pub async fn reload_controller(&self) -> Result<()> {
505 let _mutation = self.config_mutation.lock().await;
508 let controller_loader = self.controller_loader;
509 let controller = tokio::task::spawn_blocking(controller_loader)
510 .await
511 .context("daemon controller reload task panicked")??;
512 let session_count = controller.state.sessions.len();
513 *self
514 .controller
515 .lock()
516 .unwrap_or_else(PoisonError::into_inner) = controller;
517 let revision = self.publish_revision();
518 tracing::debug!(revision, session_count, "daemon controller state reloaded");
519 Ok(())
520 }
521
522 fn missing_target_record(
525 &self,
526 session_id: &str,
527 view: &ManagedSessionView,
528 ) -> Option<(String, String)> {
529 let Some(ViewError::TargetMissing(detail)) = &view.error else {
530 return None;
531 };
532 if view.connected {
533 return None;
534 }
535 if self
536 .lifecycle
537 .lock()
538 .unwrap_or_else(PoisonError::into_inner)
539 .get(session_id)
540 .is_some_and(|active| active.result.borrow().is_none())
541 {
542 return None;
543 }
544 let controller = self
545 .controller
546 .lock()
547 .unwrap_or_else(PoisonError::into_inner);
548 let session = controller.state.sessions.get(session_id)?;
549 if !matches!(
550 session.state,
551 SessionState::Running | SessionState::Disconnected
552 ) {
553 return None;
554 }
555 Some((detail.clone(), session.updated_at.clone()))
556 }
557
558 async fn persist_missing_target(
559 &self,
560 session_id: &str,
561 detail: String,
562 observed_updated_at: String,
563 ) -> Result<()> {
564 let changed = blocking({
565 let session_id = session_id.to_owned();
566 let detail = detail.clone();
567 move || {
568 crate::database::mark_session_target_missing_if_current(
569 &session_id,
570 &detail,
571 &chrono::Utc::now().to_rfc3339_opts(chrono::SecondsFormat::Secs, true),
572 &observed_updated_at,
573 )
574 }
575 })
576 .await?;
577 if changed.is_some() {
578 self.reload_controller().await?;
579 self.push_notice(session_id, detail);
580 }
581 Ok(())
582 }
583
584 async fn publish_session(&self, session_id: String, view: ManagedSessionView) -> Result<()> {
585 let connected = view.connected;
586 let has_snapshot = view.snapshot.is_some();
587 tracing::debug!(
588 %session_id,
589 connected,
590 has_snapshot,
591 "daemon received a session view"
592 );
593 if let Some(snapshot) = view.snapshot.as_ref() {
594 let controller = self
595 .controller
596 .lock()
597 .unwrap_or_else(PoisonError::into_inner);
598 if let Some(session) = controller.state.sessions.get(&session_id).cloned() {
599 let quiet =
604 view.connected && snapshot.operational.safe_to_replace(session.harness_kind);
605 if quiet {
606 self.worker_upgrade_observer
607 .observe(WorkerUpgradeObservation {
608 session: session.clone(),
609 config: controller.config.clone(),
610 worker_build: snapshot.worker_build.clone(),
611 quiet,
612 });
613 }
614 self.recovery_observer.observe(RecoveryObservation {
615 checkpoint_safe: snapshot
616 .operational
617 .safe_for_checkpoint(session.harness_kind),
618 session,
619 config: controller.config.clone(),
620 latest_completed_turn_ordinal: snapshot.latest_completed_turn_ordinal(),
621 execution: snapshot.materialized.execution,
622 });
623 }
624 }
625 self.sessions
626 .lock()
627 .unwrap_or_else(PoisonError::into_inner)
628 .insert(
629 session_id.clone(),
630 RuntimeSessionView::from_managed(session_id, view),
631 );
632 reach_test_hook("relay_projection_before_revision_publication").await?;
633 self.publish_revision();
634 Ok(())
635 }
636
637 async fn runtime_snapshot(
638 &self,
639 workspace_id: &str,
640 after_revision: u64,
641 all_workspaces: bool,
642 ) -> Result<RuntimeSnapshot> {
643 let mut revisions = self.revisions.subscribe();
644 if *revisions.borrow_and_update() <= after_revision {
645 let _ = tokio::time::timeout(Duration::from_secs(30), revisions.changed()).await;
646 }
647 let revision = self.revisions.current();
648 let moves = blocking(crate::database::load_move_operations).await?;
649 let workspace_names = blocking(crate::database::list_workspaces)
650 .await?
651 .into_iter()
652 .map(|workspace| (workspace.id, workspace.name))
653 .collect();
654 let session_ids = if all_workspaces {
655 self.controller
656 .lock()
657 .unwrap_or_else(PoisonError::into_inner)
658 .state
659 .sessions
660 .keys()
661 .cloned()
662 .collect::<BTreeSet<_>>()
663 } else {
664 let workspace_id = workspace_id.to_owned();
665 blocking(move || crate::database::session_ids_for_workspace(&workspace_id))
666 .await?
667 .into_iter()
668 .collect()
669 };
670 let sessions = self
671 .sessions
672 .lock()
673 .unwrap_or_else(PoisonError::into_inner)
674 .iter()
675 .filter(|(session_id, _)| session_ids.contains(*session_id))
676 .map(|(_, view)| view.clone())
677 .collect();
678 let controller = self
682 .controller
683 .lock()
684 .unwrap_or_else(PoisonError::into_inner);
685 let lifecycles = self
686 .lifecycle
687 .lock()
688 .unwrap_or_else(PoisonError::into_inner)
689 .iter()
690 .filter(|(session_id, active)| {
691 (all_workspaces
692 || session_ids.contains(*session_id)
693 || active.resume_workspace_id.as_deref() == Some(workspace_id))
694 && active.is_visible()
695 })
696 .map(|(session_id, active)| RuntimeLifecycleView {
697 operation_id: active.operation_id.clone(),
698 cancellable: active.is_cancellable(),
699 session_id: session_id.clone(),
700 kind: active.kind.into(),
701 started_at_epoch_seconds: active.started_at_epoch_seconds,
702 active_stages: active
703 .active_stages
704 .iter()
705 .map(|(stage, (_, started_at))| (*stage, *started_at))
706 .collect(),
707 resume_destination: active.resume_destination.clone(),
708 notice: active.notice.clone(),
709 })
710 .collect();
711 let reviews = self
712 .review_host
713 .views()
714 .into_iter()
715 .filter(|review| session_ids.contains(&review.session_id))
716 .collect();
717 let notices = self
718 .notices
719 .lock()
720 .unwrap_or_else(PoisonError::into_inner)
721 .iter()
722 .filter(|notice| session_ids.contains(¬ice.session_id))
723 .cloned()
724 .collect();
725 let records = runtime_records_for_workspace(&controller, &session_ids);
726 Ok(RuntimeSnapshot {
727 workspace_names,
728 moves: moves
729 .into_iter()
730 .filter(|operation| session_ids.contains(&operation.selection.session_id))
731 .collect(),
732 revision,
733 config: controller.config.clone(),
734 records,
735 sessions,
736 lifecycles,
737 reviews,
738 notices,
739 })
740 }
741
742 fn start_or_join_lifecycle<F, Fut>(
743 self: &Arc<Self>,
744 session_id: String,
745 kind: LifecycleKind,
746 work: F,
747 ) -> Result<
748 tokio::sync::watch::Receiver<Option<std::result::Result<DaemonLifecycleResult, String>>>,
749 >
750 where
751 F: FnOnce(Arc<Self>, String, Arc<AtomicBool>) -> Fut + Send + 'static,
752 Fut: std::future::Future<Output = Result<DaemonLifecycleResult>> + Send + 'static,
753 {
754 self.start_or_join_lifecycle_for_workspace(session_id, kind, None, work)
755 }
756
757 fn start_or_join_lifecycle_for_workspace<F, Fut>(
758 self: &Arc<Self>,
759 session_id: String,
760 kind: LifecycleKind,
761 resume_workspace_id: Option<String>,
762 work: F,
763 ) -> Result<
764 tokio::sync::watch::Receiver<Option<std::result::Result<DaemonLifecycleResult, String>>>,
765 >
766 where
767 F: FnOnce(Arc<Self>, String, Arc<AtomicBool>) -> Fut + Send + 'static,
768 Fut: std::future::Future<Output = Result<DaemonLifecycleResult>> + Send + 'static,
769 {
770 self.start_or_join_lifecycle_with_key(session_id, kind, resume_workspace_id, None, work)
771 }
772
773 fn start_or_join_lifecycle_with_key<F, Fut>(
774 self: &Arc<Self>,
775 session_id: String,
776 kind: LifecycleKind,
777 resume_workspace_id: Option<String>,
778 request_key: Option<String>,
779 work: F,
780 ) -> Result<
781 tokio::sync::watch::Receiver<Option<std::result::Result<DaemonLifecycleResult, String>>>,
782 >
783 where
784 F: FnOnce(Arc<Self>, String, Arc<AtomicBool>) -> Fut + Send + 'static,
785 Fut: std::future::Future<Output = Result<DaemonLifecycleResult>> + Send + 'static,
786 {
787 self.start_or_join_lifecycle_controlled(
788 session_id,
789 kind,
790 resume_workspace_id,
791 request_key,
792 None,
793 work,
794 )
795 }
796
797 fn start_or_join_lifecycle_controlled<F, Fut>(
798 self: &Arc<Self>,
799 session_id: String,
800 kind: LifecycleKind,
801 resume_workspace_id: Option<String>,
802 request_key: Option<String>,
803 create_control: Option<CreateSessionControl>,
804 work: F,
805 ) -> Result<
806 tokio::sync::watch::Receiver<Option<std::result::Result<DaemonLifecycleResult, String>>>,
807 >
808 where
809 F: FnOnce(Arc<Self>, String, Arc<AtomicBool>) -> Fut + Send + 'static,
810 Fut: std::future::Future<Output = Result<DaemonLifecycleResult>> + Send + 'static,
811 {
812 let mut work = Some(work);
813 ensure!(
814 matches!(kind, LifecycleKind::Move | LifecycleKind::ForceDestroy)
815 || !crate::controller::move_session::move_has_pending_queue(&session_id),
816 "Move queue admission is incomplete; retry Move on the same destination before another lifecycle operation"
817 );
818 let result = {
819 let mut lifecycle = self
820 .lifecycle
821 .lock()
822 .unwrap_or_else(PoisonError::into_inner);
823 let completed_other_kind = lifecycle
824 .get(&session_id)
825 .is_some_and(|active| active.kind != kind && active.result.borrow().is_some());
826 if completed_other_kind {
827 lifecycle.remove(&session_id);
828 }
829 if let Some(active) = lifecycle.get(&session_id) {
830 ensure!(
831 active.request_key == request_key,
832 "another lifecycle request with different selections is already running for session {session_id}"
833 );
834 ensure!(
835 active.kind == kind,
836 "another lifecycle operation is already running for session {session_id}"
837 );
838 ensure!(
839 resume_workspace_id.is_none()
840 || active.resume_workspace_id == resume_workspace_id,
841 "session {session_id} is already resuming into another workspace"
842 );
843 active.result.clone()
844 } else {
845 let cancelled = create_control
846 .as_ref()
847 .map(|control| control.cancelled.clone())
848 .unwrap_or_else(|| Arc::new(AtomicBool::new(false)));
849 let (result_tx, result_rx) = tokio::sync::watch::channel(None);
850 lifecycle.insert(
851 session_id.clone(),
852 ActiveLifecycle {
853 operation_id: new_command_id("lifecycle")?,
854 create_control,
855 kind,
856 cancelled: cancelled.clone(),
857 started_at_epoch_seconds: epoch_seconds(),
858 active_stages: BTreeMap::new(),
859 resume_workspace_id,
860 resume_destination: None,
861 notice: None,
862 request_key,
863 move_source_closed: false,
864 _move_guard: (kind == LifecycleKind::Move)
865 .then(|| MoveMutationGuard::reserve(&session_id))
866 .transpose()?,
867 result: result_rx.clone(),
868 },
869 );
870 self.publish_revision();
871 let state = Arc::clone(self);
872 let operation_session_id = session_id.clone();
873 let operation = work.take().expect("new lifecycle operation has work");
874 let completed_channel = result_rx.clone();
875 tokio::spawn(async move {
876 let operation_state = state.clone();
877 let operation_id = operation_session_id.clone();
878 let mut result = match tokio::spawn(async move {
879 operation(operation_state, operation_id, cancelled).await
880 })
881 .await
882 {
883 Ok(result) => result.map_err(|error| format!("{error:#}")),
884 Err(error) => Err(format!("daemon lifecycle task failed: {error}")),
885 };
886 if let Err(error) = state.reload_controller().await {
887 let reload_error = format!(
888 "reload daemon state after lifecycle operation for {operation_session_id}: {error:#}"
889 );
890 if result.is_ok() {
891 result = Err(reload_error);
892 } else {
893 tracing::warn!(
894 session_id = %operation_session_id,
895 error = reload_error,
896 "lifecycle failed and its durable state could not be reloaded"
897 );
898 }
899 }
900 if let Err(error) =
901 reach_test_hook("lifecycle_reservation_before_result_publication").await
902 {
903 result = Err(format!("test lifecycle publication hook failed: {error:#}"));
904 }
905 let deferred_cleanup =
906 matches!(result, Ok(DaemonLifecycleResult::DeferredCleanup));
907 result_tx.send_replace(Some(result));
908 if let Some(active) = state
912 .lifecycle
913 .lock()
914 .unwrap_or_else(PoisonError::into_inner)
915 .get_mut(&operation_session_id)
916 && active.result.same_channel(&completed_channel)
917 {
918 active._move_guard.take();
919 }
920 if deferred_cleanup
924 && let Err(error) =
925 state.start_deferred_cleanup(operation_session_id.clone())
926 {
927 tracing::warn!(session_id = %operation_session_id, %error, "could not start retained cleanup");
928 state.push_notice(&operation_session_id, "Container cleanup could not start; retry cleanup from the stopped session.");
929 state
930 .lifecycle
931 .lock()
932 .unwrap_or_else(PoisonError::into_inner)
933 .retain(|_, active| !active.result.same_channel(&completed_channel));
934 }
935 state.publish_revision();
936 });
937 result_rx
938 }
939 };
940 Ok(result)
941 }
942
943 async fn wait_lifecycle_result(
944 mut result: tokio::sync::watch::Receiver<
945 Option<std::result::Result<DaemonLifecycleResult, String>>,
946 >,
947 ) -> Result<DaemonLifecycleResult> {
948 loop {
949 if let Some(result) = result.borrow_and_update().clone() {
950 return result.map_err(anyhow::Error::msg);
951 }
952 result
953 .changed()
954 .await
955 .context("daemon lifecycle operation stopped without a result")?;
956 }
957 }
958
959 async fn run_lifecycle<F, Fut>(
960 self: &Arc<Self>,
961 session_id: String,
962 kind: LifecycleKind,
963 work: F,
964 ) -> Result<DaemonLifecycleResult>
965 where
966 F: FnOnce(Arc<Self>, String, Arc<AtomicBool>) -> Fut + Send + 'static,
967 Fut: std::future::Future<Output = Result<DaemonLifecycleResult>> + Send + 'static,
968 {
969 let result = self.start_or_join_lifecycle(session_id, kind, work)?;
970 let channel = result.clone();
971 let outcome = Self::wait_lifecycle_result(result).await;
972 self.remove_completed_lifecycle(&channel);
973 outcome
974 }
975
976 fn remove_completed_lifecycle(
977 &self,
978 channel: &tokio::sync::watch::Receiver<
979 Option<std::result::Result<DaemonLifecycleResult, String>>,
980 >,
981 ) {
982 self.lifecycle
983 .lock()
984 .unwrap_or_else(PoisonError::into_inner)
985 .retain(|_, active| !active.result.same_channel(channel) || active.is_visible());
986 }
987
988 async fn start_create_session(
989 self: &Arc<Self>,
990 request: CreateSessionRequest,
991 ) -> Result<RegisteredSession> {
992 self.start_create_session_inner(request, CreateSessionControl::default(), None)
993 .await
994 }
995
996 pub async fn start_subagent_session(
998 self: &Arc<Self>,
999 request: crate::controller::RegisterSubagentRequest,
1000 ) -> Result<mj_core::subagent::SubagentRecord> {
1001 let relation = blocking(move || {
1002 let mut controller = Controller::load()?;
1003 controller.register_subagent(request)
1004 })
1005 .await?;
1006 let session_id = relation.child_session_id.clone();
1007 self.start_or_join_lifecycle_controlled(
1008 session_id.clone(),
1009 LifecycleKind::Create,
1010 None,
1011 Some(relation.request_key.clone()),
1012 None,
1013 move |state, session_id, cancelled| async move {
1014 let mut controller = tokio::task::spawn_blocking(Controller::load)
1015 .await
1016 .context("load controller for sub-agent startup")??;
1017 let executor = DaemonStageReportingExecutor::new(
1018 CancellableProcessExecutor::new(cancelled),
1019 state,
1020 session_id.clone(),
1021 );
1022 controller
1023 .provision_subagent_session_controlled(&session_id, &executor)
1024 .await?;
1025 Ok(DaemonLifecycleResult::Done)
1026 },
1027 )?;
1028 self.reload_controller().await?;
1029 Ok(relation)
1030 }
1031
1032 pub async fn start_create_session_controlled(
1033 self: &Arc<Self>,
1034 request: CreateSessionRequest,
1035 control: CreateSessionControl,
1036 publication: tokio::sync::oneshot::Receiver<std::result::Result<(), String>>,
1037 ) -> Result<RegisteredSession> {
1038 self.start_create_session_inner(request, control, Some(publication))
1039 .await
1040 }
1041
1042 async fn start_create_session_inner(
1043 self: &Arc<Self>,
1044 request: CreateSessionRequest,
1045 control: CreateSessionControl,
1046 publication: Option<tokio::sync::oneshot::Receiver<std::result::Result<(), String>>>,
1047 ) -> Result<RegisteredSession> {
1048 let path_cancelled = control.cancelled.clone();
1049 let registered = blocking(move || {
1050 let mut controller = Controller::load()?;
1051 let path_executor = crate::targets::CancellableProcessExecutor::new(path_cancelled)
1052 .with_deadline(Duration::from_secs(30));
1053 let project_directory = request
1054 .project_directory
1055 .as_deref()
1056 .map(|path| {
1057 controller.resolve_project_directory(
1058 &request.target_template_id,
1059 path,
1060 &path_executor,
1061 )
1062 })
1063 .transpose()?;
1064 let session_id = controller.register_session_with_resources(
1065 &request.profile_id,
1066 &request.bundle_id,
1067 &request.target_template_id,
1068 request.title,
1069 SessionLaunchOptions {
1070 create_managed_worktree: request.create_managed_worktree,
1071 mjolnir_subagents: request.mjolnir_subagents,
1072 initial_prompt: request.initial_prompt,
1073 workspace_id: request.workspace_id,
1074 additional_mounts: request.additional_mounts,
1075 allow_dirty_local: request.allow_dirty_local,
1076 resource_allocation: request.resource_allocation,
1077 project_directory,
1078 session_title_override: request.session_title_override,
1079 },
1080 )?;
1081 let session = controller
1082 .state
1083 .sessions
1084 .get(&session_id)
1085 .expect("newly registered session exists")
1086 .clone();
1087 let remembered_container_size = controller
1088 .config
1089 .targets
1090 .get(&request.target_template_id)
1091 .and_then(mj_core::config::container_size_host)
1092 .and_then(|host| {
1093 controller
1094 .state
1095 .container_sizes
1096 .get(host)
1097 .copied()
1098 .map(|size| (host.to_owned(), size))
1099 });
1100 Ok(RegisteredSession {
1101 session,
1102 remembered_container_size,
1103 })
1104 })
1105 .await?;
1106 let session_id = registered.session.id.clone();
1107 self.start_or_join_lifecycle_controlled(
1108 session_id,
1109 LifecycleKind::Create,
1110 None,
1111 None,
1112 Some(control.clone()),
1113 move |state, session_id, cancelled| async move {
1114 let mut controller = tokio::task::spawn_blocking(Controller::load)
1115 .await
1116 .context("load controller for daemon create task")??;
1117 let publication_error = if let Some(publication) = publication {
1118 let published = tokio::select! {
1119 result = publication => result.context("session publication owner stopped")
1120 .and_then(|result| result.map_err(anyhow::Error::msg)),
1121 () = async {
1122 while !cancelled.load(Ordering::Acquire) {
1123 tokio::time::sleep(Duration::from_millis(25)).await;
1124 }
1125 } => Err(anyhow!("session creation cancelled before publication")),
1126 };
1127 published.err()
1128 } else {
1129 None
1130 };
1131 if publication_error.is_some() {
1132 control.request_cancel();
1133 }
1134 let executor = DaemonStageReportingExecutor::new(
1135 CancellableProcessExecutor::new(cancelled),
1136 state,
1137 session_id.clone(),
1138 );
1139 let provision = controller
1140 .provision_session_controlled_with_commit(&session_id, &executor, || {
1141 ensure!(
1142 control.grant_commit(),
1143 "session creation cancelled before commit"
1144 );
1145 Ok(())
1146 })
1147 .await;
1148 if let Some(error) = publication_error {
1149 return match provision {
1150 Ok(()) => Err(error),
1151 Err(rollback) => {
1152 Err(error.context(format!("discard unpublished session: {rollback:#}")))
1153 }
1154 };
1155 }
1156 provision?;
1157 Ok(DaemonLifecycleResult::Done)
1158 },
1159 )?;
1160 self.reload_controller().await?;
1161 Ok(registered)
1162 }
1163
1164 pub async fn wait_create_session(&self, session_id: &str) -> Result<()> {
1165 let result = {
1166 let lifecycle = self
1167 .lifecycle
1168 .lock()
1169 .unwrap_or_else(PoisonError::into_inner);
1170 let active = lifecycle
1171 .get(session_id)
1172 .with_context(|| format!("no create operation exists for session {session_id}"))?;
1173 ensure!(
1174 active.kind == LifecycleKind::Create,
1175 "session {session_id} is no longer being created"
1176 );
1177 active.result.clone()
1178 };
1179 let channel = result.clone();
1180 let outcome = Self::wait_lifecycle_result(result).await;
1181 self.remove_completed_lifecycle(&channel);
1182 match outcome? {
1183 DaemonLifecycleResult::Done => Ok(()),
1184 DaemonLifecycleResult::Move(_) => unreachable!("cleanup cannot return a move outcome"),
1185 DaemonLifecycleResult::DeferredCleanup => {
1186 unreachable!("session creation cannot schedule target cleanup")
1187 }
1188 }
1189 }
1190
1191 pub fn request_close(&self, session_id: &str) {
1192 self.close_requested
1193 .lock()
1194 .unwrap_or_else(PoisonError::into_inner)
1195 .insert(session_id.to_owned());
1196 self.publish_revision();
1197 }
1198
1199 pub fn clear_close_request(&self, session_id: &str) {
1200 self.close_requested
1201 .lock()
1202 .unwrap_or_else(PoisonError::into_inner)
1203 .remove(session_id);
1204 self.publish_revision();
1205 }
1206
1207 pub fn close_is_requested(&self, session_id: &str) -> bool {
1208 self.close_requested
1209 .lock()
1210 .unwrap_or_else(PoisonError::into_inner)
1211 .contains(session_id)
1212 }
1213
1214 pub async fn close_session(self: &Arc<Self>, session_id: String) -> Result<()> {
1215 let children = blocking({
1216 let session_id = session_id.clone();
1217 move || {
1218 let controller = Controller::load()?;
1219 Ok(controller
1220 .state
1221 .subagents
1222 .values()
1223 .filter(|child| child.parent_session_id == session_id)
1224 .filter(|child| {
1225 controller
1226 .state
1227 .sessions
1228 .get(&child.child_session_id)
1229 .is_some_and(|session| session.state.is_active())
1230 })
1231 .map(|child| child.child_session_id.clone())
1232 .collect::<Vec<_>>())
1233 }
1234 })
1235 .await?;
1236 for child_id in children {
1237 self.request_close(&child_id);
1238 let result = self.close_requested_session(child_id.clone()).await;
1239 self.clear_close_request(&child_id);
1240 result.with_context(|| format!("stop sub-agent {child_id} before its parent"))?;
1241 }
1242 self.request_close(&session_id);
1243 let result = self.close_requested_session(session_id.clone()).await;
1244 self.clear_close_request(&session_id);
1245 result
1246 }
1247
1248 async fn wait_before_close(self: &Arc<Self>, session_id: &str) -> Result<()> {
1249 let pending = {
1252 let operations = self
1253 .lifecycle
1254 .lock()
1255 .unwrap_or_else(PoisonError::into_inner);
1256 operations
1257 .get(session_id)
1258 .filter(|operation| {
1259 !matches!(
1260 operation.kind,
1261 LifecycleKind::Close | LifecycleKind::Cleanup
1262 )
1263 })
1264 .map(|operation| {
1265 operation.request_cancel();
1266 operation.result.clone()
1267 })
1268 };
1269 if let Some(pending) = pending {
1270 if let Err(error) = Self::wait_lifecycle_result(pending.clone()).await {
1271 tracing::debug!(%session_id, %error, "previous lifecycle ended before close");
1272 }
1273 self.remove_completed_lifecycle(&pending);
1274 }
1275
1276 Ok(())
1277 }
1278
1279 async fn close_requested_session(self: &Arc<Self>, session_id: String) -> Result<()> {
1280 self.wait_before_close(&session_id).await?;
1281 let (already_stopped, needs_cleanup) = blocking({
1282 let session_id = session_id.clone();
1283 move || {
1284 let controller = Controller::load()?;
1285 Ok(controller
1286 .state
1287 .sessions
1288 .get(&session_id)
1289 .map_or((false, false), |session| {
1290 (
1291 session.state == SessionState::Stopped,
1292 session.state == SessionState::Stopped && session.target.is_some(),
1293 )
1294 }))
1295 }
1296 })
1297 .await?;
1298 if already_stopped {
1299 if needs_cleanup {
1300 self.start_deferred_cleanup(session_id)?;
1301 }
1302 return Ok(());
1303 }
1304 let operation_session_id = session_id.clone();
1305 let result = self
1306 .run_lifecycle(
1307 operation_session_id,
1308 LifecycleKind::Close,
1309 |state, session_id, cancelled| async move {
1310 let _recovery_reservation = tokio::task::spawn_blocking({
1311 let observer = state.recovery_observer.clone();
1312 let session_id = session_id.clone();
1313 let cancelled = cancelled.clone();
1314 move || reserve_recovery_or_cancel(&observer, &session_id, &cancelled)
1315 })
1316 .await
1317 .context("reserve recovery for daemon close task")??;
1318 let mut controller = tokio::task::spawn_blocking(Controller::load)
1319 .await
1320 .context("load controller for daemon close task")??;
1321 let executor = DaemonStageReportingExecutor::new(
1322 CancellableProcessExecutor::new(cancelled),
1323 state.clone(),
1324 session_id.clone(),
1325 );
1326 let deferred = controller
1327 .close_session_managed_controlled(
1328 &session_id,
1329 &executor,
1330 &state.session_manager,
1331 )
1332 .await?;
1333 Ok(if deferred {
1334 DaemonLifecycleResult::DeferredCleanup
1335 } else {
1336 DaemonLifecycleResult::Done
1337 })
1338 },
1339 )
1340 .await?;
1341 let _ = result; Ok(())
1343 }
1344
1345 fn start_deferred_cleanup(
1346 self: &Arc<Self>,
1347 session_id: String,
1348 ) -> Result<
1349 tokio::sync::watch::Receiver<Option<std::result::Result<DaemonLifecycleResult, String>>>,
1350 > {
1351 let result = self.start_or_join_lifecycle(
1352 session_id.clone(),
1353 LifecycleKind::Cleanup,
1354 |state, session_id, cancelled| async move {
1355 blocking(move || {
1356 let mut controller = Controller::load()?;
1357 let executor = DaemonStageReportingExecutor::new(
1358 CancellableProcessExecutor::new(cancelled),
1359 state,
1360 session_id.clone(),
1361 );
1362 controller.cleanup_stopped_target(&session_id, &executor)?;
1363 Ok(DaemonLifecycleResult::Done)
1364 })
1365 .await
1366 },
1367 )?;
1368 let caller_result = result.clone();
1369 let channel = result.clone();
1370 let state = Arc::clone(self);
1371 tokio::spawn(async move {
1372 if let Err(error) = Self::wait_lifecycle_result(result).await {
1373 tracing::warn!(%session_id, error = format!("{error:#}"), "deferred Podman cleanup failed");
1374 state.push_notice(
1375 &session_id,
1376 "Container storage cleanup failed; the stopped session retains its target for retry.",
1377 );
1378 }
1379 state.remove_completed_lifecycle(&channel);
1380 });
1381 Ok(caller_result)
1382 }
1383
1384 fn resume_retained_cleanups(self: &Arc<Self>) {
1385 let session_ids = self
1386 .controller
1387 .lock()
1388 .unwrap_or_else(PoisonError::into_inner)
1389 .state
1390 .sessions
1391 .iter()
1392 .filter(|(_, session)| {
1393 session.state == SessionState::Stopped && session.target.is_some()
1394 })
1395 .map(|(session_id, _)| session_id.clone())
1396 .collect::<Vec<_>>();
1397 for session_id in session_ids {
1398 if crate::controller::move_session::move_owns_session(&session_id) {
1399 continue;
1400 }
1401 if let Err(error) = self.start_deferred_cleanup(session_id.clone()) {
1402 tracing::warn!(%session_id, error = format!("{error:#}"), "could not resume deferred Podman cleanup");
1403 self.push_notice(
1404 &session_id,
1405 format!("Could not resume container storage cleanup: {error:#}"),
1406 );
1407 }
1408 }
1409 }
1410
1411 async fn wait_for_deferred_cleanup(self: &Arc<Self>, session_id: &str) -> Result<()> {
1412 let existing = {
1413 let lifecycle = self
1414 .lifecycle
1415 .lock()
1416 .unwrap_or_else(PoisonError::into_inner);
1417 lifecycle.get(session_id).and_then(|active| {
1418 (active.kind == LifecycleKind::Cleanup).then(|| active.result.clone())
1419 })
1420 };
1421 let result = match existing {
1422 Some(result) => result,
1423 None => {
1424 let needs_cleanup = blocking({
1425 let session_id = session_id.to_owned();
1426 move || {
1427 let controller = Controller::load()?;
1428 Ok(controller
1429 .state
1430 .sessions
1431 .get(&session_id)
1432 .is_some_and(|session| {
1433 session.state == SessionState::Stopped && session.target.is_some()
1434 }))
1435 }
1436 })
1437 .await?;
1438 if !needs_cleanup {
1439 return Ok(());
1440 }
1441 self.start_deferred_cleanup(session_id.to_owned())?
1442 }
1443 };
1444 let channel = result.clone();
1445 let outcome = Self::wait_lifecycle_result(result).await;
1446 self.remove_completed_lifecycle(&channel);
1447 match outcome? {
1448 DaemonLifecycleResult::Done => Ok(()),
1449 DaemonLifecycleResult::Move(_) => unreachable!("cleanup cannot return a move outcome"),
1450 DaemonLifecycleResult::DeferredCleanup => {
1451 unreachable!("cleanup cannot schedule another cleanup")
1452 }
1453 }
1454 }
1455
1456 pub async fn resume_session(self: &Arc<Self>, request: ResumeSessionRequest) -> Result<()> {
1464 let session_id = request.session_id.clone();
1465 self.wait_for_deferred_cleanup(&session_id).await?;
1466 let already_running = blocking({
1469 let session_id = session_id.clone();
1470 move || {
1471 let controller = Controller::load()?;
1472 Ok(controller
1473 .state
1474 .sessions
1475 .get(&session_id)
1476 .is_some_and(|session| session.state == SessionState::Running))
1477 }
1478 })
1479 .await?;
1480 if already_running {
1481 return Ok(());
1482 }
1483 let profile_id = request.profile_id.clone();
1484 let target_template_id = request.target_template_id.clone();
1485 let workspace_id = request.workspace_id.clone();
1486 let rebind_workspace_id = workspace_id.clone();
1487 let operation_session_id = session_id.clone();
1488 let result = self.start_or_join_lifecycle_for_workspace(
1489 session_id,
1490 LifecycleKind::Resume,
1491 Some(workspace_id),
1492 move |state, session_id, cancelled| async move {
1493 let _recovery_reservation = tokio::task::spawn_blocking({
1494 let observer = state.recovery_observer.clone();
1495 let session_id = session_id.clone();
1496 let cancelled = cancelled.clone();
1497 move || reserve_recovery_or_cancel(&observer, &session_id, &cancelled)
1498 })
1499 .await
1500 .context("reserve recovery for daemon resume task")??;
1501 blocking({
1502 let session_id = session_id.clone();
1503 move || {
1504 crate::database::reassign_resumable_session_workspace(
1505 &session_id,
1506 &rebind_workspace_id,
1507 )
1508 }
1509 })
1510 .await?;
1511 let restore_request = request.clone();
1512 let mut controller = tokio::task::spawn_blocking(move || {
1513 session_move::load_controller_for_resume(&restore_request)
1514 })
1515 .await
1516 .context("load controller for daemon resume task")??;
1517 let executor = DaemonStageReportingExecutor::new(
1518 CancellableProcessExecutor::new(cancelled),
1519 state.clone(),
1520 session_id.clone(),
1521 );
1522 let materialized = controller
1523 .resume_session_controlled_with_repository_preflight(
1524 &session_id,
1525 &request.profile_id,
1526 &request.target_template_id,
1527 SessionResumeOptions {
1528 additional_mounts: request.additional_mounts,
1529 resource_allocation: request.resource_allocation,
1530 discard_queue: request.discard_queue,
1531 },
1532 request.repository_preflight,
1533 &executor,
1534 )
1535 .await?;
1536 let _ = materialized;
1540 Ok(DaemonLifecycleResult::Done)
1541 },
1542 )?;
1543 self.set_lifecycle_resume_destination(
1544 &operation_session_id,
1545 profile_id,
1546 target_template_id,
1547 );
1548 let channel = result.clone();
1549 let result = Self::wait_lifecycle_result(result).await;
1550 self.remove_completed_lifecycle(&channel);
1551 match result? {
1552 DaemonLifecycleResult::Done => {}
1553 DaemonLifecycleResult::Move(_) => unreachable!("resume cannot return a move outcome"),
1554 DaemonLifecycleResult::DeferredCleanup => {
1555 unreachable!("session resume cannot schedule target cleanup")
1556 }
1557 }
1558 blocking(move || {
1559 if let Some(mut operation) =
1560 crate::database::load_move_operation(&operation_session_id)?
1561 && !operation.queue_admission_started
1562 {
1563 operation.phase = mj_core::state::MovePhase::Cancelled;
1564 operation.queue_admission_finished = true;
1565 operation.updated_at = chrono::Utc::now().to_rfc3339();
1566 operation.error = Some("Recovered through an explicit Resume operation".into());
1567 crate::database::save_move_operation(&operation)?;
1568 }
1569 Ok(())
1570 })
1571 .await?;
1572 Ok(())
1573 }
1574
1575 async fn force_stop_session(self: &Arc<Self>, session_id: String) -> Result<()> {
1576 let children = blocking({
1577 let session_id = session_id.clone();
1578 move || {
1579 Ok(crate::database::list_subagents(&session_id)?
1580 .into_iter()
1581 .map(|child| child.child_session_id)
1582 .collect::<Vec<_>>())
1583 }
1584 })
1585 .await?;
1586 for child_id in children {
1587 Box::pin(self.force_stop_session(child_id.clone()))
1588 .await
1589 .with_context(|| format!("force-stop sub-agent {child_id} before its parent"))?;
1590 }
1591 let operation_session_id = session_id.clone();
1592 let result = self
1593 .run_lifecycle(
1594 operation_session_id,
1595 LifecycleKind::ForceStop,
1596 |state, session_id, cancelled| async move {
1597 blocking(move || {
1598 let mut controller = Controller::load()?;
1599 let executor = DaemonStageReportingExecutor::new(
1600 CancellableProcessExecutor::new(cancelled),
1601 state,
1602 session_id.clone(),
1603 );
1604 let deferred = controller.force_stop(&session_id, &executor)?;
1605 Ok(if deferred {
1606 DaemonLifecycleResult::DeferredCleanup
1607 } else {
1608 DaemonLifecycleResult::Done
1609 })
1610 })
1611 .await
1612 },
1613 )
1614 .await?;
1615 let _ = result; Ok(())
1617 }
1618
1619 async fn destroy_stopped_session(self: &Arc<Self>, session_id: String) -> Result<()> {
1620 let children = blocking({
1621 let session_id = session_id.clone();
1622 move || {
1623 Ok(crate::database::list_subagents(&session_id)?
1624 .into_iter()
1625 .map(|child| child.child_session_id)
1626 .collect::<Vec<_>>())
1627 }
1628 })
1629 .await?;
1630 for child_id in children {
1631 Box::pin(self.force_destroy_session(child_id.clone()))
1632 .await
1633 .with_context(|| format!("destroy sub-agent {child_id} before its parent"))?;
1634 }
1635 self.wait_for_deferred_cleanup(&session_id).await?;
1636 let exists = blocking({
1637 let session_id = session_id.clone();
1638 move || Ok(Controller::load()?.state.sessions.contains_key(&session_id))
1639 })
1640 .await?;
1641 if !exists {
1642 return Ok(());
1643 }
1644 self.run_lifecycle(
1645 session_id,
1646 LifecycleKind::DestroyStopped,
1647 |state, session_id, cancelled| async move {
1648 blocking(move || {
1649 let mut controller = Controller::load()?;
1650 let executor = DaemonStageReportingExecutor::new(
1651 CancellableProcessExecutor::new(cancelled),
1652 state,
1653 session_id.clone(),
1654 );
1655 controller.destroy_session_controlled(&session_id, &executor)?;
1656 Ok(DaemonLifecycleResult::Done)
1657 })
1658 .await
1659 },
1660 )
1661 .await?;
1662 Ok(())
1663 }
1664
1665 async fn preempt_active_lifecycle(self: &Arc<Self>, session_id: &str) -> Result<()> {
1676 let mut result = {
1677 let lifecycle = self
1678 .lifecycle
1679 .lock()
1680 .unwrap_or_else(PoisonError::into_inner);
1681 let Some(active) = lifecycle.get(session_id) else {
1682 return Ok(());
1683 };
1684 if !active.result.borrow().is_none() {
1685 return Ok(());
1686 }
1687 active.request_cancel();
1688 active.result.clone()
1689 };
1690 let finished = tokio::time::timeout(FORCE_DESTROY_PREEMPT_TIMEOUT, async {
1691 loop {
1692 if result.borrow().is_some() {
1693 return Ok(());
1694 }
1695 if result.changed().await.is_err() {
1696 return Err(());
1697 }
1698 }
1699 })
1700 .await;
1701 match finished {
1702 Ok(Ok(())) => Ok(()),
1705 Ok(Err(())) => bail!(
1706 "daemon lifecycle operation stopped without a result for session {session_id}"
1707 ),
1708 Err(_) => bail!(
1709 "session {session_id} still has an operation that did not stop after cancellation; try again"
1710 ),
1711 }
1712 }
1713
1714 pub async fn force_destroy_session(self: &Arc<Self>, session_id: String) -> Result<()> {
1718 let children = blocking({
1719 let session_id = session_id.clone();
1720 move || {
1721 Ok(crate::database::list_subagents(&session_id)?
1722 .into_iter()
1723 .map(|child| child.child_session_id)
1724 .collect::<Vec<_>>())
1725 }
1726 })
1727 .await?;
1728 for child_id in children {
1729 Box::pin(self.force_destroy_session(child_id.clone()))
1730 .await
1731 .with_context(|| format!("destroy sub-agent {child_id} before its parent"))?;
1732 }
1733 self.preempt_active_lifecycle(&session_id).await?;
1734 let exists = blocking({
1735 let session_id = session_id.clone();
1736 move || Ok(Controller::load()?.state.sessions.contains_key(&session_id))
1737 })
1738 .await?;
1739 if !exists {
1740 return Ok(());
1741 }
1742 self.run_lifecycle(
1743 session_id,
1744 LifecycleKind::ForceDestroy,
1745 |state, session_id, cancelled| async move {
1746 let _recovery_reservation = tokio::task::spawn_blocking({
1747 let observer = state.recovery_observer.clone();
1748 let session_id = session_id.clone();
1749 let cancelled = cancelled.clone();
1750 move || reserve_recovery_or_cancel(&observer, &session_id, &cancelled)
1751 })
1752 .await
1753 .context("reserve recovery for daemon force-destroy task")??;
1754 blocking({
1755 let session_id = session_id.clone();
1756 move || {
1757 let mut controller = Controller::load()?;
1758 let executor = DaemonStageReportingExecutor::new(
1759 CancellableProcessExecutor::new(cancelled),
1760 state,
1761 session_id.clone(),
1762 );
1763 controller.force_destroy_session(&session_id, &executor)?;
1764 crate::controller::move_session::release_move_queue_hold(&session_id);
1765 Ok(DaemonLifecycleResult::Done)
1766 }
1767 })
1768 .await
1769 },
1770 )
1771 .await?;
1772 Ok(())
1773 }
1774
1775 pub async fn force_delete_workspace(self: &Arc<Self>, workspace_id: String) -> Result<()> {
1784 ensure!(
1785 !self.workspace_has_active_resume(&workspace_id),
1786 "workspace has a session resume in progress"
1787 );
1788 let sessions = blocking({
1789 let workspace_id = workspace_id.clone();
1790 move || {
1791 let controller = Controller::load()?;
1792 Ok(active_sessions_for_force_destruction(
1793 &controller,
1794 &workspace_id,
1795 ))
1796 }
1797 })
1798 .await?;
1799 for (index, session_id) in sessions.iter().enumerate() {
1800 if let Err(error) = self.force_destroy_session(session_id.clone()).await {
1801 let remaining = sessions.len() - index - 1;
1802 bail!(
1803 "force-destroying session {session_id} failed: {error:#}; \
1804 {remaining} session(s) in the workspace remain"
1805 );
1806 }
1807 }
1808 blocking({
1809 let workspace_id = workspace_id.clone();
1810 move || crate::database::force_delete_workspace(&workspace_id)
1811 })
1812 .await?;
1813 refresh_runtime_workspaces(self).await?;
1814 Ok(())
1815 }
1816
1817 fn cancel_lifecycle(&self, session_id: &str) -> Result<()> {
1818 let lifecycle = self
1819 .lifecycle
1820 .lock()
1821 .unwrap_or_else(PoisonError::into_inner);
1822 let active = lifecycle.get(session_id).with_context(|| {
1823 format!("no lifecycle operation is running for session {session_id}")
1824 })?;
1825 ensure!(
1826 active.request_cancel(),
1827 "lifecycle operation is no longer cancellable"
1828 );
1829 self.publish_revision();
1830 Ok(())
1831 }
1832
1833 async fn cancel_and_wait_lifecycles(&self) -> Result<()> {
1838 let mut pending = {
1839 let lifecycle = self
1840 .lifecycle
1841 .lock()
1842 .unwrap_or_else(PoisonError::into_inner);
1843 lifecycle
1844 .iter()
1845 .filter(|(_, active)| active.result.borrow().is_none())
1846 .map(|(session_id, active)| {
1847 if active.kind != LifecycleKind::Cleanup {
1848 active.request_cancel();
1849 }
1850 let stage = active
1851 .active_stages
1852 .keys()
1853 .next_back()
1854 .map(|stage| stage.label())
1855 .unwrap_or_else(|| "container cleanup".to_owned());
1856 (
1857 session_id.clone(),
1858 active.kind,
1859 stage,
1860 active.cancelled.clone(),
1861 active.result.clone(),
1862 )
1863 })
1864 .collect::<Vec<_>>()
1865 };
1866 let cleanup_deadline = tokio::time::Instant::now() + Duration::from_secs(8);
1867 for (session_id, kind, stage, cancelled, result) in &mut pending {
1868 if *kind != LifecycleKind::Cleanup || result.borrow().is_some() {
1869 continue;
1870 }
1871 tracing::info!(%session_id, %stage, "daemon shutdown is waiting for deferred cleanup");
1872 self.set_lifecycle_notice(
1873 session_id,
1874 &format!("Daemon shutdown is waiting for {stage}"),
1875 );
1876 let finished = tokio::time::timeout_at(cleanup_deadline, async {
1877 while result.borrow_and_update().is_none() {
1878 result.changed().await.with_context(|| {
1879 format!("cleanup owner stopped without a result for session {session_id}")
1880 })?;
1881 }
1882 Ok::<_, anyhow::Error>(())
1883 })
1884 .await;
1885 match finished {
1886 Ok(result) => result?,
1887 Err(_) => {
1888 tracing::warn!(%session_id, %stage, "deferred cleanup exceeded the daemon shutdown drain deadline");
1889 cancelled.store(true, Ordering::Release);
1890 }
1891 }
1892 }
1893 let join_deadline = tokio::time::Instant::now() + Duration::from_secs(1);
1894 for (session_id, _, stage, cancelled, mut result) in pending {
1895 cancelled.store(true, Ordering::Release);
1896 let joined = tokio::time::timeout_at(join_deadline, async {
1897 while result.borrow_and_update().is_none() {
1898 result.changed().await.with_context(|| {
1899 format!("lifecycle owner stopped without a result for session {session_id}")
1900 })?;
1901 }
1902 Ok::<_, anyhow::Error>(())
1903 })
1904 .await;
1905 if joined.is_err() {
1906 bail!(
1907 "timed out cancelling lifecycle owner for session {session_id} while {stage}"
1908 );
1909 }
1910 joined.expect("checked timeout")?;
1911 }
1912 Ok(())
1913 }
1914
1915 pub fn active_lifecycles(&self) -> Vec<RuntimeLifecycleView> {
1924 self.lifecycle
1925 .lock()
1926 .unwrap_or_else(PoisonError::into_inner)
1927 .iter()
1928 .filter(|(_, active)| active.is_visible())
1929 .map(|(session_id, active)| RuntimeLifecycleView {
1930 operation_id: active.operation_id.clone(),
1931 cancellable: active.is_cancellable(),
1932 session_id: session_id.clone(),
1933 kind: active.kind.into(),
1934 started_at_epoch_seconds: active.started_at_epoch_seconds,
1935 active_stages: active
1936 .active_stages
1937 .iter()
1938 .map(|(stage, (_, started_at))| (*stage, *started_at))
1939 .collect(),
1940 resume_destination: active.resume_destination.clone(),
1941 notice: active.notice.clone(),
1942 })
1943 .collect()
1944 }
1945
1946 pub fn session_state(&self, session_id: &str) -> Option<mj_core::state::SessionState> {
1950 if self.close_is_requested(session_id) {
1951 return Some(SessionState::Closing);
1952 }
1953 self.controller
1954 .lock()
1955 .unwrap_or_else(PoisonError::into_inner)
1956 .state
1957 .sessions
1958 .get(session_id)
1959 .map(|record| record.state)
1960 }
1961
1962 pub fn session_record(&self, session_id: &str) -> Option<SessionRecord> {
1964 self.controller
1965 .lock()
1966 .unwrap_or_else(PoisonError::into_inner)
1967 .state
1968 .sessions
1969 .get(session_id)
1970 .cloned()
1971 }
1972
1973 pub async fn workspace_session_handle(
1974 &self,
1975 session_id: &str,
1976 ) -> Result<crate::session_manager::ManagedSessionHandle> {
1977 let record = self.session_record(session_id).context("unknown session")?;
1978 ensure!(
1979 record.target.is_some()
1980 && record.state == SessionState::Running
1981 && !self.close_is_requested(session_id),
1982 "session must have a live running target for file injection"
1983 );
1984 self.session_manager.session(session_id.to_owned()).await
1985 }
1986
1987 pub async fn checkpoint_session_now(
1994 &self,
1995 session_id: &str,
1996 ) -> Result<mj_core::state::CheckpointMetadata> {
1997 ensure_no_active_lifecycle(self)?;
1998 let mut controller = blocking(Controller::load).await?;
1999 let checkpoint = controller.checkpoint_session(session_id).await?;
2000 refresh_runtime_controller(self).await;
2001 Ok(checkpoint)
2002 }
2003
2004 pub fn session_projection(
2008 &self,
2009 ) -> (BTreeMap<String, SessionRecord>, Vec<RuntimeLifecycleView>) {
2010 let controller = self
2011 .controller
2012 .lock()
2013 .unwrap_or_else(PoisonError::into_inner);
2014 let operations = self.active_lifecycles();
2015 let mut records = controller.state.sessions.clone();
2016 for id in self
2017 .close_requested
2018 .lock()
2019 .unwrap_or_else(PoisonError::into_inner)
2020 .iter()
2021 {
2022 if let Some(record) = records.get_mut(id)
2023 && record.state != SessionState::Stopped
2024 {
2025 record.state = SessionState::Closing;
2026 }
2027 }
2028 (records, operations)
2029 }
2030
2031 pub fn cancel_lifecycle_if_active(&self, session_id: &str) {
2032 if let Some(active) = self
2033 .lifecycle
2034 .lock()
2035 .unwrap_or_else(PoisonError::into_inner)
2036 .get(session_id)
2037 {
2038 active.request_cancel();
2039 self.publish_revision();
2040 }
2041 }
2042
2043 fn set_lifecycle_resume_destination(
2044 &self,
2045 session_id: &str,
2046 profile_id: String,
2047 target_id: String,
2048 ) {
2049 if let Some(active) = self
2050 .lifecycle
2051 .lock()
2052 .unwrap_or_else(PoisonError::into_inner)
2053 .get_mut(session_id)
2054 {
2055 active.resume_destination = Some((profile_id, target_id));
2056 self.publish_revision();
2057 }
2058 }
2059
2060 fn change_lifecycle_stage(&self, session_id: &str, stage: ProvisionStage, active: bool) {
2061 let changed = {
2062 let mut lifecycle = self
2063 .lifecycle
2064 .lock()
2065 .unwrap_or_else(PoisonError::into_inner);
2066 let Some(operation) = lifecycle.get_mut(session_id) else {
2067 return;
2068 };
2069 if active {
2070 let entry = operation
2071 .active_stages
2072 .entry(stage)
2073 .or_insert_with(|| (0, epoch_seconds()));
2074 entry.0 += 1;
2075 entry.0 == 1
2076 } else {
2077 let Some((count, _)) = operation.active_stages.get_mut(&stage) else {
2078 return;
2079 };
2080 *count -= 1;
2081 if *count == 0 {
2082 operation.active_stages.remove(&stage);
2083 true
2084 } else {
2085 false
2086 }
2087 }
2088 };
2089 if changed {
2090 self.publish_revision();
2091 }
2092 }
2093
2094 fn push_notice(&self, session_id: &str, text: impl Into<String>) {
2097 const RETAINED_NOTICES: usize = 32;
2098
2099 let notice = RuntimeNotice {
2100 id: self.next_notice_id.fetch_add(1, Ordering::AcqRel),
2101 session_id: session_id.to_owned(),
2102 text: text.into(),
2103 };
2104 {
2105 let mut notices = self.notices.lock().unwrap_or_else(PoisonError::into_inner);
2106 notices.push_back(notice);
2107 while notices.len() > RETAINED_NOTICES {
2108 notices.pop_front();
2109 }
2110 }
2111 self.publish_revision();
2112 }
2113
2114 fn set_lifecycle_notice(&self, session_id: &str, notice: &str) {
2115 if let Some(active) = self
2116 .lifecycle
2117 .lock()
2118 .unwrap_or_else(PoisonError::into_inner)
2119 .get_mut(session_id)
2120 {
2121 if active.kind == LifecycleKind::Move && notice == "Preparing destination" {
2122 active.move_source_closed = true;
2123 }
2124 active.notice = Some(notice.to_owned());
2125 self.publish_revision();
2126 }
2127 }
2128}
2129
2130fn report_worker_upgrade(
2133 state: &RuntimeState,
2134 result: &crate::worker_upgrade::WorkerUpgradeResult,
2135) {
2136 use crate::controller::WorkerUpgradeOutcome;
2137
2138 let session_id = &result.session_id;
2139 if result.cancelled {
2140 tracing::debug!(%session_id, "worker upgrade was preempted");
2141 return;
2142 }
2143 match &result.outcome {
2144 Ok(WorkerUpgradeOutcome::Upgraded { build }) => {
2145 tracing::info!(%session_id, %build, "replaced the session worker with the current build");
2146 let name = state
2147 .controller
2148 .lock()
2149 .unwrap_or_else(PoisonError::into_inner)
2150 .state
2151 .sessions
2152 .get(session_id)
2153 .map_or_else(
2154 || session_id.clone(),
2155 |session| session.display_title().to_owned(),
2156 );
2157 state.push_notice(session_id, format!("Upgraded the worker for {name}."));
2158 }
2159 Ok(WorkerUpgradeOutcome::AlreadyCurrent { build }) => {
2160 tracing::debug!(%session_id, %build, "session worker already runs the current build");
2161 }
2162 Ok(WorkerUpgradeOutcome::Deferred) => {
2163 tracing::debug!(%session_id, "worker upgrade deferred: the session is working");
2164 }
2165 Err(error) => {
2166 tracing::warn!(%session_id, %error, "could not upgrade the session worker");
2167 }
2168 }
2169}
2170
2171fn runtime_records_for_workspace(
2172 controller: &Controller,
2173 session_ids: &BTreeSet<String>,
2174) -> Vec<SessionRecord> {
2175 controller
2176 .state
2177 .sessions
2178 .iter()
2179 .filter(|(session_id, session)| {
2180 !session.state.is_active() || session_ids.contains(*session_id)
2181 })
2182 .map(|(_, session)| session.clone())
2183 .collect()
2184}
2185
2186struct DaemonStageReportingExecutor<E> {
2187 inner: E,
2188 state: Arc<RuntimeState>,
2189 session_id: String,
2190}
2191
2192impl<E> DaemonStageReportingExecutor<E> {
2193 fn new(inner: E, state: Arc<RuntimeState>, session_id: String) -> Self {
2194 Self {
2195 inner,
2196 state,
2197 session_id,
2198 }
2199 }
2200}
2201
2202impl<E: CommandExecutor> CommandExecutor for DaemonStageReportingExecutor<E> {
2203 fn execute(&self, command: &CommandSpec) -> Result<CommandOutput> {
2204 let _stage = command
2205 .stage
2206 .map(|stage| ProvisionStageGuard::new(self, stage));
2207 let started = std::time::Instant::now();
2208 let result = self.inner.execute(command);
2209 tracing::info!(
2210 session_id = %self.session_id,
2211 stage = command
2212 .stage
2213 .map(ProvisionStage::label)
2214 .unwrap_or_else(|| "command".to_owned()),
2215 purpose = %command.purpose,
2216 duration_ms = started.elapsed().as_millis(),
2217 succeeded = result.as_ref().is_ok_and(|output| output.status == 0),
2218 "session command stage finished"
2219 );
2220 result
2221 }
2222
2223 fn execute_with_stdin(
2224 &self,
2225 command: &CommandSpec,
2226 input: &mut (dyn std::io::Read + Send),
2227 ) -> Result<CommandOutput> {
2228 let _stage = command
2229 .stage
2230 .map(|stage| ProvisionStageGuard::new(self, stage));
2231 let started = std::time::Instant::now();
2232 let result = self.inner.execute_with_stdin(command, input);
2233 tracing::info!(
2234 session_id = %self.session_id,
2235 stage = command
2236 .stage
2237 .map(ProvisionStage::label)
2238 .unwrap_or_else(|| "command".to_owned()),
2239 purpose = %command.purpose,
2240 duration_ms = started.elapsed().as_millis(),
2241 succeeded = result.as_ref().is_ok_and(|output| output.status == 0),
2242 "session streaming command stage finished"
2243 );
2244 result
2245 }
2246
2247 fn cancellation_requested(&self) -> bool {
2248 self.inner.cancellation_requested()
2249 }
2250
2251 fn stage_started(&self, stage: ProvisionStage) {
2252 self.state
2253 .change_lifecycle_stage(&self.session_id, stage, true);
2254 }
2255
2256 fn stage_finished(&self, stage: ProvisionStage) {
2257 self.state
2258 .change_lifecycle_stage(&self.session_id, stage, false);
2259 }
2260
2261 fn notify_notice(&self, notice: &str) {
2262 self.state.set_lifecycle_notice(&self.session_id, notice);
2263 }
2264}
2265
2266fn epoch_seconds() -> u64 {
2267 SystemTime::now()
2268 .duration_since(UNIX_EPOCH)
2269 .unwrap_or_default()
2270 .as_secs()
2271}
2272
2273fn random_hex<const N: usize>() -> Result<String> {
2274 let mut bytes = [0_u8; N];
2275 getrandom::fill(&mut bytes).map_err(|error| anyhow!("generate daemon secret: {error}"))?;
2276 Ok(bytes.iter().map(|byte| format!("{byte:02x}")).collect())
2277}
2278
2279fn write_metadata(path: &Path, metadata: &DaemonMetadata) -> Result<()> {
2280 let parent = path
2281 .parent()
2282 .context("daemon metadata path has no parent")?;
2283 fs::create_dir_all(parent)
2284 .with_context(|| format!("create daemon data directory {}", parent.display()))?;
2285 let temporary = parent.join(format!(".daemon.{}.tmp", std::process::id()));
2286 let body = serde_json::to_vec_pretty(metadata)?;
2287 let mut options = OpenOptions::new();
2288 options.write(true).create_new(true);
2289 #[cfg(unix)]
2290 {
2291 use std::os::unix::fs::OpenOptionsExt;
2292 options.mode(0o600);
2293 }
2294 let mut file = options
2295 .open(&temporary)
2296 .with_context(|| format!("create {}", temporary.display()))?;
2297 file.write_all(&body)?;
2298 file.sync_all()?;
2299 fs::rename(&temporary, path)
2300 .with_context(|| format!("publish daemon metadata {}", path.display()))?;
2301 Ok(())
2302}
2303
2304fn owner_pid_to_watch() -> Result<Option<u32>> {
2310 let Some(value) = mj_core::config::env_override("DAEMON_OWNER_PID") else {
2311 return Ok(None);
2312 };
2313 let pid: u32 = value
2314 .trim()
2315 .parse()
2316 .map_err(|_| anyhow!("MJ_DAEMON_OWNER_PID must be a process id, but it is {value:?}"))?;
2317 ensure!(
2318 process_is_alive(pid),
2319 "MJ_DAEMON_OWNER_PID names process {pid}, which is not running"
2320 );
2321 Ok(Some(pid))
2322}
2323
2324pub async fn run_daemon_process() -> Result<()> {
2325 let owner_pid = owner_pid_to_watch()?;
2328 let guard = ControllerStoreGuard::acquire()?;
2329 let database_writer = guard.start_database_writer()?;
2330 let epilogue_started = AtomicBool::new(false);
2331 let mut outcome = run_daemon_runtime(&epilogue_started, owner_pid).await;
2332 if !epilogue_started.load(Ordering::Acquire) {
2333 spawn_shutdown_watchdog();
2336 }
2337 let writer_shutdown = tokio::task::spawn_blocking(move || database_writer.shutdown())
2338 .await
2339 .context("database writer shutdown task panicked")
2340 .and_then(std::convert::identity);
2341 record_daemon_cleanup(&mut outcome, "shut down database writer", writer_shutdown);
2342 outcome
2343}
2344
2345async fn run_daemon_runtime(epilogue_started: &AtomicBool, owner_pid: Option<u32>) -> Result<()> {
2346 tokio::task::spawn_blocking(crate::controller::pin_worker_binary_sources)
2349 .await
2350 .context("worker source snapshot task failed")??;
2351 Controller::recover_config_id_rename()?;
2352 Config::migrate_legacy_localhost_target()?;
2353 let config = Config::load()?;
2354 crate::database::recover_interrupted_checkpointing_sessions(
2355 &chrono::Utc::now().to_rfc3339_opts(chrono::SecondsFormat::Secs, true),
2356 )?;
2357 crate::controller::reconcile_managed_checkpoint_archives()?;
2358
2359 let controller = Controller::load()?;
2360 let listener = TcpListener::bind(SocketAddr::new(IpAddr::V4(Ipv4Addr::LOCALHOST), 0))
2361 .await
2362 .context("bind Mjolnir daemon loopback endpoint")?;
2363 let metadata = DaemonMetadata {
2364 protocol_version: PROTOCOL_VERSION,
2365 pid: std::process::id(),
2366 address: listener.local_addr()?,
2367 token: random_hex::<32>()?,
2368 started_at: chrono::Utc::now().to_rfc3339_opts(chrono::SecondsFormat::Secs, true),
2369 build_version: env!("CARGO_PKG_VERSION").to_owned(),
2370 };
2371 let workspaces = tokio::task::spawn_blocking(crate::database::list_workspaces)
2372 .await
2373 .context("daemon workspace load task panicked")??;
2374 let mut remote = if config.phone.enabled {
2375 Some(spawn_remote_session_manager()?)
2376 } else {
2377 None
2378 };
2379
2380 let manager = spawn_session_manager()?;
2383 let manager_targets = manager.targets;
2384 manager_targets.send_replace(dashboard_worker_targets(&controller));
2385 let mut manager_updates = manager.updates;
2386 let manager_control = manager.control.clone();
2387 let manager_shutdown = manager.shutdown;
2388 let mut recovery = crate::recovery::RecoveryCoordinator::spawn(manager_control.clone());
2389 let recovery_observer = recovery.observer();
2390 let mut worker_upgrades = crate::worker_upgrade::WorkerUpgradeCoordinator::spawn(
2393 manager_control.clone(),
2394 &recovery_observer,
2395 );
2396 let state = Arc::new(RuntimeState::new(
2397 manager_control.clone(),
2398 Controller {
2399 config: controller.config.clone(),
2400 state: controller.state.clone(),
2401 },
2402 recovery_observer.clone(),
2403 worker_upgrades.observer(),
2404 workspaces,
2405 ));
2406 let move_operations = blocking(crate::database::load_move_operations).await?;
2407 let move_owned = state.recover_moves(move_operations)?;
2408 state.resume_retained_cleanups();
2409 let cancellation = crate::termination::Coordinator::install().token();
2410 let target_refresh = spawn_manager_target_refresher(
2411 manager_targets.clone(),
2412 cancellation.clone(),
2413 state.clone(),
2414 );
2415 let image_refresh = spawn_image_refresher(
2416 {
2417 let state = state.clone();
2418 move || state.with_config(crate::controller::image_refresh_plan)
2419 },
2420 cancellation.clone(),
2421 );
2422 let exit_when_idle = mj_core::config::env_override_os("DAEMON_EXIT_WHEN_IDLE").is_some();
2423 let mut idle_tick = tokio::time::interval(Duration::from_millis(100));
2424 idle_tick.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay);
2425 let mut owner_tick = tokio::time::interval(Duration::from_millis(500));
2426 owner_tick.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay);
2427 let mut recovery_tick = tokio::time::interval(Duration::from_millis(250));
2428 recovery_tick.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay);
2429 let (interrupted_close_tx, mut interrupted_close_rx) = tokio::sync::mpsc::unbounded_channel();
2430 let mut interrupted_close_cancellations = Vec::new();
2431 let mut interrupted_close_tasks = Vec::new();
2432 for session_id in interrupted_close_session_ids(&controller) {
2433 if move_owned.contains(&session_id) {
2434 continue;
2435 }
2436 let interrupted_cancellation = Arc::new(AtomicBool::new(false));
2437 let interrupted_close_task = spawn_interrupted_close_recovery(
2438 session_id,
2439 manager_control.clone(),
2440 recovery_observer.clone(),
2441 interrupted_cancellation.clone(),
2442 interrupted_close_tx.clone(),
2443 None,
2444 );
2445 interrupted_close_cancellations.push(interrupted_cancellation);
2446 interrupted_close_tasks.push(interrupted_close_task);
2447 }
2448 let mut phone_publisher: Option<RemoteSessionPublisher> = None;
2449 let mut phone_task = None;
2450 let mut remote_request_bridge = None;
2451 if let Some(remote) = remote.take() {
2452 remote
2453 .targets
2454 .send_replace(dashboard_worker_targets(&controller));
2455 phone_publisher = Some(remote.publisher.clone());
2456 remote_request_bridge = Some(spawn_remote_request_bridge(
2457 remote.requests,
2458 manager_control.clone(),
2459 ));
2460 phone_task = Some(spawn_phone_server(
2461 config.phone,
2462 cancellation.clone(),
2463 state.clone(),
2464 SessionManagerChannels {
2465 targets: remote.targets,
2466 control: remote.control,
2467 updates: remote.updates,
2468 shutdown: remote.shutdown,
2469 },
2470 ));
2471 } else {
2472 state.set_phone_status(WebViewerStatus::Disabled);
2473 state.web_viewer.publish(crate::server::WebViewerAccess::Unavailable("Web access is disabled. Enable [phone].enabled in your configuration, then restart the daemon.".into()));
2474 }
2475 let daemon_metadata_path = metadata_path();
2476 let mut client_tasks = tokio::task::JoinSet::new();
2477
2478 let mut outcome = async {
2482 write_metadata(&daemon_metadata_path, &metadata)?;
2483 reach_test_hook("daemon_metadata_before_listening").await?;
2484 loop {
2485 tokio::select! {
2486 _ = cancellation.cancelled() => break,
2487 _ = idle_tick.tick(), if exit_when_idle && state.ever_attached.load(Ordering::Acquire) => {
2488 state.prune_dead_clients();
2489 if state.attachments().is_empty() {
2490 break;
2491 }
2492 }
2493 _ = owner_tick.tick(), if owner_pid.is_some() => {
2494 if let Some(owner) = owner_pid
2495 && !process_is_alive(owner)
2496 {
2497 tracing::info!(owner_pid = owner, "daemon owner process exited; shutting down");
2498 break;
2499 }
2500 }
2501 _ = recovery_tick.tick() => {
2502 while let Some(result) = recovery.try_result() {
2503 if let Err(error) = &result.outcome {
2504 if result.deferred {
2508 tracing::info!(session_id = %result.session_id, %error, "recovery copy deferred: agent is working");
2509 } else {
2510 tracing::warn!(session_id = %result.session_id, %error, "daemon recovery checkpoint failed");
2511 }
2512 }
2513 refresh_runtime_controller(&state).await;
2514 }
2515 while let Some(result) = worker_upgrades.try_result() {
2516 report_worker_upgrade(&state, &result);
2517 }
2518 }
2519 completed = interrupted_close_rx.recv() => {
2520 if let Some(completed) = completed {
2521 let recovered = completed.result.is_ok();
2522 if let Err(error) = completed.result {
2523 tracing::warn!(session_id = %completed.session_id, %error, "daemon could not resume interrupted close");
2524 }
2525 refresh_runtime_controller(&state).await;
2526 if recovered && completed.deferred_cleanup
2527 && let Err(error) = state.start_deferred_cleanup(completed.session_id.clone())
2528 {
2529 tracing::warn!(session_id = %completed.session_id, error = format!("{error:#}"), "could not continue cleanup after interrupted close");
2530 state.push_notice(
2531 &completed.session_id,
2532 format!("Could not continue container storage cleanup: {error:#}"),
2533 );
2534 }
2535 }
2536 }
2537 accepted = listener.accept() => {
2538 let (stream, peer) = accepted.context("accept Mjolnir daemon client")?;
2539 if !peer.ip().is_loopback() {
2540 tracing::warn!(%peer, "rejected non-loopback daemon client");
2541 continue;
2542 }
2543 let metadata = metadata.clone();
2544 let state = state.clone();
2545 let cancellation = cancellation.clone();
2546 client_tasks.spawn(async move {
2547 if let Err(error) = serve_client(stream, metadata, state, cancellation).await {
2548 tracing::debug!(error = format!("{error:#}"), "daemon client disconnected");
2549 }
2550 });
2551 }
2552 completed = client_tasks.join_next(), if !client_tasks.is_empty() => {
2553 if let Some(Err(error)) = completed {
2554 tracing::warn!(%error, "daemon client task failed");
2555 }
2556 }
2557 update = manager_updates.recv() => {
2558 let Some(update) = update else {
2559 bail!("controller daemon session manager stopped");
2560 };
2561 if let Some((detail, observed_updated_at)) =
2562 state.missing_target_record(&update.session_id, &update.view)
2563 {
2564 let state = state.clone();
2565 let session_id = update.session_id.clone();
2566 client_tasks.spawn(async move {
2567 if let Err(error) = state.persist_missing_target(
2568 &session_id, detail, observed_updated_at,
2569 ).await {
2570 tracing::warn!(%session_id, %error, "could not persist missing worker target");
2571 state.push_notice(&session_id, format!("Could not record missing session target: {error:#}"));
2572 }
2573 });
2574 }
2575 if let Some(publisher) = phone_publisher.as_ref()
2576 && let Err(error) = publisher.try_publish(
2577 update.session_id.clone(),
2578 update.view.clone(),
2579 )
2580 {
2581 tracing::warn!(%error, "phone session view bridge stopped");
2582 phone_publisher = None;
2583 }
2584 state.review_host().observe(&update.session_id, &update.view);
2588 state.publish_session(update.session_id, update.view).await?;
2589 }
2590 }
2591 }
2592 Ok(())
2593 }
2594 .await;
2595
2596 epilogue_started.store(true, Ordering::Release);
2597 spawn_shutdown_watchdog();
2598 cancellation.cancel();
2601 for interrupted_cancellation in interrupted_close_cancellations {
2602 interrupted_cancellation.store(true, Ordering::Release);
2603 }
2604 drop(interrupted_close_tx);
2605 record_daemon_cleanup(
2606 &mut outcome,
2607 "remove daemon metadata",
2608 remove_daemon_metadata(&daemon_metadata_path),
2609 );
2610 record_daemon_cleanup(
2611 &mut outcome,
2612 "shut down turn review host",
2613 state
2614 .review_host()
2615 .shutdown()
2616 .await
2617 .map_err(anyhow::Error::msg),
2618 );
2619 record_daemon_cleanup(
2620 &mut outcome,
2621 "join controller target refresher",
2622 target_refresh.await.map_err(anyhow::Error::new),
2623 );
2624 record_daemon_cleanup(
2625 &mut outcome,
2626 "join container image refresher",
2627 image_refresh.await.map_err(anyhow::Error::new),
2628 );
2629 if let Some(phone_task) = phone_task {
2630 record_daemon_cleanup(
2631 &mut outcome,
2632 "join phone server",
2633 phone_task.await.map_err(anyhow::Error::new),
2634 );
2635 }
2636 if let Some(remote_request_bridge) = remote_request_bridge {
2637 record_daemon_cleanup(
2638 &mut outcome,
2639 "join phone session request bridge",
2640 remote_request_bridge.await.map_err(anyhow::Error::new),
2641 );
2642 }
2643 client_tasks.abort_all();
2644 while let Some(result) = client_tasks.join_next().await {
2645 if let Err(error) = result
2646 && !error.is_cancelled()
2647 {
2648 record_daemon_cleanup(
2649 &mut outcome,
2650 "join daemon client task",
2651 Err(anyhow::Error::new(error)),
2652 );
2653 }
2654 }
2655 record_daemon_cleanup(
2656 &mut outcome,
2657 "cancel daemon lifecycle operations",
2658 state.cancel_and_wait_lifecycles().await,
2659 );
2660 for interrupted_close_task in interrupted_close_tasks {
2661 record_daemon_cleanup(
2662 &mut outcome,
2663 "join interrupted close recovery",
2664 interrupted_close_task.await.map_err(anyhow::Error::new),
2665 );
2666 }
2667 drop(recovery);
2668 record_daemon_cleanup(
2669 &mut outcome,
2670 "shut down controller daemon session manager",
2671 manager_shutdown.shutdown().await,
2672 );
2673 outcome
2674}
2675
2676fn remove_daemon_metadata(path: &Path) -> Result<()> {
2677 match fs::remove_file(path) {
2678 Ok(()) => Ok(()),
2679 Err(error) if error.kind() == std::io::ErrorKind::NotFound => Ok(()),
2680 Err(error) => Err(error).with_context(|| format!("remove {}", path.display())),
2681 }
2682}
2683
2684fn record_daemon_cleanup(outcome: &mut Result<()>, operation: &'static str, cleanup: Result<()>) {
2688 let Err(error) = cleanup else {
2689 return;
2690 };
2691 let error = error.context(operation);
2692 if outcome.is_ok() {
2693 *outcome = Err(error);
2694 } else {
2695 tracing::warn!(error = format!("{error:#}"), "daemon cleanup step failed");
2696 }
2697}
2698
2699fn spawn_shutdown_watchdog() {
2706 tokio::spawn(async move {
2707 tokio::time::sleep(SHUTDOWN_FORCE_EXIT_TIMEOUT).await;
2708 tracing::error!(
2709 seconds = SHUTDOWN_FORCE_EXIT_TIMEOUT.as_secs(),
2710 "daemon shutdown did not finish in time; exiting"
2711 );
2712 if let Err(error) = fs::remove_file(metadata_path())
2715 && error.kind() != std::io::ErrorKind::NotFound
2716 {
2717 tracing::warn!(%error, "could not remove daemon metadata before the forced exit");
2718 }
2719 std::process::exit(1);
2721 });
2722}
2723
2724fn spawn_manager_target_refresher(
2725 targets: tokio::sync::watch::Sender<Vec<crate::session_manager::RelaySessionTarget>>,
2726 cancellation: CancellationToken,
2727 state: Arc<RuntimeState>,
2728) -> tokio::task::JoinHandle<()> {
2729 tokio::spawn(async move {
2730 let mut interval = tokio::time::interval(Duration::from_millis(500));
2731 interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay);
2732 loop {
2733 tokio::select! {
2734 _ = cancellation.cancelled() => return,
2735 _ = interval.tick() => {
2736 let _config_mutation = state.config_mutation.lock().await;
2739 match tokio::task::spawn_blocking(Controller::load).await {
2740 Ok(Ok(controller)) => {
2741 let lifecycle_sessions =
2746 state.worker_poll_exclusion_session_ids(&controller);
2747 let refreshed = dashboard_worker_targets_excluding(
2748 &controller,
2749 &lifecycle_sessions,
2750 );
2751 let changed = {
2752 let mut review = state
2753 .review_config
2754 .lock()
2755 .unwrap_or_else(PoisonError::into_inner);
2756 review.clone_from(&controller.config.review);
2757 drop(review);
2758 let mut current = state
2759 .controller
2760 .lock()
2761 .unwrap_or_else(PoisonError::into_inner);
2762 let changed = current.config != controller.config;
2763 *current = controller;
2764 changed
2765 };
2766 targets.send_replace(refreshed);
2767 if changed {
2768 state.publish_revision();
2769 }
2770 }
2771 Ok(Err(error)) => {
2772 if let Some(mismatch) = error
2781 .chain()
2782 .find_map(|cause| cause.downcast_ref::<StoreSchemaMismatch>())
2783 {
2784 tracing::error!(
2785 found = mismatch.found,
2786 supported = mismatch.supported,
2787 error = %mismatch,
2788 "daemon store schema diverged underneath the daemon; shutting down"
2789 );
2790 cancellation.cancel();
2791 return;
2792 }
2793 tracing::warn!(error = format!("{error:#}"), "could not refresh daemon session targets");
2794 }
2795 Err(error) => {
2796 tracing::error!(%error, "daemon target refresh task failed");
2797 return;
2798 }
2799 }
2800 }
2801 }
2802 }
2803 })
2804}
2805
2806async fn refresh_runtime_controller(state: &RuntimeState) {
2807 if let Err(error) = state.reload_controller().await {
2808 tracing::warn!(
2809 error = format!("{error:#}"),
2810 "could not refresh daemon controller state"
2811 );
2812 }
2813}
2814
2815async fn refresh_runtime_workspaces(state: &RuntimeState) -> Result<()> {
2816 let workspaces = tokio::task::spawn_blocking(crate::database::list_workspaces)
2817 .await
2818 .context("daemon workspace refresh task panicked")??;
2819 state.publish_workspaces(workspaces);
2820 Ok(())
2821}
2822
2823fn spawn_phone_server(
2824 config: mj_core::config::PhoneConfig,
2825 cancellation: CancellationToken,
2826 state: Arc<RuntimeState>,
2827 worker: SessionManagerChannels,
2828) -> tokio::task::JoinHandle<()> {
2829 state.set_phone_status(WebViewerStatus::Starting);
2830 let workspaces = state.workspaces();
2831 tokio::spawn(async move {
2832 match crate::server_runtime::run_server(
2833 (&config).into(),
2834 cancellation.clone(),
2835 worker,
2836 state.clone(),
2837 workspaces,
2838 )
2839 .await
2840 {
2841 Ok(()) if cancellation.is_cancelled() => {}
2842 Ok(()) => {
2843 state.set_phone_status(WebViewerStatus::Stopped);
2844 state.web_viewer.publish(crate::server::WebViewerAccess::Unavailable("The web viewer stopped unexpectedly. Restart the daemon to restore web access.".into()));
2845 }
2846 Err(error) => {
2847 tracing::warn!(error = format!("{error:#}"), "phone server stopped");
2848 state.publish_web_access(crate::server::WebViewerAccess::Unavailable(format!(
2849 "Could not start the web viewer: {error:#}"
2850 )));
2851 }
2852 }
2853 })
2854}
2855
2856fn spawn_remote_request_bridge(
2857 mut requests: crate::session_manager::RemoteSessionRequests,
2858 manager: SessionManagerControl,
2859) -> tokio::task::JoinHandle<()> {
2860 tokio::spawn(async move {
2861 let mut request_order = crate::session_manager::SessionRequestOrder::new();
2864 while let Some(request) = requests.recv().await {
2865 let manager = manager.clone();
2866 request_order.dispatch(request, move |request| {
2867 forward_in_process_session_request(request, manager)
2868 });
2869 }
2870 })
2871}
2872
2873async fn forward_in_process_session_request(
2874 request: RemoteSessionRequest,
2875 manager: SessionManagerControl,
2876) {
2877 match request {
2878 RemoteSessionRequest::Submit {
2879 session_id,
2880 command_id,
2881 command,
2882 admission,
2883 reply,
2884 } => {
2885 if admission.is_some() {
2886 let _ = reply.send(Err(
2887 "review delivery admissions cannot cross the daemon request bridge".into(),
2888 ));
2889 return;
2890 }
2891 let result = async {
2892 manager
2893 .wait_for_session(&session_id, Duration::from_secs(5))
2894 .await?
2895 .submit(command_id, command)
2896 .await
2897 }
2898 .await
2899 .map_err(|error| format!("{error:#}"));
2900 let _ = reply.send(result);
2901 }
2902 RemoteSessionRequest::Sync { session_id, reply } => {
2903 let result = async { manager.session(session_id).await?.sync_now().await }
2904 .await
2905 .map_err(|error| format!("{error:#}"));
2906 let _ = reply.send(result);
2907 }
2908 RemoteSessionRequest::RespondElicitation {
2909 session_id,
2910 elicitation_id,
2911 response,
2912 reply,
2913 } => {
2914 let result = async {
2915 manager
2916 .session(session_id)
2917 .await?
2918 .respond_elicitation(elicitation_id, response)
2919 .await
2920 }
2921 .await
2922 .map_err(|error| format!("{error:#}"));
2923 let _ = reply.send(result);
2924 }
2925 RemoteSessionRequest::StopBackgroundTask {
2926 session_id,
2927 background_task_id,
2928 reply,
2929 } => {
2930 let result = async {
2931 manager
2932 .session(session_id)
2933 .await?
2934 .stop_background_task(background_task_id)
2935 .await
2936 }
2937 .await
2938 .map_err(|error| format!("{error:#}"));
2939 let _ = reply.send(result);
2940 }
2941 RemoteSessionRequest::Reviewer {
2942 session_id,
2943 role,
2944 action,
2945 mut reply,
2946 } => {
2947 let result = tokio::select! {
2948 _ = reply.closed() => return,
2949 result = async {
2950 manager
2951 .session(session_id)
2952 .await?
2953 .reviewer_as(role, action)
2954 .await
2955 } => result,
2956 }
2957 .map_err(|error| format!("{error:#}"));
2958 let _ = reply.send(result);
2959 }
2960 }
2961}
2962
2963async fn serve_client(
2964 mut stream: TcpStream,
2965 metadata: DaemonMetadata,
2966 state: Arc<RuntimeState>,
2967 cancellation: CancellationToken,
2968) -> Result<()> {
2969 loop {
2970 let request: RequestEnvelope = match read_frame(&mut stream).await {
2971 Ok(request) => request,
2972 Err(error)
2973 if error.downcast_ref::<std::io::Error>().is_some_and(|io| {
2974 matches!(
2975 io.kind(),
2976 std::io::ErrorKind::UnexpectedEof
2977 | std::io::ErrorKind::ConnectionReset
2978 | std::io::ErrorKind::BrokenPipe
2979 )
2980 }) =>
2981 {
2982 return Ok(());
2983 }
2984 Err(error) => return Err(error),
2985 };
2986 let request_id = request.request_id;
2987 let is_management = matches!(
2991 request.action,
2992 DaemonAction::Ping | DaemonAction::Status | DaemonAction::Stop
2993 );
2994 let result = if request.token != metadata.token {
2995 Err("daemon authentication failed".to_owned())
2996 } else if request.protocol_version != PROTOCOL_VERSION && !is_management {
2997 Err(format!(
2998 "incompatible daemon protocol {}; expected {}",
2999 request.protocol_version, PROTOCOL_VERSION
3000 ))
3001 } else if cancellation.is_cancelled() && !is_management {
3002 Err("daemon is shutting down; retry to reach a fresh daemon".to_owned())
3010 } else {
3011 let reviewer = matches!(&request.action, DaemonAction::ReviewerAction { .. });
3012 if reviewer {
3013 let mut peer_probe = [0_u8; 1];
3019 let mut action = Box::pin(handle_action(
3020 request.action,
3021 &metadata,
3022 &state,
3023 &cancellation,
3024 ));
3025 tokio::select! {
3026 result = &mut action => result.map_err(|error| format!("{error:#}")),
3027 peer = stream.peek(&mut peer_probe) => {
3028 match peer {
3029 Ok(0) => return Ok(()),
3030 Ok(_) => action.await.map_err(|error| format!("{error:#}")),
3031 Err(error) => {
3032 tracing::debug!(%error, "reviewer client connection became unreadable");
3033 return Ok(());
3034 }
3035 }
3036 }
3037 }
3038 } else {
3039 handle_action(request.action, &metadata, &state, &cancellation)
3040 .await
3041 .map_err(|error| format!("{error:#}"))
3042 }
3043 };
3044 write_frame(
3047 &mut stream,
3048 &ResponseEnvelope {
3049 protocol_version: request.protocol_version,
3050 request_id,
3051 result,
3052 },
3053 )
3054 .await?;
3055 }
3056}
3057
3058async fn blocking<T: Send + 'static>(
3059 work: impl FnOnce() -> Result<T> + Send + 'static,
3060) -> Result<T> {
3061 tokio::task::spawn_blocking(work)
3062 .await
3063 .context("daemon background database task panicked")?
3064}
3065
3066async fn reach_test_hook(name: &'static str) -> Result<()> {
3067 #[cfg(feature = "test-hooks")]
3068 {
3069 tokio::task::spawn_blocking(move || mj_core::test_hooks::reach_test_hook(name))
3070 .await
3071 .context("test hook task panicked")??;
3072 }
3073 #[cfg(not(feature = "test-hooks"))]
3074 let _ = name;
3075 Ok(())
3076}
3077
3078async fn handle_action(
3079 action: DaemonAction,
3080 metadata: &DaemonMetadata,
3081 state: &Arc<RuntimeState>,
3082 cancellation: &CancellationToken,
3083) -> Result<DaemonReply> {
3084 match action {
3085 DaemonAction::Ping => Ok(DaemonReply::Pong),
3086 DaemonAction::Status => {
3087 state.prune_dead_clients();
3088 Ok(DaemonReply::Status(DaemonStatus {
3089 pid: metadata.pid,
3090 started_at: metadata.started_at.clone(),
3091 build_version: metadata.build_version.clone(),
3092 attached_clients: state.attachments().len(),
3093 phone_status: state.phone_status(),
3094 }))
3095 }
3096 DaemonAction::WebViewerAccess => {
3097 Ok(DaemonReply::WebViewerAccess(state.web_viewer.access()))
3098 }
3099 DaemonAction::RecoverWebViewer(action) => {
3100 state.web_viewer.recover(action)?;
3101 Ok(DaemonReply::Done)
3102 }
3103 DaemonAction::InspectWebListener => {
3104 let address = state.web_viewer.conflict_address()?;
3105 let processes = blocking(move || crate::web_viewer::inspect_listener(address)).await?;
3106 ensure!(
3107 state.web_viewer.conflict_address()? == address,
3108 "The viewer address changed. Inspect again."
3109 );
3110 Ok(DaemonReply::WebListeners(processes))
3111 }
3112 DaemonAction::ListWorkspaces => {
3113 state.prune_dead_clients();
3114 let workspaces = blocking(crate::database::list_workspaces).await?;
3115 Ok(DaemonReply::Workspaces(
3116 workspaces
3117 .into_iter()
3118 .map(|workspace| WorkspaceListing { workspace })
3119 .collect(),
3120 ))
3121 }
3122 DaemonAction::CreateWorkspace { name } => {
3123 let workspace =
3127 blocking(move || crate::database::create_or_get_workspace(&name)).await?;
3128 refresh_runtime_workspaces(state).await?;
3129 Ok(DaemonReply::Workspace(workspace))
3130 }
3131 DaemonAction::RenameWorkspace { workspace_id, name } => {
3132 blocking(move || crate::database::rename_workspace(&workspace_id, &name)).await?;
3133 refresh_runtime_workspaces(state).await?;
3134 Ok(DaemonReply::Done)
3135 }
3136 DaemonAction::TouchWorkspace { workspace_id } => {
3137 blocking(move || crate::database::touch_workspace(&workspace_id)).await?;
3138 refresh_runtime_workspaces(state).await?;
3139 Ok(DaemonReply::Done)
3140 }
3141 DaemonAction::DeleteWorkspace { workspace_id } => {
3142 ensure!(
3143 !state.workspace_has_active_resume(&workspace_id),
3144 "workspace has a session resume in progress"
3145 );
3146 blocking(move || crate::database::delete_workspace(&workspace_id)).await?;
3147 refresh_runtime_workspaces(state).await?;
3148 Ok(DaemonReply::Done)
3149 }
3150 DaemonAction::Attach { client_id, pid } => {
3151 state.attachments().insert(client_id, Attachment { pid });
3152 state.ever_attached.store(true, Ordering::Release);
3153 Ok(DaemonReply::Done)
3154 }
3155 DaemonAction::Detach { client_id } => {
3156 state.attachments().remove(&client_id);
3157 Ok(DaemonReply::Done)
3158 }
3159 DaemonAction::PersistReadReceipt {
3160 client_id,
3161 workspace_id,
3162 session_id,
3163 through,
3164 } => {
3165 let frontier = blocking(move || {
3166 crate::database::persist_read_receipt(
3167 &client_id,
3168 &workspace_id,
3169 &session_id,
3170 through,
3171 )
3172 })
3173 .await?;
3174 Ok(DaemonReply::Ordinal(frontier))
3175 }
3176 DaemonAction::PersistDetachedSessionState {
3177 client_id,
3178 workspace_id,
3179 session_id,
3180 through,
3181 owner_pid,
3182 draft,
3183 } => {
3184 blocking(move || {
3185 let receipt = crate::database::persist_read_receipt(
3186 &client_id,
3187 &workspace_id,
3188 &session_id,
3189 through,
3190 )
3191 .map(|_| ());
3192 let saved_draft = crate::database::save_detached_session_draft(
3195 &workspace_id,
3196 &session_id,
3197 &client_id,
3198 owner_pid,
3199 draft,
3200 )
3201 .map(|_| ());
3202 receipt.and(saved_draft)
3203 })
3204 .await?;
3205 Ok(DaemonReply::Done)
3206 }
3207 DaemonAction::SaveActiveReview { session_id, review } => {
3208 blocking(move || crate::database::save_active_review(&session_id, &review)).await?;
3209 Ok(DaemonReply::Done)
3210 }
3211 DaemonAction::ClearActiveReview { session_id } => {
3212 blocking(move || crate::database::clear_active_review(&session_id)).await?;
3213 Ok(DaemonReply::Done)
3214 }
3215 DaemonAction::RememberReviewerSelection {
3216 workspace_id,
3217 selection,
3218 } => {
3219 blocking(move || {
3220 crate::database::remember_reviewer_selection(&workspace_id, &selection)
3221 })
3222 .await?;
3223 Ok(DaemonReply::Done)
3224 }
3225 DaemonAction::SaveWorkspacePaneSizes {
3226 workspace_id,
3227 sizes,
3228 } => {
3229 blocking(move || crate::database::save_workspace_pane_sizes(&workspace_id, sizes))
3230 .await?;
3231 Ok(DaemonReply::Done)
3232 }
3233 DaemonAction::PersistImportedSession { session } => {
3234 blocking(move || crate::import::persist_imported_session_locally(&session)).await?;
3235 refresh_runtime_controller(state).await;
3236 Ok(DaemonReply::Done)
3237 }
3238 DaemonAction::SetSessionTitle { session_id, title } => {
3239 let title =
3240 blocking(move || Controller::load()?.rename_session(&session_id, &title)).await?;
3241 refresh_runtime_controller(state).await;
3242 Ok(DaemonReply::Text(title))
3243 }
3244 DaemonAction::SetSessionContainerSettings {
3245 session_id,
3246 cpus,
3247 memory,
3248 mounts,
3249 mount_history,
3250 } => {
3251 ensure!(
3252 !crate::controller::move_session::move_owns_session(&session_id),
3253 "session is moving; change container settings after Move finishes"
3254 );
3255 blocking(move || {
3256 Controller::load()?.update_session_container_settings(
3257 &session_id,
3258 cpus,
3259 memory,
3260 mounts,
3261 mount_history,
3262 )
3263 })
3264 .await?;
3265 refresh_runtime_controller(state).await;
3266 Ok(DaemonReply::Done)
3267 }
3268 DaemonAction::SetSessionAcpTitle { session_id, title } => {
3269 blocking(move || crate::database::set_session_acp_title(&session_id, title.as_deref()))
3270 .await?;
3271 refresh_runtime_controller(state).await;
3272 Ok(DaemonReply::Done)
3273 }
3274 DaemonAction::MarkSessionTargetMissing {
3275 session_id,
3276 detail,
3277 updated_at,
3278 } => {
3279 let changed = blocking(move || {
3280 crate::database::mark_session_target_missing(&session_id, &detail, &updated_at)
3281 })
3282 .await?;
3283 refresh_runtime_controller(state).await;
3284 Ok(DaemonReply::OptionalSessionState(changed))
3285 }
3286 DaemonAction::CheckpointSession { session_id } => Ok(DaemonReply::Checkpoint(
3287 state.checkpoint_session_now(&session_id).await?,
3288 )),
3289 DaemonAction::ScanRecovery => {
3290 let scan =
3291 blocking(|| Ok(Controller::load()?.scan_orphan_workers(&ProcessExecutor))).await?;
3292 Ok(DaemonReply::RecoveryScan(scan))
3293 }
3294 DaemonAction::AdoptRecovery {
3295 session_id,
3296 target_id,
3297 profile,
3298 bundle,
3299 } => {
3300 ensure_no_active_lifecycle(state)?;
3301 let mut controller = blocking(Controller::load).await?;
3302 controller
3303 .adopt_orphan_worker(
3304 &session_id,
3305 &target_id,
3306 profile.as_deref(),
3307 bundle.as_deref(),
3308 &ProcessExecutor,
3309 )
3310 .await?;
3311 refresh_runtime_controller(state).await;
3312 Ok(DaemonReply::Done)
3313 }
3314 DaemonAction::DestroyRecovery {
3315 session_id,
3316 target_id,
3317 confirmation,
3318 } => {
3319 ensure_no_active_lifecycle(state)?;
3320 blocking(move || {
3321 Controller::load()?.destroy_orphan_worker(
3322 &session_id,
3323 &target_id,
3324 &confirmation,
3325 &ProcessExecutor,
3326 )
3327 })
3328 .await?;
3329 Ok(DaemonReply::Done)
3330 }
3331 DaemonAction::Snapshot { workspace_id } => {
3332 let snapshot = blocking(move || workspace_snapshot(&workspace_id)).await?;
3333 Ok(DaemonReply::Snapshot(snapshot))
3334 }
3335 DaemonAction::RuntimeSnapshot {
3336 workspace_id,
3337 after_revision,
3338 all_workspaces,
3339 } => Ok(DaemonReply::RuntimeSnapshot(Box::new(
3340 state
3341 .runtime_snapshot(&workspace_id, after_revision, all_workspaces)
3342 .await?,
3343 ))),
3344 DaemonAction::RenameProfile { old_id, new_id } => {
3345 let _config_mutation = state.config_mutation.lock().await;
3346 ensure_no_active_lifecycle(state)?;
3347 let controller = blocking(move || {
3348 let mut controller = Controller::load()?;
3349 controller.rename_profile_id(&old_id, &new_id)?;
3350 Ok(controller)
3351 })
3352 .await?;
3353 install_renamed_controller(state, controller);
3354 Ok(DaemonReply::Done)
3355 }
3356 DaemonAction::RenameTarget { old_id, new_id } => {
3357 let _config_mutation = state.config_mutation.lock().await;
3358 ensure_no_active_lifecycle(state)?;
3359 let controller = blocking(move || {
3360 let mut controller = Controller::load()?;
3361 controller.rename_target_id(&old_id, &new_id)?;
3362 Ok(controller)
3363 })
3364 .await?;
3365 install_renamed_controller(state, controller);
3366 Ok(DaemonReply::Done)
3367 }
3368 DaemonAction::SubmitSessionCommand {
3369 inherited_draft,
3370 session_id,
3371 command_id,
3372 command,
3373 } => {
3374 let history = if let RelayCommand::Prompt { prompt } = &command {
3375 let values = serde_json::to_value(prompt)?;
3376 let values = values
3377 .as_array()
3378 .context("serialized prompt content is not an array")?;
3379 let text = mj_core::transcript::materialized_content_text(values);
3380 let bundle_id = state
3381 .controller
3382 .lock()
3383 .unwrap_or_else(PoisonError::into_inner)
3384 .state
3385 .sessions
3386 .get(&session_id)
3387 .with_context(|| format!("unknown session {session_id}"))?
3388 .bundle_id
3389 .clone();
3390 Some((bundle_id, text))
3391 } else {
3392 None
3393 };
3394 let session = state
3399 .session_manager
3400 .wait_for_session(&session_id, Duration::from_secs(5))
3401 .await?;
3402 let session_id = session.session_id().to_owned();
3403 let ordinal = session.submit(command_id, command).await?;
3404 if let Some(expected) = inherited_draft {
3405 let persisted_id = session_id.clone();
3406 let persisted_expected = expected.clone();
3407 blocking(move || {
3408 crate::database::clear_session_draft_input_if_matches(
3409 &persisted_id,
3410 &persisted_expected,
3411 )
3412 })
3413 .await?;
3414 if let Some(record) = state
3415 .controller
3416 .lock()
3417 .unwrap_or_else(PoisonError::into_inner)
3418 .state
3419 .sessions
3420 .get_mut(&session_id)
3421 && record.draft_input == expected
3422 {
3423 record.draft_input.clear();
3424 }
3425 state.publish_revision();
3426 }
3427 if let Some((bundle_id, text)) = history
3428 && let Err(error) = blocking(move || {
3429 crate::database::record_prompt(&session_id, &bundle_id, ordinal, None, &text)
3430 })
3431 .await
3432 {
3433 tracing::warn!(%error, "prompt was accepted but its history could not be stored");
3434 }
3435 Ok(DaemonReply::Ordinal(ordinal))
3436 }
3437 DaemonAction::ReviewerAction {
3438 session_id,
3439 role,
3440 action,
3441 } => {
3442 let session = state.session_manager.session(session_id).await?;
3443 Ok(DaemonReply::Reviewer(Box::new(
3444 session.reviewer_as(role, action).await?,
3445 )))
3446 }
3447 DaemonAction::StartTurnReview { session_id } => {
3448 state
3449 .review_host()
3450 .start(&session_id, true)
3451 .await
3452 .map_err(|refusal| anyhow!("{refusal}"))?;
3453 state.publish_revision();
3454 Ok(DaemonReply::Done)
3455 }
3456 DaemonAction::ResolveTurnReview {
3457 session_id,
3458 resolution,
3459 } => {
3460 state
3461 .review_host()
3462 .resolve(&session_id, resolution)
3463 .await
3464 .map_err(|error| anyhow!("{error}"))?;
3465 state.publish_revision();
3466 Ok(DaemonReply::Done)
3467 }
3468 DaemonAction::SyncSession { session_id } => {
3469 state
3470 .session_manager
3471 .session(session_id)
3472 .await?
3473 .sync_now()
3474 .await?;
3475 Ok(DaemonReply::Done)
3476 }
3477 DaemonAction::RespondElicitation {
3478 session_id,
3479 elicitation_id,
3480 response,
3481 } => {
3482 state
3483 .session_manager
3484 .session(session_id)
3485 .await?
3486 .respond_elicitation(elicitation_id, response)
3487 .await?;
3488 Ok(DaemonReply::Done)
3489 }
3490 DaemonAction::StopBackgroundTask {
3491 session_id,
3492 background_task_id,
3493 } => {
3494 state
3495 .session_manager
3496 .session(session_id)
3497 .await?
3498 .stop_background_task(background_task_id)
3499 .await?;
3500 Ok(DaemonReply::Done)
3501 }
3502 DaemonAction::CloseSession { session_id } => {
3503 state.close_session(session_id).await?;
3504 Ok(DaemonReply::Done)
3505 }
3506 DaemonAction::StartCreateSession(request) => Ok(DaemonReply::RegisteredSession(Box::new(
3507 state.start_create_session(request).await?,
3508 ))),
3509 DaemonAction::WaitCreateSession { session_id } => {
3510 state.wait_create_session(&session_id).await?;
3511 Ok(DaemonReply::Done)
3512 }
3513 DaemonAction::ResumeSession(request) => {
3514 state.resume_session(request).await?;
3515 Ok(DaemonReply::Done)
3516 }
3517 DaemonAction::PrepareMoveSession(selection) => Ok(DaemonReply::MovePreparation(Box::new(
3518 state.prepare_move_session(selection).await?,
3519 ))),
3520 DaemonAction::MoveSession(request) => {
3521 Ok(DaemonReply::MoveOutcome(state.move_session(request).await?))
3522 }
3523 DaemonAction::ForceStopSession { session_id } => {
3524 state.force_stop_session(session_id).await?;
3525 Ok(DaemonReply::Done)
3526 }
3527 DaemonAction::DestroyStoppedSession { session_id } => {
3528 state.destroy_stopped_session(session_id).await?;
3529 Ok(DaemonReply::Done)
3530 }
3531 DaemonAction::ForceDestroySession { session_id } => {
3532 state.force_destroy_session(session_id).await?;
3533 Ok(DaemonReply::Done)
3534 }
3535 DaemonAction::ForceDeleteWorkspace { workspace_id } => {
3536 state.force_delete_workspace(workspace_id).await?;
3537 Ok(DaemonReply::Done)
3538 }
3539 DaemonAction::CancelLifecycle { session_id } => {
3540 blocking({
3541 let session_id = session_id.clone();
3542 move || crate::database::request_move_cancellation(&session_id)
3543 })
3544 .await?;
3545 state.cancel_lifecycle(&session_id)?;
3546 Ok(DaemonReply::Done)
3547 }
3548 DaemonAction::RecoverDraft { draft_id } => {
3549 blocking(move || crate::database::recover_detached_draft(&draft_id)).await?;
3550 Ok(DaemonReply::Done)
3551 }
3552 DaemonAction::Stop => {
3553 cancellation.cancel();
3554 Ok(DaemonReply::Done)
3555 }
3556 }
3557}
3558
3559fn ensure_no_active_lifecycle(state: &RuntimeState) -> Result<()> {
3560 ensure!(
3561 !state
3562 .lifecycle
3563 .lock()
3564 .unwrap_or_else(PoisonError::into_inner)
3565 .values()
3566 .any(|active| active.result.borrow().is_none()),
3567 "cannot rename configuration while a session lifecycle operation is active"
3568 );
3569 Ok(())
3570}
3571
3572fn active_sessions_for_force_destruction(
3576 controller: &Controller,
3577 workspace_id: &str,
3578) -> Vec<String> {
3579 let mut sessions: Vec<&SessionRecord> = controller
3580 .state
3581 .sessions
3582 .values()
3583 .filter(|session| session.workspace_id == workspace_id && session.state.is_active())
3584 .collect();
3585 sessions.sort_by(|a, b| a.compare_by_creation(b));
3586 sessions
3587 .into_iter()
3588 .map(|session| session.id.clone())
3589 .collect()
3590}
3591
3592fn install_renamed_controller(state: &RuntimeState, controller: Controller) {
3593 *state
3594 .controller
3595 .lock()
3596 .unwrap_or_else(PoisonError::into_inner) = controller;
3597 state.publish_revision();
3598}
3599
3600fn workspace_snapshot(workspace_id: &str) -> Result<WorkspaceSnapshot> {
3601 let workspace = crate::database::list_workspaces()?
3602 .into_iter()
3603 .find(|workspace| workspace.id == workspace_id)
3604 .with_context(|| format!("unknown workspace {workspace_id:?}"))?;
3605 let ids = crate::database::session_ids_for_workspace(workspace_id)?
3606 .into_iter()
3607 .collect::<BTreeSet<_>>();
3608 let controller = Controller::load()?;
3609 let sessions = controller
3610 .state
3611 .sessions
3612 .values()
3613 .filter(|session| session.state.is_active() && ids.contains(&session.id))
3614 .map(|session| SessionPreview {
3615 id: session.id.clone(),
3616 title: session.display_title().to_owned(),
3617 project: session.project_name(&controller.config),
3618 harness: session.harness_kind.display_name().to_owned(),
3619 state: session_state_label(session.state).to_owned(),
3620 active: session.state.is_active(),
3621 updated_at: session.updated_at.clone(),
3622 })
3623 .collect();
3624 let drafts = crate::database::list_detached_drafts(workspace_id)?
3625 .into_iter()
3626 .map(|draft| DraftPreview {
3627 id: draft.id,
3628 session_id: draft.session_id,
3629 source: draft.source,
3630 owner_pid: draft.owner_pid,
3631 saved_at: draft.saved_at,
3632 })
3633 .collect();
3634 Ok(WorkspaceSnapshot {
3635 workspace,
3636 sessions,
3637 drafts,
3638 })
3639}
3640
3641fn session_state_label(state: SessionState) -> &'static str {
3642 match state {
3643 SessionState::Provisioning => "provisioning",
3644 SessionState::Running => "running",
3645 SessionState::Disconnected => "disconnected",
3646 SessionState::Checkpointing => "checkpointing",
3647 SessionState::Closing => "closing",
3648 SessionState::Destroying => "destroying",
3649 SessionState::Stopped => "stopped",
3650 SessionState::Lost => "lost",
3651 SessionState::Error => "error",
3652 SessionState::DestroyedWithDataLoss => "destroyed-with-data-loss",
3653 }
3654}
3655
3656#[cfg(test)]
3657mod tests {
3658 use super::*;
3659 use tokio::io::AsyncWriteExt;
3660
3661 #[test]
3662 fn newer_daemon_protocol_requires_updating_the_client() {
3663 assert!(ensure_supported_daemon_protocol(PROTOCOL_VERSION).is_ok());
3664 assert!(ensure_supported_daemon_protocol(PROTOCOL_VERSION - 1).is_ok());
3665 let error = ensure_supported_daemon_protocol(PROTOCOL_VERSION + 1).unwrap_err();
3666 assert!(error.to_string().contains("restart this client"));
3667 }
3668
3669 #[test]
3670 fn graceful_close_retires_worker_polling_only_during_target_teardown() {
3671 assert!(!lifecycle_owns_worker_target(
3672 LifecycleKind::Close,
3673 Some(SessionState::Running)
3674 ));
3675 assert!(!lifecycle_owns_worker_target(
3676 LifecycleKind::Close,
3677 Some(SessionState::Checkpointing)
3678 ));
3679 assert!(!lifecycle_owns_worker_target(
3680 LifecycleKind::Close,
3681 Some(SessionState::Closing)
3682 ));
3683 assert!(lifecycle_owns_worker_target(
3684 LifecycleKind::Close,
3685 Some(SessionState::Destroying)
3686 ));
3687 assert!(lifecycle_owns_worker_target(
3688 LifecycleKind::ForceStop,
3689 Some(SessionState::Running)
3690 ));
3691 assert!(lifecycle_owns_worker_target(
3692 LifecycleKind::ForceDestroy,
3693 Some(SessionState::Running)
3694 ));
3695 }
3696
3697 #[cfg(unix)]
3701 #[test]
3702 fn a_process_that_exited_but_was_not_reaped_counts_as_gone() {
3703 let mut child = std::process::Command::new("true")
3704 .spawn()
3705 .expect("spawn a process that exits immediately");
3706 let pid = child.id();
3707
3708 let deadline = std::time::Instant::now() + Duration::from_secs(5);
3709 let gone = loop {
3710 if !daemon_process_is_alive(pid) {
3711 break true;
3712 }
3713 if std::time::Instant::now() >= deadline {
3714 break false;
3715 }
3716 std::thread::sleep(Duration::from_millis(10));
3717 };
3718 assert!(
3719 gone,
3720 "an exited but unreaped process was reported as running"
3721 );
3722 let _ = child.wait();
3726 }
3727
3728 #[cfg(unix)]
3729 #[test]
3730 fn attachment_liveness_probe_does_not_reap_children() {
3731 use std::io::Read;
3732 use std::process::Stdio;
3733
3734 let mut child = std::process::Command::new("true")
3735 .stdout(Stdio::piped())
3736 .spawn()
3737 .expect("spawn a process that exits immediately");
3738 let pid = child.id();
3739 let mut output = Vec::new();
3740 child
3741 .stdout
3742 .take()
3743 .expect("capture child stdout")
3744 .read_to_end(&mut output)
3745 .expect("observe child exit");
3746
3747 let _ = process_is_alive(pid);
3748 let status = child.wait().expect("attachment probe left child waitable");
3749 assert!(status.success());
3750 }
3751
3752 #[tokio::test]
3753 async fn client_presence_is_global_and_detach_and_prune_remove_it() {
3754 let state = test_runtime_state();
3755 let metadata = test_metadata(SocketAddr::from((Ipv4Addr::LOCALHOST, 0)));
3756 let cancellation = CancellationToken::new();
3757
3758 handle_action(
3759 DaemonAction::Attach {
3760 client_id: "client-a".into(),
3761 pid: std::process::id(),
3762 },
3763 &metadata,
3764 &state,
3765 &cancellation,
3766 )
3767 .await
3768 .expect("attach presence");
3769 state
3770 .attachments()
3771 .insert("dead-client".into(), Attachment { pid: u32::MAX });
3772
3773 let DaemonReply::Status(status) =
3774 handle_action(DaemonAction::Status, &metadata, &state, &cancellation)
3775 .await
3776 .expect("status")
3777 else {
3778 panic!("status action returned a different reply");
3779 };
3780 assert_eq!(status.attached_clients, 1);
3781
3782 handle_action(
3783 DaemonAction::Detach {
3784 client_id: "client-a".into(),
3785 },
3786 &metadata,
3787 &state,
3788 &cancellation,
3789 )
3790 .await
3791 .expect("detach presence");
3792 assert!(state.attachments().is_empty());
3793 }
3794
3795 #[tokio::test]
3796 async fn workspace_deletion_guard_ignores_global_client_presence() {
3797 let state = test_runtime_state();
3798 state.attachments().insert(
3799 "client-a".into(),
3800 Attachment {
3801 pid: std::process::id(),
3802 },
3803 );
3804 assert!(!state.workspace_has_active_resume("workspace-a"));
3805
3806 let (_completed, result) = tokio::sync::watch::channel(None);
3807 state.lifecycle.lock().unwrap().insert(
3808 "session-a".into(),
3809 ActiveLifecycle {
3810 operation_id: "resume-operation".into(),
3811 create_control: None,
3812 kind: LifecycleKind::Resume,
3813 cancelled: Arc::new(AtomicBool::new(false)),
3814 started_at_epoch_seconds: 1,
3815 active_stages: BTreeMap::new(),
3816 resume_workspace_id: Some("workspace-a".into()),
3817 resume_destination: None,
3818 notice: None,
3819 request_key: None,
3820 _move_guard: None,
3821 move_source_closed: false,
3822 result,
3823 },
3824 );
3825 assert!(state.workspace_has_active_resume("workspace-a"));
3826 }
3827
3828 #[cfg(target_os = "macos")]
3829 #[test]
3830 fn zombie_only_daemon_group_counts_as_gone() {
3831 use std::io::Read;
3832 use std::os::unix::process::CommandExt;
3833 use std::process::Stdio;
3834
3835 let mut command = std::process::Command::new("true");
3836 command.process_group(0).stdout(Stdio::piped());
3837 let mut child = command
3838 .spawn()
3839 .expect("spawn process-group leader that exits immediately");
3840 let pid = libc::pid_t::try_from(child.id()).expect("child PID fits pid_t");
3841 let mut output = Vec::new();
3842 child
3843 .stdout
3844 .take()
3845 .expect("capture child stdout")
3846 .read_to_end(&mut output)
3847 .expect("observe child exit");
3848
3849 assert!(!owned_daemon_group_is_alive(pid));
3850 child.wait().expect("reap process-group leader");
3851 }
3852
3853 fn test_runtime_state() -> Arc<RuntimeState> {
3854 let remote = spawn_remote_session_manager().unwrap();
3855 let recovery = crate::recovery::RecoveryCoordinator::spawn(remote.control.clone());
3856 let upgrades = crate::worker_upgrade::WorkerUpgradeCoordinator::spawn(
3857 remote.control.clone(),
3858 &recovery.observer(),
3859 );
3860 Arc::new(RuntimeState::new_with_controller_loader(
3861 remote.control,
3862 Controller {
3863 config: Config::default(),
3864 state: mj_core::state::State::default(),
3865 },
3866 recovery.observer(),
3867 upgrades.observer(),
3868 Vec::new(),
3869 || {
3870 Ok(Controller {
3871 config: Config::default(),
3872 state: mj_core::state::State::default(),
3873 })
3874 },
3875 ))
3876 }
3877
3878 struct TestRemoteManager {
3879 control: SessionManagerControl,
3880 requests: RemoteSessionRequests,
3881 publisher: RemoteSessionPublisher,
3882 _shutdown: SessionManagerShutdown,
3883 _targets: tokio::sync::watch::Sender<Vec<RelaySessionTarget>>,
3884 }
3885
3886 impl TestRemoteManager {
3887 async fn new() -> Self {
3888 let channels = spawn_remote_session_manager().expect("remote manager");
3889 let session_id = "session-1";
3890 channels.targets.send_replace(vec![RelaySessionTarget {
3891 session_id: session_id.to_owned(),
3892 spec: CommandSpec::new("true", Vec::<String>::new()),
3893 worker_recovery: None,
3894 project_memory: None,
3895 }]);
3896 let manager = Self {
3897 control: channels.control,
3898 requests: channels.requests,
3899 publisher: channels.publisher,
3900 _shutdown: channels.shutdown,
3901 _targets: channels.targets,
3902 };
3903 manager
3904 .publisher
3905 .publish(session_id.to_owned(), ManagedSessionView::default())
3906 .await
3907 .expect("publish test session view");
3908 manager
3909 .control
3910 .wait_for_session(session_id, Duration::from_secs(5))
3911 .await
3912 .expect("remote manager creates test session");
3913 manager
3914 }
3915 }
3916
3917 fn test_runtime_state_with_manager(manager: &TestRemoteManager) -> Arc<RuntimeState> {
3918 let recovery = crate::recovery::RecoveryCoordinator::spawn(manager.control.clone());
3919 let upgrades = crate::worker_upgrade::WorkerUpgradeCoordinator::spawn(
3920 manager.control.clone(),
3921 &recovery.observer(),
3922 );
3923 Arc::new(RuntimeState::new_with_controller_loader(
3924 manager.control.clone(),
3925 Controller {
3926 config: Config::default(),
3927 state: mj_core::state::State::default(),
3928 },
3929 recovery.observer(),
3930 upgrades.observer(),
3931 Vec::new(),
3932 || {
3933 Ok(Controller {
3934 config: Config::default(),
3935 state: mj_core::state::State::default(),
3936 })
3937 },
3938 ))
3939 }
3940
3941 fn test_metadata(address: SocketAddr) -> DaemonMetadata {
3942 DaemonMetadata {
3943 protocol_version: PROTOCOL_VERSION,
3944 pid: 1,
3945 address,
3946 token: "right-token".into(),
3947 started_at: "now".into(),
3948 build_version: "test".into(),
3949 }
3950 }
3951
3952 #[tokio::test]
3953 async fn in_process_reviewer_forwarding_stops_when_the_caller_goes_away() {
3954 let mut manager = TestRemoteManager::new().await;
3955 let (reply, response) = tokio::sync::oneshot::channel();
3956 let forwarding = tokio::spawn(forward_in_process_session_request(
3957 RemoteSessionRequest::Reviewer {
3958 session_id: "session-1".into(),
3959 role: None,
3960 action: crate::session_manager::ReviewerAction::Status,
3961 reply,
3962 },
3963 manager.control.clone(),
3964 ));
3965 let request = tokio::time::timeout(Duration::from_secs(5), manager.requests.recv())
3966 .await
3967 .expect("reviewer forwarding did not reach the manager")
3968 .expect("manager request stream ended");
3969 let RemoteSessionRequest::Reviewer { mut reply, .. } = request else {
3970 panic!("expected the forwarded reviewer request")
3971 };
3972 drop(response);
3973 tokio::time::timeout(Duration::from_secs(5), reply.closed())
3974 .await
3975 .expect("in-process forwarding kept the actor reply alive");
3976 forwarding.await.expect("forwarding task panicked");
3977 }
3978
3979 #[tokio::test]
3980 async fn daemon_drops_in_flight_reviewer_work_when_the_client_eof_arrives() {
3981 let mut manager = TestRemoteManager::new().await;
3982 let state = test_runtime_state_with_manager(&manager);
3983 let listener = TcpListener::bind((Ipv4Addr::LOCALHOST, 0)).await.unwrap();
3984 let address = listener.local_addr().unwrap();
3985 let server = tokio::spawn(async move {
3986 let (stream, _) = listener.accept().await.unwrap();
3987 serve_client(
3988 stream,
3989 test_metadata(address),
3990 state,
3991 CancellationToken::new(),
3992 )
3993 .await
3994 });
3995 let mut stream = TcpStream::connect(address).await.unwrap();
3996 write_frame(
3997 &mut stream,
3998 &RequestEnvelope {
3999 protocol_version: PROTOCOL_VERSION,
4000 request_id: 1,
4001 token: "right-token".into(),
4002 action: DaemonAction::ReviewerAction {
4003 session_id: "session-1".into(),
4004 role: None,
4005 action: crate::session_manager::ReviewerAction::Status,
4006 },
4007 },
4008 )
4009 .await
4010 .unwrap();
4011 let request = tokio::time::timeout(Duration::from_secs(5), manager.requests.recv())
4012 .await
4013 .expect("reviewer request did not reach the manager")
4014 .expect("manager request stream ended");
4015 let RemoteSessionRequest::Reviewer { mut reply, .. } = request else {
4016 panic!("expected a reviewer request")
4017 };
4018 drop(stream);
4019 tokio::time::timeout(Duration::from_secs(5), reply.closed())
4020 .await
4021 .expect("daemon kept reviewer work alive after client EOF");
4022 assert!(server.await.expect("daemon task panicked").is_ok());
4023 }
4024
4025 #[tokio::test]
4026 async fn daemon_peek_keeps_a_pipelined_request_for_the_next_loop() {
4027 let mut manager = TestRemoteManager::new().await;
4028 let state = test_runtime_state_with_manager(&manager);
4029 let listener = TcpListener::bind((Ipv4Addr::LOCALHOST, 0)).await.unwrap();
4030 let address = listener.local_addr().unwrap();
4031 let server = tokio::spawn(async move {
4032 let (stream, _) = listener.accept().await.unwrap();
4033 serve_client(
4034 stream,
4035 test_metadata(address),
4036 state,
4037 CancellationToken::new(),
4038 )
4039 .await
4040 });
4041 let mut stream = TcpStream::connect(address).await.unwrap();
4042 for (request_id, action) in [
4043 (
4044 1,
4045 DaemonAction::ReviewerAction {
4046 session_id: "session-1".into(),
4047 role: None,
4048 action: crate::session_manager::ReviewerAction::Pause,
4049 },
4050 ),
4051 (2, DaemonAction::Ping),
4052 ] {
4053 write_frame(
4054 &mut stream,
4055 &RequestEnvelope {
4056 protocol_version: PROTOCOL_VERSION,
4057 request_id,
4058 token: "right-token".into(),
4059 action,
4060 },
4061 )
4062 .await
4063 .unwrap();
4064 }
4065 let request = tokio::time::timeout(Duration::from_secs(5), manager.requests.recv())
4066 .await
4067 .expect("reviewer request did not reach the manager")
4068 .expect("manager request stream ended");
4069 let RemoteSessionRequest::Reviewer { reply, .. } = request else {
4070 panic!("expected a reviewer request")
4071 };
4072 reply
4073 .send(Ok(crate::session_manager::ReviewerOutcome::Paused))
4074 .expect("daemon still awaits the reviewer reply");
4075
4076 let first: ResponseEnvelope = read_frame(&mut stream).await.unwrap();
4077 assert!(matches!(
4078 first.result,
4079 Ok(DaemonReply::Reviewer(outcome))
4080 if matches!(*outcome, crate::session_manager::ReviewerOutcome::Paused)
4081 ));
4082 let second: ResponseEnvelope = read_frame(&mut stream).await.unwrap();
4083 assert!(matches!(second.result, Ok(DaemonReply::Pong)));
4084 drop(stream);
4085 let _ = server.await.expect("daemon task panicked");
4086 }
4087
4088 #[tokio::test]
4089 async fn daemon_client_eof_does_not_cancel_a_submitted_mutation() {
4090 let mut manager = TestRemoteManager::new().await;
4091 let state = test_runtime_state_with_manager(&manager);
4092 let listener = TcpListener::bind((Ipv4Addr::LOCALHOST, 0)).await.unwrap();
4093 let address = listener.local_addr().unwrap();
4094 let server = tokio::spawn(async move {
4095 let (stream, _) = listener.accept().await.unwrap();
4096 serve_client(
4097 stream,
4098 test_metadata(address),
4099 state,
4100 CancellationToken::new(),
4101 )
4102 .await
4103 });
4104 let mut stream = TcpStream::connect(address).await.unwrap();
4105 write_frame(
4106 &mut stream,
4107 &RequestEnvelope {
4108 protocol_version: PROTOCOL_VERSION,
4109 request_id: 1,
4110 token: "right-token".into(),
4111 action: DaemonAction::SubmitSessionCommand {
4112 inherited_draft: None,
4113 session_id: "session-1".into(),
4114 command_id: "command-1".into(),
4115 command: RelayCommand::ClearQueuedPrompts,
4116 },
4117 },
4118 )
4119 .await
4120 .unwrap();
4121 let request = tokio::time::timeout(Duration::from_secs(5), manager.requests.recv())
4122 .await
4123 .expect("submit request did not reach the manager")
4124 .expect("manager request stream ended");
4125 let RemoteSessionRequest::Submit { reply, .. } = request else {
4126 panic!("expected a submitted mutation")
4127 };
4128 drop(stream);
4129 reply
4130 .send(Ok(7))
4131 .expect("daemon incorrectly cancelled a non-reviewer mutation");
4132 let _ = server.await.expect("daemon task panicked");
4133 }
4134
4135 fn runtime_test_session(id: &str, workspace_id: &str, state: SessionState) -> SessionRecord {
4136 SessionRecord {
4137 mjolnir_subagents: None,
4138 create_managed_worktree: None,
4139 id: id.into(),
4140 workspace_id: workspace_id.into(),
4141 title: id.into(),
4142 harness_kind: mj_core::config::HarnessKind::Codex,
4143 last_profile: "codex".into(),
4144 bundle_id: "project".into(),
4145 project_directory: None,
4146 managed_worktree: None,
4147 target_template_id: "local".into(),
4148 resource_allocation: None,
4149 additional_mounts: Vec::new(),
4150 container_cpus: None,
4151 container_memory: None,
4152 state,
4153 archived: false,
4154 target: None,
4155 native_session_id: None,
4156 acp_session_title: None,
4157 session_title_override: None,
4158 created_at: "2026-09-03T00:00:00Z".into(),
4159 updated_at: "2026-09-03T00:00:00Z".into(),
4160 viewed_through_event_ordinal: 0,
4161 draft_input: String::new(),
4162 last_error: None,
4163 last_checkpoint_error: None,
4164 checkpoint: None,
4165 }
4166 }
4167
4168 #[test]
4169 fn runtime_records_include_global_history_but_only_local_active_sessions() {
4170 let local = runtime_test_session("local", "workspace-a", SessionState::Running);
4171 let remote = runtime_test_session("remote", "workspace-b", SessionState::Running);
4172 let history = runtime_test_session("history", "deleted-workspace", SessionState::Stopped);
4173 let controller = Controller {
4174 config: Config::default(),
4175 state: mj_core::state::State {
4176 sessions: [local, remote, history]
4177 .into_iter()
4178 .map(|session| (session.id.clone(), session))
4179 .collect(),
4180 ..mj_core::state::State::default()
4181 },
4182 };
4183 let records =
4184 runtime_records_for_workspace(&controller, &BTreeSet::from(["local".to_owned()]));
4185 let ids = records
4186 .iter()
4187 .map(|session| session.id.as_str())
4188 .collect::<BTreeSet<_>>();
4189
4190 assert_eq!(ids, BTreeSet::from(["history", "local"]));
4191 }
4192
4193 #[tokio::test]
4194 async fn daemon_records_definitive_missing_workspaces_without_an_attached_surface() {
4195 let state = test_runtime_state();
4196 let session = runtime_test_session("missing", "workspace", SessionState::Running);
4197 state
4198 .controller
4199 .lock()
4200 .unwrap()
4201 .state
4202 .sessions
4203 .insert(session.id.clone(), session.clone());
4204 assert!(state.attachments.lock().unwrap().is_empty());
4205 let mut view = ManagedSessionView {
4206 snapshot: None,
4207 connected: false,
4208 error: Some(ViewError::Unreachable(
4209 "relay proxy disconnected during hello".into(),
4210 )),
4211 };
4212 assert!(state.missing_target_record(&session.id, &view).is_none());
4213 view.error = Some(ViewError::TargetMissing(
4214 "working directory /missing is gone".into(),
4215 ));
4216 assert_eq!(
4217 state.missing_target_record(&session.id, &view),
4218 Some((
4219 "working directory /missing is gone".into(),
4220 session.updated_at
4221 )),
4222 );
4223 state
4224 .controller
4225 .lock()
4226 .unwrap()
4227 .state
4228 .sessions
4229 .get_mut(&session.id)
4230 .unwrap()
4231 .state = SessionState::Closing;
4232 assert!(state.missing_target_record(&session.id, &view).is_none());
4233 state
4234 .controller
4235 .lock()
4236 .unwrap()
4237 .state
4238 .sessions
4239 .get_mut(&session.id)
4240 .unwrap()
4241 .state = SessionState::Error;
4242 assert!(state.missing_target_record(&session.id, &view).is_none());
4243 }
4244
4245 #[tokio::test]
4246 async fn review_host_notifier_wakes_runtime_revision_subscribers() {
4247 let revisions = RuntimeRevisions::new(40);
4248 let mut subscriber = revisions.subscribe();
4249 let notify_review_publication = revisions.notifier();
4253
4254 notify_review_publication();
4255 tokio::time::timeout(Duration::from_secs(1), subscriber.changed())
4256 .await
4257 .expect("review publication did not wake runtime subscribers")
4258 .expect("runtime revision publisher stopped");
4259
4260 assert_eq!(*subscriber.borrow_and_update(), 41);
4261 }
4262
4263 #[test]
4264 fn late_runtime_revision_publication_cannot_move_cursor_backwards() {
4265 let revisions = RuntimeRevisions::new(40);
4266 let subscriber = revisions.subscribe();
4267
4268 revisions.publish_allocated(42);
4269 revisions.publish_allocated(41);
4270
4271 assert_eq!(*subscriber.borrow(), 42);
4272 }
4273
4274 #[tokio::test]
4275 async fn workspace_publication_reaches_existing_phone_subscriber() {
4276 let state = test_runtime_state();
4277 let mut workspaces = state.workspaces();
4278 let expected = WorkspaceRecord {
4279 id: "workspace-1".into(),
4280 name: "Reliability".into(),
4281 created_at: "2026-08-30T00:00:00Z".into(),
4282 last_opened_at: "2026-08-30T00:00:00Z".into(),
4283 session_count: 0,
4284 };
4285
4286 state.publish_workspaces(vec![expected.clone()]);
4287 tokio::time::timeout(Duration::from_secs(1), workspaces.changed())
4288 .await
4289 .expect("workspace publication timed out")
4290 .expect("workspace publisher stopped");
4291
4292 assert_eq!(workspaces.borrow_and_update().as_slice(), &[expected]);
4293 }
4294
4295 #[tokio::test]
4296 async fn framing_round_trips_payloads_larger_than_a_pipe_buffer() {
4297 let listener = TcpListener::bind((Ipv4Addr::LOCALHOST, 0)).await.unwrap();
4298 let address = listener.local_addr().unwrap();
4299 let sender = tokio::spawn(async move {
4300 let mut stream = TcpStream::connect(address).await.unwrap();
4301 write_frame(&mut stream, &"x".repeat(512 * 1024))
4302 .await
4303 .unwrap();
4304 });
4305 let (mut stream, _) = listener.accept().await.unwrap();
4306 let received: String = read_frame(&mut stream).await.unwrap();
4307 sender.await.unwrap();
4308 assert_eq!(received.len(), 512 * 1024);
4309 }
4310
4311 #[tokio::test]
4312 async fn framing_rejects_an_oversized_frame_before_allocating_it() {
4313 let listener = TcpListener::bind((Ipv4Addr::LOCALHOST, 0)).await.unwrap();
4314 let address = listener.local_addr().unwrap();
4315 let sender = tokio::spawn(async move {
4316 let mut stream = TcpStream::connect(address).await.unwrap();
4317 stream
4318 .write_u32((MAX_FRAME_BYTES + 1) as u32)
4319 .await
4320 .unwrap();
4321 });
4322 let (mut stream, _) = listener.accept().await.unwrap();
4323 assert!(read_frame::<String>(&mut stream).await.is_err());
4324 sender.await.unwrap();
4325 }
4326
4327 #[tokio::test]
4331 async fn daemon_stops_serving_data_actions_once_shutdown_begins() {
4332 let state = test_runtime_state();
4333 let listener = TcpListener::bind((Ipv4Addr::LOCALHOST, 0)).await.unwrap();
4334 let address = listener.local_addr().unwrap();
4335 let metadata = DaemonMetadata {
4336 protocol_version: PROTOCOL_VERSION,
4337 pid: 1,
4338 address,
4339 token: "right-token".into(),
4340 started_at: "now".into(),
4341 build_version: "test".into(),
4342 };
4343 let cancellation = CancellationToken::new();
4344 cancellation.cancel();
4345 let server_metadata = metadata.clone();
4346 let server_cancellation = cancellation.clone();
4347 let server = tokio::spawn(async move {
4348 let (stream, _) = listener.accept().await.unwrap();
4349 serve_client(stream, server_metadata, state, server_cancellation)
4350 .await
4351 .unwrap();
4352 });
4353 let mut stream = TcpStream::connect(address).await.unwrap();
4354
4355 let mut request_id = 0;
4356 let mut ask = async |stream: &mut TcpStream, action: DaemonAction| {
4357 request_id += 1;
4358 write_frame(
4359 stream,
4360 &RequestEnvelope {
4361 protocol_version: PROTOCOL_VERSION,
4362 request_id,
4363 token: "right-token".to_owned(),
4364 action,
4365 },
4366 )
4367 .await
4368 .unwrap();
4369 read_frame::<ResponseEnvelope>(stream).await.unwrap().result
4370 };
4371
4372 let refused = ask(
4373 &mut stream,
4374 DaemonAction::Snapshot {
4375 workspace_id: "workspace-a".into(),
4376 },
4377 )
4378 .await;
4379 assert_eq!(
4380 refused.unwrap_err(),
4381 "daemon is shutting down; retry to reach a fresh daemon"
4382 );
4383 assert!(matches!(
4384 ask(&mut stream, DaemonAction::Ping).await,
4385 Ok(DaemonReply::Pong)
4386 ));
4387 assert!(matches!(
4388 ask(&mut stream, DaemonAction::Status).await,
4389 Ok(DaemonReply::Status(_))
4390 ));
4391 assert!(matches!(
4392 ask(&mut stream, DaemonAction::Stop).await,
4393 Ok(DaemonReply::Done)
4394 ));
4395
4396 drop(stream);
4397 server.await.unwrap();
4398 }
4399
4400 #[test]
4403 fn shutdown_force_exit_finishes_before_the_stop_deadline() {
4404 assert!(SHUTDOWN_FORCE_EXIT_TIMEOUT < STOP_TIMEOUT);
4405 }
4406
4407 #[tokio::test]
4408 async fn management_stop_is_bounded_when_the_daemon_never_acknowledges() {
4409 let listener = TcpListener::bind((Ipv4Addr::LOCALHOST, 0)).await.unwrap();
4410 let metadata = DaemonMetadata {
4411 protocol_version: PROTOCOL_VERSION,
4412 pid: 1,
4413 address: listener.local_addr().unwrap(),
4414 token: "test-token".into(),
4415 started_at: "now".into(),
4416 build_version: "test".into(),
4417 };
4418 let mut client = ManagementClient::new(DaemonClient::connect(metadata).await.unwrap());
4419 let (mut peer, _) = listener.accept().await.unwrap();
4420 let stop = tokio::spawn(async move { client.stop().await });
4421 let request: RequestEnvelope = read_frame(&mut peer).await.unwrap();
4422 assert!(matches!(request.action, DaemonAction::Stop));
4423 tokio::time::pause();
4425 tokio::time::advance(STOP_TIMEOUT).await;
4426 let error = stop.await.unwrap().unwrap_err();
4427 assert!(error.to_string().contains("did not acknowledge"));
4428 }
4429
4430 #[tokio::test]
4431 async fn daemon_rejects_a_request_with_the_wrong_owner_token() {
4432 let state = test_runtime_state();
4433 let listener = TcpListener::bind((Ipv4Addr::LOCALHOST, 0)).await.unwrap();
4434 let address = listener.local_addr().unwrap();
4435 let metadata = DaemonMetadata {
4436 protocol_version: PROTOCOL_VERSION,
4437 pid: 1,
4438 address,
4439 token: "right-token".into(),
4440 started_at: "now".into(),
4441 build_version: "test".into(),
4442 };
4443 let server_metadata = metadata.clone();
4444 let server = tokio::spawn(async move {
4445 let (stream, _) = listener.accept().await.unwrap();
4446 serve_client(stream, server_metadata, state, CancellationToken::new())
4447 .await
4448 .unwrap();
4449 });
4450 let mut stream = TcpStream::connect(address).await.unwrap();
4451 write_frame(
4452 &mut stream,
4453 &RequestEnvelope {
4454 protocol_version: PROTOCOL_VERSION,
4455 request_id: 42,
4456 token: "wrong-token".into(),
4457 action: DaemonAction::Ping,
4458 },
4459 )
4460 .await
4461 .unwrap();
4462 let response: ResponseEnvelope = read_frame(&mut stream).await.unwrap();
4463 assert_eq!(response.request_id, 42);
4464 assert_eq!(response.result.unwrap_err(), "daemon authentication failed");
4465 drop(stream);
4466 server.await.unwrap();
4467 }
4468
4469 #[tokio::test]
4470 async fn the_daemon_rejects_a_client_one_protocol_behind_before_dispatch() {
4471 let state = test_runtime_state();
4472 let listener = TcpListener::bind((Ipv4Addr::LOCALHOST, 0)).await.unwrap();
4473 let address = listener.local_addr().unwrap();
4474 let metadata = DaemonMetadata {
4475 protocol_version: PROTOCOL_VERSION,
4476 pid: 1,
4477 address,
4478 token: "right-token".into(),
4479 started_at: "now".into(),
4480 build_version: "test".into(),
4481 };
4482 let server_metadata = metadata.clone();
4483 let server = tokio::spawn(async move {
4484 let (stream, _) = listener.accept().await.unwrap();
4485 serve_client(stream, server_metadata, state, CancellationToken::new())
4486 .await
4487 .unwrap();
4488 });
4489 let mut stream = TcpStream::connect(address).await.unwrap();
4490 write_frame(
4491 &mut stream,
4492 &RequestEnvelope {
4493 protocol_version: PROTOCOL_VERSION - 1,
4494 request_id: 43,
4495 token: metadata.token,
4496 action: DaemonAction::PersistReadReceipt {
4497 client_id: "client-a".into(),
4498 workspace_id: "workspace-a".into(),
4499 session_id: "session-a".into(),
4500 through: 7,
4501 },
4502 },
4503 )
4504 .await
4505 .unwrap();
4506 let response: ResponseEnvelope = read_frame(&mut stream).await.unwrap();
4507 assert_eq!(response.request_id, 43);
4508 assert!(response.result.unwrap_err().contains(&format!(
4509 "incompatible daemon protocol {}; expected {PROTOCOL_VERSION}",
4510 PROTOCOL_VERSION - 1
4511 )));
4512 drop(stream);
4513 server.await.unwrap();
4514 }
4515
4516 #[test]
4520 fn management_wire_shapes_stay_frozen_across_protocol_versions() {
4521 for (action, expected) in [
4522 (DaemonAction::Ping, serde_json::json!({"action": "ping"})),
4523 (
4524 DaemonAction::Status,
4525 serde_json::json!({"action": "status"}),
4526 ),
4527 (DaemonAction::Stop, serde_json::json!({"action": "stop"})),
4528 ] {
4529 let request = RequestEnvelope {
4530 protocol_version: 3,
4531 request_id: 7,
4532 token: "tok".into(),
4533 action,
4534 };
4535 assert_eq!(
4536 serde_json::to_value(&request).unwrap(),
4537 serde_json::json!({
4538 "protocol_version": 3,
4539 "request_id": 7,
4540 "token": "tok",
4541 "action": expected,
4542 })
4543 );
4544 }
4545
4546 let response: ResponseEnvelope = serde_json::from_value(serde_json::json!({
4547 "protocol_version": 3,
4548 "request_id": 7,
4549 "result": {"Ok": {"reply": "status", "value": {
4550 "pid": 4242,
4551 "started_at": "2026-09-01T07:48:14Z",
4552 "build_version": "0.3.1",
4553 "attached_clients": 1,
4554 "phone_status": {"state": "disabled"},
4555 }}}
4556 }))
4557 .unwrap();
4558 match response.result.unwrap() {
4559 DaemonReply::Status(status) => {
4560 assert_eq!(status.pid, 4242);
4561 assert_eq!(status.build_version, "0.3.1");
4562 }
4563 reply => panic!("unexpected reply {reply:?}"),
4564 }
4565 }
4566
4567 #[test]
4568 fn client_presence_and_workspace_listing_use_global_wire_shapes() {
4569 let attach = serde_json::to_value(DaemonAction::Attach {
4570 client_id: "client-a".into(),
4571 pid: 4242,
4572 })
4573 .unwrap();
4574 assert_eq!(
4575 attach,
4576 serde_json::json!({
4577 "action": "attach",
4578 "arguments": {"client_id": "client-a", "pid": 4242},
4579 })
4580 );
4581
4582 let listing = serde_json::to_value(WorkspaceListing {
4583 workspace: WorkspaceRecord {
4584 id: "workspace-a".into(),
4585 name: "Workspace A".into(),
4586 created_at: "2026-09-01T00:00:00Z".into(),
4587 last_opened_at: "2026-09-01T00:00:00Z".into(),
4588 session_count: 0,
4589 },
4590 })
4591 .unwrap();
4592 assert!(listing.get("attached_pids").is_none());
4593 }
4594
4595 #[tokio::test]
4596 async fn daemon_serves_management_actions_for_any_protocol_version() {
4597 for version in [3_u32, 5] {
4600 let state = test_runtime_state();
4601 let listener = TcpListener::bind((Ipv4Addr::LOCALHOST, 0)).await.unwrap();
4602 let address = listener.local_addr().unwrap();
4603 let metadata = DaemonMetadata {
4604 protocol_version: PROTOCOL_VERSION,
4605 pid: 1,
4606 address,
4607 token: "right-token".into(),
4608 started_at: "now".into(),
4609 build_version: "test".into(),
4610 };
4611 let server_metadata = metadata.clone();
4612 let server = tokio::spawn(async move {
4613 let (stream, _) = listener.accept().await.unwrap();
4614 serve_client(stream, server_metadata, state, CancellationToken::new())
4615 .await
4616 .unwrap();
4617 });
4618 let mut stream = TcpStream::connect(address).await.unwrap();
4619 write_frame(
4620 &mut stream,
4621 &RequestEnvelope {
4622 protocol_version: version,
4623 request_id: 44,
4624 token: "right-token".into(),
4625 action: DaemonAction::Status,
4626 },
4627 )
4628 .await
4629 .unwrap();
4630 let response: ResponseEnvelope = read_frame(&mut stream).await.unwrap();
4631 assert_eq!(
4632 response.protocol_version, version,
4633 "reply must use the caller's dialect"
4634 );
4635 assert_eq!(response.request_id, 44);
4636 match response.result.unwrap() {
4637 DaemonReply::Status(status) => assert_eq!(status.build_version, "test"),
4638 reply => panic!("unexpected reply {reply:?}"),
4639 }
4640 drop(stream);
4641 server.await.unwrap();
4642 }
4643 }
4644
4645 struct ProtocolTranscript {
4646 protocol_version: u32,
4647 daemon_build: &'static str,
4648 expected_requests: [serde_json::Value; 2],
4650 responses: [&'static str; 2],
4653 }
4654
4655 fn released_protocol_transcripts() -> Vec<ProtocolTranscript> {
4662 let requests = |version: u32| {
4663 [
4664 serde_json::json!({
4665 "protocol_version": version,
4666 "request_id": 1,
4667 "token": "tok",
4668 "action": {"action": "status"},
4669 }),
4670 serde_json::json!({
4671 "protocol_version": version,
4672 "request_id": 2,
4673 "token": "tok",
4674 "action": {"action": "stop"},
4675 }),
4676 ]
4677 };
4678 vec![
4679 ProtocolTranscript {
4680 protocol_version: 3,
4681 daemon_build: "0.3.1",
4682 expected_requests: requests(3),
4683 responses: [
4684 r#"{"protocol_version":3,"request_id":1,"result":{"Ok":{"reply":"status","value":{"pid":4242,"started_at":"2026-09-01T07:48:14Z","build_version":"0.3.1","attached_clients":1,"phone_status":{"state":"ready","viewer_url":"https://example.test:1","viewer_code":"690451","qr_login_url":null,"fallback_reason":null}}}}}"#,
4685 r#"{"protocol_version":3,"request_id":2,"result":{"Ok":{"reply":"done"}}}"#,
4686 ],
4687 },
4688 ProtocolTranscript {
4689 protocol_version: 4,
4690 daemon_build: "0.4.1",
4691 expected_requests: requests(4),
4692 responses: [
4693 r#"{"protocol_version":4,"request_id":1,"result":{"Ok":{"reply":"status","value":{"pid":4242,"started_at":"2026-09-01T07:48:14Z","build_version":"0.4.1","attached_clients":1,"phone_status":{"state":"disabled"}}}}}"#,
4694 r#"{"protocol_version":4,"request_id":2,"result":{"Ok":{"reply":"done"}}}"#,
4695 ],
4696 },
4697 ProtocolTranscript {
4698 protocol_version: 11,
4699 daemon_build: "2.1.0",
4700 expected_requests: requests(11),
4701 responses: [
4702 r#"{"protocol_version":11,"request_id":1,"result":{"Ok":{"reply":"status","value":{"pid":4242,"started_at":"2026-09-01T07:48:14Z","build_version":"2.1.0","attached_clients":1,"phone_status":{"state":"disabled"}}}}}"#,
4703 r#"{"protocol_version":11,"request_id":2,"result":{"Ok":{"reply":"done"}}}"#,
4704 ],
4705 },
4706 ProtocolTranscript {
4707 protocol_version: 14,
4708 daemon_build: "2.1.4",
4709 expected_requests: requests(14),
4710 responses: [
4711 r#"{"protocol_version":14,"request_id":1,"result":{"Ok":{"reply":"status","value":{"pid":4242,"started_at":"2026-09-01T07:48:14Z","build_version":"2.1.4","attached_clients":1,"phone_status":{"state":"disabled"}}}}}"#,
4712 r#"{"protocol_version":14,"request_id":2,"result":{"Ok":{"reply":"done"}}}"#,
4713 ],
4714 },
4715 ProtocolTranscript {
4716 protocol_version: 15,
4717 daemon_build: "2.2.0",
4718 expected_requests: requests(15),
4719 responses: [
4720 r#"{"protocol_version":15,"request_id":1,"result":{"Ok":{"reply":"status","value":{"pid":4242,"started_at":"2026-09-01T07:48:14Z","build_version":"2.2.0","attached_clients":1,"phone_status":{"state":"disabled"}}}}}"#,
4721 r#"{"protocol_version":15,"request_id":2,"result":{"Ok":{"reply":"done"}}}"#,
4722 ],
4723 },
4724 ProtocolTranscript {
4725 protocol_version: 16,
4726 daemon_build: "2.4.0",
4727 expected_requests: requests(16),
4728 responses: [
4729 r#"{"protocol_version":16,"request_id":1,"result":{"Ok":{"reply":"status","value":{"pid":4242,"started_at":"2026-09-01T07:48:14Z","build_version":"2.4.0","attached_clients":1,"phone_status":{"state":"disabled"}}}}}"#,
4730 r#"{"protocol_version":16,"request_id":2,"result":{"Ok":{"reply":"done"}}}"#,
4731 ],
4732 },
4733 ]
4734 }
4735
4736 #[tokio::test]
4737 async fn management_client_talks_to_every_released_protocol_version() {
4738 for transcript in released_protocol_transcripts() {
4739 let protocol_version = transcript.protocol_version;
4740 let daemon_build = transcript.daemon_build;
4741 let listener = TcpListener::bind((Ipv4Addr::LOCALHOST, 0)).await.unwrap();
4742 let address = listener.local_addr().unwrap();
4743 let server = tokio::spawn(async move {
4744 let (mut stream, _) = listener.accept().await.unwrap();
4745 for (expected, response) in transcript
4746 .expected_requests
4747 .iter()
4748 .zip(transcript.responses)
4749 {
4750 let request: serde_json::Value = read_frame(&mut stream).await.unwrap();
4751 assert_eq!(
4752 &request, expected,
4753 "protocol {} daemon would reject this frame",
4754 transcript.protocol_version
4755 );
4756 stream.write_u32(response.len() as u32).await.unwrap();
4757 stream.write_all(response.as_bytes()).await.unwrap();
4758 stream.flush().await.unwrap();
4759 }
4760 });
4761 let metadata = DaemonMetadata {
4762 protocol_version,
4763 pid: 4242,
4764 address,
4765 token: "tok".into(),
4766 started_at: "2026-09-01T07:48:14Z".into(),
4767 build_version: daemon_build.into(),
4768 };
4769 let mut client = ManagementClient::new(DaemonClient::connect(metadata).await.unwrap());
4770 let status = client.status().await.unwrap();
4771 assert_eq!(status.build_version, daemon_build);
4772 assert_eq!(status.attached_clients, 1);
4773 assert_eq!(client.protocol_version(), protocol_version);
4774 client.stop().await.unwrap();
4775 server.await.unwrap();
4776 }
4777 }
4778
4779 #[tokio::test]
4780 async fn equivalent_lifecycle_requests_join_one_daemon_operation() {
4781 let state = test_runtime_state();
4782 let starts = Arc::new(std::sync::atomic::AtomicUsize::new(0));
4783 let release = Arc::new(tokio::sync::Notify::new());
4784 let first = state
4785 .start_or_join_lifecycle("session-1".into(), LifecycleKind::Close, {
4786 let starts = starts.clone();
4787 let release = release.clone();
4788 move |_state, _session_id, _cancelled| async move {
4789 starts.fetch_add(1, Ordering::AcqRel);
4790 release.notified().await;
4791 Ok(DaemonLifecycleResult::Done)
4792 }
4793 })
4794 .unwrap();
4795 tokio::task::yield_now().await;
4796 let second = state
4797 .start_or_join_lifecycle(
4798 "session-1".into(),
4799 LifecycleKind::Close,
4800 |_state, _session_id, _cancelled| async move {
4801 panic!("joined lifecycle request started duplicate work")
4802 },
4803 )
4804 .unwrap();
4805 assert_eq!(starts.load(Ordering::Acquire), 1);
4806 assert!(
4807 state
4808 .start_or_join_lifecycle(
4809 "session-1".into(),
4810 LifecycleKind::Resume,
4811 |_state, _session_id, _cancelled| async move {
4812 Ok(DaemonLifecycleResult::Done)
4813 },
4814 )
4815 .is_err()
4816 );
4817
4818 drop(first);
4820 release.notify_one();
4821 assert!(matches!(
4822 RuntimeState::wait_lifecycle_result(second).await.unwrap(),
4823 DaemonLifecycleResult::Done
4824 ));
4825 assert_eq!(starts.load(Ordering::Acquire), 1);
4826 }
4827
4828 #[tokio::test]
4829 async fn close_keeps_worker_target_available_for_checkpoint_lease() {
4830 let state = test_runtime_state();
4831 let release = Arc::new(tokio::sync::Notify::new());
4832 let result = state
4833 .start_or_join_lifecycle("session-1".into(), LifecycleKind::Close, {
4834 let release = release.clone();
4835 move |_state, _session_id, _cancelled| async move {
4836 release.notified().await;
4837 Ok(DaemonLifecycleResult::Done)
4838 }
4839 })
4840 .unwrap();
4841
4842 assert!(
4843 state
4844 .worker_poll_exclusion_session_ids(
4845 &state
4846 .controller
4847 .lock()
4848 .unwrap_or_else(PoisonError::into_inner)
4849 )
4850 .is_empty()
4851 );
4852
4853 release.notify_one();
4854 RuntimeState::wait_lifecycle_result(result).await.unwrap();
4855 }
4856
4857 #[tokio::test]
4858 async fn move_joins_only_matching_selections_and_runs_other_sessions_concurrently() {
4859 let state = test_runtime_state();
4860 let release = Arc::new(tokio::sync::Notify::new());
4861 let first = state
4862 .start_or_join_lifecycle_with_key(
4863 "move-join-one".into(),
4864 LifecycleKind::Move,
4865 None,
4866 Some("profile-a/target-a/discard".into()),
4867 {
4868 let release = release.clone();
4869 move |_, _, _| async move {
4870 release.notified().await;
4871 Ok(DaemonLifecycleResult::Done)
4872 }
4873 },
4874 )
4875 .unwrap();
4876 assert!(crate::controller::move_session::move_owns_session(
4877 "move-join-one"
4878 ));
4879 let duplicate = state
4880 .start_or_join_lifecycle_with_key(
4881 "move-join-one".into(),
4882 LifecycleKind::Move,
4883 None,
4884 Some("profile-a/target-a/discard".into()),
4885 |_, _, _| async move { panic!("duplicate move launched a second writer") },
4886 )
4887 .unwrap();
4888 assert!(
4889 state
4890 .start_or_join_lifecycle_with_key(
4891 "move-join-one".into(),
4892 LifecycleKind::Move,
4893 None,
4894 Some("profile-b/target-a/discard".into()),
4895 |_, _, _| async move { Ok(DaemonLifecycleResult::Done) },
4896 )
4897 .is_err()
4898 );
4899 let unrelated = state
4900 .start_or_join_lifecycle_with_key(
4901 "move-join-two".into(),
4902 LifecycleKind::Move,
4903 None,
4904 Some("profile-b/target-b/start".into()),
4905 |_, _, _| async move { Ok(DaemonLifecycleResult::Done) },
4906 )
4907 .unwrap();
4908 let unrelated_channel = unrelated.clone();
4909 tokio::time::timeout(
4910 Duration::from_secs(2),
4911 RuntimeState::wait_lifecycle_result(unrelated),
4912 )
4913 .await
4914 .unwrap()
4915 .unwrap();
4916 state.remove_completed_lifecycle(&unrelated_channel);
4917 assert!(!crate::controller::move_session::move_owns_session(
4918 "move-join-two"
4919 ));
4920 drop(first); release.notify_one();
4922 let channel = duplicate.clone();
4923 RuntimeState::wait_lifecycle_result(duplicate)
4924 .await
4925 .unwrap();
4926 state.remove_completed_lifecycle(&channel);
4927 assert!(!crate::controller::move_session::move_owns_session(
4928 "move-join-one"
4929 ));
4930 }
4931
4932 #[tokio::test]
4933 async fn move_task_panic_reports_failure_and_releases_mutation_hold() {
4934 let state = test_runtime_state();
4935 let result = state
4936 .start_or_join_lifecycle_with_key(
4937 "move-panics".into(),
4938 LifecycleKind::Move,
4939 None,
4940 Some("destination".into()),
4941 |_, _, _| async move { panic!("injected move task panic") },
4942 )
4943 .unwrap();
4944 let channel = result.clone();
4945 let error = tokio::time::timeout(
4946 Duration::from_secs(2),
4947 RuntimeState::wait_lifecycle_result(result),
4948 )
4949 .await
4950 .unwrap()
4951 .unwrap_err();
4952 assert!(error.to_string().contains("daemon lifecycle task failed"));
4953 state.remove_completed_lifecycle(&channel);
4954 assert!(!crate::controller::move_session::move_owns_session(
4955 "move-panics"
4956 ));
4957 }
4958
4959 #[tokio::test]
4960 async fn daemon_lifecycle_reports_balanced_concurrent_stages() {
4961 struct UnusedExecutor;
4962
4963 impl CommandExecutor for UnusedExecutor {
4964 fn execute(&self, command: &CommandSpec) -> Result<CommandOutput> {
4965 panic!("a stage notification must not run {}", command.program)
4966 }
4967 }
4968
4969 let state = test_runtime_state();
4970 let release = Arc::new(tokio::sync::Notify::new());
4971 let result = state
4972 .start_or_join_lifecycle("session-1".into(), LifecycleKind::Create, {
4973 let release = release.clone();
4974 move |_state, _session_id, _cancelled| async move {
4975 release.notified().await;
4976 Ok(DaemonLifecycleResult::Done)
4977 }
4978 })
4979 .unwrap();
4980 assert_eq!(
4981 state.worker_poll_exclusion_session_ids(
4982 &state
4983 .controller
4984 .lock()
4985 .unwrap_or_else(PoisonError::into_inner)
4986 ),
4987 BTreeSet::from(["session-1".to_owned()])
4988 );
4989 let executor =
4990 DaemonStageReportingExecutor::new(UnusedExecutor, state.clone(), "session-1".into());
4991 executor.stage_started(ProvisionStage::Cloning);
4992 executor.stage_started(ProvisionStage::Cloning);
4993 executor.stage_started(ProvisionStage::Syncing);
4994 executor.stage_finished(ProvisionStage::Cloning);
4995 {
4996 let lifecycle = state
4997 .lifecycle
4998 .lock()
4999 .unwrap_or_else(PoisonError::into_inner);
5000 let stages = &lifecycle.get("session-1").unwrap().active_stages;
5001 assert_eq!(stages.get(&ProvisionStage::Cloning).unwrap().0, 1);
5002 assert_eq!(stages.get(&ProvisionStage::Syncing).unwrap().0, 1);
5003 }
5004 executor.stage_finished(ProvisionStage::Cloning);
5005 executor.stage_finished(ProvisionStage::Syncing);
5006 assert!(
5007 state
5008 .lifecycle
5009 .lock()
5010 .unwrap_or_else(PoisonError::into_inner)
5011 .get("session-1")
5012 .unwrap()
5013 .active_stages
5014 .is_empty()
5015 );
5016 release.notify_one();
5017 assert!(matches!(
5018 RuntimeState::wait_lifecycle_result(result).await.unwrap(),
5019 DaemonLifecycleResult::Done
5020 ));
5021 }
5022
5023 #[tokio::test]
5024 async fn deferred_cleanup_is_visible_and_drains_before_shutdown_cancellation() {
5025 let state = test_runtime_state();
5026 let saw_early_cancellation = Arc::new(AtomicBool::new(false));
5027 let result = state
5028 .start_or_join_lifecycle("session-1".into(), LifecycleKind::Cleanup, {
5029 let saw_early_cancellation = saw_early_cancellation.clone();
5030 move |state, session_id, cancelled| async move {
5031 let executor = DaemonStageReportingExecutor::new(
5032 crate::targets::ProcessExecutor,
5033 state,
5034 session_id,
5035 );
5036 executor.stage_started(ProvisionStage::RemovingStorage);
5037 tokio::time::sleep(Duration::from_millis(20)).await;
5038 saw_early_cancellation
5039 .store(cancelled.load(Ordering::Acquire), Ordering::Release);
5040 executor.stage_finished(ProvisionStage::RemovingStorage);
5041 Ok(DaemonLifecycleResult::Done)
5042 }
5043 })
5044 .unwrap();
5045 tokio::task::yield_now().await;
5046
5047 let visible = state.active_lifecycles();
5048 assert_eq!(visible.len(), 1);
5049 assert_eq!(visible[0].kind, RuntimeLifecycleKind::Cleanup);
5050 assert_eq!(
5051 visible[0].active_stages[0].0,
5052 ProvisionStage::RemovingStorage
5053 );
5054
5055 state.cancel_and_wait_lifecycles().await.unwrap();
5056 assert!(!saw_early_cancellation.load(Ordering::Acquire));
5057 assert!(matches!(
5058 RuntimeState::wait_lifecycle_result(result).await.unwrap(),
5059 DaemonLifecycleResult::Done
5060 ));
5061 }
5062
5063 #[test]
5064 fn force_destruction_enumerates_only_the_workspaces_active_sessions_oldest_first() {
5065 let mut oldest = runtime_test_session("oldest", "workspace-a", SessionState::Provisioning);
5066 oldest.created_at = "2026-09-01T00:00:00Z".into();
5067 let newest = runtime_test_session("newest", "workspace-a", SessionState::Error);
5068 let elsewhere = runtime_test_session("elsewhere", "workspace-b", SessionState::Running);
5069 let history = runtime_test_session("history", "workspace-a", SessionState::Stopped);
5070 let controller = Controller {
5071 config: Config::default(),
5072 state: mj_core::state::State {
5073 sessions: [oldest, newest, elsewhere, history]
5074 .into_iter()
5075 .map(|session| (session.id.clone(), session))
5076 .collect(),
5077 ..mj_core::state::State::default()
5078 },
5079 };
5080
5081 assert_eq!(
5082 active_sessions_for_force_destruction(&controller, "workspace-a"),
5083 vec!["oldest".to_owned(), "newest".to_owned()]
5084 );
5085 assert_eq!(
5086 active_sessions_for_force_destruction(&controller, "workspace-b"),
5087 vec!["elsewhere".to_owned()]
5088 );
5089 }
5090
5091 #[test]
5092 fn force_destroy_serializes_as_its_own_lifecycle_kind() {
5093 assert_eq!(
5094 serde_json::to_string(&RuntimeLifecycleKind::ForceDestroy).unwrap(),
5095 "\"force_destroy\""
5096 );
5097 }
5098
5099 #[test]
5100 fn create_cancellation_and_commit_have_one_winner() {
5101 for _ in 0..32 {
5102 let control = CreateSessionControl::default();
5103 let barrier = Arc::new(std::sync::Barrier::new(2));
5104 let canceller = {
5105 let control = control.clone();
5106 let barrier = barrier.clone();
5107 std::thread::spawn(move || {
5108 barrier.wait();
5109 control.request_cancel()
5110 })
5111 };
5112 barrier.wait();
5113 let committed = control.grant_commit();
5114 let cancelled = canceller.join().expect("canceller panicked");
5115 assert_ne!(committed, cancelled);
5116 assert_eq!(control.cancelled.load(Ordering::Acquire), cancelled);
5117 assert!(!control.is_cancellable());
5118 assert!(!control.request_cancel());
5119 assert!(!control.grant_commit());
5120 }
5121 }
5122
5123 #[tokio::test]
5124 async fn completed_stop_stays_visible_until_cleanup_takes_ownership() {
5125 let state = test_runtime_state();
5126 let (complete, result) = tokio::sync::watch::channel(None);
5127 state.lifecycle.lock().unwrap().insert(
5128 "cleanup-gap".into(),
5129 ActiveLifecycle {
5130 operation_id: "closing-operation".into(),
5131 create_control: None,
5132 kind: LifecycleKind::Close,
5133 cancelled: Arc::new(AtomicBool::new(false)),
5134 started_at_epoch_seconds: 1,
5135 active_stages: BTreeMap::new(),
5136 resume_workspace_id: None,
5137 resume_destination: None,
5138 notice: None,
5139 request_key: None,
5140 _move_guard: None,
5141 move_source_closed: false,
5142 result: result.clone(),
5143 },
5144 );
5145 complete.send_replace(Some(Ok(DaemonLifecycleResult::DeferredCleanup)));
5146 state.remove_completed_lifecycle(&result);
5147 let view = state.active_lifecycles();
5148 assert_eq!(view.len(), 1);
5149 assert_eq!(view[0].operation_id, "closing-operation");
5150 assert!(!view[0].cancellable);
5151 let release = Arc::new(tokio::sync::Notify::new());
5152 let cleanup = state
5153 .start_or_join_lifecycle("cleanup-gap".into(), LifecycleKind::Cleanup, {
5154 let release = release.clone();
5155 move |_, _, _| async move {
5156 release.notified().await;
5157 Ok(DaemonLifecycleResult::Done)
5158 }
5159 })
5160 .unwrap();
5161 state.remove_completed_lifecycle(&result);
5162 let view = state.active_lifecycles();
5163 assert_eq!(view.len(), 1);
5164 assert_eq!(view[0].kind, RuntimeLifecycleKind::Cleanup);
5165 assert_ne!(view[0].operation_id, "closing-operation");
5166 release.notify_one();
5167 RuntimeState::wait_lifecycle_result(cleanup).await.unwrap();
5168 assert!(state.active_lifecycles().is_empty());
5169 }
5170
5171 #[tokio::test]
5172 async fn lifecycle_identity_survives_join_but_changes_for_next_operation() {
5173 let state = test_runtime_state();
5174 let release = Arc::new(tokio::sync::Notify::new());
5175 let first = state
5176 .start_or_join_lifecycle("identity".into(), LifecycleKind::Resume, {
5177 let release = release.clone();
5178 move |_, _, _| async move {
5179 release.notified().await;
5180 Ok(DaemonLifecycleResult::Done)
5181 }
5182 })
5183 .unwrap();
5184 let first_id = state.active_lifecycles()[0].operation_id.clone();
5185 let joined = state
5186 .start_or_join_lifecycle("identity".into(), LifecycleKind::Resume, |_, _, _| async {
5187 panic!("joined operation must not run twice")
5188 })
5189 .unwrap();
5190 assert_eq!(state.active_lifecycles()[0].operation_id, first_id);
5191 release.notify_one();
5192 RuntimeState::wait_lifecycle_result(first.clone())
5193 .await
5194 .unwrap();
5195 RuntimeState::wait_lifecycle_result(joined).await.unwrap();
5196 state.remove_completed_lifecycle(&first);
5197 let second = state
5198 .start_or_join_lifecycle("identity".into(), LifecycleKind::Resume, {
5199 let release = release.clone();
5200 move |_, _, _| async move {
5201 release.notified().await;
5202 Ok(DaemonLifecycleResult::Done)
5203 }
5204 })
5205 .unwrap();
5206 let second_id = state.active_lifecycles()[0].operation_id.clone();
5207 assert_ne!(second_id, first_id);
5208 state.remove_completed_lifecycle(&first);
5209 assert_eq!(state.active_lifecycles()[0].operation_id, second_id);
5210 release.notify_one();
5211 RuntimeState::wait_lifecycle_result(second).await.unwrap();
5212 }
5213
5214 #[tokio::test]
5215 async fn close_waits_for_cancelled_or_committed_provisioning_to_release_ownership() {
5216 for committed in [false, true] {
5217 let state = test_runtime_state();
5218 let control = CreateSessionControl::default();
5219 let release = Arc::new(tokio::sync::Notify::new());
5220 state
5221 .start_or_join_lifecycle_controlled(
5222 "close-race".into(),
5223 LifecycleKind::Create,
5224 None,
5225 None,
5226 Some(control.clone()),
5227 {
5228 let release = release.clone();
5229 move |_, _, _| async move {
5230 release.notified().await;
5231 Ok(DaemonLifecycleResult::Done)
5232 }
5233 },
5234 )
5235 .unwrap();
5236 if committed {
5237 assert!(control.grant_commit());
5238 }
5239 state.request_close("close-race");
5240 let waiter = {
5241 let state = state.clone();
5242 tokio::spawn(async move { state.wait_before_close("close-race").await })
5243 };
5244 tokio::task::yield_now().await;
5245 assert!(
5246 !waiter.is_finished(),
5247 "cleanup must wait for the owning operation"
5248 );
5249 assert_eq!(control.cancelled.load(Ordering::Acquire), !committed);
5250 assert_eq!(
5251 state.session_state("close-race"),
5252 Some(SessionState::Closing)
5253 );
5254 release.notify_one();
5255 tokio::time::timeout(Duration::from_secs(2), waiter)
5256 .await
5257 .unwrap()
5258 .unwrap()
5259 .unwrap();
5260 assert!(!state.lifecycle.lock().unwrap().contains_key("close-race"));
5261 }
5262 }
5263
5264 #[tokio::test]
5265 async fn committed_creation_cannot_be_cancelled_by_another_surface() {
5266 let state = test_runtime_state();
5267 let control = CreateSessionControl::default();
5268 let release = Arc::new(tokio::sync::Notify::new());
5269 let result = state
5270 .start_or_join_lifecycle_controlled(
5271 "committed".into(),
5272 LifecycleKind::Create,
5273 None,
5274 None,
5275 Some(control.clone()),
5276 {
5277 let release = release.clone();
5278 move |_, _, _| async move {
5279 release.notified().await;
5280 Ok(DaemonLifecycleResult::Done)
5281 }
5282 },
5283 )
5284 .unwrap();
5285 assert!(state.active_lifecycles()[0].cancellable);
5286 assert!(control.grant_commit());
5287 assert!(!state.active_lifecycles()[0].cancellable);
5288 assert!(state.cancel_lifecycle("committed").is_err());
5289 state.cancel_lifecycle_if_active("committed");
5290 assert!(!control.cancelled.load(Ordering::Acquire));
5291 release.notify_one();
5292 RuntimeState::wait_lifecycle_result(result).await.unwrap();
5293 }
5294
5295 #[tokio::test]
5296 async fn force_destruction_preempts_a_running_lifecycle_and_waits_for_it() {
5297 let state = test_runtime_state();
5298 let release = Arc::new(tokio::sync::Notify::new());
5299 state
5300 .start_or_join_lifecycle("session-1".into(), LifecycleKind::Create, {
5301 let release = release.clone();
5302 move |_state, _session_id, _cancelled| async move {
5303 release.notified().await;
5304 Ok(DaemonLifecycleResult::Done)
5305 }
5306 })
5307 .unwrap();
5308 tokio::task::yield_now().await;
5309
5310 let preempt_state = state.clone();
5311 let preempted =
5312 tokio::spawn(async move { preempt_state.preempt_active_lifecycle("session-1").await });
5313 tokio::task::yield_now().await;
5314 {
5315 let lifecycle = state
5316 .lifecycle
5317 .lock()
5318 .unwrap_or_else(PoisonError::into_inner);
5319 assert!(
5320 lifecycle
5321 .get("session-1")
5322 .expect("lifecycle entry")
5323 .cancelled
5324 .load(Ordering::Acquire),
5325 "preemption must cancel the running operation"
5326 );
5327 }
5328
5329 release.notify_one();
5330 preempted
5331 .await
5332 .expect("preempt task")
5333 .expect("a cancelled-and-finished lifecycle lets force destruction proceed");
5334 }
5335
5336 #[tokio::test(start_paused = true)]
5337 async fn force_destruction_preemption_times_out_without_destroying() {
5338 let state = test_runtime_state();
5339 state
5340 .start_or_join_lifecycle("session-1".into(), LifecycleKind::Create, {
5341 |_state, _session_id, _cancelled| async move {
5342 std::future::pending::<()>().await;
5344 #[allow(unreachable_code)]
5345 Ok(DaemonLifecycleResult::Done)
5346 }
5347 })
5348 .unwrap();
5349 tokio::task::yield_now().await;
5350
5351 let error = state
5352 .preempt_active_lifecycle("session-1")
5353 .await
5354 .unwrap_err();
5355 assert!(
5356 error
5357 .to_string()
5358 .contains("did not stop after cancellation"),
5359 "{error:#}"
5360 );
5361 }
5362}