Skip to main content

mj_controller/
daemon.rs

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