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, 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) 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 workspace_resume_admission: Mutex<BTreeMap<String, Arc<tokio::sync::RwLock<()>>>>,
137 harness_readiness: Mutex<HarnessReadinessWatch>,
140 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 notices: Mutex<VecDeque<RuntimeNotice>>,
154 next_notice_id: AtomicU64,
155 quota: Mutex<QuotaBoard>,
159 capacity: std::sync::OnceLock<capacity::CapacityFeed>,
162 review_config: Arc<Mutex<mj_core::config::ReviewConfig>>,
164 review_host: TurnReviewHost,
167 wiki: crate::sessionwiki::WikiIndexer,
169}
170
171#[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 ArchiveStopped,
241 ForceDestroy,
242 StopSubagent,
246 Park,
249 Unpark,
252 Cleanup,
253 StartupCleanup,
254}
255
256impl LifecycleKind {
257 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 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
291fn 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
317fn 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#[derive(Debug, Clone, Copy, PartialEq, Eq)]
339enum CloseRoute {
340 Graceful,
342 RecoverInterrupted,
344 SettleWithoutCheckpoint,
346 DeferredCleanup,
348 Done,
350}
351
352fn 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
376fn 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#[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
410struct 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 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#[derive(Debug, Clone)]
505pub(crate) struct LifecycleFailure {
506 pub(crate) detail: String,
507 pub(crate) refusal: Option<Refusal>,
508}
509
510type 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 fn internal(detail: impl Into<String>) -> Self {
526 Self {
527 detail: detail.into(),
528 refusal: None,
529 }
530 }
531
532 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 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;