Skip to main content

h2ts_client/
flow.rs

1//! HTTP/2 send-side flow-control window (RFC 7540 §6.9) — port of `flow.ts`.
2//!
3//! One direction of one window (connection- or stream-level, send side). The
4//! window may go negative when the peer lowers SETTINGS_INITIAL_WINDOW_SIZE after
5//! we were granted capacity — that is legal and must not underflow a send.
6//!
7//! The TS version returns a `Promise` from `waitPositive`; here a caller polls
8//! [`is_ready`](SendWindow::is_ready) and parks its [`Waker`] via
9//! [`register_waker`](SendWindow::register_waker) (see `connection`'s body pump).
10
11use std::task::Waker;
12
13pub struct SendWindow {
14    available: i64,
15    wakers: Vec<Waker>,
16    closed: bool,
17}
18
19impl SendWindow {
20    pub fn new(initial: i64) -> Self {
21        Self {
22            available: initial,
23            wakers: Vec::new(),
24            closed: false,
25        }
26    }
27
28    pub fn value(&self) -> i64 {
29        self.available
30    }
31
32    /// Grant more capacity (a WINDOW_UPDATE arrived).
33    pub fn update(&mut self, increment: i64) {
34        self.available += increment;
35        if self.available > 0 {
36            self.wake();
37        }
38    }
39
40    /// Adjust by a SETTINGS_INITIAL_WINDOW_SIZE change (delta may be negative).
41    pub fn adjust(&mut self, delta: i64) {
42        self.available += delta;
43        if self.available > 0 {
44            self.wake();
45        }
46    }
47
48    /// Consume capacity that a positive check already confirmed is available.
49    pub fn consume(&mut self, n: i64) {
50        self.available -= n;
51    }
52
53    pub fn is_closed(&self) -> bool {
54        self.closed
55    }
56
57    /// Ready to send (positive capacity) or torn down.
58    pub fn is_ready(&self) -> bool {
59        self.available > 0 || self.closed
60    }
61
62    /// Park a waker to be notified when capacity becomes positive (or on close).
63    pub fn register_waker(&mut self, waker: &Waker) {
64        self.wakers.push(waker.clone());
65    }
66
67    /// Abort all waiters (connection/stream tearing down).
68    pub fn close(&mut self) {
69        self.closed = true;
70        self.wake();
71    }
72
73    fn wake(&mut self) {
74        for w in self.wakers.drain(..) {
75            w.wake();
76        }
77    }
78}