use super::retire::{PendingClose, close_retired};
use super::state::{DispatchInner, binding_task_state};
use super::transport::names_session;
use super::*;
impl DispatchState {
pub fn suspend_idle_sessions(&self, now: Instant) -> Vec<String> {
let mut pending: Vec<PendingClose> = Vec::new();
let mut suspended: Vec<String> = Vec::new();
{
let mut inner = self.inner.lock();
let Some(bound) = scope::idle_bound(&inner.session_policy) else {
return Vec::new();
};
let due: Vec<String> = inner
.sessions
.iter()
.filter(|(_, slot)| {
!slot.suspended
&& !slot.read_only
&& slot.task_id.is_none()
&& slot
.idle_since
.is_some_and(|since| now.saturating_duration_since(since) >= bound)
})
.map(|(key, _)| key.clone())
.collect();
for key in due {
let session_id = inner
.sessions
.get(&key)
.map(|slot| slot.session.task_id.clone())
.unwrap_or_else(|| key.clone());
if !resumable(&inner, &key) {
tracing::info!(
session = %session_id,
idle_secs = inner
.sessions
.get(&key)
.and_then(|slot| slot.idle_since)
.map(|since| now.saturating_duration_since(since).as_secs())
.unwrap_or_default(),
"an idle session's runtime cannot resume it; the process stays and the session waits"
);
continue;
}
if suspend_locked(&mut inner, &key, &mut pending) {
suspended.push(session_id);
}
}
}
close_retired(pending);
suspended
}
pub fn is_suspended(&self, session_id: &str) -> bool {
let inner = self.inner.lock();
inner
.sessions
.iter()
.any(|(key, slot)| names_session(key, slot, session_id) && slot.suspended)
}
}
pub(super) fn resumable(inner: &DispatchInner, key: &str) -> bool {
if inner.backend.self_driven() {
return false;
}
let Some(slot) = inner.sessions.get(key) else {
return false;
};
inner
.transports
.iter()
.find(|(session, _)| names_session(key, slot, session))
.is_some_and(|(_, (_, capabilities))| capabilities.contains(&Capability::Resume))
}
fn at_rest(inner: &DispatchInner, key: &str) -> bool {
let Some(slot) = inner.sessions.get(key) else {
return false;
};
let Ok(Some(row)) = inner.store.get_session(&slot.session.task_id) else {
return false;
};
projection_of(&row, binding_task_state(inner, slot)).lifecycle == Lifecycle::Idle
}
pub(super) fn suspend_locked(
inner: &mut DispatchInner,
key: &str,
pending: &mut Vec<PendingClose>,
) -> bool {
let Some(slot) = inner.sessions.get(key) else {
return false;
};
let session = slot.session.clone();
let session_id = session.task_id.clone();
if !at_rest(inner, key) {
tracing::debug!(
session = %session_id,
"the session is not idle yet; its process stays and the family's next delivery still enters it"
);
return false;
}
if let Err(error) = feed_suspended(&inner.bridge, &inner.store, &session_id) {
tracing::warn!(
session = %session_id,
error = %error,
"a suspended session's row was not written; the process stays"
);
return false;
}
if let Err(error) = inner.store.release_binding(&session_id, &session_id) {
tracing::warn!(
session = %session_id,
error = %error,
"a suspended session's own binding was not handed back"
);
}
let now = Instant::now();
if let Some(slot) = inner.sessions.get_mut(key) {
slot.suspended = true;
slot.idle_since = None;
slot.dropped_at = None;
slot.last_beat = Some(now);
}
inner.transports.remove(&session_id);
tracing::info!(
session = %session_id,
backend = %session.backend,
"idle session suspended; its process was released and the conversation stays in the runtime"
);
pending.push(PendingClose {
backend: Arc::clone(&inner.backend),
session,
reason: crate::backend::CloseReason::Operator,
});
true
}