Skip to main content

onlyne_client/session/dispatch/
slots.rs

1use super::*;
2
3use super::outbound::store_ack;
4use super::projection::{stored_task_state, task_state_of};
5use super::retire::stored_close_reason;
6use super::state::{
7    ControlNote, ControlWord, DispatchInner, DispatchState, ToolsScope, ToolsSession,
8    due_control_settles, live_sessions, session_exited, slot_key_for_token, slot_key_named,
9    slot_key_serving_task, token_names_session, tools_session_of,
10};
11use super::transport::names_session;
12
13impl DispatchState {
14    pub fn new(
15        role: impl Into<String>,
16        workspace: impl Into<PathBuf>,
17        command: Vec<String>,
18        max_sessions: u32,
19        backend: Arc<dyn SessionBackend>,
20        store: ClientStore,
21    ) -> Self {
22        Self {
23            inner: Arc::new(Mutex::new(DispatchInner {
24                role: role.into(),
25                workspace: workspace.into(),
26                command,
27                drive: None,
28                placement: None,
29                runtime_refusal: None,
30                max_sessions,
31                session_policy: onlyne_config::SessionPolicy::default(),
32                required_targets: Vec::new(),
33                backend,
34                store,
35                bridge: Bridge::new(),
36                sessions: HashMap::new(),
37                outbox: None,
38                accept_new: Arc::new(AtomicBool::new(true)),
39                link_up: Arc::new(AtomicBool::new(false)),
40                cluster_ref: String::new(),
41                topology: String::new(),
42                transports: HashMap::new(),
43                parked: Vec::new(),
44                standing: Vec::new(),
45                stall: crate::session::stall::StallWatch::new(),
46                revived: Vec::new(),
47                control_settles: Vec::new(),
48                tools_mounts: Vec::new(),
49                in_frame: Vec::new(),
50                turn_end: super::turn_end::TurnEndWatch::default(),
51            })),
52        }
53    }
54
55    /// Adopt the role's `[client.session]` policy.
56    pub fn with_session_policy(self, policy: onlyne_config::SessionPolicy) -> Self {
57        self.inner.lock().session_policy = policy;
58        self
59    }
60
61    /// Install the placement this machine resolved. It is the other half of the
62    /// drive rule, so it is recorded even when the pair it makes is refused.
63    pub fn with_placement(self, placement: crate::backend::SessionPlacement) -> Self {
64        self.inner.lock().placement = Some(placement);
65        self
66    }
67
68    /// The drive the role's spec declares, `None` before the first `welcome`.
69    pub fn drive(&self) -> Option<onlyne_config::Drive> {
70        self.inner.lock().drive
71    }
72
73    /// The placement this machine resolved, as the registration spells it.
74    pub fn placement_name(&self) -> Option<&'static str> {
75        self.inner
76            .lock()
77            .placement
78            .map(crate::backend::SessionPlacement::as_str)
79    }
80
81    /// The name of the session backend currently installed.
82    pub fn session_backend(&self) -> &'static str {
83        self.inner.lock().backend.name()
84    }
85
86    /// The role workspace this client serves.
87    pub fn workspace(&self) -> PathBuf {
88        self.inner.lock().workspace.clone()
89    }
90
91    /// Record the drive the role's spec declares, and the sentence a delivery
92    /// meets when that drive cannot run under this machine's placement.
93    ///
94    /// The drive lands whether or not a backend could be installed for it: a
95    /// pair this machine cannot host still has to be *known*, or the next
96    /// delivery would run under the backend the previous drive left behind.
97    pub fn set_drive(&self, drive: onlyne_config::Drive, refusal: Option<String>) {
98        let mut inner = self.inner.lock();
99        inner.drive = Some(drive);
100        inner.runtime_refusal = refusal;
101    }
102
103    /// Install the backend a role's drive selects on this machine.
104    ///
105    /// Answers whether it landed. A drive cannot move while this role holds live
106    /// sessions: their panes, tabs, and children are the installed backend's to
107    /// close and to probe, and a backend that never opened a resource cannot
108    /// answer for it. So the move waits — the caller keeps the old drive
109    /// recorded, and the next `welcome` or spec reload tries again once the
110    /// sessions are gone.
111    pub fn set_backend(&self, backend: Arc<dyn SessionBackend>) -> bool {
112        let mut inner = self.inner.lock();
113        if inner.backend.name() != backend.name() && !inner.sessions.is_empty() {
114            return false;
115        }
116        if inner.backend.name() != backend.name() {
117            tracing::info!(
118                previous = inner.backend.name(),
119                selected = backend.name(),
120                "the session backend moved with the role's drive"
121            );
122        }
123        inner.backend = backend;
124        true
125    }
126
127    /// Whether this client holds a session this name answers to.
128    ///
129    /// A mount names the session it was spawned for, and a plugin that outlived
130    /// a restart still names the one it was serving when it redials — to a client
131    /// that has no memory of it, because nothing here rebuilds slots from the
132    /// store. So this is a question with a real "no", and the caller needs to be
133    /// able to ask it before it binds anything under that name.
134    pub fn knows_session(&self, session_id: &str) -> bool {
135        super::state::slot_key_named(&self.inner.lock(), session_id).is_some()
136    }
137
138    /// The role's `[client.session]` policy.
139    pub fn session_policy(&self) -> onlyne_config::SessionPolicy {
140        self.inner.lock().session_policy.clone()
141    }
142
143    pub fn session_count(&self) -> usize {
144        self.inner.lock().sessions.len()
145    }
146
147    /// Feed the turn-started fact for a self-driven session's turn, so the
148    /// row's agent phase reads `running` before the agent can report a
149    /// completion through its tools mount (`docs/v2-CONTRACT.md` §3c). A plugin
150    /// drive feeds this through its heartbeats; a self-driven drive has none,
151    /// so the dispatch path feeds it where it hands the turn to the backend.
152    pub fn feed_turn_started(&self, task_id: &str) {
153        let inner = self.inner.lock();
154        if let Err(error) =
155            crate::reconcile::feed_turn_started(&inner.bridge, &inner.store, task_id)
156        {
157            tracing::warn!(task = %task_id, error = %error, "the turn-start feed did not apply");
158        }
159    }
160
161    /// Feed the turn-ended fact for a self-driven session's turn, so the row's
162    /// agent phase reads `idle` when the turn the drive witnessed ends. The
163    /// never-ran guard reads the same column, so a completion that arrives
164    /// after this feed still passes it.
165    pub fn feed_turn_ended(&self, task_id: &str) {
166        let inner = self.inner.lock();
167        if let Err(error) = crate::reconcile::feed_turn_ended(&inner.bridge, &inner.store, task_id)
168        {
169            tracing::warn!(task = %task_id, error = %error, "the turn-end feed did not apply");
170        }
171    }
172
173    pub fn role(&self) -> String {
174        self.inner.lock().role.clone()
175    }
176    pub fn command(&self) -> Vec<String> {
177        self.inner.lock().command.clone()
178    }
179    pub fn backend_name(&self, task_id: &str) -> Option<String> {
180        self.inner
181            .lock()
182            .sessions
183            .values()
184            .find(|slot| slot.task_id.as_deref() == Some(task_id))
185            .map(|slot| slot.session.backend.clone())
186    }
187    /// The terminal-fact stream of a backend that drives its own agent.
188    pub fn outcome_feed(&self) -> Option<crate::backend::OutcomeFeed> {
189        self.inner.lock().backend.outcomes()
190    }
191
192    /// Queue a delivery ack when the plugin refuses an assignment.
193    ///
194    /// An accepted assignment is not terminal for the server row: the normal
195    /// completion path still settles that delivery. A refused assignment is a
196    /// terminal local decision, so it uses the same durable ack queue as every
197    /// other delivery settlement.
198    pub fn push_assign_ack(&self, ack: onlyne_proto::AssignAckArgs) -> bool {
199        if ack.accepted {
200            return false;
201        }
202        let mut inner = self.inner.lock();
203        let msg_id = slot_key_serving_task(&inner, &ack.task_id)
204            .and_then(|key| inner.sessions.get_mut(&key))
205            .and_then(|slot| slot.msg_id.take());
206        let Some(msg_id) = msg_id else {
207            return false;
208        };
209        store_ack(
210            &inner,
211            AckArgs {
212                msg_id,
213                op_id: None,
214                accepted: false,
215                reason: ack.reason.or_else(|| Some("assign rejected".to_string())),
216            },
217        );
218        true
219    }
220
221    /// Queue an ack the client owes the server.
222    ///
223    /// Record an ack the client owes the server.
224    ///
225    /// The ack is durable: D11's control plane is at-least-once, and a settled
226    /// session whose ack is lost leaves the row in flight forever. The intent
227    /// queue carries it across a link that is down, and the flusher is the
228    /// sender.
229    pub fn push_settled(&self, ack: AckArgs) {
230        store_ack(&self.inner.lock(), ack);
231    }
232
233    /// Note that one task's plugin has been told by this client's own `control`
234    /// command to report the ending of that task.
235    ///
236    /// `on_control` runs this before the `recycle` frame leaves, which is the
237    /// point where the client still knows the order: the plugin's completion and
238    /// the retirement that command triggers race over the session's row, and a
239    /// guard that read the row would answer the same operator action two ways. One
240    /// note per task is kept, so a command issued twice waits for one answer.
241    ///
242    /// `word` is what the operator said and `now` is when they said it, and both
243    /// are the caller's to name rather than this call's to invent: the command is
244    /// the authority on the ending it asked for, and the instant it reads is the
245    /// one the watchdog's bound runs from.
246    pub fn owe_controlled_settle(&self, task_id: &str, word: ControlWord, now: Instant) {
247        let mut inner = self.inner.lock();
248        if !inner
249            .control_settles
250            .iter()
251            .any(|owed| owed.task_id == task_id)
252        {
253            inner.control_settles.push(ControlNote {
254                task_id: task_id.to_string(),
255                noted_at: now,
256                word,
257            });
258        }
259    }
260
261    /// Whether one completion answers a command noted above, consuming the note.
262    ///
263    /// The note is spent whichever way the settle it authorises goes: a refused
264    /// verdict leaves no second answer owed, and an applied one has travelled the
265    /// command's own completion. A later report for the same task is the plugin
266    /// speaking for itself again, and reads the ordinary door.
267    pub fn take_controlled_settle(&self, task_id: &str) -> bool {
268        let mut inner = self.inner.lock();
269        let Some(at) = inner
270            .control_settles
271            .iter()
272            .position(|owed| owed.task_id == task_id)
273        else {
274            return false;
275        };
276        inner.control_settles.swap_remove(at);
277        true
278    }
279
280    /// The notes whose operator's word has gone unanswered past the bound.
281    ///
282    /// The reading a sweep takes before it acts, and it spends nothing: the
283    /// settle below goes through [`take_controlled_settle`], so a completion that
284    /// answers a word between this read and that call takes the note first.
285    ///
286    /// [`take_controlled_settle`]: DispatchState::take_controlled_settle
287    pub fn control_settles_due(&self, now: Instant) -> Vec<ControlNote> {
288        due_control_settles(&self.inner.lock(), now)
289    }
290
291    /// Settle the work one operator's word left open, with no report behind it.
292    ///
293    /// The word asks a plugin for its own ending and the completion that answers
294    /// it is a frame of the plugin's, so a plugin that never sends one — it left
295    /// with the command's frame, or implements no `recycle` at all — leaves the
296    /// task open and the delivery row this client was handed in flight, with no
297    /// later caller to answer either. This is that caller.
298    ///
299    /// The writes are the ones `retire_dropped_ghosts` makes for the task its
300    /// owed session left: the verdict through the task's own record, which
301    /// refuses to overwrite one that landed first, and the still-held delivery row
302    /// refused with the operator's word, which is terminal for that row the way
303    /// every refusal is. The publish is the caller's, because it cannot run under
304    /// this lock.
305    ///
306    /// Answers `false` when this call is not the settle: the note is already spent
307    /// by a completion that answered the word, or the task's record carries a
308    /// verdict already, and either way nothing here is written and nothing is for
309    /// the caller to publish.
310    ///
311    /// [`retire_dropped_ghosts`]: DispatchState::retire_dropped_ghosts
312    pub fn settle_unanswered_control(&self, note: &ControlNote) -> bool {
313        // The note is taken first and at once: a report that answers the word
314        // while this call waits for the lock spends it, and the task then needs no
315        // verdict from here.
316        if !self.take_controlled_settle(&note.task_id) {
317            return false;
318        }
319        let mut inner = self.inner.lock();
320        // The verdict goes through the task's own record, which keeps the first
321        // one it was handed: a row an earlier settle answered stays as that settle
322        // left it. A verdict that was not this call's is a settle that already
323        // happened — every door that writes one answers the delivery row in the
324        // same breath — so there is nothing left here to refuse or to publish.
325        let verdict = task_state_of(note.word.outcome());
326        match inner.store.settle_task(&note.task_id, verdict) {
327            Ok(true) => {}
328            Ok(false) => {
329                tracing::warn!(
330                    task = %note.task_id,
331                    ?verdict,
332                    "an unanswered control command's task was already settled; the first verdict stands"
333                );
334                return false;
335            }
336            Err(error) => {
337                // A store that refused the write must not cost the word its
338                // answer: the note goes back where it came from — through the door
339                // that records one, stamped where it was — and the next tick tries
340                // again rather than leaving the task open forever.
341                tracing::warn!(
342                    task = %note.task_id,
343                    error = %error,
344                    "the task of an unanswered control command was not settled; the word stays owed"
345                );
346                drop(inner);
347                self.owe_controlled_settle(&note.task_id, note.word, note.noted_at);
348                return false;
349            }
350        }
351        // The delivery row this client is still holding is refused, and the
352        // reason is the operator's own word: the row is answered once, by whoever
353        // still holds its handle, and a plugin's report arriving later finds no
354        // handle left to spend.
355        let held = slot_key_serving_task(&inner, &note.task_id)
356            .and_then(|key| inner.sessions.get_mut(&key))
357            .and_then(|slot| slot.msg_id.take());
358        if let Some(msg_id) = held {
359            store_ack(
360                &inner,
361                AckArgs {
362                    msg_id,
363                    op_id: None,
364                    accepted: false,
365                    reason: Some(note.word.refusal().to_string()),
366                },
367            );
368        }
369        true
370    }
371
372    /// The role slice the dispatcher currently runs.
373    pub fn role_slice(&self) -> crate::session::slice::RoleSlice {
374        let inner = self.inner.lock();
375        crate::session::slice::RoleSlice {
376            drive: inner.drive.unwrap_or_default(),
377            command: inner.command.clone(),
378            max_sessions: inner.max_sessions,
379            required_targets: inner.required_targets.clone(),
380        }
381    }
382
383    /// Task ids currently occupying a live slot.
384    pub fn live_task_ids(&self) -> std::collections::HashSet<String> {
385        let inner = self.inner.lock();
386        inner
387            .sessions
388            .values()
389            .filter_map(|slot| slot.task_id.clone())
390            .collect()
391    }
392
393    /// The sorted live sessions one hello claims: the union of the memory
394    /// slots and the DB-persisted active sessions, so a process crash does not
395    /// lose the claim.
396    ///
397    /// A store failure is not swallowed: it is logged with the memory claim
398    /// still derivable from the slots, and returned so the caller can degrade
399    /// to [`live_claim_from_slots`] deliberately instead of answering as if the
400    /// durable half were simply empty. Silently dropping that half lets the
401    /// server requeue every in_flight row a crash left behind, which is the
402    /// duplicate delivery this claim exists to prevent.
403    pub fn hello_live_sessions(&self) -> onlyne_store::StoreResult<Vec<LiveSession>> {
404        let held = self.live_claim_from_slots();
405        // Merge DB-persisted sessions that are not yet exited. A fresh process
406        // after crash has empty slots but the DB still holds the sessions it was
407        // serving, so the hello must claim them to prevent the server from
408        // requeuing work this process is still running.
409        let persisted = {
410            let inner = self.inner.lock();
411            inner.store.active_sessions()
412        };
413        match persisted {
414            Ok(persisted) => {
415                let sessions = held.into_iter().chain(persisted);
416                Ok(crate::session::claim::from_sessions(sessions))
417            }
418            Err(error) => {
419                tracing::error!(
420                    error = %error,
421                    memory_claim = ?held,
422                    "hello claim lost its durable half: active_sessions failed; the DB-persisted sessions are missing and the server may requeue them"
423                );
424                Err(error)
425            }
426        }
427    }
428
429    /// The memory half of [`hello_live_sessions`]: the claim to dial with when
430    /// the durable store cannot answer. The slots are what this process is
431    /// serving right now, so even a degraded hello keeps those rows in_flight.
432    ///
433    /// A slot answers with the session's own id — the one that stays put while a
434    /// scope hands the session delivery after delivery — the delivery it is
435    /// serving now, and whether its process has been released. A session between
436    /// deliveries claims no delivery: the server keeps the rows those sessions
437    /// are still the owners of, and invents none.
438    pub fn live_claim_from_slots(&self) -> Vec<LiveSession> {
439        let inner = self.inner.lock();
440        crate::session::claim::from_sessions(inner.sessions.values().map(|slot| LiveSession {
441            session_id: slot.session.task_id.clone(),
442            task_id: slot.task_id.clone(),
443            suspended: slot.suspended,
444        }))
445    }
446
447    /// Start the stall clock for a newly assigned task.
448    pub fn note_stall_assigned(&self, task_id: &str, now: Instant) {
449        self.inner.lock().stall.note_assigned(task_id, now);
450    }
451
452    /// Refresh the stall clock after an Applied persist.
453    pub fn note_stall_applied(&self, task_id: &str, now: Instant) {
454        self.inner.lock().stall.note_applied(task_id, now);
455    }
456
457    /// Task ids whose freeze exceeds `threshold_secs` in this episode.
458    /// Exited projections retire their remaining progress clocks.
459    pub fn stall_due(&self, now: Instant, threshold_secs: u64) -> Vec<String> {
460        let mut inner = self.inner.lock();
461        let due = inner.stall.due(now, threshold_secs);
462        let mut active = Vec::with_capacity(due.len());
463        for task_id in due {
464            if session_exited(&inner, &task_id) {
465                inner.stall.forget(&task_id);
466            } else {
467                active.push(task_id);
468            }
469        }
470        active
471    }
472
473    /// Remember that this freeze episode has been reported.
474    pub fn mark_stalled(&self, task_id: &str) {
475        self.inner.lock().stall.mark_reported(task_id);
476    }
477
478    /// Observation-only stall fault for an active task, carrying the stored
479    /// watermark. A tuple and verdict that derive `exited` retire their progress
480    /// clock before the send boundary.
481    pub fn stall_report(&self, task_id: &str) -> Option<Report> {
482        let mut inner = self.inner.lock();
483        let row = inner.store.get_session(task_id).ok().flatten();
484        let exited = row.as_ref().is_some_and(|row| {
485            projection_of(row, stored_task_state(&inner, task_id)).lifecycle == Lifecycle::Exited
486        });
487        if exited {
488            inner.stall.forget(task_id);
489            return None;
490        }
491        Some(crate::session::stall::report(
492            task_id,
493            Some(task_id.to_string()),
494            row.as_ref().map(|row| row.generation as u64),
495            row.as_ref().map(|row| row.seq as u64),
496        ))
497    }
498
499    /// Whether any adapter is currently mounted (named or parked).
500    pub fn has_mounted_adapter(&self) -> bool {
501        let inner = self.inner.lock();
502        !inner.transports.is_empty() || !inner.parked.is_empty()
503    }
504
505    /// Adopt a role slice: the one `welcome` carried, or the one a reload's
506    /// role row carries.
507    pub fn reconfigure(&self, slice: crate::session::slice::RoleSlice) {
508        let mut inner = self.inner.lock();
509        inner.command = slice.command;
510        inner.max_sessions = slice.max_sessions;
511        inner.required_targets = slice.required_targets;
512    }
513
514    /// Role prose last cached from `welcome`.
515    pub fn role_prose(&self) -> String {
516        let inner = self.inner.lock();
517        inner
518            .store
519            .prose(&inner.role)
520            .ok()
521            .flatten()
522            .map(|(prose, _)| prose)
523            .unwrap_or_default()
524    }
525
526    /// Whether one task names a session this client holds, in memory or in its
527    /// durable session rows.
528    pub fn holds_task(&self, task_id: &str) -> bool {
529        let inner = self.inner.lock();
530        inner
531            .sessions
532            .values()
533            .any(|slot| slot.task_id.as_deref() == Some(task_id))
534            || inner.store.get_session(task_id).ok().flatten().is_some()
535    }
536
537    /// Generation the reducer holds for one task, before any hand-off.
538    pub fn session_generation(&self, task_id: &str) -> Option<u64> {
539        self.inner
540            .lock()
541            .sessions
542            .values()
543            .find(|slot| slot.task_id.as_deref() == Some(task_id))
544            .map(|slot| slot.session.generation)
545    }
546
547    /// Whether a delivery has somewhere to run.
548    ///
549    /// §5's `max_sessions` caps concurrency, so a delivery that arrives at the
550    /// cap waits on the server: the row stays in flight and the next pull
551    /// offers it again once a session frees. Each task runs in its own session,
552    /// so a slot whose task has finished still spends capacity until it retires.
553    ///
554    /// A scope that hands a delivery on to a session the role already holds
555    /// spends no slot at all, so an `task` or `role` session sitting idle — or
556    /// suspended, its slot already given back — is room even at the cap. The
557    /// placement decides which delivery goes where; this only answers whether
558    /// there is somewhere for one to go. `oneshot` reuses nothing, which is the
559    /// count it has always been.
560    pub fn has_capacity(&self) -> bool {
561        let inner = self.inner.lock();
562        if !matches!(
563            inner.session_policy.scope,
564            onlyne_config::SessionScope::Oneshot
565        ) && inner.sessions.values().any(super::scope::takes_new_work)
566        {
567            return true;
568        }
569        live_sessions(&inner) < inner.max_sessions as usize
570    }
571
572    /// Whether this role already finished one task with a terminal `Done`.
573    ///
574    /// The task's own record is the account: the settle writes the verdict the
575    /// agent filed and refuses to overwrite it, so a record reading `done` means
576    /// this role answered for this task id once already. A redelivery of that
577    /// task is not new work — running it again would stage its payload on
578    /// whichever session happens to be idle, so one chain's task executes inside
579    /// another conversation and the second answer collides with the verdict the
580    /// first one settled.
581    ///
582    /// Only `done` counts. A session killed or crashed mid-flight leaves its task
583    /// open, or settles it `failed`, and the server's requeue, `repair_retry`, and
584    /// `control retry` all re-offer that task on purpose, so those deliveries
585    /// still run.
586    pub fn task_completed_here(&self, task_id: &str) -> bool {
587        let inner = self.inner.lock();
588        matches!(
589            stored_close_reason(&inner, task_id),
590            Some(crate::backend::CloseReason::Completed)
591        )
592    }
593
594    /// One staged session this role's standing runtime can be offered.
595    ///
596    /// The mirror of [`Self::staged_without_transport`], narrowed to a session
597    /// that has no transport *because* a hosting runtime will supply one. A
598    /// session another connection already serves is not offered, and a session
599    /// whose work is already in flight is not either — a runtime must never be
600    /// asked to open a second conversation for a chain that has one.
601    pub fn staged_hosting_session(&self) -> Option<String> {
602        let inner = self.inner.lock();
603        inner
604            .sessions
605            .iter()
606            .find(|(key, slot)| {
607                slot.payload.is_some()
608                    && slot.session.backend == "hosting"
609                    && !inner
610                        .transports
611                        .keys()
612                        .any(|session| super::transport::names_session(key, slot, session))
613            })
614            .map(|(key, _)| key.clone())
615    }
616
617    /// What to ask a hosting runtime for, by the name this client gave the
618    /// session.
619    ///
620    /// The resume handle of the family's previous session rides along when there
621    /// is one, so a `task`-scoped runtime hands back the conversation it already
622    /// holds rather than opening a second one for the same chain. Nothing reads
623    /// the handle here: it is the runtime's own word for where its conversation
624    /// is, and this client stores it and gives it back.
625    pub fn hosting_open_args(&self, session_id: &str) -> onlyne_proto::OpenArgs {
626        let inner = self.inner.lock();
627        let slot = inner.sessions.get(session_id);
628        let (task_id, family, prose) = match slot {
629            Some(slot) => (slot.task_id.clone(), slot.family.clone(), String::new()),
630            None => (Some(session_id.to_string()), None, String::new()),
631        };
632        let handle = inner
633            .sessions
634            .iter()
635            .filter(|(key, other)| other.family == family && *key != session_id)
636            .filter_map(|(_, other)| other.resume_handle.clone())
637            .next();
638        onlyne_proto::OpenArgs {
639            session_id: session_id.to_string(),
640            task_id: task_id.unwrap_or_else(|| session_id.to_string()),
641            scope: format!("{:?}", inner.session_policy.scope).to_lowercase(),
642            family,
643            prose,
644            resume_handle: handle,
645        }
646    }
647
648    /// The task of one session that holds a payload with no connection bound.
649    ///
650    /// A work item that arrives before its always-running agent mounts waits in
651    /// exactly this state, and the mount ends the wait. A session answers through
652    /// the transport its first task claimed, so it stays served.
653    pub fn staged_without_transport(&self) -> Option<String> {
654        let inner = self.inner.lock();
655        inner
656            .sessions
657            .iter()
658            .find(|(_, slot)| slot.payload.is_some())
659            .filter(|(key, slot)| {
660                !inner
661                    .transports
662                    .keys()
663                    .any(|session| names_session(key, slot, session))
664            })
665            .and_then(|(_, slot)| slot.task_id.clone())
666    }
667
668    /// The tools token of one live session, when this client holds it.
669    ///
670    /// The one door the token leaves this process through: the session's own
671    /// drive reads it while building the child that mounts `onlyne mcp`
672    /// (`SpawnSpec.tools_token`), and nothing else may hand it out
673    /// (`docs/v2-CONTRACT.md` §3b).
674    pub fn tools_token(&self, session_id: &str) -> Option<String> {
675        let inner = self.inner.lock();
676        let key = slot_key_named(&inner, session_id)?;
677        inner
678            .sessions
679            .get(&key)
680            .filter(|slot| !session_exited(&inner, &slot.session.task_id))
681            .map(|slot| slot.tools_token.clone())
682    }
683
684    /// The session a `tools` mount token speaks for, when it names a live one.
685    pub fn tools_mount_for(&self, token: &str) -> Option<ToolsSession> {
686        let inner = self.inner.lock();
687        let key = slot_key_for_token(&inner, token)?;
688        tools_session_of(&inner, &key)
689    }
690
691    /// Bind one `tools` connection to the session its token names.
692    ///
693    /// The hello validated the token; this writes the binding the session's
694    /// later frames are measured against, and answers the session the
695    /// connection now speaks for. A second connection presenting the same token
696    /// takes the binding over, and the earlier one then fails the per-frame
697    /// liveness check — a token belongs to one session, and the newest
698    /// connection is the one that speaks for it. `None` is a token whose
699    /// session retired between the handshake and this line.
700    pub fn bind_tools_mount(&self, token: &str, io: AdapterIo) -> Option<ToolsSession> {
701        let mut inner = self.inner.lock();
702        let key = slot_key_for_token(&inner, token)?;
703        let session = tools_session_of(&inner, &key)?;
704        inner.tools_mounts.retain(|(held, _)| held != &key);
705        inner.tools_mounts.push((key, io));
706        Some(session)
707    }
708
709    /// Whether one live tools connection still speaks for a session this client
710    /// holds.
711    ///
712    /// A token dies with its session, and a mount whose session has ended is
713    /// refused from then on: this is the per-frame half of that rule, so a
714    /// connection left open past its session's retirement answers `unauthorized`
715    /// and closes rather than speaking for work nobody holds.
716    pub fn tools_connection_live(&self, io: &AdapterIo) -> bool {
717        self.tools_scope(io).is_some()
718    }
719
720    /// What one live tools connection speaks for, read under the lock that owns
721    /// it; `None` when the connection is unbound or its session has stopped
722    /// serving.
723    ///
724    /// The token is the binding (`docs/v2-CONTRACT.md` §3b), so everything a
725    /// tools frame is stamped with comes from the session's own slot: the open
726    /// delivery its `report` and `handoff` frames name, and the session a `send`
727    /// leaves from. Reading it in one locked pass is what keeps a frame from
728    /// being stamped from a state that moved between the lookup and the stamp.
729    pub fn tools_scope(&self, io: &AdapterIo) -> Option<ToolsScope> {
730        let inner = self.inner.lock();
731        let (key, _) = inner
732            .tools_mounts
733            .iter()
734            .find(|(_, bound)| bound.same_connection(io))?;
735        inner
736            .sessions
737            .get(key)
738            .filter(|slot| token_names_session(&inner, slot))
739            .map(|slot| ToolsScope {
740                session_id: slot.session.task_id.clone(),
741                task_id: slot.task_id.clone(),
742            })
743    }
744
745    /// Drop the binding one tools connection held, whichever session it named.
746    pub fn release_tools_connection(&self, io: &AdapterIo) {
747        self.inner
748            .lock()
749            .tools_mounts
750            .retain(|(_, bound)| !bound.same_connection(io));
751    }
752}