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(¬e.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(¬e.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(¬e.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, ¬e.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}