Skip to main content

onlyne_client/session/dispatch/
delivery.rs

1use super::*;
2
3use super::env::{missing_capability, reject_unpaired_runtime, served_socket, session_env};
4use super::outbound::send_frame;
5use super::projection::{note_verdict, sync_session};
6use super::state::{
7    DispatchInner, DispatchState, SessionSlot, live_sessions, mint_tools_token, note_beat,
8    rebase_generation, render_tokens, slot_key_serving_task,
9};
10use super::transport::{
11    is_revived_connection, names_session, note_binding_locked, record_revived_connection,
12};
13use crate::delivery::{from_label, render, write_attachment};
14
15/// One delivery becomes the work of one session, or waits for the session its
16/// scope sends it to.
17///
18/// The role's scope decides which session takes the delivery and nothing else
19/// does: `oneshot` opens a session for every delivery, `task` hands the
20/// delivery to the session that already holds its family's conversation, and
21/// `role` hands it to whichever pooled session has waited longest. `Ok(None)`
22/// is the third answer and it is not a refusal: the scope's session is busy, or
23/// every slot is spent, and the row stays in flight for the pull that comes
24/// after one frees (plan §5 `max_sessions`).
25pub fn dispatch(state: &DispatchState, envelope: &Envelope) -> Result<Option<SessionRef>> {
26    let causality = envelope
27        .causality
28        .clone()
29        .context("task envelope missing causality.task")?;
30    let task_id = causality.task.clone();
31    let family = scope::family_of(&causality);
32    // The role's prose, read from its one owner before the dispatch lock is
33    // taken. A session this delivery opens needs it at spawn — a runtime with no
34    // instruction channel of its own is handed it in the workspace — and the
35    // assignment frame that follows carries the same value
36    // (`docs/v2-CONTRACT.md` §3b).
37    let prose = state.role_prose();
38    let mut inner = state.inner.lock();
39    // A task this role already serves rides its own slot, and the slot that
40    // still holds delivery rights is the one that serves it: staging the payload
41    // on a read-only revival would hand the work to an agent that may answer for
42    // it but may be handed nothing.
43    if let Some(session) = slot_key_serving_task(&inner, &task_id)
44        .and_then(|key| inner.sessions.get_mut(&key))
45        .map(|slot| {
46            if slot.payload.is_none() {
47                slot.payload = Some(envelope.clone());
48                slot.causality = causality.clone();
49            }
50            slot.session.clone()
51        })
52    {
53        inner.stall.note_assigned(&task_id, Instant::now());
54        return Ok(Some(session));
55    }
56    match scope::placement(&inner, &family) {
57        scope::Placement::Bind(key) => {
58            let session = bind_delivery(&mut inner, &key, envelope, &causality, &task_id)?;
59            return Ok(Some(session));
60        }
61        scope::Placement::Resume(key) => {
62            let session =
63                resume_delivery(&mut inner, &key, envelope, &causality, &task_id, &prose)?;
64            return Ok(Some(session));
65        }
66        scope::Placement::Wait => {
67            tracing::debug!(
68                task = %task_id,
69                family = %family,
70                "the delivery waits: the session its scope names is serving another delivery"
71            );
72            return Ok(None);
73        }
74        scope::Placement::Open => {}
75    }
76    // What is left opens a session, so this is where `max_sessions` bites.
77    if live_sessions(&inner) >= inner.max_sessions as usize {
78        tracing::debug!(
79            task = %task_id,
80            max_sessions = inner.max_sessions,
81            "no session slot is free; the delivery waits"
82        );
83        return Ok(None);
84    }
85    let session = open_session(&mut inner, envelope, &causality, &task_id, &family, &prose)?;
86    Ok(Some(session))
87}
88
89/// Open one session for a delivery: the `oneshot` path, and the first delivery
90/// of a family or a role pool.
91///
92/// The session's own id is the delivery's task id, which is what makes the row
93/// it is born onto this delivery's row: a session that goes on to serve the
94/// family's later deliveries keeps this id and gains a binding for each of them.
95fn open_session(
96    inner: &mut DispatchInner,
97    envelope: &Envelope,
98    causality: &Causality,
99    task_id: &str,
100    family: &str,
101    prose: &str,
102) -> Result<SessionRef> {
103    reject_unpaired_runtime(inner)?;
104    let session_id = task_id.to_string();
105    let command = render_tokens(&inner.command, &session_id, task_id);
106    let env = session_env(
107        &inner.role,
108        &session_id,
109        task_id,
110        &inner.topology,
111        // One tree answers both halves of this spawn: the cwd below and the
112        // socket the plugin dials, so a session whose workspace resolves to a
113        // short endpoint is handed the served path directly.
114        &served_socket(&inner.workspace),
115    );
116    // The tools token is minted before the spawn, because the child that mounts
117    // `onlyne mcp` is handed it there, and it is recorded on the slot below so
118    // the session's own state is the only other place it lives: a capability
119    // never reaches a log, a fault, or a ledger row (`docs/v2-CONTRACT.md` §3b).
120    let tools_token = mint_tools_token();
121    // A hosting runtime is already resident, so the host does not start a
122    // process for this session — it asks. The slot is staged either way: the slot
123    // is what the delivery, the session row and the scope all name. What differs
124    // is whether a `SpawnSpec` went out first.
125    //
126    // The staged `SessionRef` is the host's own name for the session and carries
127    // no command, because there is no argv to carry — the process belongs to the
128    // runtime. The connection answers `open` with whatever *it* calls that
129    // conversation, and `hosted_session_ready` puts the two together.
130    let hosted = !inner.standing.is_empty();
131    let session = if hosted {
132        SessionRef {
133            task_id: session_id.clone(),
134            backend: "hosting".into(),
135            backend_ref: serde_json::Value::Null,
136            generation: 1,
137        }
138    } else {
139        inner.backend.spawn(SpawnSpec {
140            cwd: inner.workspace.clone(),
141            task_id: session_id.clone(),
142            command: command.clone(),
143            env,
144            tools_token: tools_token.clone(),
145            prose: prose.to_string(),
146            focus: None,
147            placement: None,
148            rename: None,
149        })?
150    };
151    let family = scope::keys_on_family(inner.session_policy.scope).then(|| family.to_string());
152    // Everything that can still refuse this delivery runs inside `stage_slot`,
153    // and a refusal past this line leaves a live pane, tab, or child process
154    // behind. Nothing in `sessions` names it, so `close_all`, both reconnect
155    // sweeps, and `live_sessions` never see it again: the resource and the
156    // bridge's record of it leak for the life of the client, which is the shape
157    // a SQLite error in `open_task` or `feed_created` used to leave. The slot is
158    // given back here, and the close runs after the lock is off it.
159    if let Err(error) = stage_slot(
160        inner,
161        Spawned {
162            session: session.clone(),
163            command,
164            tools_token,
165        },
166        task_id,
167        family,
168        causality,
169        envelope,
170    ) {
171        // A hosted session never opened a resource, so there is none to hand
172        // back: the runtime owns the process and the `open` that fails undoes
173        // itself on that side.
174        let backend = Arc::clone(&inner.backend);
175        if !hosted
176            && let Err(close_error) =
177                backend.close(&session, crate::backend::CloseReason::Fault, false)
178        {
179            tracing::warn!(
180                task = %task_id,
181                backend = %session.backend,
182                resource = %session.backend_ref,
183                error = %close_error,
184                "the resource of a dispatch that did not land was not given back"
185            );
186        }
187        return Err(error);
188    }
189    inner.stall.note_assigned(task_id, Instant::now());
190    Ok(session)
191}
192
193/// Hand one delivery to a session this client already holds.
194///
195/// The session's own dimensions do not move: a new delivery changes which
196/// delivery it serves, not what it is. So the binding is taken on its own —
197/// `open_task` writes the delivery's record, the binding is opened in this
198/// client's store, and the session's row advances one sequence so the publish
199/// that follows carries the new binding to the mirror. Without that step the
200/// server's gate would take nothing at the watermark it already holds, and the
201/// row an operator reads would keep naming the delivery this session has
202/// finished.
203/// The scope word an assignment carries: the config's own spelling, so the
204/// runtime compares against what the operator wrote rather than a second
205/// vocabulary of this crate's own.
206fn scope_word(scope: &onlyne_config::SessionScope) -> String {
207    scope.as_str().to_string()
208}
209
210fn bind_delivery(
211    inner: &mut DispatchInner,
212    key: &str,
213    envelope: &Envelope,
214    causality: &Causality,
215    task_id: &str,
216) -> Result<SessionRef> {
217    let session = inner
218        .sessions
219        .get(key)
220        .map(|slot| slot.session.clone())
221        .ok_or_else(|| anyhow!("no session answers to {key}"))?;
222    let session_id = session.task_id.clone();
223    inner.store.open_task(causality, task_cause(causality))?;
224    advance_session_row(inner, &session_id);
225    inner.store.bind_task(&session_id, task_id)?;
226    if let Some(slot) = inner.sessions.get_mut(key) {
227        slot.task_id = Some(task_id.to_string());
228        slot.payload = Some(envelope.clone());
229        slot.causality = causality.clone();
230        slot.origin = Some(envelope.from.clone());
231        slot.msg_id = None;
232        slot.read_only = false;
233        slot.idle_since = None;
234    }
235    tracing::info!(
236        session = %session_id,
237        task = %task_id,
238        "the delivery joined the session its scope keeps for it"
239    );
240    inner.stall.note_assigned(task_id, Instant::now());
241    Ok(session)
242}
243
244/// Resume one suspended session for a delivery, and bind it.
245///
246/// A suspended session's conversation lives in the runtime's own store, and the
247/// client's half of resuming it is to start the command the session was born
248/// with again: the same argv — because the command carries the runtime's own
249/// session key, and it may interpolate the delivery into it — and the
250/// environment rebuilt for the delivery being served. The runtime resumes the
251/// conversation it was asked for. This client never composes a history summary
252/// to hand a model: a summary it wrote would be context pollution it also
253/// invented (plan §10).
254fn resume_delivery(
255    inner: &mut DispatchInner,
256    key: &str,
257    envelope: &Envelope,
258    causality: &Causality,
259    task_id: &str,
260    prose: &str,
261) -> Result<SessionRef> {
262    reject_unpaired_runtime(inner)?;
263    let (session_id, command, tools_token) = inner
264        .sessions
265        .get(key)
266        .map(|slot| {
267            (
268                slot.session.task_id.clone(),
269                slot.command.clone(),
270                slot.tools_token.clone(),
271            )
272        })
273        .ok_or_else(|| anyhow!("no session answers to {key}"))?;
274    let env = session_env(
275        &inner.role,
276        &session_id,
277        task_id,
278        &inner.topology,
279        &served_socket(&inner.workspace),
280    );
281    // The resumed process runs the same session, so it is handed the same tools
282    // token: the token belongs to the session and lives in its slot, and this
283    // client's binding is what the mount's first call is measured against. A
284    // session that *reopens* — a new slot for the same task — mints a new one.
285    let resumed = inner.backend.spawn(SpawnSpec {
286        cwd: inner.workspace.clone(),
287        task_id: session_id.clone(),
288        command,
289        env,
290        tools_token,
291        prose: prose.to_string(),
292        focus: None,
293        placement: None,
294        rename: None,
295    })?;
296    inner.bridge.track_live(resumed.clone());
297    // The row moves next, and the store's own binding write is handed back once
298    // it is final: `Resume` attaches the resource again while the generation
299    // stays live, and the row then names the process that has just started
300    // rather than the one the release gave back.
301    if let Err(error) = feed_resumed(&inner.bridge, &inner.store, &session_id) {
302        tracing::warn!(
303            session = %session_id,
304            error = %error,
305            "a resumed session's row was not written"
306        );
307    }
308    let _ = inner.store.release_binding(&session_id, &session_id);
309    inner.store.open_task(causality, task_cause(causality))?;
310    inner.store.bind_task(&session_id, task_id)?;
311    let now = Instant::now();
312    if let Some(slot) = inner.sessions.get_mut(key) {
313        slot.session = resumed.clone();
314        slot.suspended = false;
315        slot.task_id = Some(task_id.to_string());
316        slot.payload = Some(envelope.clone());
317        slot.causality = causality.clone();
318        slot.origin = Some(envelope.from.clone());
319        slot.msg_id = None;
320        slot.ready = false;
321        slot.idle_since = None;
322        // The runtime that resumes this session is a new process, and its plugin
323        // has to dial before any frame of this delivery can reach it: the
324        // reconnect window is the one that reads a plugin which never arrives.
325        slot.dropped_at = (!inner.backend.self_driven()).then_some(now);
326        slot.last_beat = Some(now);
327    }
328    tracing::info!(
329        session = %session_id,
330        task = %task_id,
331        backend = %resumed.backend,
332        "suspended session resumed for its scope's next delivery"
333    );
334    inner.stall.note_assigned(task_id, Instant::now());
335    Ok(resumed)
336}
337
338/// Advance one session's watermark one sequence past where it stands.
339///
340/// The row's dimensions are unchanged, so this is the write-side half of "this
341/// session was written about": the store's gate takes only a strictly newer
342/// version, and the publish that carries a new binding has to clear it. A
343/// failure here is not fatal — the binding is still taken in this client's own
344/// store, and the next publish of a real transition carries it — so it is
345/// logged and the caller goes on.
346fn advance_session_row(inner: &mut DispatchInner, session_id: &str) {
347    let row = match inner.store.get_session(session_id) {
348        Ok(Some(row)) => row,
349        _ => return,
350    };
351    let seq = row.seq.max(0) as u64 + 1;
352    if let Err(error) =
353        inner
354            .store
355            .bump_session_version(session_id, row.generation.max(0) as u64, seq)
356    {
357        tracing::warn!(
358            session = %session_id,
359            error = %error,
360            "a rebound session's row was not advanced; the binding reaches the mirror with its next write"
361        );
362    }
363}
364
365/// The resource one delivery was spawned onto, as the slot has to record it:
366/// the reference the backend answers to, the command that opened it, and the
367/// tools token this client minted for the session it will serve.
368struct Spawned {
369    session: SessionRef,
370    command: Vec<String>,
371    tools_token: String,
372}
373
374/// Land one spawned session in this role's bookkeeping: the task's own record,
375/// the session row it is born onto, the slot that holds its payload, and the
376/// stall clock that answers for its work.
377///
378/// The slot is keyed by the session's own id, which a client-held session takes
379/// from the delivery that opened it: that is the key the session keeps for every
380/// delivery the scope later hands it, and the spelling its stored row carries.
381///
382/// Split from `open_session` so a failure names itself as one: every fallible step
383/// lives here, and the caller owns the single undo that matters — a resource the
384/// backend has already opened. The steps run in the order that keeps the ledger
385/// causal: the task record opens before the session row it hosts, and the row
386/// before the slot that can report against it.
387///
388/// This function owns the in-memory half of its own undo. The bridge is tracked
389/// first because `feed_created` reads it for the session's generation, and a step
390/// that refuses would otherwise leave that track behind: the reconciler answers
391/// for sessions this client holds, and a live entry no slot addresses has nothing
392/// left to report about it. The resource is the caller's to hand back, because
393/// only the caller can give it up off the dispatch lock.
394fn stage_slot(
395    inner: &mut DispatchInner,
396    spawned: Spawned,
397    task_id: &str,
398    family: Option<String>,
399    causality: &Causality,
400    envelope: &Envelope,
401) -> Result<()> {
402    inner.bridge.track_live(spawned.session.clone());
403    let landed = (|| -> Result<()> {
404        // The task's own record opens with the session that serves it, out of the
405        // causality that named the task. A redelivery that found a slot already
406        // serving above never reaches this line, so the chain columns are the chain
407        // the session opened on; a re-dispatch after retirement refreshes them.
408        inner.store.open_task(causality, task_cause(causality))?;
409        // A row this task already carries is the record of the session that served it
410        // before, and the session staged here is born onto it: `rebase_born_session`
411        // moves that row to a generation of its own, and a task with no row keeps the
412        // plain seed below. Either way the row the feeds land on starts at
413        // `Booting`/`Detached`, under a watermark this session's own count can clear.
414        rebase_born_session(inner, task_id);
415        feed_created(&inner.bridge, &inner.store, task_id)?;
416        feed_dispatched(&inner.bridge, &inner.store, task_id);
417        // A plugin-mode session's liveness is the heartbeat its connection sends and
418        // nothing else, and that connection has not spoken yet: the window of
419        // `[client] reconnect_grace_secs` starts at birth, so a spawn whose plugin
420        // never dials is a ghost the sweep can see rather than a slot that holds its
421        // resource forever. The mount that attaches clears the stamp. A self-driven
422        // backend owns its agent and answers no adapter socket, so it never has a
423        // heartbeat to read and its lifecycle, not this clock, is what ends it.
424        let dropped_at = (!inner.backend.self_driven()).then(Instant::now);
425        // The liveness stamp starts with the slot, so a session whose plugin mounts
426        // and then never sends a frame is readable as silent rather than as a
427        // session nobody can judge.
428        let last_beat = Some(Instant::now());
429        inner.sessions.insert(
430            spawned.session.task_id.clone(),
431            SessionSlot {
432                session: spawned.session.clone(),
433                // A session the host asked a hosting runtime for is named by the
434                // host until the runtime answers `open`; nothing may resume it
435                // before that, and a session with no handle is one the runtime
436                // will open fresh next time.
437                resume_handle: None,
438                task_id: Some(task_id.to_string()),
439                ready: false,
440                payload: Some(envelope.clone()),
441                msg_id: None,
442                origin: Some(envelope.from.clone()),
443                causality: causality.clone(),
444                dropped_at,
445                last_beat,
446                read_only: false,
447                family,
448                idle_since: None,
449                suspended: false,
450                tools_token: spawned.tools_token,
451                delivered_roles: BTreeSet::new(),
452                opened_at: Instant::now(),
453                command: spawned.command,
454                keeps_idle: !matches!(
455                    inner.session_policy.scope,
456                    onlyne_config::SessionScope::Oneshot
457                ),
458            },
459        );
460        inner.stall.note_assigned(task_id, Instant::now());
461        Ok(())
462    })();
463    if landed.is_err() {
464        inner.bridge.untrack_live(&spawned.session.task_id);
465    }
466    landed
467}
468
469/// How a delivery reached this role, which is what the task record's `kind`
470/// column holds. A task with a parent above it was handed down from another
471/// session's work; one without was given to this role directly. The envelope's
472/// own message kind is not that answer: only a task-shaped delivery ever reaches
473/// a session, so the kind says nothing the chain does not.
474fn task_cause(causality: &Causality) -> &'static str {
475    if causality.parent_task.is_some() {
476        "relay"
477    } else {
478        "root"
479    }
480}
481
482/// Move the row a re-dispatched task already carries onto the generation of the
483/// session this dispatch stages.
484///
485/// This client keeps `client.db` across a restart and a session row is keyed by
486/// its task, so the session staged here is born onto whatever row that task
487/// already carries; a first dispatch is the only shape with none. That row is the
488/// record of the session that served the task before, one no slot of this process
489/// holds any more, and none of what it holds can be inherited. Its phase and its
490/// resource describe a session that no longer exists, and the feeds below would
491/// have to move them from states that refuse them: `resource_attach` from a
492/// closed resource is `UndefinedTransition`, which is how the log reads when a
493/// ghost was swept before the task came back. Its watermark is worse, because a
494/// row it stands on accepts nothing that reads older: the client's own feeds and
495/// the plugin's beats share the one counter, and the plugin is a new process
496/// whose sequence starts at its base again, below anything a session that lived a
497/// while left. Every frame the new session sends is then dropped as a stale
498/// duplicate, the turn its agent really ran never reaches the row, and the settle
499/// door refuses the completion of work that happened.
500///
501/// The new generation itself is [`rebase_generation`]'s; the body here is the
502/// born tuple of any session — the same `Observation::initial` a fresh row is
503/// seeded from, with the role's reconcile policy carried over, so the two ways a
504/// session's row comes into being cannot drift. What attests the old generation
505/// dead is this client's own bookkeeping: this call is reached only because no
506/// slot of this process serves the task, so nothing it holds speaks for that row
507/// any more.
508fn rebase_born_session(inner: &DispatchInner, task_id: &str) {
509    let verdict = rebase_generation(inner, task_id, |stored| {
510        Observation::initial(stored.isolate_after, stored.terminate_after)
511    });
512    match verdict {
513        Ok(Some(Verdict::Applied(next))) => tracing::info!(
514            task = %task_id,
515            generation = next.version.generation,
516            "the row a re-dispatched session is born onto was rebased onto a new generation"
517        ),
518        Ok(Some(verdict)) => tracing::warn!(
519            task = %task_id,
520            ?verdict,
521            "the row of a re-dispatched session was not rebased"
522        ),
523        Ok(None) => {}
524        Err(error) => tracing::warn!(
525            task = %task_id,
526            error = %error,
527            "the row of a re-dispatched session was not rebased"
528        ),
529    }
530}
531
532/// Adapter facts that make a session usable for its task.
533pub struct ReadyNotice {
534    pub task_id: String,
535    pub session_id: String,
536    pub generation: u64,
537    /// Adapter transport for plugin-driven backends. A self-driven backend owns
538    /// its agent and therefore reports ready without a socket.
539    pub io: Option<AdapterIo>,
540    pub capabilities: Vec<Capability>,
541}
542
543/// Report the session ready and hand its held payload to its agent. The `ready`
544/// row reaches the ledger before either the backend delivery or adapter frame,
545/// which is the causal order §6 requires.
546///
547/// The text both hand-over paths carry is [`crate::delivery::render`]'s answer,
548/// rendered here and nowhere else: a plugin injects it, a self-driven backend
549/// prompts with it, and neither composes a delivery of its own. The delivery's
550/// image is written into the workspace first, because the line that names it is
551/// the line the model reads.
552pub async fn on_ready(state: &DispatchState, notice: ReadyNotice, prose: &str) -> Result<()> {
553    let ReadyNotice {
554        task_id,
555        session_id,
556        generation,
557        io,
558        capabilities,
559    } = notice;
560    let (payload, target, session, backend, version, workspace) = {
561        let mut inner = state.inner.lock();
562        let backend = Arc::clone(&inner.backend);
563        // A ready report binds its connection to the session as much as a mount
564        // does, so it runs the same judgement §1 (b) hangs on: a connection that
565        // returns to a session a newer connection already serves takes nothing,
566        // leaves nothing marked ready, and is held for that task's completion.
567        if let Some(connection) = io.as_ref() {
568            if is_revived_connection(&inner, connection) {
569                return Ok(());
570            }
571            if !note_binding_locked(&mut inner, &session_id, connection) {
572                record_revived_connection(
573                    &mut inner,
574                    &session_id,
575                    connection.clone(),
576                    capabilities.clone(),
577                );
578                return Ok(());
579            }
580        }
581        let slot = inner
582            .sessions
583            .values_mut()
584            .find(|slot| {
585                // The delivery is the binding this session is serving now, and
586                // the session's own id is the other spelling a plugin may
587                // report under: a session that has served several deliveries
588                // answers to both, and the binding is what a ready report names.
589                slot.task_id.as_deref() == Some(task_id.as_str())
590                    || slot.session.task_id == task_id
591                    || slot
592                        .session
593                        .backend_ref
594                        .get("id")
595                        .and_then(|value| value.as_str())
596                        == Some(session_id.as_str())
597            })
598            .ok_or_else(|| anyhow!("unknown session for {task_id}"))?;
599        if slot.read_only {
600            return Ok(());
601        }
602        // The hand-off runs once per session: a plugin that reports ready
603        // after the assignment already left finds the payload gone.
604        let Some(payload) = slot.payload.take() else {
605            return Ok(());
606        };
607        slot.origin = Some(payload.from.clone());
608        slot.ready = true;
609        let session = slot.session.clone();
610        let verdict = feed_ready(&inner.bridge, &inner.store, &task_id)?;
611        if matches!(verdict, Verdict::Applied(_)) {
612            // The ready report is a frame this session sent that the reducer
613            // took, so it is liveness like any other: the sweep reads the stamp
614            // rather than the socket, and a session whose plugin passed the
615            // barrier and then went quiet has to be readable as quiet.
616            note_beat(&mut inner, &task_id, Instant::now());
617        }
618        let version = note_verdict(&verdict, &task_id).unwrap_or(Version::new(generation, 0));
619        (
620            payload,
621            io,
622            session,
623            backend,
624            version,
625            inner.workspace.clone(),
626        )
627    };
628    // The ready report reaches the server before the payload reaches the agent.
629    send_frame(
630        state,
631        ClientOp::Report(Report::Ready {
632            task_id: task_id.clone(),
633            session_id: session_id.clone(),
634            generation: version.generation,
635            seq: version.seq,
636            cluster_ref: None,
637        }),
638    )
639    .await?;
640    sync_session(state, &task_id).await?;
641    // The body travels as it was written, the material is quoted, and the
642    // attachments are the paths just written under the workspace. Upstream
643    // material is the one input no delivery carries yet: the field a sender
644    // fills it from is slice 3b's `complete(details)`.
645    let attachments = write_attachment(&workspace, &task_id, &payload)
646        .into_iter()
647        .collect::<Vec<_>>();
648    let text = render(
649        &from_label(&payload.from),
650        payload.body.text.as_deref().unwrap_or_default(),
651        None,
652        &attachments,
653    );
654    match (backend.self_driven(), target) {
655        (true, None) => {
656            // A self-driven drive has no heartbeats, so the dispatch path feeds
657            // the turn-started fact here — the same fact a plugin's beat would
658            // carry — before the backend starts the turn. The never-ran guard
659            // reads the row this writes, so a completion the agent files
660            // through its tools mount during the turn passes it.
661            state.feed_turn_started(&task_id);
662            backend.deliver(&session, &task_id, &text)
663        }
664        (true, Some(_)) => Err(anyhow!(
665            "self-driven session {session_id} unexpectedly has an adapter transport"
666        )),
667        (false, Some(target)) if capabilities.contains(&Capability::Inject) => {
668            let assign = AssignArgs {
669                envelope: Box::new(payload),
670                prose: prose.to_string(),
671                text,
672                attachments,
673                task_id,
674                generation,
675                session_id: Some(session_id.to_string()),
676                scope: Some(scope_word(&state.session_policy().scope)),
677                parent: None,
678            };
679            target
680                .notify(AdapterMsg::Host(HostOp::Assign(assign)))
681                .await
682                .map_err(|e| anyhow!(e))
683        }
684        (false, Some(target)) => target
685            .notify(AdapterMsg::Host(HostOp::ConfigGet(
686                onlyne_proto::ConfigGetArgs {
687                    key: format!("stdin:{text}"),
688                },
689            )))
690            .await
691            .map_err(|e| anyhow!(e)),
692        (false, None) => Err(anyhow!(
693            "adapter-backed session {session_id} reported ready without a transport"
694        )),
695    }
696}
697
698impl DispatchState {
699    /// Bind a plugin transport to one staged session and hand it the payload.
700    ///
701    /// Both hand-over paths run through here: a plugin that mounted first, and a
702    /// session staged first. The report that marks the session ready leaves
703    /// before the assignment, which is the causal order §6 line 285 fixes.
704    ///
705    /// The caller names the session. The delivery it is serving is read here,
706    /// because the two are one id only until a scope hands the session a second
707    /// delivery: the ready report and the reducer facts answer for the delivery,
708    /// and the session is the row they land on.
709    pub async fn hand_session(
710        &self,
711        session_id: &str,
712        io: AdapterIo,
713        capabilities: Vec<Capability>,
714    ) -> Result<()> {
715        let prose = self.role_prose();
716        let (task_id, generation) = self.served_delivery(session_id);
717        on_ready(
718            self,
719            ReadyNotice {
720                task_id,
721                session_id: session_id.to_string(),
722                generation,
723                io: Some(io),
724                capabilities,
725            },
726            &prose,
727        )
728        .await
729    }
730
731    /// Route one staged session to the thing that serves it.
732    ///
733    /// A self-driven backend owns its agent and takes the payload immediately,
734    /// without an adapter socket. Otherwise the session's own connection comes
735    /// first: a plugin the client spawned mounts with this session's id in
736    /// `ONLYNE_SESSION_ID`, and a plugin that reconnected mounts with it again,
737    /// so its assignment rides that socket alone. A plugin parked for the role
738    /// takes the next staged session, once. A plugin-driven session with neither
739    /// waits for its mount. Answers whether the payload had somewhere to go.
740    pub async fn hand_staged(&self, session_id: &str) -> Result<bool> {
741        let self_driven = self.inner.lock().backend.self_driven();
742        if self_driven {
743            let prose = self.role_prose();
744            let (task_id, generation) = self.served_delivery(session_id);
745            on_ready(
746                self,
747                ReadyNotice {
748                    task_id,
749                    session_id: session_id.to_string(),
750                    generation,
751                    io: None,
752                    capabilities: Vec::new(),
753                },
754                &prose,
755            )
756            .await?;
757            return Ok(true);
758        }
759        // Three sources, in the order that keeps every role's behaviour the
760        // shape it had: the session's own transport, then a parked agent — one
761        // the client spawned for work in hand, so it is the more specific match.
762        // A hosting runtime's connection is not a third source here: the session
763        // reaches it by being asked for, so a connection that has already been
764        // lent one is lent nothing by this path.
765        let transport = self
766            .session_transport(session_id)
767            .or_else(|| self.claim_parked_transport(session_id));
768        let Some((io, capabilities)) = transport else {
769            return Ok(false);
770        };
771        self.hand_session(session_id, io, capabilities).await?;
772        Ok(true)
773    }
774
775    /// The delivery one session is serving, and the generation it runs under.
776    ///
777    /// A session between deliveries answers with its own id, which is the
778    /// spelling the row it was born onto was written at: the payload is handed
779    /// over in the same breath a delivery is bound, so this is the idle-session
780    /// fallback rather than the ordinary reading.
781    fn served_delivery(&self, session_id: &str) -> (String, u64) {
782        let inner = self.inner.lock();
783        let slot = inner
784            .sessions
785            .iter()
786            .find(|(key, slot)| names_session(key, slot, session_id))
787            .map(|(_, slot)| slot);
788        match slot {
789            Some(slot) => (
790                slot.task_id
791                    .clone()
792                    .unwrap_or_else(|| slot.session.task_id.clone()),
793                slot.session.generation,
794            ),
795            None => (session_id.to_string(), 1),
796        }
797    }
798
799    /// Hand one note to the session already serving this role's work.
800    ///
801    /// A note carries no task, so it owns no session: §3's note is a message to
802    /// an agent that is already running, and the plan refuses one whose role is
803    /// offline (`note_queue` off). A role with no running agent has nothing to
804    /// answer it, which is what the caller reports. Answers whether an agent
805    /// took the note.
806    pub async fn inject_note(&self, envelope: &Envelope) -> bool {
807        let Some((task_id, session_id)) = self.ready_session() else {
808            return false;
809        };
810        let Some((io, capabilities)) = self.session_transport(&session_id) else {
811            return false;
812        };
813        if missing_capability(&capabilities, Capability::Inject) {
814            tracing::debug!(task = %task_id, "plugin takes no message mid-task");
815            return false;
816        }
817        let generation = self.session_generation(&task_id).unwrap_or(1);
818        // A note reaches a live agent the way a task does, so it travels the same
819        // one template and its image is written under the workspace first.
820        let workspace = self.inner.lock().workspace.clone();
821        let attachments = write_attachment(&workspace, &task_id, envelope)
822            .into_iter()
823            .collect::<Vec<_>>();
824        let text = render(
825            &from_label(&envelope.from),
826            envelope.body.text.as_deref().unwrap_or_default(),
827            None,
828            &attachments,
829        );
830        let assign = AssignArgs {
831            envelope: Box::new(envelope.clone()),
832            prose: self.role_prose(),
833            text,
834            attachments,
835            task_id,
836            generation,
837            session_id: Some(session_id.to_string()),
838            scope: Some(scope_word(&self.session_policy().scope)),
839            parent: None,
840        };
841        io.notify(AdapterMsg::Host(HostOp::Assign(assign)))
842            .await
843            .is_ok()
844    }
845
846    /// The session a mid-task message can join: a ready slot serving a task,
847    /// answered as (task, session key).
848    fn ready_session(&self) -> Option<(String, String)> {
849        let inner = self.inner.lock();
850        inner
851            .sessions
852            .iter()
853            .find(|(_, slot)| slot.ready && slot.task_id.is_some() && !slot.read_only)
854            .map(|(key, slot)| {
855                (
856                    slot.task_id.clone().unwrap_or_else(|| key.clone()),
857                    key.clone(),
858                )
859            })
860    }
861}