#![allow(dead_code)]
use std::collections::{HashMap, HashSet, VecDeque};
use std::future::Future;
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::Arc;
use bytes::{Buf, Bytes, BytesMut};
use tokio::sync::mpsc::error::TrySendError;
use tokio::sync::{mpsc, oneshot};
use tokio_quiche::quic::HandshakeInfo;
use tokio_quiche::quic::QuicheConnection;
use tokio_quiche::{ApplicationOverQuic, QuicResult};
use crate::buffer::{
send_from_buf, SendAccounting, SendBytesPermit, TerminalCell, WriteCompleter, MAX_CHUNK,
PKT_BUF_LEN,
};
use crate::conn::QuicConn;
use crate::error::{
classify_stream_recv_error, classify_stream_send_error, conn_terminal_from_error, CloseOrigin,
ConnTerminal, RecvEnd, SendEnd, StreamRecvClass, StreamSendClass, H3_NO_ERROR,
H3_REQUEST_CANCELLED,
};
use crate::quiche::{self, Shutdown};
const READ_BUDGET: usize = 32;
const DISCOVERY_BUDGET: usize = 32;
const RECV_RESUME_BUDGET: usize = 16;
const READABLE_BUDGET: usize = 32;
const ADMIT_BUDGET: usize = 32;
const PROMOTE_BUDGET: usize = 32;
const CHUNK_BUDGET: usize = 16;
pub(crate) const BYTE_CHANNEL_DEPTH: usize = 64;
#[derive(Clone, Copy, Debug)]
pub(crate) struct DriverBufferConfig {
pub recv_channel_depth: usize,
pub packet_buffer_size: usize,
pub max_buffered_send_bytes: Option<usize>,
}
impl Default for DriverBufferConfig {
fn default() -> Self {
Self {
recv_channel_depth: BYTE_CHANNEL_DEPTH,
packet_buffer_size: PKT_BUF_LEN,
max_buffered_send_bytes: None,
}
}
}
const CMD_BUDGET: usize = 64;
const WRITABLE_BUDGET: usize = 32;
const WRITE_BUDGET: usize = 32;
const MAX_WRITE_CHUNK: usize = MAX_CHUNK;
const RECV_ARENA: usize = 8 * MAX_CHUNK;
const REARM_THRESHOLD: usize = 1;
const OPEN_BUDGET: usize = 32;
#[inline]
fn is_bidi(id: u64) -> bool {
id & 0x2 == 0
}
#[allow(clippy::large_enum_variant)]
pub(crate) enum DriverCommand<B: Buf> {
OpenBidi {
reply: oneshot::Sender<Result<BidiHandoff<B>, Arc<ConnTerminal>>>,
},
OpenUni {
reply: oneshot::Sender<Result<SendHandoff<B>, Arc<ConnTerminal>>>,
},
Send {
id: u64,
buf: h3::quic::WriteBuf<B>,
done: WriteCompleter<SendEnd>,
permit: Option<SendBytesPermit>,
},
Finish {
id: u64,
done: oneshot::Sender<Result<(), SendEnd>>,
},
Reset { id: u64, code: u64 },
StopSending { id: u64, code: u64 },
Close { code: u64, reason: Bytes },
RecvResume { id: u64 },
AcceptBidiResume,
AcceptUniResume,
ConnectionDropped,
}
impl<B: Buf> std::fmt::Debug for DriverCommand<B> {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
DriverCommand::OpenBidi { .. } => f.write_str("OpenBidi"),
DriverCommand::OpenUni { .. } => f.write_str("OpenUni"),
DriverCommand::Send { id, .. } => write!(f, "Send {{ id: {id} }}"),
DriverCommand::Finish { id, .. } => write!(f, "Finish {{ id: {id} }}"),
DriverCommand::Reset { id, code } => write!(f, "Reset {{ id: {id}, code: {code} }}"),
DriverCommand::StopSending { id, code } => {
write!(f, "StopSending {{ id: {id}, code: {code} }}")
}
DriverCommand::Close { code, .. } => write!(f, "Close {{ code: {code} }}"),
DriverCommand::RecvResume { id } => write!(f, "RecvResume {{ id: {id} }}"),
DriverCommand::AcceptBidiResume => f.write_str("AcceptBidiResume"),
DriverCommand::AcceptUniResume => f.write_str("AcceptUniResume"),
DriverCommand::ConnectionDropped => f.write_str("ConnectionDropped"),
}
}
}
pub(crate) struct ConnShared {
pub(crate) conn_terminal: TerminalCell<Arc<ConnTerminal>>,
pub(crate) send_accounting: Arc<SendAccounting>,
}
impl ConnShared {
pub(crate) fn new(max_buffered_send_bytes: Option<usize>) -> Arc<Self> {
Arc::new(ConnShared {
conn_terminal: TerminalCell::new(),
send_accounting: SendAccounting::new(max_buffered_send_bytes),
})
}
}
pub(crate) struct RecvHandoff<B: Buf> {
pub(crate) id: u64,
pub(crate) bytes: mpsc::Receiver<Bytes>,
pub(crate) terminal: TerminalCell<RecvEnd>,
pub(crate) resume: Arc<AtomicBool>,
pub(crate) blocked: Arc<AtomicBool>,
pub(crate) cmd_tx: mpsc::UnboundedSender<DriverCommand<B>>,
pub(crate) cleanup: HandoffCleanup<B>,
}
pub(crate) struct SendHandoff<B: Buf> {
pub(crate) id: u64,
pub(crate) status: TerminalCell<SendEnd>,
pub(crate) cmd_tx: mpsc::UnboundedSender<DriverCommand<B>>,
pub(crate) send_accounting: Arc<SendAccounting>,
pub(crate) cleanup: HandoffCleanup<B>,
}
pub(crate) struct BidiHandoff<B: Buf> {
pub(crate) send: SendHandoff<B>,
pub(crate) recv: RecvHandoff<B>,
}
pub(crate) struct HandoffCleanup<B: Buf> {
id: u64,
is_recv: bool,
cmd_tx: mpsc::UnboundedSender<DriverCommand<B>>,
armed: bool,
}
impl<B: Buf> HandoffCleanup<B> {
pub(crate) fn new(
id: u64,
is_recv: bool,
cmd_tx: mpsc::UnboundedSender<DriverCommand<B>>,
) -> Self {
HandoffCleanup {
id,
is_recv,
cmd_tx,
armed: true,
}
}
pub(crate) fn disarm(mut self) {
self.armed = false;
}
}
impl<B: Buf> Drop for HandoffCleanup<B> {
fn drop(&mut self) {
if !self.armed {
return;
}
if self.is_recv {
let _ = self.cmd_tx.send(DriverCommand::StopSending {
id: self.id,
code: 0,
});
} else {
let (done, _rx) = oneshot::channel();
let _ = self
.cmd_tx
.send(DriverCommand::Finish { id: self.id, done });
}
}
}
pub(crate) struct StreamRecvState {
pub(crate) bytes: mpsc::Sender<Bytes>,
pub(crate) terminal: TerminalCell<RecvEnd>,
pub(crate) resume: Arc<AtomicBool>,
pub(crate) blocked: Arc<AtomicBool>,
}
#[allow(clippy::large_enum_variant)]
pub(crate) enum SendOp<B: Buf> {
Write {
buf: h3::quic::WriteBuf<B>,
done: WriteCompleter<SendEnd>,
permit: Option<SendBytesPermit>,
},
Finish {
done: oneshot::Sender<Result<(), SendEnd>>,
},
}
pub(crate) struct HeldWrite<B: Buf> {
buf: h3::quic::WriteBuf<B>,
}
enum HeldFlush {
Complete,
More,
Blocked,
Terminal,
}
impl<B: Buf> SendOp<B> {
fn complete(self, result: Result<(), SendEnd>) {
match self {
SendOp::Write { done, .. } => done.complete(result),
SendOp::Finish { done } => {
let _ = done.send(result);
}
}
}
}
pub(crate) struct StreamSendState<B: Buf> {
pub(crate) send_ops: VecDeque<SendOp<B>>,
pub(crate) pending_reset: Option<u64>,
pub(crate) terminal: Option<SendEnd>,
pub(crate) status: TerminalCell<SendEnd>,
pub(crate) finished: bool,
pub(crate) held: Option<HeldWrite<B>>,
}
impl<B: Buf> StreamSendState<B> {
fn new() -> Self {
StreamSendState {
send_ops: VecDeque::new(),
pending_reset: None,
terminal: None,
status: TerminalCell::new(),
finished: false,
held: None,
}
}
}
pub(crate) enum AdmitState {
Parked(PeerStream),
Registered { send_done: bool, recv_done: bool },
}
pub(crate) struct PeerStream {
pub(crate) id: u64,
pub(crate) pending_send_terminal: Option<SendEnd>,
pub(crate) pending_recv_terminal: Option<RecvEnd>,
}
impl PeerStream {
fn new(id: u64) -> Self {
PeerStream {
id,
pending_send_terminal: None,
pending_recv_terminal: None,
}
}
}
enum AdmitResult {
Registered,
Full(PeerStream),
TornDown,
}
#[derive(Debug)]
pub(crate) enum SetupFailure {
PreHandshakeWorkerExit,
}
impl std::fmt::Display for SetupFailure {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
SetupFailure::PreHandshakeWorkerExit => write!(
f,
"connection setup failed: worker exited before the handshake completed"
),
}
}
}
impl std::error::Error for SetupFailure {}
pub(crate) struct DriverHandles<B: Buf> {
pub(crate) cmd_tx: mpsc::UnboundedSender<DriverCommand<B>>,
pub(crate) accept_bidi_rx: mpsc::Receiver<BidiHandoff<B>>,
pub(crate) accept_uni_rx: mpsc::Receiver<RecvHandoff<B>>,
pub(crate) established_rx: oneshot::Receiver<Result<(), SetupFailure>>,
pub(crate) shared: Arc<ConnShared>,
pub(crate) accept_terminal_bidi: TerminalCell<Arc<ConnTerminal>>,
pub(crate) accept_terminal_uni: TerminalCell<Arc<ConnTerminal>>,
pub(crate) accept_bidi_resume: Arc<AtomicBool>,
pub(crate) accept_uni_resume: Arc<AtomicBool>,
}
impl<B: Buf + Send + 'static> DriverHandles<B> {
pub(crate) fn into_connection(self) -> crate::stream::Connection<B> {
let opener = crate::stream::StreamOpener::from_parts(self.cmd_tx, self.shared);
crate::stream::Connection::from_parts(
self.accept_bidi_rx,
self.accept_uni_rx,
self.accept_terminal_bidi,
self.accept_terminal_uni,
self.accept_bidi_resume,
self.accept_uni_resume,
opener,
)
}
pub(crate) async fn into_established_connection(
self,
) -> Result<crate::stream::Connection<B>, SetupFailure> {
let DriverHandles {
cmd_tx,
accept_bidi_rx,
accept_uni_rx,
established_rx,
shared,
accept_terminal_bidi,
accept_terminal_uni,
accept_bidi_resume,
accept_uni_resume,
} = self;
match established_rx.await {
Ok(res) => res?,
Err(_cancelled) => return Err(SetupFailure::PreHandshakeWorkerExit),
}
let opener = crate::stream::StreamOpener::from_parts(cmd_tx, shared);
Ok(crate::stream::Connection::from_parts(
accept_bidi_rx,
accept_uni_rx,
accept_terminal_bidi,
accept_terminal_uni,
accept_bidi_resume,
accept_uni_resume,
opener,
))
}
}
pub(crate) struct PendingClose {
pub(crate) code: u64,
pub(crate) reason: Bytes,
}
pub(crate) struct RecordedLocalClose {
pub(crate) code: u64,
pub(crate) reason: Bytes,
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
enum WaitDecision {
Pending,
Yield,
Recv,
}
pub(crate) struct QuicheDriver<B: Buf = Bytes> {
shared: Arc<ConnShared>,
cmd_rx: mpsc::UnboundedReceiver<DriverCommand<B>>,
cmd_tx_weak: mpsc::WeakUnboundedSender<DriverCommand<B>>,
inbox: VecDeque<DriverCommand<B>>,
accept_bidi: mpsc::Sender<BidiHandoff<B>>,
accept_uni: mpsc::Sender<RecvHandoff<B>>,
accept_bidi_resume: Arc<AtomicBool>,
accept_uni_resume: Arc<AtomicBool>,
is_server: bool,
next_bidi_id: u64,
next_uni_id: u64,
open_bidi: VecDeque<oneshot::Sender<Result<BidiHandoff<B>, Arc<ConnTerminal>>>>,
open_uni: VecDeque<oneshot::Sender<Result<SendHandoff<B>, Arc<ConnTerminal>>>>,
recv: HashMap<u64, StreamRecvState>,
send: HashMap<u64, StreamSendState<B>>,
runnable_send: VecDeque<u64>,
runnable_send_set: HashSet<u64>,
admit: HashMap<u64, AdmitState>,
pending_readable: VecDeque<u64>,
readable_set: HashSet<u64>,
pending_resume: VecDeque<u64>,
resume_set: HashSet<u64>,
pending_admit: HashMap<u64, PeerStream>,
pending_admit_bidi: VecDeque<u64>,
pending_admit_uni: VecDeque<u64>,
parked_bidi: VecDeque<u64>,
parked_uni: VecDeque<u64>,
#[cfg(test)]
phase2_pops: usize,
established: Option<oneshot::Sender<Result<(), SetupFailure>>>,
acting: bool,
pkt_buf: Vec<u8>,
recv_channel_depth: usize,
recv_buf: Option<BytesMut>,
#[cfg(test)]
recv_lookup_count: usize,
needs_iteration: bool,
graceful_close_issued: bool,
last_handle_teardown: bool,
reads_ran_this_iter: bool,
read_budget: usize,
pending_close: Option<PendingClose>,
explicit_close_attempted: bool,
local_close: Option<RecordedLocalClose>,
close_bug: Option<&'static str>,
accept_terminal_bidi: TerminalCell<Arc<ConnTerminal>>,
accept_terminal_uni: TerminalCell<Arc<ConnTerminal>>,
conn_registration: Option<crate::endpoint::ConnRegistration>,
}
impl<B: Buf + Send + 'static> QuicheDriver<B> {
pub(crate) fn new(
is_server: bool,
accept_bidi_cap: usize,
accept_uni_cap: usize,
) -> (Self, DriverHandles<B>) {
Self::with_buffers(
is_server,
accept_bidi_cap,
accept_uni_cap,
DriverBufferConfig::default(),
)
}
pub(crate) fn with_buffers(
is_server: bool,
accept_bidi_cap: usize,
accept_uni_cap: usize,
buffers: DriverBufferConfig,
) -> (Self, DriverHandles<B>) {
let (cmd_tx, cmd_rx) = mpsc::unbounded_channel();
let cmd_tx_weak = cmd_tx.downgrade();
let (accept_bidi_tx, accept_bidi_rx) = mpsc::channel(accept_bidi_cap.max(1));
let (accept_uni_tx, accept_uni_rx) = mpsc::channel(accept_uni_cap.max(1));
let (est_tx, est_rx) = oneshot::channel();
let shared = ConnShared::new(buffers.max_buffered_send_bytes);
let accept_terminal_bidi = TerminalCell::new();
let accept_terminal_uni = TerminalCell::new();
let accept_bidi_resume = Arc::new(AtomicBool::new(false));
let accept_uni_resume = Arc::new(AtomicBool::new(false));
let driver = QuicheDriver {
shared: Arc::clone(&shared),
cmd_rx,
cmd_tx_weak,
inbox: VecDeque::new(),
accept_bidi: accept_bidi_tx,
accept_uni: accept_uni_tx,
accept_bidi_resume: Arc::clone(&accept_bidi_resume),
accept_uni_resume: Arc::clone(&accept_uni_resume),
is_server,
next_bidi_id: if is_server { 1 } else { 0 },
next_uni_id: if is_server { 3 } else { 2 },
open_bidi: VecDeque::new(),
open_uni: VecDeque::new(),
recv: HashMap::new(),
send: HashMap::new(),
runnable_send: VecDeque::new(),
runnable_send_set: HashSet::new(),
admit: HashMap::new(),
pending_readable: VecDeque::new(),
readable_set: HashSet::new(),
pending_resume: VecDeque::new(),
resume_set: HashSet::new(),
pending_admit: HashMap::new(),
pending_admit_bidi: VecDeque::new(),
pending_admit_uni: VecDeque::new(),
parked_bidi: VecDeque::new(),
parked_uni: VecDeque::new(),
#[cfg(test)]
phase2_pops: 0,
established: Some(est_tx),
acting: false,
pkt_buf: vec![0u8; buffers.packet_buffer_size.max(1)],
recv_channel_depth: buffers.recv_channel_depth,
recv_buf: None,
#[cfg(test)]
recv_lookup_count: 0,
needs_iteration: false,
graceful_close_issued: false,
last_handle_teardown: false,
reads_ran_this_iter: false,
read_budget: READ_BUDGET,
pending_close: None,
explicit_close_attempted: false,
local_close: None,
close_bug: None,
accept_terminal_bidi: accept_terminal_bidi.clone(),
accept_terminal_uni: accept_terminal_uni.clone(),
conn_registration: None,
};
let handles = DriverHandles {
cmd_tx,
accept_bidi_rx,
accept_uni_rx,
established_rx: est_rx,
shared,
accept_terminal_bidi,
accept_terminal_uni,
accept_bidi_resume,
accept_uni_resume,
};
(driver, handles)
}
pub(crate) fn set_conn_registration(&mut self, reg: crate::endpoint::ConnRegistration) {
self.conn_registration = Some(reg);
}
fn wait_decision(&self) -> WaitDecision {
if self.graceful_close_issued {
WaitDecision::Pending
} else if !self.inbox.is_empty() || self.needs_iteration {
WaitDecision::Yield
} else {
WaitDecision::Recv
}
}
fn run_read_pump<C: QuicConn>(&mut self, qconn: &mut C) {
self.intake_readable(qconn);
self.drain_resumed(qconn);
self.phase1_registered_drain(qconn);
self.phase2_admission(qconn);
self.promote_parked(qconn, true);
self.promote_parked(qconn, false);
}
fn do_process_reads<C: QuicConn>(&mut self, qconn: &mut C) {
self.needs_iteration = false;
self.read_budget = READ_BUDGET;
self.reads_ran_this_iter = true;
self.run_read_pump(qconn);
}
fn do_process_writes<C: QuicConn>(&mut self, qconn: &mut C) -> QuicResult<()> {
let no_packet = !self.reads_ran_this_iter;
if no_packet {
self.needs_iteration = false;
self.read_budget = READ_BUDGET;
}
self.apply_inbox(qconn);
if no_packet {
self.run_read_pump(qconn);
}
self.stage_open(qconn);
self.stage_writable(qconn);
self.stage_send(qconn);
self.reads_ran_this_iter = false;
self.needs_iteration |= self.has_runnable_remainder();
self.apply_close_barrier(qconn)
}
fn has_runnable_remainder(&self) -> bool {
!self.pending_resume.is_empty()
|| !self.pending_readable.is_empty()
|| (!self.pending_admit_bidi.is_empty()
&& (self.parked_bidi.is_empty()
|| self.accept_bidi_resume.load(Ordering::Relaxed)))
|| (!self.pending_admit_uni.is_empty()
&& (self.parked_uni.is_empty()
|| self.accept_uni_resume.load(Ordering::Relaxed)))
|| !self.runnable_send.is_empty()
|| (!self.parked_bidi.is_empty() && self.accept_bidi_resume.load(Ordering::Relaxed))
|| (!self.parked_uni.is_empty() && self.accept_uni_resume.load(Ordering::Relaxed))
}
fn push_pending_admit(&mut self, id: u64) {
if is_bidi(id) {
self.pending_admit_bidi.push_back(id);
} else {
self.pending_admit_uni.push_back(id);
}
}
fn apply_close_barrier<C: QuicConn>(&mut self, qconn: &mut C) -> QuicResult<()> {
if self.explicit_close_attempted {
return Ok(());
}
if let Some(pc) = self.pending_close.take() {
self.explicit_close_attempted = true;
match qconn.close(true, pc.code, &pc.reason) {
Ok(()) => {
self.local_close = Some(RecordedLocalClose {
code: pc.code,
reason: pc.reason,
});
self.graceful_close_issued = true;
}
Err(quiche::Error::Done) => {
self.graceful_close_issued = true;
}
Err(_) => {
self.close_bug = Some("explicit qconn.close returned an unexpected error");
return Err(self.close_bug.unwrap().into());
}
}
return Ok(());
}
if self.last_handle_teardown
&& qconn.peer_error().is_none()
&& qconn.local_error().is_none()
&& !qconn.is_timed_out()
&& self.local_close.is_none()
{
self.explicit_close_attempted = true;
match qconn.close(true, H3_NO_ERROR, b"") {
Ok(()) => {
self.local_close = Some(RecordedLocalClose {
code: H3_NO_ERROR,
reason: Bytes::new(),
});
self.graceful_close_issued = true;
}
Err(quiche::Error::Done) => {
self.graceful_close_issued = true;
}
Err(_) => {
self.close_bug = Some("last-handle qconn.close returned an unexpected error");
return Err(self.close_bug.unwrap().into());
}
}
}
Ok(())
}
fn intake_readable<C: QuicConn>(&mut self, qconn: &mut C) {
let mut n = 0;
while n < DISCOVERY_BUDGET {
let id = match qconn.stream_readable_next() {
Some(id) => id,
None => break,
};
n += 1;
if self.recv.contains_key(&id) {
self.requeue_readable(id);
} else if self.admit.contains_key(&id) || self.pending_admit.contains_key(&id) {
} else {
self.pending_admit.insert(id, PeerStream::new(id));
self.push_pending_admit(id);
}
}
if n == DISCOVERY_BUDGET {
self.needs_iteration = true;
}
}
fn stage_writable<C: QuicConn>(&mut self, qconn: &mut C) {
let mut n = 0;
while n < WRITABLE_BUDGET {
let id = match qconn.stream_writable_next() {
Some(id) => id,
None => break,
};
n += 1;
if self.send.contains_key(&id) {
match qconn.stream_capacity(id) {
Err(quiche::Error::StreamStopped(code)) => {
let owns_reset = self
.send
.get(&id)
.map(|s| s.pending_reset.is_some() || s.terminal.is_some())
.unwrap_or(false);
if !owns_reset {
self.send_terminal_transition(
id,
SendEnd::Stopped { error_code: code },
);
}
}
_ => {
let has_work = self
.send
.get(&id)
.map(|s| !s.send_ops.is_empty() || s.pending_reset.is_some())
.unwrap_or(false);
if has_work {
self.mark_send_runnable(id);
}
}
}
continue;
}
if !is_bidi(id) {
continue;
}
if matches!(self.admit.get(&id), Some(AdmitState::Registered { .. })) {
continue;
}
let stopped = match qconn.stream_capacity(id) {
Err(quiche::Error::StreamStopped(code)) => {
Some(SendEnd::Stopped { error_code: code })
}
_ => None,
};
if let Some(peer) = self.pending_admit.get_mut(&id) {
if peer.pending_send_terminal.is_none() {
peer.pending_send_terminal = stopped;
}
} else if let Some(AdmitState::Parked(peer)) = self.admit.get_mut(&id) {
if peer.pending_send_terminal.is_none() {
peer.pending_send_terminal = stopped;
}
} else {
let mut peer = PeerStream::new(id);
peer.pending_send_terminal = stopped;
self.pending_admit.insert(id, peer);
self.push_pending_admit(id);
self.needs_iteration = true;
}
}
if n == WRITABLE_BUDGET {
self.needs_iteration = true;
}
}
fn drain_resumed<C: QuicConn>(&mut self, qconn: &mut C) {
let mut budget = RECV_RESUME_BUDGET;
while budget > 0 {
let id = match self.pending_resume.pop_front() {
Some(id) => id,
None => break,
};
self.resume_set.remove(&id);
budget -= 1;
match self.recv.get_mut(&id) {
Some(state) => {
state.resume.store(false, Ordering::Relaxed);
state.blocked.store(false, Ordering::Relaxed);
}
None => continue,
}
self.drain_stream(qconn, id);
}
}
fn phase1_registered_drain<C: QuicConn>(&mut self, qconn: &mut C) {
let mut attempts = READABLE_BUDGET;
while attempts > 0 {
if self.read_budget == 0 {
if !self.pending_readable.is_empty() {
self.needs_iteration = true;
}
break;
}
let id = match self.pending_readable.pop_front() {
Some(id) => id,
None => break,
};
self.readable_set.remove(&id);
attempts -= 1;
if !self.recv.contains_key(&id) {
continue;
}
self.drain_stream(qconn, id);
}
}
fn drain_stream<C: QuicConn>(&mut self, qconn: &mut C, id: u64) {
let tx = match self.recv.get(&id) {
Some(state) => state.bytes.clone(),
None => return,
};
#[cfg(test)]
{
self.recv_lookup_count += 1;
}
for _ in 0..CHUNK_BUDGET {
if self.read_budget == 0 {
self.requeue_readable(id);
self.needs_iteration = true;
return;
}
let permit = match tx.try_reserve() {
Ok(permit) => permit,
Err(TrySendError::Full(())) => {
if let Some(state) = self.recv.get(&id) {
state.blocked.store(true, Ordering::Release);
}
match tx.try_reserve() {
Ok(permit) => {
if let Some(state) = self.recv.get(&id) {
state.blocked.store(false, Ordering::Release);
}
permit
}
Err(TrySendError::Full(())) => return,
Err(TrySendError::Closed(())) => {
self.abandon_recv(qconn, id);
return;
}
}
}
Err(TrySendError::Closed(())) => {
self.abandon_recv(qconn, id);
return;
}
};
let read = {
let buf = self.recv_buf.get_or_insert_with(BytesMut::new);
debug_assert_eq!(
buf.len(),
0,
"recv arena must be fully carved between reads"
);
if buf.capacity() < MAX_CHUNK {
buf.reserve(RECV_ARENA);
}
buf.resize(MAX_CHUNK, 0);
match qconn.stream_recv(id, &mut buf[..]) {
Ok((len, fin)) => {
let chunk = if len > 0 {
buf.truncate(len);
Some(buf.split_to(len).freeze())
} else {
buf.truncate(0);
None
};
Ok((chunk, fin))
}
Err(err) => {
buf.truncate(0);
Err(err)
}
}
};
match read {
Ok((chunk, fin)) => {
let had_bytes = chunk.is_some();
if let Some(chunk) = chunk {
permit.send(chunk);
self.read_budget -= 1;
} else {
drop(permit);
}
if fin {
self.publish_recv_terminal(id, RecvEnd::Fin);
self.mark_recv_done(id);
return;
}
if !had_bytes || !qconn.stream_readable(id) {
return;
}
}
Err(err) => {
drop(permit);
match classify_stream_recv_error(&err) {
StreamRecvClass::Done => return,
StreamRecvClass::Reset(code) => {
self.publish_recv_terminal(id, RecvEnd::Reset { error_code: code });
self.mark_recv_done(id);
return;
}
StreamRecvClass::ConnGone => {
self.resolve_recv_via_conn(id);
return;
}
StreamRecvClass::Bug(msg) => {
self.publish_recv_terminal(
id,
RecvEnd::Conn(Arc::new(ConnTerminal::Internal(msg))),
);
self.mark_recv_done(id);
return;
}
}
}
}
}
if qconn.stream_readable(id) {
self.requeue_readable(id);
self.needs_iteration = true;
}
}
fn phase2_admission<C: QuicConn>(&mut self, qconn: &mut C) {
let mut budget = ADMIT_BUDGET;
let mut bidi_blocked = false;
let mut uni_blocked = false;
#[cfg(test)]
{
self.phase2_pops = 0;
}
let mut prefer_bidi = true;
while budget > 0 {
let bidi_ready = !bidi_blocked && !self.pending_admit_bidi.is_empty();
let uni_ready = !uni_blocked && !self.pending_admit_uni.is_empty();
let bidi = match (bidi_ready, uni_ready) {
(false, false) => break,
(true, false) => true,
(false, true) => false,
(true, true) => prefer_bidi,
};
prefer_bidi = !bidi;
let id = if bidi {
self.pending_admit_bidi.pop_front()
} else {
self.pending_admit_uni.pop_front()
};
let id = match id {
Some(id) => id,
None => continue,
};
#[cfg(test)]
{
self.phase2_pops += 1;
}
if self.admit.contains_key(&id) || self.recv.contains_key(&id) {
self.pending_admit.remove(&id);
continue;
}
let peer = match self.pending_admit.remove(&id) {
Some(peer) => peer,
None => continue,
};
budget -= 1;
match self.admit_one(qconn, peer) {
AdmitResult::Registered => {}
AdmitResult::Full(peer) => {
self.admit.insert(id, AdmitState::Parked(peer));
if bidi {
self.parked_bidi.push_back(id);
bidi_blocked = true;
} else {
self.parked_uni.push_back(id);
uni_blocked = true;
}
}
AdmitResult::TornDown => {
if bidi {
bidi_blocked = true;
} else {
uni_blocked = true;
}
}
}
}
}
fn promote_parked<C: QuicConn>(&mut self, qconn: &mut C, bidi: bool) {
let armed = if bidi {
self.accept_bidi_resume.load(Ordering::Relaxed)
} else {
self.accept_uni_resume.load(Ordering::Relaxed)
};
if !armed {
return;
}
if bidi {
self.accept_bidi_resume.store(false, Ordering::Relaxed);
} else {
self.accept_uni_resume.store(false, Ordering::Relaxed);
}
let mut budget = PROMOTE_BUDGET;
while budget > 0 {
let id = {
let queue = if bidi {
&mut self.parked_bidi
} else {
&mut self.parked_uni
};
match queue.pop_front() {
Some(id) => id,
None => break,
}
};
budget -= 1;
let peer = match self.admit.remove(&id) {
Some(AdmitState::Parked(peer)) => peer,
Some(other) => {
self.admit.insert(id, other);
continue;
}
None => continue,
};
match self.admit_one(qconn, peer) {
AdmitResult::Registered => {}
AdmitResult::Full(peer) => {
self.admit.insert(id, AdmitState::Parked(peer));
if bidi {
self.parked_bidi.push_front(id);
} else {
self.parked_uni.push_front(id);
}
break;
}
AdmitResult::TornDown => break,
}
}
}
fn admit_one<C: QuicConn>(&mut self, qconn: &mut C, mut peer: PeerStream) -> AdmitResult {
let id = peer.id;
let bidi = is_bidi(id);
if bidi {
let tx = self.accept_bidi.clone();
let permit = match tx.try_reserve() {
Ok(permit) => permit,
Err(TrySendError::Full(())) => return AdmitResult::Full(peer),
Err(TrySendError::Closed(())) => {
self.shutdown_peer_directions(qconn, id, bidi);
return AdmitResult::TornDown;
}
};
let cmd_tx = match self.cmd_tx_weak.upgrade() {
Some(cmd_tx) => cmd_tx,
None => {
self.shutdown_peer_directions(qconn, id, bidi);
return AdmitResult::TornDown;
}
};
let (recv_state, recv_handoff, recv_done) =
self.build_recv(id, cmd_tx.clone(), peer.pending_recv_terminal.take());
let (send_handoff, send_state, send_done) = build_send(
id,
cmd_tx,
Arc::clone(&self.shared.send_accounting),
peer.pending_send_terminal.take(),
);
if let Some(state) = recv_state {
self.recv.insert(id, state);
}
if let Some(state) = send_state {
self.send.insert(id, state);
}
self.admit.insert(
id,
AdmitState::Registered {
send_done,
recv_done,
},
);
permit.send(BidiHandoff {
send: send_handoff,
recv: recv_handoff,
});
} else {
let tx = self.accept_uni.clone();
let permit = match tx.try_reserve() {
Ok(permit) => permit,
Err(TrySendError::Full(())) => return AdmitResult::Full(peer),
Err(TrySendError::Closed(())) => {
self.shutdown_peer_directions(qconn, id, bidi);
return AdmitResult::TornDown;
}
};
let cmd_tx = match self.cmd_tx_weak.upgrade() {
Some(cmd_tx) => cmd_tx,
None => {
self.shutdown_peer_directions(qconn, id, bidi);
return AdmitResult::TornDown;
}
};
let (recv_state, recv_handoff, recv_done) =
self.build_recv(id, cmd_tx, peer.pending_recv_terminal.take());
if let Some(state) = recv_state {
self.recv.insert(id, state);
}
self.admit.insert(
id,
AdmitState::Registered {
send_done: true,
recv_done,
},
);
permit.send(recv_handoff);
}
self.terminal_transition(id);
if self.recv.contains_key(&id) {
self.drain_stream(qconn, id);
}
AdmitResult::Registered
}
fn build_recv(
&self,
id: u64,
cmd_tx: mpsc::UnboundedSender<DriverCommand<B>>,
retained: Option<RecvEnd>,
) -> (Option<StreamRecvState>, RecvHandoff<B>, bool) {
let (tx, rx) = mpsc::channel(self.recv_channel_depth.max(1));
let terminal = TerminalCell::new();
let resume = Arc::new(AtomicBool::new(false));
let blocked = Arc::new(AtomicBool::new(false));
let recv_done = retained.is_some();
if let Some(end) = retained {
terminal.set(end);
}
let handoff = RecvHandoff {
id,
bytes: rx,
terminal: terminal.clone(),
resume: Arc::clone(&resume),
blocked: Arc::clone(&blocked),
cmd_tx: cmd_tx.clone(),
cleanup: HandoffCleanup::new(id, true, cmd_tx),
};
let state = if recv_done {
None
} else {
Some(StreamRecvState {
bytes: tx,
terminal,
resume,
blocked,
})
};
(state, handoff, recv_done)
}
fn publish_recv_terminal(&self, id: u64, end: RecvEnd) {
if let Some(state) = self.recv.get(&id) {
state.terminal.set(end);
}
}
fn mark_recv_done(&mut self, id: u64) {
self.recv.remove(&id);
self.drop_recv_memberships(id);
if let Some(AdmitState::Registered { recv_done, .. }) = self.admit.get_mut(&id) {
*recv_done = true;
}
self.terminal_transition(id);
self.reclaim_finished_send(id);
}
fn abandon_recv<C: QuicConn>(&mut self, qconn: &mut C, id: u64) {
let _ = qconn.stream_shutdown(id, Shutdown::Read, H3_NO_ERROR);
self.mark_recv_done(id);
}
fn resolve_recv_via_conn(&mut self, id: u64) {
match self.shared.conn_terminal.get() {
Some(terminal) => {
self.publish_recv_terminal(id, RecvEnd::Conn(terminal));
self.mark_recv_done(id);
}
None => {
self.drop_recv_memberships(id);
}
}
}
fn shutdown_peer_directions<C: QuicConn>(&self, qconn: &mut C, id: u64, bidi: bool) {
let _ = qconn.stream_shutdown(id, Shutdown::Read, H3_NO_ERROR);
if bidi {
let _ = qconn.stream_shutdown(id, Shutdown::Write, H3_NO_ERROR);
}
}
fn terminal_transition(&mut self, id: u64) {
let all_terminal = match self.admit.get(&id) {
Some(AdmitState::Registered {
send_done,
recv_done,
}) => {
if is_bidi(id) {
*send_done && *recv_done
} else {
*recv_done
}
}
_ => false,
};
if all_terminal {
self.admit.remove(&id);
self.recv.remove(&id);
self.drop_recv_memberships(id);
self.drop_send_membership(id);
}
}
fn requeue_readable(&mut self, id: u64) {
if self.readable_set.insert(id) {
self.pending_readable.push_back(id);
}
}
fn drop_recv_memberships(&mut self, id: u64) {
if self.readable_set.remove(&id) {
self.pending_readable.retain(|queued| *queued != id);
}
if self.resume_set.remove(&id) {
self.pending_resume.retain(|queued| *queued != id);
}
}
fn mark_send_runnable(&mut self, id: u64) {
if self.runnable_send_set.insert(id) {
self.runnable_send.push_back(id);
}
}
fn drop_send_membership(&mut self, id: u64) {
if self.runnable_send_set.remove(&id) {
self.runnable_send.retain(|queued| *queued != id);
}
}
fn apply_inbox<C: QuicConn>(&mut self, qconn: &mut C) {
let mut budget = CMD_BUDGET;
while budget > 0 {
let cmd = match self.inbox.pop_front() {
Some(cmd) => cmd,
None => match self.cmd_rx.try_recv() {
Ok(cmd) => cmd,
Err(_) => break,
},
};
budget -= 1;
match cmd {
DriverCommand::RecvResume { id } => self.enqueue_resume(id),
DriverCommand::AcceptBidiResume => {
self.accept_bidi_resume.store(true, Ordering::Relaxed);
}
DriverCommand::AcceptUniResume => {
self.accept_uni_resume.store(true, Ordering::Relaxed);
}
DriverCommand::Send {
id,
buf,
done,
permit,
} => {
self.enqueue_send_op(id, SendOp::Write { buf, done, permit });
}
DriverCommand::Finish { id, done } => {
self.enqueue_send_op(id, SendOp::Finish { done });
}
DriverCommand::Reset { id, code } => self.apply_reset(id, code),
DriverCommand::StopSending { id, code } => {
let _ = qconn.stream_shutdown(id, Shutdown::Read, code);
self.mark_recv_done(id);
}
DriverCommand::Close { code, reason }
if self.pending_close.is_none() && !self.explicit_close_attempted =>
{
self.pending_close = Some(PendingClose { code, reason });
}
DriverCommand::ConnectionDropped => {
self.clean_undelivered_peer_streams(qconn);
}
DriverCommand::OpenBidi { reply } => {
self.open_bidi.push_back(reply);
}
DriverCommand::OpenUni { reply } => {
self.open_uni.push_back(reply);
}
DriverCommand::Close { .. } => {}
}
}
if !self.inbox.is_empty() {
self.needs_iteration = true;
}
}
fn clean_undelivered_peer_streams<C: QuicConn>(&mut self, qconn: &mut C) {
let parked: Vec<u64> = self
.parked_bidi
.drain(..)
.chain(self.parked_uni.drain(..))
.collect();
for id in parked {
if let Some(AdmitState::Parked(_)) = self.admit.get(&id) {
self.shutdown_peer_directions(qconn, id, is_bidi(id));
self.admit.remove(&id);
}
}
let pending: Vec<u64> = self
.pending_admit_bidi
.drain(..)
.chain(self.pending_admit_uni.drain(..))
.collect();
for id in pending {
if self.pending_admit.remove(&id).is_some() {
self.shutdown_peer_directions(qconn, id, is_bidi(id));
}
}
}
fn stage_open<C: QuicConn>(&mut self, qconn: &mut C) {
self.stage_open_bidi(qconn);
self.stage_open_uni(qconn);
}
fn stage_open_bidi<C: QuicConn>(&mut self, qconn: &mut C) {
let mut budget = OPEN_BUDGET;
while budget > 0 {
let reply = match self.open_bidi.pop_front() {
Some(reply) => reply,
None => break,
};
budget -= 1;
if reply.is_closed() {
continue;
}
if qconn.peer_streams_left_bidi() == 0 {
self.open_bidi.push_front(reply);
return;
}
let id = self.next_bidi_id;
let cmd_tx = match self.cmd_tx_weak.upgrade() {
Some(cmd_tx) => cmd_tx,
None => {
let _ = reply.send(Err(self.open_terminal()));
continue;
}
};
if qconn.stream_priority(id, 127, true).is_err() {
let _ = reply.send(Err(self.open_terminal()));
continue;
}
self.next_bidi_id = id.wrapping_add(4);
let (recv_state, recv_handoff, _recv_done) = self.build_recv(id, cmd_tx.clone(), None);
let (send_handoff, send_state, _send_done) =
build_send(id, cmd_tx, Arc::clone(&self.shared.send_accounting), None);
if let Some(state) = recv_state {
self.recv.insert(id, state);
}
if let Some(state) = send_state {
self.send.insert(id, state);
}
let handoff = BidiHandoff {
send: send_handoff,
recv: recv_handoff,
};
if let Err(undelivered) = reply.send(Ok(handoff)) {
drop(undelivered);
self.cleanup_undeliverable_open(qconn, id, true);
}
}
if !self.open_bidi.is_empty() {
self.needs_iteration = true;
}
}
fn stage_open_uni<C: QuicConn>(&mut self, qconn: &mut C) {
let mut budget = OPEN_BUDGET;
while budget > 0 {
let reply = match self.open_uni.pop_front() {
Some(reply) => reply,
None => break,
};
budget -= 1;
if reply.is_closed() {
continue;
}
if qconn.peer_streams_left_uni() == 0 {
self.open_uni.push_front(reply);
return;
}
let id = self.next_uni_id;
let cmd_tx = match self.cmd_tx_weak.upgrade() {
Some(cmd_tx) => cmd_tx,
None => {
let _ = reply.send(Err(self.open_terminal()));
continue;
}
};
if qconn.stream_priority(id, 127, true).is_err() {
let _ = reply.send(Err(self.open_terminal()));
continue;
}
self.next_uni_id = id.wrapping_add(4);
let (send_handoff, send_state, _send_done) =
build_send(id, cmd_tx, Arc::clone(&self.shared.send_accounting), None);
if let Some(state) = send_state {
self.send.insert(id, state);
}
if let Err(undelivered) = reply.send(Ok(send_handoff)) {
drop(undelivered);
self.cleanup_undeliverable_open(qconn, id, false);
}
}
if !self.open_uni.is_empty() {
self.needs_iteration = true;
}
}
fn open_terminal(&self) -> Arc<ConnTerminal> {
self.shared
.conn_terminal
.get()
.unwrap_or_else(|| Arc::new(ConnTerminal::Internal("open declined without a terminal")))
}
fn cleanup_undeliverable_open<C: QuicConn>(&mut self, qconn: &mut C, id: u64, bidi: bool) {
self.send.remove(&id);
let _ = qconn.stream_shutdown(id, Shutdown::Write, H3_REQUEST_CANCELLED);
if bidi {
self.recv.remove(&id);
let _ = qconn.stream_shutdown(id, Shutdown::Read, H3_REQUEST_CANCELLED);
}
}
fn reclaim_finished_send(&mut self, id: u64) {
let finished = self.send.get(&id).map(|s| s.finished).unwrap_or(false);
if finished && !self.recv.contains_key(&id) {
self.send.remove(&id);
self.admit.remove(&id);
}
}
fn enqueue_send_op(&mut self, id: u64, op: SendOp<B>) {
let state = self.send.entry(id).or_insert_with(StreamSendState::new);
if let Some(end) = state.terminal.clone() {
op.complete(Err(end));
return;
}
state.send_ops.push_back(op);
self.mark_send_runnable(id);
}
fn apply_reset(&mut self, id: u64, code: u64) {
let state = self.send.entry(id).or_insert_with(StreamSendState::new);
if state.pending_reset.is_some() || state.terminal.is_some() {
return;
}
state.pending_reset = Some(code);
let end = SendEnd::Reset { error_code: code };
state.status.set(end.clone());
if state.terminal.is_none() {
state.terminal = Some(end);
}
let terminal = state.terminal.clone().expect("terminal just set");
let ops: Vec<SendOp<B>> = state.send_ops.drain(..).collect();
state.held = None;
for op in ops {
op.complete(Err(terminal.clone()));
}
self.mark_send_runnable(id);
}
fn send_terminal_transition(&mut self, id: u64, end: SendEnd) {
let ops: Vec<SendOp<B>> = match self.send.get_mut(&id) {
Some(state) => {
state.status.set(end.clone());
if state.terminal.is_none() {
state.terminal = Some(end.clone());
}
state.held = None;
state.send_ops.drain(..).collect()
}
None => return,
};
let terminal = self
.send
.get(&id)
.and_then(|s| s.terminal.clone())
.unwrap_or(end);
for op in ops {
op.complete(Err(terminal.clone()));
}
self.mark_send_done(id);
}
fn mark_send_done(&mut self, id: u64) {
if let Some(state) = self.send.get_mut(&id) {
state.finished = true;
}
self.drop_send_membership(id);
if let Some(AdmitState::Registered { send_done, .. }) = self.admit.get_mut(&id) {
*send_done = true;
}
self.terminal_transition(id);
}
fn stage_send<C: QuicConn>(&mut self, qconn: &mut C) {
let mut turns = WRITE_BUDGET;
while turns > 0 {
let id = match self.runnable_send.pop_front() {
Some(id) => id,
None => break,
};
self.runnable_send_set.remove(&id);
turns -= 1;
match self.service_send_turn(qconn, id) {
TurnOutcome::Requeue => {
self.mark_send_runnable(id);
self.needs_iteration = true;
}
TurnOutcome::Park | TurnOutcome::Drop => {}
}
}
if !self.runnable_send.is_empty() {
self.needs_iteration = true;
}
}
fn service_send_turn<C: QuicConn>(&mut self, qconn: &mut C, id: u64) -> TurnOutcome {
let pending_reset = self.send.get(&id).and_then(|s| s.pending_reset);
if let Some(code) = pending_reset {
if let Some(state) = self.send.get_mut(&id) {
state.pending_reset = None;
}
let _ = qconn.stream_shutdown(id, Shutdown::Write, code);
self.mark_send_done(id);
return TurnOutcome::Drop;
}
let has_ops = match self.send.get(&id) {
Some(state) => state.terminal.is_none() && !state.send_ops.is_empty(),
None => false,
};
if !has_ops {
return TurnOutcome::Drop;
}
let is_write = matches!(
self.send.get(&id).and_then(|s| s.send_ops.front()),
Some(SendOp::Write { .. })
);
if is_write {
self.service_write_turn(qconn, id)
} else {
self.service_finish_turn(qconn, id)
}
}
fn flush_held_once<C: QuicConn>(&mut self, qconn: &mut C, id: u64, fin: bool) -> HeldFlush {
let remaining = match self.send.get(&id).and_then(|s| s.held.as_ref()) {
Some(held) => held.buf.remaining(),
None => return HeldFlush::Complete,
};
let mut offered = 0usize;
let result = {
let held = self
.send
.get_mut(&id)
.and_then(|s| s.held.as_mut())
.expect("held present");
send_from_buf(&mut held.buf, |chunk| {
let n = chunk.len().min(MAX_WRITE_CHUNK);
offered = n;
let last = fin && n == remaining;
qconn.stream_send(id, &chunk[..n], last)
})
};
match result {
Ok(written) => {
let drained = self
.send
.get(&id)
.and_then(|s| s.held.as_ref())
.map(|held| !held.buf.has_remaining())
.unwrap_or(true);
if drained {
if let Some(state) = self.send.get_mut(&id) {
state.held = None;
}
HeldFlush::Complete
} else if written == offered {
HeldFlush::More
} else {
match self.rearm_send(qconn, id) {
TurnOutcome::Drop => HeldFlush::Terminal,
_ => HeldFlush::Blocked,
}
}
}
Err(err) => match self.classify_send_err(qconn, id, &err) {
TurnOutcome::Park => HeldFlush::Blocked,
TurnOutcome::Requeue => HeldFlush::More,
TurnOutcome::Drop => HeldFlush::Terminal,
},
}
}
fn service_write_turn<C: QuicConn>(&mut self, qconn: &mut C, id: u64) -> TurnOutcome {
if self.send.get(&id).is_some_and(|s| s.held.is_some()) {
match self.flush_held_once(qconn, id, false) {
HeldFlush::Complete => {} HeldFlush::More => return TurnOutcome::Requeue,
HeldFlush::Blocked => return TurnOutcome::Park,
HeldFlush::Terminal => return TurnOutcome::Drop,
}
}
let remaining = match self.send.get(&id).and_then(|s| s.send_ops.front()) {
Some(SendOp::Write { buf, .. }) => buf.remaining(),
_ => return TurnOutcome::Drop,
};
let cap = match qconn.stream_capacity(id) {
Ok(cap) => cap,
Err(err) => return self.classify_send_err(qconn, id, &err),
};
if remaining > 0 && cap >= remaining {
if let Some(state) = self.send.get_mut(&id) {
if let Some(SendOp::Write { buf, done, permit }) = state.send_ops.pop_front() {
done.complete(Ok(()));
drop(permit); state.held = Some(HeldWrite { buf });
}
}
return self.runnable_after_pop(id);
}
let mut offered = 0usize;
let result = {
let state = self.send.get_mut(&id).expect("send state present");
match state.send_ops.front_mut() {
Some(SendOp::Write { buf, .. }) => send_from_buf(buf, |chunk| {
let n = chunk.len().min(MAX_WRITE_CHUNK);
offered = n;
qconn.stream_send(id, &chunk[..n], false)
}),
_ => return TurnOutcome::Drop,
}
};
match result {
Ok(written) => {
let has_remaining = self
.send
.get(&id)
.and_then(|s| s.send_ops.front())
.map(|op| match op {
SendOp::Write { buf, .. } => buf.has_remaining(),
SendOp::Finish { .. } => false,
})
.unwrap_or(false);
if !has_remaining {
if let Some(state) = self.send.get_mut(&id) {
if let Some(op) = state.send_ops.pop_front() {
op.complete(Ok(()));
}
}
return self.runnable_after_pop(id);
}
if written == offered {
TurnOutcome::Requeue
} else {
self.rearm_send(qconn, id)
}
}
Err(err) => self.classify_send_err(qconn, id, &err),
}
}
fn service_finish_turn<C: QuicConn>(&mut self, qconn: &mut C, id: u64) -> TurnOutcome {
if self.send.get(&id).is_some_and(|s| s.held.is_some()) {
match self.flush_held_once(qconn, id, true) {
HeldFlush::Complete => {
self.complete_finish_op(id);
TurnOutcome::Drop
}
HeldFlush::More => TurnOutcome::Requeue,
HeldFlush::Blocked => TurnOutcome::Park,
HeldFlush::Terminal => TurnOutcome::Drop,
}
} else {
match qconn.stream_send(id, &[], true) {
Ok(_) => {
self.complete_finish_op(id);
TurnOutcome::Drop
}
Err(err) => self.classify_send_err(qconn, id, &err),
}
}
}
fn complete_finish_op(&mut self, id: u64) {
if let Some(state) = self.send.get_mut(&id) {
if let Some(op) = state.send_ops.pop_front() {
op.complete(Ok(()));
}
}
self.mark_send_done(id);
self.reclaim_finished_send(id);
}
fn classify_send_err<C: QuicConn>(
&mut self,
qconn: &mut C,
id: u64,
err: &quiche::Error,
) -> TurnOutcome {
match classify_stream_send_error(err) {
StreamSendClass::Stopped(code) => {
self.send_terminal_transition(id, SendEnd::Stopped { error_code: code });
TurnOutcome::Drop
}
StreamSendClass::Blocked => self.rearm_send(qconn, id),
StreamSendClass::ConnGone => {
match self.shared.conn_terminal.get() {
Some(t) => self.send_terminal_transition(id, SendEnd::Conn(t)),
None => self.drop_send_membership(id),
}
TurnOutcome::Drop
}
StreamSendClass::Limit | StreamSendClass::Bug(_) => {
let end = SendEnd::Conn(Arc::new(ConnTerminal::Internal(
"unexpected stream_send error",
)));
self.send_terminal_transition(id, end);
TurnOutcome::Drop
}
}
}
fn rearm_send<C: QuicConn>(&mut self, qconn: &mut C, id: u64) -> TurnOutcome {
match qconn.stream_writable(id, REARM_THRESHOLD) {
Err(quiche::Error::StreamStopped(code)) => {
self.send_terminal_transition(id, SendEnd::Stopped { error_code: code });
TurnOutcome::Drop
}
_ => TurnOutcome::Park,
}
}
fn runnable_after_pop(&mut self, id: u64) -> TurnOutcome {
let more = self
.send
.get(&id)
.map(|s| s.terminal.is_none() && !s.send_ops.is_empty())
.unwrap_or(false);
if more {
TurnOutcome::Requeue
} else {
TurnOutcome::Drop
}
}
fn enqueue_resume(&mut self, id: u64) {
if self.resume_set.insert(id) {
self.pending_resume.push_back(id);
}
}
fn classify_conn_terminal<C: QuicConn>(&self, qconn: &C) -> ConnTerminal {
if let Some(msg) = self.close_bug {
return ConnTerminal::Internal(msg);
}
if let Some(pe) = qconn.peer_error() {
return conn_terminal_from_error(CloseOrigin::Peer, pe);
}
if qconn.is_timed_out() {
return ConnTerminal::Timeout;
}
if let Some(lc) = &self.local_close {
return ConnTerminal::AppClose {
origin: CloseOrigin::Local,
error_code: lc.code,
reason: lc.reason.clone(),
};
}
if let Some(le) = qconn.local_error() {
return conn_terminal_from_error(CloseOrigin::Local, le);
}
ConnTerminal::Internal("connection closed without a recorded terminal")
}
fn do_on_conn_close<C: QuicConn>(&mut self, qconn: &mut C) {
let terminal = Arc::new(self.classify_conn_terminal(qconn));
self.shared.conn_terminal.set(Arc::clone(&terminal));
self.accept_terminal_bidi.set(Arc::clone(&terminal));
self.accept_terminal_uni.set(Arc::clone(&terminal));
for state in self.recv.values() {
state.terminal.set(RecvEnd::Conn(Arc::clone(&terminal)));
}
let mut pending_ops: Vec<SendOp<B>> = Vec::new();
for state in self.send.values_mut() {
let end = SendEnd::Conn(Arc::clone(&terminal));
state.status.set(end.clone());
if state.terminal.is_none() {
state.terminal = Some(end);
}
state.held = None;
pending_ops.extend(state.send_ops.drain(..));
}
for op in pending_ops {
op.complete(Err(SendEnd::Conn(Arc::clone(&terminal))));
}
self.cmd_rx.close();
loop {
let cmd = match self.inbox.pop_front() {
Some(cmd) => cmd,
None => match self.cmd_rx.try_recv() {
Ok(cmd) => cmd,
Err(_) => break,
},
};
self.complete_command_on_close(cmd, &terminal);
}
for reply in self.open_bidi.drain(..) {
let _ = reply.send(Err(Arc::clone(&terminal)));
}
for reply in self.open_uni.drain(..) {
let _ = reply.send(Err(Arc::clone(&terminal)));
}
}
fn complete_command_on_close(&self, cmd: DriverCommand<B>, terminal: &Arc<ConnTerminal>) {
match cmd {
DriverCommand::Send { done, .. } => {
done.complete(Err(SendEnd::Conn(Arc::clone(terminal))));
}
DriverCommand::Finish { done, .. } => {
let _ = done.send(Err(SendEnd::Conn(Arc::clone(terminal))));
}
DriverCommand::OpenBidi { reply } => {
let _ = reply.send(Err(Arc::clone(terminal)));
}
DriverCommand::OpenUni { reply } => {
let _ = reply.send(Err(Arc::clone(terminal)));
}
DriverCommand::Reset { .. }
| DriverCommand::StopSending { .. }
| DriverCommand::RecvResume { .. }
| DriverCommand::AcceptBidiResume
| DriverCommand::AcceptUniResume
| DriverCommand::ConnectionDropped
| DriverCommand::Close { .. } => {}
}
}
}
enum TurnOutcome {
Requeue,
Park,
Drop,
}
fn build_send<B: Buf>(
id: u64,
cmd_tx: mpsc::UnboundedSender<DriverCommand<B>>,
send_accounting: Arc<SendAccounting>,
retained: Option<SendEnd>,
) -> (SendHandoff<B>, Option<StreamSendState<B>>, bool) {
let status = TerminalCell::new();
let send_done = retained.is_some();
if let Some(ref end) = retained {
status.set(end.clone());
}
let handoff = SendHandoff {
id,
status: status.clone(),
cmd_tx: cmd_tx.clone(),
send_accounting,
cleanup: HandoffCleanup::new(id, false, cmd_tx),
};
let state = if send_done {
None
} else {
Some(StreamSendState {
send_ops: VecDeque::new(),
pending_reset: None,
terminal: None,
status,
finished: false,
held: None,
})
};
(handoff, state, send_done)
}
impl<B: Buf + Send + 'static> ApplicationOverQuic for QuicheDriver<B> {
fn on_conn_established(
&mut self,
_qconn: &mut QuicheConnection,
_handshake_info: &HandshakeInfo,
) -> QuicResult<()> {
self.acting = true;
if let Some(tx) = self.established.take() {
let _ = tx.send(Ok(()));
}
Ok(())
}
fn should_act(&self) -> bool {
self.acting
}
fn buffer(&mut self) -> &mut [u8] {
&mut self.pkt_buf
}
#[allow(clippy::manual_async_fn)]
fn wait_for_data(
&mut self,
_qconn: &mut QuicheConnection,
) -> impl Future<Output = QuicResult<()>> + Send {
async move {
match self.wait_decision() {
WaitDecision::Pending => std::future::pending::<QuicResult<()>>().await,
WaitDecision::Yield => {
tokio::task::yield_now().await;
Ok(())
}
WaitDecision::Recv => match self.cmd_rx.recv().await {
Some(cmd) => {
self.inbox.push_back(cmd);
Ok(())
}
None => {
self.last_handle_teardown = true;
Ok(())
}
},
}
}
}
fn process_reads(&mut self, qconn: &mut QuicheConnection) -> QuicResult<()> {
self.do_process_reads(qconn);
Ok(())
}
fn process_writes(&mut self, qconn: &mut QuicheConnection) -> QuicResult<()> {
self.do_process_writes(qconn)
}
fn on_conn_close<M: tokio_quiche::metrics::Metrics>(
&mut self,
qconn: &mut QuicheConnection,
_metrics: &M,
_connection_result: &QuicResult<()>,
) {
self.do_on_conn_close(qconn);
}
}
impl<B: Buf> Drop for QuicheDriver<B> {
fn drop(&mut self) {
if let Some(tx) = self.established.take() {
let _ = tx.send(Err(SetupFailure::PreHandshakeWorkerExit));
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::buffer::WriteOutcome;
fn driver() -> (QuicheDriver<Bytes>, DriverHandles<Bytes>) {
QuicheDriver::<Bytes>::new(false, 4, 4)
}
#[test]
fn wait_decision_idle_awaits_commands() {
let (d, _h) = driver();
assert_eq!(d.wait_decision(), WaitDecision::Recv);
}
#[test]
fn wait_decision_yields_when_inbox_nonempty() {
let (mut d, _h) = driver();
d.inbox.push_back(DriverCommand::AcceptBidiResume);
assert_eq!(d.wait_decision(), WaitDecision::Yield);
}
#[test]
fn wait_decision_yields_when_needs_iteration() {
let (mut d, _h) = driver();
d.needs_iteration = true;
assert_eq!(d.wait_decision(), WaitDecision::Yield);
}
#[test]
fn wait_decision_pending_after_graceful_close() {
let (mut d, _h) = driver();
d.needs_iteration = true;
d.graceful_close_issued = true;
assert_eq!(d.wait_decision(), WaitDecision::Pending);
}
#[test]
fn should_act_false_until_established() {
let (d, _h) = driver();
assert!(!d.should_act());
}
#[test]
fn buffer_is_packet_sized_not_chunk_sized() {
let (mut d, _h) = driver();
assert_eq!(d.buffer().len(), PKT_BUF_LEN);
assert_ne!(PKT_BUF_LEN, MAX_CHUNK);
}
#[test]
fn sf4_sf5_default_buffer_sizes_unchanged() {
let (mut d, _h) = QuicheDriver::<Bytes>::new(false, 4, 4);
assert_eq!(d.recv_channel_depth, BYTE_CHANNEL_DEPTH);
assert_eq!(d.buffer().len(), PKT_BUF_LEN);
}
#[test]
fn sf4_sf5_buffer_overrides_take_effect() {
let (mut d, h) = QuicheDriver::<Bytes>::with_buffers(
false,
4,
4,
DriverBufferConfig {
recv_channel_depth: 8,
packet_buffer_size: 4096,
max_buffered_send_bytes: None,
},
);
assert_eq!(d.recv_channel_depth, 8);
assert_eq!(d.buffer().len(), 4096);
let cmd_tx = h.cmd_tx.clone();
let (state, _handoff, _done) = d.build_recv(0, cmd_tx, None);
let state = state.expect("live recv state");
assert_eq!(state.bytes.max_capacity(), 8);
}
#[test]
fn sf4_sf5_zero_sizes_clamped_to_one() {
let (mut d, h) = QuicheDriver::<Bytes>::with_buffers(
false,
4,
4,
DriverBufferConfig {
recv_channel_depth: 0,
packet_buffer_size: 0,
max_buffered_send_bytes: None,
},
);
assert_eq!(d.buffer().len(), 1);
let (state, _handoff, _done) = d.build_recv(0, h.cmd_tx.clone(), None);
assert_eq!(state.expect("live recv state").bytes.max_capacity(), 1);
}
#[test]
fn dropping_driver_before_handshake_reports_setup_failure() {
let (d, mut h) = driver();
drop(d);
match h.established_rx.try_recv() {
Ok(Err(SetupFailure::PreHandshakeWorkerExit)) => {}
other => panic!("expected PreHandshakeWorkerExit, got {other:?}"),
}
}
#[test]
fn last_strong_sender_drop_closes_cmd_rx() {
let (mut d, h) = driver();
assert!(d.cmd_tx_weak.upgrade().is_some());
drop(h); assert!(d.cmd_tx_weak.upgrade().is_none());
assert!(d.cmd_rx.try_recv().is_err());
}
use crate::conn::mock::{MockConn, RecvStep};
fn data(bytes: &[u8], fin: bool) -> RecvStep {
RecvStep::Data {
bytes: bytes.to_vec(),
fin,
}
}
#[test]
fn buffered_bytes_then_fin_delivers_then_seals() {
let (mut d, mut h) = driver();
let mut c = MockConn::new();
c.script_recv(0, [data(b"hello", true)]);
c.queue_readable([0]);
d.read_budget = READ_BUDGET;
d.run_read_pump(&mut c);
let mut ho = h.accept_bidi_rx.try_recv().expect("one bidi handoff");
assert_eq!(ho.recv.id, 0);
assert_eq!(
ho.recv.bytes.try_recv().unwrap(),
Bytes::from_static(b"hello")
);
assert!(matches!(ho.recv.terminal.get(), Some(RecvEnd::Fin)));
assert!(ho.recv.bytes.try_recv().is_err());
assert!(h.accept_bidi_rx.try_recv().is_err());
}
#[test]
fn sf1_delivered_chunk_shares_recv_buffer_backing() {
let (mut d, mut h) = driver();
let mut c = MockConn::new();
c.script_recv(0, [data(b"zerocopy-payload", false)]);
c.queue_readable([0]);
d.read_budget = READ_BUDGET;
d.run_read_pump(&mut c);
let mut ho = h.accept_bidi_rx.try_recv().expect("one bidi handoff");
let chunk = ho.recv.bytes.try_recv().unwrap();
assert_eq!(&chunk[..], b"zerocopy-payload");
let buf = d
.recv_buf
.as_ref()
.expect("recv buffer allocated on first read");
assert_eq!(
chunk.as_ptr() as usize + chunk.len(),
buf.as_ptr() as usize,
"delivered chunk must be contiguous with the retained buffer (shared backing)"
);
}
#[test]
fn sf1_multi_chunk_drain_no_aliasing_corruption() {
let (mut d, mut h) = driver();
let mut c = MockConn::new();
c.script_recv(0, [data(b"AAAA", false), data(b"BBBB", true)]);
c.queue_readable([0]);
d.read_budget = READ_BUDGET;
d.run_read_pump(&mut c);
let mut ho = h.accept_bidi_rx.try_recv().expect("one bidi handoff");
let c1 = ho.recv.bytes.try_recv().unwrap();
let c2 = ho.recv.bytes.try_recv().unwrap();
assert_eq!(&c1[..], b"AAAA");
assert_eq!(&c2[..], b"BBBB");
assert!(matches!(ho.recv.terminal.get(), Some(RecvEnd::Fin)));
}
#[test]
fn sf1_receive_arena_amortizes_allocation() {
let (mut d, _h) = driver();
let (tx, mut rx) = mpsc::channel::<Bytes>(BYTE_CHANNEL_DEPTH);
d.recv.insert(
0,
StreamRecvState {
bytes: tx,
terminal: TerminalCell::new(),
resume: Arc::new(AtomicBool::new(false)),
blocked: Arc::new(AtomicBool::new(false)),
},
);
let mut c = MockConn::new();
c.script_recv(
0,
[
data(b"aaaa", false),
data(b"bbbb", false),
data(b"cccc", false),
data(b"dddd", false),
],
);
d.read_budget = READ_BUDGET;
d.drain_stream(&mut c, 0);
let mut held = Vec::new();
while let Ok(chunk) = rx.try_recv() {
held.push(chunk);
}
assert_eq!(held.len(), 4, "all four chunks drained in one pass");
let mut prev_end: Option<usize> = None;
for chunk in &held {
let start = chunk.as_ptr() as usize;
if let Some(pe) = prev_end {
assert_eq!(
pe, start,
"chunk not contiguous with the previous → arena reallocated \
per chunk (allocation not amortized)"
);
}
prev_end = Some(start + chunk.len());
}
}
#[test]
fn sf7_sender_lookup_hoisted_once_per_drain() {
let (mut d, _h) = driver();
let (tx, mut rx) = mpsc::channel::<Bytes>(BYTE_CHANNEL_DEPTH);
d.recv.insert(
0,
StreamRecvState {
bytes: tx,
terminal: TerminalCell::new(),
resume: Arc::new(AtomicBool::new(false)),
blocked: Arc::new(AtomicBool::new(false)),
},
);
let mut c = MockConn::new();
c.script_recv(0, [data(b"a", false), data(b"b", false), data(b"c", false)]);
d.read_budget = READ_BUDGET;
d.drain_stream(&mut c, 0);
assert_eq!(d.recv_lookup_count, 1);
assert_eq!(rx.try_recv().unwrap(), Bytes::from_static(b"a"));
assert_eq!(rx.try_recv().unwrap(), Bytes::from_static(b"b"));
assert_eq!(rx.try_recv().unwrap(), Bytes::from_static(b"c"));
}
#[test]
fn sf5_recv_buffer_lazily_allocated() {
let (mut d, mut h) = driver();
assert!(
d.recv_buf.is_none(),
"idle connection must not allocate the receive buffer"
);
let mut c = MockConn::new();
c.script_recv(0, [data(b"x", true)]);
c.queue_readable([0]);
d.read_budget = READ_BUDGET;
d.run_read_pump(&mut c);
let _ho = h.accept_bidi_rx.try_recv().expect("one bidi handoff");
assert!(
d.recv_buf.is_some(),
"receive buffer must materialize on first stream_recv"
);
}
#[test]
fn queued_bytes_then_reset_delivers_then_seals() {
let (mut d, mut h) = driver();
let mut c = MockConn::new();
c.script_recv(
0,
[
data(b"data", false),
RecvStep::Err(crate::quiche::Error::StreamReset(7)),
],
);
c.queue_readable([0]);
d.read_budget = READ_BUDGET;
d.run_read_pump(&mut c);
let mut ho = h.accept_bidi_rx.try_recv().expect("one bidi handoff");
assert_eq!(
ho.recv.bytes.try_recv().unwrap(),
Bytes::from_static(b"data")
);
assert!(matches!(
ho.recv.terminal.get(),
Some(RecvEnd::Reset { error_code: 7 })
));
assert!(ho.recv.bytes.try_recv().is_err());
}
#[test]
fn reserve_before_read_full_channel_then_resume() {
let (mut d, _h) = driver();
let (tx, mut rx) = mpsc::channel::<Bytes>(BYTE_CHANNEL_DEPTH);
for _ in 0..BYTE_CHANNEL_DEPTH {
tx.try_send(Bytes::from_static(b"x")).unwrap();
}
let terminal = TerminalCell::new();
let resume = Arc::new(AtomicBool::new(false));
d.recv.insert(
0,
StreamRecvState {
bytes: tx,
terminal: terminal.clone(),
resume,
blocked: Arc::new(AtomicBool::new(false)),
},
);
d.admit.insert(
0,
AdmitState::Registered {
send_done: false,
recv_done: false,
},
);
d.pending_readable.push_back(0);
d.readable_set.insert(0);
let mut c = MockConn::new();
c.script_recv(0, [data(b"late", true)]);
d.read_budget = READ_BUDGET;
d.run_read_pump(&mut c);
assert!(c.recv_calls.is_empty());
assert!(d.recv.get(&0).unwrap().blocked.load(Ordering::Relaxed));
for _ in 0..BYTE_CHANNEL_DEPTH {
rx.try_recv().unwrap();
}
d.inbox.push_back(DriverCommand::RecvResume { id: 0 });
d.apply_inbox(&mut c);
d.read_budget = READ_BUDGET;
d.run_read_pump(&mut c);
assert_eq!(c.recv_calls, vec![0]);
assert_eq!(rx.try_recv().unwrap(), Bytes::from_static(b"late"));
assert!(matches!(terminal.get(), Some(RecvEnd::Fin)));
}
#[test]
fn sf2_worker_recheck_does_not_park_with_free_slot() {
let (mut d, _h) = driver();
let (tx, mut rx) = mpsc::channel::<Bytes>(BYTE_CHANNEL_DEPTH);
let blocked = Arc::new(AtomicBool::new(true));
let terminal = TerminalCell::new();
d.recv.insert(
0,
StreamRecvState {
bytes: tx,
terminal: terminal.clone(),
resume: Arc::new(AtomicBool::new(false)),
blocked: Arc::clone(&blocked),
},
);
d.admit.insert(
0,
AdmitState::Registered {
send_done: false,
recv_done: false,
},
);
d.pending_readable.push_back(0);
d.readable_set.insert(0);
let mut c = MockConn::new();
c.script_recv(0, [data(b"ok", true)]);
d.read_budget = READ_BUDGET;
d.run_read_pump(&mut c);
assert_eq!(c.recv_calls, vec![0]);
assert_eq!(rx.try_recv().unwrap(), Bytes::from_static(b"ok"));
assert!(matches!(terminal.get(), Some(RecvEnd::Fin)));
}
#[test]
fn sf2_worker_publishes_blocked_on_full_channel() {
let (mut d, _h) = driver();
let (tx, _rx) = mpsc::channel::<Bytes>(BYTE_CHANNEL_DEPTH);
for _ in 0..BYTE_CHANNEL_DEPTH {
tx.try_send(Bytes::from_static(b"x")).unwrap();
}
let blocked = Arc::new(AtomicBool::new(false));
d.recv.insert(
0,
StreamRecvState {
bytes: tx,
terminal: TerminalCell::new(),
resume: Arc::new(AtomicBool::new(false)),
blocked: Arc::clone(&blocked),
},
);
d.admit.insert(
0,
AdmitState::Registered {
send_done: false,
recv_done: false,
},
);
d.pending_readable.push_back(0);
d.readable_set.insert(0);
let mut c = MockConn::new();
c.script_recv(0, [data(b"late", true)]);
d.read_budget = READ_BUDGET;
d.run_read_pump(&mut c);
assert!(c.recv_calls.is_empty());
assert!(
blocked.load(Ordering::Acquire),
"park flag must be published"
);
}
#[test]
fn destructive_intake_admits_new_peer_once() {
let (mut d, mut h) = driver();
let mut c = MockConn::new();
c.script_recv(0, [data(b"hi", false)]);
c.queue_readable([0]);
d.read_budget = READ_BUDGET;
d.run_read_pump(&mut c);
let mut ho = h.accept_bidi_rx.try_recv().expect("admitted once");
assert_eq!(ho.recv.bytes.try_recv().unwrap(), Bytes::from_static(b"hi"));
assert!(matches!(
d.admit.get(&0),
Some(AdmitState::Registered { .. })
));
d.read_budget = READ_BUDGET;
d.run_read_pump(&mut c);
assert!(h.accept_bidi_rx.try_recv().is_err());
}
#[test]
fn parked_stream_single_admit_then_promote() {
let (mut d, mut h) = QuicheDriver::<Bytes>::new(false, 1, 1);
let mut c = MockConn::new();
c.script_recv(0, [data(b"a", false)]);
c.script_recv(4, [data(b"b", false)]);
c.queue_readable([0, 4]);
d.read_budget = READ_BUDGET;
d.run_read_pump(&mut c);
assert_eq!(d.parked_bidi.len(), 1);
assert_eq!(*d.parked_bidi.front().unwrap(), 4);
assert!(matches!(d.admit.get(&4), Some(AdmitState::Parked(_))));
assert!(matches!(
d.admit.get(&0),
Some(AdmitState::Registered { .. })
));
let ho0 = h.accept_bidi_rx.try_recv().expect("stream 0 handoff");
assert_eq!(ho0.recv.id, 0);
d.inbox.push_back(DriverCommand::AcceptBidiResume);
d.apply_inbox(&mut c);
d.read_budget = READ_BUDGET;
d.run_read_pump(&mut c);
let ho4 = h.accept_bidi_rx.try_recv().expect("stream 4 promoted");
assert_eq!(ho4.recv.id, 4);
assert!(d.parked_bidi.is_empty());
assert!(matches!(
d.admit.get(&4),
Some(AdmitState::Registered { .. })
));
assert!(h.accept_bidi_rx.try_recv().is_err()); }
#[test]
fn mf1_parked_class_bounds_scan_and_quiesces_worker() {
let (mut d, _h) = QuicheDriver::<Bytes>::new(false, 1, 1);
let mut c = MockConn::new();
d.pending_admit.insert(0, PeerStream::new(0));
d.pending_admit_bidi.push_back(0);
d.phase2_admission(&mut c);
assert!(matches!(
d.admit.get(&0),
Some(AdmitState::Registered { .. })
));
const N: u64 = 64;
for k in 1..=N {
let id = k * 4; assert!(is_bidi(id));
d.pending_admit.insert(id, PeerStream::new(id));
d.pending_admit_bidi.push_back(id);
}
assert_eq!(d.pending_admit_bidi.len() as u64, N);
d.phase2_admission(&mut c);
assert_eq!(
d.phase2_pops, 1,
"bounded scan: exactly one id examined, not the whole backlog (MF-1)"
);
assert_eq!(d.parked_bidi.len(), 1, "exactly one id parked");
assert_eq!(
*d.parked_bidi.front().unwrap(),
4,
"FIFO: id 4 parked first"
);
assert_eq!(
d.pending_admit_bidi.len() as u64,
N - 1,
"blocked-class backlog left in place, not rescanned (MF-1)"
);
assert_eq!(*d.pending_admit_bidi.front().unwrap(), 8);
assert!(
!d.has_runnable_remainder(),
"capacity-blocked class must not keep the worker runnable"
);
d.phase2_admission(&mut c);
assert_eq!(
d.phase2_pops, 1,
"repeated pumps stay O(1) while the class is accept-blocked"
);
assert!(!d.has_runnable_remainder());
d.accept_bidi_resume.store(true, Ordering::Relaxed);
assert!(
d.has_runnable_remainder(),
"AcceptBidiResume re-arms the worker for the parked class"
);
}
#[test]
fn mf1_parked_class_resumes_in_fifo_order() {
let (mut d, mut h) = QuicheDriver::<Bytes>::new(false, 1, 1);
let mut c = MockConn::new();
for id in [0u64, 4, 8] {
d.pending_admit.insert(id, PeerStream::new(id));
d.pending_admit_bidi.push_back(id);
}
d.phase2_admission(&mut c);
assert!(matches!(
d.admit.get(&0),
Some(AdmitState::Registered { .. })
));
assert_eq!(d.parked_bidi.front().copied(), Some(4));
assert_eq!(d.pending_admit_bidi.front().copied(), Some(8));
let ho0 = h.accept_bidi_rx.try_recv().expect("id 0 handoff");
assert_eq!(ho0.recv.id, 0);
d.accept_bidi_resume.store(true, Ordering::Relaxed);
d.promote_parked(&mut c, true);
let ho4 = h.accept_bidi_rx.try_recv().expect("id 4 promoted first");
assert_eq!(ho4.recv.id, 4, "FIFO: earliest-parked id promoted first");
assert!(matches!(
d.admit.get(&4),
Some(AdmitState::Registered { .. })
));
assert_eq!(d.pending_admit_bidi.front().copied(), Some(8));
}
#[test]
fn mf1_blocked_class_does_not_starve_other_class() {
let (mut d, _h) = QuicheDriver::<Bytes>::new(false, 1, 16);
let mut c = MockConn::new();
d.pending_admit.insert(0, PeerStream::new(0));
d.pending_admit_bidi.push_back(0);
d.phase2_admission(&mut c); assert!(matches!(
d.admit.get(&0),
Some(AdmitState::Registered { .. })
));
for k in 1..=32u64 {
let id = k * 4; d.pending_admit.insert(id, PeerStream::new(id));
d.pending_admit_bidi.push_back(id);
}
d.pending_admit.insert(2, PeerStream::new(2)); assert!(!is_bidi(2));
d.pending_admit_uni.push_back(2);
d.phase2_admission(&mut c);
assert!(
matches!(d.admit.get(&2), Some(AdmitState::Registered { .. })),
"unblocked uni class admitted despite blocked bidi flood"
);
assert!(d.pending_admit_uni.is_empty());
assert_eq!(d.parked_bidi.len(), 1, "bidi parked exactly once");
}
#[test]
fn mf1_fresh_never_parked_class_reports_runnable() {
let (mut d, _h) = driver();
assert!(!d.has_runnable_remainder());
d.pending_admit.insert(0, PeerStream::new(0));
d.pending_admit_bidi.push_back(0);
assert!(d.parked_bidi.is_empty());
assert!(!d.accept_bidi_resume.load(Ordering::Relaxed));
assert!(
d.has_runnable_remainder(),
"fresh never-parked class with pending ids must be runnable (MF-A)"
);
let (mut d2, _h2) = driver();
d2.pending_admit.insert(2, PeerStream::new(2));
d2.pending_admit_uni.push_back(2);
assert!(d2.has_runnable_remainder());
}
#[test]
fn mf1_cross_class_fair_servicing() {
let (mut d, _h) = QuicheDriver::<Bytes>::new(false, 128, 128);
let mut c = MockConn::new();
for k in 0..40u64 {
let bidi = k * 4; let uni = k * 4 + 2; d.pending_admit.insert(bidi, PeerStream::new(bidi));
d.pending_admit_bidi.push_back(bidi);
d.pending_admit.insert(uni, PeerStream::new(uni));
d.pending_admit_uni.push_back(uni);
}
d.phase2_admission(&mut c);
let admitted_bidi = 40 - d.pending_admit_bidi.len();
let admitted_uni = 40 - d.pending_admit_uni.len();
assert_eq!(admitted_bidi + admitted_uni, ADMIT_BUDGET);
assert_eq!(admitted_bidi, ADMIT_BUDGET / 2, "bidi got a fair half");
assert_eq!(admitted_uni, ADMIT_BUDGET / 2, "uni got a fair half");
}
#[test]
fn tombstone_contract_a_removes_bidi_at_both_terminal() {
let (mut d, mut h) = driver();
let mut c = MockConn::new();
c.script_recv(0, [data(b"x", true)]);
d.pending_admit.insert(
0,
PeerStream {
id: 0,
pending_send_terminal: Some(SendEnd::Stopped { error_code: 9 }),
pending_recv_terminal: None,
},
);
d.pending_admit_bidi.push_back(0);
d.read_budget = READ_BUDGET;
d.run_read_pump(&mut c);
let mut ho = h.accept_bidi_rx.try_recv().expect("admitted");
assert_eq!(ho.recv.bytes.try_recv().unwrap(), Bytes::from_static(b"x"));
assert!(matches!(ho.recv.terminal.get(), Some(RecvEnd::Fin)));
assert!(matches!(
ho.send.status.get(),
Some(SendEnd::Stopped { error_code: 9 })
));
assert!(!d.admit.contains_key(&0));
assert!(!d.recv.contains_key(&0));
assert!(!d.pending_admit.contains_key(&0));
assert!(!d.readable_set.contains(&0));
d.read_budget = READ_BUDGET;
d.run_read_pump(&mut c);
assert!(h.accept_bidi_rx.try_recv().is_err());
}
#[test]
fn writable_path_captures_peer_stop_sending() {
let (mut d, mut h) = driver();
let mut c = MockConn::new();
c.writable_next.push_back(0);
c.capacity.insert(0, Err(quiche::Error::StreamStopped(66)));
d.stage_writable(&mut c);
d.read_budget = READ_BUDGET;
d.run_read_pump(&mut c);
let ho = h
.accept_bidi_rx
.try_recv()
.expect("admitted via writable path");
assert_eq!(ho.recv.id, 0);
assert!(matches!(
d.admit.get(&0),
Some(AdmitState::Registered {
send_done: true,
..
})
));
assert!(matches!(
ho.send.status.get(),
Some(SendEnd::Stopped { error_code: 66 })
));
let _ = &mut h;
}
#[test]
fn read_budget_boundary_requeues_with_membership() {
let (mut d, _h) = driver();
let (tx, mut rx) = mpsc::channel::<Bytes>(BYTE_CHANNEL_DEPTH);
d.recv.insert(
0,
StreamRecvState {
bytes: tx,
terminal: TerminalCell::new(),
resume: Arc::new(AtomicBool::new(false)),
blocked: Arc::new(AtomicBool::new(false)),
},
);
d.admit.insert(
0,
AdmitState::Registered {
send_done: false,
recv_done: false,
},
);
d.pending_readable.push_back(0);
d.readable_set.insert(0);
let mut c = MockConn::new();
c.script_recv(
0,
[
data(b"a", false),
data(b"b", false),
data(b"c", false),
data(b"d", false),
data(b"e", false),
],
);
d.read_budget = 3;
d.run_read_pump(&mut c);
assert_eq!(c.recv_calls.len(), 3);
assert!(d.needs_iteration);
assert!(d.readable_set.contains(&0));
assert_eq!(d.pending_readable.len(), 1);
assert_eq!(*d.pending_readable.front().unwrap(), 0);
for expected in [b"a", b"b", b"c"] {
assert_eq!(rx.try_recv().unwrap(), Bytes::copy_from_slice(expected));
}
assert!(rx.try_recv().is_err());
}
#[test]
fn needs_iteration_survives_full_packet_iteration() {
let (mut d, mut h) = driver();
let (tx, _rx) = mpsc::channel(BYTE_CHANNEL_DEPTH);
d.recv.insert(
0,
StreamRecvState {
bytes: tx,
terminal: TerminalCell::new(),
resume: Arc::new(AtomicBool::new(false)),
blocked: Arc::new(AtomicBool::new(false)),
},
);
d.admit.insert(
0,
AdmitState::Registered {
send_done: true,
recv_done: false,
},
);
d.pending_readable.push_back(0);
d.readable_set.insert(0);
let mut c = MockConn::new();
c.script_recv(0, (0..40).map(|i| data(&[b'a' + (i % 26) as u8], false)));
d.do_process_reads(&mut c);
assert!(d.needs_iteration, "pump should defer under READ_BUDGET");
d.do_process_writes(&mut c).expect("writes ok");
assert!(
d.needs_iteration,
"deferral must survive process_writes in the same iteration"
);
let _ = &mut h;
}
use h3::quic::WriteBuf;
fn wbuf(payload: &'static [u8]) -> WriteBuf<Bytes> {
WriteBuf::from(h3::proto::frame::Frame::Data(Bytes::from_static(payload)))
}
fn wbuf_len(payload: &'static [u8]) -> usize {
wbuf(payload).remaining()
}
fn sent_len(c: &MockConn, id: u64) -> usize {
c.sent
.iter()
.filter(|(sid, _, _)| *sid == id)
.map(|(_, b, _)| b.len())
.sum()
}
fn sent_fin(c: &MockConn, id: u64) -> bool {
c.sent.iter().any(|(sid, _, fin)| *sid == id && *fin)
}
struct SendProbe {
cell: crate::buffer::WriteCompletion<SendEnd>,
generation: u64,
}
#[derive(Debug)]
enum ProbeState {
Pending,
Ok,
Err(SendEnd),
Cancelled,
}
impl SendProbe {
fn state(&self) -> ProbeState {
match self.cell.try_take(self.generation) {
None => ProbeState::Pending,
Some(WriteOutcome::Done(Ok(()))) => ProbeState::Ok,
Some(WriteOutcome::Done(Err(e))) => ProbeState::Err(e),
Some(WriteOutcome::Cancelled) => ProbeState::Cancelled,
}
}
}
fn push_send_on(
d: &mut QuicheDriver<Bytes>,
cell: &crate::buffer::WriteCompletion<SendEnd>,
id: u64,
payload: &'static [u8],
) -> SendProbe {
let generation = cell.begin();
let buf = wbuf(payload);
let permit = d.shared.send_accounting.try_reserve(buf.remaining());
d.inbox.push_back(DriverCommand::Send {
id,
buf,
done: cell.completer(generation),
permit,
});
SendProbe {
cell: cell.clone(),
generation,
}
}
fn wire_len(payload: &'static [u8]) -> usize {
wbuf(payload).remaining()
}
fn push_send(d: &mut QuicheDriver<Bytes>, id: u64, payload: &'static [u8]) -> SendProbe {
let cell = crate::buffer::WriteCompletion::new();
push_send_on(d, &cell, id, payload)
}
fn push_finish(d: &mut QuicheDriver<Bytes>, id: u64) -> oneshot::Receiver<Result<(), SendEnd>> {
let (tx, rx) = oneshot::channel();
d.inbox.push_back(DriverCommand::Finish { id, done: tx });
rx
}
#[test]
fn partial_write_then_capacity_rearms_one_ok() {
let (mut d, _h) = driver();
let mut c = MockConn::new();
let total = wbuf_len(b"hello world");
c.send_capacity.insert(0, 3);
c.capacity.insert(0, Ok(3));
let done = push_send(&mut d, 0, b"hello world");
d.apply_inbox(&mut c);
d.stage_send(&mut c);
assert!(matches!(done.state(), ProbeState::Pending));
let rearms: Vec<usize> = c
.rearms
.iter()
.filter(|(id, _)| *id == 0)
.map(|(_, len)| *len)
.collect();
assert!(!rearms.is_empty(), "blocked write must low-water re-arm");
assert_eq!(*rearms.last().unwrap(), REARM_THRESHOLD);
let after_first = sent_len(&c, 0);
assert!(after_first > 0 && after_first < total);
c.send_capacity.remove(&0);
c.writable_next.push_back(0);
loop {
d.stage_writable(&mut c);
d.stage_send(&mut c);
if sent_len(&c, 0) == total {
break;
}
c.writable_next.push_back(0);
}
assert_eq!(sent_len(&c, 0), total, "all bytes eventually accepted");
assert!(
matches!(done.state(), ProbeState::Ok),
"exactly one Ok at full acceptance"
);
}
#[test]
fn finish_accepted_at_zero_capacity_completes_once() {
let (mut d, _h) = driver();
let mut c = MockConn::new();
c.send_capacity.insert(0, 0); let mut done = push_finish(&mut d, 0);
d.apply_inbox(&mut c);
d.stage_send(&mut c);
assert!(sent_fin(&c, 0), "zero-capacity FIN accepted (Q5)");
assert!(matches!(done.try_recv(), Ok(Ok(()))));
assert!(d.runnable_send.is_empty());
d.stage_send(&mut c);
assert!(matches!(
done.try_recv(),
Err(oneshot::error::TryRecvError::Closed)
));
}
#[test]
fn coalesced_write_then_finish_emits_single_fin_frame() {
let (mut d, _h) = driver();
let mut c = MockConn::new();
let body_wire = wbuf_len(b"hi");
let done = push_send(&mut d, 0, b"hi");
let mut fin = push_finish(&mut d, 0);
d.apply_inbox(&mut c);
d.stage_send(&mut c);
assert!(
!c.sent
.iter()
.any(|(id, b, f)| *id == 0 && b.is_empty() && *f),
"no standalone empty-FIN frame is emitted (issue #10)"
);
let fin_frames: Vec<&(u64, Vec<u8>, bool)> =
c.sent.iter().filter(|(id, _, f)| *id == 0 && *f).collect();
assert_eq!(
fin_frames.len(),
1,
"exactly one FIN, coalesced onto a DATA frame — not the pre-fix pair"
);
assert!(
!fin_frames[0].1.is_empty(),
"the FIN rides a non-empty final DATA chunk"
);
assert_eq!(
sent_len(&c, 0),
body_wire,
"the whole body (frame header + payload) is sent"
);
assert!(
matches!(done.state(), ProbeState::Ok),
"the body write completes Ok exactly once"
);
assert!(
matches!(fin.try_recv(), Ok(Ok(()))),
"the finish completes Ok exactly once"
);
assert!(
!d.runnable_send_set.contains(&0),
"runnable membership released after the coalesced FIN"
);
let (mut d2, _h2) = driver();
let mut c2 = MockConn::new();
let mut lone = push_finish(&mut d2, 4);
d2.apply_inbox(&mut c2);
d2.stage_send(&mut c2);
let sent2: Vec<&(u64, Vec<u8>, bool)> =
c2.sent.iter().filter(|(id, _, _)| *id == 4).collect();
assert_eq!(sent2.len(), 1, "lone finish emits exactly one empty-FIN");
assert!(sent2[0].1.is_empty(), "empty-body FIN carries no data");
assert!(sent2[0].2, "empty-body FIN sets the FIN bit");
assert!(
matches!(lone.try_recv(), Ok(Ok(()))),
"lone finish completes Ok exactly once"
);
}
#[test]
fn reset_preempts_queued_write_keeps_earlier_ok() {
let (mut d, _h) = driver();
let mut c = MockConn::new();
let done1 = push_send(&mut d, 0, b"a");
d.apply_inbox(&mut c);
d.stage_send(&mut c);
assert!(
matches!(done1.state(), ProbeState::Ok),
"Write1 accepted before reset"
);
let done2 = push_send(&mut d, 0, b"bcde");
d.inbox.push_back(DriverCommand::Reset { id: 0, code: 42 });
d.apply_inbox(&mut c);
match done2.state() {
ProbeState::Err(SendEnd::Reset { error_code: 42 }) => {}
other => panic!("Write2 must be cancelled once with local reset, got {other:?}"),
}
assert!(matches!(
d.send.get(&0).unwrap().status.get(),
Some(SendEnd::Reset { error_code: 42 })
));
d.stage_send(&mut c);
assert!(c.shutdowns.contains(&crate::conn::mock::ShutdownCall {
id: 0,
is_write: true,
code: 42,
}));
}
#[test]
fn sf6_accounting_reserves_on_admission_releases_on_completion() {
let (mut d, _h) = driver();
let mut c = MockConn::new();
assert_eq!(d.shared.send_accounting.resident(), 0);
let done = push_send(&mut d, 0, b"hello");
assert_eq!(
d.shared.send_accounting.resident(),
wire_len(b"hello"),
"reserved when the command is admitted"
);
d.apply_inbox(&mut c);
d.stage_send(&mut c);
assert!(
matches!(done.state(), ProbeState::Ok),
"write completes once"
);
assert_eq!(
d.shared.send_accounting.resident(),
0,
"released exactly once at the completion chokepoint"
);
}
#[test]
fn sf6_accounting_released_when_reset_drains_queued_write() {
let (mut d, _h) = driver();
let mut c = MockConn::new();
let done = push_send(&mut d, 0, b"abcd");
d.apply_inbox(&mut c); assert_eq!(d.shared.send_accounting.resident(), wire_len(b"abcd"));
d.inbox.push_back(DriverCommand::Reset { id: 0, code: 7 });
d.apply_inbox(&mut c); assert!(matches!(
done.state(),
ProbeState::Err(SendEnd::Reset { error_code: 7 })
));
assert_eq!(
d.shared.send_accounting.resident(),
0,
"reset drain releases the reservation"
);
}
#[test]
fn sf3_reused_cell_across_generations_no_cross_bleed() {
let (mut d, _h) = driver();
let mut c = MockConn::new();
let cell = crate::buffer::WriteCompletion::new();
let g = push_send_on(&mut d, &cell, 0, b"hello");
d.apply_inbox(&mut c);
d.stage_send(&mut c);
assert!(matches!(g.state(), ProbeState::Ok), "g completes once");
assert_eq!(d.shared.send_accounting.resident(), 0, "g released");
assert!(matches!(g.state(), ProbeState::Pending), "g consumed once");
let g1 = push_send_on(&mut d, &cell, 0, b"world");
assert_eq!(cell.generation(), 2, "same cell reused across both writes");
d.apply_inbox(&mut c); assert_eq!(d.shared.send_accounting.resident(), wire_len(b"world"));
d.inbox.push_back(DriverCommand::Reset { id: 0, code: 7 });
d.apply_inbox(&mut c);
assert!(
matches!(
g1.state(),
ProbeState::Err(SendEnd::Reset { error_code: 7 })
),
"g+1 resolves once as the reset"
);
assert!(matches!(g.state(), ProbeState::Pending), "no bleed into g");
assert_eq!(d.shared.send_accounting.resident(), 0, "g+1 released");
}
#[test]
fn sf6_worker_cap_bounds_admitted_bytes() {
let cap = wire_len(b"abc");
let (mut d, _h) = QuicheDriver::<Bytes>::with_buffers(
false,
4,
4,
DriverBufferConfig {
recv_channel_depth: BYTE_CHANNEL_DEPTH,
packet_buffer_size: PKT_BUF_LEN,
max_buffered_send_bytes: Some(cap),
},
);
let mut c = MockConn::new();
assert_eq!(d.shared.send_accounting.cap(), Some(cap));
let done = push_send(&mut d, 0, b"abc"); assert_eq!(d.shared.send_accounting.resident(), cap);
assert!(d.shared.send_accounting.try_reserve(1).is_none());
d.apply_inbox(&mut c);
d.stage_send(&mut c);
assert!(matches!(done.state(), ProbeState::Ok));
assert_eq!(d.shared.send_accounting.resident(), 0);
assert!(d.shared.send_accounting.try_reserve(cap).is_some());
}
#[test]
fn reset_emitted_at_zero_capacity() {
let (mut d, _h) = driver();
let mut c = MockConn::new();
c.send_capacity.insert(0, 0);
d.inbox.push_back(DriverCommand::Reset { id: 0, code: 7 });
d.apply_inbox(&mut c);
d.stage_send(&mut c);
assert!(c.shutdowns.contains(&crate::conn::mock::ShutdownCall {
id: 0,
is_write: true,
code: 7,
}));
}
#[test]
fn accepted_fin_marks_send_done_and_enables_contract_a() {
let (mut d, _h) = driver();
d.admit.insert(
0,
AdmitState::Registered {
send_done: false,
recv_done: true,
},
);
d.send.insert(0, StreamSendState::new());
let mut c = MockConn::new();
let mut fin = push_finish(&mut d, 0);
d.apply_inbox(&mut c);
d.stage_send(&mut c);
assert!(matches!(fin.try_recv(), Ok(Ok(()))));
assert!(
!d.admit.contains_key(&0),
"both directions terminal → admit dropped"
);
assert!(!d.recv.contains_key(&0));
}
#[test]
fn duplicate_reset_is_idempotent() {
let (mut d, _h) = driver();
let mut c = MockConn::new();
d.inbox.push_back(DriverCommand::Reset { id: 0, code: 7 });
d.apply_inbox(&mut c);
d.stage_send(&mut c); d.inbox.push_back(DriverCommand::Reset { id: 0, code: 9 });
d.apply_inbox(&mut c);
d.stage_send(&mut c);
let resets: Vec<_> = c.shutdowns.iter().filter(|s| s.is_write).collect();
assert_eq!(resets.len(), 1, "exactly one RESET_STREAM");
assert_eq!(resets[0].code, 7, "first-effective reset code wins");
}
#[test]
fn deferred_send_after_contract_a_completes_with_sticky_terminal() {
let (mut d, _h) = driver();
d.admit.insert(
0,
AdmitState::Registered {
send_done: false,
recv_done: true,
},
);
let mut c = MockConn::new();
let w1 = push_send(&mut d, 0, b"aa");
c.capacity.insert(0, Err(quiche::Error::StreamStopped(55)));
d.apply_inbox(&mut c);
d.stage_send(&mut c);
assert!(matches!(
w1.state(),
ProbeState::Err(SendEnd::Stopped { error_code: 55 })
));
assert!(!d.admit.contains_key(&0), "contract A reclaimed admit");
assert!(d.send.contains_key(&0), "send retained for deferred ops");
let late = push_send(&mut d, 0, b"bb");
d.apply_inbox(&mut c);
assert!(matches!(
late.state(),
ProbeState::Err(SendEnd::Stopped { error_code: 55 })
));
}
#[test]
fn stop_sending_on_send_drains_all_ops_once() {
let (mut d, _h) = driver();
let mut c = MockConn::new();
let w1 = push_send(&mut d, 0, b"aa");
let w2 = push_send(&mut d, 0, b"bb");
let mut fin = push_finish(&mut d, 0);
d.apply_inbox(&mut c);
c.capacity.insert(0, Err(quiche::Error::StreamStopped(9)));
d.stage_send(&mut c);
for (label, st) in [("w1", w1.state()), ("w2", w2.state())] {
match st {
ProbeState::Err(SendEnd::Stopped { error_code: 9 }) => {}
other => panic!("{label} must complete once with Stopped, got {other:?}"),
}
}
match fin.try_recv() {
Ok(Err(SendEnd::Stopped { error_code: 9 })) => {}
other => panic!("fin must complete once with Stopped, got {other:?}"),
}
assert!(matches!(
d.send.get(&0).unwrap().status.get(),
Some(SendEnd::Stopped { error_code: 9 })
));
assert!(d.send.get(&0).unwrap().send_ops.is_empty());
assert!(
!d.runnable_send_set.contains(&0),
"runnable membership released"
);
}
#[test]
fn stop_sending_via_writable_probe_drains_ops() {
let (mut d, _h) = driver();
let mut c = MockConn::new();
let w1 = push_send(&mut d, 0, b"aa");
let w2 = push_send(&mut d, 0, b"bb");
d.apply_inbox(&mut c);
c.writable_next.push_back(0);
c.capacity.insert(0, Err(quiche::Error::StreamStopped(13)));
d.stage_writable(&mut c);
for (label, st) in [("w1", w1.state()), ("w2", w2.state())] {
match st {
ProbeState::Err(SendEnd::Stopped { error_code: 13 }) => {}
other => panic!("{label} must complete once with Stopped, got {other:?}"),
}
}
assert!(!d.runnable_send_set.contains(&0));
}
#[test]
fn round_robin_bulk_yields_to_other_stream() {
let (mut d, _h) = driver();
let mut c = MockConn::new();
static BULK: [u8; 1024 * 1024] = [b'x'; 1024 * 1024];
c.capacity.insert(0, Ok(MAX_WRITE_CHUNK));
let bulk_cell = crate::buffer::WriteCompletion::<SendEnd>::new();
let bulk_gen = bulk_cell.begin();
d.inbox.push_back(DriverCommand::Send {
id: 0,
buf: WriteBuf::from(h3::proto::frame::Frame::Data(Bytes::from_static(&BULK))),
done: bulk_cell.completer(bulk_gen),
permit: None,
});
let bulk = SendProbe {
cell: bulk_cell,
generation: bulk_gen,
};
let small_done = push_send(&mut d, 4, b"z");
d.apply_inbox(&mut c);
d.stage_send(&mut c);
assert!(
matches!(small_done.state(), ProbeState::Ok),
"small stream serviced"
);
assert!(matches!(bulk.state(), ProbeState::Pending));
assert!(d.runnable_send_set.contains(&0) || d.needs_iteration);
}
#[test]
fn send_after_terminal_completes_immediately() {
let (mut d, _h) = driver();
let mut c = MockConn::new();
d.inbox.push_back(DriverCommand::Reset { id: 0, code: 5 });
d.apply_inbox(&mut c);
d.stage_send(&mut c);
let late = push_send(&mut d, 0, b"late");
d.apply_inbox(&mut c);
match late.state() {
ProbeState::Err(SendEnd::Reset { error_code: 5 }) => {}
other => panic!("late Send must complete once with sticky terminal, got {other:?}"),
}
assert!(
d.send.get(&0).unwrap().send_ops.is_empty(),
"late op not enqueued"
);
assert!(c.sent.iter().all(|(_, b, _)| b != b"late"));
}
#[test]
fn admitted_bidi_retains_send_state_sharing_status() {
let (mut d, mut h) = driver();
let mut c = MockConn::new();
c.script_recv(0, [data(b"hi", false)]);
c.queue_readable([0]);
d.read_budget = READ_BUDGET;
d.run_read_pump(&mut c);
let ho = h.accept_bidi_rx.try_recv().expect("admitted");
assert!(d.send.contains_key(&0));
assert!(ho.send.status.get().is_none());
c.writable_next.push_back(0);
c.capacity.insert(0, Err(quiche::Error::StreamStopped(88)));
d.stage_writable(&mut c);
assert!(matches!(
ho.send.status.get(),
Some(SendEnd::Stopped { error_code: 88 })
));
}
fn conn_err(is_app: bool, code: u64, reason: &[u8]) -> quiche::ConnectionError {
quiche::ConnectionError {
is_app,
error_code: code,
reason: reason.to_vec(),
}
}
#[test]
fn last_handle_teardown_issues_h3_no_error_close() {
let (mut d, _h) = driver();
let mut c = MockConn::new();
d.last_handle_teardown = true;
d.do_process_writes(&mut c).expect("clean teardown");
assert_eq!(c.closed, Some((true, H3_NO_ERROR, b"".to_vec())));
assert!(d.explicit_close_attempted);
assert!(d.graceful_close_issued);
let lc = d.local_close.as_ref().expect("recorded last-handle close");
assert_eq!(lc.code, H3_NO_ERROR);
assert!(lc.reason.is_empty());
}
#[test]
fn explicit_close_crosses_barrier_after_saturated_batch() {
let (mut d, _h) = driver();
let mut c = MockConn::new();
for id in 0..=(WRITE_BUDGET as u64) {
let sid = id * 4; c.send_capacity.insert(sid, 1);
let _rx = push_send(&mut d, sid, b"hello world");
}
d.inbox.push_back(DriverCommand::Close {
code: 0x1234,
reason: Bytes::from_static(b"bye"),
});
d.do_process_writes(&mut c).expect("close applied");
assert!(d.needs_iteration, "write batch should be saturated");
assert_eq!(c.closed, Some((true, 0x1234, b"bye".to_vec())));
let lc = d.local_close.as_ref().expect("explicit close recorded");
assert_eq!(lc.code, 0x1234);
assert_eq!(&lc.reason[..], b"bye");
assert!(d.graceful_close_issued);
d.last_handle_teardown = true;
d.reads_ran_this_iter = false;
d.do_process_writes(&mut c).expect("no second close");
assert_eq!(
c.closed,
Some((true, 0x1234, b"bye".to_vec())),
"synthetic H3_NO_ERROR must be suppressed"
);
}
#[test]
fn first_close_wins_second_ignored() {
let (mut d, _h) = driver();
let mut c = MockConn::new();
d.inbox.push_back(DriverCommand::Close {
code: 0xaaa,
reason: Bytes::from_static(b"first"),
});
d.inbox.push_back(DriverCommand::Close {
code: 0xbbb,
reason: Bytes::from_static(b"second"),
});
d.apply_inbox(&mut c);
let pc = d.pending_close.as_ref().expect("first staged");
assert_eq!(pc.code, 0xaaa);
assert_eq!(&pc.reason[..], b"first");
d.apply_close_barrier(&mut c).expect("applied");
assert_eq!(c.closed, Some((true, 0xaaa, b"first".to_vec())));
}
#[test]
fn peer_app_close_outranks_last_handle_teardown() {
let (mut d, _h) = driver();
let mut c = MockConn::new();
c.peer_error = Some(conn_err(true, 0x99, b"peer-bye"));
d.last_handle_teardown = true;
d.do_process_writes(&mut c).expect("no synthetic close");
assert_eq!(
c.closed, None,
"synthetic close suppressed by peer terminal"
);
d.do_on_conn_close(&mut c);
match d.shared.conn_terminal.get().as_deref() {
Some(ConnTerminal::AppClose {
origin: CloseOrigin::Peer,
error_code: 0x99,
reason,
}) => assert_eq!(&reason[..], b"peer-bye"),
other => panic!("expected AppClose{{Peer}}, got {other:?}"),
}
}
#[test]
fn timeout_classifies_as_timeout() {
let (mut d, _h) = driver();
let mut c = MockConn::new();
c.timed_out = true;
d.do_on_conn_close(&mut c);
assert!(matches!(
d.shared.conn_terminal.get().as_deref(),
Some(ConnTerminal::Timeout)
));
}
#[test]
fn on_conn_close_publishes_to_all_out_of_band_cells() {
let (mut d, mut h) = driver();
let mut c = MockConn::new();
c.script_recv(0, [data(b"hi", false)]);
c.queue_readable([0]);
d.read_budget = READ_BUDGET;
d.run_read_pump(&mut c);
let ho = h.accept_bidi_rx.try_recv().expect("admitted");
c.local_error = Some(conn_err(true, 0x100, b""));
d.do_on_conn_close(&mut c);
assert!(matches!(ho.recv.terminal.get(), Some(RecvEnd::Conn(_))));
assert!(matches!(ho.send.status.get(), Some(SendEnd::Conn(_))));
assert!(h.accept_terminal_bidi.get().is_some(), "bidi accept cell");
assert!(h.accept_terminal_uni.get().is_some(), "uni accept cell");
assert!(d.shared.conn_terminal.get().is_some(), "conn cell");
}
#[test]
fn pending_send_op_drains_with_conn_terminal() {
let (mut d, _h) = driver();
let mut c = MockConn::new();
let done = push_send(&mut d, 0, b"payload");
d.apply_inbox(&mut c);
assert!(!d.send.get(&0).unwrap().send_ops.is_empty());
c.peer_error = Some(conn_err(true, 0x7, b""));
d.do_on_conn_close(&mut c);
match done.state() {
ProbeState::Err(SendEnd::Conn(_)) => {}
other => panic!("expected SendEnd::Conn, got {other:?}"),
}
}
#[test]
fn unapplied_send_at_close_completes_generation_once() {
let (mut d, _h) = driver();
let mut c = MockConn::new();
let done = push_send(&mut d, 0, b"payload");
assert!(
!d.send.contains_key(&0),
"precondition: Send not yet applied into send_ops"
);
c.peer_error = Some(conn_err(true, 0x9, b""));
d.do_on_conn_close(&mut c);
match done.state() {
ProbeState::Err(SendEnd::Conn(_)) => {}
other => panic!("expected SendEnd::Conn for unapplied Send, got {other:?}"),
}
assert!(matches!(done.state(), ProbeState::Pending));
}
#[test]
fn send_conngone_in_closing_window_defers_to_on_conn_close() {
let (mut d, _h) = driver();
let mut c = MockConn::new();
let done = push_send(&mut d, 0, b"x");
c.capacity.insert(0, Err(quiche::Error::InvalidState));
d.apply_inbox(&mut c);
d.stage_send(&mut c);
assert!(matches!(done.state(), ProbeState::Pending));
assert!(!d.send.get(&0).unwrap().send_ops.is_empty());
assert!(d.send.get(&0).unwrap().status.get().is_none());
c.peer_error = Some(conn_err(true, 0x101, b""));
d.do_on_conn_close(&mut c);
match done.state() {
ProbeState::Err(SendEnd::Conn(_)) => {}
other => panic!("expected SendEnd::Conn after close, got {other:?}"),
}
}
#[test]
fn recv_conngone_in_closing_window_defers_to_on_conn_close() {
let (mut d, _h) = driver();
let (tx, _rx) = mpsc::channel(BYTE_CHANNEL_DEPTH);
let terminal_cell = TerminalCell::new();
d.recv.insert(
0,
StreamRecvState {
bytes: tx,
terminal: terminal_cell.clone(),
resume: Arc::new(AtomicBool::new(false)),
blocked: Arc::new(AtomicBool::new(false)),
},
);
d.pending_readable.push_back(0);
d.readable_set.insert(0);
let mut c = MockConn::new();
c.readable_ids.insert(0);
c.script_recv(0, [RecvStep::Err(quiche::Error::InvalidState)]);
d.read_budget = READ_BUDGET;
d.run_read_pump(&mut c);
assert!(d.recv.contains_key(&0));
assert!(terminal_cell.get().is_none());
c.peer_error = Some(conn_err(true, 0x102, b""));
d.do_on_conn_close(&mut c);
assert!(matches!(terminal_cell.get(), Some(RecvEnd::Conn(_))));
}
#[test]
fn deferred_open_bidi_resolves_err_after_close() {
let (mut d, h) = driver();
let mut c = MockConn::new();
let (reply_tx, mut reply_rx) = oneshot::channel();
h.cmd_tx
.send(DriverCommand::OpenBidi { reply: reply_tx })
.expect("enqueue open");
c.peer_error = Some(conn_err(true, 0x2, b""));
d.do_on_conn_close(&mut c);
match reply_rx.try_recv() {
Ok(Err(t)) => assert!(matches!(
t.as_ref(),
ConnTerminal::AppClose {
origin: CloseOrigin::Peer,
..
}
)),
Ok(Ok(_)) => panic!("expected Err(terminal), got Ok(handoff)"),
Err(e) => panic!("expected Err(terminal), got {e:?}"),
}
}
#[test]
fn done_close_result_defers_to_preexisting_terminal() {
let (mut d, _h) = driver();
let mut c = MockConn::new();
c.close_result = Some(quiche::Error::Done);
c.peer_error = Some(conn_err(true, 0x55, b"peer"));
d.inbox.push_back(DriverCommand::Close {
code: 0x1,
reason: Bytes::from_static(b"local"),
});
d.apply_inbox(&mut c);
d.apply_close_barrier(&mut c)
.expect("done defers, not a bug");
assert!(d.explicit_close_attempted);
assert!(d.graceful_close_issued);
assert!(d.local_close.is_none(), "Done must not record acceptance");
d.do_on_conn_close(&mut c);
match d.shared.conn_terminal.get().as_deref() {
Some(ConnTerminal::AppClose {
origin: CloseOrigin::Peer,
error_code: 0x55,
..
}) => {}
other => panic!("expected pre-existing peer terminal, got {other:?}"),
}
}
#[test]
fn recorded_explicit_close_outranks_local_error() {
let (mut d, _h) = driver();
let mut c = MockConn::new();
d.inbox.push_back(DriverCommand::Close {
code: 0x321,
reason: Bytes::from_static(b"quit"),
});
d.apply_inbox(&mut c);
d.apply_close_barrier(&mut c).expect("applied");
c.local_error = Some(conn_err(true, 0x999, b"other"));
d.do_on_conn_close(&mut c);
match d.shared.conn_terminal.get().as_deref() {
Some(ConnTerminal::AppClose {
origin: CloseOrigin::Local,
error_code: 0x321,
reason,
}) => assert_eq!(&reason[..], b"quit"),
other => panic!("expected recorded local close, got {other:?}"),
}
}
#[test]
fn unexpected_close_error_is_internal_bug() {
let (mut d, _h) = driver();
let mut c = MockConn::new();
c.close_result = Some(quiche::Error::TlsFail);
d.last_handle_teardown = true;
let err = d.do_process_writes(&mut c);
assert!(
err.is_err(),
"unexpected close error must fail the callback"
);
assert!(d.close_bug.is_some());
d.do_on_conn_close(&mut c);
assert!(matches!(
d.shared.conn_terminal.get().as_deref(),
Some(ConnTerminal::Internal(_))
));
}
#[test]
fn connection_dropped_cleans_parked_peer_streams() {
let (mut d, _h) = driver();
let mut c = MockConn::new();
d.admit.insert(0, AdmitState::Parked(PeerStream::new(0)));
d.parked_bidi.push_back(0);
d.pending_admit.insert(2, PeerStream::new(2));
d.pending_admit_uni.push_back(2);
d.inbox.push_back(DriverCommand::ConnectionDropped);
d.apply_inbox(&mut c);
assert!(!d.admit.contains_key(&0), "parked bidi dropped");
assert!(d.parked_bidi.is_empty());
assert!(d.pending_admit.is_empty());
assert!(d.pending_admit_uni.is_empty());
assert!(c.shutdowns.iter().any(|s| s.id == 0 && s.is_write));
assert!(c.shutdowns.iter().any(|s| s.id == 0 && !s.is_write));
assert!(c.shutdowns.iter().any(|s| s.id == 2 && !s.is_write));
assert!(!c.shutdowns.iter().any(|s| s.id == 2 && s.is_write));
}
fn push_open_bidi(
d: &mut QuicheDriver<Bytes>,
) -> oneshot::Receiver<Result<BidiHandoff<Bytes>, Arc<ConnTerminal>>> {
let (tx, rx) = oneshot::channel();
d.inbox.push_back(DriverCommand::OpenBidi { reply: tx });
rx
}
fn push_open_uni(
d: &mut QuicheDriver<Bytes>,
) -> oneshot::Receiver<Result<SendHandoff<Bytes>, Arc<ConnTerminal>>> {
let (tx, rx) = oneshot::channel();
d.inbox.push_back(DriverCommand::OpenUni { reply: tx });
rx
}
#[test]
fn stage_open_bidi_allocates_one_id_and_increments_counter() {
let (mut d, _h) = driver();
let mut c = MockConn::new();
c.streams_left_bidi = 4;
let mut reply = push_open_bidi(&mut d);
d.apply_inbox(&mut c);
assert_eq!(
d.next_bidi_id, 0,
"counter unchanged before materialization"
);
d.stage_open(&mut c);
assert_eq!(c.priorities, vec![(0, 127, true)]);
assert_eq!(d.next_bidi_id, 4, "counter advances by 4 after success");
assert!(d.recv.contains_key(&0));
assert!(d.send.contains_key(&0));
match reply.try_recv() {
Ok(Ok(handoff)) => {
assert_eq!(handoff.send.id, 0);
assert_eq!(handoff.recv.id, 0);
}
_ => panic!("expected BidiHandoff"),
}
assert!(d.open_bidi.is_empty());
}
#[test]
fn stage_open_uni_allocates_send_only() {
let (mut d, _h) = driver();
let mut c = MockConn::new();
c.streams_left_uni = 4;
let mut reply = push_open_uni(&mut d);
d.apply_inbox(&mut c);
d.stage_open(&mut c);
assert_eq!(c.priorities, vec![(2, 127, true)]);
assert_eq!(d.next_uni_id, 6);
assert!(d.send.contains_key(&2));
assert!(!d.recv.contains_key(&2), "uni open has no recv half");
match reply.try_recv() {
Ok(Ok(handoff)) => assert_eq!(handoff.id, 2),
_ => panic!("expected SendHandoff"),
}
}
#[test]
fn stage_open_is_closed_skips_no_id_burned() {
let (mut d, _h) = driver();
let mut c = MockConn::new();
c.streams_left_bidi = 4;
let reply = push_open_bidi(&mut d);
drop(reply); d.apply_inbox(&mut c);
d.stage_open(&mut c);
assert!(
c.priorities.is_empty(),
"no id materialized for a dead reply"
);
assert_eq!(d.next_bidi_id, 0, "counter not advanced");
assert!(d.open_bidi.is_empty());
assert!(!d.send.contains_key(&0));
assert!(!d.recv.contains_key(&0));
}
#[test]
fn stage_open_zero_credit_defers_request() {
let (mut d, _h) = driver();
let mut c = MockConn::new();
c.streams_left_bidi = 0;
let _reply = push_open_bidi(&mut d);
d.apply_inbox(&mut c);
d.needs_iteration = false;
d.stage_open(&mut c);
assert!(c.priorities.is_empty());
assert_eq!(d.next_bidi_id, 0);
assert_eq!(d.open_bidi.len(), 1, "request stays queued");
assert!(
!d.needs_iteration,
"blocking on stream credit must not hot-spin (§5.2 progress bound)"
);
}
#[test]
fn cleanup_undeliverable_open_is_direction_aware() {
let (mut d, _h) = driver();
let mut c = MockConn::new();
let cmd_tx = d.cmd_tx_weak.upgrade().unwrap();
let (recv_state, _rh, _rd) = d.build_recv(0, cmd_tx.clone(), None);
let (_sh, send_state, _sd) =
build_send(0, cmd_tx, Arc::clone(&d.shared.send_accounting), None);
d.recv.insert(0, recv_state.unwrap());
d.send.insert(0, send_state.unwrap());
d.cleanup_undeliverable_open(&mut c, 0, true);
assert!(!d.recv.contains_key(&0));
assert!(!d.send.contains_key(&0));
let bidi_shuts: Vec<&crate::conn::mock::ShutdownCall> =
c.shutdowns.iter().filter(|s| s.id == 0).collect();
assert!(bidi_shuts
.iter()
.any(|s| s.is_write && s.code == H3_REQUEST_CANCELLED));
assert!(bidi_shuts
.iter()
.any(|s| !s.is_write && s.code == H3_REQUEST_CANCELLED));
let (_sh, send_state, _sd) = build_send(
2,
d.cmd_tx_weak.upgrade().unwrap(),
Arc::clone(&d.shared.send_accounting),
None,
);
d.send.insert(2, send_state.unwrap());
d.cleanup_undeliverable_open(&mut c, 2, false);
assert!(!d.send.contains_key(&2));
let uni_shuts: Vec<&crate::conn::mock::ShutdownCall> =
c.shutdowns.iter().filter(|s| s.id == 2).collect();
assert_eq!(uni_shuts.len(), 1);
assert!(uni_shuts[0].is_write && uni_shuts[0].code == H3_REQUEST_CANCELLED);
}
#[test]
fn stage_open_bidi_uses_server_parity() {
let (mut d, _h) = QuicheDriver::<Bytes>::new(true, 4, 4);
let mut c = MockConn::new();
c.streams_left_bidi = 4;
let _reply = push_open_bidi(&mut d);
d.apply_inbox(&mut c);
d.stage_open(&mut c);
assert_eq!(c.priorities, vec![(1, 127, true)]);
assert_eq!(d.next_bidi_id, 5);
}
#[test]
fn open_backlog_past_budget_sets_needs_iteration() {
let (mut d, _h) = driver();
let mut c = MockConn::new();
c.streams_left_bidi = u64::MAX; let _replies: Vec<_> = (0..(OPEN_BUDGET + 4))
.map(|_| push_open_bidi(&mut d))
.collect();
d.apply_inbox(&mut c);
d.needs_iteration = false;
d.stage_open(&mut c);
assert_eq!(c.priorities.len(), OPEN_BUDGET);
assert_eq!(d.open_bidi.len(), 4);
assert!(d.needs_iteration, "backlog must force another iteration");
}
#[test]
fn on_conn_close_drains_staged_open_queues() {
let (mut d, _h) = driver();
let mut c = MockConn::new();
let mut bidi_reply = push_open_bidi(&mut d);
let mut uni_reply = push_open_uni(&mut d);
d.apply_inbox(&mut c);
assert_eq!(d.open_bidi.len(), 1);
assert_eq!(d.open_uni.len(), 1);
c.peer_error = Some(conn_err(true, 0x3, b""));
d.do_on_conn_close(&mut c);
assert!(
matches!(bidi_reply.try_recv(), Ok(Err(_))),
"staged bidi open resolved"
);
assert!(
matches!(uni_reply.try_recv(), Ok(Err(_))),
"staged uni open resolved"
);
assert!(d.open_bidi.is_empty());
assert!(d.open_uni.is_empty());
}
#[test]
fn send_entry_reclaimed_when_send_finishes_after_recv() {
let (mut d, _h) = driver();
d.admit.insert(
0,
AdmitState::Registered {
send_done: false,
recv_done: false,
},
);
let (tx, _rx) = mpsc::channel(BYTE_CHANNEL_DEPTH);
d.recv.insert(
0,
StreamRecvState {
bytes: tx,
terminal: TerminalCell::new(),
resume: Arc::new(AtomicBool::new(false)),
blocked: Arc::new(AtomicBool::new(false)),
},
);
d.send.insert(0, StreamSendState::new());
let mut c = MockConn::new();
c.script_recv(0, [data(b"req", true)]);
d.pending_readable.push_back(0);
d.readable_set.insert(0);
d.read_budget = READ_BUDGET;
d.run_read_pump(&mut c);
assert!(!d.recv.contains_key(&0), "recv reclaimed on FIN");
assert!(d.send.contains_key(&0), "send retained while still open");
let mut fin = push_finish(&mut d, 0);
d.apply_inbox(&mut c);
d.stage_send(&mut c);
assert!(matches!(fin.try_recv(), Ok(Ok(()))));
assert!(
!d.send.contains_key(&0),
"send entry reclaimed on FIN-after-recv-done (no per-request leak)"
);
}
#[test]
fn packet_path_resume_forces_another_iteration() {
let (mut d, _h) = driver();
let mut c = MockConn::new();
d.reads_ran_this_iter = true;
d.needs_iteration = false;
d.inbox.push_back(DriverCommand::RecvResume { id: 0 });
d.do_process_writes(&mut c).unwrap();
assert!(
d.needs_iteration,
"a resume applied on a packet iteration must not strand"
);
}
}
#[cfg(test)]
mod loopback_tests {
use super::*;
use tokio::net::UdpSocket;
use tokio::time::{timeout, Duration};
use tokio_quiche::metrics::DefaultMetrics;
use tokio_quiche::quic::connect_with_config;
use tokio_quiche::settings::{CertificateKind, Hooks, QuicSettings, TlsCertificatePaths};
use tokio_quiche::socket::Socket;
use tokio_quiche::ConnectionParams;
use futures::StreamExt;
struct TestCerts {
cert_path: String,
key_path: String,
}
impl TestCerts {
fn generate() -> Self {
let ck = rcgen::generate_simple_self_signed(vec!["localhost".to_string()])
.expect("self-signed cert");
let dir = std::env::temp_dir();
let uniq = format!(
"quiche-h3-driver-{}-{}",
std::process::id(),
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_nanos()
);
let cert_path = dir.join(format!("{uniq}.crt"));
let key_path = dir.join(format!("{uniq}.key"));
std::fs::write(&cert_path, ck.cert.pem()).expect("write cert");
std::fs::write(&key_path, ck.signing_key.serialize_pem()).expect("write key");
Self {
cert_path: cert_path.to_string_lossy().into_owned(),
key_path: key_path.to_string_lossy().into_owned(),
}
}
}
impl Drop for TestCerts {
fn drop(&mut self) {
let _ = std::fs::remove_file(&self.cert_path);
let _ = std::fs::remove_file(&self.key_path);
}
}
fn client_params() -> ConnectionParams<'static> {
let mut settings = QuicSettings::default();
settings.verify_peer = false;
settings.max_idle_timeout = Some(Duration::from_secs(10));
ConnectionParams::new_client(settings, None, Hooks::default())
}
fn server_params(certs: &TestCerts) -> ConnectionParams<'_> {
let mut settings = QuicSettings::default();
settings.max_idle_timeout = Some(Duration::from_secs(10));
ConnectionParams::new_server(
settings,
TlsCertificatePaths {
cert: &certs.cert_path,
private_key: &certs.key_path,
kind: CertificateKind::X509,
},
Hooks::default(),
)
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
#[ignore = "loopback: binds UDP + runs a real handshake"]
async fn handshake_reaches_on_conn_established() {
let certs = TestCerts::generate();
let server_udp = UdpSocket::bind("127.0.0.1:0").await.unwrap();
let server_addr = server_udp.local_addr().unwrap();
let (server_driver, server_handles) = QuicheDriver::<Bytes>::new(true, 8, 8);
let mut listeners =
tokio_quiche::listen([server_udp], server_params(&certs), DefaultMetrics)
.expect("listen");
let server_task = tokio::spawn(async move {
let stream = &mut listeners[0];
if let Some(Ok(conn)) = stream.next().await {
let _qconn = conn.start(server_driver);
tokio::time::sleep(Duration::from_millis(500)).await;
}
});
let client_udp = UdpSocket::bind("127.0.0.1:0").await.unwrap();
client_udp.connect(server_addr).await.unwrap();
let client_socket = Socket::try_from(client_udp).expect("socket");
let (client_driver, client_handles) = QuicheDriver::<Bytes>::new(false, 8, 8);
let params = client_params();
let conn = connect_with_config(client_socket, Some("localhost"), ¶ms, client_driver)
.await
.expect("client handshake");
let DriverHandles {
cmd_tx: client_cmd_tx,
accept_bidi_rx: _c_bidi,
accept_uni_rx: _c_uni,
established_rx: client_established_rx,
shared: _c_shared,
..
} = client_handles;
let DriverHandles {
cmd_tx: server_cmd_tx,
accept_bidi_rx: _s_bidi,
accept_uni_rx: _s_uni,
established_rx: server_established_rx,
shared: _s_shared,
..
} = server_handles;
let client_est = timeout(Duration::from_secs(2), client_established_rx)
.await
.expect("client established within timeout")
.expect("client established_rx not cancelled");
assert!(client_est.is_ok(), "client establish should be Ok");
let server_est = timeout(Duration::from_secs(2), server_established_rx)
.await
.expect("server established within timeout")
.expect("server established_rx not cancelled");
assert!(server_est.is_ok(), "server establish should be Ok");
drop(client_cmd_tx);
drop(server_cmd_tx);
drop(conn);
let _ = server_task.await;
}
}