onlyne_client/session/dispatch/idle.rs
1//! Idle and suspended sessions: the two states a `task` or `role` session
2//! reaches between deliveries.
3//!
4//! A delivery that settles leaves a scoped session idle rather than ending it,
5//! and the scope's `idle_close` bound decides when this client does something
6//! about that: it releases the process, which is what `suspended` means. Whether
7//! it may release the process at all is the runtime's answer, and the degradation
8//! is part of the contract — a runtime that cannot resume keeps its process, and
9//! the session waits for its family's next delivery exactly as it stands.
10
11use super::retire::{PendingClose, close_retired};
12use super::state::{DispatchInner, binding_task_state};
13use super::transport::names_session;
14use super::*;
15
16impl DispatchState {
17 /// Release the process of every idle session whose scope bound has expired.
18 ///
19 /// A suspension is this client's own act and it is only worth doing for a
20 /// session the runtime can bring back: the conversation lives in the
21 /// runtime's own store, so the client starts the same command again when the
22 /// family's next delivery arrives and the runtime resumes the session it was
23 /// given. A runtime that declared no resume keeps its process instead — the
24 /// "process alive, session alive" degradation the scope table names — and
25 /// this sweep leaves it alone.
26 ///
27 /// The closes are host round trips, so they are collected here and run once
28 /// the dispatch lock is off them, the way the sweeps that retire sessions
29 /// already do. Answers the sessions it released, which are the ones whose
30 /// row and claim this caller has to publish.
31 pub fn suspend_idle_sessions(&self, now: Instant) -> Vec<String> {
32 let mut pending: Vec<PendingClose> = Vec::new();
33 let mut suspended: Vec<String> = Vec::new();
34 {
35 let mut inner = self.inner.lock();
36 let Some(bound) = scope::idle_bound(&inner.session_policy) else {
37 return Vec::new();
38 };
39 let due: Vec<String> = inner
40 .sessions
41 .iter()
42 .filter(|(_, slot)| {
43 !slot.suspended
44 && !slot.read_only
45 && slot.task_id.is_none()
46 && slot
47 .idle_since
48 .is_some_and(|since| now.saturating_duration_since(since) >= bound)
49 })
50 .map(|(key, _)| key.clone())
51 .collect();
52 for key in due {
53 let session_id = inner
54 .sessions
55 .get(&key)
56 .map(|slot| slot.session.task_id.clone())
57 .unwrap_or_else(|| key.clone());
58 if !resumable(&inner, &key) {
59 tracing::info!(
60 session = %session_id,
61 idle_secs = inner
62 .sessions
63 .get(&key)
64 .and_then(|slot| slot.idle_since)
65 .map(|since| now.saturating_duration_since(since).as_secs())
66 .unwrap_or_default(),
67 "an idle session's runtime cannot resume it; the process stays and the session waits"
68 );
69 continue;
70 }
71 if suspend_locked(&mut inner, &key, &mut pending) {
72 suspended.push(session_id);
73 }
74 }
75 }
76 close_retired(pending);
77 suspended
78 }
79
80 /// Whether one session's process is currently released.
81 pub fn is_suspended(&self, session_id: &str) -> bool {
82 let inner = self.inner.lock();
83 inner
84 .sessions
85 .iter()
86 .any(|(key, slot)| names_session(key, slot, session_id) && slot.suspended)
87 }
88}
89
90/// Whether one session's runtime declared that it can resume a session this
91/// client released.
92///
93/// The declaration rides the plugin's mount (`Capability::Resume`), because the
94/// runtime is the only party that knows whether its conversation survives its
95/// process: pi and DSH keep their own session files, an ACP agent declares
96/// `session/resume` or `session/load`, and a runtime that has neither cannot be
97/// brought back into a conversation it does not remember.
98///
99/// A self-driven backend owns its agent and talks to it without an adapter
100/// socket, so no mount ever declares anything for it: this client answers no,
101/// which keeps the process in place. Answering yes there would take the client's
102/// own guarantee on behalf of a runtime nobody asked.
103pub(super) fn resumable(inner: &DispatchInner, key: &str) -> bool {
104 if inner.backend.self_driven() {
105 return false;
106 }
107 let Some(slot) = inner.sessions.get(key) else {
108 return false;
109 };
110 inner
111 .transports
112 .iter()
113 .find(|(session, _)| names_session(key, slot, session))
114 .is_some_and(|(_, (_, capabilities))| capabilities.contains(&Capability::Resume))
115}
116
117/// Release one idle session's process while its row and its conversation stay.
118///
119/// The order is the store's: the tuple moves first — `Suspend` closes the
120/// resource while the generation stays live, which is what tells `project` this
121/// session is idle rather than exited — and the binding the write re-takes is
122/// handed back once the row is final. A session that keeps a binding while
123/// serving nothing reads as bound to a delivery it has finished, which is the
124/// reading `hello.live_sessions` and the mirror's `task_id` both take.
125/// Whether one session reads `idle`: the state a suspension is defined from.
126///
127/// The row is the session's, and the delivery it last served is the binding's:
128/// `projection_of` derives the lifecycle from both, which is the same reading
129/// the reducer takes when it judges the `Suspend` event this client feeds it.
130fn at_rest(inner: &DispatchInner, key: &str) -> bool {
131 let Some(slot) = inner.sessions.get(key) else {
132 return false;
133 };
134 let Ok(Some(row)) = inner.store.get_session(&slot.session.task_id) else {
135 return false;
136 };
137 projection_of(&row, binding_task_state(inner, slot)).lifecycle == Lifecycle::Idle
138}
139
140pub(super) fn suspend_locked(
141 inner: &mut DispatchInner,
142 key: &str,
143 pending: &mut Vec<PendingClose>,
144) -> bool {
145 let Some(slot) = inner.sessions.get(key) else {
146 return false;
147 };
148 let session = slot.session.clone();
149 let session_id = session.task_id.clone();
150 // The plan's state machine suspends a session that is *idle*, and the
151 // reducer reads that off the tuple: a row still projecting `working` — an
152 // agent with a turn in flight, or a delivery whose intent is unanswered —
153 // is not a session whose process may go, and its `Suspend` is an undefined
154 // transition. The sweep's own `idle_since` says how long the session has
155 // held no delivery; this says whether the runtime is at rest, which is the
156 // half only its own reports can answer.
157 if !at_rest(inner, key) {
158 tracing::debug!(
159 session = %session_id,
160 "the session is not idle yet; its process stays and the family's next delivery still enters it"
161 );
162 return false;
163 }
164 if let Err(error) = feed_suspended(&inner.bridge, &inner.store, &session_id) {
165 tracing::warn!(
166 session = %session_id,
167 error = %error,
168 "a suspended session's row was not written; the process stays"
169 );
170 return false;
171 }
172 if let Err(error) = inner.store.release_binding(&session_id, &session_id) {
173 tracing::warn!(
174 session = %session_id,
175 error = %error,
176 "a suspended session's own binding was not handed back"
177 );
178 }
179 let now = Instant::now();
180 if let Some(slot) = inner.sessions.get_mut(key) {
181 slot.suspended = true;
182 slot.idle_since = None;
183 // Nothing is expected to dial for a released session, and the two death
184 // windows read a stamp a suspended session would never refresh: the
185 // reconnect grace would retire the conversation this release just saved.
186 slot.dropped_at = None;
187 slot.last_beat = Some(now);
188 }
189 inner.transports.remove(&session_id);
190 tracing::info!(
191 session = %session_id,
192 backend = %session.backend,
193 "idle session suspended; its process was released and the conversation stays in the runtime"
194 );
195 pending.push(PendingClose {
196 backend: Arc::clone(&inner.backend),
197 session,
198 reason: crate::backend::CloseReason::Operator,
199 });
200 true
201}