diff --git a/CHANGELOG.md b/CHANGELOG.md --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -5,6 +5,7 @@ - Sanitize CRLF sequences in outgoing HTTP response headers. - Eliminate `string.lowercase` in a codebase: validate and lowercase header field names in a single pass, validate and lowercase important protocol header values (like `transfer-encoding`, `connection`, `upgrade` and more) at parse time. - Include validation of trailer header field names and values during chunked body parsing. +- Remove redundant UTF-8 validation on WebSocket `Text` frame payloads. # v3.0.5 - 15.03.2026 diff --git a/src/ewe.gleam b/src/ewe.gleam --- a/src/ewe.gleam +++ b/src/ewe.gleam @@ -423,7 +423,7 @@ Builder(..builder, ipv6: True) } -/// Enables TLS (HTTPS) support, with provided certificate and key files. +/// Enables TLS (HTTPS) support, with provided certificate and key files. /// Crashes the program if the files don't exist or are invalid. pub fn enable_tls( builder: Builder, @@ -494,25 +494,22 @@ handler.loop(handler, on_crash, factory_name, builder.idle_timeout), ) |> glisten.bind(builder.interface) - |> fn(glisten_builder) { - case builder.ipv6 { - True -> glisten.with_ipv6(glisten_builder) - False -> glisten_builder - } - } - |> fn(glisten_builder) { - case builder.tls { - Some(#(cert, key)) -> glisten.with_tls(glisten_builder, cert, key) - // Uncomment once http2 will be implemented! - // |> glisten.with_http2 - None -> glisten_builder - } - } |> glisten.with_listener_name(builder.listener_name) - |> glisten.supervised(builder.port) + + let glisten = case builder.ipv6 { + True -> glisten.with_ipv6(glisten) + False -> glisten + } + + let glisten = case builder.tls { + Some(#(cert, key)) -> glisten.with_tls(glisten, cert, key) + // Uncomment once http2 will be implemented! + // |> glisten.with_http2 + None -> glisten + } supervisor.new(supervisor.OneForAll) - |> supervisor.add(glisten) + |> supervisor.add(glisten.supervised(glisten, builder.port)) |> supervisor.add(factory_child) |> supervisor.start() |> result.map(fn(started) { @@ -554,8 +551,8 @@ pub type Request = HttpRequest(Connection) -/// Reads body from the request. Returns `BodyTooLarge` if body exceeds -/// `bytes_limit`, or `InvalidBody` if malformed. Supports both chunked and +/// Reads body from the request. Returns `BodyTooLarge` if body exceeds +/// `bytes_limit`, or `InvalidBody` if malformed. Supports both chunked and /// content-length bodies. pub fn read_body( req: Request, @@ -672,18 +669,16 @@ let transport = req.body.transport let socket = req.body.socket - let factory_name = req.body.factory_name case chunked.send_response(resp, transport, socket) { Ok(Nil) -> { - let supervisor = factory.get_by_name(factory_name) - - let start_result = - factory.start_child(supervisor, fn() { + let started = + factory.get_by_name(req.body.factory_name) + |> factory.start_child(fn() { chunked.start(transport, socket, on_init, handler, on_close) }) - case start_result { + case started { Ok(started) -> { let _ = transport.controlling_process(transport, socket, started.pid) response.new(200) |> response.set_body(Chunked) @@ -788,12 +783,15 @@ ) -> Result(WebsocketMessage(user_message), Nil) { case message { websocket.Frame(websocks.Text(payload)) -> - bit_array.to_string(payload) |> result.map(Text) + Ok(Text(unsafe_to_string(payload))) websocket.Frame(websocks.Binary(payload)) -> Ok(Binary(payload)) websocket.UserMessage(user_message) -> Ok(User(user_message)) _ -> Error(Nil) } } + +@external(erlang, "gleam_stdlib", "identity") +fn unsafe_to_string(a: BitArray) -> String /// Upgrade request to a WebSocket connection. If the initial request is not /// valid for WebSocket upgrade, 400 response is sent. @@ -828,13 +826,12 @@ let transport = req.body.transport let socket = req.body.socket - let factory_name = req.body.factory_name case ewe_http.upgrade_websocket(req, transport, socket) { Ok(#(extensions, per_message_deflate)) -> { - let supervisor = factory.get_by_name(factory_name) - let start_result = - factory.start_child(supervisor, fn() { + let started = + factory.get_by_name(req.body.factory_name) + |> factory.start_child(fn() { websocket.start( transport, socket, @@ -846,10 +843,9 @@ ) }) - case start_result { - Ok(started) -> { - let _ = transport.controlling_process(transport, socket, started.pid) - + case started { + Ok(actor.Started(pid:, ..)) -> { + let _ = transport.controlling_process(transport, socket, pid) response.new(200) |> response.set_body(Websocket) } Error(_) -> response.new(500) |> response.set_body(Empty) @@ -864,13 +860,8 @@ conn: WebsocketConnection, bits: BitArray, ) -> Result(Nil, glisten.SocketReason) { - websocket.send_frame( - websocks.encode_binary_frame, - conn.transport, - conn.socket, - conn.context, - bits, - ) + websocks.encode_binary_frame + |> websocket.send_frame(conn.transport, conn.socket, conn.context, bits) } /// Sends a text frame to the websocket client. @@ -1049,19 +1040,18 @@ let transport = req.body.transport let socket = req.body.socket - let factory_name = req.body.factory_name case sse.send_response(transport, socket) { Ok(Nil) -> { - let supervisor = factory.get_by_name(factory_name) - let start_result = - factory.start_child(supervisor, fn() { + let started = + factory.get_by_name(req.body.factory_name) + |> factory.start_child(fn() { sse.start(transport, socket, on_init, handler, on_close) }) - case start_result { - Ok(started) -> { - let _ = transport.controlling_process(transport, socket, started.pid) + case started { + Ok(actor.Started(pid:, ..)) -> { + let _ = transport.controlling_process(transport, socket, pid) response.new(200) |> response.set_body(SSE) } Error(_) -> response.new(400) |> response.set_body(Empty) diff --git a/src/ewe/internal/stream/chunked.gleam b/src/ewe/internal/stream/chunked.gleam --- a/src/ewe/internal/stream/chunked.gleam +++ b/src/ewe/internal/stream/chunked.gleam @@ -58,7 +58,7 @@ |> process.select(subject) actor.initialised(state) - |> actor.returning(subject) + |> actor.returning(Nil) |> actor.selecting(selector) |> Ok }) @@ -92,20 +92,6 @@ } }) |> actor.start() - |> result.map(after_start(_, transport, socket)) -} - -/// Maps actor's starting value to Nil. -/// -fn after_start( - started: actor.Started(Subject(user_message)), - transport: Transport, - socket: Socket, -) -> actor.Started(Nil) { - let assert Ok(pid) = process.subject_owner(started.data) - let _ = transport.controlling_process(transport, socket, pid) - - actor.Started(..started, data: Nil) } /// Sends the end marker for chunked transfer encoding. diff --git a/src/ewe/internal/stream/sse.gleam b/src/ewe/internal/stream/sse.gleam --- a/src/ewe/internal/stream/sse.gleam +++ b/src/ewe/internal/stream/sse.gleam @@ -14,7 +14,7 @@ import glisten/transport.{type Transport} /// Sends a response for a Server-Sent Events connection. -/// +/// pub fn send_response(transport: Transport, socket: Socket) -> Result(Nil, Nil) { response.new(200) |> response.set_header("content-type", "text/event-stream") @@ -26,29 +26,29 @@ } /// Represents a Server-Sent Events connection. -/// +/// pub type SSEConnection { SSEConnection(transport: Transport, socket: Socket) } /// Represents an instruction on how Server-Sent Events connection should proceed. -/// +/// pub type SSENext(user_state) { Continue(user_state) NormalStop AbnormalStop(reason: String) } -/// Represents a message that can be sent to or received from the Server-Sent +/// Represents a message that can be sent to or received from the Server-Sent /// Events connection. -/// +/// pub type SSEMessages(user_message) { User(user_message) Close } /// Starts a new Server-Sent Events connection. -/// +/// pub fn start( transport: Transport, socket: Socket, @@ -56,13 +56,15 @@ handler: fn(SSEConnection, user_state, user_message) -> SSENext(user_state), on_close: fn(SSEConnection, user_state) -> Nil, ) -> Result(actor.Started(Nil), actor.StartError) { - actor.new_with_initialiser(1000, fn(_subject) { + actor.new_with_initialiser(1000, fn(_self) { + let _ = transport.set_opts(transport, socket, [ActiveMode(Active)]) + let subject = process.new_subject() let state = on_init(subject) let selector = create_socket_selector(subject) actor.initialised(state) - |> actor.returning(subject) + |> actor.returning(Nil) |> actor.selecting(selector) |> Ok }) @@ -89,11 +91,10 @@ } }) |> actor.start() - |> result.map(after_start(_, transport, socket)) } /// Creates a selector for the Server-Sent Events connection. -/// +/// fn create_socket_selector( user_subject: Subject(user_message), ) -> Selector(SSEMessages(user_message)) { @@ -103,23 +104,8 @@ |> process.select_record(atom.create("ssl_closed"), 1, fn(_) { Close }) } -/// Maps actor's starting value to Nil. -/// -fn after_start( - started: actor.Started(Subject(user_message)), - transport: Transport, - socket: Socket, -) -> actor.Started(Nil) { - let assert Ok(pid) = process.subject_owner(started.data) - let _ = transport.controlling_process(transport, socket, pid) - - let _ = transport.set_opts(transport, socket, [ActiveMode(Active)]) - - actor.Started(..started, data: Nil) -} - /// Represents a Server-Sent Events event. -/// +/// pub type SSEEvent { SSEEvent( event: Option(String), @@ -130,7 +116,7 @@ } /// Sends an event to the client. -/// +/// pub fn send_event( transport: Transport, socket: Socket, @@ -166,7 +152,7 @@ } /// Formats a field and value for a Server-Sent Events event. -/// +/// fn format(field: String, value: String) { field <> ": " <> value <> "\n" } diff --git a/src/ewe/internal/stream/websocket.gleam b/src/ewe/internal/stream/websocket.gleam --- a/src/ewe/internal/stream/websocket.gleam +++ b/src/ewe/internal/stream/websocket.gleam @@ -3,10 +3,9 @@ import gleam/bytes_tree import gleam/dynamic import gleam/erlang/atom -import gleam/erlang/process.{type Selector, type Subject} +import gleam/erlang/process.{type Selector} import gleam/option.{type Option, None, Some} import gleam/otp/actor -import gleam/result import glisten/socket.{type Socket, type SocketReason} import glisten/socket/options.{ActiveMode, Count} import glisten/transport.{type Transport} @@ -107,7 +106,12 @@ extensions: List(String), permessage_deflate: Bool, ) -> Result(actor.Started(Nil), actor.StartError) { - actor.new_with_initialiser(1000, fn(subject) { + actor.new_with_initialiser(1000, fn(_self) { + let _ = + transport.set_opts(transport, socket, [ + ActiveMode(Count(socket_active_count)), + ]) + let context_takeovers = websocks.get_context_takeovers(extensions) let compression = case permessage_deflate { True -> Some(context_takeovers) @@ -126,7 +130,7 @@ WebsocketState(user_state:, context:) |> actor.initialised() |> actor.selecting(selector) - |> actor.returning(subject) + |> actor.returning(Nil) |> Ok }) |> actor.on_message(fn(state, msg) { @@ -160,7 +164,7 @@ } }) |> actor.start() - |> result.map(after_start(_, transport, socket)) + // |> result.map(after_start(_, transport, socket)) } // Creates selector for glisten socket events. @@ -375,21 +379,6 @@ } None -> actor.stop() } -} - -// Maps actor's starting value to Nil. -// -fn after_start( - started: actor.Started(Subject(InternalMessage(user_message))), - transport: Transport, - socket: Socket, -) -> actor.Started(Nil) { - let _ = - transport.set_opts(transport, socket, [ - ActiveMode(Count(socket_active_count)), - ]) - - actor.Started(..started, data: Nil) } /// Sends a frame to the WebSocket.