onlyne_client/session/dispatch/settle.rs
1use super::*;
2
3use super::outbound::{store_ack, transport_envelope};
4use super::projection::{note_verdict, phase, sync_session};
5use super::retire::{PendingClose, close_retired, keeps_idle, release_locked, retire_idle_locked};
6use super::state::{
7 DispatchInner, DispatchState, has_attached_transport, slot_key_named, slot_key_serving_task,
8 slot_task,
9};
10use super::transport::{names_session, serves_session};
11use onlyne_proto::ErrorCode;
12use onlyne_proto::adapter::HandoffArgs;
13
14/// Fault kind for a completion this client refused for want of a turn. The word
15/// is what `onlyne faults` and `onlyne-client status` carry, so it names the
16/// reading that refused the frame.
17pub const SETTLE_WITHOUT_TURN: &str = "settle_without_turn";
18
19/// Who asked for one settle, which is what decides whether the never-ran guard
20/// reads it.
21///
22/// The guard exists for the door where the claim and the claimant are the same
23/// party: a plugin reports its own ending, and a session whose agent never ran
24/// can report one too. The other two doors carry the client's own act, so their
25/// evidence is already in this process — a self-driven backend that watched its
26/// agent end, or an operator's `control` command that asked for the ending.
27#[derive(Clone, Copy, Debug, PartialEq, Eq)]
28pub enum SettleAuthority {
29 /// A `complete` report from a plugin connection. Guarded.
30 PluginReport,
31 /// A terminal fact this client reached through its own eyes: the ending a
32 /// self-driven backend reported through `outcome_loop`. Unguarded: this
33 /// client is the witness of the work it watched.
34 ClientOwned,
35 /// The answer to a `control recycle` or `control cancel` this client issued
36 /// for the task. Unguarded: the operator asked for this ending, and the
37 /// plugin's report and the retirement this command runs race over the row,
38 /// so the row's phase at the moment the frame lands decides nothing.
39 ControlDriven,
40}
41
42/// Whether one session's own row carries a turn, and the agent-phase word it
43/// holds for the operator's reading.
44///
45/// The row is this client's record: a beat's `agent` dimension reaches it
46/// through the `(generation, seq)` gate in `apply_persist`, and the ready
47/// barrier's `feed_ready` writes `Ready` into the same column. `Running` and
48/// `Idle` are the two phases a turn puts the tuple through
49/// (`crates/onlyne-session/src/lifecycle/state.rs:11`), so one of them is the
50/// answer. `Booting`, `Ready` and `Gone` are each a session that has run nothing
51/// this client can point at: `Gone` is written by `AgentGone` and
52/// `ResourceClosed` from any live phase, which leaves the death of an agent that
53/// never started reading exactly like the death of one that worked. A task this
54/// client holds no row for answers the same way, with its own word in the reason.
55///
56/// One read answers both halves, and both come off the row's own column: the
57/// word the operator is handed is the word `projection_of` publishes, so the
58/// refusal cannot name a phase other than the reading that caused it. The tuple
59/// inside `observed_json` is a second source for the same dimension, and a row
60/// whose bytes are unparsable rebuilds to `Booting` beside a column still
61/// reading `running` — a guard that decided on one and reported the other, and
62/// a read-only door that wrote an alert and a ledger event while answering.
63fn turn_recorded(inner: &DispatchInner, task_id: &str) -> (bool, String) {
64 let Ok(Some(row)) = inner.store.get_session(task_id) else {
65 return (false, "no session row".to_string());
66 };
67 let agent = phase(&row.agent_state, AgentPhase::Booting);
68 (
69 matches!(agent, AgentPhase::Running | AgentPhase::Idle),
70 row.agent_state,
71 )
72}
73
74/// Settle one finished task: publish the verdict, retire what answered for it,
75/// and let the receipt leave.
76///
77/// No handoff is routed here. A session's handoffs travel as its own `handoff`
78/// frames (`handoff.rs`), answered where the session sends them, so a
79/// settlement is only the completion's half of the account: the verdict, the
80/// receipt, and the resources this task was holding.
81///
82/// A replayed delivery for work whose first verdict already settled is ordinary
83/// at-least-once traffic. The first verdict stands. The replay returns the task
84/// binding. The session that took the replay lets go of its slot.
85/// The first receipt, `out_head`, and handoff relay remain attached to that
86/// first settlement.
87///
88/// A plugin report with no turn behind it is refused whole the same way, ahead of
89/// every write this call makes: the drain opens over work that never ran, so the
90/// completion intent, the verdict, the `out_head` line and the delivery ack all
91/// stay unwritten and the row is left for the server's requeue. The frame is
92/// answered as an applied one — a plugin treats a failed report as a link that
93/// died, and sends the same terminal fact again — and the fault the refusal
94/// leaves behind names the reading that refused.
95pub async fn on_out(
96 state: &DispatchState,
97 task_id: &str,
98 outcome: Outcome,
99 head: Option<String>,
100 details: Option<String>,
101 asked: SettleAuthority,
102) -> Result<()> {
103 if asked == SettleAuthority::PluginReport {
104 let inner = state.inner.lock();
105 let (turn, phase) = turn_recorded(&inner, task_id);
106 if !turn {
107 let reason = format!(
108 "no turn ran: the agent phase this client holds for the session reads {phase}"
109 );
110 crate::reconcile::record_fault(
111 &inner.store,
112 task_id,
113 SETTLE_WITHOUT_TURN,
114 "client",
115 &reason,
116 )?;
117 tracing::warn!(
118 task = %task_id,
119 ?outcome,
120 phase = %phase,
121 "a completion arrived for a session that never ran a turn; the task stays open"
122 );
123 return Ok(());
124 }
125 }
126 let (settled, session_id) = {
127 let mut inner = state.inner.lock();
128 if take_verdict(&inner, task_id, outcome)? {
129 inner
130 .store
131 .put_out_head(task_id, head.as_deref().unwrap_or(""))?;
132 // The handle and the chain this answer travels on belong to the session
133 // serving the task, not to a read-only one that came back for it.
134 let key = slot_key_serving_task(&inner, task_id);
135 let slot = key.as_deref().and_then(|key| inner.sessions.get_mut(key));
136 let origin = slot.as_ref().and_then(|slot| slot.origin.clone());
137 let causality = slot.as_ref().map(|slot| slot.causality.clone());
138 let msg_id = slot.and_then(|slot| slot.msg_id.take());
139 // The publish that follows names the session, not the delivery: a
140 // scoped session outlives this delivery and its row has to read as a
141 // live session serving nothing rather than as a session that exited.
142 let session_id = key.unwrap_or_else(|| task_id.to_string());
143 if let Some(msg_id) = msg_id {
144 store_ack(
145 &inner,
146 AckArgs {
147 msg_id,
148 op_id: None,
149 accepted: true,
150 reason: None,
151 },
152 );
153 }
154 // A settled session gives its capacity back, so a role at
155 // `max_sessions` takes the next row instead of holding finished slots.
156 release_locked(&mut inner, task_id, None)?;
157 (
158 Some((completion_envelope(
159 &inner.role,
160 origin,
161 task_id,
162 head.as_deref(),
163 details.as_deref(),
164 causality.as_ref(),
165 ),)),
166 session_id,
167 )
168 } else {
169 tracing::warn!(
170 task = %task_id,
171 ?outcome,
172 "a second verdict arrived for a settled task; the first one stands"
173 );
174 // The replayed session still owns the task binding until this
175 // release, so the standing verdict travels with the client's own
176 // post-release tuple and capacity returns to the role.
177 let session_id =
178 slot_key_serving_task(&inner, task_id).unwrap_or_else(|| task_id.to_string());
179 release_locked(&mut inner, task_id, None)?;
180 (None, session_id)
181 }
182 };
183 // The refused branch carries no receipt: the task account remains the first
184 // verdict, and the client's own row is published now that the replay session
185 // has returned its binding and completed its retirement.
186 let Some((receipt,)) = settled else {
187 return sync_session(state, &session_id).await;
188 };
189 // A connection that came back for this task has now had its ending answered,
190 // so it is retired here. The settled account above is the whole settlement:
191 // nothing here settles or releases this task a second time.
192 retire_revived(state, task_id).await;
193 // The leave request a settled delivery owes its runtime.
194 request_runtime_exit(state, &session_id).await;
195 // The terminal receipt leaves as its own envelope, so the origin — a role
196 // or a gateway conversation — learns the outcome (plan §3 `Completion`).
197 // It rides the intent queue, which is what makes a completion survive the
198 // disconnect rules of §6 line 289.
199 if let Some(envelope) = receipt {
200 transport_envelope(state, &envelope).await?;
201 }
202 sync_session(state, &session_id).await
203}
204
205/// Ask the runtime behind a settled delivery to leave, when its scope does not
206/// keep the session.
207///
208/// `complete` is the plugin's own exit handover, and a plugin that takes it
209/// leaves by itself. The other endings have no such handover — a failed report,
210/// and a delivery this client settles `blocked` at its turn end — so without
211/// this request the process behind a non-kept session stays reachable forever:
212/// the retirement keeps a session with an attached transport (rightly, that
213/// exemption is how a pool member lives), so nothing else would ever ask it to
214/// go, and the pane stays open beside a row that already reads `exited`.
215///
216/// Only a slot the scope does not keep and that serves no delivery is asked. A
217/// `task`/`role` member between deliveries is not a leftover, and a session no
218/// connection serves is the reconnect or silence sweep's to end.
219async fn request_runtime_exit(state: &DispatchState, session_id: &str) {
220 let due = {
221 let inner = state.inner.lock();
222 match slot_key_named(&inner, session_id) {
223 Some(key) => inner.sessions.get(&key).is_some_and(|slot| {
224 !keeps_idle(&inner, &key)
225 && slot.task_id.is_none()
226 && has_attached_transport(&inner, &key, slot)
227 }),
228 None => false,
229 }
230 };
231 if due {
232 state
233 .recycle_plugin(session_id, "delivery ended", None)
234 .await;
235 }
236}
237
238/// Drain one session's completion and file its task's verdict, answered with
239/// whether this report is the first verdict the task took.
240///
241/// The drain runs first and the verdict lands behind it, because the order is the
242/// row's own need: `settle` closes the completion intent, and a settled task beside
243/// a delivery that never drained is the pair `project` cannot read as `exited`. A
244/// report arriving behind a standing verdict therefore still drains the session
245/// that sent it — the retried row whose task the grace sweep had already answered,
246/// and the replayed session the caller's release retires — and only its own verdict
247/// is refused. Which verdict the task keeps is `settle_task`'s answer either way:
248/// it writes where `settled_at IS NULL` and refuses to move one that is stamped.
249///
250/// A drain with no subject is the one thing this order cannot carry. A task this
251/// client holds no row for has no intent to close, and `settle` answers the attempt
252/// with `unknown session`, which used to fail the whole report ahead of the verdict
253/// its record is still entitled to take — the shape a `client.db` replaced under a
254/// live role, or a foreign task reported into one, arrives in. The row is read for
255/// that alone, and the verdict is filed whatever the reading says.
256fn take_verdict(inner: &DispatchInner, task_id: &str, outcome: Outcome) -> Result<bool> {
257 let drain = inner
258 .store
259 .get_session(task_id)?
260 .map(|_| settle(&inner.bridge, &inner.store, task_id))
261 .transpose()?;
262 let first = inner.store.settle_task(task_id, task_state_of(outcome))?;
263 if let Some(verdict) = drain {
264 note_verdict(&verdict, task_id);
265 }
266 Ok(first)
267}
268
269/// Retire the read-only connections and slots a settled task has just answered.
270///
271/// A connection that came back for a session another connection serves is
272/// dropped from that session's record and its agent is told to leave, because
273/// the ending it reported is the whole of what it still had to say. A slot that
274/// lost its task to a newer session has the transport naming it dropped, its
275/// task binding released,
276/// and is retired as `Replaced`: the resource its agent was holding is this
277/// client's to close, and the newer session answers for the task. The account for
278/// the task is the settlement above.
279///
280/// A connection inside its own inbound frame is left alone. `adapter_socket`
281/// awaits the handler before it answers the frame, so a bye written here would
282/// leave ahead of that connection's own response, and the plugin's bye handler
283/// drops the socket and rejects every request awaiting an answer — a completion
284/// the ledger already holds would reach the agent as a failure it retries. The
285/// entry stays in `revived`: the connection's `detach` frame or its socket end
286/// retires it through `release_connection`, and the plugin that just completed
287/// ends its own session either way.
288///
289/// The name a connection mounted with is how this sweep judges which task that
290/// connection came back for, and the connection is what it removes. Two held
291/// connections can carry one name — `record_revived_connection` dedups per
292/// connection, and an agent that redials twice while another serves its session
293/// is held twice — and the one inside its own frame is left in the buffer by the
294/// rule above. Dropping held entries by name would then take that connection's
295/// entry with it: no bye reached it, nothing promoted it, and `release_connection`
296/// could no longer find the socket it still holds, which is a held connection
297/// neither silenced nor served.
298///
299/// A name that resolves to no slot is the other half of the judgement, and the
300/// answer there is the name itself: `dispatch` mints a session id from the task,
301/// so a held connection whose slot has already retired is one that came back for
302/// this task and for no other.
303async fn retire_revived(state: &DispatchState, task_id: &str) {
304 let (leaving, pending) = {
305 let mut inner = state.inner.lock();
306 let mut leaving: Vec<AdapterIo> = Vec::new();
307 let mut pending: Vec<PendingClose> = Vec::new();
308 for (session_id, io, _) in inner.revived.iter() {
309 if inner.in_frame.iter().any(|busy| busy.same_connection(io)) {
310 continue;
311 }
312 let reaches = slot_key_named(&inner, session_id)
313 .and_then(|key| {
314 inner
315 .sessions
316 .get(&key)
317 .map(|slot| slot_task(slot) == task_id)
318 })
319 .unwrap_or_else(|| session_id == task_id);
320 if reaches {
321 leaving.push(io.clone());
322 }
323 }
324 inner
325 .revived
326 .retain(|(_, revived, _)| !leaving.iter().any(|io| io.same_connection(revived)));
327 let silenced: Vec<String> = inner
328 .sessions
329 .iter()
330 .filter(|(_, slot)| slot.read_only && slot_task(slot) == task_id)
331 .map(|(key, _)| key.clone())
332 .collect();
333 for key in silenced {
334 let Some(slot) = inner.sessions.get(&key).cloned() else {
335 continue;
336 };
337 if slot.payload.is_some() {
338 tracing::warn!(
339 session = %key,
340 task = %task_id,
341 "a read-only session retires with a payload it was never handed"
342 );
343 }
344 let served: Vec<String> = inner
345 .transports
346 .keys()
347 .filter(|served| names_session(&key, &slot, served))
348 .cloned()
349 .collect();
350 for session_id in served {
351 inner.transports.remove(&session_id);
352 }
353 if let Some(current) = inner.sessions.get_mut(&key) {
354 current.task_id = None;
355 current.ready = false;
356 current.read_only = false;
357 current.dropped_at = None;
358 }
359 retire_idle_locked(
360 &mut inner,
361 &key,
362 crate::backend::CloseReason::Replaced,
363 &mut pending,
364 );
365 }
366 (leaving, pending)
367 };
368 // The replaced sessions' resources are this client's to give back, and the
369 // hosts take their time about it: the sweep already wrote every row and took
370 // every slot, so the close runs off the lock, ahead of the byes below that
371 // await the network on each held connection.
372 close_retired(pending);
373 for io in leaving {
374 let notice = AdapterMsg::Host(HostOp::Bye(onlyne_proto::ByeNotice {
375 reason: "the session that took this task answered for yours".into(),
376 }));
377 if let Err(error) = io.notify(notice).await {
378 tracing::debug!(error = %error, "the read-only connection had already left");
379 }
380 }
381}
382
383/// The receipt for one finished task, or `None` when its sender is unknown.
384///
385/// Every settled task answers its sender, the role that sent the task included:
386/// §3's `Completion` is the durable record that the work ended, and a role
387/// reading its own receipt ack is what settles the row.
388///
389/// The receipt names the task it answers and carries that task's own family
390/// figures — the family id, the hop budget, the origin, the deadline, and the
391/// labels — so a run's tasks and its completions print the same arc in
392/// `onlyne ledger`. It sits at the depth of the task it answers, and it is no
393/// link in the chain: it names no parent and replies to nothing. A task whose
394/// slot the client no longer holds, which is a row an older build opened, keeps
395/// the shape of a bare receipt.
396fn completion_envelope(
397 role: &str,
398 origin: Option<Principal>,
399 task_id: &str,
400 head: Option<&str>,
401 details: Option<&str>,
402 causality: Option<&Causality>,
403) -> Option<Envelope> {
404 let origin = origin?;
405 // The body carries the full result when the report named one, falling back
406 // to the one-line summary, and an empty string when neither is present — a
407 // turn that left no result still ends its task, and the sender still gets
408 // its answer as `text: Some("")`, which the validator accepts. The summary
409 // rides alongside in `head`, so a store that keeps a one-line preview
410 // shows it rather than the first clusters of the result.
411 let body = Body {
412 text: Some(details.or(head).unwrap_or_default().to_string()),
413 head: head.map(str::to_string),
414 image: None,
415 };
416 let mut causality = causality.cloned().unwrap_or_default();
417 causality.task = task_id.to_string();
418 causality.parent_task = None;
419 causality.reply_to = None;
420 // A receipt is written here, so it carries no redelivery count of its own.
421 causality.attempt = 0;
422 // `new_envelope` validates every protocol rule on the way out, so a receipt
423 // that cannot be addressed to its sender is the only one that goes unsent.
424 new_envelope(
425 MsgKind::Completion,
426 Principal::role(role),
427 origin,
428 body,
429 Some(causality),
430 )
431 .ok()
432}
433
434impl DispatchState {
435 /// Take one plugin `handoff` frame and answer what the plugin is told.
436 ///
437 /// The frame names the task the session is handing on and the role it goes
438 /// to. The child is minted here, through the builder the report-driven path
439 /// uses, so the family id and the family's figures ride along and the depth
440 /// grows by one hop. The envelope leaves on the queue the plugin `send` op
441 /// writes to.
442 ///
443 /// The answer names the child:
444 /// `{"task_id": "<uuid>", "hop": 3, "queued": true, "op_id": "<uuid>"}`
445 /// (`onlyne_proto::HandoffArgs`).
446 ///
447 /// A frame is answered only for the connection serving the task it names.
448 /// An unknown task and a foreign connection earn the same code and the same
449 /// field, and their messages say which of the two refused the frame.
450 pub fn plugin_handoff(&self, io: &AdapterIo, args: HandoffArgs) -> ResBody {
451 let (role, parent) = {
452 let inner = self.inner.lock();
453 let found = slot_key_serving_task(&inner, &args.task_id).and_then(|key| {
454 inner
455 .sessions
456 .get(&key)
457 .map(|slot| (key, slot.causality.clone()))
458 });
459 let Some((key, parent)) = found else {
460 return ResBody::err(
461 ErrorCode::Invalid,
462 format!("no session serves task {}", args.task_id),
463 Some("task_id".into()),
464 );
465 };
466 if !serves_session(&inner, &key, io) {
467 return ResBody::err(
468 ErrorCode::Invalid,
469 format!("this connection does not serve task {}", args.task_id),
470 Some("task_id".into()),
471 );
472 }
473 (inner.role.clone(), parent)
474 };
475 let (envelope, child) =
476 match handoff::relay(&role, &parent, &args.to, &args.text, args.image) {
477 Ok(built) => built,
478 Err(message) => return ResBody::err(ErrorCode::Invalid, message, None),
479 };
480 let queued = match self.plugin_send(io, &envelope) {
481 Ok(queued) => queued,
482 Err(error) => return ResBody::err(ErrorCode::Internal, error.to_string(), None),
483 };
484 // §4's durable class for this op: which handoff left this client, for
485 // whom, and on which hop of the chain. The enqueue above already
486 // succeeded, so the event cannot describe a handoff that never was; a
487 // failure here is the intent table's own, and the relay has left. The
488 // op goes to the server's stream, which is the fact's one owner; the
489 // durable queue is what makes the handoff event survive a crash the way
490 // the client's own row used to (`docs/v2-CONTRACT.md` §"Slice 7").
491 let op =
492 super::turn_end::record_handoff(self, &args.task_id, &args.to, child.hop, &args.text);
493 if let Err(error) = self.enqueue_op(&op) {
494 tracing::warn!(
495 task = %args.task_id,
496 to = %args.to,
497 error = %error,
498 "the handoff event was not queued"
499 );
500 }
501 // The queue path answers with the frame's `op_id`: a connection this
502 // client holds read-only serves no session, and the capability check
503 // inside `plugin_send` refuses that connection before this line.
504 ResBody::ok(serde_json::json!({
505 "task_id": child.task,
506 "hop": child.hop,
507 "queued": true,
508 "op_id": queued["op_id"],
509 }))
510 }
511}