Skip to main content

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}