spatialrust_distribute/
backpressure.rs1use crate::{DistributeError, DistributeResult, NamedTransfer};
4
5#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash)]
7pub enum BackpressureSignal {
8 Ok,
10 SoftLimit,
12 HardLimit,
14}
15
16#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash)]
18pub struct BackpressurePolicy {
19 pub soft_limit: usize,
21 pub hard_limit: usize,
23}
24
25impl BackpressurePolicy {
26 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 #[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#[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 #[must_use]
61 pub fn new(policy: BackpressurePolicy) -> Self {
62 Self { policy, items: Vec::new(), soft_trips: 0, hard_rejects: 0 }
63 }
64
65 #[must_use]
67 pub fn depth(&self) -> usize {
68 self.items.len()
69 }
70
71 #[must_use]
73 pub fn soft_trips(&self) -> u64 {
74 self.soft_trips
75 }
76
77 #[must_use]
79 pub fn hard_rejects(&self) -> u64 {
80 self.hard_rejects
81 }
82
83 #[must_use]
85 pub fn signal(&self) -> BackpressureSignal {
86 self.policy.evaluate(self.depth())
87 }
88
89 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 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}