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, 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}