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