onlyne_client/session/stall.rs
1//! Client-side zero-activity stall observation.
2//!
3//! A freeze past `stall_report_secs` is reported as `Report::Fault` kind
4//! `stalled`. The ledger row stays as stored. Complete is not sent.
5
6use onlyne_proto::Report;
7use std::collections::{HashMap, HashSet};
8use std::time::{Duration, Instant};
9
10/// Fault kind stamped on a stall observation.
11pub const STALLED: &str = "stalled";
12/// Reason text carried with the stall fault.
13pub const STALLED_REASON: &str = "no applied progress";
14
15/// Per-task progress clock and one-shot freeze reporting.
16#[derive(Debug, Default)]
17pub struct StallWatch {
18 last_progress: HashMap<String, Instant>,
19 reported: HashSet<String>,
20 /// The newest `(generation, seq)` each task's reporter has been seen at.
21 /// A heartbeat's version is a watermark: a beat at or below it is a replay,
22 /// and the row it names does not move for a beat that observed nothing.
23 beats: HashMap<String, (u64, u64)>,
24}
25
26impl StallWatch {
27 pub fn new() -> Self {
28 Self::default()
29 }
30
31 /// Start the clock when a task is assigned. A later assign of the same
32 /// task keeps the original instant.
33 pub fn note_assigned(&mut self, task_id: impl Into<String>, now: Instant) {
34 self.last_progress.entry(task_id.into()).or_insert(now);
35 }
36
37 /// An Applied persist refreshes the clock of an assigned task and clears
38 /// the freeze report bit so a later freeze can be reported again. Clock
39 /// creation belongs to [`Self::note_assigned`].
40 pub fn note_applied(&mut self, task_id: &str, now: Instant) {
41 if let Some(last_progress) = self.last_progress.get_mut(task_id) {
42 *last_progress = now;
43 self.reported.remove(task_id);
44 }
45 }
46
47 /// Drop a task that has left the live set.
48 pub fn forget(&mut self, task_id: &str) {
49 self.last_progress.remove(task_id);
50 self.reported.remove(task_id);
51 self.beats.remove(task_id);
52 }
53
54 /// Record one heartbeat's version, answering whether it is the newest one
55 /// this task's reporter has been seen at.
56 ///
57 /// The version orders the reporter's own frames: the same rule the stored
58 /// row's watermark used to apply, kept here because a beat that observed
59 /// nothing no longer moves the row. A reporter that came back with a new
60 /// generation reads as newer, which is what a rebase is for.
61 pub fn note_beat_seq(&mut self, task_id: &str, generation: u64, seq: u64) -> bool {
62 match self.beats.get(task_id) {
63 Some(&(seen_generation, seen_seq))
64 if (generation, seq) <= (seen_generation, seen_seq) =>
65 {
66 false
67 }
68 _ => {
69 self.beats.insert(task_id.to_string(), (generation, seq));
70 true
71 }
72 }
73 }
74
75 /// Task ids whose freeze exceeds `threshold_secs` and have not been
76 /// reported in this freeze episode. A threshold of zero reports nothing.
77 pub fn due(&self, now: Instant, threshold_secs: u64) -> Vec<String> {
78 if threshold_secs == 0 {
79 return Vec::new();
80 }
81 let limit = Duration::from_secs(threshold_secs);
82 let mut due: Vec<String> = self
83 .last_progress
84 .iter()
85 .filter(|(task_id, _)| !self.reported.contains(*task_id))
86 .filter(|(_, last)| now.saturating_duration_since(**last) > limit)
87 .map(|(task_id, _)| task_id.clone())
88 .collect();
89 due.sort();
90 due
91 }
92
93 /// Remember that this freeze episode has been reported.
94 pub fn mark_reported(&mut self, task_id: &str) {
95 self.reported.insert(task_id.to_string());
96 }
97}
98
99/// Observation-only fault for a stalled task. The server records a fault row
100/// and an advisory event; the ledger row is not flipped.
101pub fn report(
102 task_id: &str,
103 session_id: Option<String>,
104 generation: Option<u64>,
105 seq: Option<u64>,
106) -> Report {
107 Report::Fault {
108 task_id: Some(task_id.to_string()),
109 session_id,
110 generation,
111 seq,
112 kind: STALLED.to_string(),
113 reason: STALLED_REASON.to_string(),
114 desired: None,
115 observed: None,
116 }
117}