1mod feed;
4mod owner;
5mod record_index;
6mod session_move;
7use crate::controller::move_session::{
8 MoveMutationGuard, MoveOutcome, MovePreparation, MoveSelection, MoveSessionRequest,
9};
10pub use mj_client::daemon::*;
11use owner::RuntimeStateOwner;
12
13use std::collections::{BTreeMap, BTreeSet, VecDeque};
14use std::fs::{self, OpenOptions};
15use std::io::Write;
16use std::net::{IpAddr, Ipv4Addr, SocketAddr};
17use std::path::{Path, PathBuf};
18use std::sync::atomic::{AtomicBool, AtomicU8, AtomicU64, Ordering};
19use std::sync::{Arc, Mutex, PoisonError};
20use std::time::{Duration, SystemTime, UNIX_EPOCH};
21
22use crate::database::StoreSchemaMismatch;
23use crate::recovery_gate::RecoveryObserver;
24use crate::targets::{
25 CancellableProcessExecutor, CommandExecutor, CommandOutput, CommandSpec, ProcessExecutor,
26 ProvisionStage, ProvisionStageGuard,
27};
28use agent_client_protocol::schema::v1::{ContentBlock, TextContent};
29use anyhow::{Context, Result, anyhow, bail, ensure};
30use mj_core::config::Config;
31use mj_core::refusal::Refusal;
32use mj_core::relay::RelayCommand;
33use mj_core::state::{RecoveryObservation, SessionRecord, SessionState};
34
35use crate::controller::{
36 BeforeClose, BranchDisposition, CheckoutDisposition, Controller, ControllerStoreGuard,
37 SessionLaunchOptions, SessionResumeOptions,
38};
39use crate::review_host::TurnReviewHost;
40use crate::session_manager::{
41 ManagedSessionView, RemoteSessionPublisher, RemoteSessionRequest, SessionManagerChannels,
42 SessionManagerControl, ViewError, new_command_id, spawn_remote_session_manager,
43};
44#[cfg(test)]
45use crate::session_manager::{RelaySessionTarget, RemoteSessionRequests, SessionManagerShutdown};
46use crate::worker_upgrade::{WorkerUpgradeObservation, WorkerUpgradeObserver};
47use mj_core::workspace::WorkspaceRecord;
48use tokio::net::{TcpListener, TcpStream};
49use tokio_util::sync::CancellationToken;
50
51use crate::pollers::{
52 dashboard_worker_targets, interrupted_destroy_session_ids, interrupted_suspend_session_ids,
53 reserve_recovery_or_cancel, spawn_image_refresher, unowned_interrupted_lifecycles,
54};
55
56const SHUTDOWN_FORCE_EXIT_TIMEOUT: Duration = Duration::from_secs(10);
66
67const FORCE_DESTROY_PREEMPT_TIMEOUT: Duration = Duration::from_secs(8);
73
74#[derive(Clone, Default)]
76pub struct CreateSessionControl {
77 state: Arc<AtomicU8>,
78 pub cancelled: Arc<AtomicBool>,
79}
80
81impl CreateSessionControl {
82 pub fn request_cancel(&self) -> bool {
83 let accepted = self
84 .state
85 .compare_exchange(0, 1, Ordering::AcqRel, Ordering::Acquire)
86 .is_ok();
87 if accepted {
88 self.cancelled.store(true, Ordering::Release);
89 }
90 accepted
91 }
92
93 pub fn grant_commit(&self) -> bool {
94 self.state
95 .compare_exchange(0, 2, Ordering::AcqRel, Ordering::Acquire)
96 .is_ok()
97 }
98
99 fn is_cancellable(&self) -> bool {
100 self.state.load(Ordering::Acquire) == 0
101 }
102}
103
104#[derive(Debug, Clone)]
105struct Attachment {
106 pid: u32,
107}
108
109#[derive(Default)]
110struct QuotaBoard {
111 snapshot: mj_client::quota::QuotaSnapshot,
112 refresh: Option<tokio::sync::mpsc::Sender<()>>,
113}
114
115pub struct RuntimeState {
116 attachments: Mutex<BTreeMap<String, Attachment>>,
117 phone_status: Mutex<WebViewerStatus>,
118 pub web_viewer: crate::web_viewer::ViewerControl,
119 ever_attached: AtomicBool,
120 revisions: RuntimeRevisions,
121 workspaces_tx: tokio::sync::watch::Sender<Vec<WorkspaceRecord>>,
122 workspace_refresh: tokio::sync::Mutex<()>,
123 session_manager: SessionManagerControl,
124 owner: Mutex<RuntimeStateOwner>,
125 feed: Mutex<feed::RuntimeHistory>,
126 committed: Option<
127 tokio::sync::watch::Receiver<
128 std::result::Result<crate::database::CommittedState, Arc<str>>,
129 >,
130 >,
131 workspace_closes: Mutex<BTreeMap<String, Arc<AtomicBool>>>,
132 workspace_resume_admission: Mutex<BTreeMap<String, Arc<tokio::sync::RwLock<()>>>>,
134 harness_readiness: Mutex<HarnessReadinessWatch>,
137 startup_prompts: Mutex<BTreeMap<String, StartupQueue>>,
141 startup_enqueue: tokio::sync::Mutex<()>,
142 controller_loader: fn() -> Result<Controller>,
143 config_mutation: tokio::sync::Mutex<()>,
144 projects: Arc<crate::project_catalog::Catalog>,
145 profile_catalog:
146 std::sync::OnceLock<Arc<crate::server_runtime::profile_catalog::ProfileCatalog>>,
147 recovery_observer: RecoveryObserver,
148 worker_upgrade_observer: WorkerUpgradeObserver,
149 notices: Mutex<VecDeque<RuntimeNotice>>,
152 next_notice_id: AtomicU64,
153 quota: Mutex<QuotaBoard>,
157 review_config: Arc<Mutex<mj_core::config::ReviewConfig>>,
159 review_host: TurnReviewHost,
162 wiki: crate::sessionwiki::WikiIndexer,
164}
165
166#[derive(Clone)]
172struct RuntimeRevisions {
173 allocated: Arc<std::sync::atomic::AtomicU64>,
174 published: tokio::sync::watch::Sender<u64>,
175}
176
177impl RuntimeRevisions {
178 fn new(initial: u64) -> Self {
179 let (published, _) = tokio::sync::watch::channel(initial);
180 Self {
181 allocated: Arc::new(std::sync::atomic::AtomicU64::new(initial)),
182 published,
183 }
184 }
185
186 fn allocate(&self) -> u64 {
187 self.allocated.fetch_add(1, Ordering::AcqRel) + 1
188 }
189
190 fn publish(&self) -> u64 {
191 let revision = self.allocate();
192 self.publish_allocated(revision);
193 revision
194 }
195
196 fn publish_allocated(&self, revision: u64) {
197 self.published.send_if_modified(|visible| {
198 if revision > *visible {
199 *visible = revision;
200 true
201 } else {
202 false
203 }
204 });
205 }
206
207 fn notifier(&self) -> Arc<dyn Fn() + Send + Sync> {
208 let revisions = self.clone();
209 Arc::new(move || {
210 revisions.publish();
211 })
212 }
213
214 fn subscribe(&self) -> tokio::sync::watch::Receiver<u64> {
215 self.published.subscribe()
216 }
217
218 fn current(&self) -> u64 {
219 self.allocated.load(Ordering::Acquire)
220 }
221}
222
223#[derive(Debug, Clone, Copy, PartialEq, Eq)]
224enum LifecycleKind {
225 Create,
226 Suspend,
227 Resume,
228 Move,
229 ForceStop,
230 DestroyStopped,
231 ArchiveStopped,
235 ForceDestroy,
236 StopSubagent,
240 Park,
243 Unpark,
246 Cleanup,
247 StartupCleanup,
248}
249
250impl LifecycleKind {
251 fn is_teardown(self) -> bool {
254 matches!(
255 self,
256 LifecycleKind::ForceDestroy
257 | LifecycleKind::DestroyStopped
258 | LifecycleKind::ArchiveStopped
259 | LifecycleKind::StopSubagent
260 )
261 }
262
263 fn label(self) -> &'static str {
266 match self {
267 LifecycleKind::Create => "create",
268 LifecycleKind::Suspend => "suspend",
269 LifecycleKind::Resume => "resume",
270 LifecycleKind::Move => "move",
271 LifecycleKind::ForceStop => "force stop",
272 LifecycleKind::DestroyStopped => "destroy",
273 LifecycleKind::ArchiveStopped => "archive",
274 LifecycleKind::ForceDestroy => "force destroy",
275 LifecycleKind::StopSubagent => "stop",
276 LifecycleKind::Park => "park",
277 LifecycleKind::Unpark => "restart",
278 LifecycleKind::Cleanup => "cleanup",
279 LifecycleKind::StartupCleanup => "failed startup cleanup",
280 }
281 }
282}
283
284fn lifecycle_owns_worker_target(kind: LifecycleKind, state: Option<SessionState>) -> bool {
294 match kind {
295 LifecycleKind::Suspend => state == Some(SessionState::Destroying),
296 LifecycleKind::Park => false,
297 LifecycleKind::Move => !matches!(
298 state,
299 Some(
300 SessionState::Running
301 | SessionState::Disconnected
302 | SessionState::Checkpointing
303 | SessionState::Closing
304 )
305 ),
306 _ => true,
307 }
308}
309
310fn lifecycle_cancellable(kind: LifecycleKind, state: Option<SessionState>) -> bool {
320 !(state == Some(SessionState::StartupCleanup)
321 || kind == LifecycleKind::StartupCleanup
322 || kind == LifecycleKind::Suspend && state == Some(SessionState::Destroying)
323 || matches!(
324 kind,
325 LifecycleKind::StopSubagent | LifecycleKind::Park | LifecycleKind::Unpark
326 ))
327}
328
329#[derive(Debug, Clone, Copy, PartialEq, Eq)]
331enum CloseRoute {
332 Graceful,
334 RecoverInterrupted,
336 SettleWithoutCheckpoint,
338 DeferredCleanup,
340 Done,
342}
343
344fn close_route(session: Option<&SessionRecord>, subagent: bool) -> CloseRoute {
350 let Some(session) = session else {
351 return CloseRoute::Graceful;
352 };
353 if crate::pollers::is_interrupted_close(session) {
354 CloseRoute::RecoverInterrupted
355 } else if session.state == SessionState::Stopped {
356 if session.target.is_some() {
357 CloseRoute::DeferredCleanup
358 } else {
359 CloseRoute::Done
360 }
361 } else if crate::controller::has_nothing_to_checkpoint(session, subagent) {
362 CloseRoute::SettleWithoutCheckpoint
363 } else {
364 CloseRoute::Graceful
365 }
366}
367
368fn durable_session_state(controller: &Controller, session_id: &str) -> Option<SessionState> {
370 controller
371 .state
372 .sessions
373 .get(session_id)
374 .map(|session| session.state)
375}
376
377#[derive(serde::Serialize, serde::Deserialize)]
383pub(crate) enum StartupStep {
384 InstallHandoff(Box<mj_core::archive::CanonicalSessionSnapshot>),
385 PreparedHandoff {
386 text: String,
387 },
388 Prompt {
389 text: String,
390 inherited_draft: Option<String>,
391 },
392 Configure {
393 key: String,
394 value: String,
395 optional: bool,
396 },
397 ApiPrompt {
398 text: String,
399 },
400}
401
402struct StartupQueue {
404 identity: Arc<()>,
405 last_error: Option<String>,
406 cancel: CancellationToken,
407 task: Option<tokio::task::JoinHandle<()>>,
408}
409
410struct ActiveLifecycle {
411 upgrade_work: Arc<Mutex<Option<crate::upgrade::Work>>>,
412 phase: LifecyclePhase,
413 operation_id: String,
414 create_control: Option<CreateSessionControl>,
415 kind: LifecycleKind,
416 cancelled: Arc<AtomicBool>,
417 started_at_epoch_seconds: u64,
418 active_stages: BTreeMap<ProvisionStage, (usize, u64)>,
419 resume_workspace_id: Option<String>,
422 resume_destination: Option<(String, String)>,
423 notice: Option<String>,
424 request_key: Option<String>,
425 _move_guard: Option<MoveMutationGuard>,
426 result: LifecycleWatch,
427}
428
429enum LifecyclePhase {
430 Executing,
431 MovingDestination,
432 Cancelling,
433 CancellingMoveDestination,
434 Completed(LifecycleResult),
435}
436
437impl ActiveLifecycle {
438 fn is_running(&self) -> bool {
439 !matches!(self.phase, LifecyclePhase::Completed(_))
440 }
441
442 fn is_visible(&self) -> bool {
443 self.is_running()
444 || matches!(
445 self.phase,
446 LifecyclePhase::Completed(Ok(DaemonLifecycleResult::DeferredCleanup))
447 )
448 }
449
450 fn request_cancel(&mut self) -> bool {
451 if !matches!(
452 self.phase,
453 LifecyclePhase::Executing | LifecyclePhase::MovingDestination
454 ) {
455 return false;
456 }
457 let accepted = if let Some(control) = &self.create_control {
458 control.request_cancel()
459 } else {
460 !self.cancelled.swap(true, Ordering::AcqRel)
461 };
462 if accepted {
463 self.phase = match self.phase {
464 LifecyclePhase::MovingDestination => LifecyclePhase::CancellingMoveDestination,
465 _ => LifecyclePhase::Cancelling,
466 };
467 }
468 accepted
469 }
470
471 fn is_cancellable(&self) -> bool {
472 matches!(
473 self.phase,
474 LifecyclePhase::Executing | LifecyclePhase::MovingDestination
475 ) && self.create_control.as_ref().map_or_else(
476 || !self.cancelled.load(Ordering::Acquire),
477 CreateSessionControl::is_cancellable,
478 )
479 }
480}
481
482#[derive(Debug, Clone)]
483enum DaemonLifecycleResult {
484 Done,
485 DeferredCleanup,
486 Move(MoveOutcome),
487 Park(crate::controller::ParkOutcome),
488}
489
490#[derive(Debug, Clone)]
497pub(crate) struct LifecycleFailure {
498 pub(crate) detail: String,
499 pub(crate) refusal: Option<Refusal>,
500}
501
502type LifecycleResult = std::result::Result<DaemonLifecycleResult, LifecycleFailure>;
505type LifecycleWatch = tokio::sync::watch::Receiver<Option<LifecycleResult>>;
506
507impl LifecycleFailure {
508 fn of(error: &anyhow::Error) -> Self {
509 Self {
510 detail: format!("{error:#}"),
511 refusal: Refusal::of(error),
512 }
513 }
514
515 fn internal(detail: impl Into<String>) -> Self {
518 Self {
519 detail: detail.into(),
520 refusal: None,
521 }
522 }
523
524 fn into_error(self) -> anyhow::Error {
526 match self.refusal {
527 Some(refusal) => anyhow::Error::new(refusal).context(self.detail),
528 None => anyhow::Error::msg(self.detail),
529 }
530 }
531}
532
533impl std::fmt::Display for LifecycleFailure {
534 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
535 formatter.write_str(&self.detail)
536 }
537}
538
539impl From<LifecycleKind> for RuntimeLifecycleKind {
540 fn from(kind: LifecycleKind) -> Self {
541 match kind {
542 LifecycleKind::Create => Self::Create,
543 LifecycleKind::Suspend => Self::Suspend,
544 LifecycleKind::Resume => Self::Resume,
545 LifecycleKind::Move => Self::Move,
546 LifecycleKind::ForceStop => Self::ForceStop,
547 LifecycleKind::DestroyStopped | LifecycleKind::ArchiveStopped => Self::DestroyStopped,
548 LifecycleKind::ForceDestroy => Self::ForceDestroy,
549 LifecycleKind::StopSubagent | LifecycleKind::Park => Self::StopSubagent,
552 LifecycleKind::Unpark => Self::Resume,
553 LifecycleKind::Cleanup | LifecycleKind::StartupCleanup => Self::Cleanup,
554 }
555 }
556}
557
558mod close;
559mod close_workspace;
560mod create;
561mod readiness;
562use readiness::{HarnessReadinessWatch, ReadinessObservation, UnreadySession};
563mod lifecycle;
564mod resume;
565mod snapshot;
566pub(crate) mod startup_followup;
567mod state;
568mod subagent_park;
569mod support;
570mod views;
571use support::*;
572pub(crate) mod delegation;
573mod diagnostics;
574mod process;
575pub use process::*;
576mod serve;
577use serve::*;
578mod actions;
579use actions::*;
580mod guards;
581pub(crate) use guards::*;
582
583#[cfg(test)]
584mod tests;
585
586mod continuation;