1mod 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
54const SHUTDOWN_FORCE_EXIT_TIMEOUT: Duration = Duration::from_secs(10);
64
65const FORCE_DESTROY_PREEMPT_TIMEOUT: Duration = Duration::from_secs(8);
71
72#[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 workspace_resume_admission: Mutex<BTreeMap<String, Arc<tokio::sync::RwLock<()>>>>,
122 harness_readiness: Mutex<HarnessReadinessWatch>,
125 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 notices: Mutex<VecDeque<RuntimeNotice>>,
138 next_notice_id: AtomicU64,
139 review_config: Arc<Mutex<mj_core::config::ReviewConfig>>,
141 review_host: TurnReviewHost,
144 wiki: crate::sessionwiki::WikiIndexer,
146}
147
148#[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 ArchiveStopped,
217 ForceDestroy,
218 StopSubagent,
222 Park,
225 Unpark,
228 Cleanup,
229}
230
231impl LifecycleKind {
232 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
252fn 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
278fn 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#[derive(Debug, Clone, Copy, PartialEq, Eq)]
297enum CloseRoute {
298 Graceful,
300 RecoverInterrupted,
302 SettleWithoutCheckpoint,
304 DeferredCleanup,
306 Done,
308}
309
310fn 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
332fn 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
341pub(crate) enum StartupStep {
347 InstallHandoff(Box<mj_core::archive::CanonicalSessionSnapshot>),
348 Prompt {
349 text: String,
350 inherited_draft: Option<String>,
351 },
352}
353
354struct 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 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#[derive(Debug, Clone)]
425pub(crate) struct LifecycleFailure {
426 pub(crate) detail: String,
427 pub(crate) refusal: Option<Refusal>,
428}
429
430type 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 fn internal(detail: impl Into<String>) -> Self {
446 Self {
447 detail: detail.into(),
448 refusal: None,
449 }
450 }
451
452 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 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;