Coverage Report

Created: 2026-07-16 07:06

next uncovered line (L), next uncovered region (R), next uncovered branch (B)
/src/h2/src/proto/streams/counts.rs
Line
Count
Source
1
use super::*;
2
3
#[derive(Debug)]
4
pub(super) struct Counts {
5
    /// Acting as a client or server. This allows us to track which values to
6
    /// inc / dec.
7
    peer: peer::Dyn,
8
9
    /// Maximum number of locally initiated streams
10
    max_send_streams: usize,
11
12
    /// Current number of remote initiated streams
13
    num_send_streams: usize,
14
15
    /// Maximum number of remote initiated streams
16
    max_recv_streams: usize,
17
18
    /// Current number of locally initiated streams
19
    num_recv_streams: usize,
20
21
    /// Maximum number of pending locally reset streams
22
    max_local_reset_streams: usize,
23
24
    /// Current number of pending locally reset streams
25
    num_local_reset_streams: usize,
26
27
    /// Max number of "pending accept" streams that were remotely reset
28
    max_remote_reset_streams: usize,
29
30
    /// Current number of "pending accept" streams that were remotely reset
31
    num_remote_reset_streams: usize,
32
33
    /// Maximum number of locally reset streams due to protocol error across
34
    /// the lifetime of the connection.
35
    ///
36
    /// When this gets exceeded, we issue GOAWAYs.
37
    max_local_error_reset_streams: Option<usize>,
38
39
    /// Total number of locally reset streams due to protocol error across the
40
    /// lifetime of the connection.
41
    num_local_error_reset_streams: usize,
42
}
43
44
impl Counts {
45
    /// Create a new `Counts` using the provided configuration values.
46
12.8k
    pub fn new(peer: peer::Dyn, config: &Config) -> Self {
47
12.8k
        Counts {
48
12.8k
            peer,
49
12.8k
            max_send_streams: config.initial_max_send_streams,
50
12.8k
            num_send_streams: 0,
51
12.8k
            max_recv_streams: config.remote_max_initiated.unwrap_or(usize::MAX),
52
12.8k
            num_recv_streams: 0,
53
12.8k
            max_local_reset_streams: config.local_reset_max,
54
12.8k
            num_local_reset_streams: 0,
55
12.8k
            max_remote_reset_streams: config.remote_reset_max,
56
12.8k
            num_remote_reset_streams: 0,
57
12.8k
            max_local_error_reset_streams: config.local_max_error_reset_streams,
58
12.8k
            num_local_error_reset_streams: 0,
59
12.8k
        }
60
12.8k
    }
61
62
    /// Returns true when the next opened stream will reach capacity of outbound streams
63
    ///
64
    /// The number of client send streams is incremented in prioritize; send_request has to guess if
65
    /// it should wait before allowing another request to be sent.
66
474k
    pub fn next_send_stream_will_reach_capacity(&self) -> bool {
67
474k
        self.max_send_streams <= (self.num_send_streams + 1)
68
474k
    }
69
70
    /// Returns the current peer
71
1.07M
    pub fn peer(&self) -> peer::Dyn {
72
1.07M
        self.peer
73
1.07M
    }
74
75
2.01M
    pub fn has_streams(&self) -> bool {
76
2.01M
        self.num_send_streams != 0 || self.num_recv_streams != 0
77
2.01M
    }
78
79
    /// Returns true if we can issue another local reset due to protocol error.
80
112k
    pub fn can_inc_num_local_error_resets(&self) -> bool {
81
112k
        if let Some(max) = self.max_local_error_reset_streams {
82
112k
            max > self.num_local_error_reset_streams
83
        } else {
84
0
            true
85
        }
86
112k
    }
87
88
56.1k
    pub fn inc_num_local_error_resets(&mut self) {
89
56.1k
        assert!(self.can_inc_num_local_error_resets());
90
91
        // Increment the number of remote initiated streams
92
56.1k
        self.num_local_error_reset_streams += 1;
93
56.1k
    }
94
95
0
    pub(crate) fn max_local_error_resets(&self) -> Option<usize> {
96
0
        self.max_local_error_reset_streams
97
0
    }
98
99
    /// Returns true if the receive stream concurrency can be incremented
100
1.16k
    pub fn can_inc_num_recv_streams(&self) -> bool {
101
1.16k
        self.max_recv_streams > self.num_recv_streams
102
1.16k
    }
103
104
    /// Increments the number of concurrent receive streams.
105
    ///
106
    /// # Panics
107
    ///
108
    /// Panics on failure as this should have been validated before hand.
109
0
    pub fn inc_num_recv_streams(&mut self, stream: &mut store::Ptr) {
110
0
        assert!(self.can_inc_num_recv_streams());
111
0
        assert!(!stream.is_counted);
112
113
        // Increment the number of remote initiated streams
114
0
        self.num_recv_streams += 1;
115
0
        stream.is_counted = true;
116
0
    }
117
118
    /// Returns true if the send stream concurrency can be incremented
119
696k
    pub fn can_inc_num_send_streams(&self) -> bool {
120
696k
        self.max_send_streams > self.num_send_streams
121
696k
    }
122
123
    /// Increments the number of concurrent send streams.
124
    ///
125
    /// # Panics
126
    ///
127
    /// Panics on failure as this should have been validated before hand.
128
239k
    pub fn inc_num_send_streams(&mut self, stream: &mut store::Ptr) {
129
239k
        assert!(self.can_inc_num_send_streams());
130
239k
        assert!(!stream.is_counted);
131
132
        // Increment the number of remote initiated streams
133
239k
        self.num_send_streams += 1;
134
239k
        stream.is_counted = true;
135
239k
    }
136
137
    /// Returns true if the number of pending reset streams can be incremented.
138
173k
    pub fn can_inc_num_reset_streams(&self) -> bool {
139
173k
        self.max_local_reset_streams > self.num_local_reset_streams
140
173k
    }
141
142
    /// Increments the number of pending reset streams.
143
    ///
144
    /// # Panics
145
    ///
146
    /// Panics on failure as this should have been validated before hand.
147
54.0k
    pub fn inc_num_reset_streams(&mut self) {
148
54.0k
        assert!(self.can_inc_num_reset_streams());
149
150
54.0k
        self.num_local_reset_streams += 1;
151
54.0k
    }
152
153
0
    pub(crate) fn max_remote_reset_streams(&self) -> usize {
154
0
        self.max_remote_reset_streams
155
0
    }
156
157
    /// Returns true if the number of pending REMOTE reset streams can be
158
    /// incremented.
159
0
    pub(crate) fn can_inc_num_remote_reset_streams(&self) -> bool {
160
0
        self.max_remote_reset_streams > self.num_remote_reset_streams
161
0
    }
162
163
    /// Increments the number of pending REMOTE reset streams.
164
    ///
165
    /// # Panics
166
    ///
167
    /// Panics on failure as this should have been validated before hand.
168
0
    pub(crate) fn inc_num_remote_reset_streams(&mut self) {
169
0
        assert!(self.can_inc_num_remote_reset_streams());
170
171
0
        self.num_remote_reset_streams += 1;
172
0
    }
173
174
0
    pub(crate) fn dec_num_remote_reset_streams(&mut self) {
175
0
        assert!(self.num_remote_reset_streams > 0);
176
177
0
        self.num_remote_reset_streams -= 1;
178
0
    }
179
180
6.28k
    pub fn apply_remote_settings(&mut self, settings: &frame::Settings, is_initial: bool) {
181
6.28k
        match settings.max_concurrent_streams() {
182
590
            Some(val) => self.max_send_streams = val as usize,
183
869
            None if is_initial => self.max_send_streams = usize::MAX,
184
4.82k
            None => {}
185
        }
186
6.28k
    }
187
188
    /// Run a block of code that could potentially transition a stream's state.
189
    ///
190
    /// If the stream state transitions to closed, this function will perform
191
    /// all necessary cleanup.
192
    ///
193
    /// TODO: Is this function still needed?
194
2.22M
    pub fn transition<F, U>(&mut self, mut stream: store::Ptr, f: F) -> U
195
2.22M
    where
196
2.22M
        F: FnOnce(&mut Self, &mut store::Ptr) -> U,
197
    {
198
        // TODO: Does this need to be computed before performing the action?
199
2.22M
        let is_pending_reset = stream.is_pending_reset_expiration();
200
201
        // Run the action
202
2.22M
        let ret = f(self, &mut stream);
203
204
2.22M
        self.transition_after(stream, is_pending_reset);
205
206
2.22M
        ret
207
2.22M
    }
<h2::proto::streams::counts::Counts>::transition::<<h2::proto::streams::prioritize::Prioritize>::assign_connection_capacity<h2::proto::streams::store::Ptr>::{closure#0}, ()>
Line
Count
Source
194
33.2k
    pub fn transition<F, U>(&mut self, mut stream: store::Ptr, f: F) -> U
195
33.2k
    where
196
33.2k
        F: FnOnce(&mut Self, &mut store::Ptr) -> U,
197
    {
198
        // TODO: Does this need to be computed before performing the action?
199
33.2k
        let is_pending_reset = stream.is_pending_reset_expiration();
200
201
        // Run the action
202
33.2k
        let ret = f(self, &mut stream);
203
204
33.2k
        self.transition_after(stream, is_pending_reset);
205
206
33.2k
        ret
207
33.2k
    }
<h2::proto::streams::counts::Counts>::transition::<<h2::proto::streams::prioritize::Prioritize>::assign_connection_capacity<h2::proto::streams::store::Store>::{closure#0}, ()>
Line
Count
Source
194
41.6k
    pub fn transition<F, U>(&mut self, mut stream: store::Ptr, f: F) -> U
195
41.6k
    where
196
41.6k
        F: FnOnce(&mut Self, &mut store::Ptr) -> U,
197
    {
198
        // TODO: Does this need to be computed before performing the action?
199
41.6k
        let is_pending_reset = stream.is_pending_reset_expiration();
200
201
        // Run the action
202
41.6k
        let ret = f(self, &mut stream);
203
204
41.6k
        self.transition_after(stream, is_pending_reset);
205
206
41.6k
        ret
207
41.6k
    }
<h2::proto::streams::counts::Counts>::transition::<h2::proto::streams::streams::drop_stream_ref::{closure#0}::{closure#0}, ()>
Line
Count
Source
194
189
    pub fn transition<F, U>(&mut self, mut stream: store::Ptr, f: F) -> U
195
189
    where
196
189
        F: FnOnce(&mut Self, &mut store::Ptr) -> U,
197
    {
198
        // TODO: Does this need to be computed before performing the action?
199
189
        let is_pending_reset = stream.is_pending_reset_expiration();
200
201
        // Run the action
202
189
        let ret = f(self, &mut stream);
203
204
189
        self.transition_after(stream, is_pending_reset);
205
206
189
        ret
207
189
    }
<h2::proto::streams::counts::Counts>::transition::<<h2::proto::streams::prioritize::Prioritize>::clear_pending_capacity::{closure#0}, ()>
Line
Count
Source
194
15.9k
    pub fn transition<F, U>(&mut self, mut stream: store::Ptr, f: F) -> U
195
15.9k
    where
196
15.9k
        F: FnOnce(&mut Self, &mut store::Ptr) -> U,
197
    {
198
        // TODO: Does this need to be computed before performing the action?
199
15.9k
        let is_pending_reset = stream.is_pending_reset_expiration();
200
201
        // Run the action
202
15.9k
        let ret = f(self, &mut stream);
203
204
15.9k
        self.transition_after(stream, is_pending_reset);
205
206
15.9k
        ret
207
15.9k
    }
Unexecuted instantiation: <h2::proto::streams::counts::Counts>::transition::<<h2::proto::streams::recv::Recv>::clear_stream_window_update_queue::{closure#0}, ()>
<h2::proto::streams::counts::Counts>::transition::<h2::proto::streams::streams::drop_stream_ref::{closure#0}, ()>
Line
Count
Source
194
949k
    pub fn transition<F, U>(&mut self, mut stream: store::Ptr, f: F) -> U
195
949k
    where
196
949k
        F: FnOnce(&mut Self, &mut store::Ptr) -> U,
197
    {
198
        // TODO: Does this need to be computed before performing the action?
199
949k
        let is_pending_reset = stream.is_pending_reset_expiration();
200
201
        // Run the action
202
949k
        let ret = f(self, &mut stream);
203
204
949k
        self.transition_after(stream, is_pending_reset);
205
206
949k
        ret
207
949k
    }
<h2::proto::streams::counts::Counts>::transition::<<h2::proto::streams::streams::Inner>::recv_eof<bytes::bytes::Bytes>::{closure#0}::{closure#0}, ()>
Line
Count
Source
194
760
    pub fn transition<F, U>(&mut self, mut stream: store::Ptr, f: F) -> U
195
760
    where
196
760
        F: FnOnce(&mut Self, &mut store::Ptr) -> U,
197
    {
198
        // TODO: Does this need to be computed before performing the action?
199
760
        let is_pending_reset = stream.is_pending_reset_expiration();
200
201
        // Run the action
202
760
        let ret = f(self, &mut stream);
203
204
760
        self.transition_after(stream, is_pending_reset);
205
206
760
        ret
207
760
    }
Unexecuted instantiation: <h2::proto::streams::counts::Counts>::transition::<<h2::proto::streams::recv::Recv>::send_stream_window_updates<fuzz_e2e::MockIo, bytes::bytes::Bytes>::{closure#0}, ()>
<h2::proto::streams::counts::Counts>::transition::<<h2::proto::streams::streams::Inner>::recv_reset<bytes::bytes::Bytes>::{closure#0}, core::result::Result<(), h2::proto::error::Error>>
Line
Count
Source
194
166
    pub fn transition<F, U>(&mut self, mut stream: store::Ptr, f: F) -> U
195
166
    where
196
166
        F: FnOnce(&mut Self, &mut store::Ptr) -> U,
197
    {
198
        // TODO: Does this need to be computed before performing the action?
199
166
        let is_pending_reset = stream.is_pending_reset_expiration();
200
201
        // Run the action
202
166
        let ret = f(self, &mut stream);
203
204
166
        self.transition_after(stream, is_pending_reset);
205
206
166
        ret
207
166
    }
<h2::proto::streams::counts::Counts>::transition::<<h2::proto::streams::streams::Inner>::recv_headers<bytes::bytes::Bytes>::{closure#0}, core::result::Result<(), h2::proto::error::Error>>
Line
Count
Source
194
1.74k
    pub fn transition<F, U>(&mut self, mut stream: store::Ptr, f: F) -> U
195
1.74k
    where
196
1.74k
        F: FnOnce(&mut Self, &mut store::Ptr) -> U,
197
    {
198
        // TODO: Does this need to be computed before performing the action?
199
1.74k
        let is_pending_reset = stream.is_pending_reset_expiration();
200
201
        // Run the action
202
1.74k
        let ret = f(self, &mut stream);
203
204
1.74k
        self.transition_after(stream, is_pending_reset);
205
206
1.74k
        ret
207
1.74k
    }
<h2::proto::streams::counts::Counts>::transition::<<h2::proto::streams::streams::Inner>::recv_push_promise<bytes::bytes::Bytes>::{closure#0}, core::result::Result<core::option::Option<h2::proto::streams::store::Key>, h2::proto::error::Error>>
Line
Count
Source
194
1.16k
    pub fn transition<F, U>(&mut self, mut stream: store::Ptr, f: F) -> U
195
1.16k
    where
196
1.16k
        F: FnOnce(&mut Self, &mut store::Ptr) -> U,
197
    {
198
        // TODO: Does this need to be computed before performing the action?
199
1.16k
        let is_pending_reset = stream.is_pending_reset_expiration();
200
201
        // Run the action
202
1.16k
        let ret = f(self, &mut stream);
203
204
1.16k
        self.transition_after(stream, is_pending_reset);
205
206
1.16k
        ret
207
1.16k
    }
<h2::proto::streams::counts::Counts>::transition::<<h2::proto::streams::streams::Inner>::recv_data<bytes::bytes::Bytes>::{closure#0}, core::result::Result<(), h2::proto::error::Error>>
Line
Count
Source
194
38.2k
    pub fn transition<F, U>(&mut self, mut stream: store::Ptr, f: F) -> U
195
38.2k
    where
196
38.2k
        F: FnOnce(&mut Self, &mut store::Ptr) -> U,
197
    {
198
        // TODO: Does this need to be computed before performing the action?
199
38.2k
        let is_pending_reset = stream.is_pending_reset_expiration();
200
201
        // Run the action
202
38.2k
        let ret = f(self, &mut stream);
203
204
38.2k
        self.transition_after(stream, is_pending_reset);
205
206
38.2k
        ret
207
38.2k
    }
<h2::proto::streams::counts::Counts>::transition::<<h2::proto::streams::streams::Actions>::send_reset<bytes::bytes::Bytes>::{closure#0}, core::result::Result<(), h2::proto::error::GoAway>>
Line
Count
Source
194
55.1k
    pub fn transition<F, U>(&mut self, mut stream: store::Ptr, f: F) -> U
195
55.1k
    where
196
55.1k
        F: FnOnce(&mut Self, &mut store::Ptr) -> U,
197
    {
198
        // TODO: Does this need to be computed before performing the action?
199
55.1k
        let is_pending_reset = stream.is_pending_reset_expiration();
200
201
        // Run the action
202
55.1k
        let ret = f(self, &mut stream);
203
204
55.1k
        self.transition_after(stream, is_pending_reset);
205
206
55.1k
        ret
207
55.1k
    }
<h2::proto::streams::counts::Counts>::transition::<<h2::proto::streams::streams::Inner>::handle_error<bytes::bytes::Bytes>::{closure#0}::{closure#0}, ()>
Line
Count
Source
194
259k
    pub fn transition<F, U>(&mut self, mut stream: store::Ptr, f: F) -> U
195
259k
    where
196
259k
        F: FnOnce(&mut Self, &mut store::Ptr) -> U,
197
    {
198
        // TODO: Does this need to be computed before performing the action?
199
259k
        let is_pending_reset = stream.is_pending_reset_expiration();
200
201
        // Run the action
202
259k
        let ret = f(self, &mut stream);
203
204
259k
        self.transition_after(stream, is_pending_reset);
205
206
259k
        ret
207
259k
    }
<h2::proto::streams::counts::Counts>::transition::<<h2::proto::streams::streams::Inner>::recv_go_away<bytes::bytes::Bytes>::{closure#0}::{closure#0}, ()>
Line
Count
Source
194
134k
    pub fn transition<F, U>(&mut self, mut stream: store::Ptr, f: F) -> U
195
134k
    where
196
134k
        F: FnOnce(&mut Self, &mut store::Ptr) -> U,
197
    {
198
        // TODO: Does this need to be computed before performing the action?
199
134k
        let is_pending_reset = stream.is_pending_reset_expiration();
200
201
        // Run the action
202
134k
        let ret = f(self, &mut stream);
203
204
134k
        self.transition_after(stream, is_pending_reset);
205
206
134k
        ret
207
134k
    }
<h2::proto::streams::counts::Counts>::transition::<<h2::proto::streams::streams::Inner>::recv_eof<bytes::bytes::Bytes>::{closure#0}::{closure#0}, ()>
Line
Count
Source
194
216k
    pub fn transition<F, U>(&mut self, mut stream: store::Ptr, f: F) -> U
195
216k
    where
196
216k
        F: FnOnce(&mut Self, &mut store::Ptr) -> U,
197
    {
198
        // TODO: Does this need to be computed before performing the action?
199
216k
        let is_pending_reset = stream.is_pending_reset_expiration();
200
201
        // Run the action
202
216k
        let ret = f(self, &mut stream);
203
204
216k
        self.transition_after(stream, is_pending_reset);
205
206
216k
        ret
207
216k
    }
<h2::proto::streams::counts::Counts>::transition::<<h2::proto::streams::streams::StreamRef<bytes::bytes::Bytes>>::send_data::{closure#0}, core::result::Result<(), h2::codec::error::UserError>>
Line
Count
Source
194
473k
    pub fn transition<F, U>(&mut self, mut stream: store::Ptr, f: F) -> U
195
473k
    where
196
473k
        F: FnOnce(&mut Self, &mut store::Ptr) -> U,
197
    {
198
        // TODO: Does this need to be computed before performing the action?
199
473k
        let is_pending_reset = stream.is_pending_reset_expiration();
200
201
        // Run the action
202
473k
        let ret = f(self, &mut stream);
203
204
473k
        self.transition_after(stream, is_pending_reset);
205
206
473k
        ret
207
473k
    }
208
209
    // TODO: move this to macro?
210
3.00M
    pub fn transition_after(&mut self, mut stream: store::Ptr, is_reset_counted: bool) {
211
3.00M
        tracing::trace!(
212
0
            "transition_after; stream={:?}; state={:?}; is_closed={:?}; \
213
0
             pending_send_empty={:?}; buffered_send_data={}; \
214
0
             num_recv={}; num_send={}",
215
0
            stream.id,
216
0
            stream.state,
217
0
            stream.is_closed(),
218
0
            stream.pending_send.is_empty(),
219
0
            stream.buffered_send_data,
220
            self.num_recv_streams,
221
            self.num_send_streams
222
        );
223
224
3.00M
        if stream.is_closed() {
225
1.55M
            if !stream.is_pending_reset_expiration() {
226
1.38M
                stream.unlink();
227
1.38M
                if is_reset_counted {
228
54.0k
                    self.dec_num_reset_streams();
229
1.32M
                }
230
175k
            }
231
232
1.55M
            if !stream.state.is_scheduled_reset() && stream.is_counted {
233
239k
                tracing::trace!("dec_num_streams; stream={:?}", stream.id);
234
                // Decrement the number of active streams.
235
239k
                self.dec_num_streams(&mut stream);
236
1.31M
            }
237
1.45M
        }
238
239
        // Release the stream if it requires releasing
240
3.00M
        if stream.is_released() {
241
508k
            stream.remove();
242
2.50M
        }
243
3.00M
    }
244
245
    /// Returns the maximum number of streams that can be initiated by this
246
    /// peer.
247
0
    pub(crate) fn max_send_streams(&self) -> usize {
248
0
        self.max_send_streams
249
0
    }
250
251
    /// Returns the maximum number of streams that can be initiated by the
252
    /// remote peer.
253
0
    pub(crate) fn max_recv_streams(&self) -> usize {
254
0
        self.max_recv_streams
255
0
    }
256
257
239k
    fn dec_num_streams(&mut self, stream: &mut store::Ptr) {
258
239k
        assert!(stream.is_counted);
259
260
239k
        if self.peer.is_local_init(stream.id) {
261
239k
            assert!(self.num_send_streams > 0);
262
239k
            self.num_send_streams -= 1;
263
239k
            stream.is_counted = false;
264
        } else {
265
0
            assert!(self.num_recv_streams > 0);
266
0
            self.num_recv_streams -= 1;
267
0
            stream.is_counted = false;
268
        }
269
239k
    }
270
271
54.0k
    fn dec_num_reset_streams(&mut self) {
272
54.0k
        assert!(self.num_local_reset_streams > 0);
273
54.0k
        self.num_local_reset_streams -= 1;
274
54.0k
    }
275
}
276
277
impl Drop for Counts {
278
12.8k
    fn drop(&mut self) {
279
        use std::thread;
280
281
12.8k
        if !thread::panicking() {
282
12.8k
            debug_assert!(!self.has_streams());
283
0
        }
284
12.8k
    }
285
}