use crate::Error;
use crate::coding::{Reader, Writer};
const CONTROL_SEND_ORDER: u8 = u8::MAX;
pub struct Stream<S: web_transport_trait::Session, V> {
pub writer: Writer<S::SendStream, V>,
pub reader: Reader<S::RecvStream, V>,
}
impl<S: web_transport_trait::Session, V> Stream<S, V> {
pub async fn open(session: &S, version: V) -> Result<Self, Error>
where
V: Clone,
{
let (send, recv) = session.open_bi().await.map_err(Error::from_transport)?;
let mut writer = Writer::new(send, version.clone());
writer.set_priority(CONTROL_SEND_ORDER);
let reader = Reader::new(recv, version);
Ok(Stream { writer, reader })
}
pub async fn accept(session: &S, version: V) -> Result<Self, Error>
where
V: Clone,
{
let (send, recv) = session.accept_bi().await.map_err(Error::from_transport)?;
let mut writer = Writer::new(send, version.clone());
writer.set_priority(CONTROL_SEND_ORDER);
let reader = Reader::new(recv, version);
Ok(Stream { writer, reader })
}
pub fn with_version<V2: Clone>(self, version: V2) -> Stream<S, V2> {
Stream {
writer: self.writer.with_version(version.clone()),
reader: self.reader.with_version(version),
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::lite::{Version, test_transport::SinkSession};
#[tokio::test]
async fn open_prioritises_the_stream() {
let gate = kio::Producer::new(true);
let session = SinkSession::gated_bi(gate.consume());
let log = session.log.clone();
let _stream = Stream::open(&session, Version::Lite05).await.unwrap();
assert_eq!(
log.priorities(),
vec![CONTROL_SEND_ORDER],
"a control stream left at the transport default loses to every group",
);
}
#[tokio::test]
async fn accept_prioritises_the_stream() {
let gate = kio::Producer::new(true);
let session = SinkSession::accepted_bi(gate.consume());
let log = session.log.clone();
let _stream = Stream::accept(&session, Version::Lite05).await.unwrap();
assert_eq!(
log.priorities(),
vec![CONTROL_SEND_ORDER],
"an accepted control stream left at the transport default loses to every group",
);
}
}