mirror of
https://github.com/n0-computer/noq.git
synced 2026-10-04 05:25:48 +00:00
quinn-proto: simplify ShouldTransmit interface
This commit is contained in:
@@ -290,13 +290,13 @@ impl Streams {
|
||||
Some(rs) => rs,
|
||||
None => {
|
||||
trace!("dropping frame for closed stream");
|
||||
return Ok(ShouldTransmit::new(false));
|
||||
return Ok(ShouldTransmit(false));
|
||||
}
|
||||
};
|
||||
|
||||
if rs.is_finished() {
|
||||
trace!("dropping frame for finished stream");
|
||||
return Ok(ShouldTransmit::new(false));
|
||||
return Ok(ShouldTransmit(false));
|
||||
}
|
||||
|
||||
let new_bytes = rs.ingest(frame, self.data_recvd, self.local_max_data)?;
|
||||
@@ -304,7 +304,7 @@ impl Streams {
|
||||
|
||||
if !rs.assembler.is_stopped() {
|
||||
self.on_stream_frame(true, stream);
|
||||
return Ok(ShouldTransmit::new(false));
|
||||
return Ok(ShouldTransmit(false));
|
||||
}
|
||||
|
||||
// Stopped streams become closed instantly on FIN, so check whether we need to clean up
|
||||
@@ -338,7 +338,7 @@ impl Streams {
|
||||
Some(stream) => stream,
|
||||
None => {
|
||||
trace!("received RESET_STREAM on closed stream");
|
||||
return Ok(ShouldTransmit::new(false));
|
||||
return Ok(ShouldTransmit(false));
|
||||
}
|
||||
};
|
||||
|
||||
@@ -350,7 +350,7 @@ impl Streams {
|
||||
self.local_max_data,
|
||||
)? {
|
||||
// Redundant reset
|
||||
return Ok(ShouldTransmit::new(false));
|
||||
return Ok(ShouldTransmit(false));
|
||||
}
|
||||
let bytes_read = rs.assembler.bytes_read();
|
||||
let stopped = rs.assembler.is_stopped();
|
||||
@@ -369,7 +369,7 @@ impl Streams {
|
||||
.saturating_add(u64::from(final_offset) - end);
|
||||
self.add_read_credits(u64::from(final_offset) - bytes_read)
|
||||
} else {
|
||||
ShouldTransmit::new(false)
|
||||
ShouldTransmit(false)
|
||||
})
|
||||
}
|
||||
|
||||
@@ -446,7 +446,7 @@ impl Streams {
|
||||
return Err(UnknownStream { _private: () });
|
||||
}
|
||||
stream.assembler.stop();
|
||||
let stop_sending = ShouldTransmit::new(!stream.is_finished());
|
||||
let stop_sending = ShouldTransmit(!stream.is_finished());
|
||||
|
||||
// Issue flow control credit for unread data
|
||||
let read_credits = stream.assembler.end() - stream.assembler.bytes_read();
|
||||
@@ -879,7 +879,7 @@ impl Streams {
|
||||
self.local_max_data = self.local_max_data.saturating_add(credits);
|
||||
|
||||
if self.local_max_data > VarInt::MAX.into_inner() {
|
||||
return ShouldTransmit::new(false);
|
||||
return ShouldTransmit(false);
|
||||
}
|
||||
|
||||
// Only announce a window update if it's significant enough
|
||||
@@ -888,7 +888,7 @@ impl Streams {
|
||||
// the decision, to accomodate for connection using bigger windows requring
|
||||
// less updates.
|
||||
let diff = self.local_max_data - self.sent_max_data.into_inner();
|
||||
ShouldTransmit::new(diff >= (self.receive_window / 8))
|
||||
ShouldTransmit(diff >= (self.receive_window / 8))
|
||||
}
|
||||
|
||||
/// Records that a `MAX_DATA` announcing a certain window was sent
|
||||
@@ -965,19 +965,12 @@ pub enum StreamEvent {
|
||||
/// to prevent accidental loss of the frame transmission requirement.
|
||||
#[derive(Debug, Copy, Clone, Eq, PartialEq)]
|
||||
#[must_use = "A frame might need to be enqueued"]
|
||||
pub struct ShouldTransmit {
|
||||
should_transmit: bool,
|
||||
}
|
||||
pub struct ShouldTransmit(bool);
|
||||
|
||||
impl ShouldTransmit {
|
||||
/// Creates a new `ShouldTransmit` instance
|
||||
pub fn new(should_transmit: bool) -> Self {
|
||||
Self { should_transmit }
|
||||
}
|
||||
|
||||
/// Returns whether a frame should be transmitted
|
||||
pub fn should_transmit(self) -> bool {
|
||||
self.should_transmit
|
||||
self.0
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1024,7 +1017,7 @@ mod tests {
|
||||
data: Bytes::from_static(&[0; 2048]),
|
||||
})
|
||||
.unwrap(),
|
||||
ShouldTransmit::new(false)
|
||||
ShouldTransmit(false)
|
||||
);
|
||||
assert_eq!(client.data_recvd, 2048);
|
||||
assert_eq!(client.local_max_data - initial_max, 0);
|
||||
@@ -1038,7 +1031,7 @@ mod tests {
|
||||
final_offset: 4096u32.into(),
|
||||
})
|
||||
.unwrap(),
|
||||
ShouldTransmit::new(false)
|
||||
ShouldTransmit(false)
|
||||
);
|
||||
assert_eq!(client.data_recvd, 4096);
|
||||
assert_eq!(client.local_max_data - initial_max, 4096);
|
||||
@@ -1058,7 +1051,7 @@ mod tests {
|
||||
data: Bytes::from_static(&[0; 0]),
|
||||
})
|
||||
.unwrap(),
|
||||
ShouldTransmit::new(false)
|
||||
ShouldTransmit(false)
|
||||
);
|
||||
assert_eq!(client.data_recvd, 4096);
|
||||
assert_eq!(client.local_max_data - initial_max, 0);
|
||||
@@ -1070,7 +1063,7 @@ mod tests {
|
||||
final_offset: 4096u32.into(),
|
||||
})
|
||||
.unwrap(),
|
||||
ShouldTransmit::new(false)
|
||||
ShouldTransmit(false)
|
||||
);
|
||||
assert_eq!(client.data_recvd, 4096);
|
||||
assert_eq!(client.local_max_data - initial_max, 4096);
|
||||
@@ -1088,7 +1081,7 @@ mod tests {
|
||||
final_offset: 4096u32.into(),
|
||||
})
|
||||
.unwrap(),
|
||||
ShouldTransmit::new(false)
|
||||
ShouldTransmit(false)
|
||||
);
|
||||
assert_eq!(client.data_recvd, 4096);
|
||||
assert_eq!(
|
||||
@@ -1099,7 +1092,7 @@ mod tests {
|
||||
final_offset: 4096u32.into(),
|
||||
})
|
||||
.unwrap(),
|
||||
ShouldTransmit::new(false)
|
||||
ShouldTransmit(false)
|
||||
);
|
||||
assert_eq!(client.data_recvd, 4096);
|
||||
}
|
||||
@@ -1118,14 +1111,14 @@ mod tests {
|
||||
data: Bytes::from_static(&[0; 32]),
|
||||
})
|
||||
.unwrap(),
|
||||
ShouldTransmit::new(false)
|
||||
ShouldTransmit(false)
|
||||
);
|
||||
assert_eq!(client.local_max_data, initial_max);
|
||||
assert_eq!(
|
||||
client.stop(id).unwrap(),
|
||||
StopResult {
|
||||
max_data: ShouldTransmit::new(false),
|
||||
stop_sending: ShouldTransmit::new(true),
|
||||
max_data: ShouldTransmit(false),
|
||||
stop_sending: ShouldTransmit(true),
|
||||
}
|
||||
);
|
||||
assert!(client.stop(id).is_err());
|
||||
@@ -1141,7 +1134,7 @@ mod tests {
|
||||
data: Bytes::from_static(&[0; 16]),
|
||||
})
|
||||
.unwrap(),
|
||||
ShouldTransmit::new(false)
|
||||
ShouldTransmit(false)
|
||||
);
|
||||
assert_eq!(client.local_max_data - initial_max, 48);
|
||||
assert!(!client.recv.contains_key(&id));
|
||||
@@ -1161,14 +1154,14 @@ mod tests {
|
||||
data: Bytes::from_static(&[0; 32]),
|
||||
})
|
||||
.unwrap(),
|
||||
ShouldTransmit::new(false)
|
||||
ShouldTransmit(false)
|
||||
);
|
||||
// Client stops it
|
||||
assert_eq!(
|
||||
client.stop(id).unwrap(),
|
||||
StopResult {
|
||||
max_data: ShouldTransmit::new(false),
|
||||
stop_sending: ShouldTransmit::new(true),
|
||||
max_data: ShouldTransmit(false),
|
||||
stop_sending: ShouldTransmit(true),
|
||||
}
|
||||
);
|
||||
// Server complies
|
||||
@@ -1180,7 +1173,7 @@ mod tests {
|
||||
final_offset: 32u32.into(),
|
||||
})
|
||||
.unwrap(),
|
||||
ShouldTransmit::new(false)
|
||||
ShouldTransmit(false)
|
||||
);
|
||||
assert!(!client.recv.contains_key(&id), "stream state is freed");
|
||||
}
|
||||
|
||||
@@ -138,7 +138,7 @@ impl Recv {
|
||||
// does not get stuck.
|
||||
let diff = max_stream_data - self.sent_max_stream_data;
|
||||
let transmit = self.receiving_unknown_size() && diff >= (stream_receive_window / 8);
|
||||
(max_stream_data, ShouldTransmit::new(transmit))
|
||||
(max_stream_data, ShouldTransmit(transmit))
|
||||
}
|
||||
|
||||
/// Records that a `MAX_STREAM_DATA` announcing a certain window was sent
|
||||
|
||||
Reference in New Issue
Block a user