1mod 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
57const SHUTDOWN_FORCE_EXIT_TIMEOUT: Duration = Duration::from_secs(10);
67
68const FORCE_DESTROY_PREEMPT_TIMEOUT: Duration = Duration::from_secs(8);
74
75#[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) wait_prompt_backend:
129 std::sync::OnceLock<std::sync::Weak<crate::server_runtime::api::ApiBackend>>,
130 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 workspace_resume_admission: Mutex<BTreeMap<String, Arc<tokio::sync::RwLock<()>>>>,
145 harness_readiness: Mutex<HarnessReadinessWatch>,
148 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 notices: Mutex<VecDeque<RuntimeNotice>>,
162 next_notice_id: AtomicU64,
163 quota: Mutex<QuotaBoard>,
167 capacity: std::sync::OnceLock<capacity::CapacityFeed>,
170 review_config: Arc<Mutex<mj_core::config::ReviewConfig>>,
172 review_host: TurnReviewHost,
175 wiki: crate::sessionwiki::WikiIndexer,
177}
178
179#[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 ArchiveStopped,
249 ForceDestroy,
250 StopSubagent,
254 Park,
257 Unpark,
260 Cleanup,
261 StartupCleanup,
262}
263
264impl LifecycleKind {
265 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 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
299fn 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
325fn 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#[derive(Debug, Clone, Copy, PartialEq, Eq)]
347enum CloseRoute {
348 Graceful,
350 RecoverInterrupted,
352 SettleWithoutCheckpoint,
354 DeferredCleanup,
356 Done,
358}
359
360fn 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
384fn 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#[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
418struct 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 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#[derive(Debug, Clone)]
513pub(crate) struct LifecycleFailure {
514 pub(crate) detail: String,
515 pub(crate) refusal: Option<Refusal>,
516}
517
518type 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 fn internal(detail: impl Into<String>) -> Self {
534 Self {
535 detail: detail.into(),
536 refusal: None,
537 }
538 }
539
540 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 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;