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    BranchDisposition, CheckoutDisposition, Controller, ControllerStoreGuard, SessionLaunchOptions,
34    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    revisions: RuntimeRevisions,
114    workspaces_tx: tokio::sync::watch::Sender<Vec<WorkspaceRecord>>,
115    session_manager: SessionManagerControl,
116    lifecycle: Mutex<BTreeMap<String, ActiveLifecycle>>,
117    workspace_closes: Mutex<BTreeMap<String, Arc<AtomicBool>>>,
118    /// Resume ownership can precede its durable workspace assignment.
119    workspace_resume_admission: Mutex<BTreeMap<String, Arc<tokio::sync::RwLock<()>>>>,
120    /// The bounded wait for each live session's harness to become usable.
121    /// Driven only by the daemon's background readiness sweep.
122    harness_readiness: Mutex<HarnessReadinessWatch>,
123    /// Work waiting for one session's harness to become ready: the prompts a
124    /// person typed while it started, and the hand-off a restored session
125    /// carries. One ordered queue per session, each drained by one task.
126    startup_prompts: Mutex<BTreeMap<String, StartupQueue>>,
127    close_requested: Mutex<BTreeSet<String>>,
128    controller: Mutex<Controller>,
129    controller_loader: fn() -> Result<Controller>,
130    config_mutation: tokio::sync::Mutex<()>,
131    recovery_observer: RecoveryObserver,
132    worker_upgrade_observer: WorkerUpgradeObserver,
133    /// Recent background notices, newest last, with the id of the next one.
134    /// Bounded: a surface that never attaches must not make this grow.
135    notices: Mutex<VecDeque<RuntimeNotice>>,
136    next_notice_id: AtomicU64,
137    /// What `[review]` last said, republished by the target refresher.
138    review_config: Arc<Mutex<mj_core::config::ReviewConfig>>,
139    /// Turn review runs here, in the process that owns every session, so a
140    /// review happens whether the terminal, the phone, or nobody is attached.
141    review_host: TurnReviewHost,
142    /// Publishes checkpointed sessions into the user's SessionWiki index.
143    wiki: crate::sessionwiki::WikiIndexer,
144}
145
146/// One monotonic cursor shared by daemon snapshots and their wake-up feed.
147///
148/// Allocations can come from independent UI and daemon tasks. Publishing an
149/// older allocation after a newer one must not move the watch channel
150/// backwards, so publication compares against the last visible cursor.
151#[derive(Clone)]
152struct RuntimeRevisions {
153    allocated: Arc<std::sync::atomic::AtomicU64>,
154    published: tokio::sync::watch::Sender<u64>,
155}
156
157impl RuntimeRevisions {
158    fn new(initial: u64) -> Self {
159        let (published, _) = tokio::sync::watch::channel(initial);
160        Self {
161            allocated: Arc::new(std::sync::atomic::AtomicU64::new(initial)),
162            published,
163        }
164    }
165
166    fn allocate(&self) -> u64 {
167        self.allocated.fetch_add(1, Ordering::AcqRel) + 1
168    }
169
170    fn publish(&self) -> u64 {
171        let revision = self.allocate();
172        self.publish_allocated(revision);
173        revision
174    }
175
176    fn publish_allocated(&self, revision: u64) {
177        self.published.send_if_modified(|visible| {
178            if revision > *visible {
179                *visible = revision;
180                true
181            } else {
182                false
183            }
184        });
185    }
186
187    fn notifier(&self) -> Arc<dyn Fn() + Send + Sync> {
188        let revisions = self.clone();
189        Arc::new(move || {
190            revisions.publish();
191        })
192    }
193
194    fn subscribe(&self) -> tokio::sync::watch::Receiver<u64> {
195        self.published.subscribe()
196    }
197
198    fn current(&self) -> u64 {
199        self.allocated.load(Ordering::Acquire)
200    }
201}
202
203#[derive(Debug, Clone, Copy, PartialEq, Eq)]
204enum LifecycleKind {
205    Create,
206    Suspend,
207    Resume,
208    Move,
209    ForceStop,
210    DestroyStopped,
211    /// The archive job's destruction: the same teardown as `DestroyStopped`,
212    /// with the session's git branch kept unless another branch already
213    /// contains every one of its commits. Surfaces see it as a destroy.
214    ArchiveStopped,
215    ForceDestroy,
216    Cleanup,
217}
218
219impl LifecycleKind {
220    /// What a refusal calls this operation, so a person told that a session is
221    /// busy learns which operation is holding it (#1010).
222    fn label(self) -> &'static str {
223        match self {
224            LifecycleKind::Create => "create",
225            LifecycleKind::Suspend => "suspend",
226            LifecycleKind::Resume => "resume",
227            LifecycleKind::Move => "move",
228            LifecycleKind::ForceStop => "force stop",
229            LifecycleKind::DestroyStopped => "destroy",
230            LifecycleKind::ArchiveStopped => "archive",
231            LifecycleKind::ForceDestroy => "force destroy",
232            LifecycleKind::Cleanup => "cleanup",
233        }
234    }
235}
236
237/// Whether a lifecycle has exclusive ownership of the worker target, so the
238/// session manager must stop polling it. A graceful close needs the manager's
239/// relay lease through checkpointing and sealing; once the durable state says
240/// `Destroying`, that lease has been released and target teardown is exclusive.
241fn lifecycle_owns_worker_target(kind: LifecycleKind, state: Option<SessionState>) -> bool {
242    match kind {
243        LifecycleKind::Suspend => state == Some(SessionState::Destroying),
244        LifecycleKind::Move => !matches!(
245            state,
246            Some(
247                SessionState::Running
248                    | SessionState::Disconnected
249                    | SessionState::Checkpointing
250                    | SessionState::Closing
251            )
252        ),
253        _ => true,
254    }
255}
256
257/// Whether a running lifecycle can still be cancelled. A graceful close has a
258/// point of no return: once the durable state says `Destroying`, the verified
259/// checkpoint is sealed and the record has already committed to losing its
260/// target, so stopping the teardown only strands the target. Every other
261/// lifecycle stays cancellable while it runs.
262fn lifecycle_cancellable(kind: LifecycleKind, state: Option<SessionState>) -> bool {
263    !(kind == LifecycleKind::Suspend && state == Some(SessionState::Destroying))
264}
265
266/// How a stop request has to be carried out, given the durable record.
267#[derive(Debug, Clone, Copy, PartialEq, Eq)]
268enum CloseRoute {
269    /// Run the graceful close from the start.
270    Graceful,
271    /// A previous close stopped partway; finish it from its checkpoint.
272    RecoverInterrupted,
273    /// Nothing to checkpoint: tear down whatever target is left and settle.
274    SettleWithoutCheckpoint,
275    /// Already stopped, but the target still has to be removed.
276    DeferredCleanup,
277    /// Already stopped with nothing left to do.
278    Done,
279}
280
281/// A record mid-close with a live target cannot be closed again from the start:
282/// its worker socket is gone, so a fresh checkpoint attempt only fails on
283/// connect. Recovery finishes it from the checkpoint the first close verified.
284fn close_route(session: Option<&SessionRecord>) -> CloseRoute {
285    let Some(session) = session else {
286        return CloseRoute::Graceful;
287    };
288    if crate::pollers::is_interrupted_close(session) {
289        CloseRoute::RecoverInterrupted
290    } else if session.state == SessionState::Stopped {
291        if session.target.is_some() {
292            CloseRoute::DeferredCleanup
293        } else {
294            CloseRoute::Done
295        }
296    } else if crate::controller::has_nothing_to_checkpoint(session) {
297        CloseRoute::SettleWithoutCheckpoint
298    } else {
299        CloseRoute::Graceful
300    }
301}
302
303/// The durable state of one record as the locked controller holds it.
304fn durable_session_state(controller: &Controller, session_id: &str) -> Option<SessionState> {
305    controller
306        .state
307        .sessions
308        .get(session_id)
309        .map(|session| session.state)
310}
311
312/// One piece of work that waits for a starting session's harness.
313///
314/// Both kinds are ordered against each other on purpose: a restored session's
315/// hand-off is the hidden context its first prompt reads, so it has to be
316/// installed before any queued prompt is submitted.
317pub(crate) enum StartupStep {
318    InstallHandoff(Box<mj_core::archive::CanonicalSessionSnapshot>),
319    Prompt {
320        text: String,
321        inherited_draft: Option<String>,
322    },
323}
324
325/// The steps waiting for one session, and the task draining them.
326///
327/// The entry exists only while a drain task owns it. `in_flight` marks a step
328/// that has been popped and is running, so the queue is never treated as empty
329/// while its last step is still being carried out.
330struct StartupQueue {
331    pending: VecDeque<StartupStep>,
332    in_flight: bool,
333    cancel: CancellationToken,
334    task: Option<tokio::task::JoinHandle<()>>,
335}
336
337struct ActiveLifecycle {
338    operation_id: String,
339    create_control: Option<CreateSessionControl>,
340    kind: LifecycleKind,
341    cancelled: Arc<AtomicBool>,
342    started_at_epoch_seconds: u64,
343    active_stages: BTreeMap<ProvisionStage, (usize, u64)>,
344    /// The workspace a resume is claiming before its durable record changes.
345    /// Workspace deletion consults this so it cannot race the claim.
346    resume_workspace_id: Option<String>,
347    resume_destination: Option<(String, String)>,
348    notice: Option<String>,
349    request_key: Option<String>,
350    _move_guard: Option<MoveMutationGuard>,
351    move_source_closed: bool,
352    result: LifecycleWatch,
353}
354
355impl ActiveLifecycle {
356    fn is_visible(&self) -> bool {
357        let result = self.result.borrow();
358        result.is_none()
359            || matches!(
360                result.as_ref(),
361                Some(Ok(DaemonLifecycleResult::DeferredCleanup))
362            )
363    }
364
365    fn request_cancel(&self) -> bool {
366        if let Some(control) = &self.create_control {
367            control.request_cancel()
368        } else {
369            !self.cancelled.swap(true, Ordering::AcqRel)
370        }
371    }
372
373    fn is_cancellable(&self) -> bool {
374        self.result.borrow().is_none()
375            && self.create_control.as_ref().map_or_else(
376                || !self.cancelled.load(Ordering::Acquire),
377                CreateSessionControl::is_cancellable,
378            )
379    }
380}
381
382#[derive(Debug, Clone)]
383enum DaemonLifecycleResult {
384    Done,
385    DeferredCleanup,
386    Move(MoveOutcome),
387}
388
389/// How one lifecycle operation ended when it failed.
390///
391/// The result is broadcast to every waiter, which is why it cannot simply be
392/// the `anyhow::Error`: that is not clonable. Keeping the refusal beside the
393/// text is what lets a reason written for the caller survive the crossing; a
394/// failure rebuilt from a string alone would arrive as an internal fault.
395#[derive(Debug, Clone)]
396pub(crate) struct LifecycleFailure {
397    pub(crate) detail: String,
398    pub(crate) refusal: Option<Refusal>,
399}
400
401/// One lifecycle operation's outcome, and the channel every waiter reads it
402/// from. `None` means the operation is still running.
403type LifecycleResult = std::result::Result<DaemonLifecycleResult, LifecycleFailure>;
404type LifecycleWatch = tokio::sync::watch::Receiver<Option<LifecycleResult>>;
405
406impl LifecycleFailure {
407    fn of(error: &anyhow::Error) -> Self {
408        Self {
409            detail: format!("{error:#}"),
410            refusal: Refusal::of(error),
411        }
412    }
413
414    /// A failure with no reason written for a caller, such as a task that died
415    /// before the operation could say anything about itself.
416    fn internal(detail: impl Into<String>) -> Self {
417        Self {
418            detail: detail.into(),
419            refusal: None,
420        }
421    }
422
423    /// Rebuild the error a waiter sees, with the refusal still attached.
424    fn into_error(self) -> anyhow::Error {
425        match self.refusal {
426            Some(refusal) => anyhow::Error::new(refusal).context(self.detail),
427            None => anyhow::Error::msg(self.detail),
428        }
429    }
430}
431
432impl std::fmt::Display for LifecycleFailure {
433    fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
434        formatter.write_str(&self.detail)
435    }
436}
437
438impl From<LifecycleKind> for RuntimeLifecycleKind {
439    fn from(kind: LifecycleKind) -> Self {
440        match kind {
441            LifecycleKind::Create => Self::Create,
442            LifecycleKind::Suspend => Self::Suspend,
443            LifecycleKind::Resume => Self::Resume,
444            LifecycleKind::Move => Self::Move,
445            LifecycleKind::ForceStop => Self::ForceStop,
446            LifecycleKind::DestroyStopped | LifecycleKind::ArchiveStopped => Self::DestroyStopped,
447            LifecycleKind::ForceDestroy => Self::ForceDestroy,
448            LifecycleKind::Cleanup => Self::Cleanup,
449        }
450    }
451}
452
453mod close;
454mod close_workspace;
455mod create;
456mod readiness;
457use readiness::{HarnessReadinessWatch, ReadinessObservation, UnreadySession};
458mod lifecycle;
459mod resume;
460mod snapshot;
461mod state;
462mod support;
463mod views;
464use support::*;
465mod process;
466pub use process::*;
467mod serve;
468use serve::*;
469mod actions;
470use actions::*;
471mod guards;
472pub(crate) use guards::*;
473
474#[cfg(test)]
475mod tests;
476
477mod continuation;