Skip to main content

onlyne_client/session/dispatch/
retire.rs

1use super::*;
2
3use super::outbound::store_ack;
4use super::projection::stored_task_state;
5use super::state::{
6    DispatchInner, DispatchState, forget_tools_binding, has_attached_transport, session_exited,
7    slot_key_named, slot_key_serving_task,
8};
9use super::transport::{held_read_only, names_session};
10
11/// Reason a session that stopped answering is closed with.
12///
13/// The wire word the ledger and `onlyne sessions` read for a death this client
14/// judged: the reconnect sweep stamps it on the refusal ack that buries the
15/// session's delivery, and an operator's own
16/// `repair fail --reason session_dead` writes the same word.
17pub const SESSION_DEAD: &str = "session_dead";
18
19/// Retire one session. The stored tuple decides whether a live resource
20/// remains to close, and the caller's reason reaches the backend unchanged, so an
21/// operator cancel stops reporting itself as a completion.
22pub fn on_recycled(
23    state: &DispatchState,
24    task_id: &str,
25    reason: crate::backend::CloseReason,
26) -> Result<()> {
27    let mut inner = state.inner.lock();
28    release_locked(&mut inner, task_id, Some(reason))
29}
30
31/// Why the task of one session ended, as the retirement reason the task table
32/// records.
33///
34/// The task's own row is the only place a verdict lives: the session tuple says
35/// nothing about how its work ended, and an open task — `pending`, or no row at
36/// all, which is the same reading — has no reason to retire anything.
37pub(super) fn stored_close_reason(
38    inner: &DispatchInner,
39    task_id: &str,
40) -> Option<crate::backend::CloseReason> {
41    match stored_task_state(inner, task_id) {
42        TaskState::Pending => None,
43        TaskState::Done => Some(crate::backend::CloseReason::Completed),
44        // A blocked delivery leaves the work owed, which is what a fault reason
45        // names here, exactly as it does for a failed one.
46        TaskState::Failed | TaskState::Blocked => Some(crate::backend::CloseReason::Fault),
47        TaskState::Cancelled => Some(crate::backend::CloseReason::Cancelled),
48    }
49}
50
51/// The reason one session the reconnect grace retires is closed with, read from
52/// what this client holds rather than from a settle that never came.
53///
54/// The id is the session's own (`slot.session.task_id`), never the task binding
55/// beside it: the sweep is what feeds that session's row and closes the resource
56/// its agent was holding, and the two ids part company exactly where a slot's
57/// binding is not the session it was born for. A session still owing work closes
58/// as that task's own record reads: a `done` task is a `Completed`, a
59/// `cancelled` one is a `Cancelled`, and a `failed` task — like one that never
60/// settled at all — is a `Fault`, because the work was still owed when the agent
61/// left.
62fn grace_close_reason(inner: &DispatchInner, task_id: &str) -> crate::backend::CloseReason {
63    match stored_task_state(inner, task_id) {
64        TaskState::Done => crate::backend::CloseReason::Completed,
65        TaskState::Pending | TaskState::Failed | TaskState::Blocked => {
66            crate::backend::CloseReason::Fault
67        }
68        TaskState::Cancelled => crate::backend::CloseReason::Cancelled,
69    }
70}
71
72/// Whether one slot is due for the reconnect grace to take it.
73///
74/// The window as it always was: the connection that would have sent this
75/// session's next heartbeat has ended, and the window runs from the moment it
76/// left. A slot a live connection serves is not this arm's to end, because an
77/// attached transport is the one thing that says the agent is still reachable.
78fn dropped_past_window(
79    inner: &DispatchInner,
80    key: &str,
81    slot: &SessionSlot,
82    now: Instant,
83    window: Duration,
84) -> bool {
85    slot.dropped_at.is_some_and(|dropped| {
86        !has_attached_transport(inner, key, slot)
87            && now
88                .checked_duration_since(dropped)
89                .is_some_and(|away| away >= window)
90    })
91}
92
93/// Whether one session's own agent has gone quiet on a socket that is still up.
94///
95/// The window's other door, and it exists precisely because an attached
96/// transport — the one thing the arm above trusts — can lie. A socket that is
97/// still up proves the connection survived; it says nothing about the agent
98/// behind it, and a plugin whose event loop is blocked holds its socket and
99/// stops beating. Nothing the client reads before this could see that: no socket
100/// ends, so no drop clock ever starts, and the session keeps its slot, its
101/// projected row and its host resource for as long as the client runs.
102///
103/// So the reading is taken off the frames themselves. A session whose task is
104/// still bound and unsettled and whose last accepted frame is older than the
105/// protocol's heartbeat interval by [`HEARTBEAT_SILENCE_MARGIN`] has no agent
106/// behind its socket.
107///
108/// Why no task bound is excluded: the plugin stops its heartbeat loop with the
109/// last task it was given, so a task-free session that has gone quiet is an
110/// agent waiting for work by design — the ordinary shape between deliveries.
111/// Sweeping it would retire the very connection the next payload is staged onto.
112/// The unsettled check beside the binding says the same thing about work that
113/// already landed: a settled task has nothing left for its session to answer,
114/// and `release_locked` has already given the binding back.
115///
116/// Why the drop clock's arm never reads this stamp: the two are mutually
117/// exclusive by construction. A re-mount clears `dropped_at`, so a session that
118/// came back is judged by its beat alone, and a session whose connection went
119/// away is judged by the clock alone and never by a stamp its agent can no
120/// longer refresh.
121fn silent_past_window(inner: &DispatchInner, key: &str, slot: &SessionSlot, now: Instant) -> bool {
122    if slot.dropped_at.is_some() || !has_attached_transport(inner, key, slot) {
123        return false;
124    }
125    let Some(task_id) = slot.task_id.as_deref() else {
126        return false;
127    };
128    if stored_task_state(inner, task_id) != TaskState::Pending {
129        return false;
130    }
131    let quiet = HEARTBEAT_INTERVAL * HEARTBEAT_SILENCE_MARGIN;
132    slot.last_beat.is_some_and(|beat| {
133        now.checked_duration_since(beat)
134            .is_some_and(|away| away >= quiet)
135    })
136}
137
138/// Stop holding a session's connection as its transport.
139///
140/// The act a socket ending performs, run here on the client's own verdict
141/// instead: a session judged dead by its silence still holds its socket, and the
142/// connection is not going to end on its own while the agent behind it is
143/// blocked. The binding goes now, so the retirement below runs as the ordinary
144/// one, and a frame from that connection afterwards is refused the way every
145/// frame from a connection no session answers for is — it is served no state,
146/// which is the same door a stale reporter already comes to.
147fn unbind_transports(inner: &mut DispatchInner, key: &str, slot: &SessionSlot) {
148    inner
149        .transports
150        .retain(|served, _| !names_session(key, slot, served));
151}
152
153/// One backend close a retirement owed, carried off the dispatch lock.
154///
155/// A close is a host round trip — a pane kill, a terminal close, an agent reap —
156/// and it can take seconds. The dispatch lock is the one lock every adapter frame,
157/// every report, and every slot of this role queues behind, so a sweep that
158/// retires more than one session at a time must not spend that lock on the hosts
159/// it is calling. A retirement therefore records what it took down here and runs
160/// the closes once it is off the lock.
161pub(super) struct PendingClose {
162    pub(super) backend: Arc<dyn SessionBackend>,
163    pub(super) session: SessionRef,
164    pub(super) reason: crate::backend::CloseReason,
165}
166
167/// Run the closes a retirement collected.
168///
169/// Four of the five callers hold the dispatch lock and can put it down first — the
170/// two sweeps, the goodbye path, and the merged-handoff retire — and each does,
171/// because a close that blocks must not hold every adapter frame of this role
172/// behind it. The fifth is `release_locked`, which runs inside a caller that
173/// already holds the lock for the store work above it and cannot leave; it takes
174/// the close where it stands, which is the behavior that path has always had.
175///
176/// A failed close is reported and not propagated. The slot is gone from this
177/// client's books either way and its row already reads closed, so the caller has
178/// nothing left to undo; an error that travelled back would replace the answer
179/// the caller is waiting for — which sessions left, and therefore must be
180/// published — with a failure to say so.
181pub(super) fn close_retired(pending: Vec<PendingClose>) {
182    for close in pending {
183        if let Err(error) = close.backend.close(&close.session, close.reason, false) {
184            tracing::warn!(
185                task = %close.session.task_id,
186                backend = %close.session.backend,
187                resource = %close.session.backend_ref,
188                error = %error,
189                "session resource retirement failed"
190            );
191        }
192    }
193}
194
195/// Whether this role's scope keeps one session alive after its delivery settles.
196///
197/// `oneshot` is the rule the client has always had: the session served its one
198/// delivery and the slot is done with it. A `task` or `role` session outlives
199/// the delivery that opened it, which is the whole point of the scope, so its
200/// slot stays serving nothing until the scope sends it the next delivery or
201/// `idle_close` releases its process.
202///
203/// A read-only slot is never kept: it is a connection that came back for a task
204/// a newer session already took, and it owns no delivery of this role.
205pub(super) fn keeps_idle(inner: &DispatchInner, key: &str) -> bool {
206    inner
207        .sessions
208        .get(key)
209        .is_some_and(|slot| slot.keeps_idle && !slot.read_only)
210}
211
212/// Retire one task-free session after its transport set becomes empty.
213///
214/// This is the `oneshot` ending, and the ending of every session whose scope
215/// does not keep it: its resource closes because the agent able to run another
216/// task in it has left. [`keeps_idle`] answers for the scopes that keep a
217/// session between deliveries, and those never reach here with work behind them.
218///
219/// The idle slot releases its backend resource because the agent able to run
220/// another task in it has left. An attached transport keeps the resource because
221/// that agent remains reachable. The dispatch lock serializes the final transport
222/// check, reference refresh, lifecycle projection, and slot removal with adapter
223/// binding. The backend close is deliberately NOT one of them: it is recorded in
224/// `pending` for the caller to run once it is off the lock, so a sweep of sessions
225/// cannot hold every frame of this role behind a host round trip. See
226/// [`close_retired`] for the one path that cannot put the lock down.
227pub(super) fn retire_idle_locked(
228    inner: &mut DispatchInner,
229    key: &str,
230    reason: crate::backend::CloseReason,
231    pending: &mut Vec<PendingClose>,
232) -> bool {
233    let Some(slot) = inner.sessions.get(key) else {
234        return false;
235    };
236    if slot.task_id.is_some() || has_attached_transport(inner, key, slot) {
237        return false;
238    }
239
240    let original = slot.session.clone();
241    let task_id = original.task_id.clone();
242    let resource = inner
243        .store
244        .get_session(&task_id)
245        .ok()
246        .flatten()
247        .map(|row| row.resource_state)
248        .unwrap_or_else(|| "detached".to_string());
249    if resource != "detached" && resource != "closed" {
250        let session = match inner.backend.attach(&original) {
251            Ok(refreshed) => {
252                if refreshed != original {
253                    inner.bridge.track_live(refreshed.clone());
254                    if let Some(slot) = inner.sessions.get_mut(key) {
255                        slot.session = refreshed.clone();
256                    }
257                }
258                refreshed
259            }
260            Err(_) => original,
261        };
262        tracing::info!(
263            task = %task_id,
264            backend = %session.backend,
265            resource = %session.backend_ref,
266            ?reason,
267            "retiring idle session resource"
268        );
269        if let Err(error) = feed_resource_closed(&inner.bridge, &inner.store, &task_id) {
270            tracing::warn!(
271                task = %task_id,
272                backend = %session.backend,
273                resource = %session.backend_ref,
274                error = %error,
275                "session resource close projection failed"
276            );
277        }
278        pending.push(PendingClose {
279            backend: Arc::clone(&inner.backend),
280            session,
281            reason,
282        });
283    }
284    if reason == crate::backend::CloseReason::Completed {
285        if let Err(error) = feed_agent_gone(&inner.bridge, &inner.store, &task_id) {
286            tracing::warn!(
287                task = %task_id,
288                error = %error,
289                "agent-gone projection failed for a completed session"
290            );
291        }
292    }
293    inner.bridge.untrack_live(&task_id);
294    forget_tools_binding(inner, key);
295    inner.sessions.remove(key);
296    true
297}
298
299/// Give one session's task slot back. Settled tasks enter idle retirement, and
300/// explicit reasons drive the control-close path.
301pub(super) fn release_locked(
302    inner: &mut DispatchInner,
303    task_id: &str,
304    reason: Option<crate::backend::CloseReason>,
305) -> Result<()> {
306    let resource = inner
307        .store
308        .get_session(task_id)?
309        .map(|row| row.resource_state)
310        .unwrap_or_else(|| "detached".to_string());
311    // A client-held session is named by the delivery it was born for, and once
312    // that delivery settled the slot serves no task. An operator's close still
313    // has to reach the process it left behind, so the lookup falls back to the
314    // session that id names — but only for a session its scope does not keep: a
315    // pool member between deliveries is not a leftover, and the operator's word
316    // on an old settled task of one is not its ending.
317    let key = slot_key_serving_task(inner, task_id)
318        .or_else(|| slot_key_named(inner, task_id).filter(|key| !keeps_idle(inner, key)));
319    if let Some((key, slot)) =
320        key.and_then(|key| inner.sessions.get(&key).map(|slot| (key, slot.clone())))
321    {
322        if let Some(reason) = reason {
323            if resource != "detached" && resource != "closed" {
324                inner.backend.close(&slot.session, reason, false)?;
325                // Persist the close only after the host accepted it. A failed
326                // pane close must leave the row attached and the slot present,
327                // so a later recycle or shutdown can retry the same resource.
328                // Writing `closed` first makes every later path skip the only
329                // close attempt and strands the host pane permanently.
330                feed_resource_closed(&inner.bridge, &inner.store, task_id)?;
331            }
332            // The agent goes with the resource: this path closes a session whose work an
333            // operator ended or whose backend faulted, and the slot below leaves the map in
334            // the same breath. Without this feed the tuple keeps the agent phase its last
335            // beat reported, and `project` answers `working` for a `cancelled` or `failed`
336            // task whenever the agent is not `Gone` — so the mirrored row read `working` for
337            // a session the client had already closed, which is what `onlyne sessions` and
338            // the board showed beside a ledger row that had settled. The reconnect sweep
339            // feeds the same event for the same ending.
340            if let Err(error) = feed_agent_gone(&inner.bridge, &inner.store, task_id) {
341                tracing::warn!(
342                    task = %task_id,
343                    error = %error,
344                    "agent-gone projection failed for a closed session"
345                );
346            }
347            inner.bridge.untrack_live(task_id);
348            // The delivery row this session is still holding is answered here,
349            // while the slot that holds its handle is still in the map. The close
350            // takes the slot, and the handle goes with it: a row this client
351            // never answered is a row the server still reads as owed, so the
352            // pull that would have taken it passes and the release of a session
353            // the server judges gone hands it back to the queue — which
354            // dispatches the task a second time and runs it again, as the live
355            // run showed for a task whose work the close had already ended.
356            //
357            // The handle is taken, so the row is answered once: the control
358            // settle watchdog and the reconnect sweep refuse the same handle
359            // the same way, and whichever of the three runs first spends it and
360            // the others find nothing left to answer.
361            let held = inner
362                .sessions
363                .get_mut(&key)
364                .and_then(|slot| slot.msg_id.take());
365            if let Some(msg_id) = held {
366                store_ack(
367                    inner,
368                    AckArgs {
369                        msg_id,
370                        op_id: None,
371                        accepted: false,
372                        reason: Some(close_refusal(reason).to_string()),
373                    },
374                );
375            }
376            forget_tools_binding(inner, &key);
377            inner.sessions.remove(&key);
378        } else {
379            let session_id = inner
380                .sessions
381                .get(&key)
382                .map(|slot| slot.session.task_id.clone())
383                .unwrap_or_else(|| task_id.to_string());
384            let keeps = keeps_idle(inner, &key);
385            let now = Instant::now();
386            if let Some(session) = inner.sessions.get_mut(&key) {
387                session.task_id = None;
388                session.ready = false;
389                session.idle_since = Some(now);
390            }
391            // The binding goes with the delivery. A session that serves several
392            // deliveries in turn serves none between them, and both the hello
393            // claim and the mirrored row read the open binding to say which
394            // delivery a session is on.
395            if let Err(error) = inner.store.release_binding(&session_id, task_id) {
396                tracing::warn!(
397                        session = %session_id,
398                        task = %task_id,
399                        error = %error,
400                        "a settled session's binding was not handed back"
401                );
402            }
403            if keeps {
404                // A `task` or `role` session outlives the delivery that opened
405                // it: that is what the scope is for. It stays idle, with its
406                // process and its row, until the scope sends it the next
407                // delivery or `idle_close` releases the process.
408                tracing::info!(
409                        session = %session_id,
410                        task = %task_id,
411                        "delivery settled; the session stays idle for its scope's next one"
412                );
413            } else {
414                let mut pending = Vec::new();
415                retire_idle_locked(
416                    inner,
417                    &key,
418                    crate::backend::CloseReason::Completed,
419                    &mut pending,
420                );
421                // This function answers to its caller's lock, which is already held
422                // across the store work above, so the close it owes cannot wait for an
423                // unlock this scope cannot perform. Running it here is the same
424                // under-lock call the path has always made; the sweeps that CAN leave
425                // the lock — `retire_dropped_ghosts`, `reclaim_exited_resources`,
426                // `close_all` — are the ones that now defer it.
427                close_retired(pending);
428            }
429        }
430    }
431    inner.stall.forget(task_id);
432    if reason.is_some() && resource == "detached" {
433        inner
434            .store
435            .note_alert(format!("session recycled {task_id}"));
436    }
437    Ok(())
438}
439
440/// The operator's word a control close stands for, as the refusal that answers
441/// the row its session was still holding.
442///
443/// The words are the ones the settle fallback already writes for the same
444/// commands ([`ControlWord::refusal`]), so a row reads the same thing whichever
445/// door refused it, and each names the command that was given rather than the
446/// verdict that command left behind.
447fn close_refusal(reason: crate::backend::CloseReason) -> &'static str {
448    match reason {
449        // The close an operator's `cancel` runs: the `ControlOp::Cancel` arm of
450        // `on_control`, which is also how the server asks for `repair close` and
451        // `repair fail` to reach this client.
452        crate::backend::CloseReason::Cancelled => ControlWord::Cancel.refusal(),
453        // The close an operator's `recycle` runs: the `ControlOp::Recycle` arm
454        // of `on_control`.
455        crate::backend::CloseReason::Operator => ControlWord::Recycle.refusal(),
456        // No control command reaches this branch with another reason: a
457        // `completed`, `fault` or `replaced` close retires an idle slot, and a
458        // `shutdown` close runs `close_all`, neither of which comes through
459        // here. The word is the server's own for a close that named no command
460        // — `repair close` without a `--reason` writes it on the task's row.
461        _ => "operator close",
462    }
463}
464
465/// Close every live session's resource with `reason` and forget the slots.
466///
467/// This is the shutdown path: a stopped client must not leave resources behind
468/// that only it can address, and each backend's own record of the resource —
469/// the Orca tab map included — ends with the session. `budget` bounds the whole
470/// sweep, because an operator's SIGTERM must not turn into a hang while a slow
471/// backend CLI exits; whatever the budget cuts off is reported and dropped
472/// anyway.
473///
474/// The order is deliberate: every slot is taken off this client's books under the
475/// dispatch lock — untracked and removed, so nothing still reads a live session —
476/// and the closes it collected run once the lock is let go. A shutdown that
477/// closed each pane while holding that lock spent the whole budget on the host,
478/// with the lock no adapter frame could reach.
479pub fn close_all(state: &DispatchState, reason: crate::backend::CloseReason, budget: Duration) {
480    let started = Instant::now();
481    let pending: Vec<PendingClose> = {
482        let mut inner = state.inner.lock();
483        let sessions: Vec<(String, SessionRef)> = inner
484            .sessions
485            .iter()
486            .map(|(key, slot)| (key.clone(), slot.session.clone()))
487            .collect();
488        let mut pending = Vec::with_capacity(sessions.len());
489        for (key, session) in sessions {
490            inner.bridge.untrack_live(&session.task_id);
491            forget_tools_binding(&mut inner, &key);
492            inner.sessions.remove(&key);
493            pending.push(PendingClose {
494                backend: Arc::clone(&inner.backend),
495                session,
496                reason,
497            });
498        }
499        pending
500    };
501    for close in pending {
502        if started.elapsed() > budget {
503            tracing::warn!(
504                task = %close.session.task_id,
505                "shutdown close budget reached; the resource is left behind"
506            );
507            continue;
508        }
509        if let Err(error) = close.backend.close(&close.session, close.reason, false) {
510            tracing::warn!(
511                task = %close.session.task_id,
512                error = %error,
513                "session close failed during shutdown"
514            );
515        }
516    }
517}
518
519/// Which way a session's window closed.
520///
521/// Two arms reach the same verdict through different facts, and an operator reading one
522/// aggregate log line cannot tell them apart: one means the connection ended and the agent
523/// stayed away, the other means the connection is still up while nothing the client accepts
524/// arrives over it. The words exist so the log can name which reading retired a session.
525#[derive(Clone, Copy, Debug, PartialEq, Eq)]
526pub enum RetirementArm {
527    /// The plugin connection ended and `[client] reconnect_grace_secs` expired.
528    Dropped,
529    /// The connection stayed up while the session went quiet past the heartbeat window.
530    Silent,
531}
532
533impl RetirementArm {
534    pub fn word(self) -> &'static str {
535        match self {
536            Self::Dropped => "reconnect_grace",
537            Self::Silent => "heartbeat_silence",
538        }
539    }
540}
541
542/// One session the sweep retired, with what it read on the way.
543pub struct Retired {
544    /// The retired session's own id, which a client-held session shares with its task.
545    pub session_id: String,
546    /// The arm that decided it.
547    pub arm: RetirementArm,
548    /// Seconds since this session's last accepted frame.
549    pub quiet_secs: u64,
550    /// Seconds since its connection ended, when one did.
551    pub away_secs: Option<u64>,
552}
553
554/// How long a session has gone without a frame this client accepted.
555fn quiet_secs(slot: &SessionSlot, now: Instant) -> u64 {
556    slot.last_beat
557        .map(|beat| now.saturating_duration_since(beat).as_secs())
558        .unwrap_or(u64::MAX)
559}
560
561/// How long a session's connection has been gone, for a session still waiting on one.
562fn away_secs(slot: &SessionSlot, now: Instant) -> Option<u64> {
563    slot.dropped_at
564        .map(|left| now.saturating_duration_since(left).as_secs())
565}
566
567impl DispatchState {
568    /// Retire tracked resources whose stored lifecycle has reached `Exited`.
569    ///
570    /// The periodic readiness tick calls this after completed work becomes an
571    /// idle slot. Task-free sessions with an attached transport stay bound to
572    /// their host resource, and task-free sessions whose agent has left release it.
573    ///
574    /// The session ids come back because the retirement wrote each one of those
575    /// rows and the server mirrors only what this client reports: the resource
576    /// close and, for a completed session, the agent's exit both moved the row
577    /// this tick found, and a publish cannot run under this lock. The caller is
578    /// handed what to publish, the same answer [`DispatchState::retire_dropped_ghosts`]
579    /// gives the sweep above.
580    pub fn reclaim_exited_resources(&self) -> Vec<String> {
581        let mut inner = self.inner.lock();
582        let candidates: Vec<(String, crate::backend::CloseReason)> = inner
583            .sessions
584            .iter()
585            .filter(|(key, slot)| {
586                // A suspended session is not an exited one. Its process is gone
587                // *because* it was released, its conversation is what the
588                // family's next delivery is for, and the row it published says
589                // `idle` for exactly that reason — `binding_task_state` answers
590                // `Pending` for a slot that serves nothing and whose scope keeps
591                // it, so the two derivations of "is this over" disagree here.
592                //
593                // This sweep asks the other one, and it asks about the *work*:
594                // a suspended session's delivery is finished, so it reads
595                // `Exited` and the slot is retired within a tick. The family's
596                // next delivery then finds nothing to resume and opens a second
597                // conversation for one chain — which is the failure this filter
598                // existed to prevent, reached through the sweep that was meant
599                // to clean up after a session that was already gone.
600                !slot.suspended
601                    && slot.task_id.is_none()
602                    && session_exited(&inner, &slot.session.task_id)
603                    && !has_attached_transport(&inner, key, slot)
604            })
605            .filter_map(|(key, slot)| {
606                stored_close_reason(&inner, &slot.session.task_id)
607                    .map(|reason| (key.clone(), reason))
608            })
609            .collect();
610        let mut retired: Vec<String> = Vec::new();
611        let mut pending: Vec<PendingClose> = Vec::new();
612        for (key, reason) in candidates {
613            // The id that travels is the session's own, the one whose row the
614            // retirement below is about to write.
615            let Some(task_id) = inner
616                .sessions
617                .get(&key)
618                .map(|slot| slot.session.task_id.clone())
619            else {
620                continue;
621            };
622            if retire_idle_locked(&mut inner, &key, reason, &mut pending) {
623                retired.push(task_id);
624            }
625        }
626        // The tick that found a batch of finished sessions owes a host close for
627        // each, and this sweep runs every 250 ms: on the lock, one slow backend
628        // would spend the whole window and every adapter frame queued behind it.
629        drop(inner);
630        close_retired(pending);
631        retired
632    }
633
634    /// Retire the sessions whose plugin connection dropped and never came back, or
635    /// whose connection stayed up while they went quiet, and answer which ones left
636    /// and why.
637    ///
638    /// A connection that ends without a `detach` frame leaves its session tracked
639    /// so an agent that restarts inside `[client] reconnect_grace_secs` finds the
640    /// resource it was using. That promise has to expire: a process that is
641    /// simply gone would otherwise hold a slot, a projected `idle` row, and a live
642    /// host resource forever, and on a role with `max_sessions = 1` it stops every
643    /// later delivery. The window answers for the agent itself, so a session still
644    /// bound to a task goes with it: the plugin connection that would have
645    /// reported the ending is the one that dropped. The agent-gone feed is what
646    /// says the process left — the session's own tuple reaches `Exited` through
647    /// `AgentPhase::Gone` rather than through a task result — and the reason the
648    /// backend is handed is the one `grace_close_reason` reads off what the slot
649    /// still owes.
650    ///
651    /// What the slot owed is settled too: the task a bound session was serving
652    /// ends `failed` here, because the agent that would have reported its ending
653    /// is the one that left. A task with no verdict stays open for the server to
654    /// re-offer and for `open_tasks` to keep reading, and no later caller exists
655    /// to write one.
656    ///
657    /// A slot this client holds read-only is not this sweep's to end, agent gone
658    /// or not: the session id it would feed is the task id, so the ghost's death
659    /// would take the live session's mirror and its delivery row down with it.
660    /// That retirement belongs to `retire_revived`, which runs when the
661    /// completion that answers the held connection merges.
662    ///
663    /// The window has a second way to open, and it is the one a socket cannot
664    /// report: a plugin whose event loop is blocked keeps its connection and
665    /// stops beating, so no socket ends and no clock this sweep could read
666    /// before moved. What such a session leaves behind is a stamp going stale
667    /// while its task stays bound and unsettled, and that is the reading this
668    /// sweep takes now. It is the same window and the same verdict — one clock,
669    /// one retirement, no second threshold beside `[client]
670    /// reconnect_grace_secs` and no fault row of the kind `stall_report_secs`
671    /// records and leaves behind.
672    ///
673    /// The sessions' own ids come back rather than a count, because a retirement
674    /// still owes the server the session's own ending: it is the only writer left
675    /// for that task, and a mirror nobody tells keeps that session's last reading —
676    /// `working`, for one that had beaten — until the server's own observer records
677    /// a fault about it. The publish is
678    /// [`sync_session`](crate::session::dispatch::sync_session)'s, which is the
679    /// report an ordinary ending travels on, and it cannot run under this lock —
680    /// so the caller is handed what to publish instead of a second writer being
681    /// invented here.
682    pub fn retire_dropped_ghosts(&self, now: Instant, grace_secs: u64) -> Vec<Retired> {
683        if grace_secs == 0 {
684            return Vec::new();
685        }
686        let window = Duration::from_secs(grace_secs);
687        let mut inner = self.inner.lock();
688        // The arm is decided from the pair of readings here, while both still
689        // describe the slot: which fact closed the window is what the operator
690        // has to be able to tell apart once the retirement itself is a line in
691        // the log.
692        let due: Vec<(String, RetirementArm)> = inner
693            .sessions
694            .iter()
695            .filter_map(|(key, slot)| {
696                let arm = if dropped_past_window(&inner, key, slot, now, window) {
697                    RetirementArm::Dropped
698                } else if silent_past_window(&inner, key, slot, now) {
699                    RetirementArm::Silent
700                } else {
701                    return None;
702                };
703                Some((key.clone(), arm))
704            })
705            .collect();
706        let mut retired: Vec<Retired> = Vec::new();
707        let mut pending: Vec<PendingClose> = Vec::new();
708        for (key, arm) in due {
709            let Some(slot) = inner.sessions.get(&key).cloned() else {
710                continue;
711            };
712            // A slot this client holds read-only is not this sweep's to end, for
713            // the reason the doc above gives. The demotion alone does not decide
714            // it: the held connection that owns the slot can go without the task
715            // ever completing, and a slot nothing owns any more is what this
716            // window is for.
717            if held_read_only(&inner, &key, &slot) {
718                continue;
719            }
720            // A session judged dead on its silence is the one case where the
721            // connection is still there: the socket has not ended and will not
722            // while the agent behind it is blocked, so the client's own verdict
723            // has to take the binding the way the death of the socket would have.
724            // The retirement below refuses a slot an attached transport still
725            // serves, and that refusal is what this unbinding answers: the
726            // verdict has already been reached here, so the binding goes rather
727            // than the death of the socket that would normally take it.
728            if slot.dropped_at.is_none() {
729                unbind_transports(&mut inner, &key, &slot);
730            }
731            let task_id = slot.session.task_id.clone();
732            // Both the feed and the reason name the session's own task, so the id
733            // that travels is the one the row and the resource are keyed by.
734            let reason = grace_close_reason(&inner, &task_id);
735            // The work this session still owed ends here, and this sweep is the
736            // only writer left to say so: the plugin connection that would have
737            // reported the ending is the one that dropped. A task nobody answers
738            // stays `settled_at IS NULL` forever, so `open_tasks` keeps reading
739            // it and the server keeps re-offering a delivery no client can take.
740            // The binding is what the slot owed — a slot past its window with no
741            // task bound owes nothing — and `failed` is the verdict the close
742            // reason above already carries for it. A row an earlier verdict
743            // settled keeps that one: `settle_task` updates only where
744            // `settled_at IS NULL` and answers `false`.
745            //
746            // The write runs before the agent-gone feed and before the binding
747            // hand-back, both of which end this slot's turn through the sweep:
748            // the id is captured here, and the verdict is on disk before the
749            // only handle on it goes away.
750            if let Some(owed) = slot.task_id.clone() {
751                // A verdict written here answers the task for good, so a
752                // `control` command that is still waiting for its plugin's report
753                // has nothing left to authorise: the note goes with the verdict
754                // that outranks it.
755                inner.control_settles.retain(|noted| noted.task_id != owed);
756                if let Err(error) = inner.store.settle_task(&owed, TaskState::Failed) {
757                    tracing::warn!(
758                        task = %owed,
759                        error = %error,
760                        "the task of a retired ghost was not settled"
761                    );
762                }
763                // The delivery handle this session was holding is spent as a
764                // refusal that names the death, and it is the ledger half of the
765                // verdict above. A row left `in_flight` is handed to a pull no
766                // longer — `pull` passes by a row whose ticket is armed, and a
767                // role-level pull's ticket carries no session id for the release
768                // path to match — so nothing would answer for this task until the
769                // link dropped, and an operator reading `onlyne ledger` would see
770                // a session that has been buried as one still holding its
771                // delivery. The reason is this client's own word for a death
772                // (`SESSION_DEAD`), and a refusal is terminal: the work comes back
773                // through `repair retry`, not by itself.
774                let handle = inner
775                    .sessions
776                    .get_mut(&key)
777                    .and_then(|slot| slot.msg_id.take());
778                if let Some(msg_id) = handle {
779                    store_ack(
780                        &inner,
781                        AckArgs {
782                            msg_id,
783                            op_id: None,
784                            accepted: false,
785                            reason: Some(SESSION_DEAD.to_string()),
786                        },
787                    );
788                }
789            }
790            if let Err(error) = feed_agent_gone(&inner.bridge, &inner.store, &task_id) {
791                tracing::warn!(
792                    task = %task_id,
793                    error = %error,
794                    "agent-gone projection failed for a retired ghost"
795                );
796            }
797            // The agent left, so the session owes no task any more: the binding
798            // goes back before the idle retirement takes the slot.
799            if let Some(slot) = inner.sessions.get_mut(&key) {
800                slot.task_id = None;
801            }
802            if retire_idle_locked(&mut inner, &key, reason, &mut pending) {
803                // The id that travels is the session's own, the one whose row was
804                // just fed agent-gone and resource-closed: that row is what the
805                // server mirrors, and its ending is what the caller publishes. The
806                // ages are read off the slot as it stood before the retirement.
807                retired.push(Retired {
808                    session_id: task_id,
809                    arm,
810                    quiet_secs: quiet_secs(&slot, now),
811                    away_secs: away_secs(&slot, now),
812                });
813            }
814        }
815        // Every session this sweep took is off the slot map and fully written down
816        // by now; what remains is the host's own work, and the lock that answers an
817        // agent's next frame is not to be held through it.
818        drop(inner);
819        close_retired(pending);
820        retired
821    }
822}