onlyne_client/session/dispatch/transport.rs
1use super::*;
2
3use super::env::missing_capability;
4use super::idle::{resumable, suspend_locked};
5use super::outbound::queue_outbound_locked;
6use super::retire::{
7 PendingClose, close_retired, keeps_idle, retire_idle_locked, stored_close_reason,
8};
9use super::state::{
10 DispatchInner, DispatchState, FrameGuard, SessionSlot, has_attached_transport,
11 rebase_generation, slot_key_named, slot_key_serving_task, slot_task,
12};
13
14/// The one sentence a session that ended a turn without a completion reads.
15///
16/// The wording is a contract (`docs/v2-CONTRACT.md` §3c): the plugin composes
17/// none of it and keeps no copy, so this client owns the bytes and hands them
18/// over verbatim through [`DispatchState::nudge_plugin`].
19pub const NUDGE_TEXT: &str =
20 "If this task is finished, report it with onlyne_complete; if something is missing, say what.";
21
22/// Whether one session slot is the session an adapter mount named.
23///
24/// The mount carries the id the client spawned the plugin with
25/// (`ONLYNE_SESSION_ID`), which is the slot's key and its stored reference.
26/// Each task gets its own session, so those spellings all name one session; the
27/// extra checks stay because a slot keeps the id its session was born with even
28/// after its task binding changes.
29pub(super) fn names_session(key: &str, slot: &SessionSlot, session_id: &str) -> bool {
30 key == session_id
31 || slot.session.task_id == session_id
32 || slot.task_id.as_deref() == Some(session_id)
33}
34
35/// Whether a connection other than `io` is the one serving one slot.
36fn attached_to_other(inner: &DispatchInner, key: &str, slot: &SessionSlot, io: &AdapterIo) -> bool {
37 inner.transports.iter().any(|(session_id, (live, _))| {
38 !live.same_connection(io) && names_session(key, slot, session_id)
39 })
40}
41
42/// Whether `io` is the connection serving one session, as the binding rules left
43/// it.
44///
45/// A connection this client holds read-only is never the answer: a mount that
46/// finds its session already served lands in `revived` and never in
47/// `transports`, which is the map this reads. Nothing else about the connection
48/// is consulted — a frame carries its sender, so a connection serving one
49/// session cannot answer for another by naming it.
50pub(super) fn serves_session(inner: &DispatchInner, session_id: &str, io: &AdapterIo) -> bool {
51 let Some(key) = slot_key_named(inner, session_id) else {
52 return false;
53 };
54 let Some(slot) = inner.sessions.get(&key) else {
55 return false;
56 };
57 slot_is_served_by(inner, &key, slot, io)
58}
59
60/// Whether one slot's session is served by `io`.
61///
62/// The body of [`serves_session`] and [`serves_task`]: the question is never
63/// which name the frame carried, but whether this connection is the transport of
64/// the slot that name resolves to.
65fn slot_is_served_by(inner: &DispatchInner, key: &str, slot: &SessionSlot, io: &AdapterIo) -> bool {
66 inner.transports.iter().any(|(served, (transport, _))| {
67 transport.same_connection(io) && names_session(key, slot, served)
68 }) || tools_bound_to(inner, key, io)
69}
70
71/// Whether `io` is the tools mount bound to one session.
72///
73/// A tools mount holds no process and is no transport, so the map the binding
74/// rules fill says nothing about it: the handshake's own record is the whole
75/// answer, and it is what lets a tools connection's `handoff` and `report`
76/// frames travel the same authority door an agent's frames do
77/// (`docs/v2-CONTRACT.md` §3b).
78fn tools_bound_to(inner: &DispatchInner, key: &str, io: &AdapterIo) -> bool {
79 inner
80 .tools_mounts
81 .iter()
82 .any(|(held, bound)| held == key && bound.same_connection(io))
83}
84
85/// Whether `io` is the connection serving the session that answers for one task.
86///
87/// The task-spelling companion of [`serves_session`], and the authority every
88/// state frame a plugin sends is measured against: a report names a task, and the
89/// frame carries its sender, so a connection serving one session cannot end,
90/// fault, or move another session by naming its task.
91///
92/// Unlike [`serves_session`] this walks every slot the task names rather than the
93/// first one a map reach finds. Two slots answer to one task exactly when a
94/// session came back for a task a newer session already took, and there the
95/// first match is whichever one `HashMap` order happens to yield — an authority
96/// check that reads that way refuses the same frame twice and applies it a third
97/// time. The extra walk costs nothing on the ordinary path, where one slot
98/// answers to the task, and the demoted half can never pass anyway: a read-only
99/// slot is served by no transport at all.
100pub(super) fn serves_task(inner: &DispatchInner, task_id: &str, io: &AdapterIo) -> bool {
101 inner.sessions.iter().any(|(key, slot)| {
102 names_session(key, slot, task_id) && slot_is_served_by(inner, key, slot, io)
103 })
104}
105
106/// Whether `io` is any session's transport right now.
107///
108/// Weaker than [`serves_task`] — it says the connection serves *a* session of
109/// this role, not the one a frame names — and it exists for the one window where
110/// the stronger question cannot be asked: an operator's `recycle` or `cancel`
111/// retires the slot before the completion it ordered arrives, and a retired slot
112/// is out of `sessions` for `serves_task` to resolve. See
113/// `ending_is_authorised` in `reports.rs`.
114pub(super) fn is_bound_transport(inner: &DispatchInner, io: &AdapterIo) -> bool {
115 inner
116 .transports
117 .values()
118 .any(|(transport, _)| transport.same_connection(io))
119}
120
121/// Whether one connection is held read-only: it mounted a session this client
122/// already serves through a different live connection.
123pub(super) fn is_revived_connection(inner: &DispatchInner, io: &AdapterIo) -> bool {
124 inner
125 .revived
126 .iter()
127 .any(|(_, revived, _)| revived.same_connection(io))
128}
129
130/// Whether the connection this client holds read-only answers for one task.
131///
132/// The held half of the pair [`serves_task`] reads live. A demoted connection is
133/// no session's transport, so the strict rule refuses its observation — and the
134/// design still lets it answer for the task its own agent finished: `plugin_send`
135/// routes what it sends onto the wire like any other session's, and
136/// `retire_revived` retires the demoted slot with the completion that answers it.
137/// A held connection is therefore the connection that may report one task ended,
138/// and no other: the name it mounted with resolves to the slot the task answers
139/// to.
140fn held_for_task(inner: &DispatchInner, task_id: &str, io: &AdapterIo) -> bool {
141 inner.revived.iter().any(|(name, revived, _)| {
142 revived.same_connection(io)
143 && slot_key_named(inner, name)
144 .and_then(|key| inner.sessions.get(&key))
145 .is_some_and(|slot| slot_task(slot) == task_id)
146 })
147}
148
149/// Whether `io` may end or fault the session that answers for one task.
150///
151/// [`serves_task`] is the whole rule for a frame that moves state: a beat is an
152/// observation of a session, and only its transport has one. An ending is a
153/// different fact — the agent behind either connection ran the work, and the
154/// client may have moved the binding while it was running — so the two windows
155/// where the strict rule cannot see the sender are read here instead:
156///
157/// * the connection is held read-only for this very task ([`held_for_task`]), and
158/// * this client ordered the task to its ending with `control`, which retired the
159/// slot before the answer it asked for arrived. The note names the task, and a
160/// transport still bound somewhere in this role names the sender; a forged
161/// ending has neither, and a foreign connection answering for a task someone
162/// else's command opened has the note without the binding.
163pub(super) fn serves_ending(inner: &DispatchInner, task_id: &str, io: &AdapterIo) -> bool {
164 serves_task(inner, task_id, io)
165 || held_for_task(inner, task_id, io)
166 || (inner
167 .control_settles
168 .iter()
169 .any(|note| note.task_id == task_id)
170 && is_bound_transport(inner, io))
171}
172
173/// Whether one slot is held by a connection this client keeps read-only.
174///
175/// A slot demoted by [`note_binding_locked`] — its task taken by a newer session
176/// — is served by no transport at all: the connection that came back for it
177/// waits in `revived` under the name it mounted with, which is the name that
178/// names this slot. While that connection is there the slot has an owner:
179/// `retire_revived` retires it with the completion that answers what the held
180/// connection wrote. A demoted slot whose held connection has since gone has no
181/// such owner left, so the answer is read off the held names rather than off the
182/// demotion alone — and it is read by asking which slots a name names, never by
183/// asking which slot a name resolves to, because two slots can answer to one
184/// name and only one of them is the held connection's.
185pub(super) fn held_read_only(inner: &DispatchInner, key: &str, slot: &SessionSlot) -> bool {
186 slot.read_only
187 && inner
188 .revived
189 .iter()
190 .any(|(name, _, _)| names_session(key, slot, name))
191}
192
193/// Record a returning connection as read-only, once per connection.
194pub(super) fn record_revived_connection(
195 inner: &mut DispatchInner,
196 session_id: &str,
197 io: AdapterIo,
198 capabilities: Vec<Capability>,
199) {
200 if !is_revived_connection(inner, &io) {
201 inner
202 .revived
203 .push((session_id.to_string(), io, capabilities));
204 }
205}
206
207/// Give one session back to the oldest connection that was held read-only for it.
208///
209/// A held connection is read-only only while another connection serves its
210/// session, so the moment that connection goes is the moment the held one becomes
211/// the session's only transport. Without this, a plugin that redials while the
212/// client still holds the dead socket behind it is silenced for the rest of the
213/// session: the assignment it came for is written to a connection nobody reads,
214/// and nothing promotes it later. The demotion lifts with the promotion, so the
215/// reconnected agent keeps both its session and its delivery rights.
216fn promote_held_connection(inner: &mut DispatchInner, key: &str) {
217 let attached = inner
218 .sessions
219 .get(key)
220 .is_some_and(|slot| has_attached_transport(inner, key, slot));
221 if attached {
222 return;
223 }
224 let held = inner.revived.iter().position(|(name, _, _)| {
225 slot_key_named(inner, name).is_some_and(|held_key| held_key == key)
226 });
227 let Some(index) = held else { return };
228 let (name, io, capabilities) = inner.revived.remove(index);
229 if let Some(slot) = inner.sessions.get_mut(key) {
230 slot.read_only = false;
231 }
232 tracing::info!(
233 session = %name,
234 "a held connection takes the session its predecessor left"
235 );
236 attach_transport_locked(inner, &name, io, capabilities);
237}
238
239/// Settle what one mounting connection means for the session it names.
240///
241/// A mount that finds nothing serving its session takes it and clears the clock
242/// [`DispatchState::release_connection`] started: that is the agent that came
243/// back inside the reconnect grace, and the always-running agent serving task
244/// after task lives in this path. A mount that finds the session already served
245/// takes nothing — either another connection holds that very slot, or the task it
246/// names now answers from a slot of its own, which is the case where a newer
247/// session was spawned to retry the work while the old agent's process came back.
248/// Such a connection is recorded read-only, and the slot it names is demoted too
249/// when it owns a slot of its own.
250///
251/// Every binding path runs this one judgement, including the ready report, which
252/// reaches an agent without writing a transport. Concurrency is what decides, not
253/// the drop clock: the retry that claims an unclaimed session is served, and the
254/// connection that returns to a session already served is held, whichever of the
255/// two mounted first. [`DispatchState::release_connection`] promotes a held
256/// connection when the live one it waited behind goes away, so a plugin that
257/// redials over a socket the client has not yet seen die still gets its session.
258pub(super) fn note_binding_locked(
259 inner: &mut DispatchInner,
260 session_id: &str,
261 io: &AdapterIo,
262) -> bool {
263 if is_revived_connection(inner, io) {
264 return false;
265 }
266 let Some(key) = slot_key_named(inner, session_id) else {
267 return true;
268 };
269 let Some(slot) = inner.sessions.get(&key) else {
270 return true;
271 };
272 let task = slot_task(slot);
273 let taken = attached_to_other(inner, &key, slot, io);
274 let moved_on = !taken
275 && inner.sessions.iter().any(|(other, other_slot)| {
276 *other != key
277 && slot_task(other_slot) == task
278 && attached_to_other(inner, other, other_slot, io)
279 });
280 let revived = taken || moved_on;
281 // Whether this very connection already stands as the transport of the slot the
282 // name resolves to. The answer decides what this judgement is allowed to
283 // renew: a connection that already serves the session takes nothing new from
284 // the mount, and the mount proves nothing the transport record did not
285 // already prove. Without the distinction the stamp is a clock any caller can
286 // wind, because `hand_staged` runs this judgement through `on_ready` for the
287 // connection already serving a session — so every re-delivery of a task, and
288 // every zombie mount that pushes one through, bought a dead agent a fresh
289 // heartbeat window without one frame from the agent whose silence is what the
290 // sweep reads. The stamp a returning agent does earn is written by the frames
291 // it sends: the ready report of §6 stamps it when the reducer takes the row
292 // (`note_beat`), and so does every accepted beat afterwards.
293 let already_serving = slot_is_served_by(inner, &key, slot, io);
294 // A mount that takes a session whose death clock was running is the agent
295 // coming back inside the reconnect grace: the barrier had already passed, so
296 // `ready` says the plugin spoke once, and the clock says the connection it
297 // spoke through has since ended. That is the only shape this rebase is for.
298 // A first mount has no clock running over a session that ever spoke, and a
299 // settled session owes no work its reporter would answer for.
300 let returning = !revived && slot.dropped_at.is_some() && slot.ready && slot.task_id.is_some();
301 if let Some(slot) = inner.sessions.get_mut(&key) {
302 if revived {
303 // Only the spelling where this name's own slot is still served by
304 // the newer connection leaves a slot of its own to silence; when it
305 // is, the slot belongs to that live connection and keeps its rights.
306 slot.read_only = moved_on;
307 } else {
308 slot.dropped_at = None;
309 slot.read_only = false;
310 // The mount is a frame this session's agent sent, and taking it is
311 // what says the agent is here now. Without this the liveness stamp
312 // would still read the moment before the drop, and the sweep's
313 // silence arm would judge a returning agent that has not beaten yet
314 // on the age of a frame from the process before it.
315 if !already_serving {
316 slot.last_beat = Some(Instant::now());
317 }
318 }
319 }
320 if revived {
321 tracing::warn!(
322 session = %session_id,
323 task = %task,
324 "a plugin mounted a session this client already serves; it is held read-only"
325 );
326 return false;
327 }
328 if returning {
329 rebase_returned_reporter(inner, &key);
330 }
331 true
332}
333
334/// Rebase the watermark of a session whose agent came back, so its next frame
335/// lands.
336///
337/// A plugin restarts its own sequence at its base and its generation is a
338/// constant, while the watermark the row holds is the last sequence the process
339/// that left reached. Left alone, every frame the returning reporter sends reads
340/// at or below that watermark and is dropped as a stale duplicate until its
341/// sequence climbs past it — for a session that had been running a while, the
342/// whole rest of its work, reported into a tuple that never moves.
343///
344/// The rebase itself is [`rebase_generation`]'s: a new generation over the tuple
345/// the client already holds, because the generation is the half the reporter
346/// cannot be talked out of — its `generation` field is a constant it never
347/// raises, so stamping the beat with the session's own generation is what lets a
348/// frame from the new generation through at all, and the sequence starts again
349/// under it.
350///
351/// The body is the stored tuple with the agent dimension put back to `Booting`
352/// and the recovery line beside it dropped. That is not a guess about the agent:
353/// the connection that witnessed the last agent fact is the one that ended, and
354/// the plugin that just mounted has reported no turn fact yet, which is exactly
355/// what `Booting` means. `recovery` goes because the reducer's own coupling
356/// refuses a recovery substate beside a booting agent, and an accepted receipt
357/// goes to `Pending` for the same reason — the repair `compose_observation`
358/// already makes for a plugin that reports a booting process over a closed
359/// drain. Everything else the client owns — the delivery drain, the reconcile
360/// tuning, the counter — rides forward untouched, and so does the resource.
361///
362/// The generation being replaced is the one whose connection ended: this client
363/// is the only authority on its own connection bookkeeping, and it carries the
364/// observation content forward rather than replacing it, which is what the
365/// attestation the reducer asks for is guarding against. A reporter from a
366/// connection this client holds read-only never reaches the reducer at all —
367/// `serves_session` refuses its frames before the gate — so lowering the
368/// watermark here admits a returning reporter and nothing else.
369fn rebase_returned_reporter(inner: &mut DispatchInner, key: &str) {
370 let Some(slot) = inner.sessions.get(key) else {
371 return;
372 };
373 let task_id = slot.session.task_id.clone();
374 let verdict = rebase_generation(inner, &task_id, |stored| {
375 let mut body = stored.clone();
376 body.agent = AgentPhase::Booting;
377 body.recovery = RecoveryPhase::NoRecovery;
378 if body.delivery == DeliveryPhase::Accepted {
379 body.delivery = DeliveryPhase::Pending;
380 }
381 body
382 });
383 match verdict {
384 Ok(Some(Verdict::Applied(next))) => tracing::info!(
385 task = %task_id,
386 generation = next.version.generation,
387 "a returning plugin's watermark was rebased onto a new generation"
388 ),
389 Ok(Some(verdict)) => tracing::warn!(
390 task = %task_id,
391 ?verdict,
392 "the returning plugin's watermark was not rebased"
393 ),
394 Ok(None) => tracing::warn!(
395 task = %task_id,
396 "the returning plugin's row was not there to rebase"
397 ),
398 Err(error) => tracing::warn!(
399 task = %task_id,
400 error = %error,
401 "the returning plugin's watermark was not rebased"
402 ),
403 }
404}
405
406/// Attach one plugin connection to the session it names, or hold it read-only.
407///
408/// This is where a mount that named a session becomes a transport, so the
409/// read-only connection of §1 (b) never lands in `transports` and never steals
410/// the assignment, delivery, or note addressed to the connection that serves the
411/// session now. The one other writer of that map is the parked claim, which runs
412/// the same judgement and keeps a refused agent in the park instead of recording
413/// it read-only for a session it never named. Answers whether the connection took
414/// the session.
415fn attach_transport_locked(
416 inner: &mut DispatchInner,
417 session_id: &str,
418 io: AdapterIo,
419 capabilities: Vec<Capability>,
420) -> bool {
421 if !note_binding_locked(inner, session_id, &io) {
422 record_revived_connection(inner, session_id, io, capabilities);
423 return false;
424 }
425 inner
426 .transports
427 .insert(session_id.to_string(), (io, capabilities));
428 true
429}
430
431impl DispatchState {
432 /// Hold `io` for as long as one of its inbound frames is being handled.
433 pub fn hold_frame(&self, io: &AdapterIo) -> FrameGuard<'_> {
434 self.inner.lock().in_frame.push(io.clone());
435 FrameGuard {
436 state: self,
437 io: io.clone(),
438 }
439 }
440
441 /// Bind one adapter connection to the session it named.
442 ///
443 /// The name is the session id the client spawned the plugin with, which is
444 /// enough on its own: a plugin that mounts before the client staged its
445 /// session is remembered here and takes the payload the moment it is
446 /// staged, and a plugin that mounts after finds its session waiting.
447 pub fn bind_adapter(&self, session_id: &str, io: AdapterIo, capabilities: Vec<Capability>) {
448 attach_transport_locked(&mut self.inner.lock(), session_id, io, capabilities);
449 }
450
451 /// Remember the delivery handle for one task.
452 ///
453 /// The handle goes to the session serving the task, not to a read-only slot
454 /// that came back for it, so the ack this earns answers the live delivery.
455 pub fn attach_msg_id(&self, task_id: &str, msg_id: &str) {
456 let mut inner = self.inner.lock();
457 let Some(key) = slot_key_serving_task(&inner, task_id) else {
458 return;
459 };
460 if let Some(slot) = inner.sessions.get_mut(&key) {
461 slot.msg_id = Some(msg_id.to_string());
462 }
463 }
464
465 /// Take one plugin `send` frame and answer what the plugin is told.
466 ///
467 /// One path, whichever connection sent the frame. A handoff is a real
468 /// delivery to the role it names, and this client is the only process
469 /// holding a link that could carry it, so a connection that came back for a
470 /// task another session now serves routes here exactly like the serving one.
471 /// Nothing is held for a later merge (§3c): a session's handoffs travel as
472 /// its own frames, answered where it sends them, so a settlement is only
473 /// the completion's half of the account.
474 pub fn plugin_send(&self, io: &AdapterIo, envelope: &Envelope) -> Result<serde_json::Value> {
475 let mut inner = self.inner.lock();
476 // The envelope leaves on this client's authenticated link, so the server
477 // reads what it carries as this role's own message.
478 send_is_authorised(&inner, envelope)?;
479 let op_id = queue_outbound_locked(&mut inner, envelope)?;
480 // A send the client carried is a delivery to the role it names, and it is
481 // the evidence the relay guard reads at this session's next completion
482 // (`guards.rs`). Only a connection that serves a session records: a
483 // connection that came back for a task another session holds speaks for
484 // no session of this client's.
485 if let Some(key) = super::guards::session_key_of_connection(&inner, io) {
486 super::guards::record_delivery(&mut inner, &key, &envelope.to);
487 }
488 Ok(serde_json::json!({"queued": true, "op_id": op_id}))
489 }
490
491 /// Park one plugin connection as this role's waiting agent.
492 ///
493 /// Only a mount that names no session parks: it is a plugin that attached
494 /// before any work existed, so it takes the next session this role stages
495 /// (plan §6 line 285). A role can host more than one such agent, and each
496 /// that arrives joins the back of the queue: the record is a queue and not
497 /// one slot, because overwriting it dropped the connection that was already
498 /// waiting with no accounting of any kind — no release, no log, and no word
499 /// to the plugin, which is how a quiet role lost a worker.
500 ///
501 /// A connection already in the queue is refreshed where it stands rather
502 /// than sent to the back: a mount that names nothing twice over is the same
503 /// agent re-helloing, and its place in line is that agent's due.
504 pub fn park_transport(&self, io: AdapterIo, capabilities: Vec<Capability>) {
505 let mut inner = self.inner.lock();
506 let waiting = inner
507 .parked
508 .iter()
509 .position(|(parked, _)| parked.same_connection(&io));
510 match waiting {
511 Some(index) => inner.parked[index].1 = capabilities,
512 None => inner.parked.push((io, capabilities)),
513 }
514 }
515
516 /// Claim the longest-waiting agent of this role for one staged session.
517 ///
518 /// An always-running plugin mounts naming no session, so the park holds the
519 /// connection that can serve the session staged next (plan §6 line 285). The
520 /// queue answers in the order it filled, so the agent that has waited longest
521 /// is the one that takes the work. A claim left unbound strands the staged
522 /// work: the session has a payload and this client holds no record of the
523 /// socket that serves it, so a claim the binding judgement refuses returns the
524 /// connection to the back of the park instead of spending it.
525 ///
526 /// A refused parked connection is not a mount that named a session and found
527 /// it served, and the difference is what it is kept for: recording it
528 /// read-only under a session it never mounted would hold it answerable to
529 /// that session's task, and a role with more work to stage would lose the one
530 /// waiting agent it has. It waits, and the socket's own end takes it out of
531 /// the queue through `release_connection`.
532 pub(super) fn claim_parked_transport(
533 &self,
534 session_id: &str,
535 ) -> Option<(AdapterIo, Vec<Capability>)> {
536 let mut inner = self.inner.lock();
537 if inner.parked.is_empty() {
538 return None;
539 }
540 // The queue is served oldest first.
541 let (io, capabilities) = inner.parked.remove(0);
542 // The judgement of §1 (b) runs first, and the transport is written only
543 // once it answers that this connection takes the session: a refusal here
544 // must leave no read-only record behind, which is why this path does not
545 // run `attach_transport_locked`.
546 if note_binding_locked(&mut inner, session_id, &io) {
547 inner
548 .transports
549 .insert(session_id.to_string(), (io.clone(), capabilities.clone()));
550 return Some((io, capabilities));
551 }
552 tracing::warn!(
553 session = %session_id,
554 "the longest-waiting agent took nothing from this session; it waits for the next"
555 );
556 inner.parked.push((io, capabilities));
557 None
558 }
559
560 /// Hold `io` as this role's **standing** transport.
561 ///
562 /// A hosting runtime's connection is not a worker waiting for one job, so it
563 /// does not join [`Self::park_transport`]'s queue: a claim takes the oldest
564 /// entry and spends it, and a role whose hosting runtime serves four
565 /// sessions would have nothing left for the other three. It is held here
566 /// instead, where [`Self::staged_hosting_session`] finds it: the connection
567 /// is asked for each session the role opens, and one connection serves as
568 /// many as the role needs.
569 ///
570 /// A connection already standing is refreshed where it stands, the rule the
571 /// park uses too: a mount that says this twice is one runtime re-helloing.
572 pub fn stand_transport(&self, io: AdapterIo, capabilities: Vec<Capability>) {
573 let mut inner = self.inner.lock();
574 let standing = inner
575 .standing
576 .iter()
577 .position(|(held, _)| held.same_connection(&io));
578 match standing {
579 Some(index) => inner.standing[index].1 = capabilities,
580 None => inner.standing.push((io, capabilities)),
581 }
582 }
583
584 /// Take the answer to an `open` a hosting runtime gave, and bind the
585 /// connection that asked.
586 ///
587 /// The host named the session when it staged it, and the runtime names the
588 /// conversation it opened. Both are kept: the slot stays under the host's name
589 /// — that is the key every other lookup uses, and `session_tasks` binds through
590 /// it — while the runtime's own name and its resume handle go beside it, so the
591 /// next `open` for this family hands the handle back and the runtime resumes
592 /// rather than starting a second conversation for one chain.
593 ///
594 /// The transport is written here because the judgement has already run: the
595 /// runtime answered an `open`, and a refusal now would have to leave no
596 /// binding behind.
597 pub(crate) fn hosted_session_ready(
598 &self,
599 session_id: &str,
600 opened: &onlyne_proto::OpenedArgs,
601 io: AdapterIo,
602 capabilities: Vec<Capability>,
603 ) -> bool {
604 let mut inner = self.inner.lock();
605 let Some(slot) = inner.sessions.get_mut(session_id) else {
606 return false;
607 };
608 if !opened.session_id.is_empty() {
609 slot.session.backend_ref = serde_json::Value::String(opened.session_id.clone());
610 }
611 slot.resume_handle = opened.resume_handle.clone();
612 inner
613 .transports
614 .insert(session_id.to_string(), (io, capabilities));
615 true
616 }
617
618 /// The connection that serves one session, when its plugin is attached.
619 ///
620 /// A plugin names the session it was spawned for, and the slot's key is the
621 /// other spelling worth trying.
622 pub fn session_transport(&self, session_id: &str) -> Option<(AdapterIo, Vec<Capability>)> {
623 let inner = self.inner.lock();
624 if let Some(transport) = inner.transports.get(session_id) {
625 return Some(transport.clone());
626 }
627 let key = inner
628 .sessions
629 .iter()
630 .find(|(key, slot)| names_session(key, slot, session_id))
631 .map(|(key, _)| key.clone())?;
632 inner.transports.get(&key).cloned()
633 }
634
635 /// Tell the plugin serving one session to tear itself down, when that plugin
636 /// implements `recycle`. A plugin without the capability is skipped: the
637 /// caller's backend close stops the process either way.
638 pub async fn recycle_plugin(&self, task_id: &str, reason: &str, outcome: Option<Outcome>) {
639 let Some((io, capabilities)) = self.session_transport(task_id) else {
640 return;
641 };
642 if missing_capability(&capabilities, Capability::Recycle) {
643 tracing::debug!(
644 task = %task_id,
645 "plugin does not implement recycle; the host closes the resource"
646 );
647 return;
648 }
649 let args = RecycleArgs {
650 task_id: task_id.to_string(),
651 reason: reason.to_string(),
652 outcome,
653 };
654 if let Err(error) = io.notify(AdapterMsg::Host(HostOp::Recycle(args))).await {
655 tracing::warn!(error = %error, task = %task_id, "recycle frame did not reach the plugin");
656 }
657 }
658
659 /// Ask the plugin serving one session for a fresh observation.
660 ///
661 /// The plugin answers with a heartbeat report, which is the reducer's
662 /// evidence and the projection the operator reads. Answers whether a probe
663 /// frame actually went out, and `false` says nothing was asked: a session no
664 /// connection serves has no plugin to put the question to, and a `notify` that
665 /// failed left the frame in this process. The caller must not record either
666 /// as a probe that landed, because the answer a live plugin would have given
667 /// never existed — the projection that says a plugin is gone is written when
668 /// the session ends, never by this call.
669 pub async fn probe_plugin(&self, task_id: &str) -> bool {
670 let Some((io, _)) = self.session_transport(task_id) else {
671 tracing::warn!(task = %task_id, "probe found no plugin connection to ask");
672 return false;
673 };
674 let request = serde_json::json!({"task_id": task_id});
675 if let Err(error) = io.notify(AdapterMsg::Host(HostOp::Probe(request))).await {
676 tracing::warn!(error = %error, task = %task_id, "probe frame did not reach the plugin");
677 return false;
678 }
679 true
680 }
681
682 /// Hand one session's own sentence to its agent through the plugin's input.
683 ///
684 /// The turn-end rule has one owner, this client, so the client is what tells
685 /// a plugin-driven session that its turn ended without a completion
686 /// (`docs/v2-CONTRACT.md` §3c). The plugin injects [`NUDGE_TEXT`] through the
687 /// same channel an assignment's text takes and keeps no copy of it; the
688 /// frame carries no envelope, no prose, and no attachments, and it is not a
689 /// delivery: the task stays open and nothing in the plugin's turn
690 /// bookkeeping is reset by it.
691 ///
692 /// Answers whether the frame went out. `false` is the honest answer for a
693 /// session with no plugin to ask and for a plugin that declared no `inject`:
694 /// a drive that cannot be nudged must not be told it was, and the caller
695 /// settles the delivery at that turn end instead.
696 pub async fn nudge_plugin(&self, task_id: &str) -> bool {
697 let Some((io, capabilities)) = self.session_transport(task_id) else {
698 tracing::debug!(
699 task = %task_id,
700 "nudge found no plugin connection to hand the sentence to"
701 );
702 return false;
703 };
704 if missing_capability(&capabilities, Capability::Inject) {
705 tracing::debug!(
706 task = %task_id,
707 "plugin does not implement inject; the turn end settles the delivery"
708 );
709 return false;
710 }
711 let frame = AdapterMsg::Host(HostOp::Nudge {
712 task_id: task_id.to_string(),
713 text: NUDGE_TEXT.to_string(),
714 });
715 if let Err(error) = io.notify(frame).await {
716 tracing::warn!(error = %error, task = %task_id, "nudge frame did not reach the plugin");
717 return false;
718 }
719 true
720 }
721
722 /// Release the bindings served by one plugin connection.
723 ///
724 /// A graceful detach retires each idle session because the agent that owned
725 /// it has left. An attached transport preserves the idle resource because
726 /// the same agent is still reachable. A connection ending through
727 /// another path preserves the slot and resource for an agent reconnection and
728 /// starts the reconnect clock on it, which is what bounds how long a session
729 /// waits for an agent that is never coming back. Every released binding
730 /// retires its task progress clock. A slot carrying work remains under
731 /// lifecycle ownership, and it carries that same clock: a goodbye and a
732 /// silent drop both leave no heartbeat coming for the task it owes, and the
733 /// window is what ends a session whose agent never returns.
734 ///
735 /// The session ids a goodbye retired come back, because that retirement wrote
736 /// their rows: the resource close and, for a completed session, the agent's
737 /// exit. The server mirrors only what this client reports, and a publish
738 /// cannot run under this lock, so the caller is handed what to publish — the
739 /// same answer [`DispatchState::retire_dropped_ghosts`] gives its sweep.
740 pub fn release_connection(
741 &self,
742 session_id: Option<&str>,
743 io: &AdapterIo,
744 graceful_detach: bool,
745 ) -> Vec<String> {
746 let mut inner = self.inner.lock();
747 // A read-only connection ending is not the session losing its agent: the
748 // live connection still serves it, and its drop clock stays untouched.
749 let revived_connection = {
750 let before = inner.revived.len();
751 inner
752 .revived
753 .retain(|(_, revived, _)| !revived.same_connection(io));
754 before != inner.revived.len()
755 };
756 // A waiting agent whose socket ended leaves the queue and nothing else:
757 // the agents behind it are other connections, still waiting for work.
758 inner
759 .parked
760 .retain(|(parked, _)| !parked.same_connection(io));
761 let released: Vec<String> = inner
762 .transports
763 .iter()
764 .filter(|(session, (transport, _))| {
765 transport.same_connection(io)
766 && session_id.is_none_or(|mounted| mounted == session.as_str())
767 })
768 .map(|(session, _)| session.clone())
769 .collect();
770 let served_tasks: Vec<String> = released
771 .iter()
772 .map(|session| {
773 inner
774 .sessions
775 .iter()
776 .find(|(key, slot)| names_session(key, slot, session))
777 .map(|(_, slot)| {
778 slot.task_id
779 .clone()
780 .unwrap_or_else(|| slot.session.task_id.clone())
781 })
782 .unwrap_or_else(|| session.clone())
783 })
784 .collect();
785 for task_id in served_tasks {
786 inner.stall.forget(&task_id);
787 }
788 for session in &released {
789 inner.transports.remove(session);
790 }
791 if !revived_connection {
792 // The window of `[client] reconnect_grace_secs` starts wherever a
793 // session loses the connection that would have sent its next
794 // heartbeat and no other one is attached: the agent left without
795 // saying so — every session that connection served — or it said
796 // goodbye while its session still owed a task, which leaves no beat
797 // coming either. The session keeps its slot and its resource for the
798 // window, and the mount that returns inside it clears the stamp. A
799 // gracefully detached session holding no task needs no window: the
800 // idle retirement below takes that slot now.
801 let now = Instant::now();
802 for session in &released {
803 if let Some((_, slot)) = inner.sessions.iter_mut().find(|(key, slot)| {
804 names_session(key, slot, session)
805 && (!graceful_detach || slot.task_id.is_some())
806 }) {
807 slot.dropped_at = Some(now);
808 }
809 }
810 }
811 for session in &released {
812 // Nothing serves this name any more, so the first connection that
813 // mounted it read-only behind the one that just went becomes its
814 // transport; a session with no such connection keeps waiting out the
815 // reconnect grace, which is the sweep's to answer.
816 if let Some(key) = slot_key_named(&inner, session) {
817 promote_held_connection(&mut inner, &key);
818 }
819 }
820 let mut retired: Vec<String> = Vec::new();
821 let mut pending: Vec<PendingClose> = Vec::new();
822 if graceful_detach {
823 let idle: Vec<String> = released
824 .iter()
825 .filter_map(|session| {
826 inner
827 .sessions
828 .iter()
829 .find(|(key, slot)| {
830 names_session(key, slot, session) && slot.task_id.is_none()
831 })
832 .map(|(key, _)| key.clone())
833 })
834 .collect();
835 for key in idle {
836 // The id that travels is the session's own, the one whose row the
837 // retirement below is about to write.
838 let Some(task_id) = inner
839 .sessions
840 .get(&key)
841 .map(|slot| slot.session.task_id.clone())
842 else {
843 continue;
844 };
845 // A scoped session whose runtime can resume is not this goodbye's
846 // to end: the conversation is in the runtime's own store, and the
847 // delivery that comes next starts the command again. Releasing the
848 // process now is the same act the idle bound performs, and the
849 // session is resumed the same way.
850 if keeps_idle(&inner, &key) && resumable(&inner, &key) {
851 if suspend_locked(&mut inner, &key, &mut pending) {
852 retired.push(task_id);
853 }
854 continue;
855 }
856 let reason = inner
857 .sessions
858 .get(&key)
859 .and_then(|slot| stored_close_reason(&inner, &slot.session.task_id))
860 .unwrap_or(crate::backend::CloseReason::Completed);
861 if retire_idle_locked(&mut inner, &key, reason, &mut pending) {
862 retired.push(task_id);
863 }
864 }
865 }
866 // The goodbye above took the idle slots; the host work their sessions owe
867 // runs with the dispatch lock off it, so a plugin leaving does not make
868 // every other session of this role wait out a pane close.
869 drop(inner);
870 close_retired(pending);
871 retired
872 }
873}
874
875/// Whether this client may carry one plugin `send` to the server as its own.
876///
877/// Both rules are the client's to enforce because both are about what this
878/// process is willing to sign for: the server's acl answers for the wire, and it
879/// answers for a link this client fills with whichever `from` a plugin thought to
880/// write.
881///
882/// * A role principal besides this role is another session's voice. Every
883/// legitimate sender stamps its own name: the plugin's `send` carries the role
884/// its `hello` answer gave it, which is this client's role, and
885/// [`DispatchState::plugin_handoff`] builds its envelope from the stored role.
886/// A principal naming no role at all — a gateway or cluster sender — is left
887/// alone: no agent plugin mounts with one, and the host paths that build such
888/// envelopes never come through this door.
889/// * `control` is the operator's plane. The server mints every command and admits
890/// one from an admin link alone; a plugin that could queue one here would be
891/// handing its own orders to `on_control` through the back of a message send.
892///
893/// A refusal is an error the sender reads, and reads correctly: unlike a state
894/// report, a `send` the client will not carry has to be answered as a failure, or
895/// the plugin records a handoff that never left.
896///
897/// [`DispatchState::plugin_handoff`]: DispatchState::plugin_handoff
898fn send_is_authorised(inner: &DispatchInner, envelope: &Envelope) -> Result<()> {
899 if envelope.kind == MsgKind::Control {
900 tracing::warn!(
901 role = %inner.role,
902 to = %envelope.to,
903 "a plugin send naming a control command was refused"
904 );
905 return Err(anyhow!("this connection may not queue a control command"));
906 }
907 if let Some(from) = envelope.from.role_name()
908 && from != inner.role
909 {
910 tracing::warn!(
911 role = %inner.role,
912 from,
913 "a plugin send written as another role was refused"
914 );
915 return Err(anyhow!("this role may not send as {from}"));
916 }
917 Ok(())
918}
919
920#[cfg(test)]
921mod tests;