mod duplex;
mod request;
mod response;
use std::sync::Arc;
pub use duplex::H3DuplexStream;
pub use request::RequestBodyWriter;
use request::{WriteGuard, request_headers};
use response::ReadGuard;
pub use response::{H3ResponseBody, ResponseFut};
use crate::{
h3::client::{
app::{Http3ClientApp, register_stream},
error::RequestError,
},
quic::connection::{ConnectionHandle, QuicScionConn, WeakConnectionHandle},
};
pub(crate) fn initiate_request(
handle: &ConnectionHandle<Http3ClientApp>,
req: http::Request<()>,
) -> Result<(ResponseFut, RequestBodyWriter), RequestError> {
let (parts, ()) = req.into_parts();
let headers = request_headers(&parts);
let stream_id = {
let mut guard = handle.lock();
let QuicScionConn { inner, app, .. } = &mut *guard;
let Some(h3) = app.h3.as_mut() else {
return Err(RequestError::ConnectionClosed);
};
match h3.send_request(inner, &headers, false) {
Ok(stream_id) => {
register_stream(app, stream_id);
stream_id
}
Err(squiche::h3::Error::StreamBlocked) => return Err(RequestError::StreamBlocked),
Err(err) => return Err(RequestError::H3(err)),
}
};
handle.notify();
let stream_ref = StreamRef::new(handle.downgrade(), stream_id);
let response = ResponseFut::new(ReadGuard::new(stream_ref.clone()));
let writer = RequestBodyWriter::new(WriteGuard::new(stream_ref));
Ok((response, writer))
}
pub(crate) struct StreamRef {
handle: WeakConnectionHandle<Http3ClientApp>,
stream_id: u64,
}
impl StreamRef {
pub(crate) fn new(handle: WeakConnectionHandle<Http3ClientApp>, stream_id: u64) -> Arc<Self> {
Arc::new(Self { handle, stream_id })
}
pub(crate) fn handle(&self) -> WeakConnectionHandle<Http3ClientApp> {
self.handle.clone()
}
pub(crate) fn stream_id(&self) -> u64 {
self.stream_id
}
}
impl Drop for StreamRef {
fn drop(&mut self) {
let Some(handle) = self.handle.upgrade() else {
return;
};
let mut guard = handle.lock_recovering();
guard.app.streams.remove(&self.stream_id);
guard.app.response_heads.remove(&self.stream_id);
}
}