Skip to main content

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}