use std::task::{Context, Poll, ready};
use crate::Error;
use crate::coding::{Reader, StreamCodes, Writer};
const CONTROL_SEND_ORDER: u8 = u8::MAX;
pub struct Stream<S: crate::transport::poll::Session, V> {
pub writer: Writer<S::SendStream, V>,
pub reader: Reader<S::RecvStream, V>,
}
impl<S: crate::transport::poll::Session, V: StreamCodes> Stream<S, V> {
pub fn poll_open(session: &mut S, version: V, cx: &mut Context<'_>) -> Poll<Result<Self, Error>>
where
V: Clone,
{
let (send, recv) = ready!(session.poll_open_bi(cx)).map_err(Error::from_transport)?;
Poll::Ready(Ok(Self::build(send, recv, version)))
}
pub async fn open(session: &mut S, version: V) -> Result<Self, Error>
where
V: Clone,
{
let (send, recv) = session.open_bi().await.map_err(Error::from_transport)?;
Ok(Self::build(send, recv, version))
}
pub fn poll_accept(session: &mut S, version: V, cx: &mut Context<'_>) -> Poll<Result<Self, Error>>
where
V: Clone,
{
let (send, recv) = ready!(session.poll_accept_bi(cx)).map_err(Error::from_transport)?;
Poll::Ready(Ok(Self::build(send, recv, version)))
}
pub async fn accept(session: &mut S, version: V) -> Result<Self, Error>
where
V: Clone,
{
let (send, recv) = session.accept_bi().await.map_err(Error::from_transport)?;
Ok(Self::build(send, recv, version))
}
fn build(send: S::SendStream, recv: S::RecvStream, version: V) -> Self
where
V: Clone,
{
let mut writer = Writer::new(send, version.clone());
writer.set_priority(CONTROL_SEND_ORDER);
Self {
writer,
reader: Reader::new(recv, version),
}
}
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 mut session = SinkSession::gated_bi(gate.consume());
let log = session.log.clone();
let _stream = Stream::open(&mut 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 mut session = SinkSession::accepted_bi(gate.consume());
let log = session.log.clone();
let _stream = Stream::accept(&mut 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",
);
}
}