diff --git a/CHANGELOG.md b/CHANGELOG.md index 9721ae1..321358d 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -27,7 +27,8 @@ - Add `Http1Options`, `default_http1_options` and `with_http1` to set the limits and timeouts for every HTTP/1 connection. The request line, header line and header count limits, chunk size line limit, idle and body read timeouts, and - the auto drain limit and chunk size. + the auto drain limit and chunk size. A value outside the range a field accepts + is replaced with the default and logged as a warning when the server starts. - Move `idle_timeout` from a builder function to a field of `Http1Options`. ### HTTP/2 @@ -40,7 +41,8 @@ their refill marks, frame and header list sizes, the HPACK table size, the CONTINUATION and header block caps, the Rapid Reset window and threshold, the handshake, drain and body read timeouts, and the file read threshold. A value - the protocol does not allow is replaced with the default. + outside the range a field accepts is replaced with the default and logged as a + warning when the server starts. ### Requests and responses diff --git a/examples/src/client_info.gleam b/examples/src/client_info.gleam index 87b9872..917fcf6 100644 --- a/examples/src/client_info.gleam +++ b/examples/src/client_info.gleam @@ -47,9 +47,9 @@ fn handle_request( } // Behind a proxy this is the proxy's address and not the browser's. The address -// the proxy puts in `x-forwarded-for` is the one to use there but only when -// the proxy is yours, since any client can send that header itself. -// +// the proxy puts in `x-forwarded-for` is the one to use there. +// +// See https://developer.mozilla.org/en-US/docs/Web/HTTP/Reference/Headers/X-Forwarded-For fn describe_client(connection: ewe.Connection) -> String { case ewe.get_client_info(connection) { Ok(ewe.TcpSocketAddress(ip_address:, port:)) -> { diff --git a/src/ewe.gleam b/src/ewe.gleam index c33e727..524135f 100644 --- a/src/ewe.gleam +++ b/src/ewe.gleam @@ -371,6 +371,9 @@ fn to_internal_tls_key_type(key_type: TlsKeyType) -> options.TlsKeyType { /// Sizes are in bytes and timeouts in milliseconds. Build one by updating /// `default_http1_options` so you only state the ones you care about. /// +/// A value outside the range a field accepts is replaced with the default and +/// logged as a warning when the server starts. +/// /// # Examples /// /// ```gleam @@ -406,7 +409,7 @@ pub type Http1Options { /// Get the default HTTP/1 limits and timeouts to be adjusted and given to /// `with_http1`. pub fn default_http1_options() -> Http1Options { - let http1.Config( + let http1.Options( max_request_line:, max_header_line:, max_headers:, @@ -415,7 +418,7 @@ pub fn default_http1_options() -> Http1Options { body_read_timeout:, auto_drain_limit:, auto_drain_chunk_bytes:, - ) = http1.default_config() + ) = http1.default_options() Http1Options( max_request_line:, @@ -429,27 +432,113 @@ pub fn default_http1_options() -> Http1Options { ) } -fn to_internal_http1_options(options: Http1Options) -> http1.Config { - let Http1Options( - max_request_line:, - max_header_line:, - max_headers:, - max_chunk_size_line:, - idle_timeout:, - body_read_timeout:, - auto_drain_limit:, - auto_drain_chunk_bytes:, - ) = options +fn out_of_range(field: String, value: String, default: String) -> Nil { + logging.log( + logging.Warning, + field <> " of " <> value <> " is out of range, using " <> default, + ) +} - http1.Config( - max_request_line:, - max_header_line:, - max_headers:, - max_chunk_size_line:, - idle_timeout:, - body_read_timeout:, - auto_drain_limit:, - auto_drain_chunk_bytes:, +fn at_least(value: Int, minimum: Int, default: Int, field: String) -> Int { + case value >= minimum { + True -> value + False -> { + out_of_range(field, int.to_string(value), int.to_string(default)) + default + } + } +} + +fn within( + value: Int, + minimum: Int, + maximum: Int, + default: Int, + field: String, +) -> Int { + case value >= minimum && value <= maximum { + True -> value + False -> { + out_of_range(field, int.to_string(value), int.to_string(default)) + default + } + } +} + +fn optional_at_least( + value: Option(Int), + minimum: Int, + default: Option(Int), + field: String, +) -> Option(Int) { + case value { + Some(limit) if limit < minimum -> { + out_of_range(field, int.to_string(limit), limit_to_string(default)) + default + } + _value -> value + } +} + +fn limit_to_string(limit: Option(Int)) -> String { + case limit { + Some(limit) -> int.to_string(limit) + None -> "no limit" + } +} + +fn to_internal_http1_options(options: Http1Options) -> http1.Options { + let defaults = http1.default_options() + + http1.Options( + max_request_line: at_least( + options.max_request_line, + 1, + defaults.max_request_line, + "max_request_line", + ), + max_header_line: at_least( + options.max_header_line, + 1, + defaults.max_header_line, + "max_header_line", + ), + max_headers: at_least( + options.max_headers, + 1, + defaults.max_headers, + "max_headers", + ), + max_chunk_size_line: at_least( + options.max_chunk_size_line, + 1, + defaults.max_chunk_size_line, + "max_chunk_size_line", + ), + idle_timeout: at_least( + options.idle_timeout, + 1, + defaults.idle_timeout, + "idle_timeout", + ), + body_read_timeout: at_least( + options.body_read_timeout, + 1, + defaults.body_read_timeout, + "body_read_timeout", + ), + auto_drain_limit: at_least( + options.auto_drain_limit, + 0, + defaults.auto_drain_limit, + "auto_drain_limit", + ), + auto_drain_chunk_bytes: at_least( + options.auto_drain_chunk_bytes, + 1, + defaults.auto_drain_chunk_bytes, + "auto_drain_chunk_bytes", + ), ) } @@ -458,9 +547,10 @@ const max_window_size = 2_147_483_647 /// The limits and timeouts applied to every HTTP/2 connection. /// /// Sizes are in bytes and timeouts in milliseconds. Build one by updating -/// `default_http2_options` so you only state the ones you care about. A value -/// the protocol does not allow is replaced with the default rather than being -/// sent to a peer. +/// `default_http2_options` so you only state the ones you care about. +/// +/// A value outside the range a field accepts is replaced with the default and +/// logged as a warning when the server starts. /// /// # Examples /// @@ -514,7 +604,7 @@ pub type Http2Options { /// Get the default HTTP/2 limits and timeouts to be adjusted and given to /// `with_http2`. pub fn default_http2_options() -> Http2Options { - let http2.Config( + let http2.Options( max_concurrent_streams:, initial_window_size:, max_frame_size:, @@ -530,7 +620,7 @@ pub fn default_http2_options() -> Http2Options { recv_window_high_water_mark:, file_read_threshold:, body_read_timeout:, - ) = http2.default_config() + ) = http2.default_options() Http2Options( max_concurrent_streams:, @@ -551,64 +641,113 @@ pub fn default_http2_options() -> Http2Options { ) } -/// Anything the protocol rules out would break connections so it is dropped -/// for the default here instead of reaching a peer. -fn to_internal_http2_options(options: Http2Options) -> http2.Config { - let defaults = http2.default_config() - - let max_concurrent_streams = case options.max_concurrent_streams { - Some(limit) if limit <= 0 -> None - limit -> limit - } - - let initial_window_size = case options.initial_window_size { - size if size < 0 || size > max_window_size -> defaults.initial_window_size - size -> size - } - - let max_frame_size = case options.max_frame_size { - size if size < 16_384 || size > 16_777_215 -> defaults.max_frame_size - size -> size - } - - let drain_timeout_ms = case options.drain_timeout { - timeout if timeout <= 0 -> defaults.drain_timeout_ms - timeout -> timeout - } +fn to_internal_http2_options(options: Http2Options) -> http2.Options { + let defaults = http2.default_options() - let file_read_threshold = case options.file_read_threshold { - threshold if threshold < 0 -> defaults.file_read_threshold - threshold -> threshold - } - - // The marks only mean anything as a pair so a bad one replaces both. let #(recv_window_low_water_mark, recv_window_high_water_mark) = case options.recv_window_low_water_mark, options.recv_window_high_water_mark { low, high if low > 0 && low < high && high <= max_window_size -> #(low, high) - _low, _high -> #( - defaults.recv_window_low_water_mark, - defaults.recv_window_high_water_mark, - ) + low, high -> { + out_of_range( + "recv window water marks", + int.to_string(low) <> " and " <> int.to_string(high), + int.to_string(defaults.recv_window_low_water_mark) + <> " and " + <> int.to_string(defaults.recv_window_high_water_mark), + ) + + #( + defaults.recv_window_low_water_mark, + defaults.recv_window_high_water_mark, + ) + } } - http2.Config( - max_concurrent_streams:, - initial_window_size:, - max_frame_size:, - max_header_list_size: options.max_header_list_size, - header_table_size: options.header_table_size, - max_continuation_frames: options.max_continuation_frames, - max_header_block_bytes: options.max_header_block_bytes, - rapid_reset_window_ms: options.rapid_reset_window, - rapid_reset_threshold: options.rapid_reset_threshold, - handshake_timeout_ms: options.handshake_timeout, - drain_timeout_ms:, + http2.Options( + max_concurrent_streams: optional_at_least( + options.max_concurrent_streams, + 1, + defaults.max_concurrent_streams, + "max_concurrent_streams", + ), + initial_window_size: within( + options.initial_window_size, + 0, + max_window_size, + defaults.initial_window_size, + "initial_window_size", + ), + max_frame_size: within( + options.max_frame_size, + 16_384, + 16_777_215, + defaults.max_frame_size, + "max_frame_size", + ), + max_header_list_size: optional_at_least( + options.max_header_list_size, + 1, + defaults.max_header_list_size, + "max_header_list_size", + ), + header_table_size: at_least( + options.header_table_size, + 0, + defaults.header_table_size, + "header_table_size", + ), + max_continuation_frames: at_least( + options.max_continuation_frames, + 1, + defaults.max_continuation_frames, + "max_continuation_frames", + ), + max_header_block_bytes: at_least( + options.max_header_block_bytes, + 1, + defaults.max_header_block_bytes, + "max_header_block_bytes", + ), + rapid_reset_window_ms: at_least( + options.rapid_reset_window, + 1, + defaults.rapid_reset_window_ms, + "rapid_reset_window", + ), + rapid_reset_threshold: at_least( + options.rapid_reset_threshold, + 1, + defaults.rapid_reset_threshold, + "rapid_reset_threshold", + ), + handshake_timeout_ms: at_least( + options.handshake_timeout, + 1, + defaults.handshake_timeout_ms, + "handshake_timeout", + ), + drain_timeout_ms: at_least( + options.drain_timeout, + 1, + defaults.drain_timeout_ms, + "drain_timeout", + ), recv_window_low_water_mark:, recv_window_high_water_mark:, - file_read_threshold:, - body_read_timeout: options.body_read_timeout, + file_read_threshold: at_least( + options.file_read_threshold, + 0, + defaults.file_read_threshold, + "file_read_threshold", + ), + body_read_timeout: at_least( + options.body_read_timeout, + 1, + defaults.body_read_timeout, + "body_read_timeout", + ), ) } @@ -874,6 +1013,7 @@ pub fn start( ), loop: handler_.loop, ) + |> glisten.with_http2 let pool = case builder.tls { Some(Disk(cert:, key:)) -> @@ -889,13 +1029,6 @@ pub fn start( None -> pool } - // h2c needs nothing announced but over TLS a client only knows HTTP/2 is on - // offer if ALPN says so. - let pool = case builder.tls { - Some(_tls) -> glisten.with_http2(pool) - None -> pool - } - let pool = case builder.client_verification { Some(ca_cert) -> glisten.with_client_verification( @@ -1105,6 +1238,8 @@ pub fn read_body_chunk( max_chunk_bytes max_chunk_bytes: Int, limit limit: Int, ) -> Result(ReadEvent, BodyError) { + let max_chunk_bytes = int.max(max_chunk_bytes, 1) + case req.body { connection.Http1(connection) -> { case http1_body.read_body_chunk(connection, max_chunk_bytes:, limit:) { @@ -1224,8 +1359,6 @@ pub fn socket_reason_to_string(reason: SocketReason) -> String { fn from_interrupted(interrupted: http2.Interrupted) -> SendError { case interrupted { http2.StreamReset -> StreamReset - // A write is never given a deadline, so the only way one reports a - // timeout is the connection having stopped answering at all. http2.ConnectionClosed | http2.TimedOut -> ConnectionClosed } } @@ -1291,8 +1424,6 @@ pub fn stream_response( response: response.Response(a), handler: fn(ResponseWriter) -> Result(Nil, SendError), ) -> response.Response(Body) { - // What the handler was left holding when a write failed is its own business, - // ewe already learns whether the stream finished from the writer. let stream = fn(writer) { let _sent = handler(writer) Nil diff --git a/src/ewe/internal/connection.gleam b/src/ewe/internal/connection.gleam index aa2f6c7..3843ee6 100644 --- a/src/ewe/internal/connection.gleam +++ b/src/ewe/internal/connection.gleam @@ -26,8 +26,6 @@ pub type Sse { SseMetadata(handler: fn(SseConnection) -> Outcome) } -/// The context is built during the handshake where the negotiated extensions -/// are known and handed to whichever protocol goes on to run the socket. pub type Websocket { WebsocketMetadata( context: websocks.Context, @@ -39,13 +37,8 @@ pub type Streaming { StreamingMetadata(handler: fn(ResponseWriter) -> Nil) } -/// A raw descriptor belongs to the process that opened it, so whether the -/// handler can carry one depends on the protocol. An HTTP/1 handler runs in the -/// process that writes the socket, an HTTP/2 stream handler does not. pub type File { - /// Already open, and closed by whoever writes or drops the response. OpenFile(handle: FileDescriptor, offset: Int, length: Int) - /// Sized but not yet open; the connection process opens it at send time. PendingFile(path: String, offset: Int, length: Int) } @@ -61,8 +54,6 @@ pub type SseConnection { Http2Sse(http2.SseConnection(Body)) } -/// HTTP/2 carries WebSockets over extended CONNECT (RFC 8441), which ewe does -/// not negotiate yet. pub type WebsocketConnection { Http1Websocket(http1.WebsocketConnection) } @@ -78,12 +69,9 @@ pub type Message { Http2Stream(http2.Reply(Body)) Http2Exit(process.ExitMessage) Http2Drain - /// A stream process that outstayed the grace it was given after a reset. Http2StreamClose(pid: process.Pid) } -/// Concatenating onto an empty buffer would copy the incoming bytes for -/// nothing, which is the common case on a connection with no pipelining. pub fn append_buffer(buffer: BitArray, data: BitArray) -> BitArray { case buffer { <<>> -> data diff --git a/src/ewe/internal/ewe_ffi.erl b/src/ewe/internal/ewe_ffi.erl index 333646a..2919a78 100644 --- a/src/ewe/internal/ewe_ffi.erl +++ b/src/ewe/internal/ewe_ffi.erl @@ -72,7 +72,6 @@ skip_ascii(<>) when skip_ascii(Rest); skip_ascii(<>) when Word band ?HIGH_BITS =:= 0 -> skip_ascii(Rest); -%% Tails shorter than a word, each masked to its own width. skip_ascii(<>) when Word band 16#808080808080 =:= 0 -> <<>>; skip_ascii(<>) when Word band 16#8080808080 =:= 0 -> <<>>; skip_ascii(<>) when Word band 16#80808080 =:= 0 -> <<>>; diff --git a/src/ewe/internal/file.gleam b/src/ewe/internal/file.gleam index 0dde35b..cdda44b 100644 --- a/src/ewe/internal/file.gleam +++ b/src/ewe/internal/file.gleam @@ -64,7 +64,6 @@ fn range( } } -/// Hands back a descriptor. pub fn release(file: connection.File) -> Nil { case file { connection.OpenFile(handle:, ..) -> close(handle) @@ -92,14 +91,11 @@ pub fn send( case file { connection.OpenFile(handle:, offset:, length:) -> send_handle(transport, socket, handle, offset, length) - // Only HTTP/2 leaves a file unopened and its connection process opens one - // itself rather than writing it through here. connection.PendingFile(..) -> panic as "an unopened file cannot be written to an HTTP/1 socket" } } -/// Owns the descriptor from here on so it is closed however the write ends. fn send_handle( transport: transport.Transport, socket: socket.Socket, @@ -113,9 +109,6 @@ fn send_handle( sent } -/// Writes one range of an open file leaving the descriptor open. HTTP/2 sizes -/// each range to the stream's send window and comes back for the next one so -/// the descriptor has to outlive the individual write. pub fn send_chunk( transport: transport.Transport, socket: socket.Socket, diff --git a/src/ewe/internal/file_ffi.erl b/src/ewe/internal/file_ffi.erl index d0a5e1e..6a423b5 100644 --- a/src/ewe/internal/file_ffi.erl +++ b/src/ewe/internal/file_ffi.erl @@ -44,8 +44,6 @@ pread(Fd, Offset, Length) -> read_range(_Path, _Offset, 0) -> {ok, <<>>}; -%% A short read means the file shrank between being sized and being read which -%% leaves no correct body to send so it is reported rather than padded over. read_range(Path, Offset, Length) -> case open(Path) of {ok, Fd} -> diff --git a/src/ewe/internal/handler.gleam b/src/ewe/internal/handler.gleam index 0a68e63..0a4dad5 100644 --- a/src/ewe/internal/handler.gleam +++ b/src/ewe/internal/handler.gleam @@ -12,10 +12,8 @@ import glisten/socket/options import glisten/transport import logging -/// The connection's protocol, still undecided until the HTTP/2 preface has been -/// ruled in or out. pub type State { - Initialised(http1.State, http2_connection.Config) + Initialised(http1.State, http2_connection.Options) Http1(http1.State) Http2(http2.State) } @@ -23,8 +21,8 @@ pub type State { pub fn on_init( handler: fn(request.Request(connection.Connection)) -> response.Response(connection.Body), - http1_config: http1_connection.Config, - http2_config: http2_connection.Config, + http1_options: http1_connection.Options, + http2_options: http2_connection.Options, ) { fn(connection: glisten.Connection(connection.Message)) -> #( State, @@ -36,12 +34,12 @@ pub fn on_init( buffer: <<>>, idle_timer: connection.start_idle_timer( connection, - http1_config.idle_timeout, + http1_options.idle_timeout, ), - config: http1_config, + options: http1_options, ) - #(Initialised(state, http2_config), option.None) + #(Initialised(state, http2_options), option.None) } } @@ -51,7 +49,7 @@ pub fn loop( connection: glisten.Connection(connection.Message), ) -> glisten.Next(State, glisten.Message(connection.Message)) { case state, message { - Initialised(state, http2_config), glisten.Packet(data) -> { + Initialised(state, http2_options), glisten.Packet(data) -> { connection.cancel_idle_timer(state.idle_timer) let buffer = connection.append_buffer(state.buffer, data) @@ -62,13 +60,13 @@ pub fn loop( buffer:, idle_timer: connection.start_idle_timer( connection, - state.config.idle_timeout, + state.options.idle_timeout, ), ) - |> Initialised(http2_config) + |> Initialised(http2_options) |> glisten.continue Http2Preface(remaining:) -> - start_http2(connection, state.handler, http2_config, remaining) + start_http2(connection, state.handler, http2_options, remaining) NotHttp2(buffer:) -> http1.State(..state, buffer:, idle_timer: option.None) |> http1.handle_message(connection) @@ -94,13 +92,11 @@ pub fn loop( } } -/// Takes the connection over for HTTP/2. The client's already waiting on our -/// SETTINGS by now, so that goes out first. fn start_http2( connection: glisten.Connection(connection.Message), handler: fn(request.Request(connection.Connection)) -> response.Response(connection.Body), - config: http2_connection.Config, + options: http2_connection.Options, remaining: BitArray, ) -> glisten.Next(State, glisten.Message(connection.Message)) { process.trap_exits(True) @@ -109,8 +105,9 @@ fn start_http2( let replies = process.new_subject() let peer = transport.peername(connection.transport, connection.socket) + let parent = http2_connection.parent_pid() - let state = http2.init(handler, config, self, replies, peer) + let state = http2.init(handler, options, self, replies, peer, parent) let selector = process.new_selector() diff --git a/src/ewe/internal/http1.gleam b/src/ewe/internal/http1.gleam index 6e870c9..0001f10 100644 --- a/src/ewe/internal/http1.gleam +++ b/src/ewe/internal/http1.gleam @@ -23,7 +23,7 @@ pub type State { response.Response(connection.Body), buffer: BitArray, idle_timer: option.Option(process.Timer), - config: http1.Config, + options: http1.Options, ) } @@ -45,7 +45,7 @@ pub fn handle_message( ) -> Next { connection.cancel_idle_timer(state.idle_timer) - case parser.parse(state.buffer, state.config) { + case parser.parse(state.buffer, state.options) { Ok(parser.Complete(head, metadata, remaining)) -> { let self = process.new_subject() @@ -58,7 +58,7 @@ pub fn handle_message( framing: metadata.framing, read: 0, chunk_remaining: 0, - config: state.config, + options: state.options, upgrade: metadata.upgrade, ) @@ -69,7 +69,7 @@ pub fn handle_message( Ok(response) -> { let drained = drain_messages(self) let ResolvedBody(buffer, body_keep_alive) = - resolve_body(body_connection, drained.body, state.config) + resolve_body(body_connection, drained.body, state.options) let keep_alive = http1.and_keep_alive(metadata.keep_alive, body_keep_alive) @@ -115,7 +115,7 @@ pub fn handle_message( } Ok(parser.Incomplete) -> { let idle_timer = - connection.start_idle_timer(connection, state.config.idle_timeout) + connection.start_idle_timer(connection, state.options.idle_timeout) Continue(State(..state, idle_timer:)) } Error(error) -> { @@ -136,8 +136,6 @@ pub fn handle_message( } } -/// A handler that crashed has said nothing about what it read or meant to -/// send so the connection is answered and dropped rather than handed back. fn crashed( connection: glisten.Connection(connection.Message), details: String, @@ -157,8 +155,6 @@ fn crashed( Close } -/// A pipelining client sends its next request without waiting for this -/// response so anything left over has to be handled now. fn await_next_request( state: State, buffer: BitArray, @@ -167,7 +163,7 @@ fn await_next_request( case buffer { <<>> -> { let idle_timer = - connection.start_idle_timer(connection, state.config.idle_timeout) + connection.start_idle_timer(connection, state.options.idle_timeout) Continue(State(..state, buffer:, idle_timer:)) } _buffer -> @@ -219,8 +215,6 @@ fn send_response( )) Ok(to_sent(keep_alive)) } - // `file.send` owns the descriptor once it is reached, so only a head that - // never made it to the socket leaves one to hand back. encoder.RemainderFile(data) -> case transport.send(transport, socket, head) { Error(reason) -> { @@ -245,7 +239,6 @@ fn send_response( )) case rescue.handler(fn() { stream_handler(writer) }) { - // The head is already on the wire so no other answer can be given. Error(details) -> { logging.log( logging.Error, @@ -258,9 +251,6 @@ fn send_response( case drained.stream { option.Some(http1.StreamFinished(keep_alive:)) -> Ok(to_sent(keep_alive)) - // A handler that returns without finishing left the body - // unterminated, so close it out here and drop a connection we can - // no longer reuse. option.None -> { let _ = encoder.end_stream(transport, socket, framing) Ok(SentClose) @@ -269,8 +259,6 @@ fn send_response( } } } - // Once the handshake is written the connection has stopped being HTTP, so - // it never goes back to the request loop however the socket ends. encoder.RemainderWebsocket(context:, handler: websocket_handler) -> { use Nil <- result.try(transport.send(transport, socket, head)) @@ -294,8 +282,6 @@ fn send_response( let _ = encoder.end_stream(transport, socket, framing) - // The stream reports whether it left the socket at a point another - // request could start from. let drained = drain_messages(self) let stream_keep_alive = case drained.stream { option.Some(http1.StreamFinished(keep_alive:)) -> keep_alive @@ -320,9 +306,6 @@ fn to_sent(keep_alive: http1.KeepAlive) -> Sent { } } -/// What the request body left behind: the bytes after it, which begin the next -/// pipelined request, and whether it was consumed cleanly enough to reuse the -/// connection at all. type ResolvedBody { ResolvedBody(leftover: BitArray, keep_alive: http1.KeepAlive) } @@ -330,7 +313,7 @@ type ResolvedBody { fn resolve_body( conn: http1.Connection, drained: option.Option(http1.BodySignal), - config: http1.Config, + options: http1.Options, ) -> ResolvedBody { case drained { option.Some(http1.BodyDrained(leftover)) -> @@ -339,8 +322,8 @@ fn resolve_body( ResolvedBody(<<>>, http1.CloseAfterResponse) option.Some(http1.BodyProgress(buffer:, read:, chunk_remaining:)) -> http1.Connection(..conn, buffer:, read:, chunk_remaining:) - |> drain_remaining(config) - option.None -> drain_remaining(conn, config) + |> drain_remaining(options) + option.None -> drain_remaining(conn, options) } } @@ -370,19 +353,20 @@ fn do_drain_messages( fn drain_remaining( conn: http1.Connection, - config: http1.Config, + options: http1.Options, ) -> ResolvedBody { let http1.Connection(read:, ..) = conn - do_drain_remaining(conn, read + config.auto_drain_limit, config) + do_drain_remaining(conn, read + options.auto_drain_limit, options) } fn do_drain_remaining( conn: http1.Connection, limit: Int, - config: http1.Config, + options: http1.Options, ) -> ResolvedBody { - case body.pull_chunk(conn, config.auto_drain_chunk_bytes, limit) { - Ok(body.PulledChunk(_data, next)) -> do_drain_remaining(next, limit, config) + case body.pull_chunk(conn, options.auto_drain_chunk_bytes, limit) { + Ok(body.PulledChunk(_data, next)) -> + do_drain_remaining(next, limit, options) Ok(body.PulledDone(_trailers, leftover)) -> ResolvedBody(leftover, http1.KeepAlive) Error(_reason) -> ResolvedBody(<<>>, http1.CloseAfterResponse) diff --git a/src/ewe/internal/http1/body.gleam b/src/ewe/internal/http1/body.gleam index 7d605a9..e9dbfc1 100644 --- a/src/ewe/internal/http1/body.gleam +++ b/src/ewe/internal/http1/body.gleam @@ -25,8 +25,6 @@ pub fn read_body( send_body_signal(self, http1.BodyDrained(leftover:)) Ok(#(body, trailers)) } - // An oversized fixed body was rejected before reading anything, so the - // outer loop can still drain it and reuse the connection. http1.Fixed(_length), Error(BodyTooLarge) -> Error(BodyTooLarge) _framing, Error(error) -> { send_body_signal(self, http1.BodyAbandoned) @@ -141,7 +139,7 @@ fn pull_fixed_chunk( socket, buffer, want, - conn.config.body_read_timeout, + conn.options.body_read_timeout, )) let conn = http1.Connection(..conn, buffer: leftover, read: read + want) Ok(PulledChunk(data, conn)) @@ -180,7 +178,7 @@ fn pull_chunked_chunk( chunk_remaining: Int, max_chunk_bytes: Int, ) -> Result(Pulled, parser.ParseError) { - let http1.Connection(transport:, socket:, buffer:, config:, ..) = conn + let http1.Connection(transport:, socket:, buffer:, options:, ..) = conn case chunk_remaining { 0 -> { @@ -189,8 +187,8 @@ fn pull_chunked_chunk( transport, socket, buffer, - config.body_read_timeout, - parse_chunk_line(_, config), + options.body_read_timeout, + parse_chunk_line(_, options), ), ) @@ -201,14 +199,14 @@ fn pull_chunked_chunk( transport, socket, buffer, - config.body_read_timeout, + options.body_read_timeout, ) parser.parse_headers( buffer, [], 0, parser.initial_header_state(), - config, + options, ) }) @@ -243,7 +241,7 @@ fn take_chunk_slice( transport, socket, buffer, - conn.config.body_read_timeout, + conn.options.body_read_timeout, ) take_chunk_prefix(buffer, want, slice) }) @@ -285,11 +283,11 @@ fn pull_until( fn parse_chunk_line( buffer: BitArray, - config: http1.Config, + options: http1.Options, ) -> parser.Step(#(Int, BitArray)) { use #(line, remaining) <- parser.try_step(parser.extract_line( buffer, - config.max_chunk_size_line, + options.max_chunk_size_line, parser.ChunkSizeLineTooLong, parser.BadChunkSize, )) @@ -317,8 +315,6 @@ fn parse_hex_digits(bits: BitArray, acc: Int, any: Bool) -> Result(Int, Nil) { } } -/// Whether a slice reaches the end of the current chunk, and so must be -/// followed by the chunk's trailing CRLF. type Slice { LastSlice PartialSlice diff --git a/src/ewe/internal/http1/connection.gleam b/src/ewe/internal/http1/connection.gleam index 15a5df8..85ff6b4 100644 --- a/src/ewe/internal/http1/connection.gleam +++ b/src/ewe/internal/http1/connection.gleam @@ -4,9 +4,8 @@ import glisten/socket import glisten/transport import websocks -/// The limits and timeouts an HTTP/1 connection is held to. -pub type Config { - Config( +pub type Options { + Options( max_request_line: Int, max_header_line: Int, max_headers: Int, @@ -18,8 +17,8 @@ pub type Config { ) } -pub fn default_config() -> Config { - Config( +pub fn default_options() -> Options { + Options( max_request_line: 8192, max_header_line: 8192, max_headers: 100, @@ -40,26 +39,19 @@ pub type Connection { framing: Framing, read: Int, chunk_remaining: Int, - config: Config, + options: Options, upgrade: option.Option(Upgrade), ) } -/// How many socket messages a stream is delivered before it has to ask for -/// more. Bounded so a peer that keeps sending cannot grow the mailbox faster -/// than the loop drains it. pub const active_count = 100 -/// How the request body declares its length. pub type Framing { Fixed(length: Int) Chunked NoBody } -/// What a request asked to become instead of HTTP/1, gathered as the headers -/// go past rather than looked up again afterwards. Whether the fields amount to -/// a handshake is for whoever answers it to decide. pub type Upgrade { WebsocketUpgrade( key: option.Option(String), @@ -69,8 +61,6 @@ pub type Upgrade { OtherUpgrade(name: String) } -/// Handlers run inside the connection process, so they report what they did to -/// the request body and the response stream by messaging it. pub type Signal { BodySignal(BodySignal) StreamSignal(StreamSignal) @@ -96,13 +86,11 @@ pub type ResponseWriter { ) } -/// Whether the connection survives the response, or is closed once it is done. pub type KeepAlive { KeepAlive CloseAfterResponse } -/// The connection is only reusable when every party to the exchange agrees. pub fn and_keep_alive(left: KeepAlive, right: KeepAlive) -> KeepAlive { case left { KeepAlive -> right @@ -110,8 +98,6 @@ pub fn and_keep_alive(left: KeepAlive, right: KeepAlive) -> KeepAlive { } } -/// How a streamed response body delimits itself: `chunked` transfer encoding on -/// HTTP/1.1, or by closing the connection on HTTP/1.0. pub type StreamFraming { ChunkedStream CloseDelimitedStream @@ -126,8 +112,6 @@ pub type SseConnection { ) } -/// The connection has stopped being HTTP by this point so there is nothing to -/// keep alive and nothing to report back about reuse. pub type WebsocketConnection { WebsocketConnection( transport: transport.Transport, diff --git a/src/ewe/internal/http1/encoder.gleam b/src/ewe/internal/http1/encoder.gleam index caaa9cc..76d212c 100644 --- a/src/ewe/internal/http1/encoder.gleam +++ b/src/ewe/internal/http1/encoder.gleam @@ -24,7 +24,6 @@ type EncodeState { EncodeState(tree: bytes_tree.BytesTree, keep_alive: http1.KeepAlive) } -/// Whatever still has to reach the socket once the head has been written. pub type Remainder { NoRemainder RemainderInline(bytes_tree.BytesTree) @@ -72,15 +71,12 @@ pub fn encode_response( body, False -> encode_body(state, status, keep_alive, version, body) } - // A HEAD response keeps the framing headers it would have had minus the body. Ok(case method { http.Head -> drop_body(encoded) _method -> encoded }) } -/// These statuses are defined as carrying no body so they get neither one nor -/// a framing header for the client to wait on. fn is_bodyless(status: Int) -> Bool { status == 204 || status == 304 || { status >= 100 && status < 200 } } @@ -95,9 +91,6 @@ fn bodyless( Encoded(build_head(state, status, keep_alive, <<>>), keep_alive, NoRemainder) } -/// The handshake's own `connection` and `upgrade` headers are the whole point -/// of the response so no framing header is written and the head is closed off -/// without one. fn switching_protocols( state: EncodeState, status: Int, @@ -151,7 +144,6 @@ fn encode_body( } } -/// Anything the body was holding is let go here. fn drop_body(encoded: Encoded) -> Encoded { case encoded.remainder { RemainderFile(data) -> file.release(data) @@ -197,7 +189,6 @@ fn encode_stream( keep_alive, RemainderStream(handler:, framing: http1.ChunkedStream), ) - // HTTP/1.0 has no chunked encoding, so the close delimits the body instead. parser.Http10 -> close_delimited( state, @@ -207,16 +198,6 @@ fn encode_stream( } } -/// On HTTP/1.1 the stream is framed as chunked, which proxies handle far better -/// than one delimited only by the close, and which leaves the socket sitting at -/// a known point afterwards. The connection is advertised as reusable on that -/// basis; whether it is handed back is settled once the stream ends. HTTP/1.0 -/// has no chunked encoding, so there the close is the framing and the -/// connection cannot survive it. -/// -/// The content type is fixed by the format and the no-cache is what keeps -/// intermediaries from buffering the stream, so both are written from constants -/// here rather than built into the handler's header list. fn encode_sse( state: EncodeState, status: Int, @@ -274,7 +255,6 @@ pub type ResponseWriter = const last_chunk = <<"0\r\n\r\n":utf8>> -/// Wraps one piece of a streamed body in whatever delimits it on the wire. pub fn frame( chunk: bytes_tree.BytesTree, framing: http1.StreamFraming, @@ -303,7 +283,6 @@ pub fn finish_chunk( writer: ResponseWriter, chunk: BitArray, ) -> Result(Nil, socket.SocketReason) { - // The terminator rides along with the last chunk to save a write. let bytes = case writer.framing { http1.ChunkedStream -> bytes_tree.append( @@ -342,9 +321,6 @@ pub fn end_stream( } } -/// The stream is over either way so the connection process is told so even -/// when the last write never landed. One that ended on a failed write leaves -/// nothing to hand back. fn finished( writer: ResponseWriter, sent: Result(Nil, socket.SocketReason), @@ -361,8 +337,6 @@ fn finished( sent } -/// Which headers the encoder writes itself for a body, and so drops from the -/// handler's list rather than emitting twice. type Reserved { Framing FramingAndSse @@ -388,7 +362,6 @@ fn encode_headers( let initial = EncodeState(bytes_tree.new(), http1.KeepAlive) use state, #(name, value) <- list.try_fold(headers, initial) - // TODO: just trust the handler? case name, reserved { "date", _reserved -> Ok(state) _name, Handshake -> append_header(state, name, value) @@ -514,8 +487,6 @@ fn status_line(status: Int) -> BitArray { } } -/// A bare response for a request that never reached a handler, sent on a -/// connection that is closed straight after. pub fn error_response(status: Int) -> bytes_tree.BytesTree { EncodeState(bytes_tree.new(), http1.CloseAfterResponse) |> build_head(status, http1.CloseAfterResponse, << diff --git a/src/ewe/internal/http1/parser.gleam b/src/ewe/internal/http1/parser.gleam index 36ca208..7f1511e 100644 --- a/src/ewe/internal/http1/parser.gleam +++ b/src/ewe/internal/http1/parser.gleam @@ -78,8 +78,6 @@ pub fn error_to_string(error: ParseError) -> String { } } -/// The status a rejected request is answered with, so the client is told why -/// rather than left to work it out from a closed socket. pub fn error_to_status(error: ParseError) -> Int { case error { RequestLineTooLong -> 414 @@ -111,12 +109,12 @@ pub type Parsed { pub fn parse( buffer: BitArray, - config: http1.Config, + options: http1.Options, ) -> Result(Parsed, ParseError) { let step = { use #(method, target, version, remaining) <- try_step(parse_request_line( buffer, - config, + options, )) use #(headers, state, remaining) <- try_step(parse_headers( @@ -124,7 +122,7 @@ pub fn parse( [], 0, initial_header_state(), - config, + options, )) use #(host, port, path, query) <- try_step(resolve_target( @@ -166,11 +164,11 @@ pub fn try_step(step: Step(a), next: fn(a) -> Step(b)) -> Step(b) { fn parse_request_line( buffer: BitArray, - config: http1.Config, + options: http1.Options, ) -> Step(#(http.Method, BitArray, Version, BitArray)) { use #(line, remaining) <- try_step(extract_line( buffer, - config.max_request_line, + options.max_request_line, RequestLineTooLong, BadRequestLine, )) @@ -391,17 +389,12 @@ fn parse_decimal_digits(bits: BitArray, acc: Int) -> Result(Int, Nil) { } } -/// What the request's `Connection` header asked for, before the version's -/// default is applied. pub type ConnectionIntent { RequestedKeepAlive RequestedClose NothingRequested } -/// The final transfer coding the request declared. Only `chunked` delimits a -/// body, and it is the only coding decoded here, so anything else leaves a -/// length that cannot be determined. pub type TransferEncoding { NoTransferEncoding ChunkedFinal @@ -442,8 +435,6 @@ fn resolve_metadata(state: HeaderState, version: Version) -> Step(Metadata) { option.Some(_length), ChunkedFinal -> ParseError(AmbiguousFraming) option.Some(length), NoTransferEncoding -> StepDone(complete_metadata(state, version, http1.Fixed(length))) - // HTTP/1.0 has no chunked coding, so a body claiming it has no framing at - // all and cannot be told apart from the next request. option.None, ChunkedFinal -> case version { Http11 -> StepDone(complete_metadata(state, version, http1.Chunked)) @@ -468,8 +459,6 @@ fn complete_metadata( Metadata(framing:, keep_alive:, upgrade: resolve_upgrade(state)) } -/// An `Upgrade` header only asks for anything if the `Connection` header named -/// it, so one without the other is nothing. fn resolve_upgrade(state: HeaderState) -> option.Option(http1.Upgrade) { case state.connection_upgrade, state.upgrade { True, option.Some("websocket") -> @@ -488,21 +477,21 @@ pub fn parse_headers( acc: List(#(String, String)), count: Int, state: HeaderState, - config: http1.Config, + options: http1.Options, ) -> Step(#(List(#(String, String)), HeaderState, BitArray)) { use #(line, remaining) <- try_step(extract_line( buffer, - config.max_header_line, + options.max_header_line, HeaderLineTooLong, BadHeader, )) case line { <<>> -> StepDone(#(list.reverse(acc), state, remaining)) - _line if count >= config.max_headers -> ParseError(TooManyHeaders) + _line if count >= options.max_headers -> ParseError(TooManyHeaders) _line -> { use #(header, state) <- try_step(parse_header_line(line, state)) - parse_headers(remaining, [header, ..acc], count + 1, state, config) + parse_headers(remaining, [header, ..acc], count + 1, state, options) } } } @@ -550,8 +539,6 @@ fn classify( Error(Nil) -> ParseError(BadContentLength) } } - // Repeated headers concatenate into one coding list, so the last one seen - // carries the final coding. "transfer-encoding" -> { let transfer_encoding = case value |> lowercase_ascii |> tokens { [<<"chunked":utf8>>] -> ChunkedFinal @@ -577,8 +564,6 @@ fn classify( let lowered = value |> lowercase_ascii |> unsafe_to_string StepDone(HeaderState(..state, upgrade: option.Some(lowered))) } - // The key is base64 and is echoed back as sent, so unlike the rest it is - // kept with its case. "sec-websocket-key" -> StepDone( HeaderState( @@ -626,8 +611,6 @@ pub fn extract_line( False -> More } Ok(0) -> ParseError(malformed) - // The check above only bounds what is buffered while the line is still - // arriving, so a line that turns up whole in one packet is measured here. Ok(position) -> case position - 1 > max_len { True -> ParseError(too_long) @@ -667,8 +650,6 @@ fn trim_trailing_ows(bits: BitArray) -> BitArray { } } -/// The comma separated list a header value carries, one trimmed token per -/// element. pub fn tokens(value: BitArray) -> List(BitArray) { split_comma(value) |> list.map(trim_ows) } diff --git a/src/ewe/internal/http1/sse.gleam b/src/ewe/internal/http1/sse.gleam index 92e73f6..6c0e6db 100644 --- a/src/ewe/internal/http1/sse.gleam +++ b/src/ewe/internal/http1/sse.gleam @@ -11,8 +11,6 @@ import glisten/socket/options import glisten/transport import logging -/// Runs a Server-Sent Events stream, reporting through `conn.self` whether the -/// connection can carry another request afterwards. pub fn run( conn: http1.SseConnection, on_init: fn(process.Subject(user_message)) -> user_state, @@ -31,8 +29,6 @@ pub fn run( } } -/// Every way a stream ends runs the handler's `on_close`, reports whether the -/// connection survived it, and answers with the outcome. fn ended( conn: http1.SseConnection, handle: connection.SseConnection, @@ -41,15 +37,13 @@ fn ended( keep_alive: http1.KeepAlive, outcome: connection.Outcome, ) -> connection.Outcome { - // A bug in `on_close` is still a bug but it must not take the connection - // down on the way out of a stream that has already ended. - // TODO: log for a user? - let _crashed = rescue.handler(fn() { on_close(handle, state) }) + rescue.logged("server-sent events close handler", fn() { + on_close(handle, state) + }) finished(conn, keep_alive) outcome } -/// A stream the client hung up on or one the handler could no longer write to. fn dropped( conn: http1.SseConnection, handle: connection.SseConnection, @@ -66,9 +60,6 @@ fn dropped( ) } -/// A handler that crashed cannot be asked what to do next so the stream is -/// ended for it and the connection given up rather than the crash taking the -/// whole process with it. fn crashed( conn: http1.SseConnection, handle: connection.SseConnection, @@ -85,7 +76,6 @@ fn crashed( |> ended(conn, handle, state, on_close, http1.CloseAfterResponse, _) } -/// The socket gave out, which ends the stream whatever it was doing. fn socket_failed( conn: http1.SseConnection, handle: connection.SseConnection, @@ -98,8 +88,6 @@ fn socket_failed( |> ended(conn, handle, state, on_close, http1.CloseAfterResponse, _) } -/// Whether anything happened during the stream that rules out handing the -/// connection back for another request. type Reuse { Clean Spoiled @@ -116,8 +104,6 @@ fn loop( on_close: fn(connection.SseConnection, user_state) -> Nil, ) -> connection.Outcome { case process.selector_receive_forever(selector) { - // Whatever the client sent has been taken off the socket and cannot be put - // back, so the connection is no longer safe to reuse. ClientData -> loop(conn, handle, selector, state, Spoiled, step, on_close) Exhausted -> case activate(conn) { @@ -146,8 +132,6 @@ fn loop( } } -/// Only a stream the handler ended itself, on a connection nothing else has -/// touched, leaves the socket sitting exactly at the end of the response. fn keep_alive(outcome: connection.Outcome, reuse: Reuse) -> http1.KeepAlive { case outcome, reuse { connection.Stopped, Clean -> http1.KeepAlive @@ -171,15 +155,12 @@ pub fn send( |> transport.send(conn.transport, conn.socket, _) } -/// glisten rearms `{active, once}` only once its loop callback returns and an -/// SSE stream does not return until it is over. fn activate(conn: http1.SseConnection) -> Result(Nil, socket.SocketReason) { transport.set_opts(conn.transport, conn.socket, [ options.ActiveMode(options.Count(http1.active_count)), ]) } -/// What woke the stream while it was waiting on the handler's subject. type Received(user_message) { Message(user_message) Disconnected @@ -188,8 +169,6 @@ type Received(user_message) { Exhausted } -/// glisten handles socket messages in its own loop, which is blocked for the -/// duration of the stream, so the stream has to match them itself. fn selector( subject: process.Subject(user_message), ) -> process.Selector(Received(user_message)) { @@ -219,7 +198,6 @@ fn failed(record: dynamic.Dynamic) -> Received(user_message) { |> Failed } -/// Clients are not expected to send anything once the stream is open. fn client_data(_record: dynamic.Dynamic) -> Received(user_message) { ClientData } diff --git a/src/ewe/internal/http1/websocket.gleam b/src/ewe/internal/http1/websocket.gleam index ab771ea..b27df1f 100644 --- a/src/ewe/internal/http1/websocket.gleam +++ b/src/ewe/internal/http1/websocket.gleam @@ -33,8 +33,6 @@ pub fn handshake_error_to_string(error: HandshakeError) -> String { } } -/// What the handshake settled on. The key the client checks the reply against -/// and the compression it asked for. pub type Handshake { Handshake( accept: String, @@ -42,8 +40,6 @@ pub type Handshake { ) } -/// Everything this needs was picked up while the headers were parsed so the -/// request is not walked again here. pub fn handshake( method: http.Method, conn: http1.Connection, @@ -99,8 +95,6 @@ pub fn run( } } -/// Every way a socket ends runs the handler's `on_close` and then frees the -/// compression resources the context holds. fn ended( conn: http1.WebsocketConnection, state: user_state, @@ -108,10 +102,7 @@ fn ended( outcome: connection.Outcome, ) -> connection.Outcome { let handle = connection.Http1Websocket(conn) - // A bug in `on_close` is still a bug but it must not take the connection - // down on the way out of a socket that has already ended. - // TODO: log for the user? - let _crashed = rescue.handler(fn() { on_close(handle, state) }) + rescue.logged("websocket close handler", fn() { on_close(handle, state) }) websocks.close_context(conn.context) outcome @@ -125,8 +116,6 @@ fn stopped( ended(conn, state, on_close, connection.Stopped) } -/// A handler that crashed cannot be asked what to do next so the socket is -/// ended for it rather than the crash taking the whole process with it. fn crashed( conn: http1.WebsocketConnection, state: user_state, @@ -146,7 +135,6 @@ fn crashed( ) } -/// The socket gave out which ends the connection whatever it was doing. fn socket_failed( conn: http1.WebsocketConnection, state: user_state, @@ -158,7 +146,6 @@ fn socket_failed( |> ended(conn, state, on_close, _) } -/// Where the loop carries on once the handler has seen a message. type Resume { DrainBuffer AwaitSocket @@ -188,15 +175,12 @@ fn loop( websocks.push_data(conn.context, data) |> with_context(conn, _) |> drain(selector, state, step, on_close) - // Nothing reached the socket, so there is nothing new to decode. Received(message) -> websocket.UserMessage(message) |> deliver(conn, selector, state, step, on_close, AwaitSocket, _) } } -/// One read can carry several frames so the buffer is drained before the loop -/// waits on the socket again. fn drain( conn: http1.WebsocketConnection, selector: process.Selector(Received(user_message)), @@ -219,8 +203,6 @@ fn drain( let conn = with_context(conn, context) case frame { - // Answered here rather than handed on since a peer's keepalive is not - // the handler's business. websocks.Control(websocks.Ping(payload)) -> case write( @@ -233,8 +215,6 @@ fn drain( } websocks.Control(websocks.Pong(_payload)) -> drain(conn, selector, state, step, on_close) - // The peer started the closing handshake so it is echoed back and the - // socket is done. websocks.Control(websocks.Close(reason)) -> close(conn, reason) |> resolve(conn, state, on_close, _) websocks.Text(payload) -> @@ -243,7 +223,6 @@ fn drain( websocks.Binary(payload) -> websocket.BinaryFrame(payload) |> deliver(conn, selector, state, step, on_close, DrainBuffer, _) - // Fragments are reassembled by the decoder so one never surfaces. websocks.Continuation(_payload) -> drain(conn, selector, state, step, on_close) } @@ -283,8 +262,6 @@ fn deliver( } } -/// A close the server sends ends the socket either way, so only whether the -/// frame reached the peer decides how it is reported. fn resolve( conn: http1.WebsocketConnection, state: user_state, @@ -351,8 +328,6 @@ fn write( |> transport.send(conn.transport, conn.socket, _) } -/// glisten rearms the socket only once its loop callback returns and a socket -/// does not return until it is over, so the frames have to be asked for here. fn activate( conn: http1.WebsocketConnection, ) -> Result(Nil, socket.SocketReason) { diff --git a/src/ewe/internal/http1_ffi.erl b/src/ewe/internal/http1_ffi.erl index 780e9ec..3c70c1d 100644 --- a/src/ewe/internal/http1_ffi.erl +++ b/src/ewe/internal/http1_ffi.erl @@ -47,12 +47,9 @@ find(Bin, Key) -> {Pos, _Len} -> {ok, Pos} end. -%% Splits on every comma in one call. split_comma(Bin) -> binary:split(Bin, persistent_term:get({?MODULE, comma}), [global]). -%% Scans for uppercase natively so an already lowercase binary is returned -%% untouched, and rewrites in one pass otherwise. lowercase_ascii(Bin) -> case binary:match(Bin, persistent_term:get({?MODULE, upper})) of nomatch -> Bin; @@ -62,11 +59,9 @@ lowercase_ascii(Bin) -> lower(Byte) when Byte >= $A, Byte =< $Z -> Byte + 32; lower(Byte) -> Byte. -%% Reason carried by a `{tcp_error, Socket, Reason}` message. socket_error_reason({_Tag, _Socket, Reason}) -> Reason. -%% Validates UTF-8 and returns the bytes unchanged. bit_array_to_string(Bin) -> case ewe_ffi:is_valid_utf8(Bin) of true -> {ok, Bin}; diff --git a/src/ewe/internal/http2.gleam b/src/ewe/internal/http2.gleam index 83bb095..804ee84 100644 --- a/src/ewe/internal/http2.gleam +++ b/src/ewe/internal/http2.gleam @@ -61,7 +61,6 @@ fn apply_setting( } } -/// What the connection does with the socket once a message has been handled. pub type Next { Continue(State) Close @@ -119,8 +118,6 @@ pub type Stream { ) } -// Opaque handle to a `binary:compile_pattern/1` result. Compiled once per -// node and fetched once per connection. type Pattern pub opaque type HeaderPatterns { @@ -176,10 +173,11 @@ pub type State { draining: Bool, drain_subject: process.Subject(connection.Message), drain_timer: Option(process.Timer), - config: http2.Config, + options: http2.Options, settings_frame: bytes_tree.BytesTree, patterns: HeaderPatterns, peer: Result(socket.SockName, Nil), + parent: Result(process.Pid, Nil), ) } @@ -187,32 +185,30 @@ const default_send_window = 65_535 const max_window_size = 2_147_483_647 -/// How many socket messages arrive before the connection has to ask for more. -/// Caps how fast a peer can grow the mailbox. pub const socket_active_batch_size = 32 -fn build_settings_frame(config: http2.Config) -> bytes_tree.BytesTree { - let params = case config.header_table_size { +fn build_settings_frame(options: http2.Options) -> bytes_tree.BytesTree { + let params = case options.header_table_size { 4096 -> [] - _size -> [frame.HeaderTableSize(config.header_table_size)] + _size -> [frame.HeaderTableSize(options.header_table_size)] } - let params = case config.initial_window_size { + let params = case options.initial_window_size { 65_535 -> params - _size -> [frame.InitialWindowSize(config.initial_window_size), ..params] + _size -> [frame.InitialWindowSize(options.initial_window_size), ..params] } - let params = case config.max_frame_size { + let params = case options.max_frame_size { 16_384 -> params - _size -> [frame.MaxFrameSize(config.max_frame_size), ..params] + _size -> [frame.MaxFrameSize(options.max_frame_size), ..params] } - let params = case config.max_concurrent_streams { + let params = case options.max_concurrent_streams { Some(value) -> [frame.MaxConcurrentStreams(value), ..params] None -> params } - let params = case config.max_header_list_size { + let params = case options.max_header_list_size { Some(value) -> [frame.MaxHeaderListSize(value), ..params] None -> params } @@ -222,8 +218,6 @@ fn build_settings_frame(config: http2.Config) -> bytes_tree.BytesTree { |> bytes_tree.from_bit_array } -/// Wired to glisten's close callback so a connection dropped underneath its -/// streams still releases what they were holding. pub fn kill_live_workers(state: State) -> Nil { use _stream_id, entry <- dict.each(state.streams) @@ -238,21 +232,20 @@ pub fn kill_live_workers(state: State) -> Nil { } } -/// Builds the state for a connection whose preface is already read. The caller -/// makes the subjects because the caller is what selects on them. pub fn init( handler: fn(Request(connection.Connection)) -> Response(connection.Body), - config: http2.Config, + options: http2.Options, self: process.Subject(connection.Message), reply_subject: process.Subject(http2.Reply(connection.Body)), peer: Result(socket.SockName, Nil), + parent: Result(process.Pid, Nil), ) -> State { - let table = alpacki.new_dynamic(config.header_table_size) + let table = alpacki.new_dynamic(options.header_table_size) let timer = process.send_after( self, - config.handshake_timeout_ms, + options.handshake_timeout_ms, connection.Http2Handshake, ) @@ -276,10 +269,11 @@ pub fn init( draining: False, drain_subject: self, drain_timer: None, - config:, - settings_frame: build_settings_frame(config), + options:, + settings_frame: build_settings_frame(options), patterns: header_patterns(), peer:, + parent:, ) } @@ -299,8 +293,6 @@ pub fn handle_message( glisten.User(connection.Http2Drain) -> stop_connection(state) glisten.User(connection.Http2StreamClose(pid)) -> finish_or_continue(handle_stream_close_timeout(state, pid)) - // HTTP/1's idle timer is cancelled before a connection becomes HTTP/2 - // which has its own handshake and drain deadlines instead. glisten.User(connection.Timeout) -> Continue(state) } } @@ -335,7 +327,7 @@ fn process_frames( state: State, connection: glisten.Connection(connection.Message), ) -> Next { - case frame.decode(state.buffer, state.config.max_frame_size) { + case frame.decode(state.buffer, state.options.max_frame_size) { Ok(#(frame, remaining)) -> case handle_frame(State(..state, buffer: remaining), frame, connection) { Proceed(state) -> process_frames(state, connection) @@ -380,8 +372,6 @@ fn reset_and_remove_stream( case entry.status { Computing(pid) -> { process.send_abnormal_exit(pid, "stream_reset") - // Backstop for a long-lived handler that doesn't react to the exit - // signal promptly process.send_after( state.drain_subject, stream_close_grace_ms, @@ -406,7 +396,7 @@ fn handle_stream_close_timeout(state: State, pid: process.Pid) -> State { Error(Nil) -> state Ok(_stream_id) -> { process.kill(pid) - clear_stream_pid(state, pid) + state } } } @@ -492,8 +482,6 @@ fn handle_ping( } } -/// A client cancelling a stream. Enough of them in one window means Rapid Reset -/// (CVE-2023-44487), not ordinary cancelling, and the connection goes. @internal pub fn handle_client_reset(state: State, stream_id: Int) -> FrameResult { let #(state, tripped) = record_reset(state) @@ -517,7 +505,7 @@ fn record_reset(state: State) -> #(State, Bool) { let new_window = state.reset_count == 0 - || now - state.reset_window_start > state.config.rapid_reset_window_ms + || now - state.reset_window_start > state.options.rapid_reset_window_ms let #(reset_window_start, reset_count) = case new_window { True -> #(now, 1) @@ -526,7 +514,7 @@ fn record_reset(state: State) -> #(State, Bool) { #( State(..state, reset_window_start:, reset_count:), - reset_count > state.config.rapid_reset_threshold, + reset_count > state.options.rapid_reset_threshold, ) } @@ -601,7 +589,7 @@ fn start_header_assembly( ) let oversized = - bit_array.byte_size(payload) > state.config.max_header_block_bytes + bit_array.byte_size(payload) > state.options.max_header_block_bytes case oversized, end_headers, trailers { True, _end_headers, _trailers -> Terminate(Some(frame.EnhanceYourCalm)) @@ -612,7 +600,6 @@ fn start_header_assembly( } } -// Clients only ever open odd-numbered streams and always in increasing order. fn is_new_client_stream_id(state: State, stream_id: Int) -> Bool { stream_id % 2 == 1 && stream_id > state.highest_client_stream_id_seen } @@ -640,10 +627,6 @@ pub fn handle_data( } } -// The client sent this DATA frame before finding out the stream is closed on -// our end so it still counts against the connection window. Flow control has -// to be here even though the stream itself is already gone otherwise our view -// of the window drifts from the client's. fn reject_closed_stream_data(state: State, size: Int) -> FrameResult { case state.conn_recv_window - size < 0 { True -> Terminate(Some(frame.FlowControlError)) @@ -651,8 +634,6 @@ fn reject_closed_stream_data(state: State, size: Int) -> FrameResult { } } -// Same here, still gotta debit the connection window before we can reject the -// stream. fn reject_data_after_half_close( state: State, stream_id: Int, @@ -695,7 +676,7 @@ fn apply_data( |> deliver_data(payload, end_stream, []) let #(entry, stream_increment) = case delivered_to_reader { - True -> stream_recv_credit(entry, state.config) + True -> stream_recv_credit(entry, state.options) False -> #(entry, 0) } @@ -739,18 +720,27 @@ fn deliver_data( let request_half_closed = entry.request_half_closed || end_stream let entry = Stream(..entry, request_half_closed:) - case entry.parked_reader { - Some(reply_to) -> { - case payload, end_stream { - <<>>, _half_closed -> process.send(reply_to, http2.DoneEvent(trailers)) - _payload, True -> - process.send(reply_to, http2.LastChunkEvent(payload, trailers)) - _payload, False -> process.send(reply_to, http2.ChunkEvent(payload)) - } + case entry.parked_reader, payload, end_stream { + _reader, <<>>, False -> #(entry, False) + Some(reply_to), <<>>, True -> { + process.send(reply_to, http2.DoneEvent(trailers)) + #(Stream(..entry, parked_reader: None), True) + } + Some(reply_to), _payload, True -> { + process.send(reply_to, http2.LastChunkEvent(payload, trailers)) + #(Stream(..entry, parked_reader: None), True) + } + Some(reply_to), _payload, False -> { + process.send(reply_to, http2.ChunkEvent(payload)) #(Stream(..entry, parked_reader: None), True) } - None -> { + None, _payload, _end_stream -> { let recv_buffer = bytes_tree.append(entry.recv_buffer, payload) + let trailers = case trailers { + [] -> entry.trailers + _received -> trailers + } + #(Stream(..entry, recv_buffer:, trailers:), False) } } @@ -766,7 +756,7 @@ pub fn handle_continuation( frame.Continuation(stream_id, end_headers, payload) if stream_id == assembly.stream_id -> - case add_fragment(assembly, payload, state.config), end_headers { + case add_fragment(assembly, payload, state.options), end_headers { Error(code), _end_headers -> Terminate(Some(code)) Ok(updated), False -> Proceed(State(..state, header_assembly: Some(updated))) @@ -783,14 +773,14 @@ pub fn handle_continuation( pub fn add_fragment( assembly: HeaderAssembly, fragment: BitArray, - config: http2.Config, + options: http2.Options, ) -> Result(HeaderAssembly, frame.ErrorCode) { let fragment_count = assembly.fragment_count + 1 let block = <> case - fragment_count > config.max_continuation_frames - || bit_array.byte_size(block) > config.max_header_block_bytes + fragment_count > options.max_continuation_frames + || bit_array.byte_size(block) > options.max_header_block_bytes { True -> Error(frame.EnhanceYourCalm) False -> Ok(HeaderAssembly(..assembly, fragment_count:, block:)) @@ -809,7 +799,7 @@ fn decode_and_validate_header_block( dynamic_table:, remaining:, )) -> { - let header_list_size_exceeded = case state.config.max_header_list_size { + let header_list_size_exceeded = case state.options.max_header_list_size { None -> False Some(limit) -> decoded_size > limit } @@ -817,7 +807,7 @@ fn decode_and_validate_header_block( case remaining != <<>> || alpacki.dynamic_max_size(dynamic_table) - > state.config.header_table_size, + > state.options.header_table_size, header_list_size_exceeded { True, _header_list_size_exceeded -> @@ -849,7 +839,7 @@ pub fn complete_header_block( pending: <<>>, pending_trailers: None, read: 0, - body_read_timeout: state.config.body_read_timeout, + body_read_timeout: state.options.body_read_timeout, peer: state.peer, )) @@ -949,7 +939,7 @@ fn spawn_stream( end_stream: Bool, content_length: Option(Int), ) -> FrameResult { - let concurrent_streams_exceeded = case state.config.max_concurrent_streams { + let concurrent_streams_exceeded = case state.options.max_concurrent_streams { None -> False Some(limit) -> dict.size(state.streams) >= limit } @@ -978,7 +968,7 @@ fn track_stream( pending: PendingBytes(<<>>), pending_end_stream: True, write_ack: None, - recv_window: state.config.initial_window_size, + recv_window: state.options.initial_window_size, recv_buffer: bytes_tree.new(), request_half_closed: end_stream, parked_reader: None, @@ -1320,8 +1310,6 @@ fn handle_client_settings( } } -/// SETTINGS_INITIAL_WINDOW_SIZE applies retroactively, so every stream already -/// open gets the delta too. @internal pub fn adjust_stream_windows(state: State, delta: Int) -> State { case delta { @@ -1354,7 +1342,11 @@ fn terminate( ) -> Next { case code { Some(error_code) -> { - let _sent = send_frame(connection, frame.Goaway(0, 0, error_code, <<>>)) + let _sent = + send_frame( + connection, + frame.Goaway(0, state.highest_client_stream_id_seen, error_code, <<>>), + ) Nil } None -> Nil @@ -1382,7 +1374,7 @@ fn begin_drain( let timer = process.send_after( state.drain_subject, - state.config.drain_timeout_ms, + state.options.drain_timeout_ms, connection.Http2Drain, ) @@ -1471,7 +1463,7 @@ fn handle_read_body( process.send(reply_to, http2.ChunkEvent(buffered)) let drained_entry = Stream(..entry, recv_buffer: bytes_tree.new()) let #(drained_entry, stream_increment) = - stream_recv_credit(drained_entry, state.config) + stream_recv_credit(drained_entry, state.options) let streams = dict.insert(state.streams, stream_id, drained_entry) let state = State(..state, streams:) @@ -1494,19 +1486,17 @@ fn handle_read_body( } } -/// Tops a stream's receive window back up once it's fallen far enough to be -/// worth a frame. Beats trickling out an update per chunk read. @internal pub fn stream_recv_credit( entry: Stream, - config: http2.Config, + options: http2.Options, ) -> #(Stream, Int) { - case entry.recv_window <= config.recv_window_low_water_mark { + case entry.recv_window <= options.recv_window_low_water_mark { True -> { - let increment = config.recv_window_high_water_mark - entry.recv_window + let increment = options.recv_window_high_water_mark - entry.recv_window #( - Stream(..entry, recv_window: config.recv_window_high_water_mark), + Stream(..entry, recv_window: options.recv_window_high_water_mark), increment, ) } @@ -1514,18 +1504,17 @@ pub fn stream_recv_credit( } } -/// The same for the connection's own window which every stream draws from. @internal pub fn conn_recv_credit(state: State) -> #(State, Int) { - case state.conn_recv_window <= state.config.recv_window_low_water_mark { + case state.conn_recv_window <= state.options.recv_window_low_water_mark { True -> { let increment = - state.config.recv_window_high_water_mark - state.conn_recv_window + state.options.recv_window_high_water_mark - state.conn_recv_window #( State( ..state, - conn_recv_window: state.config.recv_window_high_water_mark, + conn_recv_window: state.options.recv_window_high_water_mark, ), increment, ) @@ -1591,8 +1580,6 @@ fn response_body_size(body: connection.Body) -> Int { connection.Empty -> 0 connection.File(connection.OpenFile(length:, ..)) | connection.File(connection.PendingFile(length:, ..)) -> length - // A streamed body never reaches here as its stream process writes the - // response itself connection.Streaming(_metadata) | connection.Sse(_metadata) | connection.Websocket(_metadata) -> @@ -1613,7 +1600,6 @@ fn open_pending( connection.Empty -> Ok(PendingBytes(<<>>)) connection.File(connection.OpenFile(handle:, offset:, length:)) -> Ok(PendingFile(handle, offset, length)) - // Small enough to answer from memory. connection.File(connection.PendingFile(path:, offset:, length:)) if length <= file_read_threshold -> file.read_range(path, offset, length) |> result.map(PendingBytes) @@ -1633,7 +1619,7 @@ fn respond( response: Response(connection.Body), connection: glisten.Connection(connection.Message), ) -> Next { - case open_pending(response.body, state.config.file_read_threshold) { + case open_pending(response.body, state.options.file_read_threshold) { Ok(pending) -> send_response( state, @@ -1724,9 +1710,6 @@ fn send_response( } } -/// A stream sets these itself, so the handler's copies get dropped. HTTP/1 -/// reserves `transfer-encoding` too, which HTTP/2 has no use for and must -/// never send. fn drop_sse_headers( headers: List(#(String, String)), ) -> List(#(String, String)) { @@ -1854,8 +1837,6 @@ fn resolve_frame_result( } } -/// Frames a header block, splitting into CONTINUATION frames when it's longer -/// than the peer takes in one. @internal pub fn append_header_frames( acc: bytes_tree.BytesTree, @@ -2151,8 +2132,6 @@ fn send_file_chunk( } } -/// Writes what a stream has pending as far as its window and the connection's -/// allow. @internal pub fn flush_stream( state: State, @@ -2260,33 +2239,47 @@ fn handle_stream_exit( exit: process.ExitMessage, connection: glisten.Connection(connection.Message), ) -> Next { - case dict.get(state.stream_pids, exit.pid) { - // Not a stream this connection started so it is the parent asking it to - // shut down. - Error(Nil) -> begin_drain(state, connection) + case state.parent == Ok(exit.pid), exit.reason { + True, process.Abnormal(reason) -> + case http2.is_shutdown(reason) { + True -> begin_drain(state, connection) + False -> stop_connection(state) + } + True, process.Normal | True, process.Killed -> stop_connection(state) + False, _reason -> handle_child_exit(state, exit.pid, connection) + } +} + +fn handle_child_exit( + state: State, + pid: process.Pid, + connection: glisten.Connection(connection.Message), +) -> Next { + case dict.get(state.stream_pids, pid) { + Error(Nil) -> finish_or_continue(state) Ok(stream_id) -> case dict.get(state.streams, stream_id) { - // Still computing when its process ended means the handler returned - // without a response being completed. - Ok(Stream(status: Computing(pid), ..) as entry) if pid == exit.pid -> + Ok(Stream(status: Computing(stream_pid), ..) as entry) + if stream_pid == pid + -> remove_stream(state, stream_id, entry) |> reject_stream(connection, _, stream_id, frame.InternalError) - _entry -> finish_or_continue(clear_stream_pid(state, exit.pid)) + _entry -> finish_or_continue(clear_stream_pid(state, pid)) } } } @internal pub fn test_state() -> State { - let config = http2.default_config() + let options = http2.default_options() State( buffer: <<>>, handshake: Connected, peer_settings: default_peer_settings, timer: None, - hpack_decoder: alpacki.new_dynamic(config.header_table_size), - hpack_encoder: alpacki.new_dynamic(config.header_table_size), + hpack_decoder: alpacki.new_dynamic(options.header_table_size), + hpack_encoder: alpacki.new_dynamic(options.header_table_size), header_assembly: None, reply_subject: process.new_subject(), handler: fn(_request) { @@ -2302,9 +2295,10 @@ pub fn test_state() -> State { draining: False, drain_subject: process.new_subject(), drain_timer: None, - config:, - settings_frame: build_settings_frame(config), + options:, + settings_frame: build_settings_frame(options), patterns: header_patterns(), peer: Error(Nil), + parent: Error(Nil), ) } diff --git a/src/ewe/internal/http2/body.gleam b/src/ewe/internal/http2/body.gleam index a3fa2e4..714494d 100644 --- a/src/ewe/internal/http2/body.gleam +++ b/src/ewe/internal/http2/body.gleam @@ -36,7 +36,6 @@ fn read_all( } } -/// The result of one `read_body_chunk` call. pub type ReadEvent(body) { Chunk(data: BitArray, connection: http2.Connection(body)) Done(trailers: List(#(String, String))) @@ -58,7 +57,6 @@ pub fn read_body_chunk( } } -/// Leftovers past `max_chunk_bytes` wait on the connection for the next call. fn split( connection: http2.Connection(body), data: BitArray, @@ -78,7 +76,6 @@ fn next_chunk( False, _pending, _trailers -> Ok(Done([])) True, <<>>, option.Some(trailers) -> Ok(Done(trailers)) True, <<>>, option.None -> pull(connection) - // Owed from the last split. Hand these back before asking for more. True, pending, _trailers -> Ok(Chunk(pending, http2.Connection(..connection, pending: <<>>))) } diff --git a/src/ewe/internal/http2/connection.gleam b/src/ewe/internal/http2/connection.gleam index 0c129a7..d90bb8e 100644 --- a/src/ewe/internal/http2/connection.gleam +++ b/src/ewe/internal/http2/connection.gleam @@ -5,9 +5,8 @@ import gleam/http/response import gleam/option import glisten/socket -/// The limits and timeouts an HTTP/2 connection is held to. -pub type Config { - Config( +pub type Options { + Options( max_concurrent_streams: option.Option(Int), initial_window_size: Int, max_frame_size: Int, @@ -26,8 +25,8 @@ pub type Config { ) } -pub fn default_config() -> Config { - Config( +pub fn default_options() -> Options { + Options( max_concurrent_streams: option.None, initial_window_size: 2_097_152, max_frame_size: 16_384, @@ -46,30 +45,19 @@ pub fn default_config() -> Config { ) } -/// What a handler holds for one stream. The handler gets its own process, not -/// the connection's. Reading the body and writing the response are messages, -/// not socket writes. -/// -/// The body type is a parameter to break the import cycle with the module that -/// defines it. pub type Connection(body) { Connection( connection: process.Subject(Reply(body)), stream_id: Int, has_body: Bool, - /// Body bytes handed over but not yet returned to the caller. pending: BitArray, pending_trailers: option.Option(List(#(String, String))), - /// Body bytes read so far. `read_body_chunk` caps its limit against this. read: Int, body_read_timeout: Int, - /// Resolved once for the connection. A stream process has no socket to - /// ask. peer: Result(socket.SockName, Nil), ) } -/// What a stream process asks of the connection process. pub type Reply(body) { Respond(stream_id: Int, response: response.Response(body)) ReadBody(stream_id: Int, reply_to: process.Subject(BodyEvent)) @@ -88,8 +76,6 @@ pub type Reply(body) { ) } -/// Which headers the connection sets itself for a streamed body. They get -/// dropped from the handler's list so nothing goes out twice. pub type Reserved { Nothing SseHeaders @@ -105,8 +91,6 @@ pub type WriteAck { WriteAck } -/// A handle for writing a response a frame at a time. The connection tags its -/// acks with the reference so a stream waits on it directly, no selector. pub type ResponseWriter(body) { ResponseWriter( connection: process.Subject(Reply(body)), @@ -120,16 +104,12 @@ pub type SseConnection(body) { SseConnection(writer: ResponseWriter(body)) } -/// Why a stream process stopped waiting. Only a body read times out. It is the -/// one wait that hangs on the client. pub type Interrupted { StreamReset ConnectionClosed TimedOut } -/// Waits for the connection to answer. Gives up if the stream resets or the -/// connection goes. @external(erlang, "http2_ffi", "recv_or_exit") pub fn receive_reply(tag: reference.Reference) -> Result(message, Interrupted) @@ -139,8 +119,11 @@ pub fn receive_reply_within( timeout: Int, ) -> Result(message, Interrupted) -/// `unsafely_create_subject` wants the tag the messages carry and the -/// connection tags replies with a reference, so the two meet as `Dynamic`. -/// Waiting on the reference directly saves building a selector per chunk. @external(erlang, "ewe_ffi", "identity") pub fn tag(reference: reference.Reference) -> dynamic.Dynamic + +@external(erlang, "http2_ffi", "parent_pid") +pub fn parent_pid() -> Result(process.Pid, Nil) + +@external(erlang, "http2_ffi", "is_shutdown") +pub fn is_shutdown(reason: dynamic.Dynamic) -> Bool diff --git a/src/ewe/internal/http2/frame.gleam b/src/ewe/internal/http2/frame.gleam index 765f60b..e45fdd9 100644 --- a/src/ewe/internal/http2/frame.gleam +++ b/src/ewe/internal/http2/frame.gleam @@ -144,8 +144,6 @@ fn decode_payload( 0x2 -> case payload { <<_exclusive:1, stream_dependency:31, _weight:8>> -> - // Stream 0 can't be prioritized and a stream depending on itself - // is a cycle of one which both are protocol errors. case stream_id == 0, stream_dependency == stream_id { True, _self_dependent | _on_stream_zero, True -> Error(Violation(ProtocolError)) @@ -161,7 +159,6 @@ fn decode_payload( False, _payload -> Error(Violation(FrameSizeError)) } 0x4 -> - // SETTINGS ACKs never carry params. case end_stream_or_ack, payload { True, <<>> -> Ok(Settings(stream_id, True, [])) True, _payload -> Error(Violation(FrameSizeError)) @@ -215,7 +212,6 @@ fn strip_padding( case payload { <> -> { let content_length = bit_array.byte_size(remaining) - pad_length - // A client can claim more padding than bytes actually follow case content_length >= 0 { True -> case remaining { @@ -382,8 +378,6 @@ fn encode_flags(end_stream_or_ack: Bool, end_headers: Bool) -> BitArray { <<0:5, bit(end_headers):1, 0:1, bit(end_stream_or_ack):1>> } -/// Frames a DATA payload without copying it into the header, so a body already -/// held as a `BytesTree` can be written straight after this. pub fn encode_data_header( stream_id: Int, end_stream: Bool, diff --git a/src/ewe/internal/http2/sse.gleam b/src/ewe/internal/http2/sse.gleam index 009381e..4fd5331 100644 --- a/src/ewe/internal/http2/sse.gleam +++ b/src/ewe/internal/http2/sse.gleam @@ -9,8 +9,6 @@ import gleam/erlang/reference import gleam/result import logging -/// Runs a Server-Sent Events stream. The response head is already written by -/// the time this is reached so there is nothing to negotiate. pub fn run( conn: http2.SseConnection(connection.Body), on_init: fn(process.Subject(user_message)) -> user_state, @@ -35,7 +33,6 @@ fn loop( on_close: fn(connection.SseConnection, user_state) -> Nil, ) -> connection.Outcome { case http2.receive_reply(tag) { - // The stream was reset or the connection went. Error(_interrupted) -> ended(handle, state, on_close, connection.Stopped) Ok(message) -> case rescue.handler(fn() { step(handle, state, message) }) { @@ -44,8 +41,6 @@ fn loop( Ok(sse.Halt(outcome)) -> { let outcome = ended(handle, state, on_close, outcome) - // A stream the handler ended gets a clean close. An abnormal one - // resets. case outcome { connection.Stopped -> { let _sent = stream.finish_response(conn.writer) @@ -60,22 +55,18 @@ fn loop( } } -/// Every ending runs `on_close` and reports the outcome. fn ended( handle: connection.SseConnection, state: user_state, on_close: fn(connection.SseConnection, user_state) -> Nil, outcome: connection.Outcome, ) -> connection.Outcome { - // A bug in `on_close` is still a bug. It just must not kill the process on - // the way out of a stream that already ended. - // TODO: logging here? - let _crashed = rescue.handler(fn() { on_close(handle, state) }) + rescue.logged("server-sent events close handler", fn() { + on_close(handle, state) + }) outcome } -/// A crashed handler cannot say what to do next. End the stream for it and -/// reset. fn crashed( handle: connection.SseConnection, state: user_state, diff --git a/src/ewe/internal/http2/stream.gleam b/src/ewe/internal/http2/stream.gleam index 6da0a23..f2402ec 100644 --- a/src/ewe/internal/http2/stream.gleam +++ b/src/ewe/internal/http2/stream.gleam @@ -53,8 +53,6 @@ fn deliver( case response.body { connection.Streaming(connection.StreamingMetadata(handler:)) -> case begin(reply_to, stream_id, response, http2.Nothing) { - // Nothing was written and nothing can be. No stream left to produce - // a body for. Error(_interrupted) -> Nil Ok(writer) -> handler(connection.Http2Writer(writer)) } @@ -67,8 +65,6 @@ fn deliver( connection.StoppedAbnormal(reason) -> abort(reason) } } - // A WebSocket body never gets here. The handshake needs extended CONNECT - // and is refused long before a handler returns one. connection.Websocket(_metadata) -> { logging.log( logging.Error, @@ -85,8 +81,6 @@ fn deliver( } } -/// Writes the response head. A handler only gets a writer once the client has -/// something to hang the body off. fn begin( reply_to: process.Subject(http2.Reply(connection.Body)), stream_id: Int, @@ -111,8 +105,6 @@ fn begin( http2.ResponseWriter(connection: reply_to, stream_id:, ack:, ack_ref:) } -/// Sends one body chunk. Returns once the connection has it on the wire. An -/// empty chunk gets no frame. pub fn send_chunk( writer: http2.ResponseWriter(connection.Body), chunk: BitArray, @@ -136,9 +128,6 @@ pub fn finish_response( finish_chunk(writer, <<>>) } -/// Waiting on the ack is how flow control reaches the handler. The connection -/// answers once the bytes are gone, so a handler outrunning the client blocks -/// here. fn write( writer: http2.ResponseWriter(connection.Body), chunk: BitArray, diff --git a/src/ewe/internal/http2_ffi.erl b/src/ewe/internal/http2_ffi.erl index e494523..5992e8f 100644 --- a/src/ewe/internal/http2_ffi.erl +++ b/src/ewe/internal/http2_ffi.erl @@ -14,7 +14,9 @@ colon_pattern/0, exit_self/1, recv_or_exit/1, - recv_or_exit/2 + recv_or_exit/2, + parent_pid/0, + is_shutdown/1 ]). @@ -43,6 +45,15 @@ monotonic_ms() -> exit_self(Reason) -> erlang:exit(Reason). +parent_pid() -> + case erlang:process_info(self(), parent) of + {parent, Pid} when is_pid(Pid) -> {ok, Pid}; + _Other -> {error, nil} + end. + +is_shutdown(shutdown) -> true; +is_shutdown(_Reason) -> false. + recv_or_exit(Ref) -> receive {Ref, Message} -> {ok, Message}; diff --git a/src/ewe/internal/rescue.gleam b/src/ewe/internal/rescue.gleam index e27336d..69833b6 100644 --- a/src/ewe/internal/rescue.gleam +++ b/src/ewe/internal/rescue.gleam @@ -1,2 +1,15 @@ +import logging + @external(erlang, "ewe_ffi", "rescue_handler") pub fn handler(handler: fn() -> a) -> Result(a, String) + +pub fn logged(what: String, callback: fn() -> Nil) -> Nil { + case handler(callback) { + Ok(Nil) -> Nil + Error(details) -> + logging.log( + logging.Error, + "Caught a crash in the " <> what <> ": " <> details, + ) + } +} diff --git a/src/ewe/internal/sse.gleam b/src/ewe/internal/sse.gleam index ae34ed0..ef15ebd 100644 --- a/src/ewe/internal/sse.gleam +++ b/src/ewe/internal/sse.gleam @@ -14,7 +14,6 @@ pub type Event { ) } -/// What a protocol's stream loop does once the handler has seen a message. pub type Step(user_state) { Proceed(user_state) Halt(connection.Outcome) @@ -42,8 +41,6 @@ pub fn encode(event: Event) -> bytes_tree.BytesTree { const newline = <<"\n":utf8>> -/// A break in a single line field would be read as the start of the next field, -/// so it is dropped rather than allowed to forge one. fn append_field( tree: bytes_tree.BytesTree, prefix: String, @@ -58,7 +55,6 @@ fn append_field( } } -/// A break splits the value over repeated fields, which the client rejoins. fn append_lines( tree: bytes_tree.BytesTree, prefix: String, diff --git a/src/ewe/internal/websocket.gleam b/src/ewe/internal/websocket.gleam index 6aaf149..b6008d8 100644 --- a/src/ewe/internal/websocket.gleam +++ b/src/ewe/internal/websocket.gleam @@ -2,7 +2,6 @@ import ewe/internal/connection import gleam/erlang/process import gleam/option -/// What a protocol's socket loop does once the handler has seen a message. pub type Step(user_state, user_message) { Proceed( user_state: user_state, diff --git a/src/ewe/internal/websocket_ffi.erl b/src/ewe/internal/websocket_ffi.erl index 84f2eea..e8ccf3f 100644 --- a/src/ewe/internal/websocket_ffi.erl +++ b/src/ewe/internal/websocket_ffi.erl @@ -2,6 +2,5 @@ -export([socket_payload/1]). -%% Payload carried by a `{tcp, Socket, Data}` message. socket_payload({_Tag, _Socket, Data}) -> Data. diff --git a/test/ewe/internal/http1/parser_test.gleam b/test/ewe/internal/http1/parser_test.gleam index 7e68dd2..1c4be08 100644 --- a/test/ewe/internal/http1/parser_test.gleam +++ b/test/ewe/internal/http1/parser_test.gleam @@ -4,7 +4,7 @@ import gleam/http import gleam/option.{None, Some} fn parse(buffer: BitArray) -> Result(parser.Parsed, parser.ParseError) { - parser.parse(buffer, http1.default_config()) + parser.parse(buffer, http1.default_options()) } pub fn simple_get_test() { @@ -170,9 +170,9 @@ pub fn configured_max_headers_is_applied_test() { let buffer = << "GET / HTTP/1.1\r\nHost: example.com\r\nX-A: 1\r\nX-B: 2\r\n\r\n":utf8, >> - let config = http1.Config(..http1.default_config(), max_headers: 2) + let options = http1.Options(..http1.default_options(), max_headers: 2) - assert parser.parse(buffer, config) == Error(parser.TooManyHeaders) + assert parser.parse(buffer, options) == Error(parser.TooManyHeaders) let assert Ok(parser.Complete(..)) = parse(buffer) as "the same request is fine under the default limit" } @@ -181,16 +181,16 @@ pub fn configured_max_request_line_is_applied_test() { let buffer = << "GET /a/fairly/long/path HTTP/1.1\r\nHost: example.com\r\n":utf8, >> - let config = http1.Config(..http1.default_config(), max_request_line: 8) + let options = http1.Options(..http1.default_options(), max_request_line: 8) - assert parser.parse(buffer, config) == Error(parser.RequestLineTooLong) + assert parser.parse(buffer, options) == Error(parser.RequestLineTooLong) } pub fn configured_max_header_line_is_applied_test() { let buffer = <<"GET / HTTP/1.1\r\nHost: example.com\r\n":utf8>> - let config = http1.Config(..http1.default_config(), max_header_line: 4) + let options = http1.Options(..http1.default_options(), max_header_line: 4) - assert parser.parse(buffer, config) == Error(parser.HeaderLineTooLong) + assert parser.parse(buffer, options) == Error(parser.HeaderLineTooLong) } pub fn rejections_carry_a_status_test() { @@ -429,7 +429,6 @@ pub fn websocket_handshake_fields_collected_test() { assert metadata.upgrade == Some(http1.WebsocketUpgrade( - // Base64 is case sensitive, so the key alone keeps the case it arrived in. key: Some("dGhlIHNhbXBsZSBub25jZQ=="), version: Some("13"), extensions: Some("permessage-deflate; client_max_window_bits"), diff --git a/test/ewe/internal/http2/connection_test.gleam b/test/ewe/internal/http2/connection_test.gleam index a396aaa..e0f27bd 100644 --- a/test/ewe/internal/http2/connection_test.gleam +++ b/test/ewe/internal/http2/connection_test.gleam @@ -228,13 +228,13 @@ pub fn non_continuation_mid_assembly_is_protocol_error_test() { pub fn add_fragment_within_limits_test() { let assembly = connection.HeaderAssembly(1, True, 1, <<"a":utf8>>, False) let assert Ok(updated) = - connection.add_fragment(assembly, <<"b":utf8>>, http2.default_config()) + connection.add_fragment(assembly, <<"b":utf8>>, http2.default_options()) assert updated == connection.HeaderAssembly(1, True, 2, <<"ab":utf8>>, False) } pub fn add_fragment_over_count_cap_is_enhance_your_calm_test() { let assembly = connection.HeaderAssembly(1, True, 100, <<>>, False) - let result = connection.add_fragment(assembly, <<>>, http2.default_config()) + let result = connection.add_fragment(assembly, <<>>, http2.default_options()) assert result == Error(frame.EnhanceYourCalm) } @@ -242,7 +242,7 @@ pub fn add_fragment_over_byte_cap_is_enhance_your_calm_test() { let assembly = connection.HeaderAssembly(1, True, 1, <<0:size({ 65_536 * 8 })>>, False) let result = - connection.add_fragment(assembly, <<"x":utf8>>, http2.default_config()) + connection.add_fragment(assembly, <<"x":utf8>>, http2.default_options()) assert result == Error(frame.EnhanceYourCalm) } @@ -258,9 +258,9 @@ pub fn complete_header_block_oversized_list_is_enhance_your_calm_test() { let field = alpacki.HeaderField(<<"x":utf8>>, big_value, alpacki.WithoutIndexing) let assembly = connection.HeaderAssembly(1, True, 1, encode([field]), False) - let config = - http2.Config(..http2.default_config(), max_header_list_size: Some(16_384)) - let state = connection.State(..connection.test_state(), config:) + let options = + http2.Options(..http2.default_options(), max_header_list_size: Some(16_384)) + let state = connection.State(..connection.test_state(), options:) let result = connection.complete_header_block(state, assembly) assert result == connection.Terminate(Some(frame.EnhanceYourCalm)) } @@ -811,6 +811,47 @@ pub fn handle_data_sends_done_to_parked_reader_on_empty_end_stream_test() { let assert Ok(http2.DoneEvent([])) = process.receive(reply_to, 100) } +pub fn handle_data_keeps_reader_parked_on_empty_open_frame_test() { + let reply_to = process.new_subject() + let entry = inbound_stream(2_097_152, Some(reply_to)) + let state = + connection.State( + ..connection.test_state(), + streams: dict.from_list([#(1, entry)]), + highest_client_stream_id_seen: 1, + conn_recv_window: 2_097_152, + ) + + let assert connection.Proceed(state) = + connection.handle_data(state, 1, False, <<>>, 0) + + assert process.receive(reply_to, 0) == Error(Nil) + + let assert Ok(updated) = dict.get(state.streams, 1) + assert updated.parked_reader == Some(reply_to) + assert updated.request_half_closed == False +} + +pub fn handle_data_keeps_trailers_left_by_an_earlier_block_test() { + let entry = + connection.Stream(..inbound_stream(2_097_152, None), trailers: [ + #("x-checksum", "deadbeef"), + ]) + let state = + connection.State( + ..connection.test_state(), + streams: dict.from_list([#(1, entry)]), + highest_client_stream_id_seen: 1, + conn_recv_window: 2_097_152, + ) + + let assert connection.Proceed(state) = + connection.handle_data(state, 1, False, <<"abc":utf8>>, 3) + + let assert Ok(updated) = dict.get(state.streams, 1) + assert updated.trailers == [#("x-checksum", "deadbeef")] +} + pub fn handle_data_stream_window_violation_rejects_stream_test() { let entry = inbound_stream(5, None) let state = @@ -938,7 +979,7 @@ pub fn stream_recv_credit_above_low_water_mark_does_not_emit_test() { let entry = inbound_stream(300_000, None) let #(entry, increment) = - connection.stream_recv_credit(entry, http2.default_config()) + connection.stream_recv_credit(entry, http2.default_options()) assert increment == 0 assert entry.recv_window == 300_000 @@ -948,7 +989,7 @@ pub fn stream_recv_credit_at_low_water_mark_refills_to_high_test() { let entry = inbound_stream(262_144, None) let #(entry, increment) = - connection.stream_recv_credit(entry, http2.default_config()) + connection.stream_recv_credit(entry, http2.default_options()) assert increment == 2_097_152 - 262_144 assert entry.recv_window == 2_097_152 @@ -958,7 +999,7 @@ pub fn stream_recv_credit_below_low_water_mark_refills_to_high_test() { let entry = inbound_stream(1000, None) let #(entry, increment) = - connection.stream_recv_credit(entry, http2.default_config()) + connection.stream_recv_credit(entry, http2.default_options()) assert increment == 2_097_152 - 1000 assert entry.recv_window == 2_097_152 diff --git a/test/ewe_test.gleam b/test/ewe_test.gleam index b02a3c7..4e2721a 100644 --- a/test/ewe_test.gleam +++ b/test/ewe_test.gleam @@ -1,3 +1,10 @@ +import ewe +import ewe/internal/connection +import ewe/internal/http2/connection as http2 +import gleam/erlang/process +import gleam/http +import gleam/http/request +import gleam/option import gleeunit import logging @@ -7,3 +14,70 @@ pub fn main() -> Nil { gleeunit.main() } + +fn answer_reads( + conn_subject: process.Subject(http2.Reply(connection.Body)), + script: List(http2.BodyEvent), +) -> Nil { + case script { + [] -> Nil + [event, ..remaining] -> { + let assert http2.ReadBody(_stream_id, reply_to) = + process.receive_forever(conn_subject) + process.send(reply_to, event) + answer_reads(conn_subject, remaining) + } + } +} + +fn request_reading( + script: List(http2.BodyEvent), +) -> request.Request(ewe.Connection) { + let init_subject = process.new_subject() + process.spawn(fn() { + let conn_subject = process.new_subject() + process.send(init_subject, conn_subject) + answer_reads(conn_subject, script) + }) + + let body = + connection.Http2(http2.Connection( + connection: process.receive_forever(init_subject), + stream_id: 1, + has_body: True, + pending: <<>>, + pending_trailers: option.None, + read: 0, + body_read_timeout: 1000, + peer: Error(Nil), + )) + + request.Request( + method: http.Post, + headers: [], + body:, + scheme: http.Http, + host: "example.com", + port: option.None, + path: "/", + query: option.None, + ) +} + +pub fn read_body_chunk_reads_a_zero_chunk_size_as_one_test() { + let request = request_reading([http2.ChunkEvent(<<"abc":utf8>>)]) + + let assert Ok(ewe.Chunk(data, _request)) = + ewe.read_body_chunk(request, max_chunk_bytes: 0, limit: 1000) + + assert data == <<"a":utf8>> +} + +pub fn read_body_chunk_reads_a_negative_chunk_size_as_one_test() { + let request = request_reading([http2.ChunkEvent(<<"abc":utf8>>)]) + + let assert Ok(ewe.Chunk(data, _request)) = + ewe.read_body_chunk(request, max_chunk_bytes: -5, limit: 1000) + + assert data == <<"a":utf8>> +}