Skip to main content

h2/proto/streams/
streams.rs

1use super::recv::RecvHeaderBlockError;
2use super::store::{self, Entry, Resolve, Store};
3use super::{Buffer, BufferStatus, Config, Counts, Prioritized, Recv, Send, Stream, StreamId};
4use crate::codec::{Codec, SendError, UserError};
5use crate::ext::Protocol;
6use crate::frame::{self, Frame, Reason};
7use crate::proto::{peer, Error, Initiator, Open, Peer, WindowSize};
8use crate::{client, proto, server};
9
10use bytes::{Buf, Bytes};
11use http::{HeaderMap, Request, Response};
12use std::task::{Context, Poll, Waker};
13use tokio::io::AsyncWrite;
14
15use std::sync::{Arc, Mutex};
16use std::{fmt, io};
17
18#[derive(Debug)]
19pub(crate) struct Streams<B, P>
20where
21    P: Peer,
22{
23    /// Holds most of the connection and stream related state for processing
24    /// HTTP/2 frames associated with streams.
25    inner: Arc<Mutex<Inner>>,
26
27    /// This is the queue of frames to be written to the wire. This is split out
28    /// to avoid requiring a `B` generic on all public API types even if `B` is
29    /// not technically required.
30    ///
31    /// Currently, splitting this out requires a second `Arc` + `Mutex`.
32    /// However, it should be possible to avoid this duplication with a little
33    /// bit of unsafe code. This optimization has been postponed until it has
34    /// been shown to be necessary.
35    send_buffer: Arc<SendBuffer<B>>,
36
37    _p: ::std::marker::PhantomData<P>,
38}
39
40// Like `Streams` but with a `peer::Dyn` field instead of a static `P: Peer` type parameter.
41// Ensures that the methods only get one instantiation, instead of two (client and server)
42#[derive(Debug)]
43pub(crate) struct DynStreams<'a, B> {
44    inner: &'a Mutex<Inner>,
45
46    send_buffer: &'a SendBuffer<B>,
47
48    peer: peer::Dyn,
49}
50
51/// Reference to the stream state
52#[derive(Debug)]
53pub(crate) struct StreamRef<B> {
54    opaque: OpaqueStreamRef,
55    send_buffer: Arc<SendBuffer<B>>,
56}
57
58/// Reference to the stream state that hides the send data chunk generic
59pub(crate) struct OpaqueStreamRef {
60    inner: Arc<Mutex<Inner>>,
61    key: store::Key,
62}
63
64/// Fields needed to manage state related to managing the set of streams. This
65/// is mostly split out to make ownership happy.
66///
67/// TODO: better name
68#[derive(Debug)]
69struct Inner {
70    /// Tracks send & recv stream concurrency.
71    counts: Counts,
72
73    /// Connection level state and performs actions on streams
74    actions: Actions,
75
76    /// Stores stream state
77    store: Store,
78
79    /// The number of stream refs to this shared state.
80    refs: usize,
81}
82
83#[derive(Debug)]
84struct Actions {
85    /// Manages state transitions initiated by receiving frames
86    recv: Recv,
87
88    /// Manages state transitions initiated by sending frames
89    send: Send,
90
91    /// Task that calls `poll_complete`.
92    task: Option<Waker>,
93
94    /// If the connection errors, a copy is kept for any StreamRefs.
95    conn_error: Option<proto::Error>,
96}
97
98/// Contains the buffer of frames to be written to the wire.
99#[derive(Debug)]
100struct SendBuffer<B> {
101    inner: Mutex<Buffer<Frame<B>>>,
102}
103
104// ===== impl Streams =====
105
106impl<B, P> Streams<B, P>
107where
108    B: Buf,
109    P: Peer,
110{
111    pub fn new(config: Config) -> Self {
112        let peer = P::r#dyn();
113
114        Streams {
115            inner: Inner::new(peer, config),
116            send_buffer: Arc::new(SendBuffer::new()),
117            _p: ::std::marker::PhantomData,
118        }
119    }
120
121    pub fn set_target_connection_window_size(&mut self, size: WindowSize) -> Result<(), Reason> {
122        let mut me = self.inner.lock().unwrap();
123        let me = &mut *me;
124
125        me.actions
126            .recv
127            .set_target_connection_window(size, &mut me.actions.task)
128    }
129
130    pub fn next_incoming(&mut self) -> Option<StreamRef<B>> {
131        let mut me = self.inner.lock().unwrap();
132        let me = &mut *me;
133        me.actions.recv.next_incoming(&mut me.store).map(|key| {
134            let stream = &mut me.store.resolve(key);
135            tracing::trace!(
136                "next_incoming; id={:?}, state={:?}",
137                stream.id,
138                stream.state
139            );
140            // TODO: ideally, OpaqueStreamRefs::new would do this, but we're holding
141            // the lock, so it can't.
142            me.refs += 1;
143
144            // Pending-accepted remotely-reset streams are counted.
145            if stream.state.is_remote_reset() {
146                me.counts.dec_num_remote_reset_streams();
147            }
148
149            StreamRef {
150                opaque: OpaqueStreamRef::new(self.inner.clone(), stream),
151                send_buffer: self.send_buffer.clone(),
152            }
153        })
154    }
155
156    pub fn send_pending_refusal<T>(
157        &mut self,
158        cx: &mut Context,
159        dst: &mut Codec<T, Prioritized<B>>,
160    ) -> Poll<io::Result<()>>
161    where
162        T: AsyncWrite + Unpin,
163    {
164        loop {
165            let status = {
166                let mut me = self.inner.lock().unwrap();
167                let me = &mut *me;
168                me.actions.recv.send_pending_refusal(dst)?
169            };
170
171            match status {
172                BufferStatus::Complete => return Poll::Ready(Ok(())),
173                BufferStatus::CodecFull => ready!(dst.poll_ready(cx))?,
174            }
175        }
176    }
177
178    pub fn clear_expired_reset_streams(&mut self) {
179        let mut me = self.inner.lock().unwrap();
180        let me = &mut *me;
181        me.actions
182            .recv
183            .clear_expired_reset_streams(&mut me.store, &mut me.counts);
184    }
185
186    pub fn poll_complete<T>(
187        &mut self,
188        cx: &mut Context,
189        dst: &mut Codec<T, Prioritized<B>>,
190    ) -> Poll<io::Result<()>>
191    where
192        T: AsyncWrite + Unpin,
193    {
194        loop {
195            // Make any required socket progress before taking stream locks.
196            ready!(dst.poll_ready(cx))?;
197
198            let status = {
199                let mut me = self.inner.lock().unwrap();
200                let status = me.buffer_pending(&self.send_buffer, dst)?;
201
202                // Register the task while holding the same lock used to
203                // observe that all pending frames have been buffered. A
204                // producer that queues another frame while the codec is being
205                // flushed will then take and wake this task.
206                if status == BufferStatus::Complete {
207                    me.actions.task = Some(cx.waker().clone());
208                }
209
210                status
211            };
212
213            match status {
214                BufferStatus::Complete => {}
215                BufferStatus::CodecFull => continue,
216            }
217
218            // Flush any frames staged by `buffer_pending` without holding the
219            // stream-state or send-buffer mutexes.
220            ready!(dst.flush(cx))?;
221
222            let reclaimed = {
223                let mut me = self.inner.lock().unwrap();
224                me.reclaim_written_frame(&self.send_buffer, dst)
225            };
226
227            if !reclaimed {
228                return Poll::Ready(Ok(()));
229            }
230        }
231    }
232
233    pub fn apply_remote_settings(
234        &mut self,
235        frame: &frame::Settings,
236        is_initial: bool,
237    ) -> Result<(), Error> {
238        let mut me = self.inner.lock().unwrap();
239        let me = &mut *me;
240
241        let mut send_buffer = self.send_buffer.inner.lock().unwrap();
242        let send_buffer = &mut *send_buffer;
243
244        me.counts.apply_remote_settings(frame, is_initial);
245
246        me.actions.send.apply_remote_settings(
247            frame,
248            send_buffer,
249            &mut me.store,
250            &mut me.counts,
251            &mut me.actions.task,
252        )
253    }
254
255    pub fn apply_local_settings(&mut self, frame: &frame::Settings) -> Result<(), Error> {
256        let mut me = self.inner.lock().unwrap();
257        let me = &mut *me;
258
259        me.actions.recv.apply_local_settings(frame, &mut me.store)
260    }
261
262    pub fn send_request(
263        &mut self,
264        mut request: Request<()>,
265        end_of_stream: bool,
266        pending: Option<&OpaqueStreamRef>,
267    ) -> Result<(StreamRef<B>, bool), SendError> {
268        use super::stream::ContentLength;
269        use http::Method;
270
271        let protocol = request.extensions_mut().remove::<Protocol>();
272
273        // Clear before taking lock, incase extensions contain a StreamRef.
274        request.extensions_mut().clear();
275
276        // TODO: There is a hazard with assigning a stream ID before the
277        // prioritize layer. If prioritization reorders new streams, this
278        // implicitly closes the earlier stream IDs.
279        //
280        // See: hyperium/h2#11
281        let mut me = self.inner.lock().unwrap();
282        let me = &mut *me;
283
284        let mut send_buffer = self.send_buffer.inner.lock().unwrap();
285        let send_buffer = &mut *send_buffer;
286
287        me.actions.ensure_no_conn_error()?;
288        me.actions.send.ensure_next_stream_id()?;
289
290        // The `pending` argument is provided by the `Client`, and holds
291        // a store `Key` of a `Stream` that may have been not been opened
292        // yet.
293        //
294        // If that stream is still pending, the Client isn't allowed to
295        // queue up another pending stream. They should use `poll_ready`.
296        if let Some(stream) = pending {
297            if me.store.resolve(stream.key).is_pending_open {
298                return Err(UserError::Rejected.into());
299            }
300        }
301
302        if me.counts.peer().is_server() {
303            // Servers cannot open streams. PushPromise must first be reserved.
304            return Err(UserError::UnexpectedFrameType.into());
305        }
306
307        let stream_id = me.actions.send.open()?;
308
309        let mut stream = Stream::new(
310            stream_id,
311            me.actions.send.init_window_sz(),
312            me.actions.recv.init_window_sz(),
313        );
314
315        if *request.method() == Method::HEAD {
316            stream.content_length = ContentLength::Head;
317        }
318
319        // Convert the message
320        let headers =
321            client::Peer::convert_send_message(stream_id, request, protocol, end_of_stream)?;
322
323        let mut stream = me.store.insert(stream.id, stream);
324
325        let sent = me.actions.send.send_headers(
326            headers,
327            send_buffer,
328            &mut stream,
329            &mut me.counts,
330            &mut me.actions.task,
331        );
332
333        // send_headers can return a UserError, if it does,
334        // we should forget about this stream.
335        if let Err(err) = sent {
336            stream.unlink();
337            stream.remove();
338            return Err(err.into());
339        }
340
341        // Given that the stream has been initialized, it should not be in the
342        // closed state.
343        debug_assert!(!stream.state.is_closed());
344
345        // TODO: ideally, OpaqueStreamRefs::new would do this, but we're holding
346        // the lock, so it can't.
347        me.refs += 1;
348
349        let is_full = me.counts.next_send_stream_will_reach_capacity();
350        Ok((
351            StreamRef {
352                opaque: OpaqueStreamRef::new(self.inner.clone(), &mut stream),
353                send_buffer: self.send_buffer.clone(),
354            },
355            is_full,
356        ))
357    }
358
359    pub(crate) fn is_extended_connect_protocol_enabled(&self) -> bool {
360        self.inner
361            .lock()
362            .unwrap()
363            .actions
364            .send
365            .is_extended_connect_protocol_enabled()
366    }
367
368    pub fn current_max_send_streams(&self) -> usize {
369        let me = self.inner.lock().unwrap();
370        me.counts.max_send_streams()
371    }
372
373    pub fn current_max_recv_streams(&self) -> usize {
374        let me = self.inner.lock().unwrap();
375        me.counts.max_recv_streams()
376    }
377}
378
379impl<B> DynStreams<'_, B> {
380    pub fn is_buffer_empty(&self) -> bool {
381        self.send_buffer.is_empty()
382    }
383
384    pub fn is_server(&self) -> bool {
385        self.peer.is_server()
386    }
387
388    pub fn recv_headers(&mut self, frame: frame::Headers) -> Result<(), Error> {
389        let mut me = self.inner.lock().unwrap();
390
391        me.recv_headers(self.peer, self.send_buffer, frame)
392    }
393
394    pub fn recv_data(&mut self, frame: frame::Data) -> Result<(), Error> {
395        let mut me = self.inner.lock().unwrap();
396        me.recv_data(self.peer, self.send_buffer, frame)
397    }
398
399    pub fn recv_reset(&mut self, frame: frame::Reset) -> Result<(), Error> {
400        let mut me = self.inner.lock().unwrap();
401
402        me.recv_reset(self.send_buffer, frame)
403    }
404
405    /// Notify all streams that a connection-level error happened.
406    pub fn handle_error(&mut self, err: proto::Error) -> StreamId {
407        let mut me = self.inner.lock().unwrap();
408        me.handle_error(self.send_buffer, err)
409    }
410
411    pub fn recv_go_away(&mut self, frame: &frame::GoAway) -> Result<(), Error> {
412        let mut me = self.inner.lock().unwrap();
413        me.recv_go_away(self.send_buffer, frame)
414    }
415
416    pub fn last_processed_id(&self) -> StreamId {
417        self.inner.lock().unwrap().actions.recv.last_processed_id()
418    }
419
420    pub fn recv_window_update(&mut self, frame: frame::WindowUpdate) -> Result<(), Error> {
421        let mut me = self.inner.lock().unwrap();
422        me.recv_window_update(self.send_buffer, frame)
423    }
424
425    pub fn recv_push_promise(&mut self, frame: frame::PushPromise) -> Result<(), Error> {
426        let mut me = self.inner.lock().unwrap();
427        me.recv_push_promise(self.send_buffer, frame)
428    }
429
430    pub fn recv_eof(&mut self, clear_pending_accept: bool) -> Result<(), ()> {
431        let mut me = self.inner.lock().map_err(|_| ())?;
432        me.recv_eof(self.send_buffer, clear_pending_accept)
433    }
434
435    pub fn send_reset(
436        &mut self,
437        id: StreamId,
438        reason: Reason,
439    ) -> Result<(), crate::proto::error::GoAway> {
440        let mut me = self.inner.lock().unwrap();
441        me.send_reset(self.send_buffer, id, reason)
442    }
443
444    pub fn send_go_away(&mut self, last_processed_id: StreamId) {
445        let mut me = self.inner.lock().unwrap();
446        me.actions.recv.go_away(last_processed_id);
447    }
448}
449
450impl Inner {
451    fn new(peer: peer::Dyn, config: Config) -> Arc<Mutex<Self>> {
452        Arc::new(Mutex::new(Inner {
453            counts: Counts::new(peer, &config),
454            actions: Actions {
455                recv: Recv::new(peer, &config),
456                send: Send::new(&config),
457                task: None,
458                conn_error: None,
459            },
460            store: Store::new(),
461            refs: 1,
462        }))
463    }
464
465    fn recv_headers<B>(
466        &mut self,
467        peer: peer::Dyn,
468        send_buffer: &SendBuffer<B>,
469        frame: frame::Headers,
470    ) -> Result<(), Error> {
471        let id = frame.stream_id();
472
473        // The GOAWAY process has begun. All streams with a greater ID than
474        // specified as part of GOAWAY should be ignored.
475        if id > self.actions.recv.max_stream_id() {
476            tracing::trace!(
477                "id ({:?}) > max_stream_id ({:?}), ignoring HEADERS",
478                id,
479                self.actions.recv.max_stream_id()
480            );
481            return Ok(());
482        }
483
484        let key = match self.store.find_entry(id) {
485            Entry::Occupied(e) => e.key(),
486            Entry::Vacant(e) => {
487                // Client: it's possible to send a request, and then send
488                // a RST_STREAM while the response HEADERS were in transit.
489                //
490                // Server: we can't reset a stream before having received
491                // the request headers, so don't allow.
492                if !peer.is_server() {
493                    // This may be response headers for a stream we've already
494                    // forgotten about...
495                    if self.actions.may_have_forgotten_stream(peer, id) {
496                        tracing::debug!(
497                            "recv_headers for old stream={:?}, sending STREAM_CLOSED",
498                            id,
499                        );
500                        return Err(Error::library_reset(id, Reason::STREAM_CLOSED));
501                    }
502                }
503
504                match self
505                    .actions
506                    .recv
507                    .open(id, Open::Headers, &mut self.counts)?
508                {
509                    Some(stream_id) => {
510                        let stream = Stream::new(
511                            stream_id,
512                            self.actions.send.init_window_sz(),
513                            self.actions.recv.init_window_sz(),
514                        );
515
516                        e.insert(stream)
517                    }
518                    None => return Ok(()),
519                }
520            }
521        };
522
523        let stream = self.store.resolve(key);
524
525        if stream.is_pending_open {
526            proto_err!(conn: "recv_headers: received frame on idle stream {:?}", id);
527            return Err(Error::library_go_away(Reason::PROTOCOL_ERROR));
528        }
529
530        if stream.state.is_local_error() {
531            // Locally reset streams must ignore frames "for some time".
532            // This is because the remote may have sent trailers before
533            // receiving the RST_STREAM frame.
534            tracing::trace!("recv_headers; ignoring trailers on {:?}", stream.id);
535            return Ok(());
536        }
537
538        let actions = &mut self.actions;
539        let mut send_buffer = send_buffer.inner.lock().unwrap();
540        let send_buffer = &mut *send_buffer;
541
542        self.counts.transition(stream, |counts, stream| {
543            tracing::trace!(
544                "recv_headers; stream={:?}; state={:?}",
545                stream.id,
546                stream.state
547            );
548
549            let res = if stream.state.is_recv_headers() {
550                match actions.recv.recv_headers(frame, stream, counts) {
551                    Ok(()) => Ok(()),
552                    Err(RecvHeaderBlockError::Oversize(resp)) => {
553                        if let Some(resp) = resp {
554                            let sent = actions.send.send_headers(
555                                resp, send_buffer, stream, counts, &mut actions.task);
556                            debug_assert!(sent.is_ok(), "oversize response should not fail");
557
558                            actions.send.schedule_implicit_reset(
559                                stream,
560                                Reason::PROTOCOL_ERROR,
561                                counts,
562                                &mut actions.task);
563
564                            actions.recv.enqueue_reset_expiration(stream, counts);
565
566                            Ok(())
567                        } else {
568                            Err(Error::library_reset(stream.id, Reason::PROTOCOL_ERROR))
569                        }
570                    },
571                    Err(RecvHeaderBlockError::State(err)) => Err(err),
572                }
573            } else {
574                if !frame.is_end_stream() {
575                    // Receiving trailers that don't set EOS is a "malformed"
576                    // message. Malformed messages are a stream error.
577                    proto_err!(stream: "recv_headers: trailers frame was not EOS; stream={:?}", stream.id);
578                    return Err(Error::library_reset(stream.id, Reason::PROTOCOL_ERROR));
579                }
580
581                actions.recv.recv_trailers(frame, stream)
582            };
583
584            actions.reset_on_recv_stream_err(send_buffer, stream, counts, res)
585        })
586    }
587
588    fn recv_data<B>(
589        &mut self,
590        peer: peer::Dyn,
591        send_buffer: &SendBuffer<B>,
592        frame: frame::Data,
593    ) -> Result<(), Error> {
594        let id = frame.stream_id();
595
596        let stream = match self.store.find_mut(&id) {
597            Some(stream) => stream,
598            None => {
599                // The GOAWAY process has begun. All streams with a greater ID
600                // than specified as part of GOAWAY should be ignored.
601                if id > self.actions.recv.max_stream_id() {
602                    tracing::trace!(
603                        "id ({:?}) > max_stream_id ({:?}), ignoring DATA",
604                        id,
605                        self.actions.recv.max_stream_id()
606                    );
607
608                    // We still need to account for connection-level flow control.
609                    let sz = frame.flow_controlled_len();
610                    assert!(sz <= super::MAX_WINDOW_SIZE as usize);
611                    let sz = sz as WindowSize;
612                    self.actions.recv.ignore_data(sz)?;
613
614                    return Ok(());
615                }
616
617                if self.actions.may_have_forgotten_stream(peer, id) {
618                    tracing::debug!("recv_data for old stream={:?}, sending STREAM_CLOSED", id,);
619
620                    let sz = frame.flow_controlled_len();
621                    // This should have been enforced at the codec::FramedRead layer, so
622                    // this is just a sanity check.
623                    assert!(sz <= super::MAX_WINDOW_SIZE as usize);
624                    let sz = sz as WindowSize;
625                    self.actions.recv.ignore_data(sz)?;
626
627                    return Err(Error::library_reset(id, Reason::STREAM_CLOSED));
628                }
629
630                proto_err!(conn: "recv_data: stream not found; id={:?}", id);
631                return Err(Error::library_go_away(Reason::PROTOCOL_ERROR));
632            }
633        };
634
635        let actions = &mut self.actions;
636        let mut send_buffer = send_buffer.inner.lock().unwrap();
637        let send_buffer = &mut *send_buffer;
638
639        self.counts.transition(stream, |counts, stream| {
640            let sz = frame.flow_controlled_len();
641            let is_end_stream = frame.is_end_stream();
642            let payload_len = frame.payload().len();
643            let mut res = actions.recv.recv_data(frame, stream);
644            // A stream can receive at most one final DATA frame, so it cannot
645            // be used to create unbounded framing overhead on that stream.
646            if res.is_ok() && !is_end_stream {
647                res = counts.record_data_frame(payload_len).map_err(|_| {
648                    tracing::debug!("too many small DATA frames");
649                    Error::library_go_away_data(Reason::ENHANCE_YOUR_CALM, "too_many_data_frames")
650                });
651            }
652
653            // Any stream error after receiving a DATA frame means
654            // we won't give the data to the user, and so they can't
655            // release the capacity. We do it automatically.
656            if let Err(Error::Reset(..)) = res {
657                actions
658                    .recv
659                    .release_connection_capacity(sz as WindowSize, &mut None);
660            }
661            actions.reset_on_recv_stream_err(send_buffer, stream, counts, res)
662        })
663    }
664
665    fn recv_reset<B>(
666        &mut self,
667        send_buffer: &SendBuffer<B>,
668        frame: frame::Reset,
669    ) -> Result<(), Error> {
670        let id = frame.stream_id();
671
672        if id.is_zero() {
673            proto_err!(conn: "recv_reset: invalid stream ID 0");
674            return Err(Error::library_go_away(Reason::PROTOCOL_ERROR));
675        }
676
677        // The GOAWAY process has begun. All streams with a greater ID than
678        // specified as part of GOAWAY should be ignored.
679        if id > self.actions.recv.max_stream_id() {
680            tracing::trace!(
681                "id ({:?}) > max_stream_id ({:?}), ignoring RST_STREAM",
682                id,
683                self.actions.recv.max_stream_id()
684            );
685            return Ok(());
686        }
687
688        let stream = match self.store.find_mut(&id) {
689            Some(stream) => stream,
690            None => {
691                // TODO: Are there other error cases?
692                self.actions
693                    .ensure_not_idle(self.counts.peer(), id)
694                    .map_err(Error::library_go_away)?;
695
696                return Ok(());
697            }
698        };
699
700        if stream.is_pending_open {
701            proto_err!(conn: "recv_reset: received frame on idle stream {:?}", id);
702            return Err(Error::library_go_away(Reason::PROTOCOL_ERROR));
703        }
704
705        let mut send_buffer = send_buffer.inner.lock().unwrap();
706        let send_buffer = &mut *send_buffer;
707
708        let actions = &mut self.actions;
709
710        self.counts.transition(stream, |counts, stream| {
711            actions.recv.recv_reset(frame, stream, counts)?;
712            actions.send.handle_error(send_buffer, stream, counts);
713            assert!(stream.state.is_closed());
714            Ok(())
715        })
716    }
717
718    fn recv_window_update<B>(
719        &mut self,
720        send_buffer: &SendBuffer<B>,
721        frame: frame::WindowUpdate,
722    ) -> Result<(), Error> {
723        let id = frame.stream_id();
724
725        let mut send_buffer = send_buffer.inner.lock().unwrap();
726        let send_buffer = &mut *send_buffer;
727
728        if id.is_zero() {
729            self.actions
730                .send
731                .recv_connection_window_update(frame, &mut self.store, &mut self.counts)
732                .map_err(Error::library_go_away)?;
733        } else {
734            // The remote may send window updates for streams that the local now
735            // considers closed. It's ok...
736            if let Some(mut stream) = self.store.find_mut(&id) {
737                if stream.is_pending_open {
738                    proto_err!(conn: "recv_window_update: received frame on idle stream {:?}", id);
739                    return Err(Error::library_go_away(Reason::PROTOCOL_ERROR));
740                }
741
742                let res = self
743                    .actions
744                    .send
745                    .recv_stream_window_update(
746                        frame.size_increment(),
747                        send_buffer,
748                        &mut stream,
749                        &mut self.counts,
750                        &mut self.actions.task,
751                    )
752                    .map_err(|reason| Error::library_reset(id, reason));
753
754                return self.actions.reset_on_recv_stream_err(
755                    send_buffer,
756                    &mut stream,
757                    &mut self.counts,
758                    res,
759                );
760            } else {
761                self.actions
762                    .ensure_not_idle(self.counts.peer(), id)
763                    .map_err(Error::library_go_away)?;
764            }
765        }
766
767        Ok(())
768    }
769
770    fn handle_error<B>(&mut self, send_buffer: &SendBuffer<B>, err: proto::Error) -> StreamId {
771        let actions = &mut self.actions;
772        let counts = &mut self.counts;
773        let mut send_buffer = send_buffer.inner.lock().unwrap();
774        let send_buffer = &mut *send_buffer;
775
776        let last_processed_id = actions.recv.last_processed_id();
777
778        self.store.for_each(|stream| {
779            counts.transition(stream, |counts, stream| {
780                actions.recv.handle_error(&err, &mut *stream);
781                actions.send.handle_error(send_buffer, stream, counts);
782            })
783        });
784
785        actions.conn_error = Some(err);
786
787        last_processed_id
788    }
789
790    fn recv_go_away<B>(
791        &mut self,
792        send_buffer: &SendBuffer<B>,
793        frame: &frame::GoAway,
794    ) -> Result<(), Error> {
795        let actions = &mut self.actions;
796        let counts = &mut self.counts;
797        let mut send_buffer = send_buffer.inner.lock().unwrap();
798        let send_buffer = &mut *send_buffer;
799
800        let last_stream_id = frame.last_stream_id();
801
802        actions.send.recv_go_away(last_stream_id)?;
803
804        let err = Error::remote_go_away(frame.debug_data().clone(), frame.reason());
805
806        let peer = counts.peer();
807        self.store.for_each(|stream| {
808            if stream.id > last_stream_id && peer.is_local_init(stream.id) {
809                counts.transition(stream, |counts, stream| {
810                    actions.recv.handle_error(&err, &mut *stream);
811                    actions.send.handle_error(send_buffer, stream, counts);
812                })
813            }
814        });
815
816        actions.conn_error = Some(err);
817
818        Ok(())
819    }
820
821    fn recv_push_promise<B>(
822        &mut self,
823        send_buffer: &SendBuffer<B>,
824        frame: frame::PushPromise,
825    ) -> Result<(), Error> {
826        let id = frame.stream_id();
827        let promised_id = frame.promised_id();
828
829        // First, ensure that the initiating stream is still in a valid state.
830        let parent_key = match self.store.find_mut(&id) {
831            Some(stream) => {
832                // The GOAWAY process has begun. All streams with a greater ID
833                // than specified as part of GOAWAY should be ignored.
834                if id > self.actions.recv.max_stream_id() {
835                    tracing::trace!(
836                        "id ({:?}) > max_stream_id ({:?}), ignoring PUSH_PROMISE",
837                        id,
838                        self.actions.recv.max_stream_id()
839                    );
840                    return Ok(());
841                }
842
843                // The stream must be receive open
844                if !stream.state.ensure_recv_open()? {
845                    proto_err!(conn: "recv_push_promise: initiating stream is not opened");
846                    return Err(Error::library_go_away(Reason::PROTOCOL_ERROR));
847                }
848
849                stream.key()
850            }
851            None => {
852                proto_err!(conn: "recv_push_promise: initiating stream is in an invalid state");
853                return Err(Error::library_go_away(Reason::PROTOCOL_ERROR));
854            }
855        };
856
857        // TODO: Streams in the reserved states do not count towards the concurrency
858        // limit. However, it seems like there should be a cap otherwise this
859        // could grow in memory indefinitely.
860
861        // Ensure that we can reserve streams
862        self.actions.recv.ensure_can_reserve()?;
863
864        // Next, open the stream.
865        //
866        // If `None` is returned, then the stream is being refused. There is no
867        // further work to be done.
868        if self
869            .actions
870            .recv
871            .open(promised_id, Open::PushPromise, &mut self.counts)?
872            .is_none()
873        {
874            return Ok(());
875        }
876
877        // Try to handle the frame and create a corresponding key for the pushed stream
878        // this requires a bit of indirection to make the borrow checker happy.
879        let child_key: Option<store::Key> = {
880            // Create state for the stream
881            let stream = self.store.insert(promised_id, {
882                Stream::new(
883                    promised_id,
884                    self.actions.send.init_window_sz(),
885                    self.actions.recv.init_window_sz(),
886                )
887            });
888
889            let actions = &mut self.actions;
890
891            self.counts.transition(stream, |counts, stream| {
892                let stream_valid = actions.recv.recv_push_promise(frame, stream);
893
894                match stream_valid {
895                    Ok(()) => Ok(Some(stream.key())),
896                    _ => {
897                        let mut send_buffer = send_buffer.inner.lock().unwrap();
898                        actions
899                            .reset_on_recv_stream_err(
900                                &mut *send_buffer,
901                                stream,
902                                counts,
903                                stream_valid,
904                            )
905                            .map(|()| None)
906                    }
907                }
908            })?
909        };
910        // If we're successful, push the headers and stream...
911        if let Some(child) = child_key {
912            let mut ppp = self.store[parent_key].pending_push_promises.take();
913            ppp.push(&mut self.store.resolve(child));
914
915            let parent = &mut self.store.resolve(parent_key);
916            parent.pending_push_promises = ppp;
917            parent.notify_push();
918        };
919
920        Ok(())
921    }
922
923    fn recv_eof<B>(
924        &mut self,
925        send_buffer: &SendBuffer<B>,
926        clear_pending_accept: bool,
927    ) -> Result<(), ()> {
928        let actions = &mut self.actions;
929        let counts = &mut self.counts;
930        let mut send_buffer = send_buffer.inner.lock().unwrap();
931        let send_buffer = &mut *send_buffer;
932
933        if actions.conn_error.is_none() {
934            actions.conn_error = Some(
935                io::Error::new(
936                    io::ErrorKind::BrokenPipe,
937                    "connection closed because of a broken pipe",
938                )
939                .into(),
940            );
941        }
942
943        tracing::trace!("Streams::recv_eof");
944
945        self.store.for_each(|stream| {
946            counts.transition(stream, |counts, stream| {
947                actions.recv.recv_eof(stream);
948
949                // This handles resetting send state associated with the
950                // stream
951                actions.send.handle_error(send_buffer, stream, counts);
952            })
953        });
954
955        actions.clear_queues(clear_pending_accept, &mut self.store, counts);
956        Ok(())
957    }
958
959    fn buffer_pending<T, B>(
960        &mut self,
961        send_buffer: &SendBuffer<B>,
962        dst: &mut Codec<T, Prioritized<B>>,
963    ) -> io::Result<BufferStatus>
964    where
965        T: AsyncWrite + Unpin,
966        B: Buf,
967    {
968        let mut send_buffer = send_buffer.inner.lock().unwrap();
969        let send_buffer = &mut *send_buffer;
970
971        // Send WINDOW_UPDATE frames first
972        //
973        // TODO: It would probably be better to interleave updates w/ data
974        // frames.
975        if self
976            .actions
977            .recv
978            .buffer_pending(&mut self.store, &mut self.counts, dst)?
979            == BufferStatus::CodecFull
980        {
981            return Ok(BufferStatus::CodecFull);
982        }
983
984        // Send any other pending frames
985        if self
986            .actions
987            .send
988            .buffer_pending(send_buffer, &mut self.store, &mut self.counts, dst)?
989            == BufferStatus::CodecFull
990        {
991            return Ok(BufferStatus::CodecFull);
992        }
993
994        Ok(BufferStatus::Complete)
995    }
996
997    fn reclaim_written_frame<T, B>(
998        &mut self,
999        send_buffer: &SendBuffer<B>,
1000        dst: &mut Codec<T, Prioritized<B>>,
1001    ) -> bool
1002    where
1003        B: Buf,
1004    {
1005        let mut send_buffer = send_buffer.inner.lock().unwrap();
1006        let send_buffer = &mut *send_buffer;
1007
1008        self.actions
1009            .send
1010            .reclaim_written_frame(send_buffer, &mut self.store, dst)
1011    }
1012
1013    fn send_reset<B>(
1014        &mut self,
1015        send_buffer: &SendBuffer<B>,
1016        id: StreamId,
1017        reason: Reason,
1018    ) -> Result<(), crate::proto::error::GoAway> {
1019        let key = match self.store.find_entry(id) {
1020            Entry::Occupied(e) => e.key(),
1021            Entry::Vacant(e) => {
1022                // Resetting a stream we don't know about? That could be OK...
1023                //
1024                // 1. As a server, we just received a request, but that request
1025                //    was bad, so we're resetting before even accepting it.
1026                //    This is totally fine.
1027                //
1028                // 2. The remote may have sent us a frame on new stream that
1029                //    it's *not* supposed to have done, and thus, we don't know
1030                //    the stream. In that case, sending a reset will "open" the
1031                //    stream in our store. Maybe that should be a connection
1032                //    error instead? At least for now, we need to update what
1033                //    our vision of the next stream is.
1034                if self.counts.peer().is_local_init(id) {
1035                    // We normally would open this stream, so update our
1036                    // next-send-id record.
1037                    self.actions.send.maybe_reset_next_stream_id(id);
1038                } else {
1039                    // We normally would recv this stream, so update our
1040                    // next-recv-id record.
1041                    self.actions.recv.maybe_reset_next_stream_id(id);
1042                }
1043
1044                let stream = Stream::new(id, 0, 0);
1045
1046                e.insert(stream)
1047            }
1048        };
1049
1050        let stream = self.store.resolve(key);
1051        let mut send_buffer = send_buffer.inner.lock().unwrap();
1052        let send_buffer = &mut *send_buffer;
1053        self.actions.send_reset(
1054            stream,
1055            reason,
1056            Initiator::Library,
1057            &mut self.counts,
1058            send_buffer,
1059        )
1060    }
1061}
1062
1063impl<B> Streams<B, client::Peer>
1064where
1065    B: Buf,
1066{
1067    pub fn poll_pending_open(
1068        &mut self,
1069        cx: &Context,
1070        pending: Option<&OpaqueStreamRef>,
1071    ) -> Poll<Result<(), crate::Error>> {
1072        let mut me = self.inner.lock().unwrap();
1073        let me = &mut *me;
1074
1075        me.actions.ensure_no_conn_error()?;
1076        me.actions.send.ensure_next_stream_id()?;
1077
1078        if let Some(pending) = pending {
1079            let mut stream = me.store.resolve(pending.key);
1080            tracing::trace!("poll_pending_open; stream = {:?}", stream.is_pending_open);
1081            if stream.is_pending_open {
1082                stream.wait_send(cx);
1083                return Poll::Pending;
1084            }
1085        }
1086        Poll::Ready(Ok(()))
1087    }
1088}
1089
1090impl<B, P> Streams<B, P>
1091where
1092    P: Peer,
1093{
1094    pub fn as_dyn(&self) -> DynStreams<'_, B> {
1095        let Self {
1096            inner,
1097            send_buffer,
1098            _p,
1099        } = self;
1100        DynStreams {
1101            inner,
1102            send_buffer,
1103            peer: P::r#dyn(),
1104        }
1105    }
1106
1107    /// This function is safe to call multiple times.
1108    ///
1109    /// A `Result` is returned to avoid panicking if the mutex is poisoned.
1110    pub fn recv_eof(&mut self, clear_pending_accept: bool) -> Result<(), ()> {
1111        self.as_dyn().recv_eof(clear_pending_accept)
1112    }
1113
1114    pub(crate) fn max_send_streams(&self) -> usize {
1115        self.inner.lock().unwrap().counts.max_send_streams()
1116    }
1117
1118    pub(crate) fn max_recv_streams(&self) -> usize {
1119        self.inner.lock().unwrap().counts.max_recv_streams()
1120    }
1121
1122    #[cfg(feature = "unstable")]
1123    pub fn num_active_streams(&self) -> usize {
1124        let me = self.inner.lock().unwrap();
1125        me.store.num_active_streams()
1126    }
1127
1128    pub fn has_streams(&self) -> bool {
1129        let me = self.inner.lock().unwrap();
1130        me.counts.has_streams()
1131    }
1132
1133    pub fn has_streams_or_other_references(&self) -> bool {
1134        let me = self.inner.lock().unwrap();
1135        me.counts.has_streams() || me.refs > 1
1136    }
1137
1138    #[cfg(feature = "unstable")]
1139    pub fn num_wired_streams(&self) -> usize {
1140        let me = self.inner.lock().unwrap();
1141        me.store.num_wired_streams()
1142    }
1143}
1144
1145// no derive because we don't need B and P to be Clone.
1146impl<B, P> Clone for Streams<B, P>
1147where
1148    P: Peer,
1149{
1150    fn clone(&self) -> Self {
1151        self.inner.lock().unwrap().refs += 1;
1152        Streams {
1153            inner: self.inner.clone(),
1154            send_buffer: self.send_buffer.clone(),
1155            _p: ::std::marker::PhantomData,
1156        }
1157    }
1158}
1159
1160impl<B, P> Drop for Streams<B, P>
1161where
1162    P: Peer,
1163{
1164    fn drop(&mut self) {
1165        if let Ok(mut inner) = self.inner.lock() {
1166            inner.refs -= 1;
1167            if inner.refs == 1 {
1168                if let Some(task) = inner.actions.task.take() {
1169                    task.wake();
1170                }
1171            }
1172        }
1173    }
1174}
1175
1176// ===== impl StreamRef =====
1177
1178impl<B> StreamRef<B> {
1179    pub fn send_data(&mut self, data: B, end_stream: bool) -> Result<(), UserError>
1180    where
1181        B: Buf,
1182    {
1183        let mut me = self.opaque.inner.lock().unwrap();
1184        let me = &mut *me;
1185
1186        let stream = me.store.resolve(self.opaque.key);
1187        let actions = &mut me.actions;
1188        let mut send_buffer = self.send_buffer.inner.lock().unwrap();
1189        let send_buffer = &mut *send_buffer;
1190
1191        me.counts.transition(stream, |counts, stream| {
1192            // Create the data frame
1193            let mut frame = frame::Data::new(stream.id, data);
1194            frame.set_end_stream(end_stream);
1195
1196            // Send the data frame
1197            actions
1198                .send
1199                .send_data(frame, send_buffer, stream, counts, &mut actions.task)
1200        })
1201    }
1202
1203    pub fn send_trailers(&mut self, trailers: HeaderMap) -> Result<(), UserError> {
1204        let mut me = self.opaque.inner.lock().unwrap();
1205        let me = &mut *me;
1206
1207        let stream = me.store.resolve(self.opaque.key);
1208        let actions = &mut me.actions;
1209        let mut send_buffer = self.send_buffer.inner.lock().unwrap();
1210        let send_buffer = &mut *send_buffer;
1211
1212        me.counts.transition(stream, |counts, stream| {
1213            // Create the trailers frame
1214            let frame = frame::Headers::trailers(stream.id, trailers);
1215
1216            // Send the trailers frame
1217            actions
1218                .send
1219                .send_trailers(frame, send_buffer, stream, counts, &mut actions.task)
1220        })
1221    }
1222
1223    pub fn send_reset(&mut self, reason: Reason) {
1224        let mut me = self.opaque.inner.lock().unwrap();
1225        let me = &mut *me;
1226
1227        let stream = me.store.resolve(self.opaque.key);
1228        let mut send_buffer = self.send_buffer.inner.lock().unwrap();
1229        let send_buffer = &mut *send_buffer;
1230
1231        match me
1232            .actions
1233            .send_reset(stream, reason, Initiator::User, &mut me.counts, send_buffer)
1234        {
1235            Ok(()) => (),
1236            Err(crate::proto::error::GoAway { .. }) => {
1237                // this should never happen, because Initiator::User resets do
1238                // not count toward the local limit.
1239                // we could perhaps make this state impossible, if we made the
1240                // initiator argument a generic, and so this could return
1241                // Infallible instead of an impossible GoAway, but oh well.
1242                unreachable!("Initiator::User should not error sending reset");
1243            }
1244        }
1245    }
1246
1247    pub fn send_informational_headers(&mut self, frame: frame::Headers) -> Result<(), UserError> {
1248        let mut me = self.opaque.inner.lock().unwrap();
1249        let me = &mut *me;
1250
1251        let stream = me.store.resolve(self.opaque.key);
1252        let actions = &mut me.actions;
1253        let mut send_buffer = self.send_buffer.inner.lock().unwrap();
1254        let send_buffer = &mut *send_buffer;
1255
1256        me.counts.transition(stream, |counts, stream| {
1257            // For informational responses (1xx), we need to send headers without
1258            // changing the stream state. This allows multiple informational responses
1259            // to be sent before the final response.
1260
1261            // Validate that this is actually an informational response
1262            debug_assert!(
1263                frame.is_informational(),
1264                "Frame must be informational after conversion from informational response"
1265            );
1266
1267            // Ensure the frame is not marked as end_stream for informational responses
1268            if frame.is_end_stream() {
1269                return Err(UserError::UnexpectedFrameType);
1270            }
1271
1272            // Send the interim informational headers directly to the buffer without state changes
1273            // This bypasses the normal send_headers flow that would transition the stream state
1274            actions.send.send_interim_informational_headers(
1275                frame,
1276                send_buffer,
1277                stream,
1278                counts,
1279                &mut actions.task,
1280            )
1281        })
1282    }
1283
1284    pub fn send_response(
1285        &mut self,
1286        mut response: Response<()>,
1287        end_of_stream: bool,
1288    ) -> Result<(), UserError> {
1289        // Clear before taking lock, incase extensions contain a StreamRef.
1290        response.extensions_mut().clear();
1291        let mut me = self.opaque.inner.lock().unwrap();
1292        let me = &mut *me;
1293
1294        let stream = me.store.resolve(self.opaque.key);
1295        let actions = &mut me.actions;
1296        let mut send_buffer = self.send_buffer.inner.lock().unwrap();
1297        let send_buffer = &mut *send_buffer;
1298
1299        me.counts.transition(stream, |counts, stream| {
1300            let frame = server::Peer::convert_send_message(stream.id, response, end_of_stream);
1301
1302            actions
1303                .send
1304                .send_headers(frame, send_buffer, stream, counts, &mut actions.task)
1305        })
1306    }
1307
1308    pub fn send_push_promise(
1309        &mut self,
1310        mut request: Request<()>,
1311    ) -> Result<StreamRef<B>, UserError> {
1312        // Clear before taking lock, incase extensions contain a StreamRef.
1313        request.extensions_mut().clear();
1314        let mut me = self.opaque.inner.lock().unwrap();
1315        let me = &mut *me;
1316
1317        let mut send_buffer = self.send_buffer.inner.lock().unwrap();
1318        let send_buffer = &mut *send_buffer;
1319
1320        let actions = &mut me.actions;
1321        let promised_id = actions.send.reserve_local()?;
1322
1323        let child_key = {
1324            let mut child_stream = me.store.insert(
1325                promised_id,
1326                Stream::new(
1327                    promised_id,
1328                    actions.send.init_window_sz(),
1329                    actions.recv.init_window_sz(),
1330                ),
1331            );
1332            child_stream.state.reserve_local()?;
1333            child_stream.is_pending_push = true;
1334            child_stream.key()
1335        };
1336
1337        let pushed = {
1338            let mut stream = me.store.resolve(self.opaque.key);
1339
1340            let frame = crate::server::Peer::convert_push_message(stream.id, promised_id, request)?;
1341
1342            actions
1343                .send
1344                .send_push_promise(frame, send_buffer, &mut stream, &mut actions.task)
1345        };
1346
1347        if let Err(err) = pushed {
1348            let mut child_stream = me.store.resolve(child_key);
1349            child_stream.unlink();
1350            child_stream.remove();
1351            return Err(err);
1352        }
1353
1354        me.refs += 1;
1355        let opaque =
1356            OpaqueStreamRef::new(self.opaque.inner.clone(), &mut me.store.resolve(child_key));
1357
1358        Ok(StreamRef {
1359            opaque,
1360            send_buffer: self.send_buffer.clone(),
1361        })
1362    }
1363
1364    /// Called by the server after the stream is accepted. Given that clients
1365    /// initialize streams by sending HEADERS, the request will always be
1366    /// available.
1367    ///
1368    /// # Panics
1369    ///
1370    /// This function panics if the request isn't present.
1371    pub fn take_request(&self) -> Request<()> {
1372        let mut me = self.opaque.inner.lock().unwrap();
1373        let me = &mut *me;
1374
1375        let mut stream = me.store.resolve(self.opaque.key);
1376        me.actions.recv.take_request(&mut stream)
1377    }
1378
1379    /// Called by a client to see if the current stream is pending open
1380    pub fn is_pending_open(&self) -> bool {
1381        let mut me = self.opaque.inner.lock().unwrap();
1382        me.store.resolve(self.opaque.key).is_pending_open
1383    }
1384
1385    /// Request capacity to send data
1386    pub fn reserve_capacity(&mut self, capacity: WindowSize) {
1387        let mut me = self.opaque.inner.lock().unwrap();
1388        let me = &mut *me;
1389
1390        let mut stream = me.store.resolve(self.opaque.key);
1391
1392        me.actions
1393            .send
1394            .reserve_capacity(capacity, &mut stream, &mut me.counts)
1395    }
1396
1397    /// Returns the stream's current send capacity.
1398    pub fn capacity(&self) -> WindowSize {
1399        let mut me = self.opaque.inner.lock().unwrap();
1400        let me = &mut *me;
1401
1402        let mut stream = me.store.resolve(self.opaque.key);
1403
1404        me.actions.send.capacity(&mut stream)
1405    }
1406
1407    /// Request to be notified when the stream's capacity increases
1408    pub fn poll_capacity(&mut self, cx: &Context) -> Poll<Option<Result<WindowSize, UserError>>> {
1409        let mut me = self.opaque.inner.lock().unwrap();
1410        let me = &mut *me;
1411
1412        let mut stream = me.store.resolve(self.opaque.key);
1413
1414        me.actions.send.poll_capacity(cx, &mut stream)
1415    }
1416
1417    /// Request to be notified for if a `RST_STREAM` is received for this stream.
1418    pub(crate) fn poll_reset(
1419        &mut self,
1420        cx: &Context,
1421        mode: proto::PollReset,
1422    ) -> Poll<Result<Reason, crate::Error>> {
1423        let mut me = self.opaque.inner.lock().unwrap();
1424        let me = &mut *me;
1425
1426        let mut stream = me.store.resolve(self.opaque.key);
1427
1428        me.actions.send.poll_reset(cx, &mut stream, mode)
1429    }
1430
1431    pub fn clone_to_opaque(&self) -> OpaqueStreamRef {
1432        self.opaque.clone()
1433    }
1434
1435    pub fn stream_id(&self) -> StreamId {
1436        self.opaque.stream_id()
1437    }
1438}
1439
1440impl<B> Clone for StreamRef<B> {
1441    fn clone(&self) -> Self {
1442        StreamRef {
1443            opaque: self.opaque.clone(),
1444            send_buffer: self.send_buffer.clone(),
1445        }
1446    }
1447}
1448
1449// ===== impl OpaqueStreamRef =====
1450
1451impl OpaqueStreamRef {
1452    fn new(inner: Arc<Mutex<Inner>>, stream: &mut store::Ptr) -> OpaqueStreamRef {
1453        stream.ref_inc();
1454        OpaqueStreamRef {
1455            inner,
1456            key: stream.key(),
1457        }
1458    }
1459    /// Called by a client to check for a received response.
1460    pub fn poll_response(&mut self, cx: &Context) -> Poll<Result<Response<()>, proto::Error>> {
1461        let mut me = self.inner.lock().unwrap();
1462        let me = &mut *me;
1463
1464        let mut stream = me.store.resolve(self.key);
1465
1466        me.actions.recv.poll_response(cx, &mut stream)
1467    }
1468
1469    /// Called by a client to check for informational responses (1xx status codes)
1470    pub fn poll_informational(
1471        &mut self,
1472        cx: &Context,
1473    ) -> Poll<Option<Result<Response<()>, proto::Error>>> {
1474        let mut me = self.inner.lock().unwrap();
1475        let me = &mut *me;
1476
1477        let mut stream = me.store.resolve(self.key);
1478
1479        me.actions.recv.poll_informational(cx, &mut stream)
1480    }
1481    /// Called by a client to check for a pushed request.
1482    pub fn poll_pushed(
1483        &mut self,
1484        cx: &Context,
1485    ) -> Poll<Option<Result<(Request<()>, OpaqueStreamRef), proto::Error>>> {
1486        let mut me = self.inner.lock().unwrap();
1487        let me = &mut *me;
1488
1489        let mut stream = me.store.resolve(self.key);
1490        me.actions
1491            .recv
1492            .poll_pushed(cx, &mut stream)
1493            .map_ok(|(h, key)| {
1494                me.refs += 1;
1495                let opaque_ref =
1496                    OpaqueStreamRef::new(self.inner.clone(), &mut me.store.resolve(key));
1497                (h, opaque_ref)
1498            })
1499    }
1500
1501    pub fn is_end_stream(&self) -> bool {
1502        let mut me = self.inner.lock().unwrap();
1503        let me = &mut *me;
1504
1505        let stream = me.store.resolve(self.key);
1506
1507        me.actions.recv.is_end_stream(&stream)
1508    }
1509
1510    pub fn poll_data(&mut self, cx: &Context) -> Poll<Option<Result<Bytes, proto::Error>>> {
1511        let mut me = self.inner.lock().unwrap();
1512        let me = &mut *me;
1513
1514        let mut stream = me.store.resolve(self.key);
1515
1516        me.actions
1517            .recv
1518            .poll_data(cx, &mut stream)
1519            .map(|result| match result {
1520                Some(Ok(data)) => {
1521                    if data.is_budgeted {
1522                        me.counts.release_data_frame(data.payload.len());
1523                    }
1524                    Some(Ok(data.payload))
1525                }
1526                Some(Err(err)) => Some(Err(err)),
1527                None => None,
1528            })
1529    }
1530
1531    pub fn poll_trailers(&mut self, cx: &Context) -> Poll<Option<Result<HeaderMap, proto::Error>>> {
1532        let mut me = self.inner.lock().unwrap();
1533        let me = &mut *me;
1534
1535        let mut stream = me.store.resolve(self.key);
1536
1537        me.actions.recv.poll_trailers(cx, &mut stream)
1538    }
1539
1540    pub(crate) fn available_recv_capacity(&self) -> isize {
1541        let me = self.inner.lock().unwrap();
1542        let me = &*me;
1543
1544        let stream = &me.store[self.key];
1545        stream.recv_flow.available().into()
1546    }
1547
1548    pub(crate) fn used_recv_capacity(&self) -> WindowSize {
1549        let me = self.inner.lock().unwrap();
1550        let me = &*me;
1551
1552        let stream = &me.store[self.key];
1553        stream.in_flight_recv_data
1554    }
1555
1556    /// Releases recv capacity back to the peer. This may result in sending
1557    /// WINDOW_UPDATE frames on both the stream and connection.
1558    pub fn release_capacity(&mut self, capacity: WindowSize) -> Result<(), UserError> {
1559        let mut me = self.inner.lock().unwrap();
1560        let me = &mut *me;
1561
1562        let mut stream = me.store.resolve(self.key);
1563
1564        me.actions
1565            .recv
1566            .release_capacity(capacity, &mut stream, &mut me.actions.task)
1567    }
1568
1569    /// Clear the receive queue and set the status to no longer receive data frames.
1570    pub(crate) fn clear_recv_buffer(&mut self) {
1571        let mut me = self.inner.lock().unwrap();
1572        let me = &mut *me;
1573
1574        let mut stream = me.store.resolve(self.key);
1575        stream.is_recv = false;
1576        me.actions
1577            .recv
1578            .clear_recv_buffer(&mut stream, &mut me.actions.task, &mut me.counts);
1579    }
1580
1581    pub fn stream_id(&self) -> StreamId {
1582        self.inner.lock().unwrap().store[self.key].id
1583    }
1584}
1585
1586impl fmt::Debug for OpaqueStreamRef {
1587    fn fmt(&self, fmt: &mut fmt::Formatter) -> fmt::Result {
1588        use std::sync::TryLockError::*;
1589
1590        match self.inner.try_lock() {
1591            Ok(me) => {
1592                let stream = &me.store[self.key];
1593                fmt.debug_struct("OpaqueStreamRef")
1594                    .field("stream_id", &stream.id)
1595                    .field("ref_count", &stream.ref_count)
1596                    .finish()
1597            }
1598            Err(Poisoned(_)) => fmt
1599                .debug_struct("OpaqueStreamRef")
1600                .field("inner", &"<Poisoned>")
1601                .finish(),
1602            Err(WouldBlock) => fmt
1603                .debug_struct("OpaqueStreamRef")
1604                .field("inner", &"<Locked>")
1605                .finish(),
1606        }
1607    }
1608}
1609
1610impl Clone for OpaqueStreamRef {
1611    fn clone(&self) -> Self {
1612        // Increment the ref count
1613        let mut inner = self.inner.lock().unwrap();
1614        inner.store.resolve(self.key).ref_inc();
1615        inner.refs += 1;
1616
1617        OpaqueStreamRef {
1618            inner: self.inner.clone(),
1619            key: self.key,
1620        }
1621    }
1622}
1623
1624impl Drop for OpaqueStreamRef {
1625    fn drop(&mut self) {
1626        drop_stream_ref(&self.inner, self.key);
1627    }
1628}
1629
1630// TODO: Move back in fn above
1631fn drop_stream_ref(inner: &Mutex<Inner>, key: store::Key) {
1632    let mut me = match inner.lock() {
1633        Ok(inner) => inner,
1634        Err(_) => {
1635            if ::std::thread::panicking() {
1636                tracing::trace!("StreamRef::drop; mutex poisoned");
1637                return;
1638            } else {
1639                panic!("StreamRef::drop; mutex poisoned");
1640            }
1641        }
1642    };
1643
1644    let me = &mut *me;
1645    me.refs -= 1;
1646    let mut stream = me.store.resolve(key);
1647
1648    tracing::trace!("drop_stream_ref; stream={:?}", stream);
1649
1650    // decrement the stream's ref count by 1.
1651    stream.ref_dec();
1652
1653    let actions = &mut me.actions;
1654
1655    // If the stream is not referenced and it is already
1656    // closed (does not have to go through logic below
1657    // of canceling the stream), we should notify the task
1658    // (connection) so that it can close properly
1659    if stream.ref_count == 0 && stream.is_closed() {
1660        if let Some(task) = actions.task.take() {
1661            task.wake();
1662        }
1663    }
1664
1665    me.counts.transition(stream, |counts, stream| {
1666        maybe_cancel(stream, actions, counts);
1667
1668        if stream.ref_count == 0 {
1669            // Release any recv window back to connection, no one can access
1670            // it anymore.
1671            actions
1672                .recv
1673                .release_closed_capacity(stream, &mut actions.task, counts);
1674
1675            // We won't be able to reach our push promises anymore
1676            let mut ppp = stream.pending_push_promises.take();
1677            while let Some(promise) = ppp.pop(stream.store_mut()) {
1678                counts.transition(promise, |counts, stream| {
1679                    maybe_cancel(stream, actions, counts);
1680                });
1681            }
1682        }
1683    });
1684}
1685
1686fn maybe_cancel(stream: &mut store::Ptr, actions: &mut Actions, counts: &mut Counts) {
1687    if stream.is_canceled_interest() {
1688        // Server is allowed to early respond without fully consuming the client input stream
1689        // But per the RFC, must send a RST_STREAM(NO_ERROR) in such cases. https://www.rfc-editor.org/rfc/rfc7540#section-8.1
1690        // Some other http2 implementation may interpret other error code as fatal if not respected (i.e: nginx https://trac.nginx.org/nginx/ticket/2376)
1691        let reason = if counts.peer().is_server()
1692            && stream.state.is_send_closed()
1693            && stream.state.is_recv_streaming()
1694        {
1695            Reason::NO_ERROR
1696        } else {
1697            Reason::CANCEL
1698        };
1699
1700        actions
1701            .send
1702            .schedule_implicit_reset(stream, reason, counts, &mut actions.task);
1703        actions.recv.enqueue_reset_expiration(stream, counts);
1704    }
1705}
1706
1707// ===== impl SendBuffer =====
1708
1709impl<B> SendBuffer<B> {
1710    fn new() -> Self {
1711        let inner = Mutex::new(Buffer::new());
1712        SendBuffer { inner }
1713    }
1714
1715    pub fn is_empty(&self) -> bool {
1716        let buf = self.inner.lock().unwrap();
1717        buf.is_empty()
1718    }
1719}
1720
1721// ===== impl Actions =====
1722
1723impl Actions {
1724    fn send_reset<B>(
1725        &mut self,
1726        stream: store::Ptr,
1727        reason: Reason,
1728        initiator: Initiator,
1729        counts: &mut Counts,
1730        send_buffer: &mut Buffer<Frame<B>>,
1731    ) -> Result<(), crate::proto::error::GoAway> {
1732        counts.transition(stream, |counts, stream| {
1733            if initiator.is_library() {
1734                if counts.can_inc_num_local_error_resets() {
1735                    counts.inc_num_local_error_resets();
1736                } else {
1737                    tracing::warn!(
1738                        "locally-reset streams reached limit ({:?})",
1739                        counts.max_local_error_resets().unwrap(),
1740                    );
1741                    return Err(crate::proto::error::GoAway {
1742                        reason: Reason::ENHANCE_YOUR_CALM,
1743                        debug_data: "too_many_internal_resets".into(),
1744                    });
1745                }
1746            }
1747
1748            self.send.send_reset(
1749                reason,
1750                initiator,
1751                send_buffer,
1752                stream,
1753                counts,
1754                &mut self.task,
1755            );
1756            self.recv.enqueue_reset_expiration(stream, counts);
1757            // if a RecvStream is parked, ensure it's notified
1758            stream.notify_recv();
1759
1760            Ok(())
1761        })
1762    }
1763
1764    fn reset_on_recv_stream_err<B>(
1765        &mut self,
1766        buffer: &mut Buffer<Frame<B>>,
1767        stream: &mut store::Ptr,
1768        counts: &mut Counts,
1769        res: Result<(), Error>,
1770    ) -> Result<(), Error> {
1771        if let Err(Error::Reset(stream_id, reason, initiator)) = res {
1772            debug_assert_eq!(stream_id, stream.id);
1773
1774            if counts.can_inc_num_local_error_resets() {
1775                counts.inc_num_local_error_resets();
1776
1777                // Reset the stream.
1778                self.send
1779                    .send_reset(reason, initiator, buffer, stream, counts, &mut self.task);
1780                self.recv.enqueue_reset_expiration(stream, counts);
1781                // if a RecvStream is parked, ensure it's notified
1782                stream.notify_recv();
1783                Ok(())
1784            } else {
1785                tracing::warn!(
1786                    "reset_on_recv_stream_err; locally-reset streams reached limit ({:?})",
1787                    counts.max_local_error_resets().unwrap(),
1788                );
1789                Err(Error::library_go_away_data(
1790                    Reason::ENHANCE_YOUR_CALM,
1791                    "too_many_internal_resets",
1792                ))
1793            }
1794        } else {
1795            res
1796        }
1797    }
1798
1799    fn ensure_not_idle(&mut self, peer: peer::Dyn, id: StreamId) -> Result<(), Reason> {
1800        if peer.is_local_init(id) {
1801            self.send.ensure_not_idle(id)
1802        } else {
1803            self.recv.ensure_not_idle(id)
1804        }
1805    }
1806
1807    fn ensure_no_conn_error(&self) -> Result<(), proto::Error> {
1808        if let Some(ref err) = self.conn_error {
1809            Err(err.clone())
1810        } else {
1811            Ok(())
1812        }
1813    }
1814
1815    /// Check if we possibly could have processed and since forgotten this stream.
1816    ///
1817    /// If we send a RST_STREAM for a stream, we will eventually "forget" about
1818    /// the stream to free up memory. It's possible that the remote peer had
1819    /// frames in-flight, and by the time we receive them, our own state is
1820    /// gone. We *could* tear everything down by sending a GOAWAY, but it
1821    /// is more likely to be latency/memory constraints that caused this,
1822    /// and not a bad actor. So be less catastrophic, the spec allows
1823    /// us to send another RST_STREAM of STREAM_CLOSED.
1824    fn may_have_forgotten_stream(&self, peer: peer::Dyn, id: StreamId) -> bool {
1825        if id.is_zero() {
1826            return false;
1827        }
1828        if peer.is_local_init(id) {
1829            self.send.may_have_created_stream(id)
1830        } else {
1831            self.recv.may_have_created_stream(id)
1832        }
1833    }
1834
1835    fn clear_queues(&mut self, clear_pending_accept: bool, store: &mut Store, counts: &mut Counts) {
1836        self.recv.clear_queues(clear_pending_accept, store, counts);
1837        self.send.clear_queues(store, counts);
1838    }
1839}