Skip to main content

onlyne_client/runtime/runloop/
sessions.rs

1use super::config::{OUTCOME_POLL_MS, RunState};
2use super::run::settle_control;
3use crate::backend::SessionOutcome;
4use crate::session::accept::AcceptPath;
5use crate::session::dispatch::{self, ClientLink};
6use anyhow::{Result, anyhow};
7use onlyne_proto::{AckArgs, ClientOp, Delivery, QueryRolesArgs, RoleInfo};
8use std::sync::atomic::Ordering;
9use std::time::{Duration, Instant};
10use tokio::time::sleep;
11
12/// Drain terminal facts emitted by a backend that owns its agent.
13///
14/// The backend queue is synchronous and destructive. Each fact is moved out
15/// before this task awaits the ordinary settlement path, so neither its queue
16/// lock nor the dispatch lock can survive into session teardown.
17pub(super) async fn outcome_loop(state: RunState) -> Result<()> {
18    loop {
19        // The feed belongs to the backend the role's drive installed, and that
20        // backend moves with the drive: at startup this client holds the
21        // default drive's backend, and the role's own drive arrives later with
22        // `welcome`. So the feed is read every round instead of being latched
23        // once. A backend that reports no endings of its own — a pane host, a
24        // child process, the in-process test runtime — answers `None`, which is
25        // a round with nothing to drain and not the end of this task: waiting
26        // on the first answer forever is how a session whose drive turns out to
27        // be `acp` would run a whole turn with nobody listening for its ending.
28        let Some(feed) = state.dispatch.outcome_feed() else {
29            sleep(Duration::from_millis(OUTCOME_POLL_MS)).await;
30            continue;
31        };
32        while let Some(outcome) = feed.try_recv() {
33            settle_session_outcome(&state, outcome).await?;
34        }
35        sleep(Duration::from_millis(OUTCOME_POLL_MS)).await;
36    }
37}
38
39/// Feed one self-driven ending through the same fault and settlement paths an
40/// adapter report uses.
41///
42/// A backend that owns its agent reports two shapes of ending. A verdict — how
43/// the work ended — settles here, through the same door a plugin's completion
44/// takes. A turn that simply ended settles nothing: an agent that stopped its
45/// turn without calling `onlyne_complete` is the case §3c's client-owned rule
46/// answers, and `None` is exactly that shape.
47pub(super) async fn settle_session_outcome(
48    state: &RunState,
49    outcome: SessionOutcome,
50) -> Result<()> {
51    let SessionOutcome {
52        task_id,
53        outcome,
54        head,
55        note,
56        refusals,
57    } = outcome;
58    if let Some(reason) = refusals.as_deref() {
59        crate::reconcile::record_fault(&state.store, &task_id, "permission", "acp", reason)?;
60    }
61    // The turn ended and its drive asked for nothing: §3c's rule owns the rest
62    // — the one nudge, and the settlement when a second turn ends the same way.
63    // The agent's own closing line rides along, because the settlement this call
64    // may reach writes it into the task's head.
65    //
66    // The turn-ended fact is fed here for the same reason the turn-started fact
67    // is fed where the dispatch path hands a turn to the backend: a self-driven
68    // drive has no heartbeats, so the row's agent phase would otherwise never
69    // reach `idle`, and the never-ran guard that reads it would refuse a
70    // completion the agent files through its tools mount.
71    state.dispatch.feed_turn_ended(&task_id);
72    let Some(outcome) = outcome else {
73        return dispatch::on_turn_end(&state.dispatch, &task_id, head).await;
74    };
75    // A self-driven backend reports in the task's own vocabulary, and the
76    // settlement travels in the wire's. `pending` is the absence of a verdict,
77    // which is nothing this loop can settle: a terminal outcome is what an
78    // ending reports.
79    let terminal = dispatch::task_outcome_of(outcome).ok_or_else(|| {
80        anyhow!("self-driven backend reported a non-terminal outcome for task {task_id}")
81    })?;
82    if outcome == onlyne_proto::TaskState::Failed
83        && let Some(reason) = note.as_deref()
84    {
85        crate::reconcile::record_fault(&state.store, &task_id, "acp", "acp", reason)?;
86    }
87    dispatch::on_out(
88        &state.dispatch,
89        &task_id,
90        terminal,
91        head,
92        None,
93        // The ending came from this client's own backend, which watched the agent
94        // it is reporting: the never-ran guard belongs to the plugin's door, where
95        // the claimant and the claim are the same party.
96        dispatch::SettleAuthority::ClientOwned,
97    )
98    .await
99}
100
101/// One delivery becomes a session, or an immediate refusal ack.
102///
103/// A plugin mounted before any work existed is parked in the dispatcher, so the
104/// session staged here hands straight over to it. That is the order an
105/// always-running agent takes: it attaches first and receives its assignment
106/// when a task arrives (plan §6 line 285).
107///
108/// Only a `Task` costs a session slot, so only a `Task` waits at the capacity
109/// gate. A `Completion` is a terminal receipt and a `Note` wakes a session that
110/// already exists; neither opens a slot, and refusing one because the role is
111/// full would hold back work that is already done (the verdict of a task) or
112/// traffic that has no session to create.
113pub(super) async fn accept_delivery(state: &RunState, delivery: &Delivery) {
114    // A control command acts on the work the role already holds, so it answers
115    // before the capacity gate and before the `accept_new` gate: a role at
116    // `max_sessions` is exactly the role whose operator wants to free.
117    if delivery.envelope.kind == onlyne_proto::MsgKind::Control {
118        settle_control(state, delivery).await;
119        return;
120    }
121    // A task this role already finished is not new work. The server re-offers an
122    // unacknowledged row after a link flap or an operator repair, and a
123    // completion that was in flight when the link dropped can land after the
124    // requeue, so this row's task may already be `Done` here. Dispatching it
125    // again would stage its payload on whichever session is idle — one chain's
126    // task running inside another conversation, with a second answer aimed at
127    // the ledger row the first one settled. Acknowledge the row and run nothing.
128    if delivery.envelope.kind == onlyne_proto::MsgKind::Task
129        && let Some(task_id) = delivery.envelope.task_id()
130        && state.dispatch.task_completed_here(task_id)
131    {
132        tracing::warn!(
133            msg_id = %delivery.msg_id,
134            task = %task_id,
135            "redelivery of a finished task settled without running it"
136        );
137        state.dispatch.push_settled(AckArgs {
138            msg_id: delivery.msg_id.clone(),
139            op_id: None,
140            accepted: true,
141            reason: Some("task already completed by this role".to_string()),
142        });
143        return;
144    }
145    // A `Completion` is a terminal receipt, so it settles the row it names and
146    // starts no session (plan §3 line 152's `Completion`).
147    if delivery.envelope.kind == onlyne_proto::MsgKind::Completion {
148        state.dispatch.push_settled(AckArgs {
149            msg_id: delivery.msg_id.clone(),
150            op_id: None,
151            accepted: true,
152            reason: None,
153        });
154        return;
155    }
156    // A `Note` names no task, so it starts no session: it is the wake-up a role
157    // sends to a running agent (§3), and an agent that does not exist yet has
158    // nothing to wake. §5's `note_queue` keeps one out of the queue when its
159    // role is offline, and this is the matching half on the receiving side.
160    if delivery.envelope.kind == onlyne_proto::MsgKind::Note {
161        let injected = state.dispatch.inject_note(&delivery.envelope).await;
162        state.dispatch.push_settled(AckArgs {
163            msg_id: delivery.msg_id.clone(),
164            op_id: None,
165            accepted: injected,
166            reason: (!injected).then(|| "note has no live session to wake".to_string()),
167        });
168        return;
169    }
170    // What is left creates a session, so this is where `max_sessions` bites.
171    if !state.dispatch.has_capacity() {
172        // The row stays in flight on the server, which offers it again when a
173        // session frees (plan §5 `max_sessions`).
174        tracing::debug!(msg_id = %delivery.msg_id, "delivery waits for a free session");
175        return;
176    }
177    let accept_new = state.accept_new.load(Ordering::SeqCst);
178    let path = AcceptPath::new(state.dispatch.clone(), state.dispatch.role_prose());
179    match path.accept_new(delivery, accept_new) {
180        Ok(Some(session)) => {
181            if let Some(task_id) = delivery.envelope.task_id() {
182                state.dispatch.attach_msg_id(task_id, &delivery.msg_id);
183            }
184            // A plugin attached to this session takes the payload now, or the
185            // one parked for the role does; a session whose own plugin is
186            // still starting waits for its mount to hand it over.
187            if let Err(error) = state.dispatch.hand_staged(&session.task_id).await {
188                tracing::warn!(error = %error, task = %session.task_id, "staged hand-off refused");
189            }
190        }
191        // The gate is the connection's own (`watch_readiness` shuts it when the
192        // link leaves `Ready` and opens it when the redial lands), so a delivery
193        // the pull already had in hand when the link flapped arrives here with the
194        // gate shut. That answer is not this client's to give: a refusal settles
195        // the row `rejected`, which is terminal, and the row the teardown's
196        // requeue would have brought back is destroyed instead. The row stays in
197        // flight — unanswered is not a decision — and the next `hello` that does
198        // not claim it is what puts it back on the queue.
199        //
200        // The same answer covers the scope's own wait: a delivery whose family's
201        // session is mid-delivery, or one that arrives with every slot spent,
202        // belongs to a session this role will free. It waits the same way.
203        Ok(None) => tracing::debug!(
204            msg_id = %delivery.msg_id,
205            "no session takes the delivery yet; it stays in flight"
206        ),
207        Err(error) => {
208            tracing::warn!(error = %error, msg_id = %delivery.msg_id, "delivery refused");
209            state.dispatch.push_settled(AckArgs {
210                msg_id: delivery.msg_id.clone(),
211                op_id: None,
212                accepted: false,
213                reason: Some(error.to_string()),
214            });
215        }
216    }
217}
218
219/// Report running sessions whose Applied clock has exceeded the stall
220/// threshold. The fault is observation-only; the ledger row stays as stored.
221pub(super) async fn scan_stalls(state: &RunState) {
222    if state.stall_report_secs == 0 {
223        return;
224    }
225    let due = state
226        .dispatch
227        .stall_due(Instant::now(), state.stall_report_secs);
228    for task_id in due {
229        let Some(report) = state.dispatch.stall_report(&task_id) else {
230            continue;
231        };
232        match dispatch::send_frame(&state.dispatch, ClientOp::Report(report)).await {
233            Ok(()) => state.dispatch.mark_stalled(&task_id),
234            Err(error) => {
235                tracing::warn!(error = %error, task = %task_id, "stall fault was not sent")
236            }
237        }
238    }
239}
240
241/// Reclaim the resources of sessions this client has already put past their work,
242/// and publish each one's exit.
243///
244/// The reclaim runs on every readiness tick, ahead of the reconnect window, so a
245/// completed session whose plugin left without a goodbye is this sweep's to end:
246/// its resource is still open, its stored lifecycle already reads `Exited`, and
247/// the sweep below waits out a grace the completion itself did not ask for. What
248/// the reclaim writes is the agent's exit and the resource close, and the server
249/// mirrors only what this client reports — so each session it retired travels the
250/// report an ordinary ending travels, once the lock has been given back and the
251/// stored row is final. Without it the mirror keeps the reading the settle
252/// published, `exited` beside an agent still `running` and a resource still
253/// `attached`, which is what a peer's census of completed sessions found.
254///
255/// A published exit runs the server's `release_exited_delivery` for that task, and
256/// the task's own delivery row was answered when the completion settled it, so
257/// there is no in-flight row of that task for the release to hand back. The
258/// publish is also after the reclaim rather than around it because the row the
259/// server mirrors is the row the reclaim writes.
260pub(super) async fn scan_reclaimed_resources(state: &RunState) {
261    for session_id in state.dispatch.reclaim_exited_resources() {
262        if let Err(error) = dispatch::sync_session(&state.dispatch, &session_id).await {
263            tracing::warn!(
264                session = %session_id,
265                error = %error,
266                "a reclaimed session's exit was not published"
267            );
268        }
269    }
270}
271
272/// Retire the sessions whose plugin connection dropped and did not come back
273/// within `[client] reconnect_grace_secs`, and publish each one's exit. The tick
274/// sweeps every session the window expired on, bound to a task or not: a plugin
275/// that never came back is an agent that is gone, whether or not its work was
276/// still owed.
277///
278/// The publish is the half of the ending the sweep cannot write: the retirement
279/// feeds the session's own tuple to `Exited` and files the verdict, and the server
280/// only learns either from what this client reports. Without it the mirrored row
281/// keeps reading `working` until the server's own observer records a
282/// `stale_working` or `heartbeat_missing` fault — a reader waits for a fault, which
283/// names the silence and moves no row, to hear what this client already knew. So
284/// each session that left travels the report an ordinary ending already travels,
285/// once the lock has been given back and the stored row is final. Only the durable
286/// queue refusing the frame reaches the log: a live send that gave up is
287/// `sync_session`'s own fallback to that queue, not a lost publish.
288pub(super) async fn scan_reconnect_grace(state: &RunState) {
289    if state.reconnect_grace_secs == 0 {
290        return;
291    }
292    let retired = state
293        .dispatch
294        .retire_dropped_ghosts(Instant::now(), state.reconnect_grace_secs);
295    if retired.is_empty() {
296        return;
297    }
298    // One line per retirement, and it names the arm: two readings close this window — the
299    // connection ended and stayed away, or the connection held while nothing the client
300    // accepted arrived — and a single count with one threshold made an operator reading the
301    // log guess. The ages are the sweep's own inputs, so the line settles whether the agent
302    // left or merely stopped reporting.
303    for retired in retired {
304        tracing::info!(
305            session = %retired.session_id,
306            arm = retired.arm.word(),
307            quiet_secs = retired.quiet_secs,
308            away_secs = retired.away_secs,
309            silence_window_secs =
310                dispatch::HEARTBEAT_INTERVAL.as_secs() * dispatch::HEARTBEAT_SILENCE_MARGIN as u64,
311            grace_secs = state.reconnect_grace_secs,
312            "session retired past its window"
313        );
314        if let Err(error) = dispatch::sync_session(&state.dispatch, &retired.session_id).await {
315            tracing::warn!(
316                session = %retired.session_id,
317                error = %error,
318                "a retired session's exit was not published"
319            );
320        }
321    }
322}
323
324/// Settle the work an operator's word left open unanswered, and publish each
325/// one's exit.
326///
327/// `recycle` and `cancel` ask a session's plugin for its own ending, and the
328/// completion that answers the command is a frame of the plugin's. A plugin that
329/// never sends one — it left with the command's frame, or implements no
330/// `recycle` at all — leaves the task open, the mirrored row reading `working`,
331/// and the delivery row this client was handed in flight, and nothing in this
332/// process is left to answer any of the three: the close the command ran is what
333/// ended the session's own row already, so the sweep above finds no window left
334/// open on it and no work of it to settle.
335///
336/// What answers the word is the client's own record of it, held past
337/// `dispatch::CONTROL_SETTLE_BOUND`. The note is read under the dispatch lock and
338/// spent one at a time behind it, so a completion that arrives in between
339/// settles the task through the report it came on and this sweep writes nothing;
340/// the verdict, the refusal of the delivery row and the report behind it are
341/// [`DispatchState::settle_unanswered_control`]'s. The publish is this sweep's,
342/// for the reason the retirement above publishes: the server mirrors what this
343/// client reports, and that is what moves the row an operator is reading.
344///
345/// [`DispatchState::settle_unanswered_control`]:
346///     crate::session::dispatch::DispatchState::settle_unanswered_control
347pub(super) async fn scan_control_settles(state: &RunState) {
348    let now = Instant::now();
349    for note in state.dispatch.control_settles_due(now) {
350        // The note is spent through the one door that spends it: a completion
351        // that answered this word between the reading above and this call takes
352        // the note first, and its verdict is the one that stands.
353        if !state.dispatch.settle_unanswered_control(&note) {
354            continue;
355        }
356        tracing::info!(
357            task = %note.task_id,
358            outcome = ?note.word.outcome(),
359            waited_secs = now.saturating_duration_since(note.noted_at).as_secs(),
360            "a task was settled on an operator's word no plugin answered"
361        );
362        if let Err(error) = dispatch::sync_session(&state.dispatch, &note.task_id).await {
363            tracing::warn!(
364                task = %note.task_id,
365                error = %error,
366                "a settled task's exit was not published"
367            );
368        }
369    }
370}
371
372/// Release the process of every idle session whose scope bound has expired, and
373/// publish each one's row.
374///
375/// The suspension is [`DispatchState::suspend_idle_sessions`]'s: the session is
376/// marked suspended, its row moves through the `Suspend` event, and one backend
377/// close is collected per session. What this sweep adds is the report — a
378/// released process sends no frame of its own, so the row it just left would sit
379/// on the mirror as the last thing the living session published. A session whose
380/// runtime cannot resume is left exactly where it is, which is the "process
381/// alive, session alive" degradation the scope table names: this client will not
382/// release a process it cannot bring back.
383pub(super) async fn scan_idle_sessions(state: &RunState) {
384    for session_id in state.dispatch.suspend_idle_sessions(Instant::now()) {
385        if let Err(error) = dispatch::sync_session(&state.dispatch, &session_id).await {
386            tracing::warn!(
387                session = %session_id,
388                error = %error,
389                "a suspended session's row was not published"
390            );
391        }
392    }
393}
394
395pub(super) async fn refresh_role_slice(link: &ClientLink, state: &RunState) -> Result<()> {
396    let role = state.dispatch.role();
397    let reply = link
398        .request(ClientOp::QueryRoles(QueryRolesArgs { role: Some(role) }))
399        .await?;
400    if !reply.ok {
401        tracing::warn!(error = ?reply.error, "role slice refresh query refused");
402        return Ok(());
403    }
404    let rows: Vec<RoleInfo> = reply
405        .data
406        .as_ref()
407        .and_then(|value| value.get("roles"))
408        .cloned()
409        .map(serde_json::from_value)
410        .transpose()?
411        .unwrap_or_default();
412    if let Some(info) = rows.first() {
413        apply_role_info(state, info);
414    }
415    Ok(())
416}
417
418pub(super) fn apply_role_info(state: &RunState, info: &RoleInfo) -> Vec<&'static str> {
419    let current = state.dispatch.role_slice();
420    let next = crate::session::slice::RoleSlice::from_role_info(info, &current);
421    let Some((applied, fields)) = crate::session::slice::apply_if_changed(&current, next) else {
422        return Vec::new();
423    };
424    // A drive can move under a live link, and the backend it selects has to move
425    // with it before the new command is handed to the old one. A drive this
426    // machine cannot host is recorded as a refusal rather than silently kept.
427    state.install_runtime(applied.drive);
428    state.dispatch.reconfigure(applied);
429    fields
430}
431
432/// The local accept path for the current role slice.
433pub fn accept_path(state: &RunState) -> Result<AcceptPath> {
434    Ok(AcceptPath::new(
435        state.dispatch.clone(),
436        state.dispatch.role_prose(),
437    ))
438}