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