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