Skip to main content

mj_controller/
daemon.rs

1//! Persistent per-user controller daemon and its authenticated local protocol.
2
3mod feed;
4mod owner;
5mod record_index;
6mod restart;
7mod session_move;
8use crate::controller::move_session::{
9    MoveMutationGuard, MoveOutcome, MovePreparation, MoveSelection, MoveSessionRequest,
10};
11pub use mj_client::daemon::*;
12use owner::RuntimeStateOwner;
13
14use std::collections::{BTreeMap, BTreeSet, VecDeque};
15use std::fs::{self, OpenOptions};
16use std::io::Write;
17use std::net::{IpAddr, Ipv4Addr, SocketAddr};
18use std::path::{Path, PathBuf};
19use std::sync::atomic::{AtomicBool, AtomicU8, AtomicU64, Ordering};
20use std::sync::{Arc, Mutex, PoisonError};
21use std::time::{Duration, Instant, SystemTime, UNIX_EPOCH};
22
23use crate::database::StoreSchemaMismatch;
24use crate::recovery_gate::RecoveryObserver;
25use crate::targets::{
26    CancellableProcessExecutor, CommandExecutor, CommandOutput, CommandSpec, ProcessExecutor,
27    ProvisionStage, ProvisionStageGuard,
28};
29use agent_client_protocol::schema::v1::{ContentBlock, TextContent};
30use anyhow::{Context, Result, anyhow, bail, ensure};
31use mj_core::config::Config;
32use mj_core::refusal::Refusal;
33use mj_core::relay::RelayCommand;
34use mj_core::state::{RecoveryObservation, SessionRecord, SessionState};
35
36use crate::controller::{
37    BeforeClose, BranchDisposition, CheckoutDisposition, Controller, ControllerStoreGuard,
38    SessionLaunchOptions, SessionResumeOptions,
39};
40use crate::review_host::TurnReviewHost;
41use crate::session_manager::{
42    ManagedSessionView, RemoteSessionPublisher, RemoteSessionRequest, SessionManagerChannels,
43    SessionManagerControl, ViewError, new_command_id, spawn_remote_session_manager,
44};
45#[cfg(test)]
46use crate::session_manager::{RelaySessionTarget, RemoteSessionRequests, SessionManagerShutdown};
47use crate::worker_upgrade::{WorkerUpgradeObservation, WorkerUpgradeObserver};
48use mj_core::workspace::WorkspaceRecord;
49use tokio::net::{TcpListener, TcpStream};
50use tokio_util::sync::CancellationToken;
51
52use crate::pollers::{
53    dashboard_worker_targets, interrupted_destroy_session_ids, interrupted_suspend_session_ids,
54    reserve_recovery_or_cancel, spawn_image_refresher, unowned_interrupted_lifecycles,
55};
56
57// Move preparation now reports whether source state must be recovered without its harness.
58
59/// How long the epilogue is given before the process leaves anyway.
60///
61/// Every daemon exit -- stop, SIGTERM, idle, a store that moved underneath it
62/// -- unwinds through the same epilogue, and every step of it is bounded in
63/// practice. This makes "the daemon did not stop" impossible rather than
64/// unlikely, and it must stay well inside [`STOP_TIMEOUT`] so a client waiting
65/// on a stop sees the exit rather than its own deadline.
66const SHUTDOWN_FORCE_EXIT_TIMEOUT: Duration = Duration::from_secs(10);
67
68/// How long force destruction waits for a cancelled lifecycle to actually
69/// stop before refusing to destroy under it. Cancellation kills the
70/// operation's child process groups and unwinds its persistence, which is
71/// fast in practice; an operation that outlives this bound is wedged in a
72/// way destruction must not paper over.
73const FORCE_DESTROY_PREEMPT_TIMEOUT: Duration = Duration::from_secs(8);
74
75/// Cancellation and committing a newly started session are one atomic decision.
76#[derive(Clone, Default)]
77pub struct CreateSessionControl {
78    state: Arc<AtomicU8>,
79    pub cancelled: Arc<AtomicBool>,
80}
81
82impl CreateSessionControl {
83    pub fn request_cancel(&self) -> bool {
84        let accepted = self
85            .state
86            .compare_exchange(0, 1, Ordering::AcqRel, Ordering::Acquire)
87            .is_ok();
88        if accepted {
89            self.cancelled.store(true, Ordering::Release);
90        }
91        accepted
92    }
93
94    pub fn grant_commit(&self) -> bool {
95        self.state
96            .compare_exchange(0, 2, Ordering::AcqRel, Ordering::Acquire)
97            .is_ok()
98    }
99
100    fn is_cancellable(&self) -> bool {
101        self.state.load(Ordering::Acquire) == 0
102    }
103}
104
105#[derive(Debug, Clone)]
106struct Attachment {
107    pid: u32,
108}
109
110#[derive(Default)]
111struct QuotaBoard {
112    snapshot: mj_client::quota::QuotaSnapshot,
113    refresh: Option<tokio::sync::mpsc::Sender<()>>,
114}
115
116pub struct RuntimeState {
117    attachments: Mutex<BTreeMap<String, Attachment>>,
118    phone_status: Mutex<WebViewerStatus>,
119    pub web_viewer: crate::web_viewer::ViewerControl,
120    ever_attached: AtomicBool,
121    revisions: RuntimeRevisions,
122    workspaces_tx: tokio::sync::watch::Sender<Vec<WorkspaceRecord>>,
123    workspace_refresh: tokio::sync::Mutex<()>,
124    session_manager: SessionManagerControl,
125    owner: Mutex<RuntimeStateOwner>,
126    /// Rebuilt by the delegation service at daemon startup; weak so the
127    /// backend's ExportRuntime reference does not form an Arc cycle.
128    pub(crate) wait_prompt_backend:
129        std::sync::OnceLock<std::sync::Weak<crate::server_runtime::api::ApiBackend>>,
130    /// Parents whose level-triggered wait-prompt reconcile failed. The
131    /// delegation coordinator retries these; startup also rebuilds the set
132    /// from child records after a daemon restart.
133    pub(crate) wait_prompt_retries: Mutex<BTreeMap<String, Instant>>,
134    pub(crate) credential_targets:
135        Arc<tokio::sync::watch::Sender<Vec<mj_core::credentials::CredentialSyncTarget>>>,
136    feed: Mutex<feed::RuntimeHistory>,
137    committed: Option<
138        tokio::sync::watch::Receiver<
139            std::result::Result<Arc<crate::database::CommittedState>, Arc<str>>,
140        >,
141    >,
142    workspace_closes: Mutex<BTreeMap<String, Arc<AtomicBool>>>,
143    /// Resume ownership can precede its durable workspace assignment.
144    workspace_resume_admission: Mutex<BTreeMap<String, Arc<tokio::sync::RwLock<()>>>>,
145    /// The bounded wait for each live session's harness to become usable.
146    /// Driven only by the daemon's background readiness sweep.
147    harness_readiness: Mutex<HarnessReadinessWatch>,
148    /// Work waiting for one session's harness to become ready: the prompts a
149    /// person typed while it started, and the hand-off a restored session
150    /// carries. One ordered queue per session, each drained by one task.
151    startup_prompts: Mutex<BTreeMap<String, StartupQueue>>,
152    startup_enqueue: tokio::sync::Mutex<()>,
153    controller_loader: fn() -> Result<Controller>,
154    config_mutation: tokio::sync::Mutex<()>,
155    projects: Arc<crate::project_catalog::Catalog>,
156    profile_catalog: crate::review_host::SharedProfileCatalog,
157    recovery_observer: RecoveryObserver,
158    worker_upgrade_observer: WorkerUpgradeObserver,
159    /// Recent background notices, newest last, with the id of the next one.
160    /// Bounded: a surface that never attaches must not make this grow.
161    notices: Mutex<VecDeque<RuntimeNotice>>,
162    next_notice_id: AtomicU64,
163    /// What the quota poller last published, and how to wake it. The poller
164    /// is the only process that asks a provider for quota; everything else
165    /// reads this.
166    quota: Mutex<QuotaBoard>,
167    /// The one capacity poller and what it shares; set once the daemon
168    /// starts it.
169    capacity: std::sync::OnceLock<capacity::CapacityFeed>,
170    /// What `[review]` last said, republished by the target refresher.
171    review_config: Arc<Mutex<mj_core::config::ReviewConfig>>,
172    /// Turn review runs here, in the process that owns every session, so a
173    /// review happens whether the terminal, the phone, or nobody is attached.
174    review_host: TurnReviewHost,
175    /// Publishes checkpointed sessions into the user's SessionWiki index.
176    wiki: crate::sessionwiki::WikiIndexer,
177}
178
179/// One monotonic cursor shared by daemon snapshots and their wake-up feed.
180///
181/// Allocations can come from independent UI and daemon tasks. Publishing an
182/// older allocation after a newer one must not move the watch channel
183/// backwards, so publication compares against the last visible cursor.
184#[derive(Clone)]
185struct RuntimeRevisions {
186    allocated: Arc<std::sync::atomic::AtomicU64>,
187    published: tokio::sync::watch::Sender<u64>,
188}
189
190impl RuntimeRevisions {
191    fn new(initial: u64) -> Self {
192        let (published, _) = tokio::sync::watch::channel(initial);
193        Self {
194            allocated: Arc::new(std::sync::atomic::AtomicU64::new(initial)),
195            published,
196        }
197    }
198
199    fn allocate(&self) -> u64 {
200        self.allocated.fetch_add(1, Ordering::AcqRel) + 1
201    }
202
203    fn publish(&self) -> u64 {
204        let revision = self.allocate();
205        self.publish_allocated(revision);
206        revision
207    }
208
209    fn publish_allocated(&self, revision: u64) {
210        self.published.send_if_modified(|visible| {
211            if revision > *visible {
212                *visible = revision;
213                true
214            } else {
215                false
216            }
217        });
218    }
219
220    fn notifier(&self) -> Arc<dyn Fn() + Send + Sync> {
221        let revisions = self.clone();
222        Arc::new(move || {
223            revisions.publish();
224        })
225    }
226
227    fn subscribe(&self) -> tokio::sync::watch::Receiver<u64> {
228        self.published.subscribe()
229    }
230
231    fn current(&self) -> u64 {
232        self.allocated.load(Ordering::Acquire)
233    }
234}
235
236#[derive(Debug, Clone, Copy, PartialEq, Eq)]
237enum LifecycleKind {
238    Create,
239    Suspend,
240    Resume,
241    Restart,
242    Move,
243    ForceStop,
244    DestroyStopped,
245    /// The archive job's destruction: the same teardown as `DestroyStopped`,
246    /// with the session's git branch kept unless another branch already
247    /// contains every one of its commits. Surfaces see it as a destroy.
248    ArchiveStopped,
249    ForceDestroy,
250    /// A sub-agent stopped because its parent is being suspended: the
251    /// teardown of `ForceDestroy`, owned by the parent's suspend, so it is
252    /// shown as a stop and cannot be cancelled on its own.
253    StopSubagent,
254    /// A sub-agent whose turn ended and whose parent was told is having its
255    /// worker stopped, keeping everything else (#1161).
256    Park,
257    /// A parked sub-agent's worker is being started again for its parent's
258    /// `send_input`.
259    Unpark,
260    Cleanup,
261    StartupCleanup,
262}
263
264impl LifecycleKind {
265    /// Whether this operation is a destroy's teardown of the session, which a
266    /// second destroy of the same session waits for instead of cancelling.
267    fn is_teardown(self) -> bool {
268        matches!(
269            self,
270            LifecycleKind::ForceDestroy
271                | LifecycleKind::DestroyStopped
272                | LifecycleKind::ArchiveStopped
273                | LifecycleKind::StopSubagent
274        )
275    }
276
277    /// What a refusal calls this operation, so a person told that a session is
278    /// busy learns which operation is holding it (#1010).
279    fn label(self) -> &'static str {
280        match self {
281            LifecycleKind::Create => "create",
282            LifecycleKind::Suspend => "suspend",
283            LifecycleKind::Resume => "resume",
284            LifecycleKind::Restart => "restart",
285            LifecycleKind::Move => "move",
286            LifecycleKind::ForceStop => "force stop",
287            LifecycleKind::DestroyStopped => "destroy",
288            LifecycleKind::ArchiveStopped => "archive",
289            LifecycleKind::ForceDestroy => "force destroy",
290            LifecycleKind::StopSubagent => "stop",
291            LifecycleKind::Park => "park",
292            LifecycleKind::Unpark => "restart",
293            LifecycleKind::Cleanup => "cleanup",
294            LifecycleKind::StartupCleanup => "failed startup cleanup",
295        }
296    }
297}
298
299/// Whether a lifecycle has exclusive ownership of the worker target, so the
300/// session manager must stop polling it. A graceful close needs the manager's
301/// relay lease through checkpointing and sealing; once the durable state says
302/// `Destroying`, that lease has been released and target teardown is exclusive.
303///
304/// A park needs the session actor: it takes the actor's connection to admit
305/// the park only while the worker is idle. Once the record says `Parked` the
306/// manager drops the session on its own. An unpark owns the target, so no
307/// actor races the worker it is starting.
308fn lifecycle_owns_worker_target(kind: LifecycleKind, state: Option<SessionState>) -> bool {
309    match kind {
310        LifecycleKind::Suspend => state == Some(SessionState::Destroying),
311        LifecycleKind::Park => false,
312        LifecycleKind::Move | LifecycleKind::Restart => !matches!(
313            state,
314            Some(
315                SessionState::Running
316                    | SessionState::Disconnected
317                    | SessionState::Checkpointing
318                    | SessionState::Closing
319            )
320        ),
321        _ => true,
322    }
323}
324
325/// Whether a running lifecycle can still be cancelled. A graceful close has a
326/// point of no return: once the durable state says `Destroying`, the verified
327/// checkpoint is sealed and the record has already committed to losing its
328/// target, so stopping the teardown only strands the target. A sub-agent's
329/// stop belongs to its parent's suspend, which is what a person cancels.
330/// Every other lifecycle stays cancellable while it runs.
331///
332/// A park and an unpark belong to the parent's sub-agent tools: a person ends
333/// either by closing the child, which waits for them.
334fn lifecycle_cancellable(kind: LifecycleKind, state: Option<SessionState>) -> bool {
335    !(state == Some(SessionState::StartupCleanup)
336        || kind == LifecycleKind::StartupCleanup
337        || matches!(kind, LifecycleKind::Suspend | LifecycleKind::Restart)
338            && state == Some(SessionState::Destroying)
339        || matches!(
340            kind,
341            LifecycleKind::StopSubagent | LifecycleKind::Park | LifecycleKind::Unpark
342        ))
343}
344
345/// How a stop request has to be carried out, given the durable record.
346#[derive(Debug, Clone, Copy, PartialEq, Eq)]
347enum CloseRoute {
348    /// Run the graceful close from the start.
349    Graceful,
350    /// A previous close stopped partway; finish it from its checkpoint.
351    RecoverInterrupted,
352    /// Nothing to checkpoint: tear down whatever target is left and settle.
353    SettleWithoutCheckpoint,
354    /// Already stopped, but the target still has to be removed.
355    DeferredCleanup,
356    /// Already stopped with nothing left to do.
357    Done,
358}
359
360/// A record mid-close with a live target cannot be closed again from the start:
361/// its worker socket is gone, so a fresh checkpoint attempt only fails on
362/// connect. Recovery finishes it from the checkpoint the first close verified.
363/// `subagent` says the session is a Mjolnir sub-agent; see
364/// [`crate::controller::has_nothing_to_checkpoint`].
365fn close_route(session: Option<&SessionRecord>, subagent: bool) -> CloseRoute {
366    let Some(session) = session else {
367        return CloseRoute::Graceful;
368    };
369    if crate::pollers::is_interrupted_close(session) {
370        CloseRoute::RecoverInterrupted
371    } else if session.state == SessionState::Stopped {
372        if session.target.is_some() {
373            CloseRoute::DeferredCleanup
374        } else {
375            CloseRoute::Done
376        }
377    } else if crate::controller::has_nothing_to_checkpoint(session, subagent) {
378        CloseRoute::SettleWithoutCheckpoint
379    } else {
380        CloseRoute::Graceful
381    }
382}
383
384/// The durable state of one record as the locked controller holds it.
385fn durable_session_state(controller: &Controller, session_id: &str) -> Option<SessionState> {
386    controller
387        .state
388        .sessions
389        .get(session_id)
390        .map(|session| session.state)
391}
392
393/// One piece of work that waits for a starting session's harness.
394///
395/// Both kinds are ordered against each other on purpose: a restored session's
396/// hand-off is the hidden context its first prompt reads, so it has to be
397/// installed before any queued prompt is submitted.
398#[derive(serde::Serialize, serde::Deserialize)]
399pub(crate) enum StartupStep {
400    InstallHandoff(Box<mj_core::archive::CanonicalSessionSnapshot>),
401    PreparedHandoff {
402        text: String,
403    },
404    Prompt {
405        text: String,
406        inherited_draft: Option<String>,
407    },
408    Configure {
409        key: String,
410        value: String,
411        optional: bool,
412    },
413    ApiPrompt {
414        text: String,
415    },
416}
417
418/// One supervised drain per session. Pending payloads stay in SQLite.
419struct StartupQueue {
420    identity: Arc<()>,
421    last_error: Option<String>,
422    cancel: CancellationToken,
423    task: Option<tokio::task::JoinHandle<()>>,
424}
425
426struct ActiveLifecycle {
427    upgrade_work: Arc<Mutex<Option<crate::upgrade::Work>>>,
428    phase: LifecyclePhase,
429    operation_id: String,
430    create_control: Option<CreateSessionControl>,
431    kind: LifecycleKind,
432    cancelled: Arc<AtomicBool>,
433    started_at_epoch_seconds: u64,
434    active_stages: BTreeMap<ProvisionStage, (usize, u64)>,
435    /// The workspace a resume is claiming before its durable record changes.
436    /// Workspace deletion consults this so it cannot race the claim.
437    resume_workspace_id: Option<String>,
438    resume_destination: Option<(String, String)>,
439    notice: Option<String>,
440    request_key: Option<String>,
441    _move_guard: Option<MoveMutationGuard>,
442    result: LifecycleWatch,
443}
444
445enum LifecyclePhase {
446    Executing,
447    MovingDestination,
448    Cancelling,
449    CancellingMoveDestination,
450    Completed(LifecycleResult),
451}
452
453impl ActiveLifecycle {
454    fn is_running(&self) -> bool {
455        !matches!(self.phase, LifecyclePhase::Completed(_))
456    }
457
458    fn is_visible(&self) -> bool {
459        self.is_running()
460            || matches!(
461                self.phase,
462                LifecyclePhase::Completed(Ok(DaemonLifecycleResult::DeferredCleanup))
463            )
464    }
465
466    fn request_cancel(&mut self) -> bool {
467        if !matches!(
468            self.phase,
469            LifecyclePhase::Executing | LifecyclePhase::MovingDestination
470        ) {
471            return false;
472        }
473        let accepted = if let Some(control) = &self.create_control {
474            control.request_cancel()
475        } else {
476            !self.cancelled.swap(true, Ordering::AcqRel)
477        };
478        if accepted {
479            self.phase = match self.phase {
480                LifecyclePhase::MovingDestination => LifecyclePhase::CancellingMoveDestination,
481                _ => LifecyclePhase::Cancelling,
482            };
483        }
484        accepted
485    }
486
487    fn is_cancellable(&self) -> bool {
488        matches!(
489            self.phase,
490            LifecyclePhase::Executing | LifecyclePhase::MovingDestination
491        ) && self.create_control.as_ref().map_or_else(
492            || !self.cancelled.load(Ordering::Acquire),
493            CreateSessionControl::is_cancellable,
494        )
495    }
496}
497
498#[derive(Debug, Clone)]
499enum DaemonLifecycleResult {
500    Done,
501    DeferredCleanup,
502    Move(MoveOutcome),
503    Park(crate::controller::ParkOutcome),
504}
505
506/// How one lifecycle operation ended when it failed.
507///
508/// The result is broadcast to every waiter, which is why it cannot simply be
509/// the `anyhow::Error`: that is not clonable. Keeping the refusal beside the
510/// text is what lets a reason written for the caller survive the crossing; a
511/// failure rebuilt from a string alone would arrive as an internal fault.
512#[derive(Debug, Clone)]
513pub(crate) struct LifecycleFailure {
514    pub(crate) detail: String,
515    pub(crate) refusal: Option<Refusal>,
516}
517
518/// One lifecycle operation's outcome, and the channel every waiter reads it
519/// from. `None` means the operation is still running.
520type LifecycleResult = std::result::Result<DaemonLifecycleResult, LifecycleFailure>;
521type LifecycleWatch = tokio::sync::watch::Receiver<Option<LifecycleResult>>;
522
523impl LifecycleFailure {
524    fn of(error: &anyhow::Error) -> Self {
525        Self {
526            detail: format!("{error:#}"),
527            refusal: Refusal::of(error),
528        }
529    }
530
531    /// A failure with no reason written for a caller, such as a task that died
532    /// before the operation could say anything about itself.
533    fn internal(detail: impl Into<String>) -> Self {
534        Self {
535            detail: detail.into(),
536            refusal: None,
537        }
538    }
539
540    /// Rebuild the error a waiter sees, with the refusal still attached.
541    fn into_error(self) -> anyhow::Error {
542        match self.refusal {
543            Some(refusal) => anyhow::Error::new(refusal).context(self.detail),
544            None => anyhow::Error::msg(self.detail),
545        }
546    }
547}
548
549impl std::fmt::Display for LifecycleFailure {
550    fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
551        formatter.write_str(&self.detail)
552    }
553}
554
555impl From<LifecycleKind> for RuntimeLifecycleKind {
556    fn from(kind: LifecycleKind) -> Self {
557        match kind {
558            LifecycleKind::Create => Self::Create,
559            LifecycleKind::Suspend => Self::Suspend,
560            LifecycleKind::Resume => Self::Resume,
561            LifecycleKind::Restart => Self::Resume,
562            LifecycleKind::Move => Self::Move,
563            LifecycleKind::ForceStop => Self::ForceStop,
564            LifecycleKind::DestroyStopped | LifecycleKind::ArchiveStopped => Self::DestroyStopped,
565            LifecycleKind::ForceDestroy => Self::ForceDestroy,
566            // A park is shown as the sub-agent stop it is, and an unpark as
567            // the resume it is; neither needs a wire kind of its own.
568            LifecycleKind::StopSubagent | LifecycleKind::Park => Self::StopSubagent,
569            LifecycleKind::Unpark => Self::Resume,
570            LifecycleKind::Cleanup | LifecycleKind::StartupCleanup => Self::Cleanup,
571        }
572    }
573}
574
575mod capacity;
576mod close;
577mod close_workspace;
578mod create;
579mod readiness;
580use readiness::{HarnessReadinessWatch, ReadinessObservation, UnreadySession};
581mod lifecycle;
582mod resume;
583mod snapshot;
584pub(crate) mod startup_followup;
585mod state;
586mod subagent_park;
587mod support;
588mod views;
589use support::*;
590pub(crate) mod delegation;
591mod diagnostics;
592mod process;
593pub use process::*;
594mod serve;
595use serve::*;
596mod actions;
597use actions::*;
598mod guards;
599pub(crate) use guards::*;
600
601#[cfg(test)]
602mod tests;
603
604mod continuation;