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