#![allow(dead_code)]
use moq_net::{Client, Server, Session, Version, origin};
use super::mock::create_mock_session_pair;
pub use moq_net::time::run;
pub fn now() -> moq_net::time::Instant {
tokio::time::Instant::now().into_std()
}
pub struct MockConnectOptions {
pub version: Version,
pub client_publish: Option<origin::Producer>,
pub client_subscribe: Option<origin::Producer>,
pub server_publish: Option<origin::Producer>,
pub server_subscribe: Option<origin::Producer>,
}
impl MockConnectOptions {
pub fn new(version: Version) -> Self {
Self {
version,
client_publish: None,
client_subscribe: None,
server_publish: None,
server_subscribe: None,
}
}
}
pub struct MockPair {
pub client: Session,
pub server: Session,
}
pub async fn connect_mock(opts: MockConnectOptions) -> MockPair {
let protocol = opts.version.alpn();
let (client_transport, server_transport) = create_mock_session_pair(Some(protocol));
let mut client = Client::new().with_versions(opts.version.into());
if let Some(publish) = &opts.client_publish {
client = client.with_publisher(publish);
}
if let Some(subscribe) = opts.client_subscribe {
client = client.with_subscriber(subscribe);
}
let mut server = Server::new().with_versions(opts.version.into());
if let Some(publish) = &opts.server_publish {
server = server.with_publisher(publish);
}
if let Some(subscribe) = opts.server_subscribe {
server = server.with_subscriber(subscribe);
}
let client_fut = async {
let (session, driver) = client
.connect(now(), client_transport)
.await
.expect("client handshake failed");
tokio::spawn(run(driver));
session
};
let server_fut = async {
let (session, driver) = server
.accept(now(), server_transport)
.await
.expect("server handshake failed");
tokio::spawn(run(driver));
session
};
let (client_session, server_session) = tokio::join!(client_fut, server_fut);
MockPair {
client: client_session,
server: server_session,
}
}