use super::*;
use super::outbound::store_ack;
use super::projection::stored_task_state;
use super::state::{
DispatchInner, DispatchState, forget_tools_binding, has_attached_transport, session_exited,
slot_key_serving_task,
};
use super::transport::{held_read_only, names_session};
pub const SESSION_DEAD: &str = "session_dead";
pub fn on_recycled(
state: &DispatchState,
task_id: &str,
reason: crate::backend::CloseReason,
) -> Result<()> {
let mut inner = state.inner.lock();
release_locked(&mut inner, task_id, Some(reason))
}
pub(super) fn stored_close_reason(
inner: &DispatchInner,
task_id: &str,
) -> Option<crate::backend::CloseReason> {
match stored_task_state(inner, task_id) {
TaskState::Pending => None,
TaskState::Done => Some(crate::backend::CloseReason::Completed),
TaskState::Failed | TaskState::Blocked => Some(crate::backend::CloseReason::Fault),
TaskState::Cancelled => Some(crate::backend::CloseReason::Cancelled),
}
}
fn grace_close_reason(inner: &DispatchInner, task_id: &str) -> crate::backend::CloseReason {
match stored_task_state(inner, task_id) {
TaskState::Done => crate::backend::CloseReason::Completed,
TaskState::Pending | TaskState::Failed | TaskState::Blocked => {
crate::backend::CloseReason::Fault
}
TaskState::Cancelled => crate::backend::CloseReason::Cancelled,
}
}
fn dropped_past_window(
inner: &DispatchInner,
key: &str,
slot: &SessionSlot,
now: Instant,
window: Duration,
) -> bool {
slot.dropped_at.is_some_and(|dropped| {
!has_attached_transport(inner, key, slot)
&& now
.checked_duration_since(dropped)
.is_some_and(|away| away >= window)
})
}
fn silent_past_window(inner: &DispatchInner, key: &str, slot: &SessionSlot, now: Instant) -> bool {
if slot.dropped_at.is_some() || !has_attached_transport(inner, key, slot) {
return false;
}
let Some(task_id) = slot.task_id.as_deref() else {
return false;
};
if stored_task_state(inner, task_id) != TaskState::Pending {
return false;
}
let quiet = HEARTBEAT_INTERVAL * HEARTBEAT_SILENCE_MARGIN;
slot.last_beat.is_some_and(|beat| {
now.checked_duration_since(beat)
.is_some_and(|away| away >= quiet)
})
}
fn unbind_transports(inner: &mut DispatchInner, key: &str, slot: &SessionSlot) {
inner
.transports
.retain(|served, _| !names_session(key, slot, served));
}
pub(super) struct PendingClose {
pub(super) backend: Arc<dyn SessionBackend>,
pub(super) session: SessionRef,
pub(super) reason: crate::backend::CloseReason,
}
pub(super) fn close_retired(pending: Vec<PendingClose>) {
for close in pending {
if let Err(error) = close.backend.close(&close.session, close.reason, false) {
tracing::warn!(
task = %close.session.task_id,
backend = %close.session.backend,
resource = %close.session.backend_ref,
error = %error,
"session resource retirement failed"
);
}
}
}
pub(super) fn keeps_idle(inner: &DispatchInner, key: &str) -> bool {
inner
.sessions
.get(key)
.is_some_and(|slot| slot.keeps_idle && !slot.read_only)
}
pub(super) fn retire_idle_locked(
inner: &mut DispatchInner,
key: &str,
reason: crate::backend::CloseReason,
pending: &mut Vec<PendingClose>,
) -> bool {
let Some(slot) = inner.sessions.get(key) else {
return false;
};
if slot.task_id.is_some() || has_attached_transport(inner, key, slot) {
return false;
}
let original = slot.session.clone();
let task_id = original.task_id.clone();
let resource = inner
.store
.get_session(&task_id)
.ok()
.flatten()
.map(|row| row.resource_state)
.unwrap_or_else(|| "detached".to_string());
if resource != "detached" && resource != "closed" {
let session = match inner.backend.attach(&original) {
Ok(refreshed) => {
if refreshed != original {
inner.bridge.track_live(refreshed.clone());
if let Some(slot) = inner.sessions.get_mut(key) {
slot.session = refreshed.clone();
}
}
refreshed
}
Err(_) => original,
};
tracing::info!(
task = %task_id,
backend = %session.backend,
resource = %session.backend_ref,
?reason,
"retiring idle session resource"
);
if let Err(error) = feed_resource_closed(&inner.bridge, &inner.store, &task_id) {
tracing::warn!(
task = %task_id,
backend = %session.backend,
resource = %session.backend_ref,
error = %error,
"session resource close projection failed"
);
}
pending.push(PendingClose {
backend: Arc::clone(&inner.backend),
session,
reason,
});
}
if reason == crate::backend::CloseReason::Completed {
if let Err(error) = feed_agent_gone(&inner.bridge, &inner.store, &task_id) {
tracing::warn!(
task = %task_id,
error = %error,
"agent-gone projection failed for a completed session"
);
}
}
inner.bridge.untrack_live(&task_id);
forget_tools_binding(inner, key);
inner.sessions.remove(key);
true
}
pub(super) fn release_locked(
inner: &mut DispatchInner,
task_id: &str,
reason: Option<crate::backend::CloseReason>,
) -> Result<()> {
let resource = inner
.store
.get_session(task_id)?
.map(|row| row.resource_state)
.unwrap_or_else(|| "detached".to_string());
if let Some((key, slot)) = slot_key_serving_task(inner, task_id)
.and_then(|key| inner.sessions.get(&key).map(|slot| (key, slot.clone())))
{
if let Some(reason) = reason {
if resource != "detached" && resource != "closed" {
feed_resource_closed(&inner.bridge, &inner.store, task_id)?;
inner.backend.close(&slot.session, reason, false)?;
}
if let Err(error) = feed_agent_gone(&inner.bridge, &inner.store, task_id) {
tracing::warn!(
task = %task_id,
error = %error,
"agent-gone projection failed for a closed session"
);
}
inner.bridge.untrack_live(task_id);
let held = inner
.sessions
.get_mut(&key)
.and_then(|slot| slot.msg_id.take());
if let Some(msg_id) = held {
store_ack(
inner,
AckArgs {
msg_id,
op_id: None,
accepted: false,
reason: Some(close_refusal(reason).to_string()),
},
);
}
forget_tools_binding(inner, &key);
inner.sessions.remove(&key);
} else {
let session_id = inner
.sessions
.get(&key)
.map(|slot| slot.session.task_id.clone())
.unwrap_or_else(|| task_id.to_string());
let keeps = keeps_idle(inner, &key);
let now = Instant::now();
if let Some(session) = inner.sessions.get_mut(&key) {
session.task_id = None;
session.ready = false;
session.idle_since = Some(now);
}
if let Err(error) = inner.store.release_binding(&session_id, task_id) {
tracing::warn!(
session = %session_id,
task = %task_id,
error = %error,
"a settled session's binding was not handed back"
);
}
if keeps {
tracing::info!(
session = %session_id,
task = %task_id,
"delivery settled; the session stays idle for its scope's next one"
);
} else {
let mut pending = Vec::new();
retire_idle_locked(
inner,
&key,
crate::backend::CloseReason::Completed,
&mut pending,
);
close_retired(pending);
}
}
}
inner.stall.forget(task_id);
if reason.is_some() && resource == "detached" {
inner
.store
.note_alert(format!("session recycled {task_id}"));
}
Ok(())
}
fn close_refusal(reason: crate::backend::CloseReason) -> &'static str {
match reason {
crate::backend::CloseReason::Cancelled => ControlWord::Cancel.refusal(),
crate::backend::CloseReason::Operator => ControlWord::Recycle.refusal(),
_ => "operator close",
}
}
pub fn close_all(state: &DispatchState, reason: crate::backend::CloseReason, budget: Duration) {
let started = Instant::now();
let pending: Vec<PendingClose> = {
let mut inner = state.inner.lock();
let sessions: Vec<(String, SessionRef)> = inner
.sessions
.iter()
.map(|(key, slot)| (key.clone(), slot.session.clone()))
.collect();
let mut pending = Vec::with_capacity(sessions.len());
for (key, session) in sessions {
inner.bridge.untrack_live(&session.task_id);
forget_tools_binding(&mut inner, &key);
inner.sessions.remove(&key);
pending.push(PendingClose {
backend: Arc::clone(&inner.backend),
session,
reason,
});
}
pending
};
for close in pending {
if started.elapsed() > budget {
tracing::warn!(
task = %close.session.task_id,
"shutdown close budget reached; the resource is left behind"
);
continue;
}
if let Err(error) = close.backend.close(&close.session, close.reason, false) {
tracing::warn!(
task = %close.session.task_id,
error = %error,
"session close failed during shutdown"
);
}
}
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub enum RetirementArm {
Dropped,
Silent,
}
impl RetirementArm {
pub fn word(self) -> &'static str {
match self {
Self::Dropped => "reconnect_grace",
Self::Silent => "heartbeat_silence",
}
}
}
pub struct Retired {
pub session_id: String,
pub arm: RetirementArm,
pub quiet_secs: u64,
pub away_secs: Option<u64>,
}
fn quiet_secs(slot: &SessionSlot, now: Instant) -> u64 {
slot.last_beat
.map(|beat| now.saturating_duration_since(beat).as_secs())
.unwrap_or(u64::MAX)
}
fn away_secs(slot: &SessionSlot, now: Instant) -> Option<u64> {
slot.dropped_at
.map(|left| now.saturating_duration_since(left).as_secs())
}
impl DispatchState {
pub fn reclaim_exited_resources(&self) -> Vec<String> {
let mut inner = self.inner.lock();
let candidates: Vec<(String, crate::backend::CloseReason)> = inner
.sessions
.iter()
.filter(|(key, slot)| {
!slot.suspended
&& slot.task_id.is_none()
&& session_exited(&inner, &slot.session.task_id)
&& !has_attached_transport(&inner, key, slot)
})
.filter_map(|(key, slot)| {
stored_close_reason(&inner, &slot.session.task_id)
.map(|reason| (key.clone(), reason))
})
.collect();
let mut retired: Vec<String> = Vec::new();
let mut pending: Vec<PendingClose> = Vec::new();
for (key, reason) in candidates {
let Some(task_id) = inner
.sessions
.get(&key)
.map(|slot| slot.session.task_id.clone())
else {
continue;
};
if retire_idle_locked(&mut inner, &key, reason, &mut pending) {
retired.push(task_id);
}
}
drop(inner);
close_retired(pending);
retired
}
pub fn retire_dropped_ghosts(&self, now: Instant, grace_secs: u64) -> Vec<Retired> {
if grace_secs == 0 {
return Vec::new();
}
let window = Duration::from_secs(grace_secs);
let mut inner = self.inner.lock();
let due: Vec<(String, RetirementArm)> = inner
.sessions
.iter()
.filter_map(|(key, slot)| {
let arm = if dropped_past_window(&inner, key, slot, now, window) {
RetirementArm::Dropped
} else if silent_past_window(&inner, key, slot, now) {
RetirementArm::Silent
} else {
return None;
};
Some((key.clone(), arm))
})
.collect();
let mut retired: Vec<Retired> = Vec::new();
let mut pending: Vec<PendingClose> = Vec::new();
for (key, arm) in due {
let Some(slot) = inner.sessions.get(&key).cloned() else {
continue;
};
if held_read_only(&inner, &key, &slot) {
continue;
}
if slot.dropped_at.is_none() {
unbind_transports(&mut inner, &key, &slot);
}
let task_id = slot.session.task_id.clone();
let reason = grace_close_reason(&inner, &task_id);
if let Some(owed) = slot.task_id.clone() {
inner.control_settles.retain(|noted| noted.task_id != owed);
if let Err(error) = inner.store.settle_task(&owed, TaskState::Failed) {
tracing::warn!(
task = %owed,
error = %error,
"the task of a retired ghost was not settled"
);
}
let handle = inner
.sessions
.get_mut(&key)
.and_then(|slot| slot.msg_id.take());
if let Some(msg_id) = handle {
store_ack(
&inner,
AckArgs {
msg_id,
op_id: None,
accepted: false,
reason: Some(SESSION_DEAD.to_string()),
},
);
}
}
if let Err(error) = feed_agent_gone(&inner.bridge, &inner.store, &task_id) {
tracing::warn!(
task = %task_id,
error = %error,
"agent-gone projection failed for a retired ghost"
);
}
if let Some(slot) = inner.sessions.get_mut(&key) {
slot.task_id = None;
}
if retire_idle_locked(&mut inner, &key, reason, &mut pending) {
retired.push(Retired {
session_id: task_id,
arm,
quiet_secs: quiet_secs(&slot, now),
away_secs: away_secs(&slot, now),
});
}
}
drop(inner);
close_retired(pending);
retired
}
}