/src/openssl36/ssl/quic/quic_stream_map.c
Line | Count | Source |
1 | | /* |
2 | | * Copyright 2022-2026 The OpenSSL Project Authors. All Rights Reserved. |
3 | | * |
4 | | * Licensed under the Apache License 2.0 (the "License"). You may not use |
5 | | * this file except in compliance with the License. You can obtain a copy |
6 | | * in the file LICENSE in the source distribution or at |
7 | | * https://www.openssl.org/source/license.html |
8 | | */ |
9 | | |
10 | | #include "internal/quic_stream_map.h" |
11 | | #include "internal/nelem.h" |
12 | | #include "internal/quic_channel.h" |
13 | | |
14 | | /* |
15 | | * QUIC Stream Map |
16 | | * =============== |
17 | | */ |
18 | | DEFINE_LHASH_OF_EX(QUIC_STREAM); |
19 | | |
20 | | static void shutdown_flush_done(QUIC_STREAM_MAP *qsm, QUIC_STREAM *qs); |
21 | | |
22 | | /* Circular list management. */ |
23 | | static void list_insert_tail(QUIC_STREAM_LIST_NODE *l, |
24 | | QUIC_STREAM_LIST_NODE *n) |
25 | 161k | { |
26 | | /* Must not be in list. */ |
27 | 161k | assert(n->prev == NULL && n->next == NULL |
28 | 161k | && l->prev != NULL && l->next != NULL); |
29 | | |
30 | 161k | n->prev = l->prev; |
31 | 161k | n->prev->next = n; |
32 | 161k | l->prev = n; |
33 | 161k | n->next = l; |
34 | 161k | } |
35 | | |
36 | | static void list_remove(QUIC_STREAM_LIST_NODE *l, |
37 | | QUIC_STREAM_LIST_NODE *n) |
38 | 161k | { |
39 | 161k | assert(n->prev != NULL && n->next != NULL |
40 | 161k | && n->prev != n && n->next != n); |
41 | | |
42 | 161k | n->prev->next = n->next; |
43 | 161k | n->next->prev = n->prev; |
44 | 161k | n->next = n->prev = NULL; |
45 | 161k | } |
46 | | |
47 | | static QUIC_STREAM *list_next(QUIC_STREAM_LIST_NODE *l, QUIC_STREAM_LIST_NODE *n, |
48 | | size_t off) |
49 | 71.5M | { |
50 | 71.5M | assert(n->prev != NULL && n->next != NULL |
51 | 71.5M | && (n == l || (n->prev != n && n->next != n)) |
52 | 71.5M | && l->prev != NULL && l->next != NULL); |
53 | | |
54 | 71.5M | n = n->next; |
55 | | |
56 | 71.5M | if (n == l) |
57 | 71.5M | n = n->next; |
58 | 71.5M | if (n == l) |
59 | 71.5M | return NULL; |
60 | | |
61 | 71.5M | assert(n != NULL); |
62 | | |
63 | 43.0k | return (QUIC_STREAM *)(((char *)n) - off); |
64 | 43.0k | } |
65 | | |
66 | 42.6k | #define active_next(l, s) list_next((l), &(s)->active_node, \ |
67 | 42.6k | offsetof(QUIC_STREAM, active_node)) |
68 | 0 | #define accept_next(l, s) list_next((l), &(s)->accept_node, \ |
69 | 0 | offsetof(QUIC_STREAM, accept_node)) |
70 | 461 | #define accept_head(l) list_next((l), (l), \ |
71 | 461 | offsetof(QUIC_STREAM, accept_node)) |
72 | 71.5M | #define ready_for_gc_head(l) list_next((l), (l), \ |
73 | 71.5M | offsetof(QUIC_STREAM, ready_for_gc_node)) |
74 | | |
75 | | static unsigned long hash_stream(const QUIC_STREAM *s) |
76 | 17.0M | { |
77 | 17.0M | return (unsigned long)s->id; |
78 | 17.0M | } |
79 | | |
80 | | static int cmp_stream(const QUIC_STREAM *a, const QUIC_STREAM *b) |
81 | 320k | { |
82 | 320k | if (a->id < b->id) |
83 | 0 | return -1; |
84 | 320k | if (a->id > b->id) |
85 | 0 | return 1; |
86 | 320k | return 0; |
87 | 320k | } |
88 | | |
89 | | int ossl_quic_stream_map_init(QUIC_STREAM_MAP *qsm, |
90 | | uint64_t (*get_stream_limit_cb)(int uni, void *arg), |
91 | | void *get_stream_limit_cb_arg, |
92 | | QUIC_RXFC *max_streams_bidi_rxfc, |
93 | | QUIC_RXFC *max_streams_uni_rxfc, |
94 | | QUIC_CHANNEL *ch) |
95 | 41.4k | { |
96 | 41.4k | qsm->map = lh_QUIC_STREAM_new(hash_stream, cmp_stream); |
97 | 41.4k | qsm->active_list.prev = qsm->active_list.next = &qsm->active_list; |
98 | 41.4k | qsm->accept_list.prev = qsm->accept_list.next = &qsm->accept_list; |
99 | 41.4k | qsm->ready_for_gc_list.prev = qsm->ready_for_gc_list.next |
100 | 41.4k | = &qsm->ready_for_gc_list; |
101 | 41.4k | qsm->rr_stepping = 1; |
102 | 41.4k | qsm->rr_counter = 0; |
103 | 41.4k | qsm->rr_cur = NULL; |
104 | | |
105 | 41.4k | qsm->num_accept_bidi = 0; |
106 | 41.4k | qsm->num_accept_uni = 0; |
107 | 41.4k | qsm->num_shutdown_flush = 0; |
108 | | |
109 | 41.4k | qsm->get_stream_limit_cb = get_stream_limit_cb; |
110 | 41.4k | qsm->get_stream_limit_cb_arg = get_stream_limit_cb_arg; |
111 | 41.4k | qsm->max_streams_bidi_rxfc = max_streams_bidi_rxfc; |
112 | 41.4k | qsm->max_streams_uni_rxfc = max_streams_uni_rxfc; |
113 | 41.4k | qsm->ch = ch; |
114 | 41.4k | return 1; |
115 | 41.4k | } |
116 | | |
117 | | static void release_each(QUIC_STREAM *stream, void *arg) |
118 | 143k | { |
119 | 143k | QUIC_STREAM_MAP *qsm = arg; |
120 | | |
121 | 143k | ossl_quic_stream_map_release(qsm, stream); |
122 | 143k | } |
123 | | |
124 | | void ossl_quic_stream_map_cleanup(QUIC_STREAM_MAP *qsm) |
125 | 41.4k | { |
126 | 41.4k | lh_QUIC_STREAM_set_down_load(qsm->map, 0); |
127 | 41.4k | ossl_quic_stream_map_visit(qsm, release_each, qsm); |
128 | | |
129 | 41.4k | lh_QUIC_STREAM_free(qsm->map); |
130 | 41.4k | qsm->map = NULL; |
131 | 41.4k | } |
132 | | |
133 | | void ossl_quic_stream_map_visit(QUIC_STREAM_MAP *qsm, |
134 | | void (*visit_cb)(QUIC_STREAM *stream, void *arg), |
135 | | void *visit_cb_arg) |
136 | 198k | { |
137 | 198k | lh_QUIC_STREAM_doall_arg(qsm->map, visit_cb, visit_cb_arg); |
138 | 198k | } |
139 | | |
140 | | QUIC_STREAM *ossl_quic_stream_map_alloc(QUIC_STREAM_MAP *qsm, |
141 | | uint64_t stream_id, |
142 | | int type) |
143 | 84.6k | { |
144 | 84.6k | QUIC_STREAM *s; |
145 | 84.6k | QUIC_STREAM key; |
146 | | |
147 | 84.6k | key.id = stream_id; |
148 | | |
149 | 84.6k | s = lh_QUIC_STREAM_retrieve(qsm->map, &key); |
150 | 84.6k | if (s != NULL) |
151 | 0 | return NULL; |
152 | | |
153 | 84.6k | s = OPENSSL_zalloc(sizeof(*s)); |
154 | 84.6k | if (s == NULL) |
155 | 0 | return NULL; |
156 | | |
157 | 84.6k | s->id = stream_id; |
158 | 84.6k | s->type = type; |
159 | 84.6k | s->as_server = ossl_quic_channel_is_server(qsm->ch); |
160 | 84.6k | s->send_state = (ossl_quic_stream_is_local_init(s) |
161 | 82.1k | || ossl_quic_stream_is_bidi(s)) |
162 | 84.6k | ? QUIC_SSTREAM_STATE_READY |
163 | 84.6k | : QUIC_SSTREAM_STATE_NONE; |
164 | 84.6k | s->recv_state = (!ossl_quic_stream_is_local_init(s) |
165 | 2.43k | || ossl_quic_stream_is_bidi(s)) |
166 | 84.6k | ? QUIC_RSTREAM_STATE_RECV |
167 | 84.6k | : QUIC_RSTREAM_STATE_NONE; |
168 | | |
169 | 84.6k | s->send_final_size = UINT64_MAX; |
170 | | |
171 | 84.6k | lh_QUIC_STREAM_insert(qsm->map, s); |
172 | 84.6k | if (lh_QUIC_STREAM_error(qsm->map)) { |
173 | 0 | OPENSSL_free(s); |
174 | 0 | return NULL; |
175 | 0 | } |
176 | 84.6k | return s; |
177 | 84.6k | } |
178 | | |
179 | | void ossl_quic_stream_map_release(QUIC_STREAM_MAP *qsm, QUIC_STREAM *stream) |
180 | 150k | { |
181 | 150k | if (stream == NULL) |
182 | 6.97k | return; |
183 | | |
184 | 143k | if (stream->active_node.next != NULL) |
185 | 12.2k | list_remove(&qsm->active_list, &stream->active_node); |
186 | 143k | if (stream->accept_node.next != NULL) |
187 | 131k | list_remove(&qsm->accept_list, &stream->accept_node); |
188 | 143k | if (stream->ready_for_gc_node.next != NULL) |
189 | 0 | list_remove(&qsm->ready_for_gc_list, &stream->ready_for_gc_node); |
190 | | |
191 | 143k | ossl_quic_sstream_free(stream->sstream); |
192 | 143k | stream->sstream = NULL; |
193 | | |
194 | 143k | ossl_quic_rstream_free(stream->rstream); |
195 | 143k | stream->rstream = NULL; |
196 | | |
197 | 143k | lh_QUIC_STREAM_delete(qsm->map, stream); |
198 | 143k | OPENSSL_free(stream); |
199 | 143k | } |
200 | | |
201 | | QUIC_STREAM *ossl_quic_stream_map_get_by_id(QUIC_STREAM_MAP *qsm, |
202 | | uint64_t stream_id) |
203 | 16.6M | { |
204 | 16.6M | QUIC_STREAM key; |
205 | | |
206 | 16.6M | key.id = stream_id; |
207 | | |
208 | 16.6M | return lh_QUIC_STREAM_retrieve(qsm->map, &key); |
209 | 16.6M | } |
210 | | |
211 | | static void stream_map_mark_active(QUIC_STREAM_MAP *qsm, QUIC_STREAM *s) |
212 | 83.4k | { |
213 | 83.4k | if (s->active) |
214 | 60.5k | return; |
215 | | |
216 | 22.8k | list_insert_tail(&qsm->active_list, &s->active_node); |
217 | | |
218 | 22.8k | if (qsm->rr_cur == NULL) |
219 | 16.4k | qsm->rr_cur = s; |
220 | | |
221 | 22.8k | s->active = 1; |
222 | 22.8k | } |
223 | | |
224 | | static void stream_map_mark_inactive(QUIC_STREAM_MAP *qsm, QUIC_STREAM *s) |
225 | 1.04M | { |
226 | 1.04M | if (!s->active) |
227 | 1.03M | return; |
228 | | |
229 | 10.6k | if (qsm->rr_cur == s) |
230 | 9.41k | qsm->rr_cur = active_next(&qsm->active_list, s); |
231 | 10.6k | if (qsm->rr_cur == s) |
232 | 8.15k | qsm->rr_cur = NULL; |
233 | | |
234 | 10.6k | list_remove(&qsm->active_list, &s->active_node); |
235 | | |
236 | 10.6k | s->active = 0; |
237 | 10.6k | } |
238 | | |
239 | | void ossl_quic_stream_map_set_rr_stepping(QUIC_STREAM_MAP *qsm, size_t stepping) |
240 | 0 | { |
241 | 0 | qsm->rr_stepping = stepping; |
242 | 0 | qsm->rr_counter = 0; |
243 | 0 | } |
244 | | |
245 | | static int stream_has_data_to_send(QUIC_STREAM *s) |
246 | 1.04M | { |
247 | 1.04M | OSSL_QUIC_FRAME_STREAM shdr; |
248 | 1.04M | OSSL_QTX_IOVEC iov[2]; |
249 | 1.04M | size_t num_iov; |
250 | 1.04M | uint64_t fc_credit, fc_swm, fc_limit; |
251 | | |
252 | 1.04M | switch (s->send_state) { |
253 | 91.3k | case QUIC_SSTREAM_STATE_READY: |
254 | 129k | case QUIC_SSTREAM_STATE_SEND: |
255 | 129k | case QUIC_SSTREAM_STATE_DATA_SENT: |
256 | | /* |
257 | | * We can still have data to send in DATA_SENT due to retransmissions, |
258 | | * etc. |
259 | | */ |
260 | 129k | break; |
261 | 917k | default: |
262 | 917k | return 0; /* Nothing to send. */ |
263 | 1.04M | } |
264 | | |
265 | | /* |
266 | | * We cannot determine if we have data to send simply by checking if |
267 | | * ossl_quic_txfc_get_credit() is zero, because we may also have older |
268 | | * stream data we need to retransmit. The SSTREAM returns older data first, |
269 | | * so we do a simple comparison of the next chunk the SSTREAM wants to send |
270 | | * against the TXFC CWM. |
271 | | */ |
272 | 129k | num_iov = OSSL_NELEM(iov); |
273 | 129k | if (!ossl_quic_sstream_get_stream_frame(s->sstream, 0, &shdr, iov, |
274 | 129k | &num_iov)) |
275 | 101k | return 0; |
276 | | |
277 | 27.8k | fc_credit = ossl_quic_txfc_get_credit(&s->txfc, 0); |
278 | 27.8k | fc_swm = ossl_quic_txfc_get_swm(&s->txfc); |
279 | 27.8k | fc_limit = fc_swm + fc_credit; |
280 | | |
281 | 27.8k | return (shdr.is_fin && shdr.len == 0) || shdr.offset < fc_limit; |
282 | 129k | } |
283 | | |
284 | | static ossl_unused int qsm_send_part_permits_gc(const QUIC_STREAM *qs) |
285 | 0 | { |
286 | 0 | switch (qs->send_state) { |
287 | 0 | case QUIC_SSTREAM_STATE_NONE: |
288 | 0 | case QUIC_SSTREAM_STATE_DATA_RECVD: |
289 | 0 | case QUIC_SSTREAM_STATE_RESET_RECVD: |
290 | 0 | return 1; |
291 | 0 | default: |
292 | 0 | return 0; |
293 | 0 | } |
294 | 0 | } |
295 | | |
296 | | static int qsm_ready_for_gc(QUIC_STREAM_MAP *qsm, QUIC_STREAM *qs) |
297 | 1.12M | { |
298 | 1.12M | int recv_stream_fully_drained = 0; /* TODO(QUIC FUTURE): Optimisation */ |
299 | | |
300 | | /* |
301 | | * If sstream has no FIN, we auto-reset it at marked-for-deletion time, so |
302 | | * we don't need to worry about that here. |
303 | | */ |
304 | 1.12M | assert(!qs->deleted |
305 | 1.12M | || !ossl_quic_stream_has_send(qs) |
306 | 1.12M | || ossl_quic_stream_send_is_reset(qs) |
307 | 1.12M | || ossl_quic_stream_send_get_final_size(qs, NULL)); |
308 | | |
309 | 1.12M | return qs->deleted |
310 | 12.0k | && (!ossl_quic_stream_has_recv(qs) |
311 | 12.0k | || recv_stream_fully_drained |
312 | 12.0k | || qs->acked_stop_sending) |
313 | 0 | && (!ossl_quic_stream_has_send(qs) |
314 | 0 | || qs->send_state == QUIC_SSTREAM_STATE_DATA_RECVD |
315 | 0 | || qs->send_state == QUIC_SSTREAM_STATE_RESET_RECVD); |
316 | 1.12M | } |
317 | | |
318 | | int ossl_quic_stream_map_is_local_allowed_by_stream_limit(QUIC_STREAM_MAP *qsm, |
319 | | uint64_t stream_ordinal, |
320 | | int is_uni) |
321 | 66.5k | { |
322 | 66.5k | uint64_t stream_limit; |
323 | | |
324 | 66.5k | if (qsm->get_stream_limit_cb == NULL) |
325 | 0 | return 1; |
326 | | |
327 | 66.5k | stream_limit = qsm->get_stream_limit_cb(is_uni, qsm->get_stream_limit_cb_arg); |
328 | 66.5k | return stream_ordinal < stream_limit; |
329 | 66.5k | } |
330 | | |
331 | | void ossl_quic_stream_map_update_state(QUIC_STREAM_MAP *qsm, QUIC_STREAM *s) |
332 | 1.12M | { |
333 | 1.12M | int should_be_active, allowed_by_stream_limit = 1; |
334 | | |
335 | 1.12M | if (ossl_quic_stream_is_server_init(s) == ossl_quic_channel_is_server(qsm->ch)) { |
336 | 58.0k | int is_uni = !ossl_quic_stream_is_bidi(s); |
337 | 58.0k | uint64_t stream_ordinal = s->id >> 2; |
338 | | |
339 | 58.0k | allowed_by_stream_limit |
340 | 58.0k | = ossl_quic_stream_map_is_local_allowed_by_stream_limit(qsm, |
341 | 58.0k | stream_ordinal, |
342 | 58.0k | is_uni); |
343 | 58.0k | } |
344 | | |
345 | 1.12M | if (s->send_state == QUIC_SSTREAM_STATE_DATA_SENT |
346 | 0 | && ossl_quic_sstream_is_totally_acked(s->sstream)) |
347 | 0 | ossl_quic_stream_map_notify_totally_acked(qsm, s); |
348 | 1.12M | else if (s->shutdown_flush |
349 | 0 | && s->send_state == QUIC_SSTREAM_STATE_SEND |
350 | 0 | && ossl_quic_sstream_is_totally_acked(s->sstream)) |
351 | 0 | shutdown_flush_done(qsm, s); |
352 | | |
353 | 1.12M | if (!s->ready_for_gc) { |
354 | 1.12M | s->ready_for_gc = qsm_ready_for_gc(qsm, s); |
355 | 1.12M | if (s->ready_for_gc) |
356 | 0 | list_insert_tail(&qsm->ready_for_gc_list, &s->ready_for_gc_node); |
357 | 1.12M | } |
358 | | |
359 | 1.12M | should_be_active |
360 | 1.12M | = allowed_by_stream_limit |
361 | 1.12M | && !s->ready_for_gc |
362 | 1.12M | && ((ossl_quic_stream_has_recv(s) |
363 | 1.12M | && !ossl_quic_stream_recv_is_reset(s) |
364 | 1.10M | && (s->recv_state == QUIC_RSTREAM_STATE_RECV |
365 | 1.03M | && (s->want_max_stream_data |
366 | 1.02M | || ossl_quic_rxfc_has_cwm_changed(&s->rxfc, 0)))) |
367 | 1.12M | || s->want_stop_sending |
368 | 1.09M | || s->want_reset_stream |
369 | 1.06M | || (!s->peer_stop_sending && stream_has_data_to_send(s))); |
370 | | |
371 | 1.12M | if (should_be_active) |
372 | 83.4k | stream_map_mark_active(qsm, s); |
373 | 1.04M | else |
374 | 1.04M | stream_map_mark_inactive(qsm, s); |
375 | 1.12M | } |
376 | | |
377 | | /* |
378 | | * Stream Send Part State Management |
379 | | * ================================= |
380 | | */ |
381 | | |
382 | | int ossl_quic_stream_map_ensure_send_part_id(QUIC_STREAM_MAP *qsm, |
383 | | QUIC_STREAM *qs) |
384 | 17.0k | { |
385 | 17.0k | switch (qs->send_state) { |
386 | 0 | case QUIC_SSTREAM_STATE_NONE: |
387 | | /* Stream without send part - caller error. */ |
388 | 0 | return 0; |
389 | | |
390 | 17.0k | case QUIC_SSTREAM_STATE_READY: |
391 | | /* |
392 | | * We always allocate a stream ID upfront, so we don't need to do it |
393 | | * here. |
394 | | */ |
395 | 17.0k | qs->send_state = QUIC_SSTREAM_STATE_SEND; |
396 | 17.0k | return 1; |
397 | | |
398 | 0 | default: |
399 | | /* Nothing to do. */ |
400 | 0 | return 1; |
401 | 17.0k | } |
402 | 17.0k | } |
403 | | |
404 | | int ossl_quic_stream_map_notify_all_data_sent(QUIC_STREAM_MAP *qsm, |
405 | | QUIC_STREAM *qs) |
406 | 0 | { |
407 | 0 | switch (qs->send_state) { |
408 | 0 | default: |
409 | | /* Wrong state - caller error. */ |
410 | 0 | case QUIC_SSTREAM_STATE_NONE: |
411 | | /* Stream without send part - caller error. */ |
412 | 0 | return 0; |
413 | | |
414 | 0 | case QUIC_SSTREAM_STATE_SEND: |
415 | 0 | if (!ossl_quic_sstream_get_final_size(qs->sstream, &qs->send_final_size)) |
416 | 0 | return 0; |
417 | | |
418 | 0 | qs->send_state = QUIC_SSTREAM_STATE_DATA_SENT; |
419 | 0 | return 1; |
420 | 0 | } |
421 | 0 | } |
422 | | |
423 | | static void shutdown_flush_done(QUIC_STREAM_MAP *qsm, QUIC_STREAM *qs) |
424 | 6.07k | { |
425 | 6.07k | if (!qs->shutdown_flush) |
426 | 6.07k | return; |
427 | | |
428 | 6.07k | assert(qsm->num_shutdown_flush > 0); |
429 | 0 | qs->shutdown_flush = 0; |
430 | 0 | --qsm->num_shutdown_flush; |
431 | | |
432 | | /* |
433 | | * when num_shutdown_flush becomes zero we need to poke |
434 | | * SSL_poll() it's time to poke to SSL_shutdown() to proceed |
435 | | * with shutdown process as all streams are gone (flushed). |
436 | | */ |
437 | 0 | if (qsm->num_shutdown_flush == 0) |
438 | 0 | ossl_quic_channel_notify_flush_done(qsm->ch); |
439 | 0 | } |
440 | | |
441 | | int ossl_quic_stream_map_notify_totally_acked(QUIC_STREAM_MAP *qsm, |
442 | | QUIC_STREAM *qs) |
443 | 0 | { |
444 | 0 | switch (qs->send_state) { |
445 | 0 | default: |
446 | | /* Wrong state - caller error. */ |
447 | 0 | case QUIC_SSTREAM_STATE_NONE: |
448 | | /* Stream without send part - caller error. */ |
449 | 0 | return 0; |
450 | | |
451 | 0 | case QUIC_SSTREAM_STATE_DATA_SENT: |
452 | 0 | qs->send_state = QUIC_SSTREAM_STATE_DATA_RECVD; |
453 | | /* |
454 | | * Remember final size in case SSL_get_stream_write_state() |
455 | | * gets called. |
456 | | */ |
457 | 0 | qs->have_final_size = ossl_quic_sstream_get_final_size(qs->sstream, |
458 | 0 | NULL); |
459 | | |
460 | | /* We no longer need a QUIC_SSTREAM in this state. */ |
461 | 0 | ossl_quic_sstream_free(qs->sstream); |
462 | 0 | qs->sstream = NULL; |
463 | |
|
464 | 0 | shutdown_flush_done(qsm, qs); |
465 | 0 | return 1; |
466 | 0 | } |
467 | 0 | } |
468 | | |
469 | | int ossl_quic_stream_map_reset_stream_send_part(QUIC_STREAM_MAP *qsm, |
470 | | QUIC_STREAM *qs, |
471 | | uint64_t aec) |
472 | 94.4k | { |
473 | 94.4k | switch (qs->send_state) { |
474 | 0 | default: |
475 | 0 | case QUIC_SSTREAM_STATE_NONE: |
476 | | /* |
477 | | * RESET_STREAM pertains to sending part only, so we cannot reset a |
478 | | * receive-only stream. |
479 | | */ |
480 | 0 | case QUIC_SSTREAM_STATE_DATA_RECVD: |
481 | | /* |
482 | | * RFC 9000 s. 3.3: A sender MUST NOT [...] send RESET_STREAM from a |
483 | | * terminal state. If the stream has already finished normally and the |
484 | | * peer has acknowledged this, we cannot reset it. |
485 | | */ |
486 | 0 | return 0; |
487 | | |
488 | 13.8k | case QUIC_SSTREAM_STATE_READY: |
489 | 13.8k | if (!ossl_quic_stream_map_ensure_send_part_id(qsm, qs)) |
490 | 0 | return 0; |
491 | | |
492 | | /* FALLTHROUGH */ |
493 | 17.0k | case QUIC_SSTREAM_STATE_SEND: |
494 | | /* |
495 | | * If we already have a final size (e.g. because we are coming from |
496 | | * DATA_SENT), we have to be consistent with that, so don't change it. |
497 | | * If we don't already have a final size, determine a final size value. |
498 | | * This is the value which we will end up using for a RESET_STREAM frame |
499 | | * for flow control purposes. We could send the stream size (total |
500 | | * number of bytes appended to QUIC_SSTREAM by the application), but it |
501 | | * is in our interest to exclude any bytes we have not actually |
502 | | * transmitted yet, to avoid unnecessarily consuming flow control |
503 | | * credit. We can get this from the TXFC. |
504 | | */ |
505 | 17.0k | qs->send_final_size = ossl_quic_txfc_get_swm(&qs->txfc); |
506 | | |
507 | | /* FALLTHROUGH */ |
508 | 17.0k | case QUIC_SSTREAM_STATE_DATA_SENT: |
509 | 17.0k | qs->reset_stream_aec = aec; |
510 | 17.0k | qs->want_reset_stream = 1; |
511 | 17.0k | qs->send_state = QUIC_SSTREAM_STATE_RESET_SENT; |
512 | | |
513 | 17.0k | ossl_quic_sstream_free(qs->sstream); |
514 | 17.0k | qs->sstream = NULL; |
515 | | |
516 | 17.0k | shutdown_flush_done(qsm, qs); |
517 | 17.0k | ossl_quic_stream_map_update_state(qsm, qs); |
518 | 17.0k | return 1; |
519 | | |
520 | 77.3k | case QUIC_SSTREAM_STATE_RESET_SENT: |
521 | 77.3k | case QUIC_SSTREAM_STATE_RESET_RECVD: |
522 | | /* |
523 | | * Idempotent - no-op. In any case, do not send RESET_STREAM again - as |
524 | | * mentioned, we must not send it from a terminal state. |
525 | | */ |
526 | 77.3k | return 1; |
527 | 94.4k | } |
528 | 94.4k | } |
529 | | |
530 | | int ossl_quic_stream_map_notify_reset_stream_acked(QUIC_STREAM_MAP *qsm, |
531 | | QUIC_STREAM *qs) |
532 | 0 | { |
533 | 0 | switch (qs->send_state) { |
534 | 0 | default: |
535 | | /* Wrong state - caller error. */ |
536 | 0 | case QUIC_SSTREAM_STATE_NONE: |
537 | | /* Stream without send part - caller error. */ |
538 | 0 | return 0; |
539 | | |
540 | 0 | case QUIC_SSTREAM_STATE_RESET_SENT: |
541 | 0 | qs->send_state = QUIC_SSTREAM_STATE_RESET_RECVD; |
542 | 0 | return 1; |
543 | | |
544 | 0 | case QUIC_SSTREAM_STATE_RESET_RECVD: |
545 | | /* Already in the correct state. */ |
546 | 0 | return 1; |
547 | 0 | } |
548 | 0 | } |
549 | | |
550 | | /* |
551 | | * Stream Receive Part State Management |
552 | | * ==================================== |
553 | | */ |
554 | | |
555 | | int ossl_quic_stream_map_notify_size_known_recv_part(QUIC_STREAM_MAP *qsm, |
556 | | QUIC_STREAM *qs, |
557 | | uint64_t final_size) |
558 | 7.86k | { |
559 | 7.86k | switch (qs->recv_state) { |
560 | 0 | default: |
561 | | /* Wrong state - caller error. */ |
562 | 0 | case QUIC_RSTREAM_STATE_NONE: |
563 | | /* Stream without receive part - caller error. */ |
564 | 0 | return 0; |
565 | | |
566 | 7.86k | case QUIC_RSTREAM_STATE_RECV: |
567 | 7.86k | qs->recv_state = QUIC_RSTREAM_STATE_SIZE_KNOWN; |
568 | 7.86k | return 1; |
569 | 7.86k | } |
570 | 7.86k | } |
571 | | |
572 | | int ossl_quic_stream_map_notify_totally_received(QUIC_STREAM_MAP *qsm, |
573 | | QUIC_STREAM *qs) |
574 | 3.10k | { |
575 | 3.10k | switch (qs->recv_state) { |
576 | 0 | default: |
577 | | /* Wrong state - caller error. */ |
578 | 0 | case QUIC_RSTREAM_STATE_NONE: |
579 | | /* Stream without receive part - caller error. */ |
580 | 0 | return 0; |
581 | | |
582 | 3.10k | case QUIC_RSTREAM_STATE_SIZE_KNOWN: |
583 | 3.10k | qs->recv_state = QUIC_RSTREAM_STATE_DATA_RECVD; |
584 | 3.10k | qs->want_stop_sending = 0; |
585 | 3.10k | return 1; |
586 | 3.10k | } |
587 | 3.10k | } |
588 | | |
589 | | int ossl_quic_stream_map_notify_totally_read(QUIC_STREAM_MAP *qsm, |
590 | | QUIC_STREAM *qs) |
591 | 649 | { |
592 | 649 | switch (qs->recv_state) { |
593 | 0 | default: |
594 | | /* Wrong state - caller error. */ |
595 | 0 | case QUIC_RSTREAM_STATE_NONE: |
596 | | /* Stream without receive part - caller error. */ |
597 | 0 | return 0; |
598 | | |
599 | 649 | case QUIC_RSTREAM_STATE_DATA_RECVD: |
600 | 649 | qs->recv_state = QUIC_RSTREAM_STATE_DATA_READ; |
601 | | |
602 | | /* QUIC_RSTREAM is no longer needed */ |
603 | 649 | ossl_quic_rstream_free(qs->rstream); |
604 | 649 | qs->rstream = NULL; |
605 | 649 | return 1; |
606 | 649 | } |
607 | 649 | } |
608 | | |
609 | | int ossl_quic_stream_map_notify_reset_recv_part(QUIC_STREAM_MAP *qsm, |
610 | | QUIC_STREAM *qs, |
611 | | uint64_t app_error_code, |
612 | | uint64_t final_size) |
613 | 6.95k | { |
614 | 6.95k | uint64_t prev_final_size; |
615 | | |
616 | 6.95k | switch (qs->recv_state) { |
617 | 0 | default: |
618 | 0 | case QUIC_RSTREAM_STATE_NONE: |
619 | | /* Stream without receive part - caller error. */ |
620 | 0 | return 0; |
621 | | |
622 | 1.74k | case QUIC_RSTREAM_STATE_RECV: |
623 | 1.75k | case QUIC_RSTREAM_STATE_SIZE_KNOWN: |
624 | 1.77k | case QUIC_RSTREAM_STATE_DATA_RECVD: |
625 | 1.77k | if (ossl_quic_stream_recv_get_final_size(qs, &prev_final_size) |
626 | 30 | && prev_final_size != final_size) |
627 | | /* Cannot change previous final size. */ |
628 | 0 | return 0; |
629 | | |
630 | 1.77k | qs->recv_state = QUIC_RSTREAM_STATE_RESET_RECVD; |
631 | 1.77k | qs->peer_reset_stream_aec = app_error_code; |
632 | | |
633 | | /* RFC 9000 s. 3.3: No point sending STOP_SENDING if already reset. */ |
634 | 1.77k | qs->want_stop_sending = 0; |
635 | | |
636 | | /* QUIC_RSTREAM is no longer needed */ |
637 | 1.77k | ossl_quic_rstream_free(qs->rstream); |
638 | 1.77k | qs->rstream = NULL; |
639 | | |
640 | 1.77k | ossl_quic_stream_map_update_state(qsm, qs); |
641 | 1.77k | return 1; |
642 | | |
643 | 0 | case QUIC_RSTREAM_STATE_DATA_READ: |
644 | | /* |
645 | | * If we already retired the FIN to the application this is moot |
646 | | * - just ignore. |
647 | | */ |
648 | 5.18k | case QUIC_RSTREAM_STATE_RESET_RECVD: |
649 | 5.18k | case QUIC_RSTREAM_STATE_RESET_READ: |
650 | | /* Could be a reordered/retransmitted frame - just ignore. */ |
651 | 5.18k | return 1; |
652 | 6.95k | } |
653 | 6.95k | } |
654 | | |
655 | | int ossl_quic_stream_map_notify_app_read_reset_recv_part(QUIC_STREAM_MAP *qsm, |
656 | | QUIC_STREAM *qs) |
657 | 169 | { |
658 | 169 | switch (qs->recv_state) { |
659 | 0 | default: |
660 | | /* Wrong state - caller error. */ |
661 | 0 | case QUIC_RSTREAM_STATE_NONE: |
662 | | /* Stream without receive part - caller error. */ |
663 | 0 | return 0; |
664 | | |
665 | 169 | case QUIC_RSTREAM_STATE_RESET_RECVD: |
666 | 169 | qs->recv_state = QUIC_RSTREAM_STATE_RESET_READ; |
667 | 169 | return 1; |
668 | 169 | } |
669 | 169 | } |
670 | | |
671 | | int ossl_quic_stream_map_stop_sending_recv_part(QUIC_STREAM_MAP *qsm, |
672 | | QUIC_STREAM *qs, |
673 | | uint64_t aec) |
674 | 10.9k | { |
675 | 10.9k | if (qs->stop_sending) |
676 | 0 | return 0; |
677 | | |
678 | 10.9k | switch (qs->recv_state) { |
679 | 0 | default: |
680 | 0 | case QUIC_RSTREAM_STATE_NONE: |
681 | | /* Send-only stream, so this makes no sense. */ |
682 | 0 | case QUIC_RSTREAM_STATE_DATA_RECVD: |
683 | 0 | case QUIC_RSTREAM_STATE_DATA_READ: |
684 | | /* |
685 | | * Not really any point in STOP_SENDING if we already received all data. |
686 | | */ |
687 | 0 | case QUIC_RSTREAM_STATE_RESET_RECVD: |
688 | 0 | case QUIC_RSTREAM_STATE_RESET_READ: |
689 | | /* |
690 | | * RFC 9000 s. 3.5: "STOP_SENDING SHOULD only be sent for a stream that |
691 | | * has not been reset by the peer." |
692 | | * |
693 | | * No point in STOP_SENDING if the peer already reset their send part. |
694 | | */ |
695 | 0 | return 0; |
696 | | |
697 | 9.22k | case QUIC_RSTREAM_STATE_RECV: |
698 | 10.9k | case QUIC_RSTREAM_STATE_SIZE_KNOWN: |
699 | | /* |
700 | | * RFC 9000 s. 3.5: "If the stream is in the Recv or Size Known state, |
701 | | * the transport SHOULD signal this by sending a STOP_SENDING frame to |
702 | | * prompt closure of the stream in the opposite direction." |
703 | | * |
704 | | * Note that it does make sense to send STOP_SENDING for a receive part |
705 | | * of a stream which has a known size (because we have received a FIN) |
706 | | * but which still has other (previous) stream data yet to be received. |
707 | | */ |
708 | 10.9k | break; |
709 | 10.9k | } |
710 | | |
711 | 10.9k | qs->stop_sending = 1; |
712 | 10.9k | qs->stop_sending_aec = aec; |
713 | 10.9k | return ossl_quic_stream_map_schedule_stop_sending(qsm, qs); |
714 | 10.9k | } |
715 | | |
716 | | /* Called to mark STOP_SENDING for generation, or regeneration after loss. */ |
717 | | int ossl_quic_stream_map_schedule_stop_sending(QUIC_STREAM_MAP *qsm, QUIC_STREAM *qs) |
718 | 10.9k | { |
719 | 10.9k | if (!qs->stop_sending) |
720 | 0 | return 0; |
721 | | |
722 | | /* |
723 | | * Ignore the call as a no-op if already scheduled, or in a state |
724 | | * where it makes no sense to send STOP_SENDING. |
725 | | */ |
726 | 10.9k | if (qs->want_stop_sending) |
727 | 0 | return 1; |
728 | | |
729 | 10.9k | switch (qs->recv_state) { |
730 | 0 | default: |
731 | 0 | return 1; /* ignore */ |
732 | 9.22k | case QUIC_RSTREAM_STATE_RECV: |
733 | 10.9k | case QUIC_RSTREAM_STATE_SIZE_KNOWN: |
734 | | /* |
735 | | * RFC 9000 s. 3.5: "An endpoint is expected to send another |
736 | | * STOP_SENDING frame if a packet containing a previous STOP_SENDING is |
737 | | * lost. However, once either all stream data or a RESET_STREAM frame |
738 | | * has been received for the stream -- that is, the stream is in any |
739 | | * state other than "Recv" or "Size Known" -- sending a STOP_SENDING |
740 | | * frame is unnecessary." |
741 | | */ |
742 | 10.9k | break; |
743 | 10.9k | } |
744 | | |
745 | 10.9k | qs->want_stop_sending = 1; |
746 | 10.9k | ossl_quic_stream_map_update_state(qsm, qs); |
747 | 10.9k | return 1; |
748 | 10.9k | } |
749 | | |
750 | | QUIC_STREAM *ossl_quic_stream_map_peek_accept_queue(QUIC_STREAM_MAP *qsm) |
751 | 461 | { |
752 | 461 | return accept_head(&qsm->accept_list); |
753 | 461 | } |
754 | | |
755 | | QUIC_STREAM *ossl_quic_stream_map_find_in_accept_queue(QUIC_STREAM_MAP *qsm, |
756 | | int is_uni) |
757 | 0 | { |
758 | 0 | QUIC_STREAM *qs; |
759 | |
|
760 | 0 | if (ossl_quic_stream_map_get_accept_queue_len(qsm, is_uni) == 0) |
761 | 0 | return NULL; |
762 | | |
763 | 0 | qs = ossl_quic_stream_map_peek_accept_queue(qsm); |
764 | 0 | while (qs != NULL) { |
765 | 0 | if ((is_uni && !ossl_quic_stream_is_bidi(qs)) |
766 | 0 | || (!is_uni && ossl_quic_stream_is_bidi(qs))) |
767 | 0 | break; |
768 | 0 | qs = accept_next(&qsm->accept_list, qs); |
769 | 0 | } |
770 | 0 | return qs; |
771 | 0 | } |
772 | | |
773 | | void ossl_quic_stream_map_push_accept_queue(QUIC_STREAM_MAP *qsm, |
774 | | QUIC_STREAM *s) |
775 | 138k | { |
776 | 138k | list_insert_tail(&qsm->accept_list, &s->accept_node); |
777 | 138k | if (ossl_quic_stream_is_bidi(s)) |
778 | 37.6k | ++qsm->num_accept_bidi; |
779 | 100k | else |
780 | 100k | ++qsm->num_accept_uni; |
781 | 138k | } |
782 | | |
783 | | static QUIC_RXFC *qsm_get_max_streams_rxfc(QUIC_STREAM_MAP *qsm, QUIC_STREAM *s) |
784 | 6.99k | { |
785 | 6.99k | return ossl_quic_stream_is_bidi(s) |
786 | 6.99k | ? qsm->max_streams_bidi_rxfc |
787 | 6.99k | : qsm->max_streams_uni_rxfc; |
788 | 6.99k | } |
789 | | |
790 | | void ossl_quic_stream_map_remove_from_accept_queue(QUIC_STREAM_MAP *qsm, |
791 | | QUIC_STREAM *s, |
792 | | OSSL_TIME rtt) |
793 | 6.99k | { |
794 | 6.99k | QUIC_RXFC *max_streams_rxfc; |
795 | | |
796 | 6.99k | list_remove(&qsm->accept_list, &s->accept_node); |
797 | 6.99k | if (ossl_quic_stream_is_bidi(s)) |
798 | 6.42k | --qsm->num_accept_bidi; |
799 | 573 | else |
800 | 573 | --qsm->num_accept_uni; |
801 | | |
802 | 6.99k | if ((max_streams_rxfc = qsm_get_max_streams_rxfc(qsm, s)) != NULL) |
803 | 6.99k | (void)ossl_quic_rxfc_on_retire(max_streams_rxfc, 1, rtt); |
804 | 6.99k | } |
805 | | |
806 | | size_t ossl_quic_stream_map_get_accept_queue_len(QUIC_STREAM_MAP *qsm, int is_uni) |
807 | 22.7k | { |
808 | 22.7k | return is_uni ? qsm->num_accept_uni : qsm->num_accept_bidi; |
809 | 22.7k | } |
810 | | |
811 | | size_t ossl_quic_stream_map_get_total_accept_queue_len(QUIC_STREAM_MAP *qsm) |
812 | 11.3k | { |
813 | 11.3k | return ossl_quic_stream_map_get_accept_queue_len(qsm, /*is_uni=*/0) |
814 | 11.3k | + ossl_quic_stream_map_get_accept_queue_len(qsm, /*is_uni=*/1); |
815 | 11.3k | } |
816 | | |
817 | | void ossl_quic_stream_map_gc(QUIC_STREAM_MAP *qsm) |
818 | 71.5M | { |
819 | 71.5M | QUIC_STREAM *qs; |
820 | | |
821 | 71.5M | while ((qs = ready_for_gc_head(&qsm->ready_for_gc_list)) != NULL) { |
822 | 0 | ossl_quic_stream_map_release(qsm, qs); |
823 | 0 | } |
824 | 71.5M | } |
825 | | |
826 | | static int eligible_for_shutdown_flush(QUIC_STREAM *qs) |
827 | 0 | { |
828 | | /* |
829 | | * We only care about servicing the send part of a stream (if any) during |
830 | | * shutdown flush. We make sure we flush a stream if it is either |
831 | | * non-terminated or was terminated normally such as via |
832 | | * SSL_stream_conclude. A stream which was terminated via a reset is not |
833 | | * flushed, and we will have thrown away the send buffer in that case |
834 | | * anyway. |
835 | | */ |
836 | 0 | switch (qs->send_state) { |
837 | 0 | case QUIC_SSTREAM_STATE_SEND: |
838 | 0 | case QUIC_SSTREAM_STATE_DATA_SENT: |
839 | 0 | return !ossl_quic_sstream_is_totally_acked(qs->sstream); |
840 | 0 | default: |
841 | 0 | return 0; |
842 | 0 | } |
843 | 0 | } |
844 | | |
845 | | static void begin_shutdown_flush_each(QUIC_STREAM *qs, void *arg) |
846 | 0 | { |
847 | 0 | QUIC_STREAM_MAP *qsm = arg; |
848 | |
|
849 | 0 | if (!eligible_for_shutdown_flush(qs) || qs->shutdown_flush) |
850 | 0 | return; |
851 | | |
852 | 0 | qs->shutdown_flush = 1; |
853 | 0 | ++qsm->num_shutdown_flush; |
854 | 0 | } |
855 | | |
856 | | void ossl_quic_stream_map_begin_shutdown_flush(QUIC_STREAM_MAP *qsm) |
857 | 0 | { |
858 | 0 | qsm->num_shutdown_flush = 0; |
859 | |
|
860 | 0 | ossl_quic_stream_map_visit(qsm, begin_shutdown_flush_each, qsm); |
861 | 0 | } |
862 | | |
863 | | int ossl_quic_stream_map_is_shutdown_flush_finished(QUIC_STREAM_MAP *qsm) |
864 | 0 | { |
865 | 0 | return qsm->num_shutdown_flush == 0; |
866 | 0 | } |
867 | | |
868 | | /* |
869 | | * QUIC Stream Iterator |
870 | | * ==================== |
871 | | */ |
872 | | void ossl_quic_stream_iter_init(QUIC_STREAM_ITER *it, QUIC_STREAM_MAP *qsm, |
873 | | int advance_rr) |
874 | 8.07M | { |
875 | 8.07M | it->qsm = qsm; |
876 | 8.07M | it->stream = it->first_stream = qsm->rr_cur; |
877 | 8.07M | if (advance_rr && it->stream != NULL |
878 | 14.9k | && ++qsm->rr_counter >= qsm->rr_stepping) { |
879 | 14.9k | qsm->rr_counter = 0; |
880 | 14.9k | qsm->rr_cur = active_next(&qsm->active_list, qsm->rr_cur); |
881 | 14.9k | } |
882 | 8.07M | } |
883 | | |
884 | | void ossl_quic_stream_iter_next(QUIC_STREAM_ITER *it) |
885 | 18.2k | { |
886 | 18.2k | if (it->stream == NULL) |
887 | 0 | return; |
888 | | |
889 | 18.2k | it->stream = active_next(&it->qsm->active_list, it->stream); |
890 | 18.2k | if (it->stream == it->first_stream) |
891 | 14.8k | it->stream = NULL; |
892 | 18.2k | } |