mirror of
https://github.com/n0-computer/noq.git
synced 2026-10-04 05:25:48 +00:00
quinn-proto: move post_read() logic into Retransmits
This commit is contained in:
@@ -50,7 +50,7 @@ pub use stats::ConnectionStats;
|
||||
|
||||
mod streams;
|
||||
pub use streams::Streams;
|
||||
pub use streams::{FinishError, ReadError, StreamEvent, UnknownStream, WriteError};
|
||||
pub use streams::{FinishError, ReadError, ShouldTransmit, StreamEvent, UnknownStream, WriteError};
|
||||
|
||||
mod timer;
|
||||
use timer::{Timer, TimerTable};
|
||||
@@ -922,23 +922,16 @@ where
|
||||
}
|
||||
|
||||
fn post_read<T>(&mut self, id: StreamId, result: &streams::ReadResult<T>) {
|
||||
if let Ok(Some(ref result)) = *result {
|
||||
let pending = &mut self.spaces[SpaceId::Data].pending;
|
||||
if result.max_data.should_transmit() {
|
||||
pending.max_data = true;
|
||||
}
|
||||
if result.max_stream_data.should_transmit() {
|
||||
// Only bother issuing stream credit if the peer wants to send more
|
||||
pending.max_stream_data.insert(id);
|
||||
}
|
||||
}
|
||||
let (did_read, max_data, max_stream_data) = match result {
|
||||
Ok(Some(did)) => (true, did.max_data, did.max_stream_data),
|
||||
_ => (false, ShouldTransmit::default(), ShouldTransmit::default()),
|
||||
};
|
||||
|
||||
if self.streams.take_max_streams_dirty(id.dir()) {
|
||||
let pending = &mut self.spaces[SpaceId::Data].pending;
|
||||
match id.dir() {
|
||||
Dir::Uni => pending.max_uni_stream_id = true,
|
||||
Dir::Bi => pending.max_bi_stream_id = true,
|
||||
}
|
||||
let max_dirty = self.streams.take_max_streams_dirty(id.dir());
|
||||
if did_read || max_dirty {
|
||||
self.spaces[SpaceId::Data]
|
||||
.pending
|
||||
.post_read(id, max_data, max_stream_data, max_dirty);
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -7,9 +7,10 @@ use std::{
|
||||
};
|
||||
|
||||
use super::assembler::Assembler;
|
||||
use super::streams::ShouldTransmit;
|
||||
use crate::{
|
||||
crypto, crypto::Keys, frame, packet::SpaceId, range_set::RangeSet, shared::IssuedCid, StreamId,
|
||||
VarInt,
|
||||
crypto, crypto::Keys, frame, packet::SpaceId, range_set::RangeSet, shared::IssuedCid, Dir,
|
||||
StreamId, VarInt,
|
||||
};
|
||||
|
||||
pub(crate) struct PacketSpace<S>
|
||||
@@ -204,6 +205,28 @@ pub struct Retransmits {
|
||||
}
|
||||
|
||||
impl Retransmits {
|
||||
pub(crate) fn post_read(
|
||||
&mut self,
|
||||
id: StreamId,
|
||||
max_data: ShouldTransmit,
|
||||
max_stream_data: ShouldTransmit,
|
||||
max_dirty: bool,
|
||||
) {
|
||||
if max_data.should_transmit() {
|
||||
self.max_data = true;
|
||||
}
|
||||
if max_stream_data.should_transmit() {
|
||||
// Only bother issuing stream credit if the peer wants to send more
|
||||
self.max_stream_data.insert(id);
|
||||
}
|
||||
if max_dirty {
|
||||
match id.dir() {
|
||||
Dir::Uni => self.max_uni_stream_id = true,
|
||||
Dir::Bi => self.max_bi_stream_id = true,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
pub fn is_empty(&self) -> bool {
|
||||
!self.max_data
|
||||
&& !self.max_uni_stream_id
|
||||
|
||||
@@ -19,8 +19,8 @@ use crate::{
|
||||
|
||||
mod recv;
|
||||
pub use recv::ReadError;
|
||||
pub(super) use recv::ReadResult;
|
||||
use recv::{BytesRead, DidRead, ReadChunks, Recv, StreamReadResult};
|
||||
use recv::{BytesRead, ReadChunks, Recv, StreamReadResult};
|
||||
pub(super) use recv::{DidRead, ReadResult};
|
||||
|
||||
mod send;
|
||||
pub use send::{FinishError, WriteError};
|
||||
@@ -963,7 +963,7 @@ pub enum StreamEvent {
|
||||
///
|
||||
/// This type wraps around bool and uses the `#[must_use]` attribute in order
|
||||
/// to prevent accidental loss of the frame transmission requirement.
|
||||
#[derive(Debug, Copy, Clone, Eq, PartialEq)]
|
||||
#[derive(Copy, Clone, Debug, Default, Eq, PartialEq)]
|
||||
#[must_use = "A frame might need to be enqueued"]
|
||||
pub struct ShouldTransmit(bool);
|
||||
|
||||
|
||||
Reference in New Issue
Block a user