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