Skip to main content

spatialrust_distribute/
backpressure.rs

1//! Backpressure contracts and bounded transfer queues.
2
3use crate::{DistributeError, DistributeResult, NamedTransfer};
4
5/// Signal indicating consumer pressure.
6#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash)]
7pub enum BackpressureSignal {
8    /// Accept more work.
9    Ok,
10    /// Soft warn at high watermark.
11    SoftLimit,
12    /// Hard reject.
13    HardLimit,
14}
15
16/// Watermark policy for explicit backpressure.
17#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash)]
18pub struct BackpressurePolicy {
19    /// Soft watermark.
20    pub soft_limit: usize,
21    /// Hard watermark.
22    pub hard_limit: usize,
23}
24
25impl BackpressurePolicy {
26    /// Creates a validated policy.
27    pub fn try_new(soft_limit: usize, hard_limit: usize) -> DistributeResult<Self> {
28        if soft_limit == 0 || hard_limit < soft_limit {
29            return Err(DistributeError::InvalidConfiguration(
30                "require 0 < soft_limit <= hard_limit".into(),
31            ));
32        }
33        Ok(Self { soft_limit, hard_limit })
34    }
35
36    /// Evaluates queue depth against watermarks.
37    #[must_use]
38    pub fn evaluate(self, depth: usize) -> BackpressureSignal {
39        if depth >= self.hard_limit {
40            BackpressureSignal::HardLimit
41        } else if depth >= self.soft_limit {
42            BackpressureSignal::SoftLimit
43        } else {
44            BackpressureSignal::Ok
45        }
46    }
47}
48
49/// FIFO queue of named transfers with watermark-driven admissions.
50#[derive(Clone, Debug)]
51pub struct BoundedTransferQueue {
52    policy: BackpressurePolicy,
53    items: Vec<NamedTransfer>,
54    soft_trips: u64,
55    hard_rejects: u64,
56}
57
58impl BoundedTransferQueue {
59    /// Creates an empty queue.
60    #[must_use]
61    pub fn new(policy: BackpressurePolicy) -> Self {
62        Self { policy, items: Vec::new(), soft_trips: 0, hard_rejects: 0 }
63    }
64
65    /// Current depth.
66    #[must_use]
67    pub fn depth(&self) -> usize {
68        self.items.len()
69    }
70
71    /// Soft-limit trip count.
72    #[must_use]
73    pub fn soft_trips(&self) -> u64 {
74        self.soft_trips
75    }
76
77    /// Hard-reject count.
78    #[must_use]
79    pub fn hard_rejects(&self) -> u64 {
80        self.hard_rejects
81    }
82
83    /// Current pressure signal for `depth()`.
84    #[must_use]
85    pub fn signal(&self) -> BackpressureSignal {
86        self.policy.evaluate(self.depth())
87    }
88
89    /// Attempts to enqueue a transfer; hard-limit depths are rejected.
90    pub fn try_push(&mut self, transfer: NamedTransfer) -> DistributeResult<BackpressureSignal> {
91        let signal = self.policy.evaluate(self.depth());
92        match signal {
93            BackpressureSignal::HardLimit => {
94                self.hard_rejects += 1;
95                Err(DistributeError::CapacityExceeded {
96                    queue: transfer.name,
97                    depth: self.depth(),
98                    hard_limit: self.policy.hard_limit,
99                })
100            }
101            BackpressureSignal::SoftLimit => {
102                self.soft_trips += 1;
103                self.items.push(transfer);
104                Ok(BackpressureSignal::SoftLimit)
105            }
106            BackpressureSignal::Ok => {
107                self.items.push(transfer);
108                Ok(self.policy.evaluate(self.depth()))
109            }
110        }
111    }
112
113    /// Pops the oldest transfer.
114    pub fn pop(&mut self) -> Option<NamedTransfer> {
115        if self.items.is_empty() {
116            None
117        } else {
118            Some(self.items.remove(0))
119        }
120    }
121}
122
123#[cfg(test)]
124mod tests {
125    use super::{BackpressurePolicy, BackpressureSignal, BoundedTransferQueue};
126    use crate::{NamedTransfer, TransferDirection, TransferKind};
127
128    fn sample(name: &str) -> NamedTransfer {
129        NamedTransfer::try_new(
130            name,
131            TransferDirection::HostToNetwork,
132            TransferKind::ExplicitCopy,
133            "a",
134            "b",
135            8,
136        )
137        .unwrap()
138    }
139
140    #[test]
141    fn soft_and_hard_limits() {
142        let policy = BackpressurePolicy::try_new(1, 2).unwrap();
143        let mut queue = BoundedTransferQueue::new(policy);
144        assert_eq!(queue.try_push(sample("t0")).unwrap(), BackpressureSignal::SoftLimit);
145        assert_eq!(queue.try_push(sample("t1")).unwrap(), BackpressureSignal::SoftLimit);
146        assert!(queue.try_push(sample("t2")).is_err());
147        assert_eq!(queue.hard_rejects(), 1);
148        assert_eq!(queue.pop().unwrap().name, "t0");
149    }
150}