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