Skip to content
Merged
2 changes: 1 addition & 1 deletion protocols/turnloop-http/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -8,7 +8,7 @@ Own HTTP/1.1 and HTTP/2 wire engines, with no runtime or transport dependency.
- `hpack` provides the independent bounded RFC 7541 codec. The encoder never indexes credentials/cookies. Decode failures poison a context.
- `client::Pool` reserves connections before DNS/connect so parallel commands respect per-origin/proxy limits. `Route` emits resolution/TLS requests and builds proxy CONNECT or absolute-form heads. `Resolver` belongs to the host. `Request::redirect` applies Fetch redirects; use `DEFAULT_MAX_REDIRECTS` (20). JS conversions and promise delivery remain with Perry. `Pool::contains`/`forget` check or drop a stale `ConnectionId`. `Pool::next_timeout` is one O(1) deadline over idle connections and every request registered with `set_request_deadline`; drain `handle_timeout` and `handle_request_timeout`. `client::Deadlines` is the same heap for host-keyed timers.
- `multipart::Form` encodes in-memory `multipart/form-data` text and file parts. The host supplies 16 bytes of entropy; the boundary is regenerated until it occurs in no part.
- `compression::StreamingDecoder` accepts input and caller-owned output. Reuse it with `reset` to retain scratch buffers across bodies. Hosts may cache one per content encoding. `decode` is the convenience whole-body path and constructs algorithm state each time.
- `compression::StreamingDecoder` accepts input and caller-owned output. Each step says whether it stopped for input (`needs_input`), for output space, or because the body finished. `new` takes a `Content-Encoding` value, list included (`gzip, br`, decoded innermost-last); `from_codings(head.values("content-encoding"), limit)` merges repeated header lines. Reuse it with `reset` to retain scratch buffers across bodies. Hosts may cache one per content encoding. `decode` is the convenience whole-body path and constructs algorithm state each time.

No engine samples a clock. The host passes `Instant` deadlines and invokes timeout handlers. Hold output storage stable until a completion-shaped write finishes: do not mutate the engine while an I/O operation borrows its output. Error codes are transport causes; Perry creates the JS error objects and detailed OS diagnostics.

