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