onlyne_client/session/dispatch/state.rs
1use super::*;
2
3use super::outbound::Outbox;
4use super::projection::{projection_of, stored_task_state};
5use super::transport::names_session;
6use super::turn_end::TurnEndWatch;
7
8#[derive(Clone)]
9pub struct DispatchState {
10 pub(super) inner: Arc<Mutex<DispatchInner>>,
11}
12pub(super) struct DispatchInner {
13 pub(super) role: String,
14 pub(super) workspace: PathBuf,
15 pub(super) command: Vec<String>,
16 /// The drive the role's spec declares (`[client.runtime] drive`). `None`
17 /// until the first `welcome` names it.
18 pub(super) drive: Option<onlyne_config::Drive>,
19 /// The placement this machine resolved, the other half of the drive rule.
20 /// `None` until the run resolves one; a pair the rule refuses is what
21 /// [`super::env::reject_unpaired_runtime`] refuses a session over.
22 pub(super) placement: Option<crate::backend::SessionPlacement>,
23 /// Why a session may not open under this role's drive on this machine, when
24 /// the pair is one the rule refuses. Every delivery the role is offered is
25 /// refused with it: a backend the previous drive left behind must not serve
26 /// work under a policy nobody set.
27 pub(super) runtime_refusal: Option<String>,
28 pub(super) max_sessions: u32,
29 /// The roles a session of this role owes a delivery to, read off the
30 /// server's spec slice (`allowed_targets`): the same list the server gates
31 /// the ACL on is the obligation the completion guard measures a session
32 /// against. Empty is the default and means the role owes nothing.
33 pub(super) required_targets: Vec<String>,
34 /// The role's workspace session policy: which deliveries one session serves,
35 /// and how long an idle one may wait before its process is released (§10).
36 pub(super) session_policy: onlyne_config::SessionPolicy,
37 pub(super) backend: Arc<dyn SessionBackend>,
38 pub(super) store: ClientStore,
39 pub(super) bridge: Bridge,
40 pub(super) sessions: HashMap<String, SessionSlot>,
41 /// Live link, installed by the runloop while the connection is up.
42 pub(super) outbox: Option<Arc<dyn Outbox>>,
43 /// Flag the runloop and the dispatcher share while the link is down.
44 pub(super) accept_new: Arc<AtomicBool>,
45 /// Whether the role holds a ready server link. The runloop owns it, and the
46 /// adapter socket reports it to the `status` verb.
47 pub(super) link_up: Arc<AtomicBool>,
48 /// Aggregate name this role supervises, empty for a plain role.
49 pub(super) cluster_ref: String,
50 /// The server's topology name, read from `welcome.cluster` (the server's own
51 /// `spec.toml [server] name`). A host backend uses it as the address of the
52 /// tree it puts sessions into: a pane host keeps one tree per server root,
53 /// labelled after this name. Empty until the first welcome arrives.
54 pub(super) topology: String,
55 /// The adapter connection serving each session of this role, keyed by the
56 /// session id the plugin mounted with (`ONLYNE_SESSION_ID`). A plugin the
57 /// client spawned names the one session it was spawned for, so a task
58 /// never rides the connection of an earlier one.
59 pub(super) transports: HashMap<String, (AdapterIo, Vec<Capability>)>,
60 /// A plugin that mounted naming no session: an always-running agent
61 /// waiting for this role's next assignment (plan §6 line 285), oldest mount
62 /// first.
63 ///
64 /// A queue and not a single slot, because a role can hold more than one
65 /// always-running agent and the second mount that arrived naming no session
66 /// has an open socket either way. Overwriting the first dropped its
67 /// `AdapterIo` with no accounting of any kind: no release, no log, and no
68 /// bye, so the role silently lost a worker that was waiting to be told.
69 /// `revived` has held several connections the same way from the start.
70 pub(super) parked: Vec<(AdapterIo, Vec<Capability>)>,
71 /// Connections from **hosting** runtimes: ones that declared `open`,
72 /// `suspend` or `close` and own their sessions rather than serving the one
73 /// the client started them for.
74 ///
75 /// This is not the park with a different name. A parked connection is a
76 /// *worker* waiting for one job — `claim_parked_transport` takes the oldest
77 /// off the queue and that connection is spent, which is right for an agent
78 /// the client spawned for a single session. A hosting connection is a shared
79 /// resource: it takes no job off a queue because it is not waiting for one,
80 /// and it serves whatever sessions the role opens. Putting one in `parked`
81 /// would let the first staged session consume it, and the role's other three
82 /// sessions would have nothing to run on.
83 ///
84 /// A role can hold both, and the order they are consulted in is
85 /// `hand_staged`'s: a parked agent was spawned for the work in hand, so it
86 /// is the more specific match, and a hosting connection is the fallback.
87 pub(super) standing: Vec<(AdapterIo, Vec<Capability>)>,
88 /// Zero-activity clock for running tasks. Applied persists refresh it.
89 pub(super) stall: crate::session::stall::StallWatch,
90 /// A plugin connection that mounted a session a live connection already
91 /// serves: the agent that dropped came back after a newer session took the
92 /// task. It is served no state, and the ending it reports is answered like
93 /// any other connection's; the name it mounted with is what that judgement
94 /// reads. The capabilities that mount came with are kept beside it, because
95 /// a session whose live connection goes away hands itself to the first
96 /// connection that was holding for it.
97 pub(super) revived: Vec<(String, AdapterIo, Vec<Capability>)>,
98 /// Tasks whose ending this client asked for with a `control` command.
99 ///
100 /// `recycle` and `cancel` reach the agent as a `notify`, so the completion
101 /// that answers the command races the retirement the same command runs, and
102 /// the row's phase at the moment the frame lands is that race's answer. The
103 /// note is written before the frame leaves, so both orders of the race read
104 /// one authority: the operator asked for this ending, and the settle door
105 /// (`settle.rs`) takes the note as its `SettleAuthority::ControlDriven`.
106 /// `on_out` consumes it, and the reconnect sweep drops the note of every task
107 /// it settles — the other way a command's answer stops coming.
108 ///
109 /// A note nothing answers is the watchdog's, and it is the same note: the
110 /// record of the word, held to [`CONTROL_SETTLE_BOUND`] and settled by the
111 /// tick's own sweep when no report has come to settle it instead.
112 pub(super) control_settles: Vec<ControlNote>,
113 /// 3c's turn-end bookkeeping, keyed by the task whose delivery it belongs to.
114 ///
115 /// Which delivery already spent its one nudge, and which settled at that
116 /// door. It stays here: no column holds it, no frame carries it, and no log
117 /// line reads it, because the ending it remembers is one this client
118 /// witnessed for itself (`docs/v2-CONTRACT.md` §3c).
119 pub(super) turn_end: TurnEndWatch,
120 /// Live `tools` mounts, one binding per session they speak for.
121 ///
122 /// A tools mount holds no process, so it is no transport and no parked
123 /// agent: it lands in neither of those tables. This one exists to answer
124 /// which session a connection speaks for, and nothing else consults it. The
125 /// token itself stays in the session's own slot; an entry here carries the
126 /// slot's key and the connection, never the token
127 /// (`docs/v2-CONTRACT.md` §3b).
128 pub(super) tools_mounts: Vec<(String, AdapterIo)>,
129 /// Connections inside one of their own inbound frames right now.
130 ///
131 /// A frame handler runs to completion before `adapter_socket` answers the
132 /// frame, so a host frame written on the same connection during that handler
133 /// leaves first. The bye sweep in `retire_revived` skips these connections.
134 pub(super) in_frame: Vec<AdapterIo>,
135}
136
137#[derive(Clone)]
138pub struct SessionSlot {
139 pub(super) session: SessionRef,
140 pub(super) task_id: Option<String>,
141 pub(super) ready: bool,
142 /// Payload held until the adapter reports ready, which keeps the ready
143 /// barrier of §6 ahead of the `assign` frame.
144 pub(super) payload: Option<Envelope>,
145 /// Delivery handle, owed back to the server as one `ack`.
146 pub(super) msg_id: Option<String>,
147 /// Sender of the payload this session serves, kept for its `Completion`.
148 pub(super) origin: Option<Principal>,
149 /// The causality of the task this slot serves, read off the envelope that
150 /// arrived with it. A handoff the session reports afterwards is its child,
151 /// which is what `onlyne handoff` computes too. The whole link is kept here
152 /// so the family id and the family's figures travel with the child, the
153 /// depth included.
154 pub(super) causality: Causality,
155 /// When the connection that would have sent this session's next heartbeat
156 /// last left it: at birth for a session a plugin still has to mount, and
157 /// again each time a connection ends without a `detach` frame or says
158 /// goodbye while the session still owes a task. `None` while a connection is
159 /// attached, which is what clears it, and for a self-driven session, which
160 /// answers no adapter socket at all. The reconnect grace of `[client]
161 /// reconnect_grace_secs` reads it: an agent that comes back inside the window
162 /// clears it and keeps its session.
163 pub(super) dropped_at: Option<Instant>,
164 /// When a frame this session sent was last accepted.
165 ///
166 /// The other half of the same death window, and the half a socket cannot
167 /// answer for: a plugin whose event loop is blocked keeps its connection and
168 /// stops beating, so `dropped_at` stays `None` and no socket ever ends. Set
169 /// when the session is staged and refreshed wherever a frame of its own is
170 /// accepted, so it reads as the moment the agent last proved it was alive.
171 /// The reconnect sweep compares it against the protocol's heartbeat cadence.
172 pub(super) last_beat: Option<Instant>,
173 /// Whether this session's task has been taken by a newer session, leaving
174 /// this slot served only by a connection that came back for it. A read-only
175 /// slot is handed no assignment and no note: the work belongs to the session
176 /// that took the task, and this one answers only for how its own turn ended.
177 pub(super) read_only: bool,
178 /// The task family this session was opened for, under the `task` scope. The
179 /// scope hands that family's later deliveries here instead of opening a
180 /// second conversation for one chain, so this is the key the resolution
181 /// reads and the reason a session outlives the delivery that opened it.
182 pub(super) family: Option<String>,
183 /// When this session stopped serving a delivery, while it is still alive.
184 /// The role's `idle_close` bound runs from here, and a slot serving a
185 /// delivery has none.
186 pub(super) idle_since: Option<Instant>,
187 /// A session whose process has been released and the conversation lives in
188 /// the runtime's own store: the next delivery bound to this session resumes
189 /// it. A suspended slot spends no capacity, holds no transport, and answers
190 /// no frame.
191 pub(super) suspended: bool,
192 /// The capability token a `tools` mount presents to speak for this session.
193 ///
194 /// Minted when the session opens and handed to the session's own drive
195 /// through its spawn spec; it lives here, in the session's own state, so
196 /// nothing outside the session may hand it out, it dies with the slot, and
197 /// a session that reopens gets a new one. It is a capability, so it never
198 /// reaches a log line, a fault, or a ledger row
199 /// (`docs/v2-CONTRACT.md` §3b, `AGENTS.md` §8).
200 pub(super) tools_token: String,
201 /// The roles this session delivered to since it opened.
202 ///
203 /// A send the client carried is a handoff, whatever envelope kind it was,
204 /// and this set is the evidence the relay guard reads at the session's next
205 /// completion. It belongs to the session rather than to one delivery: the
206 /// family's obligation outlives the task that opened it
207 /// (`plugins/onlyne-agent-pi`, `guards.rs`).
208 pub(super) delivered_roles: BTreeSet<String>,
209 /// When this session was opened. The order a `role` pool hands its sessions
210 /// out in: the one that has waited longest takes the delivery.
211 pub(super) opened_at: Instant,
212 /// The argv this session's runtime was started with, rendered when the
213 /// session was born. Resuming it starts this command again rather than a
214 /// freshly rendered one, because the command carries the runtime's own key
215 /// for the conversation and may interpolate the delivery into it.
216 pub(super) command: Vec<String>,
217 /// What a hosting runtime called the conversation it opened, opaque to this
218 /// client: stored beside the session, handed back on the next `open` for the
219 /// same family, and read by nothing here. A session without one is a
220 /// conversation the runtime cannot find again, and the client starts a fresh
221 /// one rather than composing a history summary to stand in for it — a summary
222 /// this process wrote is context the model did not produce, which is the rot
223 /// the whole delivery design exists to avoid.
224 pub(super) resume_handle: Option<String>,
225 /// This session's scope keeps it alive after a delivery settles.
226 ///
227 /// `oneshot` does not: that session's own id is the delivery that opened it,
228 /// and when that delivery is answered the session is over. A `task` or
229 /// `role` session serving nothing between deliveries is idle, which is a
230 /// live session rather than an exit, and the verdict of a delivery it has
231 /// already finished says nothing about it.
232 pub(super) keeps_idle: bool,
233}
234
235/// One operator's word this client is still waiting to see answered.
236///
237/// The command is a `notify` the plugin may or may not live to answer, so the
238/// note is what this client knows on its own: which task the word named, when it
239/// was given, and what the operator said. It is not a second settle path — the
240/// note is the authority of the one settle door, which `settle.rs` reads as
241/// `SettleAuthority::ControlDriven` — and it is the record the watchdog holds a
242/// word to once no report arrives at all.
243#[derive(Clone, Debug, PartialEq, Eq)]
244pub struct ControlNote {
245 /// The task the word named.
246 pub task_id: String,
247 /// When the word was given. The watchdog's bound runs from here.
248 pub noted_at: Instant,
249 /// What the operator said.
250 pub word: ControlWord,
251}
252
253/// The operator's word one note is the record of.
254///
255/// `cancel` and `recycle` are the two commands that reach a live task, and each
256/// word carries both halves of what a note is for: the ending it gives the work,
257/// and the string a refusal of that task's delivery row is written with. The note
258/// holds the word rather than the verdict it stands for, because the verdict does
259/// not name the word back — `failed` is only the fallback a `recycle` leaves when
260/// nothing answers it — and the column the refusal lands in has to read what the
261/// operator actually said.
262#[derive(Clone, Copy, Debug, PartialEq, Eq)]
263pub enum ControlWord {
264 /// `cancel`: the work ends now, and the task reads `cancelled`.
265 Cancel,
266 /// `recycle`: the plugin is asked for its own ending. It prescribes no
267 /// outcome, so a word nothing answers leaves the task `failed`.
268 Recycle,
269}
270
271impl ControlWord {
272 /// The outcome this word stands for: the verdict a settle answering the note
273 /// files when the plugin reports none of its own.
274 pub fn outcome(self) -> Outcome {
275 match self {
276 Self::Cancel => Outcome::Cancelled,
277 Self::Recycle => Outcome::Failed,
278 }
279 }
280
281 /// The word as the refusal that names it reads on the wire.
282 ///
283 /// `operator cancel` and `operator recycle` stand in the same column as the
284 /// `operator close` and `operator ack` an operator's other verbs already wrote
285 /// there, and each names the command that was given rather than the verdict it
286 /// left behind.
287 pub fn refusal(self) -> &'static str {
288 match self {
289 Self::Cancel => "operator cancel",
290 Self::Recycle => "operator recycle",
291 }
292 }
293}
294
295/// How long this client holds an operator's word open before settling it.
296///
297/// The words this bounds are `cancel` and `recycle`, and each one is a `notify`
298/// the plugin answers with a frame on the connection it already serves: a plugin
299/// that ends its turn to answer has answered inside one
300/// [`HEARTBEAT_INTERVAL`](crate::session::dispatch::HEARTBEAT_INTERVAL), and the
301/// request round trip the adapter bounds itself with
302/// ([`REQUEST_TIMEOUT`](crate::session::dispatch::REQUEST_TIMEOUT)) is well past
303/// that. Three intervals leaves a plugin that is stalled but still alive two
304/// missed beats before this client decides the word went unanswered — it is the
305/// window the reconnect sweep reads an agent's silence through, so the two
306/// readings agree — and it is half the sixty-second default of `[client]
307/// reconnect_grace_secs`. An operator watching a stuck row gave up on the live
308/// run in seconds and reached for `onlyne repair fail`; a minute would lose to
309/// that, and this does not.
310///
311/// A constant rather than a config key on purpose: what it bounds is not a policy
312/// an operator tunes, it is the point past which this client's own record of the
313/// word outlives the plugin that was asked to answer it.
314pub const CONTROL_SETTLE_BOUND: Duration = Duration::from_secs(HEARTBEAT_INTERVAL.as_secs() * 3);
315
316/// The notes whose operator's word has gone unanswered past
317/// [`CONTROL_SETTLE_BOUND`].
318///
319/// The reading consumes nothing. The caller settles each note through
320/// [`take_controlled_settle`](DispatchState::take_controlled_settle), the one
321/// door that spends a note, so a completion that answers a word between this read
322/// and that call takes the note first and the task needs no verdict from the
323/// sweep. A note stamped ahead of `now` is not due: an elapsed window is the only
324/// reading this makes.
325pub(super) fn due_control_settles(inner: &DispatchInner, now: Instant) -> Vec<ControlNote> {
326 inner
327 .control_settles
328 .iter()
329 .filter(|note| now.saturating_duration_since(note.noted_at) >= CONTROL_SETTLE_BOUND)
330 .cloned()
331 .collect()
332}
333
334/// Stamp the moment one session last had a frame of its own accepted.
335///
336/// The stamp is the liveness half of the reconnect sweep: a socket that is still
337/// up and still attached proves the connection survived, not that the agent
338/// behind it did. A plugin whose event loop is blocked keeps its socket and
339/// stops beating, and nothing but this stamp says so.
340///
341/// Called where a frame is accepted rather than where one arrives — the mount
342/// that binds the connection, the ready barrier, and each beat the reducer took
343/// — because a frame the client refused moved no state and is the evidence of
344/// nothing. A frame this client refused on other grounds still buys the stamp:
345/// the silence arm asks whether an agent lives behind the slot, and it reads
346/// this stamp only for a session whose task is still bound and unsettled, so a
347/// settled or taken task is never the stamp's question and a demoted slot still
348/// retires on its own schedule. What a refused frame never buys is a write: the
349/// dimensions stay where the connection serving the session left them.
350pub(super) fn note_beat(inner: &mut DispatchInner, task_id: &str, now: Instant) {
351 let Some(key) = slot_key_named(inner, task_id) else {
352 return;
353 };
354 if let Some(slot) = inner.sessions.get_mut(&key) {
355 slot.last_beat = Some(now);
356 }
357}
358
359pub(super) fn has_attached_transport(inner: &DispatchInner, key: &str, slot: &SessionSlot) -> bool {
360 inner
361 .transports
362 .keys()
363 .any(|session_id| names_session(key, slot, session_id))
364}
365
366/// Move one session's row onto the generation after the one it holds.
367///
368/// Two shapes leave a row behind a generation that is gone, and both close the
369/// gap the same way. A plugin that comes back to a session this client still
370/// serves reports from a new process, so its sequence starts at its base again.
371/// A session this client stages onto a task that already carries a row is born
372/// onto the record the session that last served the task left, watermark and
373/// all. The watermark is the counter this client's own feeds and the plugin's
374/// beats share — `upsert_session` writes the event's version into it, and
375/// `next_version` allocates the stored one plus one — so an inherited one drops
376/// every frame of the new reporter as a stale duplicate, and no dimension ever
377/// moves on that row again. A same-generation rebase is not available: the
378/// reducer's no-op detection compares the version-free tuple, so an event that
379/// moved only the watermark would be ignored and the row would keep the old one.
380/// The generation is what moves, and the sequence starts again under it.
381///
382/// The caller composes `body`, the content of the new generation, and the caller
383/// attests the old one dead: `Supersede` is the reducer's vocabulary for that
384/// pair, and each caller is the authority on the fact it asserts. Answers `None`
385/// when the task carries no row yet — a first dispatch, which has nothing to
386/// rebase because `feed_created` seeds the row instead.
387pub(super) fn rebase_generation(
388 inner: &DispatchInner,
389 task_id: &str,
390 body: impl FnOnce(&Observation) -> Observation,
391) -> anyhow::Result<Option<Verdict>> {
392 let Some(row) = inner.store.get_session(task_id)? else {
393 return Ok(None);
394 };
395 let stored = stored_observation(&inner.store, Some(&row));
396 let event = LifecycleEvent::Supersede {
397 v: Version::new(stored.version.generation.saturating_add(1), 0),
398 old_generation_dead: true,
399 body: body(&stored),
400 };
401 Ok(Some(apply_persist(
402 &inner.bridge,
403 &inner.store,
404 task_id,
405 &event,
406 )?))
407}
408
409/// The task one slot answers for: its current binding, or the task its session
410/// was spawned for while no task is bound.
411pub(super) fn slot_task(slot: &SessionSlot) -> String {
412 slot.task_id
413 .clone()
414 .unwrap_or_else(|| slot.session.task_id.clone())
415}
416
417/// The slot one adapter mount names, answered as its key.
418pub(super) fn slot_key_named(inner: &DispatchInner, session_id: &str) -> Option<String> {
419 inner
420 .sessions
421 .iter()
422 .find(|(key, slot)| names_session(key, slot, session_id))
423 .map(|(key, _)| key.clone())
424}
425
426/// The slot serving one task, preferring the one that still holds delivery
427/// rights.
428///
429/// Two slots answer to one task only when a session came back for a task a newer
430/// session already took, and `HashMap` order decides which one a search reaches
431/// first. Every handle that belongs to the task — its delivery `msg_id` above
432/// all — has to land on the session that is actually serving it, or the ack the
433/// retry earns would be written into a slot that can never answer and the server
434/// row would sit in flight.
435pub(super) fn slot_key_serving_task(inner: &DispatchInner, task_id: &str) -> Option<String> {
436 let mut bound = None;
437 for (key, slot) in inner.sessions.iter() {
438 if slot.task_id.as_deref() != Some(task_id) {
439 continue;
440 }
441 if !slot.read_only {
442 return Some(key.clone());
443 }
444 if bound.is_none() {
445 bound = Some(key.clone());
446 }
447 }
448 bound
449}
450
451/// One plugin connection held inside an inbound frame it is still being answered.
452///
453/// The bye sweep in `retire_revived` leaves such a connection alone, which keeps a
454/// bye behind the response to the frame. A plugin's bye handler drops the socket and
455/// rejects every request awaiting an answer, so a bye that overtakes the response
456/// turns work the ledger already holds into a failure the agent reports again.
457#[must_use = "the connection stops being held as soon as the guard is dropped"]
458pub struct FrameGuard<'a> {
459 pub(super) state: &'a DispatchState,
460 pub(super) io: AdapterIo,
461}
462
463impl Drop for FrameGuard<'_> {
464 fn drop(&mut self) {
465 self.state
466 .inner
467 .lock()
468 .in_frame
469 .retain(|held| !held.same_connection(&self.io));
470 }
471}
472
473pub(super) fn render_tokens(tokens: &[String], session: &str, task: &str) -> Vec<String> {
474 tokens
475 .iter()
476 .map(|token| token.replace("{session}", session).replace("{task}", task))
477 .collect()
478}
479
480/// The session a `tools` mount token spoke for, as the mount needs it.
481///
482/// The token *is* the binding (`docs/v2-CONTRACT.md` §3b): the role this mount
483/// speaks for is the client's own, and the session and its generation come from
484/// the client's record rather than from a field the caller supplies. A caller
485/// that could name its own session could speak for one it never held.
486#[derive(Clone, Debug, PartialEq, Eq)]
487pub struct ToolsSession {
488 /// The session's own id, the spelling the mount answers for.
489 pub session_id: String,
490 /// The generation the session was opened under, the third half of the
491 /// `(role, session_id, generation)` record.
492 pub generation: u64,
493 /// The delivery the session serves right now, when one is open.
494 pub task_id: Option<String>,
495}
496
497/// Mint the capability token one session's `tools` mount presents.
498///
499/// The value is random and per-session: it is a capability, so it is handed to
500/// the session's own drive alone and never logged, faulted, or written down.
501pub(super) fn mint_tools_token() -> String {
502 uuid::Uuid::new_v4().to_string()
503}
504
505/// The slot one live session's tools token names.
506///
507/// A token names a session while that session lives and serves: the slot
508/// carries both, so a token whose slot has retired — or whose tuple and verdict
509/// project `Exited` — names nothing, and neither does one whose session is
510/// suspended (no process to mount from) or read-only (a newer session took the
511/// task, and a read-only slot serves no state). The mount that presents such a
512/// token is refused. The scan is small (one role's slots) and the token never
513/// leaves this process.
514pub(super) fn slot_key_for_token(inner: &DispatchInner, token: &str) -> Option<String> {
515 if token.is_empty() {
516 return None;
517 }
518 inner
519 .sessions
520 .iter()
521 .find(|(_, slot)| slot.tools_token == token && token_names_session(inner, slot))
522 .map(|(key, _)| key.clone())
523}
524
525/// Whether one slot is a session a tools token may still name.
526///
527/// A token names a session while that session lives and serves: an `Exited`
528/// projection is the session over, a suspended session has no process to mount
529/// from, and a read-only slot is a session a newer one took the task from, which
530/// serves no state. The three refusals are one predicate so the handshake, the
531/// per-frame gate, and the stamped scope cannot disagree about which tokens
532/// still name anything.
533pub(super) fn token_names_session(inner: &DispatchInner, slot: &SessionSlot) -> bool {
534 !slot.suspended && !slot.read_only && !slot_exited(inner, slot)
535}
536
537/// What one live tools connection speaks for.
538///
539/// Everything a tools frame is stamped with comes from here, so the fields are
540/// the session's own record rather than anything the mount supplied: the
541/// delivery its `report` and `handoff` frames name, and the session a `send`
542/// frame leaves from (`docs/v2-CONTRACT.md` §3b). A `send` starts a family of
543/// its own, so no chain of this session's rides out with it.
544#[derive(Clone, Debug)]
545pub struct ToolsScope {
546 /// The session's own id, the key its slot and row are held under.
547 pub session_id: String,
548 /// The delivery the session serves right now, when one is open.
549 pub task_id: Option<String>,
550}
551
552/// The mount-side facts of one session slot.
553pub(super) fn tools_session_of(inner: &DispatchInner, key: &str) -> Option<ToolsSession> {
554 inner.sessions.get(key).map(|slot| ToolsSession {
555 session_id: slot.session.task_id.clone(),
556 generation: slot.session.generation,
557 task_id: slot.task_id.clone(),
558 })
559}
560
561/// Drop the tools binding one session held, when the session goes away.
562///
563/// A binding holds a live connection handle, and a slot that leaves this
564/// client's books leaves nothing for the mount to speak for: the entry goes
565/// with it so a retired session's socket is not kept open by this table.
566pub(super) fn forget_tools_binding(inner: &mut DispatchInner, key: &str) {
567 inner.tools_mounts.retain(|(held, _)| held != key);
568}
569
570/// Whether the session of one task is over.
571///
572/// The answer is derived, never read: `client.db` holds the session tuple the
573/// reducer wrote under the `(generation, seq)` gate, the task table holds the
574/// verdict the agent filed, and `project` is the only thing that says whether
575/// that pair is `exited` — so a session whose work finished cleanly leaves here,
576/// and one whose agent went away leaves here too. The rows stay readable after
577/// the session ends. A session with no row yet is live: it exists as a spawned
578/// resource alone, with nothing derived from it.
579pub(super) fn session_exited(inner: &DispatchInner, task_id: &str) -> bool {
580 let Some(row) = inner.store.get_session(task_id).ok().flatten() else {
581 return false;
582 };
583 projection_of(&row, stored_task_state(inner, task_id)).lifecycle == Lifecycle::Exited
584}
585
586/// How the delivery one slot is serving ended, read from the task's own record.
587///
588/// A slot serving nothing answers `Pending`. That is not a guess about work in
589/// flight: no delivery of this session is open for a verdict right now, and
590/// `project` reads the pair as a live session rather than as an exit — which is
591/// what an idle `task` or `role` session is. The verdict of a delivery it has
592/// already finished belongs to that delivery, and says nothing about a session
593/// the scope kept open for the family's next one.
594pub(super) fn binding_task_state(inner: &DispatchInner, slot: &SessionSlot) -> TaskState {
595 match slot.task_id.as_deref() {
596 Some(task_id) => stored_task_state(inner, task_id),
597 None if slot.keeps_idle => TaskState::Pending,
598 // A session no scope keeps has one delivery behind it and no next one:
599 // its own id is that delivery's record, and its verdict is the answer.
600 None => stored_task_state(inner, &slot.session.task_id),
601 }
602}
603
604/// Whether the session one slot holds is over.
605///
606/// The slot's own id addresses the row: a session that has served several
607/// deliveries has one row and a moving binding, and the row is the session's.
608/// A slot with no row at all has not been written about yet and is live.
609pub(super) fn slot_exited(inner: &DispatchInner, slot: &SessionSlot) -> bool {
610 let Some(row) = inner
611 .store
612 .get_session(&slot.session.task_id)
613 .ok()
614 .flatten()
615 else {
616 return false;
617 };
618 projection_of(&row, binding_task_state(inner, slot)).lifecycle == Lifecycle::Exited
619}
620
621/// Sessions that hold the role's concurrency: the staged slots whose tuple and
622/// task verdict have not projected to `Exited`.
623///
624/// §5's `max_sessions` caps concurrent sessions, and a session the reducer has
625/// ended answers no task, so it stops spending capacity the moment its tuple and
626/// verdict derive `exited` — whether that came from a completion report the task
627/// table settled or from the observation a plugin sends as its last heartbeat.
628/// The rows stay in `client.db` and stay queryable; a session with no row at all
629/// is live.
630///
631/// A suspended session spends nothing: the scope's rules count active sessions,
632/// and a session whose process has been released is what freeing a slot means.
633pub(super) fn live_sessions(inner: &DispatchInner) -> usize {
634 inner
635 .sessions
636 .values()
637 .filter(|slot| !slot.suspended && !slot_exited(inner, slot))
638 .count()
639}