From 05971c247d238ba1924c9219c716149f04b18b9a Mon Sep 17 00:00:00 2001 From: vshakitskiy Date: Wed, 5 Aug 2026 10:47:08 +0300 Subject: [PATCH] further fixing changes --- CHANGELOG.md | 2 + src/ewe.gleam | 18 ++-- src/ewe/internal/ewe_ffi.erl | 17 +++- src/ewe/internal/http1.gleam | 109 ++++++++++++++++-------- src/ewe/internal/http1/connection.gleam | 5 ++ src/ewe/internal/http1/sse.gleam | 81 +++++++++++------- src/ewe/internal/http1/websocket.gleam | 93 ++++++++++++-------- 7 files changed, 211 insertions(+), 114 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 3cfad8c..8164285 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -15,6 +15,8 @@ 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. +- In replacement of the `idle_timeout` builder function there is now the + `idle_timeout` field of `Http1Config`. - Switch from `erlang:decode_packet` to self implemented parsing solution. - Responses are now framed by the server, which computes `content-length` or `transfer-encoding: chunked` for a streamed body and drops the handler's own diff --git a/src/ewe.gleam b/src/ewe.gleam index cececb2..0b5936e 100644 --- a/src/ewe.gleam +++ b/src/ewe.gleam @@ -702,10 +702,10 @@ pub fn send_event(conn: SseConnection, event: SseEvent) -> Nil { /// reusable afterwards as long as the handler ended the stream itself and the /// client sent nothing during it. /// -/// `on_init` is called once, with a subject the rest of your program uses to -/// push messages at the client, and returns the starting state. `handler` is -/// called for each message sent to that subject. `on_close` is called once -/// however the stream ended. +/// - `on_init` is called once, with a subject the rest of your program uses to +/// push messages at the client, and returns the starting state. +/// - `handler` is called for each message sent to that subject. +/// - `on_close` is called once however the stream ended. pub fn sse( response: response.Response(a), on_init on_init: fn(process.Subject(user_message)) -> user_state, @@ -897,12 +897,12 @@ pub fn send_close_frame( /// A request that is not a valid handshake is answered with a 400 and the /// handler is never run. /// -/// `on_init` is called once, with an empty selector to add whatever the rest of -/// your program sends this connection to and returns the starting state along -/// with that selector. -/// `handler` is called for each frame from the client and each message the +/// - `on_init` is called once with an empty selector to add whatever the rest +/// of your program sends this connection to and returns the starting state +/// along with that selector. +/// - `handler` is called for each frame from the client and each message the /// selector picks up. -/// `on_close` is called once however the WebSocket ended. +/// - `on_close` is called once however the WebSocket ended. pub fn upgrade_websocket( request request: request.Request(Connection), on_init on_init: fn(WebsocketConnection, process.Selector(user_message)) -> diff --git a/src/ewe/internal/ewe_ffi.erl b/src/ewe/internal/ewe_ffi.erl index e9f1455..e696159 100644 --- a/src/ewe/internal/ewe_ffi.erl +++ b/src/ewe/internal/ewe_ffi.erl @@ -1,10 +1,25 @@ -module(ewe_ffi). --export([identity/1, now_datetime/0, set_http_date/1, get_http_date/0]). +-export([ + identity/1, + now_datetime/0, + set_http_date/1, + get_http_date/0, + rescue_handler/1 +]). identity(X) -> X. +rescue_handler(Func) -> + try + {ok, Func()} + catch + Class:Reason:Stacktrace -> + Formatted = erl_error:format_exception(Class, Reason, Stacktrace), + {error, unicode:characters_to_binary(Formatted)} + end. + now_datetime() -> {Date, Time} = calendar:universal_time(), Weekday = calendar:day_of_the_week(Date), diff --git a/src/ewe/internal/http1.gleam b/src/ewe/internal/http1.gleam index 430051c..c0f6068 100644 --- a/src/ewe/internal/http1.gleam +++ b/src/ewe/internal/http1.gleam @@ -62,41 +62,55 @@ pub fn handle_message( upgrade: metadata.upgrade, ) - let response = - state.handler(to_request(head, connection, body_connection)) + let request = to_request(head, connection, body_connection) - let drained = drain_messages(self) - let ResolvedBody(buffer, body_keep_alive) = - resolve_body(body_connection, drained.body, state.config) - let keep_alive = - http1.and_keep_alive(metadata.keep_alive, body_keep_alive) - - let sent = case - encoder.encode_response(response, head.method, head.version, keep_alive) - { - Ok(encoded) -> - send_response(encoded, connection.transport, connection.socket, self) - Error(encoder.UnsafeHeader(name)) -> { - logging.log( - logging.Error, - "Handler produced an unsafe response header: " <> name, - ) - file.release_body(response.body) - - transport.send( - connection.transport, - connection.socket, - encoder.internal_server_error(), - ) - |> result.replace(SentClose) - } - } + case rescue_handler(fn() { state.handler(request) }) { + Error(details) -> crashed(connection, details) + Ok(response) -> { + let drained = drain_messages(self) + let ResolvedBody(buffer, body_keep_alive) = + resolve_body(body_connection, drained.body, state.config) + let keep_alive = + http1.and_keep_alive(metadata.keep_alive, body_keep_alive) + + let sent = case + encoder.encode_response( + response, + head.method, + head.version, + keep_alive, + ) + { + Ok(encoded) -> + send_response( + encoded, + connection.transport, + connection.socket, + self, + ) + Error(encoder.UnsafeHeader(name)) -> { + logging.log( + logging.Error, + "Handler produced an unsafe response header: " <> name, + ) + file.release_body(response.body) + + transport.send( + connection.transport, + connection.socket, + encoder.internal_server_error(), + ) + |> result.replace(SentClose) + } + } - case sent { - Ok(SentKeepAlive) -> await_next_request(state, buffer, connection) - Ok(SentClose) -> Close - Ok(SentAbnormal(reason)) -> CloseAbnormal(reason) - Error(_reason) -> Close + case sent { + Ok(SentKeepAlive) -> await_next_request(state, buffer, connection) + Ok(SentClose) -> Close + Ok(SentAbnormal(reason)) -> CloseAbnormal(reason) + Error(_reason) -> Close + } + } } } Ok(parser.Incomplete) -> { @@ -122,9 +136,29 @@ pub fn handle_message( } } -/// A pipelining client sends its next request without waiting for this -/// response, so anything left over has to be handled now. Waiting on the socket -/// for it would deadlock: those bytes have already arrived. +/// 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, +) -> Next { + logging.log( + logging.Error, + "Caught a crash in the request handler: " <> details, + ) + + let _sent = + transport.send( + connection.transport, + connection.socket, + encoder.internal_server_error(), + ) + + 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, @@ -349,3 +383,6 @@ fn do_drain_remaining( Error(_reason) -> ResolvedBody(<<>>, http1.CloseAfterResponse) } } + +@external(erlang, "ewe_ffi", "rescue_handler") +fn rescue_handler(handler: fn() -> a) -> Result(a, String) diff --git a/src/ewe/internal/http1/connection.gleam b/src/ewe/internal/http1/connection.gleam index cc71aa4..15a5df8 100644 --- a/src/ewe/internal/http1/connection.gleam +++ b/src/ewe/internal/http1/connection.gleam @@ -45,6 +45,11 @@ pub type Connection { ) } +/// 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) diff --git a/src/ewe/internal/http1/sse.gleam b/src/ewe/internal/http1/sse.gleam index 0222349..8fedb6e 100644 --- a/src/ewe/internal/http1/sse.gleam +++ b/src/ewe/internal/http1/sse.gleam @@ -26,10 +26,7 @@ pub fn run( case activate(conn) { Ok(Nil) -> loop(conn, handle, selector(subject), state, Clean, step, on_close) - Error(reason) -> - socket.reason_to_string(reason) - |> connection.StoppedAbnormal - |> ended(conn, handle, state, on_close, http1.CloseAfterResponse, _) + Error(reason) -> socket_failed(conn, handle, state, on_close, reason) } } @@ -43,11 +40,41 @@ fn ended( keep_alive: http1.KeepAlive, outcome: connection.Outcome, ) -> connection.Outcome { - on_close(handle, state) + let _dead = stream.rescue_dead(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, + state: user_state, + on_close: fn(connection.SseConnection, user_state) -> Nil, +) -> connection.Outcome { + ended( + conn, + handle, + state, + on_close, + http1.CloseAfterResponse, + connection.Stopped, + ) +} + +/// The socket gave out, which ends the stream whatever it was doing. +fn socket_failed( + conn: http1.SseConnection, + handle: connection.SseConnection, + state: user_state, + on_close: fn(connection.SseConnection, user_state) -> Nil, + reason: socket.SocketReason, +) -> connection.Outcome { + socket.reason_to_string(reason) + |> connection.StoppedAbnormal + |> 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 { @@ -69,15 +96,12 @@ fn loop( // 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) - Disconnected -> - ended( - conn, - handle, - state, - on_close, - http1.CloseAfterResponse, - connection.Stopped, - ) + Exhausted -> + case activate(conn) { + Ok(Nil) -> loop(conn, handle, selector, state, reuse, step, on_close) + Error(reason) -> socket_failed(conn, handle, state, on_close, reason) + } + Disconnected -> dropped(conn, handle, state, on_close) Failed(reason) -> connection.StoppedAbnormal(reason) |> ended(conn, handle, state, on_close, http1.CloseAfterResponse, _) @@ -85,15 +109,7 @@ fn loop( // A send inside the handler can find the client gone before the socket // has told us, so both routes out land on the same teardown. case stream.rescue_dead(fn() { step(handle, state, message) }) { - Error(_reason) -> - ended( - conn, - handle, - state, - on_close, - http1.CloseAfterResponse, - connection.Stopped, - ) + Error(_reason) -> dropped(conn, handle, state, on_close) Ok(sse.Proceed(state)) -> loop(conn, handle, selector, state, reuse, step, on_close) Ok(sse.Halt(outcome)) -> @@ -134,13 +150,11 @@ pub fn send(conn: http1.SseConnection, event: sse.Event) -> Nil { } } -/// glisten re arms `{active, once}` only once its loop callback returns, and an -/// SSE stream does not return until it is over. Without switching to full -/// active mode the socket goes quiet for the whole stream, `tcp_closed` -/// included, and a client hanging up would never be noticed. +/// 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.Active), + options.ActiveMode(options.Count(http1.active_count)), ]) } @@ -150,6 +164,7 @@ type Received(user_message) { Disconnected Failed(reason: String) ClientData + Exhausted } /// glisten handles socket messages in its own loop, which is blocked for the @@ -165,6 +180,12 @@ fn selector( |> process.select_record(atom.create("ssl_error"), 2, failed) |> process.select_record(atom.create("tcp"), 2, client_data) |> process.select_record(atom.create("ssl"), 2, client_data) + |> process.select_record(atom.create("tcp_passive"), 1, exhausted) + |> process.select_record(atom.create("ssl_passive"), 1, exhausted) +} + +fn exhausted(_record: dynamic.Dynamic) -> Received(user_message) { + Exhausted } fn disconnected(_record: dynamic.Dynamic) -> Received(user_message) { @@ -177,9 +198,7 @@ fn failed(record: dynamic.Dynamic) -> Received(user_message) { |> Failed } -/// Clients are not expected to send anything once the stream is open. Matching -/// it anyway keeps a chatty one from growing the mailbox without bound, at the -/// cost of giving up on reusing the connection. +/// 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 d10d8ea..1921f3f 100644 --- a/src/ewe/internal/http1/websocket.gleam +++ b/src/ewe/internal/http1/websocket.gleam @@ -77,9 +77,6 @@ pub fn handshake( } } -/// How many socket messages are delivered before the loop rearms. -const active_count = 100 - pub fn run( conn: http1.WebsocketConnection, on_init: fn(connection.WebsocketConnection, process.Selector(user_message)) -> @@ -97,26 +94,51 @@ pub fn run( case activate(conn) { Ok(Nil) -> loop(conn, merge_socket_selector(messages), state, step, on_close) - Error(reason) -> - socket.reason_to_string(reason) - |> connection.StoppedAbnormal - |> ended(conn, state, on_close, _) + Error(reason) -> socket_failed(conn, state, on_close, reason) } } -/// Every way a socket ends frees the compression resources the context holds -/// and runs the handler's `on_close`. +/// 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, on_close: fn(connection.WebsocketConnection, user_state) -> Nil, outcome: connection.Outcome, ) -> connection.Outcome { + let handle = connection.Http1Websocket(conn) + let _dead = stream.rescue_dead(fn() { on_close(handle, state) }) + websocks.close_context(conn.context) - on_close(connection.Http1Websocket(conn), state) outcome } +fn stopped( + conn: http1.WebsocketConnection, + state: user_state, + on_close: fn(connection.WebsocketConnection, user_state) -> Nil, +) -> connection.Outcome { + ended(conn, state, on_close, connection.Stopped) +} + +/// The socket gave out which ends the connection whatever it was doing. +fn socket_failed( + conn: http1.WebsocketConnection, + state: user_state, + on_close: fn(connection.WebsocketConnection, user_state) -> Nil, + reason: socket.SocketReason, +) -> connection.Outcome { + socket.reason_to_string(reason) + |> connection.StoppedAbnormal + |> ended(conn, state, on_close, _) +} + +/// Where the loop carries on once the handler has seen a message. +type Resume { + DrainBuffer + AwaitSocket +} + fn loop( conn: http1.WebsocketConnection, selector: process.Selector(Received(user_message)), @@ -129,24 +151,22 @@ fn loop( on_close: fn(connection.WebsocketConnection, user_state) -> Nil, ) -> connection.Outcome { case process.selector_receive_forever(selector) { - Closed -> ended(conn, state, on_close, connection.Stopped) + Closed -> stopped(conn, state, on_close) Failed(reason) -> ended(conn, state, on_close, connection.StoppedAbnormal(reason)) Exhausted -> case activate(conn) { Ok(Nil) -> loop(conn, selector, state, step, on_close) - Error(reason) -> - socket.reason_to_string(reason) - |> connection.StoppedAbnormal - |> ended(conn, state, on_close, _) + Error(reason) -> socket_failed(conn, state, on_close, reason) } Packet(data) -> 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, _) + |> deliver(conn, selector, state, step, on_close, AwaitSocket, _) } } @@ -178,13 +198,13 @@ fn drain( // the handler's business. websocks.Control(websocks.Ping(payload)) -> case - write(conn, websocks.encode_pong_frame(payload:, masking: none)) + write( + conn, + websocks.encode_pong_frame(payload:, masking: option.None), + ) { Ok(Nil) -> drain(conn, selector, state, step, on_close) - Error(reason) -> - socket.reason_to_string(reason) - |> connection.StoppedAbnormal - |> ended(conn, state, on_close, _) + Error(reason) -> socket_failed(conn, state, on_close, reason) } websocks.Control(websocks.Pong(_payload)) -> drain(conn, selector, state, step, on_close) @@ -194,10 +214,10 @@ fn drain( close(conn, reason) |> resolve(conn, state, on_close, _) websocks.Text(payload) -> websocket.TextFrame(unsafe_to_string(payload)) - |> deliver(conn, selector, state, step, on_close, _) + |> deliver(conn, selector, state, step, on_close, DrainBuffer, _) websocks.Binary(payload) -> websocket.BinaryFrame(payload) - |> deliver(conn, selector, state, step, on_close, _) + |> 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) @@ -216,19 +236,23 @@ fn deliver( websocket.Message(user_message), ) -> websocket.Step(user_state, user_message), on_close: fn(connection.WebsocketConnection, user_state) -> Nil, + resume: Resume, message: websocket.Message(user_message), ) -> connection.Outcome { let handle = connection.Http1Websocket(conn) case stream.rescue_dead(fn() { step(handle, state, message) }) { - Error(_reason) -> ended(conn, state, on_close, connection.Stopped) + Error(_reason) -> stopped(conn, state, on_close) Ok(websocket.Proceed(user_state: state, messages:)) -> { let selector = case messages { option.Some(messages) -> merge_socket_selector(messages) option.None -> selector } - drain(conn, selector, state, step, on_close) + case resume { + DrainBuffer -> drain(conn, selector, state, step, on_close) + AwaitSocket -> loop(conn, selector, state, step, on_close) + } } Ok(websocket.Halt(outcome)) -> ended(conn, state, on_close, outcome) } @@ -243,11 +267,8 @@ fn resolve( sent: Result(Nil, socket.SocketReason), ) -> connection.Outcome { case sent { - Ok(Nil) -> ended(conn, state, on_close, connection.Stopped) - Error(reason) -> - socket.reason_to_string(reason) - |> connection.StoppedAbnormal - |> ended(conn, state, on_close, _) + Ok(Nil) -> stopped(conn, state, on_close) + Error(reason) -> socket_failed(conn, state, on_close, reason) } } @@ -255,7 +276,7 @@ pub fn send_text(conn: http1.WebsocketConnection, text: String) -> Nil { websocks.encode_text_frame( payload: bit_array_from_string(text), context: conn.context, - masking: none, + masking: option.None, ) |> write_or_die(conn, _) } @@ -264,7 +285,7 @@ pub fn send_binary(conn: http1.WebsocketConnection, data: BitArray) -> Nil { websocks.encode_binary_frame( payload: data, context: conn.context, - masking: none, + masking: option.None, ) |> write_or_die(conn, _) } @@ -273,12 +294,10 @@ pub fn send_close( conn: http1.WebsocketConnection, reason: websocks.CloseReason, ) -> Nil { - websocks.encode_close_frame(reason:, masking: none) + websocks.encode_close_frame(reason:, masking: option.None) |> write_or_die(conn, _) } -const none = option.None - fn with_context( conn: http1.WebsocketConnection, context: websocks.Context, @@ -290,7 +309,7 @@ fn close( conn: http1.WebsocketConnection, reason: websocks.CloseReason, ) -> Result(Nil, socket.SocketReason) { - write(conn, websocks.encode_close_frame(reason:, masking: none)) + write(conn, websocks.encode_close_frame(reason:, masking: option.None)) } fn write( @@ -314,7 +333,7 @@ fn activate( conn: http1.WebsocketConnection, ) -> Result(Nil, socket.SocketReason) { transport.set_opts(conn.transport, conn.socket, [ - options.ActiveMode(options.Count(active_count)), + options.ActiveMode(options.Count(http1.active_count)), ]) } -- 2.51.2