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, 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    pub(crate) credential_targets:
127        Arc<tokio::sync::watch::Sender<Vec<mj_core::credentials::CredentialSyncTarget>>>,
128    feed: Mutex<feed::RuntimeHistory>,
129    committed: Option<
130        tokio::sync::watch::Receiver<
131            std::result::Result<Arc<crate::database::CommittedState>, Arc<str>>,
132        >,
133    >,
134    workspace_closes: Mutex<BTreeMap<String, Arc<AtomicBool>>>,
135    /// Resume ownership can precede its durable workspace assignment.
136    workspace_resume_admission: Mutex<BTreeMap<String, Arc<tokio::sync::RwLock<()>>>>,
137    /// The bounded wait for each live session's harness to become usable.
138    /// Driven only by the daemon's background readiness sweep.
139    harness_readiness: Mutex<HarnessReadinessWatch>,
140    /// Work waiting for one session's harness to become ready: the prompts a
141    /// person typed while it started, and the hand-off a restored session
142    /// carries. One ordered queue per session, each drained by one task.
143    startup_prompts: Mutex<BTreeMap<String, StartupQueue>>,
144    startup_enqueue: tokio::sync::Mutex<()>,
145    controller_loader: fn() -> Result<Controller>,
146    config_mutation: tokio::sync::Mutex<()>,
147    projects: Arc<crate::project_catalog::Catalog>,
148    profile_catalog: crate::review_host::SharedProfileCatalog,
149    recovery_observer: RecoveryObserver,
150    worker_upgrade_observer: WorkerUpgradeObserver,
151    /// Recent background notices, newest last, with the id of the next one.
152    /// Bounded: a surface that never attaches must not make this grow.
153    notices: Mutex<VecDeque<RuntimeNotice>>,
154    next_notice_id: AtomicU64,
155    /// What the quota poller last published, and how to wake it. The poller
156    /// is the only process that asks a provider for quota; everything else
157    /// reads this.
158    quota: Mutex<QuotaBoard>,
159    /// The one capacity poller and what it shares; set once the daemon
160    /// starts it.
161    capacity: std::sync::OnceLock<capacity::CapacityFeed>,
162    /// What `[review]` last said, republished by the target refresher.
163    review_config: Arc<Mutex<mj_core::config::ReviewConfig>>,
164    /// Turn review runs here, in the process that owns every session, so a
165    /// review happens whether the terminal, the phone, or nobody is attached.
166    review_host: TurnReviewHost,
167    /// Publishes checkpointed sessions into the user's SessionWiki index.
168    wiki: crate::sessionwiki::WikiIndexer,
169}
170
171/// One monotonic cursor shared by daemon snapshots and their wake-up feed.
172///
173/// Allocations can come from independent UI and daemon tasks. Publishing an
174/// older allocation after a newer one must not move the watch channel
175/// backwards, so publication compares against the last visible cursor.
176#[derive(Clone)]
177struct RuntimeRevisions {
178    allocated: Arc<std::sync::atomic::AtomicU64>,
179    published: tokio::sync::watch::Sender<u64>,
180}
181
182impl RuntimeRevisions {
183    fn new(initial: u64) -> Self {
184        let (published, _) = tokio::sync::watch::channel(initial);
185        Self {
186            allocated: Arc::new(std::sync::atomic::AtomicU64::new(initial)),
187            published,
188        }
189    }
190
191    fn allocate(&self) -> u64 {
192        self.allocated.fetch_add(1, Ordering::AcqRel) + 1
193    }
194
195    fn publish(&self) -> u64 {
196        let revision = self.allocate();
197        self.publish_allocated(revision);
198        revision
199    }
200
201    fn publish_allocated(&self, revision: u64) {
202        self.published.send_if_modified(|visible| {
203            if revision > *visible {
204                *visible = revision;
205                true
206            } else {
207                false
208            }
209        });
210    }
211
212    fn notifier(&self) -> Arc<dyn Fn() + Send + Sync> {
213        let revisions = self.clone();
214        Arc::new(move || {
215            revisions.publish();
216        })
217    }
218
219    fn subscribe(&self) -> tokio::sync::watch::Receiver<u64> {
220        self.published.subscribe()
221    }
222
223    fn current(&self) -> u64 {
224        self.allocated.load(Ordering::Acquire)
225    }
226}
227
228#[derive(Debug, Clone, Copy, PartialEq, Eq)]
229enum LifecycleKind {
230    Create,
231    Suspend,
232    Resume,
233    Restart,
234    Move,
235    ForceStop,
236    DestroyStopped,
237    /// The archive job's destruction: the same teardown as `DestroyStopped`,
238    /// with the session's git branch kept unless another branch already
239    /// contains every one of its commits. Surfaces see it as a destroy.
240    ArchiveStopped,
241    ForceDestroy,
242    /// A sub-agent stopped because its parent is being suspended: the
243    /// teardown of `ForceDestroy`, owned by the parent's suspend, so it is
244    /// shown as a stop and cannot be cancelled on its own.
245    StopSubagent,
246    /// A sub-agent whose turn ended and whose parent was told is having its
247    /// worker stopped, keeping everything else (#1161).
248    Park,
249    /// A parked sub-agent's worker is being started again for its parent's
250    /// `send_input`.
251    Unpark,
252    Cleanup,
253    StartupCleanup,
254}
255
256impl LifecycleKind {
257    /// Whether this operation is a destroy's teardown of the session, which a
258    /// second destroy of the same session waits for instead of cancelling.
259    fn is_teardown(self) -> bool {
260        matches!(
261            self,
262            LifecycleKind::ForceDestroy
263                | LifecycleKind::DestroyStopped
264                | LifecycleKind::ArchiveStopped
265                | LifecycleKind::StopSubagent
266        )
267    }
268
269    /// What a refusal calls this operation, so a person told that a session is
270    /// busy learns which operation is holding it (#1010).
271    fn label(self) -> &'static str {
272        match self {
273            LifecycleKind::Create => "create",
274            LifecycleKind::Suspend => "suspend",
275            LifecycleKind::Resume => "resume",
276            LifecycleKind::Restart => "restart",
277            LifecycleKind::Move => "move",
278            LifecycleKind::ForceStop => "force stop",
279            LifecycleKind::DestroyStopped => "destroy",
280            LifecycleKind::ArchiveStopped => "archive",
281            LifecycleKind::ForceDestroy => "force destroy",
282            LifecycleKind::StopSubagent => "stop",
283            LifecycleKind::Park => "park",
284            LifecycleKind::Unpark => "restart",
285            LifecycleKind::Cleanup => "cleanup",
286            LifecycleKind::StartupCleanup => "failed startup cleanup",
287        }
288    }
289}
290
291/// Whether a lifecycle has exclusive ownership of the worker target, so the
292/// session manager must stop polling it. A graceful close needs the manager's
293/// relay lease through checkpointing and sealing; once the durable state says
294/// `Destroying`, that lease has been released and target teardown is exclusive.
295///
296/// A park needs the session actor: it takes the actor's connection to admit
297/// the park only while the worker is idle. Once the record says `Parked` the
298/// manager drops the session on its own. An unpark owns the target, so no
299/// actor races the worker it is starting.
300fn lifecycle_owns_worker_target(kind: LifecycleKind, state: Option<SessionState>) -> bool {
301    match kind {
302        LifecycleKind::Suspend => state == Some(SessionState::Destroying),
303        LifecycleKind::Park => false,
304        LifecycleKind::Move | LifecycleKind::Restart => !matches!(
305            state,
306            Some(
307                SessionState::Running
308                    | SessionState::Disconnected
309                    | SessionState::Checkpointing
310                    | SessionState::Closing
311            )
312        ),
313        _ => true,
314    }
315}
316
317/// Whether a running lifecycle can still be cancelled. A graceful close has a
318/// point of no return: once the durable state says `Destroying`, the verified
319/// checkpoint is sealed and the record has already committed to losing its
320/// target, so stopping the teardown only strands the target. A sub-agent's
321/// stop belongs to its parent's suspend, which is what a person cancels.
322/// Every other lifecycle stays cancellable while it runs.
323///
324/// A park and an unpark belong to the parent's sub-agent tools: a person ends
325/// either by closing the child, which waits for them.
326fn lifecycle_cancellable(kind: LifecycleKind, state: Option<SessionState>) -> bool {
327    !(state == Some(SessionState::StartupCleanup)
328        || kind == LifecycleKind::StartupCleanup
329        || matches!(kind, LifecycleKind::Suspend | LifecycleKind::Restart)
330            && state == Some(SessionState::Destroying)
331        || matches!(
332            kind,
333            LifecycleKind::StopSubagent | LifecycleKind::Park | LifecycleKind::Unpark
334        ))
335}
336
337/// How a stop request has to be carried out, given the durable record.
338#[derive(Debug, Clone, Copy, PartialEq, Eq)]
339enum CloseRoute {
340    /// Run the graceful close from the start.
341    Graceful,
342    /// A previous close stopped partway; finish it from its checkpoint.
343    RecoverInterrupted,
344    /// Nothing to checkpoint: tear down whatever target is left and settle.
345    SettleWithoutCheckpoint,
346    /// Already stopped, but the target still has to be removed.
347    DeferredCleanup,
348    /// Already stopped with nothing left to do.
349    Done,
350}
351
352/// A record mid-close with a live target cannot be closed again from the start:
353/// its worker socket is gone, so a fresh checkpoint attempt only fails on
354/// connect. Recovery finishes it from the checkpoint the first close verified.
355/// `subagent` says the session is a Mjolnir sub-agent; see
356/// [`crate::controller::has_nothing_to_checkpoint`].
357fn close_route(session: Option<&SessionRecord>, subagent: bool) -> CloseRoute {
358    let Some(session) = session else {
359        return CloseRoute::Graceful;
360    };
361    if crate::pollers::is_interrupted_close(session) {
362        CloseRoute::RecoverInterrupted
363    } else if session.state == SessionState::Stopped {
364        if session.target.is_some() {
365            CloseRoute::DeferredCleanup
366        } else {
367            CloseRoute::Done
368        }
369    } else if crate::controller::has_nothing_to_checkpoint(session, subagent) {
370        CloseRoute::SettleWithoutCheckpoint
371    } else {
372        CloseRoute::Graceful
373    }
374}
375
376/// The durable state of one record as the locked controller holds it.
377fn durable_session_state(controller: &Controller, session_id: &str) -> Option<SessionState> {
378    controller
379        .state
380        .sessions
381        .get(session_id)
382        .map(|session| session.state)
383}
384
385/// One piece of work that waits for a starting session's harness.
386///
387/// Both kinds are ordered against each other on purpose: a restored session's
388/// hand-off is the hidden context its first prompt reads, so it has to be
389/// installed before any queued prompt is submitted.
390#[derive(serde::Serialize, serde::Deserialize)]
391pub(crate) enum StartupStep {
392    InstallHandoff(Box<mj_core::archive::CanonicalSessionSnapshot>),
393    PreparedHandoff {
394        text: String,
395    },
396    Prompt {
397        text: String,
398        inherited_draft: Option<String>,
399    },
400    Configure {
401        key: String,
402        value: String,
403        optional: bool,
404    },
405    ApiPrompt {
406        text: String,
407    },
408}
409
410/// One supervised drain per session. Pending payloads stay in SQLite.
411struct StartupQueue {
412    identity: Arc<()>,
413    last_error: Option<String>,
414    cancel: CancellationToken,
415    task: Option<tokio::task::JoinHandle<()>>,
416}
417
418struct ActiveLifecycle {
419    upgrade_work: Arc<Mutex<Option<crate::upgrade::Work>>>,
420    phase: LifecyclePhase,
421    operation_id: String,
422    create_control: Option<CreateSessionControl>,
423    kind: LifecycleKind,
424    cancelled: Arc<AtomicBool>,
425    started_at_epoch_seconds: u64,
426    active_stages: BTreeMap<ProvisionStage, (usize, u64)>,
427    /// The workspace a resume is claiming before its durable record changes.
428    /// Workspace deletion consults this so it cannot race the claim.
429    resume_workspace_id: Option<String>,
430    resume_destination: Option<(String, String)>,
431    notice: Option<String>,
432    request_key: Option<String>,
433    _move_guard: Option<MoveMutationGuard>,
434    result: LifecycleWatch,
435}
436
437enum LifecyclePhase {
438    Executing,
439    MovingDestination,
440    Cancelling,
441    CancellingMoveDestination,
442    Completed(LifecycleResult),
443}
444
445impl ActiveLifecycle {
446    fn is_running(&self) -> bool {
447        !matches!(self.phase, LifecyclePhase::Completed(_))
448    }
449
450    fn is_visible(&self) -> bool {
451        self.is_running()
452            || matches!(
453                self.phase,
454                LifecyclePhase::Completed(Ok(DaemonLifecycleResult::DeferredCleanup))
455            )
456    }
457
458    fn request_cancel(&mut self) -> bool {
459        if !matches!(
460            self.phase,
461            LifecyclePhase::Executing | LifecyclePhase::MovingDestination
462        ) {
463            return false;
464        }
465        let accepted = if let Some(control) = &self.create_control {
466            control.request_cancel()
467        } else {
468            !self.cancelled.swap(true, Ordering::AcqRel)
469        };
470        if accepted {
471            self.phase = match self.phase {
472                LifecyclePhase::MovingDestination => LifecyclePhase::CancellingMoveDestination,
473                _ => LifecyclePhase::Cancelling,
474            };
475        }
476        accepted
477    }
478
479    fn is_cancellable(&self) -> bool {
480        matches!(
481            self.phase,
482            LifecyclePhase::Executing | LifecyclePhase::MovingDestination
483        ) && self.create_control.as_ref().map_or_else(
484            || !self.cancelled.load(Ordering::Acquire),
485            CreateSessionControl::is_cancellable,
486        )
487    }
488}
489
490#[derive(Debug, Clone)]
491enum DaemonLifecycleResult {
492    Done,
493    DeferredCleanup,
494    Move(MoveOutcome),
495    Park(crate::controller::ParkOutcome),
496}
497
498/// How one lifecycle operation ended when it failed.
499///
500/// The result is broadcast to every waiter, which is why it cannot simply be
501/// the `anyhow::Error`: that is not clonable. Keeping the refusal beside the
502/// text is what lets a reason written for the caller survive the crossing; a
503/// failure rebuilt from a string alone would arrive as an internal fault.
504#[derive(Debug, Clone)]
505pub(crate) struct LifecycleFailure {
506    pub(crate) detail: String,
507    pub(crate) refusal: Option<Refusal>,
508}
509
510/// One lifecycle operation's outcome, and the channel every waiter reads it
511/// from. `None` means the operation is still running.
512type LifecycleResult = std::result::Result<DaemonLifecycleResult, LifecycleFailure>;
513type LifecycleWatch = tokio::sync::watch::Receiver<Option<LifecycleResult>>;
514
515impl LifecycleFailure {
516    fn of(error: &anyhow::Error) -> Self {
517        Self {
518            detail: format!("{error:#}"),
519            refusal: Refusal::of(error),
520        }
521    }
522
523    /// A failure with no reason written for a caller, such as a task that died
524    /// before the operation could say anything about itself.
525    fn internal(detail: impl Into<String>) -> Self {
526        Self {
527            detail: detail.into(),
528            refusal: None,
529        }
530    }
531
532    /// Rebuild the error a waiter sees, with the refusal still attached.
533    fn into_error(self) -> anyhow::Error {
534        match self.refusal {
535            Some(refusal) => anyhow::Error::new(refusal).context(self.detail),
536            None => anyhow::Error::msg(self.detail),
537        }
538    }
539}
540
541impl std::fmt::Display for LifecycleFailure {
542    fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
543        formatter.write_str(&self.detail)
544    }
545}
546
547impl From<LifecycleKind> for RuntimeLifecycleKind {
548    fn from(kind: LifecycleKind) -> Self {
549        match kind {
550            LifecycleKind::Create => Self::Create,
551            LifecycleKind::Suspend => Self::Suspend,
552            LifecycleKind::Resume => Self::Resume,
553            LifecycleKind::Restart => Self::Resume,
554            LifecycleKind::Move => Self::Move,
555            LifecycleKind::ForceStop => Self::ForceStop,
556            LifecycleKind::DestroyStopped | LifecycleKind::ArchiveStopped => Self::DestroyStopped,
557            LifecycleKind::ForceDestroy => Self::ForceDestroy,
558            // A park is shown as the sub-agent stop it is, and an unpark as
559            // the resume it is; neither needs a wire kind of its own.
560            LifecycleKind::StopSubagent | LifecycleKind::Park => Self::StopSubagent,
561            LifecycleKind::Unpark => Self::Resume,
562            LifecycleKind::Cleanup | LifecycleKind::StartupCleanup => Self::Cleanup,
563        }
564    }
565}
566
567mod capacity;
568mod close;
569mod close_workspace;
570mod create;
571mod readiness;
572use readiness::{HarnessReadinessWatch, ReadinessObservation, UnreadySession};
573mod lifecycle;
574mod resume;
575mod snapshot;
576pub(crate) mod startup_followup;
577mod state;
578mod subagent_park;
579mod support;
580mod views;
581use support::*;
582pub(crate) mod delegation;
583mod diagnostics;
584mod process;
585pub use process::*;
586mod serve;
587use serve::*;
588mod actions;
589use actions::*;
590mod guards;
591pub(crate) use guards::*;
592
593#[cfg(test)]
594mod tests;
595
596mod continuation;