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 inner: Arc<Mutex<Inner>>,
26
27 send_buffer: Arc<SendBuffer<B>>,
36
37 _p: ::std::marker::PhantomData<P>,
38}
39
40#[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#[derive(Debug)]
53pub(crate) struct StreamRef<B> {
54 opaque: OpaqueStreamRef,
55 send_buffer: Arc<SendBuffer<B>>,
56}
57
58pub(crate) struct OpaqueStreamRef {
60 inner: Arc<Mutex<Inner>>,
61 key: store::Key,
62}
63
64#[derive(Debug)]
69struct Inner {
70 counts: Counts,
72
73 actions: Actions,
75
76 store: Store,
78
79 refs: usize,
81}
82
83#[derive(Debug)]
84struct Actions {
85 recv: Recv,
87
88 send: Send,
90
91 task: Option<Waker>,
93
94 conn_error: Option<proto::Error>,
96}
97
98#[derive(Debug)]
100struct SendBuffer<B> {
101 inner: Mutex<Buffer<Frame<B>>>,
102}
103
104impl<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 me.refs += 1;
143
144 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 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 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 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 request.extensions_mut().clear();
275
276 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 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 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 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 if let Err(err) = sent {
336 stream.unlink();
337 stream.remove();
338 return Err(err.into());
339 }
340
341 debug_assert!(!stream.state.is_closed());
344
345 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 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 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 if !peer.is_server() {
493 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 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 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 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 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 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 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 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 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 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 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 let parent_key = match self.store.find_mut(&id) {
831 Some(stream) => {
832 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 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 self.actions.recv.ensure_can_reserve()?;
863
864 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 let child_key: Option<store::Key> = {
880 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 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 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 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 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 if self.counts.peer().is_local_init(id) {
1035 self.actions.send.maybe_reset_next_stream_id(id);
1038 } else {
1039 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 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
1145impl<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
1176impl<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 let mut frame = frame::Data::new(stream.id, data);
1194 frame.set_end_stream(end_stream);
1195
1196 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 let frame = frame::Headers::trailers(stream.id, trailers);
1215
1216 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 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 debug_assert!(
1263 frame.is_informational(),
1264 "Frame must be informational after conversion from informational response"
1265 );
1266
1267 if frame.is_end_stream() {
1269 return Err(UserError::UnexpectedFrameType);
1270 }
1271
1272 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 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 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 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 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 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 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 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 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
1449impl 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 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 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 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 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 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 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
1630fn 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 stream.ref_dec();
1652
1653 let actions = &mut me.actions;
1654
1655 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 actions
1672 .recv
1673 .release_closed_capacity(stream, &mut actions.task, counts);
1674
1675 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 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
1707impl<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
1721impl 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 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 self.send
1779 .send_reset(reason, initiator, buffer, stream, counts, &mut self.task);
1780 self.recv.enqueue_reset_expiration(stream, counts);
1781 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 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}