mod session_move;
use crate::controller::move_session::{
MoveMutationGuard, MoveOutcome, MovePreparation, MoveSelection, MoveSessionRequest,
};
pub use mj_client::daemon::*;
use std::collections::{BTreeMap, BTreeSet, VecDeque};
use std::fs::{self, OpenOptions};
use std::io::Write;
use std::net::{IpAddr, Ipv4Addr, SocketAddr};
use std::path::Path;
use std::sync::atomic::{AtomicBool, AtomicU8, AtomicU64, Ordering};
use std::sync::{Arc, Mutex, PoisonError};
use std::time::{Duration, SystemTime, UNIX_EPOCH};
use crate::database::StoreSchemaMismatch;
use crate::recovery_gate::RecoveryObserver;
use crate::targets::{
CancellableProcessExecutor, CommandExecutor, CommandOutput, CommandSpec, ProcessExecutor,
ProvisionStage, ProvisionStageGuard,
};
use anyhow::{Context, Result, anyhow, bail, ensure};
use mj_core::config::Config;
use mj_core::relay::RelayCommand;
use mj_core::state::{RecoveryObservation, SessionRecord, SessionState};
use mj_core::subagent::SubagentRecord;
use crate::controller::{
BranchDisposition, Controller, ControllerStoreGuard, SessionLaunchOptions, SessionResumeOptions,
};
use crate::review_host::TurnReviewHost;
use crate::session_manager::{
ManagedSessionView, RemoteSessionPublisher, RemoteSessionRequest, SessionManagerChannels,
SessionManagerControl, ViewError, new_command_id, spawn_remote_session_manager,
spawn_session_manager,
};
#[cfg(test)]
use crate::session_manager::{RelaySessionTarget, RemoteSessionRequests, SessionManagerShutdown};
use crate::worker_upgrade::{WorkerUpgradeObservation, WorkerUpgradeObserver};
use mj_core::workspace::WorkspaceRecord;
use tokio::net::{TcpListener, TcpStream};
use tokio_util::sync::CancellationToken;
use crate::pollers::{
dashboard_worker_targets, dashboard_worker_targets_excluding, interrupted_close_session_ids,
reserve_recovery_or_cancel, spawn_image_refresher, spawn_interrupted_close_recovery,
};
const SHUTDOWN_FORCE_EXIT_TIMEOUT: Duration = Duration::from_secs(10);
const FORCE_DESTROY_PREEMPT_TIMEOUT: Duration = Duration::from_secs(8);
#[derive(Clone, Default)]
pub struct CreateSessionControl {
state: Arc<AtomicU8>,
pub cancelled: Arc<AtomicBool>,
}
impl CreateSessionControl {
pub fn request_cancel(&self) -> bool {
let accepted = self
.state
.compare_exchange(0, 1, Ordering::AcqRel, Ordering::Acquire)
.is_ok();
if accepted {
self.cancelled.store(true, Ordering::Release);
}
accepted
}
pub fn grant_commit(&self) -> bool {
self.state
.compare_exchange(0, 2, Ordering::AcqRel, Ordering::Acquire)
.is_ok()
}
fn is_cancellable(&self) -> bool {
self.state.load(Ordering::Acquire) == 0
}
}
#[derive(Debug, Clone)]
struct Attachment {
pid: u32,
}
pub struct RuntimeState {
attachments: Mutex<BTreeMap<String, Attachment>>,
phone_status: Mutex<WebViewerStatus>,
pub web_viewer: crate::web_viewer::ViewerControl,
ever_attached: AtomicBool,
sessions: Mutex<BTreeMap<String, RuntimeSessionView>>,
revisions: RuntimeRevisions,
workspaces_tx: tokio::sync::watch::Sender<Vec<WorkspaceRecord>>,
session_manager: SessionManagerControl,
lifecycle: Mutex<BTreeMap<String, ActiveLifecycle>>,
close_requested: Mutex<BTreeSet<String>>,
controller: Mutex<Controller>,
controller_loader: fn() -> Result<Controller>,
config_mutation: tokio::sync::Mutex<()>,
recovery_observer: RecoveryObserver,
worker_upgrade_observer: WorkerUpgradeObserver,
notices: Mutex<VecDeque<RuntimeNotice>>,
next_notice_id: AtomicU64,
review_config: Arc<Mutex<mj_core::config::ReviewConfig>>,
review_host: TurnReviewHost,
wiki: crate::sessionwiki::WikiIndexer,
}
#[derive(Clone)]
struct RuntimeRevisions {
allocated: Arc<std::sync::atomic::AtomicU64>,
published: tokio::sync::watch::Sender<u64>,
}
impl RuntimeRevisions {
fn new(initial: u64) -> Self {
let (published, _) = tokio::sync::watch::channel(initial);
Self {
allocated: Arc::new(std::sync::atomic::AtomicU64::new(initial)),
published,
}
}
fn allocate(&self) -> u64 {
self.allocated.fetch_add(1, Ordering::AcqRel) + 1
}
fn publish(&self) -> u64 {
let revision = self.allocate();
self.publish_allocated(revision);
revision
}
fn publish_allocated(&self, revision: u64) {
self.published.send_if_modified(|visible| {
if revision > *visible {
*visible = revision;
true
} else {
false
}
});
}
fn notifier(&self) -> Arc<dyn Fn() + Send + Sync> {
let revisions = self.clone();
Arc::new(move || {
revisions.publish();
})
}
fn subscribe(&self) -> tokio::sync::watch::Receiver<u64> {
self.published.subscribe()
}
fn current(&self) -> u64 {
self.allocated.load(Ordering::Acquire)
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum LifecycleKind {
Create,
Close,
Resume,
Move,
ForceStop,
DestroyStopped,
ArchiveStopped,
ForceDestroy,
Cleanup,
}
fn lifecycle_owns_worker_target(kind: LifecycleKind, state: Option<SessionState>) -> bool {
match kind {
LifecycleKind::Close => state == Some(SessionState::Destroying),
LifecycleKind::Move => !matches!(
state,
Some(
SessionState::Running
| SessionState::Disconnected
| SessionState::Checkpointing
| SessionState::Closing
)
),
_ => true,
}
}
fn lifecycle_cancellable(kind: LifecycleKind, state: Option<SessionState>) -> bool {
!(kind == LifecycleKind::Close && state == Some(SessionState::Destroying))
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum CloseRoute {
Graceful,
RecoverInterrupted,
DeferredCleanup,
Done,
}
fn close_route(session: Option<&SessionRecord>) -> CloseRoute {
let Some(session) = session else {
return CloseRoute::Graceful;
};
if crate::pollers::is_interrupted_close(session) {
CloseRoute::RecoverInterrupted
} else if session.state == SessionState::Stopped {
if session.target.is_some() {
CloseRoute::DeferredCleanup
} else {
CloseRoute::Done
}
} else {
CloseRoute::Graceful
}
}
fn durable_session_state(controller: &Controller, session_id: &str) -> Option<SessionState> {
controller
.state
.sessions
.get(session_id)
.map(|session| session.state)
}
struct ActiveLifecycle {
operation_id: String,
create_control: Option<CreateSessionControl>,
kind: LifecycleKind,
cancelled: Arc<AtomicBool>,
started_at_epoch_seconds: u64,
active_stages: BTreeMap<ProvisionStage, (usize, u64)>,
resume_workspace_id: Option<String>,
resume_destination: Option<(String, String)>,
notice: Option<String>,
request_key: Option<String>,
_move_guard: Option<MoveMutationGuard>,
move_source_closed: bool,
result:
tokio::sync::watch::Receiver<Option<std::result::Result<DaemonLifecycleResult, String>>>,
}
impl ActiveLifecycle {
fn is_visible(&self) -> bool {
let result = self.result.borrow();
result.is_none()
|| matches!(
result.as_ref(),
Some(Ok(DaemonLifecycleResult::DeferredCleanup))
)
}
fn request_cancel(&self) -> bool {
if let Some(control) = &self.create_control {
control.request_cancel()
} else {
!self.cancelled.swap(true, Ordering::AcqRel)
}
}
fn is_cancellable(&self) -> bool {
self.result.borrow().is_none()
&& self.create_control.as_ref().map_or_else(
|| !self.cancelled.load(Ordering::Acquire),
CreateSessionControl::is_cancellable,
)
}
}
#[derive(Debug, Clone)]
enum DaemonLifecycleResult {
Done,
DeferredCleanup,
Move(MoveOutcome),
}
impl From<LifecycleKind> for RuntimeLifecycleKind {
fn from(kind: LifecycleKind) -> Self {
match kind {
LifecycleKind::Create => Self::Create,
LifecycleKind::Close => Self::Close,
LifecycleKind::Resume => Self::Resume,
LifecycleKind::Move => Self::Move,
LifecycleKind::ForceStop => Self::ForceStop,
LifecycleKind::DestroyStopped | LifecycleKind::ArchiveStopped => Self::DestroyStopped,
LifecycleKind::ForceDestroy => Self::ForceDestroy,
LifecycleKind::Cleanup => Self::Cleanup,
}
}
}
mod close;
mod create;
mod lifecycle;
mod resume;
mod snapshot;
mod state;
mod support;
mod views;
use support::*;
mod process;
pub use process::*;
mod serve;
use serve::*;
mod actions;
use actions::*;
mod guards;
pub(crate) use guards::*;
#[cfg(test)]
mod tests;