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}