Skip to main content

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}