Expand Down
78 changes: 76 additions & 2 deletions protocols/turnloop-http/src/asynchronous/client.rs
Original file line number Diff line number Diff line change
Expand Up @@ -652,8 +652,24 @@ impl Decoders {
if !options.decompress {
return Ok(());
}
let encoding = std::str::from_utf8(head.get("content-encoding").unwrap_or(b"identity"))
.map_err(io::Error::other)?;
// Repeated Content-Encoding lines are one list (RFC 9110 section 5.3).
// The common single-line case borrows; only a repeated field pays for
// joining the lines into one cache key.
let joined;
let encoding = if head.values("content-encoding").nth(1).is_some() {
let mut key = Vec::new();
for value in head.values("content-encoding") {
if !key.is_empty() {
key.extend_from_slice(b", ");
}
key.extend_from_slice(value);
}
joined = String::from_utf8(key).map_err(io::Error::other)?;
joined.as_str()
} else {
std::str::from_utf8(head.get("content-encoding").unwrap_or(b"identity"))
.map_err(io::Error::other)?
};
if encoding == "identity" {
return Ok(());
}
Expand Down Expand Up @@ -714,3 +730,61 @@ impl Decoders {
Ok(())
}
}
#[cfg(test)]
mod tests {
use super::*;
use std::io::Write;

fn gzip(bytes: &[u8]) -> Vec<u8> {
let mut encoder = flate2::write::GzEncoder::new(Vec::new(), flate2::Compression::fast());
encoder.write_all(bytes).expect("gzip into a Vec");
encoder.finish().expect("gzip into a Vec")
}
fn head(encodings: &[&str]) -> Head {
Head {
method: String::new(),
target: String::new(),
status: 200,
version: 1,
headers: encodings
.iter()
.map(|value| Header::new("Content-Encoding", value))
.collect(),
keep_alive: true,
}
}
fn decode(decoders: &mut Decoders, head: &Head, wire: &[u8]) -> io::Result<Vec<u8>> {
decoders.start(head, &Options::default())?;
let mut out = Vec::new();
decoders.feed(wire, true, &mut |bytes| {
out.extend_from_slice(bytes);
Ok(())
})?;
Ok(out)
}

/// PerryTS/turnloop#79: repeated Content-Encoding lines are one list, so
/// two `gzip` lines decode twice rather than once.
#[test]
fn repeated_content_encoding_lines_decode_as_one_list() {
let body = b"a body compressed twice, one coding per header line".repeat(50);
let wire = gzip(&gzip(&body));
let mut decoders = Decoders::default();
assert_eq!(
decode(&mut decoders, &head(&["gzip", "gzip"]), &wire).unwrap(),
body
);
// The same chain on one line shares the joined cache entry.
assert_eq!(
decode(&mut decoders, &head(&["gzip, gzip"]), &wire).unwrap(),
body
);
assert_eq!(decoders.cache.len(), 1);
// A single line still decodes once.
let once = gzip(&body);
assert_eq!(
decode(&mut decoders, &head(&["gzip"]), &once).unwrap(),
body
);
}
}
5 changes: 4 additions & 1 deletion protocols/turnloop-http/src/asynchronous/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -28,15 +28,18 @@ impl<S: Stream> Http1<S> {
upgraded: false,
}
}
/// A request that asked to upgrade ends the message too, and a server that
/// declined it may reuse the connection. A response upgrade never can.
pub fn reusable(&self) -> bool {
self.stream.is_some() && self.ended && self.decoder.reusable()
self.stream.is_some() && (self.ended || self.upgraded) && self.decoder.reusable()
}
pub fn abort(&mut self) {
self.stream.take();
}
pub fn reset(&mut self) -> io::Result<()> {
self.decoder.reset().map_err(io::Error::other)?;
self.ended = false;
self.upgraded = false;
Ok(())
}
pub fn response_to(&mut self, method: &str) {
Expand Down
22 changes: 20 additions & 2 deletions protocols/turnloop-http/src/asynchronous/server.rs
Original file line number Diff line number Diff line change
Expand Up @@ -253,6 +253,7 @@ pub struct Response {
output: Vec<u8>,
finished: bool,
keep_alive: bool,
status: u16,
}
impl Default for Response {
fn default() -> Self {
Expand All @@ -261,6 +262,7 @@ impl Default for Response {
output: Vec::with_capacity(65536),
finished: false,
keep_alive: true,
status: 0,
}
}
}
Expand All @@ -275,6 +277,7 @@ impl Response {
}
self.encoder = Some(super::encode_head(head, length, &mut self.output)?);
self.keep_alive = head.keep_alive && !head.token("connection", "close");
self.status = head.status;
Ok(())
}
pub fn body(&mut self, bytes: &[u8]) -> io::Result<()> {
Expand All @@ -299,6 +302,12 @@ impl Response {
/// A connection that ends after a response (shutdown, `connection: close` or a
/// non-reusable request) and an idle connection at shutdown close with a lingering
/// close bounded by [`Options::linger_timeout`] and [`Shutdown::stop_by`].
///
/// A request asking to upgrade (or a CONNECT) ends with
/// [`Event::Upgrade`](crate::http1::Event::Upgrade) in place of `End`. A service
/// that declines answers normally and the connection stays HTTP/1. This driver
/// cannot hand the transport over, so a `101` (or a `2xx` to CONNECT) closes the
/// connection after the response. Serve upgrades from [`super::Http1`] directly.
pub async fn http1<S: HalfClose>(
stream: S,
shutdown: Shutdown,
Expand Down Expand Up @@ -326,18 +335,23 @@ pub async fn http1<S: HalfClose>(
return Ok(());
};
let head = head?;
let connect = head.method == "CONNECT";
service(crate::http1::Event::Head(head), &mut response)?;
turnloop_io::write_all(
conn.stream.as_mut().ok_or_else(super::closed)?,
&response.output,
)
.await?;
response.output.clear();
let mut upgrade = false;
loop {
let mut ended = false;
let received = conn
.event(|event| {
ended = matches!(event, crate::http1::Event::End);
// An upgrade request ends like any other; the service
// decides whether to switch by the status it answers.
upgrade = matches!(event, crate::http1::Event::Upgrade);
ended = upgrade || matches!(event, crate::http1::Event::End);
service(event, &mut response)
})
.await?;
Expand All @@ -357,7 +371,11 @@ pub async fn http1<S: HalfClose>(
if !response.finished {
return Err(io::Error::other("service did not finish response"));
}
if shutdown.is_stopped() || !conn.reusable() || !response.keep_alive {
// This driver has no handoff, so an accepted upgrade ends the connection
// rather than parsing the next protocol's bytes as HTTP/1.
let switched =
upgrade && (response.status == 101 || connect && (200..300).contains(&response.status));
if shutdown.is_stopped() || switched || !conn.reusable() || !response.keep_alive {
linger(&mut conn.stream, &mut conn.input, &shutdown).await;
return Ok(());
}
Expand Down
Loading