Skip to main content

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;