use super::*;
impl RuntimeState {
pub(super) fn cancel_lifecycle(&self, session_id: &str) -> Result<()> {
let mut controller_owner = self.owner();
let durable = durable_session_state(controller_owner.controller(), session_id);
let active = controller_owner
.lifecycle
.get_mut(session_id)
.with_context(|| {
format!("no lifecycle operation is running for session {session_id}")
})?;
ensure!(
lifecycle_cancellable(active.kind, durable),
"stop of {session_id} has passed its verified checkpoint and is removing the target; \
it cannot be cancelled"
);
ensure!(
active.request_cancel(),
"lifecycle operation is no longer cancellable"
);
drop(controller_owner);
self.publish_revision();
Ok(())
}
pub(super) async fn cancel_and_wait_lifecycles(&self) -> Result<()> {
let mut pending = {
let mut lifecycle_owner = self.owner();
let lifecycle = &mut lifecycle_owner.lifecycle;
lifecycle
.iter_mut()
.filter(|(_, active)| active.is_running())
.map(|(session_id, active)| {
if active.kind != LifecycleKind::Cleanup {
active.request_cancel();
}
let stage = active
.active_stages
.keys()
.next_back()
.map(|stage| stage.label())
.unwrap_or_else(|| "container cleanup".to_owned());
(
session_id.clone(),
active.operation_id.clone(),
active.kind,
stage,
active.result.clone(),
)
})
.collect::<Vec<_>>()
};
let cleanup_deadline = tokio::time::Instant::now() + Duration::from_secs(8);
for (session_id, operation_id, kind, stage, result) in &mut pending {
if *kind != LifecycleKind::Cleanup || result.borrow().is_some() {
continue;
}
tracing::info!(%session_id, %stage, "daemon shutdown is waiting for deferred cleanup");
self.set_lifecycle_notice(
session_id,
operation_id,
&format!("Daemon shutdown is waiting for {stage}"),
);
let finished = tokio::time::timeout_at(cleanup_deadline, async {
while result.borrow_and_update().is_none() {
result.changed().await.with_context(|| {
format!("cleanup owner stopped without a result for session {session_id}")
})?;
}
Ok::<_, anyhow::Error>(())
})
.await;
match finished {
Ok(result) => result?,
Err(_) => {
tracing::warn!(%session_id, %stage, "deferred cleanup exceeded the daemon shutdown drain deadline");
self.cancel_operation(session_id, operation_id);
}
}
}
let join_deadline = tokio::time::Instant::now() + Duration::from_secs(1);
for (session_id, operation_id, _, stage, mut result) in pending {
self.cancel_operation(&session_id, &operation_id);
let joined = tokio::time::timeout_at(join_deadline, async {
while result.borrow_and_update().is_none() {
result.changed().await.with_context(|| {
format!("lifecycle owner stopped without a result for session {session_id}")
})?;
}
Ok::<_, anyhow::Error>(())
})
.await;
if joined.is_err() {
bail!(
"timed out cancelling lifecycle owner for session {session_id} while {stage}"
);
}
joined.expect("checked timeout")?;
}
Ok(())
}
pub fn active_lifecycles(&self) -> Vec<RuntimeLifecycleView> {
let controller_owner = self.owner();
Self::active_lifecycles_with(&controller_owner)
}
pub(super) fn active_lifecycles_with(owner: &RuntimeStateOwner) -> Vec<RuntimeLifecycleView> {
let controller = owner.controller();
owner
.lifecycle
.iter()
.filter(|(_, active)| active.is_visible())
.map(|(session_id, active)| RuntimeLifecycleView {
operation_id: active.operation_id.clone(),
cancellable: active.is_cancellable()
&& lifecycle_cancellable(
active.kind,
durable_session_state(controller, session_id),
),
session_id: session_id.clone(),
kind: active.kind.into(),
started_at_epoch_seconds: active.started_at_epoch_seconds,
active_stages: active
.active_stages
.iter()
.map(|(stage, (_, started_at))| (*stage, *started_at))
.collect(),
resume_destination: active.resume_destination.clone(),
notice: active.notice.clone(),
})
.collect()
}
pub fn session_state(&self, session_id: &str) -> Option<mj_core::state::SessionState> {
let owner = self.owner();
if owner.close_requested.contains(session_id) {
return Some(SessionState::Closing);
}
owner
.controller()
.state
.sessions
.get(session_id)
.map(|record| record.state)
}
pub fn session_record(&self, session_id: &str) -> Option<SessionRecord> {
self.owner()
.controller()
.state
.sessions
.get(session_id)
.cloned()
}
pub async fn workspace_session_handle(
&self,
session_id: &str,
) -> Result<crate::session_manager::ManagedSessionHandle> {
let record = self.session_record(session_id).context("unknown session")?;
ensure!(
record.target.is_some()
&& record.state == SessionState::Running
&& !self.close_is_requested(session_id),
"session must have a live running target for file injection"
);
self.session_manager.session(session_id.to_owned()).await
}
pub async fn checkpoint_session_now(
&self,
session_id: &str,
) -> Result<mj_core::state::CheckpointMetadata> {
let _upgrade_work = crate::upgrade::activity("requested checkpoint")?;
if let Some(busy) = self.session_lifecycle_busy(session_id) {
return Err(anyhow::Error::new(busy));
}
let mut controller = blocking(Controller::load).await?;
let checkpoint = controller.checkpoint_session(session_id).await?;
refresh_runtime_controller(self).await;
Ok(checkpoint)
}
pub(super) fn session_lifecycle_busy(&self, session_id: &str) -> Option<SessionLifecycleBusy> {
let lifecycle_owner = self.owner();
let lifecycle = &lifecycle_owner.lifecycle;
let active = lifecycle.get(session_id)?;
active
.is_running()
.then(|| describe_lifecycle_busy(session_id, active))
}
pub(super) fn any_lifecycle_busy(&self) -> Option<SessionLifecycleBusy> {
self.owner()
.lifecycle
.iter()
.find(|(_, active)| active.is_running())
.map(|(session_id, active)| describe_lifecycle_busy(session_id, active))
}
pub fn session_projection(
&self,
) -> (
mj_core::snapshot_map::SnapshotMap<String, SessionRecord>,
Vec<RuntimeLifecycleView>,
) {
let controller_owner = self.owner();
let operations = Self::active_lifecycles_with(&controller_owner);
(controller_owner.projected_records(), operations)
}
pub(crate) fn worker_controller_projection(&self) -> Controller {
self.owner().pollable_worker_inputs().controller()
}
pub(crate) fn controller_projection(&self) -> Controller {
let owner = self.owner();
let mut state = owner.controller().state.clone();
state.sessions = owner.projected_records();
Controller {
config: owner.controller().config.clone(),
state,
}
}
pub(crate) fn active_controller_projection(&self) -> Controller {
let owner = self.owner();
Controller {
config: owner.controller().config.clone(),
state: mj_core::state::State {
sessions: owner
.indexes
.active
.keys()
.filter_map(|id| {
owner
.controller()
.state
.sessions
.get(id)
.map(|record| (id.clone(), record.clone()))
})
.collect(),
..Default::default()
},
}
}
pub fn cancel_lifecycle_if_active(&self, session_id: &str) {
if let Some(active) = self.owner().lifecycle.get_mut(session_id) {
active.request_cancel();
self.publish_revision();
}
}
pub(super) fn set_lifecycle_resume_destination(
&self,
session_id: &str,
profile_id: String,
target_id: String,
) {
if let Some(active) = self.owner().lifecycle.get_mut(session_id) {
active.resume_destination = Some((profile_id, target_id));
self.publish_revision();
}
}
pub(super) fn change_lifecycle_stage(
&self,
session_id: &str,
operation_id: &str,
stage: ProvisionStage,
active: bool,
) {
let changed = {
let mut lifecycle_owner = self.owner();
let lifecycle = &mut lifecycle_owner.lifecycle;
let Some(operation) = lifecycle.get_mut(session_id) else {
return;
};
if operation.operation_id != operation_id || !operation.is_running() {
return;
}
if active {
let entry = operation
.active_stages
.entry(stage)
.or_insert_with(|| (0, epoch_seconds()));
entry.0 += 1;
entry.0 == 1
} else {
let Some((count, _)) = operation.active_stages.get_mut(&stage) else {
return;
};
*count -= 1;
if *count == 0 {
operation.active_stages.remove(&stage);
true
} else {
false
}
}
};
if changed {
self.publish_revision();
}
}
pub(crate) fn push_notice(&self, session_id: &str, text: impl Into<String>) {
const RETAINED_NOTICES: usize = 32;
let notice = RuntimeNotice {
id: self.next_notice_id.fetch_add(1, Ordering::AcqRel),
session_id: session_id.to_owned(),
text: text.into(),
};
{
let mut notices = self.notices.lock().unwrap_or_else(PoisonError::into_inner);
notices.push_back(notice);
while notices.len() > RETAINED_NOTICES {
notices.pop_front();
}
}
self.publish_revision();
}
pub(super) fn reserve_move_destination(&self, session_id: &str, operation_id: &str) {
let mut owner = self.owner();
if let Some(active) = owner.lifecycle.get_mut(session_id)
&& active.operation_id == operation_id
&& active.kind == LifecycleKind::Move
{
active.phase = match active.phase {
LifecyclePhase::Executing => LifecyclePhase::MovingDestination,
LifecyclePhase::Cancelling => LifecyclePhase::CancellingMoveDestination,
_ => return,
};
self.publish_revision();
}
}
pub(super) fn set_lifecycle_notice(&self, session_id: &str, operation_id: &str, notice: &str) {
if let Some(active) = self.owner().lifecycle.get_mut(session_id)
&& active.operation_id == operation_id
&& active.is_running()
{
active.notice = Some(notice.to_owned());
self.publish_revision();
}
}
fn cancel_operation(&self, session_id: &str, operation_id: &str) {
if let Some(active) = self.owner().lifecycle.get_mut(session_id)
&& active.operation_id == operation_id
&& active.request_cancel()
{
self.publish_revision();
}
}
}