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