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(¬e) {
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, ¬e.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, ¤t);
421 let Some((applied, fields)) = crate::session::slice::apply_if_changed(¤t, 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}