Skip to main content

mj_controller/
daemon.rs

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