diff --git a/CHANGELOG.md b/CHANGELOG.md index e1c1fc9..3ed9905 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -5,9 +5,22 @@ - Add support for unix sockets. - Add support for in-memory certificate files for the tls. - Change some of the function argument lables. +- Rename `enable_tls` to `with_tls`. +- Remove `with_name`. +- Drop the library's own `Request`, `Response` aliases. +- `ResponseBody` is now named `Body`. - Builder's `new` now requires listener and connection factory names provided in order to prevent possible atom table exhaustion. - Switch from `erlang:decode_packet` to self implemented parsing solution. +- In replacement of `stream_body` there is now `read_body_chunk`, no consumer + api anymore. +- Request loop no longer ignore request's body on the socket if it was not + consumed inside user's handler. We now flush remaining body bytes allowing + correct reuse of the socket for the next request. If the size of the bytes is + more than 1mb, we close the connection instead. +- In replacement of `chunked_body`/`send_chunk`/`chunked_continue`/`chunked_stop` + there is now `stream_response`/`send_chunk`/`finish_chunk`/`finish_response` + and no init/loop callback for response streaming anymore. ## v4.0.1 - 04.06.2026 diff --git a/benchmark/mist/src/app.gleam b/benchmark/mist/src/app.gleam index 18f8f06..3bab9bf 100644 --- a/benchmark/mist/src/app.gleam +++ b/benchmark/mist/src/app.gleam @@ -29,9 +29,11 @@ fn handle_request( |> response.set_body(mist.Bytes(bytes_tree.from_string("Hello, Joe!"))) "/whoami" -> response.new(200) - |> response.set_body(mist.Bytes(bytes_tree.from_string( - "method=" <> method_to_string(request.method), - ))) + |> response.set_body( + mist.Bytes(bytes_tree.from_string( + "method=" <> method_to_string(request.method), + )), + ) "/echo" -> case mist.read_body(request, 10_000_000) { Ok(req) -> diff --git a/dev/priv/chunked.lua b/dev/priv/chunked@echo.lua similarity index 100% rename from dev/priv/chunked.lua rename to dev/priv/chunked@echo.lua diff --git a/dev/priv/chunked@echo_chunked.lua b/dev/priv/chunked@echo_chunked.lua new file mode 100644 index 0000000..fea55c0 --- /dev/null +++ b/dev/priv/chunked@echo_chunked.lua @@ -0,0 +1,22 @@ +wrk.method = "POST" + +local chunks = { + string.rep("a", 256), + string.rep("b", 256), + string.rep("c", 256), +} + +local encoded = "" +for _, chunk in ipairs(chunks) do + encoded = encoded .. string.format("%x", #chunk) .. "\r\n" .. chunk .. "\r\n" +end +encoded = encoded .. "0\r\n\r\n" + +request = function() + return "POST /echo/chunked HTTP/1.1\r\n" .. + "Host: 127.0.0.1\r\n" .. + "Transfer-Encoding: chunked\r\n" .. + "Content-Type: application/octet-stream\r\n" .. + "Connection: keep-alive\r\n\r\n" .. + encoded +end \ No newline at end of file diff --git a/dev/serve.gleam b/dev/serve.gleam index 9f3e0db..df64527 100644 --- a/dev/serve.gleam +++ b/dev/serve.gleam @@ -34,6 +34,12 @@ fn handle_request( |> response.set_body(ewe.Bytes(bytes_tree.from_bit_array(req.body))) Error(_error) -> response.new(400) |> response.set_body(ewe.Empty) } + "/echo/chunked" -> echo_chunked(request, bytes_tree.new()) + "/stream" -> { + use writer <- ewe.stream_response(response.new(200)) + let writer = ewe.send_chunk(writer, <<"hello ":utf8>>) + ewe.finish_chunk(writer, <<"world":utf8>>) + } "/file/small" -> { // head -c 100K /dev/urandom > file_100kb.bin let assert Ok(file) = @@ -69,3 +75,18 @@ fn handle_request( |> response.set_body(ewe.Empty) } } + +fn echo_chunked( + request: request.Request(ewe.Connection), + acc: bytes_tree.BytesTree, +) -> response.Response(ewe.Body) { + case + echo ewe.read_body_chunk(request, max_chunk_bytes: 4096, limit: 10_000_000) + { + Ok(ewe.Chunk(data, request)) -> + echo_chunked(request, bytes_tree.append(acc, data)) + Ok(ewe.Done(_request)) -> + response.new(200) |> response.set_body(ewe.Bytes(acc)) + Error(_error) -> response.new(400) |> response.set_body(ewe.Empty) + } +} diff --git a/src/ewe.gleam b/src/ewe.gleam index db04422..da9cc72 100644 --- a/src/ewe.gleam +++ b/src/ewe.gleam @@ -31,6 +31,7 @@ pub type Body { Text(String) Empty File(connection.File) + Streaming(handler: fn(ResponseWriter) -> Nil) } pub type IpAddress { @@ -412,3 +413,92 @@ pub fn read_body( connection.Http2 -> todo as "HTTP/2 is not implemented yet!" } } + +// http1.Connection and connection.Connection are structurally identical. +@external(erlang, "gleam_stdlib", "identity") +fn unsafe_from_http1_connection(conn: http1.Connection) -> Connection + +/// The result of one `read_body_chunk` call. +pub type ReadEvent { + /// Up to `max_chunk_bytes` of body data. Feed `request` into the next call. + Chunk(data: BitArray, request: request.Request(Connection)) + /// The body is fully consumed. Any chunked trailer fields have been appended + /// to the returned request's `headers`. + Done(request: request.Request(Nil)) +} + +/// Pulls up to `max_chunk_bytes` of body per call instead of buffering the +/// whole body, capped overall at `limit`. Feed the request carried by `Chunk` +/// into the next call. +pub fn read_body_chunk( + req: request.Request(Connection), + max_chunk_bytes max_chunk_bytes: Int, + limit limit: Int, +) -> Result(ReadEvent, BodyError) { + case req.body { + connection.Http1(..) -> { + let conn = http1.unsafe_to_http1_connection(req.body) + case http1.read_body_chunk(conn, max_chunk_bytes:, limit:) { + Ok(http1.Chunk(data, connection)) -> { + let body = unsafe_from_http1_connection(connection) + Ok(Chunk(data, request.set_body(req, body))) + } + Ok(http1.Done(trailers)) -> { + let headers = list.append(req.headers, trailers) + Ok(Done(request.Request(..req, headers:, body: Nil))) + } + Error(error) -> Error(unsafe_from_internal_body_error(error)) + } + } + connection.Http2 -> todo as "HTTP/2 is not implemented yet!" + } +} + +/// A handle for writing a streamed response's body, obtained from +/// `stream_response`. +pub type ResponseWriter = + connection.ResponseWriter + +// http1.ResponseWriter and ResponseWriter are structurally identical. +@external(erlang, "gleam_stdlib", "identity") +fn unsafe_from_http1_writer(writer: http1.ResponseWriter) -> ResponseWriter + +/// Starts a streamed response. `handler` must end by calling `finish_chunk` or +/// `finish_response` on it, since that's what closes the stream. +pub fn stream_response( + response: response.Response(a), + handler: fn(ResponseWriter) -> Nil, +) -> response.Response(Body) { + response.set_body(response, Streaming(handler)) +} + +/// Sends one response body chunk, threading the writer through so it can be +/// piped. For the last chunk use `finish_chunk` instead, it closes the +/// stream in the same round trip. +pub fn send_chunk(writer: ResponseWriter, chunk: BitArray) -> ResponseWriter { + case writer { + connection.Http1Writer(..) -> + http1.send_chunk(http1.unsafe_to_http1_writer(writer), chunk) + |> unsafe_from_http1_writer + connection.Http2Writer -> todo as "HTTP/2 is not implemented yet!" + } +} + +/// Sends `chunk` as the final response body chunk and closes the stream. +pub fn finish_chunk(writer: ResponseWriter, chunk: BitArray) -> Nil { + case writer { + connection.Http1Writer(..) -> + http1.finish_chunk(http1.unsafe_to_http1_writer(writer), chunk) + connection.Http2Writer -> todo as "HTTP/2 is not implemented yet!" + } +} + +/// Closes the stream with no further data. Use `finish_chunk` instead if +/// there's one last chunk to send. +pub fn finish_response(writer: ResponseWriter) -> Nil { + case writer { + connection.Http1Writer(..) -> + http1.finish_response(http1.unsafe_to_http1_writer(writer)) + connection.Http2Writer -> todo as "HTTP/2 is not implemented yet!" + } +} diff --git a/src/ewe/internal/connection.gleam b/src/ewe/internal/connection.gleam index 385a290..c6111ec 100644 --- a/src/ewe/internal/connection.gleam +++ b/src/ewe/internal/connection.gleam @@ -11,12 +11,18 @@ pub type Connection { self: process.Subject(handler.Message(Message)), buffer: BitArray, framing: Framing, + // Body bytes delivered to the caller so far, via `read_body_chunk`. + read: Int, + // Bytes left in the chunked-encoding chunk currently being delivered; + // 0 means the next pull starts at a chunk boundary. Unused for `Fixed` + // and `NoBody`. + chunk_remaining: Int, ) Http2 } /// How to find the end of a request body on the wire, derived once from -/// `Content-Length`/`Transfer-Encoding` at parse time. +/// `Content-Length` or `Transfer-Encoding` at parse time. pub type Framing { Fixed(length: Int) Chunked @@ -28,6 +34,7 @@ pub type Body { Text(String) Empty File(File) + Streaming(handler: fn(ResponseWriter) -> Nil) } pub type File { @@ -36,10 +43,39 @@ pub type File { pub type Message { Timeout - /// Sent by `http1.read_body` to itself once the request body has been - /// fully consumed, carrying whatever bytes came after it. + /// Sent by `http1.read_body` and `http1.read_body_chunk` to themselves + /// once the request body has been fully consumed, carrying whatever bytes + /// came after it. BodyDrained(leftover: BitArray) - /// Sent by `http1.read_body` to itself when it gave up on the body - /// part-way through, leaving the connection in an unknown position. + /// Sent by `http1.read_body` and `http1.read_body_chunk` to themselves + /// when they gave up on the body part-way through, leaving the connection + /// in an unknown position. BodyAbandoned + /// Sent by `http1.read_body_chunk` to itself after every chunk it + /// delivers, so that if the caller stops reading before `BodyDrained`, + /// the connection can still resume draining from here instead of from the + /// start of the body. + BodyProgress(buffer: BitArray, read: Int, chunk_remaining: Int) + /// Sent by `http1.finish_chunk`/`http1.finish_response` to themselves once + /// a streamed response's terminator has been written. + StreamFinished(keep_alive: Bool) +} + +/// How a streamed response writes its body chunks to the wire. +pub type ResponseWriter { + Http1Writer( + transport: transport.Transport, + socket: socket.Socket, + self: process.Subject(handler.Message(Message)), + // `True` on HTTP/1.1. Frame each chunk with its hex size and a trailing + // `0\r\n\r\n` terminator. + // + // `False` on HTTP/1.0, which has no chunked encoding write raw bytes and + // let the body end when the connection closes. + chunked: Bool, + // Whether the connection should stay alive once the stream finishes. Always + // `False` when `chunked` is `False`. + keep_alive: Bool, + ) + Http2Writer } diff --git a/src/ewe/internal/handler.gleam b/src/ewe/internal/handler.gleam index f361479..984f432 100644 --- a/src/ewe/internal/handler.gleam +++ b/src/ewe/internal/handler.gleam @@ -5,6 +5,7 @@ import gleam/erlang/process import gleam/http/request import gleam/http/response import gleam/option +import gleam/string import glisten import glisten/internal/handler import logging @@ -60,12 +61,15 @@ pub fn loop( logging.log(logging.Debug, "Connection idled for too long, closing.") glisten.stop() } + // TODO: just use other subject brah glisten.User(connection.BodyDrained(..)) - | glisten.User(connection.BodyAbandoned) -> { - // http1.handle_message always drains these itself before returning + | glisten.User(connection.BodyAbandoned) + | glisten.User(connection.BodyProgress(..)) + | glisten.User(connection.StreamFinished(..)) -> { logging.log( logging.Critical, - "Web server loop received a message that should not be reached to the handler! That means there is a bug somewhere within the implementation. Please open an issue addressing the alert, thank you! https://github.com/vshakitskiy/ewe/issues/new", + "Web server loop received a message that should not be reached to the handler! That means there is a bug somewhere within the implementation. Please open an issue addressing the alert, thank you! https://github.com/vshakitskiy/ewe/issues/new\n\nMessage received: " + <> string.inspect(message), ) glisten.continue(state) diff --git a/src/ewe/internal/http1.gleam b/src/ewe/internal/http1.gleam index 4a5e36b..277d31b 100644 --- a/src/ewe/internal/http1.gleam +++ b/src/ewe/internal/http1.gleam @@ -57,6 +57,8 @@ pub fn handle_message( self: connection.subject, buffer: remaining, framing: metadata.framing, + read: 0, + chunk_remaining: 0, ) let request = @@ -73,26 +75,77 @@ pub fn handle_message( let response = state.handler(request) + let drained = drain_messages(connection.subject) + let #(buffer, body_drained) = - resolve_body(unsafe_to_http1_connection(body_connection)) + resolve_body(unsafe_to_http1_connection(body_connection), drained.body) + let metadata = Metadata(..metadata, keep_alive: metadata.keep_alive && body_drained) - let sent = case encode_response(response, head.method, metadata) { - Ok(Encoded(bytes:, keep_alive:, file: file_body)) -> { + let sent = case + encode_response(response, head.method, head.version, metadata) + { + Ok(Encoded(bytes:, keep_alive:, file: file_body, stream:)) -> { use Nil <- result.try(transport.send( connection.transport, connection.socket, bytes, )) - case file_body { - option.None -> Ok(keep_alive) - option.Some(data) -> + case file_body, stream { + option.None, option.None -> Ok(keep_alive) + option.Some(data), option.None -> case file.send(connection.transport, connection.socket, data) { Ok(Nil) -> Ok(keep_alive) Error(reason) -> Error(reason) } + option.None, option.Some(Stream(handler: stream_handler, chunked:)) + -> { + connection.Http1Writer( + transport: connection.transport, + socket: connection.socket, + self: connection.subject, + chunked:, + keep_alive:, + ) + |> stream_handler + + let stream_drained = drain_messages(connection.subject) + let stream_keep_alive = case stream_drained.stream { + option.Some(connection.StreamFinished(keep_alive:)) -> + keep_alive + option.None -> { + // The handler never called finish_chunk or finish_response. + // Every chunk sent so far already closed its own framing, + // so writing the terminator now still leaves the wire in + // a valid state... But the connection can't be trusted + // enough to reuse, so it closes regardless. + case chunked { + True -> { + let _ = + transport.send( + connection.transport, + connection.socket, + bytes_tree.from_bit_array(<<"0\r\n\r\n":utf8>>), + ) + Nil + } + False -> Nil + } + False + } + option.Some(connection.Timeout) + | option.Some(connection.BodyDrained(..)) + | option.Some(connection.BodyAbandoned) + | option.Some(connection.BodyProgress(..)) -> + panic as "drain_messages routes body messages to the other slot" + } + + Ok(keep_alive && stream_keep_alive) + } + option.Some(_file), option.Some(_stream) -> + panic as "encode_response never sets both file and stream" } } Error(UnsafeHeader(name)) -> { @@ -166,6 +219,8 @@ pub type Connection { self: process.Subject(handler.Message(connection.Message)), buffer: BitArray, framing: connection.Framing, + read: Int, + chunk_remaining: Int, ) } @@ -186,7 +241,7 @@ pub fn read_body( conn: Connection, limit: Int, ) -> Result(#(BitArray, List(#(String, String))), BodyError) { - let Http1(transport:, socket:, self:, buffer:, framing:) = conn + let Http1(transport:, socket:, self:, buffer:, framing:, ..) = conn case framing, consume_body(transport, socket, buffer, framing, limit) { _framing, Ok(#(body, trailers, leftover)) -> { @@ -201,23 +256,120 @@ pub fn read_body( } } -// Reconciles what the handler did with the request body so the connection can -// be safely reused. If the handler called `read_body`, picks up the leftover -// bytes it reported via a message to `conn.self`. Otherwise drains the declared -// body itself, up to `auto_drain_limit`, so a well behaved client isn't -// punished with a dropped connection just because the handler didn't care about -// its body. -fn resolve_body(conn: Connection) -> #(BitArray, Bool) { - let Http1(transport:, socket:, self:, buffer:, framing:) = conn +/// The result of one `read_body_chunk` call. +pub type ChunkRead { + /// `max_chunk_bytes` of body data. Feed `connection` into the next call. + Chunk(data: BitArray, connection: Connection) + /// The body is fully consumed. Carries any chunked trailer fields. `Fixed` + /// and `NoBody` requests never have trailers. + Done(trailers: List(#(String, String))) +} + +/// Pulls up to `max_chunk_bytes` of body per call instead of buffering the +/// whole body, capped overall at `limit`. Feed the connection carried by +/// `Chunk` into the next call. Ends with `Done` once the body, and any chunked +/// trailers, are fully consumed. +pub fn read_body_chunk( + conn: Connection, + max_chunk_bytes max_chunk_bytes: Int, + limit limit: Int, +) -> Result(ChunkRead, BodyError) { + let Http1(self:, buffer:, read:, chunk_remaining:, ..) = conn + + case pull_chunk(conn, max_chunk_bytes, limit) { + Ok(PulledChunk(data, next)) -> { + handler.User(connection.BodyProgress(buffer:, read:, chunk_remaining:)) + |> process.send(self, _) + + Ok(Chunk(data, next)) + } + Ok(PulledDone(trailers, leftover)) -> { + process.send(self, handler.User(connection.BodyDrained(leftover:))) + + Ok(Done(trailers)) + } + Error(error) -> { + process.send(self, handler.User(connection.BodyAbandoned)) + Error(error) + } + } +} + +// Cap on how much is pulled off the wire per step while blind draining an +// abandoned body. Unrelated to any caller supplied `max_chunk_bytes`. +const auto_drain_chunk_bytes = 65_536 + +// Reconciles what the handler did with the request body so the connection can +// be safely reused. Replays whatever `read_body`/`read_body_chunk` last +// reported about themselves via a message to `conn.self`, drained by +// `drain_messages` into `drained.body`: +// - `BodyDrained`: the body was fully consumed, use its leftover directly. +// - `BodyAbandoned`: the read failed mid-body, the connection can't be reused. +// - `BodyProgress` or nothing at all: the handler stopped reading or never +// started, so drain whatever's left from that point, up to `auto_drain_limit` +// more bytes, so a well behaved client isn't punished with a dropped +// connection just because the handler didn't care about its body. +fn resolve_body( + conn: Connection, + drained: option.Option(connection.Message), +) -> #(BitArray, Bool) { + case drained { + option.Some(connection.BodyDrained(leftover)) -> #(leftover, True) + option.Some(connection.BodyAbandoned) -> #(<<>>, False) + option.Some(connection.BodyProgress(buffer:, read:, chunk_remaining:)) -> + drain_remaining(Http1(..conn, buffer:, read:, chunk_remaining:)) + option.Some(connection.Timeout) | option.None -> drain_remaining(conn) + option.Some(connection.StreamFinished(..)) -> + panic as "drain_messages routes stream messages to the other slot" + } +} + +// The latest body read and response stream outcome messages currently queued +// for a connection's own subject, drained in one pass so a handler that both +// reads a request body and streams a response doesn't lose one outcome to the +// other. +type Drained { + Drained( + body: option.Option(connection.Message), + stream: option.Option(connection.Message), + ) +} + +// Selectively receives every message already queued for `self`, keeping only +// the last one of each kind. +fn drain_messages( + self: process.Subject(handler.Message(connection.Message)), +) -> Drained { + do_drain_messages(self, Drained(body: option.None, stream: option.None)) +} + +fn do_drain_messages( + self: process.Subject(handler.Message(connection.Message)), + acc: Drained, +) -> Drained { case process.receive(self, 0) { - Ok(handler.User(connection.BodyDrained(leftover))) -> #(leftover, True) - Ok(handler.User(connection.BodyAbandoned)) -> #(<<>>, False) - _nothing_pending -> - case consume_body(transport, socket, buffer, framing, auto_drain_limit) { - Ok(#(_discarded, _trailers, leftover)) -> #(leftover, True) - Error(_reason) -> #(<<>>, False) - } + Ok(handler.User(connection.StreamFinished(..) as message)) -> + do_drain_messages(self, Drained(..acc, stream: option.Some(message))) + Ok(handler.User(message)) -> + do_drain_messages(self, Drained(..acc, body: option.Some(message))) + Ok(handler.Internal(_message)) -> do_drain_messages(self, acc) + Error(Nil) -> acc + } +} + +// Drains whatever's left of `conn`'s body, up to `auto_drain_limit` more bytes +// past however much has already been read. +fn drain_remaining(conn: Connection) -> #(BitArray, Bool) { + let Http1(read:, ..) = conn + do_drain_remaining(conn, read + auto_drain_limit) +} + +fn do_drain_remaining(conn: Connection, limit: Int) -> #(BitArray, Bool) { + case pull_chunk(conn, auto_drain_chunk_bytes, limit) { + Ok(PulledChunk(_data, next)) -> do_drain_remaining(next, limit) + Ok(PulledDone(_trailers, leftover)) -> #(leftover, True) + Error(_reason) -> #(<<>>, False) } } @@ -242,9 +394,7 @@ fn consume_body( } } -fn to_body_result( - result: Result(#(BitArray, List(#(String, String)), BitArray), ParseError), -) -> Result(#(BitArray, List(#(String, String)), BitArray), BodyError) { +fn to_body_result(result: Result(a, ParseError)) -> Result(a, BodyError) { case result { Ok(value) -> Ok(value) Error(ChunkTooLarge) -> Error(BodyTooLarge) @@ -279,9 +429,9 @@ fn read_fixed( } } -// A `Transfer-Encoding: chunked` body. Each call advances by one chunk (or -// trailers), pulling more from the socket only for whichever piece is -// currently short. +// A `Transfer-Encoding: chunked` body. Each call advances by one chunk or +// trailers, pulling more from the socket only for whichever piece is currently +// short. fn read_chunked( transport: transport.Transport, socket: socket.Socket, @@ -299,11 +449,11 @@ fn read_chunked( case size { 0 -> { - use #(trailers, _state, remaining) <- result.try( - pull_until(transport, socket, remaining, fn(buffer) { - parse_headers(buffer, [], 0, initial_header_state()) - }), - ) + use #(trailers, _state, remaining) <- result.try({ + use buffer <- pull_until(transport, socket, remaining) + parse_headers(buffer, [], 0, initial_header_state()) + }) + Ok(#(bytes_tree.to_bit_array(acc), trailers, remaining)) } size -> { @@ -311,9 +461,11 @@ fn read_chunked( case total > limit { True -> Error(ChunkTooLarge) False -> { - use #(data, remaining) <- result.try( - pull_until(transport, socket, remaining, take_chunk_data(_, size)), - ) + use #(data, remaining) <- result.try({ + use buffer <- pull_until(transport, socket, remaining) + take_chunk_prefix(buffer, size, True) + }) + read_chunked( transport, socket, @@ -328,6 +480,125 @@ fn read_chunked( } } +// The result of pulling one step of a streamed body. +type Pulled { + PulledChunk(data: BitArray, connection: Connection) + PulledDone(trailers: List(#(String, String)), leftover: BitArray) +} + +// Pulls at most `max_chunk_bytes` of body off `conn`, capped overall at `limit`. +// Mirrors `consume_body`, but advances by one bounded slice instead of reading +// the whole declared body. +fn pull_chunk( + conn: Connection, + max_chunk_bytes: Int, + limit: Int, +) -> Result(Pulled, BodyError) { + let Http1(buffer:, framing:, read:, chunk_remaining:, ..) = conn + + case framing { + connection.NoBody -> Ok(PulledDone([], buffer)) + connection.Fixed(length) if length > limit -> Error(BodyTooLarge) + connection.Fixed(length) -> + pull_fixed_chunk(conn, length, read, max_chunk_bytes) |> to_body_result + connection.Chunked -> + pull_chunked_chunk(conn, limit, read, chunk_remaining, max_chunk_bytes) + |> to_body_result + } +} + +// A `Content-Length` body. Takes `min(remaining, max_chunk_bytes)`, since +// there's no wire framing inside the payload itself to bound a single pull to +// less than that. +fn pull_fixed_chunk( + conn: Connection, + length: Int, + read: Int, + max_chunk_bytes: Int, +) -> Result(Pulled, ParseError) { + let Http1(transport:, socket:, buffer:, ..) = conn + + case length - read { + 0 -> Ok(PulledDone([], buffer)) + remaining -> { + let want = int.min(remaining, max_chunk_bytes) + use #(data, _trailers, leftover) <- result.try(read_fixed( + transport, + socket, + buffer, + want, + )) + Ok(PulledChunk(data, Http1(..conn, buffer: leftover, read: read + want))) + } + } +} + +// A `Transfer-Encoding: chunked` body. At a chunk boundary, parses the next +// chunk-size line or the trailer section. Otherwise resumes taking a bounded +// slice out of the chunk already in progress. +fn pull_chunked_chunk( + conn: Connection, + limit: Int, + read: Int, + chunk_remaining: Int, + max_chunk_bytes: Int, +) -> Result(Pulled, ParseError) { + let Http1(transport:, socket:, buffer:, ..) = conn + + case chunk_remaining { + 0 -> { + use #(size, buffer) <- result.try(pull_until( + transport, + socket, + buffer, + parse_chunk_line, + )) + + case size { + 0 -> { + use #(trailers, _state, buffer) <- result.try({ + use buffer <- pull_until(transport, socket, buffer) + parse_headers(buffer, [], 0, initial_header_state()) + }) + + Ok(PulledDone(trailers, buffer)) + } + size if read + size > limit -> Error(ChunkTooLarge) + size -> + take_chunk_slice(Http1(..conn, buffer:), read, size, max_chunk_bytes) + } + } + remaining -> take_chunk_slice(conn, read, remaining, max_chunk_bytes) + } +} + +fn take_chunk_slice( + conn: Connection, + read: Int, + chunk_remaining: Int, + max_chunk_bytes: Int, +) -> Result(Pulled, ParseError) { + let Http1(transport:, socket:, buffer:, ..) = conn + + let want = int.min(chunk_remaining, max_chunk_bytes) + let final_slice = want == chunk_remaining + + use #(data, buffer) <- result.try({ + use buffer <- pull_until(transport, socket, buffer) + take_chunk_prefix(buffer, want, final_slice) + }) + + Ok(PulledChunk( + data, + Http1( + ..conn, + buffer:, + read: read + want, + chunk_remaining: chunk_remaining - want, + ), + )) +} + // Runs `step` against `buffer`, pulling more bytes from the socket only when // `step` itself reports it's short. fn pull_until( @@ -337,7 +608,7 @@ fn pull_until( step: fn(BitArray) -> Step(a), ) -> Result(a, ParseError) { case step(buffer) { - Done(value) -> Ok(value) + StepDone(value) -> Ok(value) ParseError(error) -> Error(error) More -> case transport.receive_timeout(transport, socket, 0, body_read_timeout) { @@ -356,14 +627,14 @@ fn parse_chunk_line(buffer: BitArray) -> Step(#(Int, BitArray)) { BadChunkSize, )) use size <- try_step(parse_chunk_size(line)) - Done(#(size, remaining)) + StepDone(#(size, remaining)) } // Chunk-size lines may carry `;extensions`. Only the hex size before them // matters here fn parse_chunk_size(line: BitArray) -> Step(Int) { case parse_hex_digits(line, 0, False) { - Ok(size) -> Done(size) + Ok(size) -> StepDone(size) Error(Nil) -> ParseError(BadChunkSize) } } @@ -381,14 +652,31 @@ fn parse_hex_digits(bits: BitArray, acc: Int, any: Bool) -> Result(Int, Nil) { } } -fn take_chunk_data(buffer: BitArray, size: Int) -> Step(#(BitArray, BitArray)) { - case buffer { - <> -> - Done(#(data, remaining)) - _buffer -> - case bit_array.byte_size(buffer) < size + 2 { - True -> More - False -> ParseError(BadChunkFraming) +// Takes `want` bytes off the front of a chunk's data. `final_slice` marks +// whether `want` reaches the end of the chunk, in which case the trailing CRLF +// is expected right after and consumed too. Otherwise `want` is just a prefix +// of a bigger chunk still being streamed across calls. +fn take_chunk_prefix( + buffer: BitArray, + want: Int, + final_slice: Bool, +) -> Step(#(BitArray, BitArray)) { + case final_slice { + True -> + case buffer { + <> -> + StepDone(#(data, remaining)) + _buffer -> + case bit_array.byte_size(buffer) < want + 2 { + True -> More + False -> ParseError(BadChunkFraming) + } + } + False -> + case buffer { + <> -> + StepDone(#(data, remaining)) + _buffer -> More } } } @@ -411,49 +699,209 @@ fn initial_encode_state() -> EncodeState { EncodeState(tree: bytes_tree.new(), force_close: False) } -/// The result of encoding a response: `bytes` is ready for a single -/// `transport.send`, and `file`, when present, still needs to be streamed -/// separately since it was never loaded into `bytes`. +/// A `Streaming` body's handler, plus whether its chunks get +/// `Transfer-Encoding: chunked` framing (HTTP/1.1) or written raw +/// (HTTP/1.0, which has no chunked encoding). +pub type Stream { + Stream(handler: fn(connection.ResponseWriter) -> Nil, chunked: Bool) +} + +/// The result of encoding a response. `bytes` is ready for a single +/// `transport.send`. `file` and `stream`, when present, still need handling +/// separately since neither is loaded into `bytes`; `encode_response` never +/// sets both. pub type Encoded { Encoded( bytes: bytes_tree.BytesTree, keep_alive: Bool, file: option.Option(connection.File), + stream: option.Option(Stream), ) } /// Builds the response as a `BytesTree`. `Bytes`, `Text` and `Empty` bodies /// are folded straight into it, so the whole response goes out in a single -/// `transport.send`. +/// `transport.send`. A `Streaming` body's handler isn't run here -- like +/// `File`, the caller runs it separately, only after sending `bytes`. pub fn encode_response( response: response.Response(connection.Body), method: http.Method, + version: Version, metadata: Metadata, ) -> Result(Encoded, EncodeError) { use state <- result.try(encode_headers(response.headers)) - let length = body_length(response.body) let keep_alive = metadata.keep_alive && !state.force_close - let head = - state.tree - |> append_date() - |> append_connection(keep_alive) - |> bytes_tree.prepend(status_line(response.status)) - |> bytes_tree.append_string("content-length: " <> int.to_string(length)) - |> bytes_tree.append(<<"\r\n\r\n":utf8>>) - - case method, response.body { - http.Head, _body -> Ok(Encoded(head, keep_alive, option.None)) - _method, connection.File(data) -> - Ok(Encoded(head, keep_alive, option.Some(data))) - _method, connection.Bytes(tree) -> - Ok(Encoded(bytes_tree.append_tree(head, tree), keep_alive, option.None)) - _method, connection.Text(text) -> - Ok(Encoded(bytes_tree.append_string(head, text), keep_alive, option.None)) - _method, connection.Empty -> Ok(Encoded(head, keep_alive, option.None)) + case response.body { + connection.Streaming(handler) -> + Ok(encode_stream(state, response, keep_alive, version, method, handler)) + _other_body -> { + let length = body_length(response.body) + let head = + state.tree + |> append_date() + |> append_connection(keep_alive) + |> bytes_tree.prepend(status_line(response.status)) + |> bytes_tree.append_string("content-length: " <> int.to_string(length)) + |> bytes_tree.append(<<"\r\n\r\n":utf8>>) + + case method, response.body { + http.Head, _body -> + Ok(Encoded(head, keep_alive, option.None, option.None)) + _method, connection.File(data) -> + Ok(Encoded(head, keep_alive, option.Some(data), option.None)) + _method, connection.Bytes(tree) -> + Ok(Encoded( + bytes_tree.append_tree(head, tree), + keep_alive, + option.None, + option.None, + )) + _method, connection.Text(text) -> + Ok(Encoded( + bytes_tree.append_string(head, text), + keep_alive, + option.None, + option.None, + )) + _method, connection.Empty -> + Ok(Encoded(head, keep_alive, option.None, option.None)) + _method, connection.Streaming(..) -> + panic as "the Streaming body is handled above" + } + } } } +// Builds the head for a `Streaming` body: `Transfer-Encoding: chunked` on +// HTTP/1.1, or a close-delimited body with no chunk framing on HTTP/1.0, +// which has no chunked encoding. `HEAD` never gets a body, mirroring how +// `encode_response` skips one for every other body type, so `handler` is +// simply never run. +fn encode_stream( + state: EncodeState, + response: response.Response(connection.Body), + keep_alive: Bool, + version: Version, + method: http.Method, + handler: fn(connection.ResponseWriter) -> Nil, +) -> Encoded { + case version { + Http11 -> { + let head = + state.tree + |> append_date() + |> append_connection(keep_alive) + |> bytes_tree.prepend(status_line(response.status)) + |> bytes_tree.append_string("transfer-encoding: chunked") + |> bytes_tree.append(<<"\r\n\r\n":utf8>>) + + case method { + http.Head -> Encoded(head, keep_alive, option.None, option.None) + _method -> + Encoded( + head, + keep_alive, + option.None, + option.Some(Stream(handler, True)), + ) + } + } + Http10 -> { + let head = + state.tree + |> append_date() + |> append_connection(False) + |> bytes_tree.prepend(status_line(response.status)) + |> bytes_tree.append(<<"\r\n":utf8>>) + + case method { + http.Head -> Encoded(head, False, option.None, option.None) + _method -> + Encoded(head, False, option.None, option.Some(Stream(handler, False))) + } + } + } +} + +/// How a streamed response writes its body chunks to the wire. +pub type ResponseWriter { + ResponseWriter( + transport: transport.Transport, + socket: socket.Socket, + self: process.Subject(handler.Message(connection.Message)), + // `True` on HTTP/1.1: frame each chunk with its hex size and a trailing + // `0\r\n\r\n` terminator. `False` on HTTP/1.0, which has no chunked + // encoding: write raw bytes and let the body end when the connection + // closes. + chunked: Bool, + // Whether the connection should stay alive once the stream finishes. + // Always `False` when `chunked` is `False`. + keep_alive: Bool, + ) +} + +// ResponseWriter and connection.ResponseWriter are structurally identical. +@external(erlang, "gleam_stdlib", "identity") +pub fn unsafe_to_http1_writer( + writer: connection.ResponseWriter, +) -> ResponseWriter + +fn chunk_frame(chunk: BitArray) -> bytes_tree.BytesTree { + bytes_tree.new() + |> bytes_tree.append_string(int.to_base16(bit_array.byte_size(chunk))) + |> bytes_tree.append(<<"\r\n":utf8>>) + |> bytes_tree.append(chunk) + |> bytes_tree.append(<<"\r\n":utf8>>) +} + +/// Sends one response body chunk. Framed as one `Transfer-Encoding: chunked` +/// piece on HTTP/1.1, written raw otherwise. For the last chunk, use +/// `finish_chunk` instead: it closes the stream in the same round trip. +pub fn send_chunk(writer: ResponseWriter, chunk: BitArray) -> ResponseWriter { + let bytes = case writer.chunked { + True -> chunk_frame(chunk) + False -> bytes_tree.from_bit_array(chunk) + } + let _ = transport.send(writer.transport, writer.socket, bytes) + writer +} + +/// Sends `chunk` as the final response body chunk and closes the stream. +pub fn finish_chunk(writer: ResponseWriter, chunk: BitArray) -> Nil { + let bytes = case writer.chunked { + True -> bytes_tree.append(chunk_frame(chunk), <<"0\r\n\r\n":utf8>>) + False -> bytes_tree.from_bit_array(chunk) + } + let _ = transport.send(writer.transport, writer.socket, bytes) + finish(writer) +} + +/// Closes the stream with no further data. Use `finish_chunk` instead if +/// there's one last chunk to send. +pub fn finish_response(writer: ResponseWriter) -> Nil { + case writer.chunked { + True -> { + let _ = + transport.send( + writer.transport, + writer.socket, + bytes_tree.from_bit_array(<<"0\r\n\r\n":utf8>>), + ) + Nil + } + False -> Nil + } + finish(writer) +} + +fn finish(writer: ResponseWriter) -> Nil { + process.send( + writer.self, + handler.User(connection.StreamFinished(keep_alive: writer.keep_alive)), + ) +} + // Builds the header block. fn encode_headers( headers: List(#(String, String)), @@ -586,6 +1034,7 @@ fn body_length(body: connection.Body) -> Int { connection.Text(text) -> string.byte_size(text) connection.Empty -> 0 connection.File(data) -> data.length + connection.Streaming(..) -> panic as "the Streaming body is handled above" } } @@ -718,7 +1167,7 @@ pub fn parse(buffer: BitArray) -> Result(Parsed, ParseError) { use metadata <- try_step(resolve_metadata(state, version)) - Done(Complete( + StepDone(Complete( Head(method:, host:, port:, path:, query:, version:, headers:), metadata, remaining, @@ -726,21 +1175,21 @@ pub fn parse(buffer: BitArray) -> Result(Parsed, ParseError) { } case step { - Done(parsed) -> Ok(parsed) + StepDone(parsed) -> Ok(parsed) More -> Ok(Incomplete) ParseError(error) -> Error(error) } } type Step(a) { - Done(a) + StepDone(a) More ParseError(ParseError) } fn try_step(step: Step(a), next: fn(a) -> Step(b)) -> Step(b) { case step { - Done(value) -> next(value) + StepDone(value) -> next(value) More -> More ParseError(error) -> ParseError(error) } @@ -760,20 +1209,20 @@ fn parse_request_line( use #(target, version) <- try_step(parse_target_version(target_and_version)) - Done(#(method, target, version, remaining)) + StepDone(#(method, target, version, remaining)) } fn parse_method(line: BitArray) -> Step(#(http.Method, BitArray)) { case line { - <<"GET ":utf8, remaining:bits>> -> Done(#(http.Get, remaining)) - <<"POST ":utf8, remaining:bits>> -> Done(#(http.Post, remaining)) - <<"PUT ":utf8, remaining:bits>> -> Done(#(http.Put, remaining)) - <<"DELETE ":utf8, remaining:bits>> -> Done(#(http.Delete, remaining)) - <<"HEAD ":utf8, remaining:bits>> -> Done(#(http.Head, remaining)) - <<"OPTIONS ":utf8, remaining:bits>> -> Done(#(http.Options, remaining)) - <<"PATCH ":utf8, remaining:bits>> -> Done(#(http.Patch, remaining)) - <<"TRACE ":utf8, remaining:bits>> -> Done(#(http.Trace, remaining)) - <<"CONNECT ":utf8, remaining:bits>> -> Done(#(http.Connect, remaining)) + <<"GET ":utf8, remaining:bits>> -> StepDone(#(http.Get, remaining)) + <<"POST ":utf8, remaining:bits>> -> StepDone(#(http.Post, remaining)) + <<"PUT ":utf8, remaining:bits>> -> StepDone(#(http.Put, remaining)) + <<"DELETE ":utf8, remaining:bits>> -> StepDone(#(http.Delete, remaining)) + <<"HEAD ":utf8, remaining:bits>> -> StepDone(#(http.Head, remaining)) + <<"OPTIONS ":utf8, remaining:bits>> -> StepDone(#(http.Options, remaining)) + <<"PATCH ":utf8, remaining:bits>> -> StepDone(#(http.Patch, remaining)) + <<"TRACE ":utf8, remaining:bits>> -> StepDone(#(http.Trace, remaining)) + <<"CONNECT ":utf8, remaining:bits>> -> StepDone(#(http.Connect, remaining)) _other -> parse_other_method(line) } } @@ -788,7 +1237,7 @@ fn parse_other_method(line: BitArray) -> Step(#(http.Method, BitArray)) { Error(Nil) -> ParseError(BadMethod) Ok(name) -> case http.parse_method(name) { - Ok(method) -> Done(#(method, remaining)) + Ok(method) -> StepDone(#(method, remaining)) Error(Nil) -> ParseError(BadMethod) } } @@ -806,9 +1255,9 @@ fn parse_target_version(bits: BitArray) -> Step(#(BitArray, Version)) { let target_size = size - 9 case bits { <> -> - Done(#(target, Http11)) + StepDone(#(target, Http11)) <> -> - Done(#(target, Http10)) + StepDone(#(target, Http10)) _bits -> ParseError(BadVersion) } } @@ -820,7 +1269,7 @@ fn split_target(target: BitArray) -> Step(#(String, option.Option(String))) { case find_question(target) { Error(Nil) -> { use path <- try_step(decode_component(target, BadTarget)) - Done(#(path, option.None)) + StepDone(#(path, option.None)) } Ok(position) -> { let size = bit_array.byte_size(target) @@ -833,7 +1282,7 @@ fn split_target(target: BitArray) -> Step(#(String, option.Option(String))) { >> -> { use path <- try_step(decode_component(path, BadTarget)) use query <- try_step(decode_component(query, BadTarget)) - Done(#(path, option.Some(query))) + StepDone(#(path, option.Some(query))) } _target -> ParseError(BadTarget) } @@ -843,7 +1292,7 @@ fn split_target(target: BitArray) -> Step(#(String, option.Option(String))) { fn decode_component(bits: BitArray, on_error: ParseError) -> Step(String) { case bit_array_to_string(bits) { - Ok(value) -> Done(value) + Ok(value) -> StepDone(value) Error(Nil) -> ParseError(on_error) } } @@ -852,8 +1301,8 @@ fn decode_component(bits: BitArray, on_error: ParseError) -> Step(String) { // valid for OPTIONS (RFC 9112 ยง3.2). fn validate_path(method: http.Method, path: String) -> Step(String) { case path, method { - "*", http.Options -> Done(path) - "/" <> _remaining, _method -> Done(path) + "*", http.Options -> StepDone(path) + "/" <> _remaining, _method -> StepDone(path) _path, _method -> ParseError(BadTarget) } } @@ -873,7 +1322,7 @@ fn resolve_target( case split_host_port(target) { Ok(#(host, option.Some(_port) as port)) -> case bit_array_to_string(host) { - Ok(host) -> Done(#(host, port, "", option.None)) + Ok(host) -> StepDone(#(host, port, "", option.None)) Error(Nil) -> ParseError(BadTarget) } Ok(#(_host, option.None)) -> ParseError(BadTarget) @@ -883,7 +1332,7 @@ fn resolve_target( use #(path, query) <- try_step(split_target(target)) use path <- try_step(validate_path(method, path)) use #(host, port) <- try_step(resolve_host(version, header_host)) - Done(#(host, port, path, query)) + StepDone(#(host, port, path, query)) } } } @@ -895,8 +1344,8 @@ fn resolve_host( header_host: option.Option(#(String, option.Option(Int))), ) -> Step(#(String, option.Option(Int))) { case header_host, version { - option.Some(host_port), _version -> Done(host_port) - option.None, Http10 -> Done(#("", option.None)) + option.Some(host_port), _version -> StepDone(host_port) + option.None, Http10 -> StepDone(#("", option.None)) option.None, Http11 -> ParseError(MissingHost) } } @@ -1028,7 +1477,7 @@ fn resolve_metadata(state: HeaderState, version: Version) -> Step(Metadata) { True -> state.upgrade False -> option.None } - Done(Metadata(framing:, keep_alive:, upgrade:)) + StepDone(Metadata(framing:, keep_alive:, upgrade:)) } } } @@ -1047,7 +1496,7 @@ fn parse_headers( )) case line { - <<>> -> Done(#(list.reverse(acc), state, remaining)) + <<>> -> StepDone(#(list.reverse(acc), state, remaining)) _line if count >= max_headers -> ParseError(TooManyHeaders) _line -> { use #(header, state) <- try_step(parse_header_line(line, state)) @@ -1071,7 +1520,7 @@ fn parse_header_line( case bit_array_to_string(name), bit_array_to_string(value) { Ok(name), Ok(value) -> { use state <- try_step(classify(name, value, state)) - Done(#(#(name, value), state)) + StepDone(#(#(name, value), state)) } _other, _other -> ParseError(BadHeader) } @@ -1095,14 +1544,16 @@ fn classify( option.None -> case int.parse(value) { Ok(length) if length >= 0 -> - Done(HeaderState(..state, content_length: option.Some(length))) + StepDone( + HeaderState(..state, content_length: option.Some(length)), + ) _bad -> ParseError(BadContentLength) } } "transfer-encoding" -> { let lowered = value |> bit_array.from_string |> lowercase_ascii let chunked = state.chunked || has_token(lowered, <<"chunked":utf8>>) - Done(HeaderState(..state, chunked:)) + StepDone(HeaderState(..state, chunked:)) } "connection" -> { let lowered = value |> bit_array.from_string |> lowercase_ascii @@ -1116,7 +1567,7 @@ fn classify( } let connection_upgrade = state.connection_upgrade || has_token(lowered, <<"upgrade":utf8>>) - Done(HeaderState(..state, connection:, connection_upgrade:)) + StepDone(HeaderState(..state, connection:, connection_upgrade:)) } "upgrade" -> { // ASCII-lowered bytes of a validated `String` are valid UTF-8. @@ -1125,7 +1576,7 @@ fn classify( |> bit_array.from_string |> lowercase_ascii |> unsafe_to_string - Done(HeaderState(..state, upgrade: option.Some(lowered))) + StepDone(HeaderState(..state, upgrade: option.Some(lowered))) } "host" -> case state.host { @@ -1137,12 +1588,12 @@ fn classify( case split_host_port(bit_array.from_string(value)) { Ok(#(host, port)) -> { let host = option.Some(#(unsafe_to_string(host), port)) - Done(HeaderState(..state, host:)) + StepDone(HeaderState(..state, host:)) } Error(Nil) -> ParseError(BadHost) } } - _other -> Done(state) + _other -> StepDone(state) } } @@ -1163,7 +1614,7 @@ fn extract_line( Ok(position) -> case buffer { <> -> - Done(#(line, remaining)) + StepDone(#(line, remaining)) _bad -> ParseError(malformed) } }