pub struct DispatchState { /* private fields */ }Implementations§
Source§impl DispatchState
impl DispatchState
Sourcepub async fn hand_session(
&self,
session_id: &str,
io: AdapterIo,
capabilities: Vec<Capability>,
) -> Result<()>
pub async fn hand_session( &self, session_id: &str, io: AdapterIo, capabilities: Vec<Capability>, ) -> Result<()>
Bind a plugin transport to one staged session and hand it the payload.
Both hand-over paths run through here: a plugin that mounted first, and a session staged first. The report that marks the session ready leaves before the assignment, which is the causal order §6 line 285 fixes.
The caller names the session. The delivery it is serving is read here, because the two are one id only until a scope hands the session a second delivery: the ready report and the reducer facts answer for the delivery, and the session is the row they land on.
Sourcepub async fn hand_staged(&self, session_id: &str) -> Result<bool>
pub async fn hand_staged(&self, session_id: &str) -> Result<bool>
Route one staged session to the thing that serves it.
A self-driven backend owns its agent and takes the payload immediately,
without an adapter socket. Otherwise the session’s own connection comes
first: a plugin the client spawned mounts with this session’s id in
ONLYNE_SESSION_ID, and a plugin that reconnected mounts with it again,
so its assignment rides that socket alone. A plugin parked for the role
takes the next staged session, once. A plugin-driven session with neither
waits for its mount. Answers whether the payload had somewhere to go.
Sourcepub async fn inject_note(&self, envelope: &Envelope) -> bool
pub async fn inject_note(&self, envelope: &Envelope) -> bool
Hand one note to the session already serving this role’s work.
A note carries no task, so it owns no session: §3’s note is a message to
an agent that is already running, and the plan refuses one whose role is
offline (note_queue off). A role with no running agent has nothing to
answer it, which is what the caller reports. Answers whether an agent
took the note.
Source§impl DispatchState
A tools mount’s own bookkeeping, stamped from this client’s record.
impl DispatchState
A tools mount’s own bookkeeping, stamped from this client’s record.
The scope is DispatchState::tools_scope’s answer — one locked pass over
the connection’s binding and the session’s slot — and the stamps below are
pure over it, so a frame cannot be attributed from a state that moved between
the lookup and the write.
Sourcepub const TOOLS_GONE_MESSAGE: &'static str = "token names no live session"
pub const TOOLS_GONE_MESSAGE: &'static str = "token names no live session"
The one sentence a tools mount reads when the session its token names is
gone.
The handshake and the per-frame gate answer it. The adapter’s welcome
refusal carries a code and a message and has no field slot, so the
sentence is where the field name goes; the per-frame Self::tools_gone
names the same field in the slot it has.
Sourcepub fn tools_gone() -> ResBody
pub fn tools_gone() -> ResBody
The refusal a tools frame earns when its connection speaks for no live
session.
The field is named and the token’s own value is not: the token is a
capability the caller already holds, and a sentence that echoed it would
put one in a fault a model reads (AGENTS.md §8). A session retires
between one frame’s liveness check and its stamping, and both doors read
one answer because it is one fact.
Sourcepub fn stamp_tools_task(
&self,
io: &AdapterIo,
task_id: &mut String,
) -> Result<(), ResBody>
pub fn stamp_tools_task( &self, io: &AdapterIo, task_id: &mut String, ) -> Result<(), ResBody>
Stamp one report or handoff frame from a tools mount with the task
the session serves.
An empty task_id is this path’s spelling of “the session’s own open
task” and is replaced with it. A non-empty one that disagrees with the
session’s own binding is refused rather than silently overwritten — a
caller wrong about the session it speaks for is a bug in the caller, and
a correction that hides it is the wrong answer. A session holding no open
task answers invalid on the same field, because a completion for a task
nobody holds cannot be recorded.
Sourcepub fn stamp_tools_send(
&self,
io: &AdapterIo,
envelope: &mut Envelope,
) -> Result<(), ResBody>
pub fn stamp_tools_send( &self, io: &AdapterIo, envelope: &mut Envelope, ) -> Result<(), ResBody>
The sender one send frame from a tools mount leaves as, and the chain
that frame starts.
The bridge supplies the recipient, the body, the kind, and the image; the
role the frame leaves as comes from this client’s own record, so nothing
here is read off the frame or off welcome. A session serving no
delivery is refused: a send it made would belong to no work at all.
handoff continues a family and send starts one, so a task send is a
root — a fresh task id at hop 0, attempt 0, with a fresh op_id — and
nothing of the served delivery’s family, budget, origin, or deadline rides
along. The hop budget belongs to the family, which is why a spent budget
stops a forward and never new work. A reader who “simplifies” this into
Causality::child_of would make a new task spend its sender’s hop and
hand the server an envelope that plugins/onlyne-agent-pi’s own
sendEnvelope never mints: two drives, one obligation, two shapes. A
note joins no family: no chain and no op_id.
Source§impl DispatchState
impl DispatchState
Sourcepub fn completion_refusal(
&self,
from: Option<&AdapterIo>,
report: &Report,
) -> Option<ResBody>
pub fn completion_refusal( &self, from: Option<&AdapterIo>, report: &Report, ) -> Option<ResBody>
The refusal one incoming completion meets, when a client-side constraint
refuses it; None when the frame may settle.
from is the connection the frame arrived on. A completion that arrives
with no connection — the local operator surface, and the fault a refused
focus files — is measured against the shape rule alone: the relay guard
reads the session’s own delivery record, and a door with no sender has no
session whose record it could read.
Sourcepub fn handoff_refusal(
&self,
io: &AdapterIo,
args: &HandoffArgs,
) -> Option<ResBody>
pub fn handoff_refusal( &self, io: &AdapterIo, args: &HandoffArgs, ) -> Option<ResBody>
The refusal one handoff frame meets when the family’s hop budget is
spent; None when the frame may be built.
A child sits one hop below the task that hands it on, and a family’s
hop_budget is the depth the chain may reach: the frame that would sit
over it is refused here and names the budget it would break, rather than
being minted for the server to sort out. A frame from a connection that
does not serve the task is left for plugin_handoff’s own authority
answer, so this check reveals nothing to a foreign connection.
Source§impl DispatchState
impl DispatchState
Sourcepub fn suspend_idle_sessions(&self, now: Instant) -> Vec<String>
pub fn suspend_idle_sessions(&self, now: Instant) -> Vec<String>
Release the process of every idle session whose scope bound has expired.
A suspension is this client’s own act and it is only worth doing for a session the runtime can bring back: the conversation lives in the runtime’s own store, so the client starts the same command again when the family’s next delivery arrives and the runtime resumes the session it was given. A runtime that declared no resume keeps its process instead — the “process alive, session alive” degradation the scope table names — and this sweep leaves it alone.
The closes are host round trips, so they are collected here and run once the dispatch lock is off them, the way the sweeps that retire sessions already do. Answers the sessions it released, which are the ones whose row and claim this caller has to publish.
Sourcepub fn is_suspended(&self, session_id: &str) -> bool
pub fn is_suspended(&self, session_id: &str) -> bool
Whether one session’s process is currently released.
Source§impl DispatchState
impl DispatchState
Sourcepub fn enqueue_outbound(&self, envelope: &Envelope) -> Result<String>
pub fn enqueue_outbound(&self, envelope: &Envelope) -> Result<String>
Queue an outbound envelope before its first write and answer its op_id.
The queue keys every row by an op_id, and the proto requires that key
only for the non-note kinds, so a note that arrives without one gets a
fresh client-minted id here: the row is keyed and what it replays is the
whole stamped envelope. A non-note keeps the id it brought, so a
re-delivered task still dedups on its original one.
Sourcepub fn accept_new(&self) -> Arc<AtomicBool> ⓘ
pub fn accept_new(&self) -> Arc<AtomicBool> ⓘ
The flag the runloop and the dispatcher share.
Sourcepub fn runtime_name(&self) -> String
pub fn runtime_name(&self) -> String
The name of the session runtime hosting this role’s sessions.
The registration beside the client socket carries it, and an external runtime’s plugin matches on it to find the client its own session belongs to, so the answer is the backend’s own name rather than a spelling invented here.
Sourcepub fn set_link_up(&self, up: bool)
pub fn set_link_up(&self, up: bool)
Record that the server link came up or went down.
Sourcepub fn cluster_ref(&self) -> String
pub fn cluster_ref(&self) -> String
Aggregate name this role supervises, empty for a plain role.
Sourcepub fn set_topology(&self, cluster: &str)
pub fn set_topology(&self, cluster: &str)
Record the server’s topology name, read from welcome.cluster.
Each spawned session carries it as ONLYNE_CLUSTER, which is how a host
backend addresses the tree it splits panes into. The runloop calls
this on every welcome, so a server that reloads under a new name is
followed by the sessions spawned after that point.
Sourcepub fn topology(&self) -> String
pub fn topology(&self) -> String
The topology name recorded from welcome, empty before the first welcome.
Sourcepub fn set_cluster_ref(&self, aggregate: impl Into<String>)
pub fn set_cluster_ref(&self, aggregate: impl Into<String>)
Record the aggregate name once, so every report keeps the same value across a reconnect.
Sourcepub async fn request(&self, op: ClientOp) -> Result<ResBody, NetError>
pub async fn request(&self, op: ClientOp) -> Result<ResBody, NetError>
Ask the server one question over the live link.
Err means the link is down, never a refusal: a refusal arrives as an
Ok body carrying ok: false, which is what the local CLI shows.
Sourcepub fn attach_outbox(&self, outbox: Arc<dyn Outbox>)
pub fn attach_outbox(&self, outbox: Arc<dyn Outbox>)
Install the live link as the outbound path.
Sourcepub fn detach_outbox(&self)
pub fn detach_outbox(&self)
Remove the outbound path; lifecycle frames then queue as intents.
Sourcepub fn enqueue_op(&self, op: &ClientOp) -> Result<String>
pub fn enqueue_op(&self, op: &ClientOp) -> Result<String>
Queue one client op in the durable intent table and answer its op_id.
Source§impl DispatchState
impl DispatchState
Sourcepub fn reclaim_exited_resources(&self) -> Vec<String>
pub fn reclaim_exited_resources(&self) -> Vec<String>
Retire tracked resources whose stored lifecycle has reached Exited.
The periodic readiness tick calls this after completed work becomes an idle slot. Task-free sessions with an attached transport stay bound to their host resource, and task-free sessions whose agent has left release it.
The session ids come back because the retirement wrote each one of those
rows and the server mirrors only what this client reports: the resource
close and, for a completed session, the agent’s exit both moved the row
this tick found, and a publish cannot run under this lock. The caller is
handed what to publish, the same answer DispatchState::retire_dropped_ghosts
gives the sweep above.
Sourcepub fn retire_dropped_ghosts(
&self,
now: Instant,
grace_secs: u64,
) -> Vec<Retired>
pub fn retire_dropped_ghosts( &self, now: Instant, grace_secs: u64, ) -> Vec<Retired>
Retire the sessions whose plugin connection dropped and never came back, or whose connection stayed up while they went quiet, and answer which ones left and why.
A connection that ends without a detach frame leaves its session tracked
so an agent that restarts inside [client] reconnect_grace_secs finds the
resource it was using. That promise has to expire: a process that is
simply gone would otherwise hold a slot, a projected idle row, and a live
host resource forever, and on a role with max_sessions = 1 it stops every
later delivery. The window answers for the agent itself, so a session still
bound to a task goes with it: the plugin connection that would have
reported the ending is the one that dropped. The agent-gone feed is what
says the process left — the session’s own tuple reaches Exited through
AgentPhase::Gone rather than through a task result — and the reason the
backend is handed is the one grace_close_reason reads off what the slot
still owes.
What the slot owed is settled too: the task a bound session was serving
ends failed here, because the agent that would have reported its ending
is the one that left. A task with no verdict stays open for the server to
re-offer and for open_tasks to keep reading, and no later caller exists
to write one.
A slot this client holds read-only is not this sweep’s to end, agent gone
or not: the session id it would feed is the task id, so the ghost’s death
would take the live session’s mirror and its delivery row down with it.
That retirement belongs to retire_revived, which runs when the
completion that answers the held connection merges.
The window has a second way to open, and it is the one a socket cannot
report: a plugin whose event loop is blocked keeps its connection and
stops beating, so no socket ends and no clock this sweep could read
before moved. What such a session leaves behind is a stamp going stale
while its task stays bound and unsettled, and that is the reading this
sweep takes now. It is the same window and the same verdict — one clock,
one retirement, no second threshold beside [client] reconnect_grace_secs and no fault row of the kind stall_report_secs
records and leaves behind.
The sessions’ own ids come back rather than a count, because a retirement
still owes the server the session’s own ending: it is the only writer left
for that task, and a mirror nobody tells keeps that session’s last reading —
working, for one that had beaten — until the server’s own observer records
a fault about it. The publish is
sync_session’s, which is the
report an ordinary ending travels on, and it cannot run under this lock —
so the caller is handed what to publish instead of a second writer being
invented here.
Source§impl DispatchState
impl DispatchState
Sourcepub fn plugin_handoff(&self, io: &AdapterIo, args: HandoffArgs) -> ResBody
pub fn plugin_handoff(&self, io: &AdapterIo, args: HandoffArgs) -> ResBody
Take one plugin handoff frame and answer what the plugin is told.
The frame names the task the session is handing on and the role it goes
to. The child is minted here, through the builder the report-driven path
uses, so the family id and the family’s figures ride along and the depth
grows by one hop. The envelope leaves on the queue the plugin send op
writes to.
The answer names the child:
{"task_id": "<uuid>", "hop": 3, "queued": true, "op_id": "<uuid>"}
(onlyne_proto::HandoffArgs).
A frame is answered only for the connection serving the task it names. An unknown task and a foreign connection earn the same code and the same field, and their messages say which of the two refused the frame.
Source§impl DispatchState
impl DispatchState
pub fn new( role: impl Into<String>, workspace: impl Into<PathBuf>, command: Vec<String>, max_sessions: u32, backend: Arc<dyn SessionBackend>, store: ClientStore, ) -> Self
Sourcepub fn with_session_policy(self, policy: SessionPolicy) -> Self
pub fn with_session_policy(self, policy: SessionPolicy) -> Self
Adopt the role’s [client.session] policy.
Sourcepub fn with_placement(self, placement: SessionPlacement) -> Self
pub fn with_placement(self, placement: SessionPlacement) -> Self
Install the placement this machine resolved. It is the other half of the drive rule, so it is recorded even when the pair it makes is refused.
Sourcepub fn drive(&self) -> Option<Drive>
pub fn drive(&self) -> Option<Drive>
The drive the role’s spec declares, None before the first welcome.
Sourcepub fn placement_name(&self) -> Option<&'static str>
pub fn placement_name(&self) -> Option<&'static str>
The placement this machine resolved, as the registration spells it.
Sourcepub fn session_backend(&self) -> &'static str
pub fn session_backend(&self) -> &'static str
The name of the session backend currently installed.
Sourcepub fn set_drive(&self, drive: Drive, refusal: Option<String>)
pub fn set_drive(&self, drive: Drive, refusal: Option<String>)
Record the drive the role’s spec declares, and the sentence a delivery meets when that drive cannot run under this machine’s placement.
The drive lands whether or not a backend could be installed for it: a pair this machine cannot host still has to be known, or the next delivery would run under the backend the previous drive left behind.
Sourcepub fn set_backend(&self, backend: Arc<dyn SessionBackend>) -> bool
pub fn set_backend(&self, backend: Arc<dyn SessionBackend>) -> bool
Install the backend a role’s drive selects on this machine.
Answers whether it landed. A drive cannot move while this role holds live
sessions: their panes, tabs, and children are the installed backend’s to
close and to probe, and a backend that never opened a resource cannot
answer for it. So the move waits — the caller keeps the old drive
recorded, and the next welcome or spec reload tries again once the
sessions are gone.
Sourcepub fn knows_session(&self, session_id: &str) -> bool
pub fn knows_session(&self, session_id: &str) -> bool
Whether this client holds a session this name answers to.
A mount names the session it was spawned for, and a plugin that outlived a restart still names the one it was serving when it redials — to a client that has no memory of it, because nothing here rebuilds slots from the store. So this is a question with a real “no”, and the caller needs to be able to ask it before it binds anything under that name.
Sourcepub fn session_policy(&self) -> SessionPolicy
pub fn session_policy(&self) -> SessionPolicy
The role’s [client.session] policy.
pub fn session_count(&self) -> usize
Sourcepub fn feed_turn_started(&self, task_id: &str)
pub fn feed_turn_started(&self, task_id: &str)
Feed the turn-started fact for a self-driven session’s turn, so the
row’s agent phase reads running before the agent can report a
completion through its tools mount (docs/v2-CONTRACT.md §3c). A plugin
drive feeds this through its heartbeats; a self-driven drive has none,
so the dispatch path feeds it where it hands the turn to the backend.
Sourcepub fn feed_turn_ended(&self, task_id: &str)
pub fn feed_turn_ended(&self, task_id: &str)
Feed the turn-ended fact for a self-driven session’s turn, so the row’s
agent phase reads idle when the turn the drive witnessed ends. The
never-ran guard reads the same column, so a completion that arrives
after this feed still passes it.
pub fn role(&self) -> String
pub fn command(&self) -> Vec<String>
pub fn backend_name(&self, task_id: &str) -> Option<String>
Sourcepub fn outcome_feed(&self) -> Option<OutcomeFeed>
pub fn outcome_feed(&self) -> Option<OutcomeFeed>
The terminal-fact stream of a backend that drives its own agent.
Sourcepub fn push_assign_ack(&self, ack: AssignAckArgs) -> bool
pub fn push_assign_ack(&self, ack: AssignAckArgs) -> bool
Queue a delivery ack when the plugin refuses an assignment.
An accepted assignment is not terminal for the server row: the normal completion path still settles that delivery. A refused assignment is a terminal local decision, so it uses the same durable ack queue as every other delivery settlement.
Sourcepub fn push_settled(&self, ack: AckArgs)
pub fn push_settled(&self, ack: AckArgs)
Queue an ack the client owes the server.
Record an ack the client owes the server.
The ack is durable: D11’s control plane is at-least-once, and a settled session whose ack is lost leaves the row in flight forever. The intent queue carries it across a link that is down, and the flusher is the sender.
Sourcepub fn owe_controlled_settle(
&self,
task_id: &str,
word: ControlWord,
now: Instant,
)
pub fn owe_controlled_settle( &self, task_id: &str, word: ControlWord, now: Instant, )
Note that one task’s plugin has been told by this client’s own control
command to report the ending of that task.
on_control runs this before the recycle frame leaves, which is the
point where the client still knows the order: the plugin’s completion and
the retirement that command triggers race over the session’s row, and a
guard that read the row would answer the same operator action two ways. One
note per task is kept, so a command issued twice waits for one answer.
word is what the operator said and now is when they said it, and both
are the caller’s to name rather than this call’s to invent: the command is
the authority on the ending it asked for, and the instant it reads is the
one the watchdog’s bound runs from.
Sourcepub fn take_controlled_settle(&self, task_id: &str) -> bool
pub fn take_controlled_settle(&self, task_id: &str) -> bool
Whether one completion answers a command noted above, consuming the note.
The note is spent whichever way the settle it authorises goes: a refused verdict leaves no second answer owed, and an applied one has travelled the command’s own completion. A later report for the same task is the plugin speaking for itself again, and reads the ordinary door.
Sourcepub fn control_settles_due(&self, now: Instant) -> Vec<ControlNote>
pub fn control_settles_due(&self, now: Instant) -> Vec<ControlNote>
The notes whose operator’s word has gone unanswered past the bound.
The reading a sweep takes before it acts, and it spends nothing: the
settle below goes through take_controlled_settle, so a completion that
answers a word between this read and that call takes the note first.
Sourcepub fn settle_unanswered_control(&self, note: &ControlNote) -> bool
pub fn settle_unanswered_control(&self, note: &ControlNote) -> bool
Settle the work one operator’s word left open, with no report behind it.
The word asks a plugin for its own ending and the completion that answers
it is a frame of the plugin’s, so a plugin that never sends one — it left
with the command’s frame, or implements no recycle at all — leaves the
task open and the delivery row this client was handed in flight, with no
later caller to answer either. This is that caller.
The writes are the ones retire_dropped_ghosts makes for the task its
owed session left: the verdict through the task’s own record, which
refuses to overwrite one that landed first, and the still-held delivery row
refused with the operator’s word, which is terminal for that row the way
every refusal is. The publish is the caller’s, because it cannot run under
this lock.
Answers false when this call is not the settle: the note is already spent
by a completion that answered the word, or the task’s record carries a
verdict already, and either way nothing here is written and nothing is for
the caller to publish.
Sourcepub fn role_slice(&self) -> RoleSlice
pub fn role_slice(&self) -> RoleSlice
The role slice the dispatcher currently runs.
Sourcepub fn live_task_ids(&self) -> HashSet<String>
pub fn live_task_ids(&self) -> HashSet<String>
Task ids currently occupying a live slot.
Sourcepub fn hello_live_sessions(&self) -> StoreResult<Vec<LiveSession>>
pub fn hello_live_sessions(&self) -> StoreResult<Vec<LiveSession>>
The sorted live sessions one hello claims: the union of the memory slots and the DB-persisted active sessions, so a process crash does not lose the claim.
A store failure is not swallowed: it is logged with the memory claim
still derivable from the slots, and returned so the caller can degrade
to [live_claim_from_slots] deliberately instead of answering as if the
durable half were simply empty. Silently dropping that half lets the
server requeue every in_flight row a crash left behind, which is the
duplicate delivery this claim exists to prevent.
Sourcepub fn live_claim_from_slots(&self) -> Vec<LiveSession>
pub fn live_claim_from_slots(&self) -> Vec<LiveSession>
The memory half of [hello_live_sessions]: the claim to dial with when
the durable store cannot answer. The slots are what this process is
serving right now, so even a degraded hello keeps those rows in_flight.
A slot answers with the session’s own id — the one that stays put while a scope hands the session delivery after delivery — the delivery it is serving now, and whether its process has been released. A session between deliveries claims no delivery: the server keeps the rows those sessions are still the owners of, and invents none.
Sourcepub fn note_stall_assigned(&self, task_id: &str, now: Instant)
pub fn note_stall_assigned(&self, task_id: &str, now: Instant)
Start the stall clock for a newly assigned task.
Sourcepub fn note_stall_applied(&self, task_id: &str, now: Instant)
pub fn note_stall_applied(&self, task_id: &str, now: Instant)
Refresh the stall clock after an Applied persist.
Sourcepub fn stall_due(&self, now: Instant, threshold_secs: u64) -> Vec<String>
pub fn stall_due(&self, now: Instant, threshold_secs: u64) -> Vec<String>
Task ids whose freeze exceeds threshold_secs in this episode.
Exited projections retire their remaining progress clocks.
Sourcepub fn mark_stalled(&self, task_id: &str)
pub fn mark_stalled(&self, task_id: &str)
Remember that this freeze episode has been reported.
Sourcepub fn stall_report(&self, task_id: &str) -> Option<Report>
pub fn stall_report(&self, task_id: &str) -> Option<Report>
Observation-only stall fault for an active task, carrying the stored
watermark. A tuple and verdict that derive exited retire their progress
clock before the send boundary.
Sourcepub fn has_mounted_adapter(&self) -> bool
pub fn has_mounted_adapter(&self) -> bool
Whether any adapter is currently mounted (named or parked).
Sourcepub fn reconfigure(&self, slice: RoleSlice)
pub fn reconfigure(&self, slice: RoleSlice)
Adopt a role slice: the one welcome carried, or the one a reload’s
role row carries.
Sourcepub fn role_prose(&self) -> String
pub fn role_prose(&self) -> String
Role prose last cached from welcome.
Sourcepub fn holds_task(&self, task_id: &str) -> bool
pub fn holds_task(&self, task_id: &str) -> bool
Whether one task names a session this client holds, in memory or in its durable session rows.
Sourcepub fn session_generation(&self, task_id: &str) -> Option<u64>
pub fn session_generation(&self, task_id: &str) -> Option<u64>
Generation the reducer holds for one task, before any hand-off.
Sourcepub fn has_capacity(&self) -> bool
pub fn has_capacity(&self) -> bool
Whether a delivery has somewhere to run.
§5’s max_sessions caps concurrency, so a delivery that arrives at the
cap waits on the server: the row stays in flight and the next pull
offers it again once a session frees. Each task runs in its own session,
so a slot whose task has finished still spends capacity until it retires.
A scope that hands a delivery on to a session the role already holds
spends no slot at all, so an task or role session sitting idle — or
suspended, its slot already given back — is room even at the cap. The
placement decides which delivery goes where; this only answers whether
there is somewhere for one to go. oneshot reuses nothing, which is the
count it has always been.
Sourcepub fn task_completed_here(&self, task_id: &str) -> bool
pub fn task_completed_here(&self, task_id: &str) -> bool
Whether this role already finished one task with a terminal Done.
The task’s own record is the account: the settle writes the verdict the
agent filed and refuses to overwrite it, so a record reading done means
this role answered for this task id once already. A redelivery of that
task is not new work — running it again would stage its payload on
whichever session happens to be idle, so one chain’s task executes inside
another conversation and the second answer collides with the verdict the
first one settled.
Only done counts. A session killed or crashed mid-flight leaves its task
open, or settles it failed, and the server’s requeue, repair_retry, and
control retry all re-offer that task on purpose, so those deliveries
still run.
Sourcepub fn staged_hosting_session(&self) -> Option<String>
pub fn staged_hosting_session(&self) -> Option<String>
One staged session this role’s standing runtime can be offered.
The mirror of Self::staged_without_transport, narrowed to a session
that has no transport because a hosting runtime will supply one. A
session another connection already serves is not offered, and a session
whose work is already in flight is not either — a runtime must never be
asked to open a second conversation for a chain that has one.
Sourcepub fn hosting_open_args(&self, session_id: &str) -> OpenArgs
pub fn hosting_open_args(&self, session_id: &str) -> OpenArgs
What to ask a hosting runtime for, by the name this client gave the session.
The resume handle of the family’s previous session rides along when there
is one, so a task-scoped runtime hands back the conversation it already
holds rather than opening a second one for the same chain. Nothing reads
the handle here: it is the runtime’s own word for where its conversation
is, and this client stores it and gives it back.
Sourcepub fn staged_without_transport(&self) -> Option<String>
pub fn staged_without_transport(&self) -> Option<String>
The task of one session that holds a payload with no connection bound.
A work item that arrives before its always-running agent mounts waits in exactly this state, and the mount ends the wait. A session answers through the transport its first task claimed, so it stays served.
Sourcepub fn tools_token(&self, session_id: &str) -> Option<String>
pub fn tools_token(&self, session_id: &str) -> Option<String>
The tools token of one live session, when this client holds it.
The one door the token leaves this process through: the session’s own
drive reads it while building the child that mounts onlyne mcp
(SpawnSpec.tools_token), and nothing else may hand it out
(docs/v2-CONTRACT.md §3b).
Sourcepub fn tools_mount_for(&self, token: &str) -> Option<ToolsSession>
pub fn tools_mount_for(&self, token: &str) -> Option<ToolsSession>
The session a tools mount token speaks for, when it names a live one.
Sourcepub fn bind_tools_mount(
&self,
token: &str,
io: AdapterIo,
) -> Option<ToolsSession>
pub fn bind_tools_mount( &self, token: &str, io: AdapterIo, ) -> Option<ToolsSession>
Bind one tools connection to the session its token names.
The hello validated the token; this writes the binding the session’s
later frames are measured against, and answers the session the
connection now speaks for. A second connection presenting the same token
takes the binding over, and the earlier one then fails the per-frame
liveness check — a token belongs to one session, and the newest
connection is the one that speaks for it. None is a token whose
session retired between the handshake and this line.
Sourcepub fn tools_connection_live(&self, io: &AdapterIo) -> bool
pub fn tools_connection_live(&self, io: &AdapterIo) -> bool
Whether one live tools connection still speaks for a session this client holds.
A token dies with its session, and a mount whose session has ended is
refused from then on: this is the per-frame half of that rule, so a
connection left open past its session’s retirement answers unauthorized
and closes rather than speaking for work nobody holds.
Sourcepub fn tools_scope(&self, io: &AdapterIo) -> Option<ToolsScope>
pub fn tools_scope(&self, io: &AdapterIo) -> Option<ToolsScope>
What one live tools connection speaks for, read under the lock that owns
it; None when the connection is unbound or its session has stopped
serving.
The token is the binding (docs/v2-CONTRACT.md §3b), so everything a
tools frame is stamped with comes from the session’s own slot: the open
delivery its report and handoff frames name, and the session a send
leaves from. Reading it in one locked pass is what keeps a frame from
being stamped from a state that moved between the lookup and the stamp.
Sourcepub fn release_tools_connection(&self, io: &AdapterIo)
pub fn release_tools_connection(&self, io: &AdapterIo)
Drop the binding one tools connection held, whichever session it named.
Source§impl DispatchState
impl DispatchState
Sourcepub fn hold_frame(&self, io: &AdapterIo) -> FrameGuard<'_>
pub fn hold_frame(&self, io: &AdapterIo) -> FrameGuard<'_>
Hold io for as long as one of its inbound frames is being handled.
Sourcepub fn bind_adapter(
&self,
session_id: &str,
io: AdapterIo,
capabilities: Vec<Capability>,
)
pub fn bind_adapter( &self, session_id: &str, io: AdapterIo, capabilities: Vec<Capability>, )
Bind one adapter connection to the session it named.
The name is the session id the client spawned the plugin with, which is enough on its own: a plugin that mounts before the client staged its session is remembered here and takes the payload the moment it is staged, and a plugin that mounts after finds its session waiting.
Sourcepub fn attach_msg_id(&self, task_id: &str, msg_id: &str)
pub fn attach_msg_id(&self, task_id: &str, msg_id: &str)
Remember the delivery handle for one task.
The handle goes to the session serving the task, not to a read-only slot that came back for it, so the ack this earns answers the live delivery.
Sourcepub fn plugin_send(&self, io: &AdapterIo, envelope: &Envelope) -> Result<Value>
pub fn plugin_send(&self, io: &AdapterIo, envelope: &Envelope) -> Result<Value>
Take one plugin send frame and answer what the plugin is told.
One path, whichever connection sent the frame. A handoff is a real delivery to the role it names, and this client is the only process holding a link that could carry it, so a connection that came back for a task another session now serves routes here exactly like the serving one. Nothing is held for a later merge (§3c): a session’s handoffs travel as its own frames, answered where it sends them, so a settlement is only the completion’s half of the account.
Sourcepub fn park_transport(&self, io: AdapterIo, capabilities: Vec<Capability>)
pub fn park_transport(&self, io: AdapterIo, capabilities: Vec<Capability>)
Park one plugin connection as this role’s waiting agent.
Only a mount that names no session parks: it is a plugin that attached before any work existed, so it takes the next session this role stages (plan §6 line 285). A role can host more than one such agent, and each that arrives joins the back of the queue: the record is a queue and not one slot, because overwriting it dropped the connection that was already waiting with no accounting of any kind — no release, no log, and no word to the plugin, which is how a quiet role lost a worker.
A connection already in the queue is refreshed where it stands rather than sent to the back: a mount that names nothing twice over is the same agent re-helloing, and its place in line is that agent’s due.
Sourcepub fn stand_transport(&self, io: AdapterIo, capabilities: Vec<Capability>)
pub fn stand_transport(&self, io: AdapterIo, capabilities: Vec<Capability>)
Hold io as this role’s standing transport.
A hosting runtime’s connection is not a worker waiting for one job, so it
does not join Self::park_transport’s queue: a claim takes the oldest
entry and spends it, and a role whose hosting runtime serves four
sessions would have nothing left for the other three. It is held here
instead, where Self::staged_hosting_session finds it: the connection
is asked for each session the role opens, and one connection serves as
many as the role needs.
A connection already standing is refreshed where it stands, the rule the park uses too: a mount that says this twice is one runtime re-helloing.
Sourcepub fn session_transport(
&self,
session_id: &str,
) -> Option<(AdapterIo, Vec<Capability>)>
pub fn session_transport( &self, session_id: &str, ) -> Option<(AdapterIo, Vec<Capability>)>
The connection that serves one session, when its plugin is attached.
A plugin names the session it was spawned for, and the slot’s key is the other spelling worth trying.
Sourcepub async fn recycle_plugin(
&self,
task_id: &str,
reason: &str,
outcome: Option<Outcome>,
)
pub async fn recycle_plugin( &self, task_id: &str, reason: &str, outcome: Option<Outcome>, )
Tell the plugin serving one session to tear itself down, when that plugin
implements recycle. A plugin without the capability is skipped: the
caller’s backend close stops the process either way.
Sourcepub async fn probe_plugin(&self, task_id: &str) -> bool
pub async fn probe_plugin(&self, task_id: &str) -> bool
Ask the plugin serving one session for a fresh observation.
The plugin answers with a heartbeat report, which is the reducer’s
evidence and the projection the operator reads. Answers whether a probe
frame actually went out, and false says nothing was asked: a session no
connection serves has no plugin to put the question to, and a notify that
failed left the frame in this process. The caller must not record either
as a probe that landed, because the answer a live plugin would have given
never existed — the projection that says a plugin is gone is written when
the session ends, never by this call.
Sourcepub async fn nudge_plugin(&self, task_id: &str) -> bool
pub async fn nudge_plugin(&self, task_id: &str) -> bool
Hand one session’s own sentence to its agent through the plugin’s input.
The turn-end rule has one owner, this client, so the client is what tells
a plugin-driven session that its turn ended without a completion
(docs/v2-CONTRACT.md §3c). The plugin injects [NUDGE_TEXT] through the
same channel an assignment’s text takes and keeps no copy of it; the
frame carries no envelope, no prose, and no attachments, and it is not a
delivery: the task stays open and nothing in the plugin’s turn
bookkeeping is reset by it.
Answers whether the frame went out. false is the honest answer for a
session with no plugin to ask and for a plugin that declared no inject:
a drive that cannot be nudged must not be told it was, and the caller
settles the delivery at that turn end instead.
Sourcepub fn release_connection(
&self,
session_id: Option<&str>,
io: &AdapterIo,
graceful_detach: bool,
) -> Vec<String>
pub fn release_connection( &self, session_id: Option<&str>, io: &AdapterIo, graceful_detach: bool, ) -> Vec<String>
Release the bindings served by one plugin connection.
A graceful detach retires each idle session because the agent that owned it has left. An attached transport preserves the idle resource because the same agent is still reachable. A connection ending through another path preserves the slot and resource for an agent reconnection and starts the reconnect clock on it, which is what bounds how long a session waits for an agent that is never coming back. Every released binding retires its task progress clock. A slot carrying work remains under lifecycle ownership, and it carries that same clock: a goodbye and a silent drop both leave no heartbeat coming for the task it owes, and the window is what ends a session whose agent never returns.
The session ids a goodbye retired come back, because that retirement wrote
their rows: the resource close and, for a completed session, the agent’s
exit. The server mirrors only what this client reports, and a publish
cannot run under this lock, so the caller is handed what to publish — the
same answer DispatchState::retire_dropped_ghosts gives its sweep.