diff --git a/CHANGELOG.md b/CHANGELOG.md index a79dde6..9d7f027 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -2,8 +2,10 @@ # Unreleased +- Bump dependency package versions - Move examples to appropriate folder - Improve examples with useful documentation lines +- Replace active socket tcp message decoding for WebSocket selector # v2.1.2 - 19.11.2025 diff --git a/Makefile b/Makefile index 4ab73bf..7dd27c1 100644 --- a/Makefile +++ b/Makefile @@ -16,4 +16,7 @@ autobahn_docker: -v "${PWD}/autobahn:/reports" \ --network host \ crossbario/autobahn-testsuite \ - wstest -m fuzzingclient -s /autobahn.json \ No newline at end of file + wstest -m fuzzingclient -s /autobahn.json + +autobahn_serve: + cd "${PWD}/autobahn" && bun run serve.ts diff --git a/autobahn/config.json b/autobahn/config.json index 10f44b0..7c4551f 100644 --- a/autobahn/config.json +++ b/autobahn/config.json @@ -6,12 +6,10 @@ "agent": "ewe" } ], - "cases": [ - "*" - ], + "cases": ["*"], "exclude-cases": [], "exclude-agent-cases": {}, "optioms": { "failByDrop": false } -} \ No newline at end of file +} diff --git a/autobahn/serve.ts b/autobahn/serve.ts new file mode 100644 index 0000000..1f3fe79 --- /dev/null +++ b/autobahn/serve.ts @@ -0,0 +1,19 @@ +const server = Bun.serve({ + async fetch(req) { + const url = new URL(req.url); + const filePath = `./server${url.pathname}`; + + const file = Bun.file(filePath); + const exists = await file.exists(); + if (exists) { + return new Response(file); + } + + return new Response("Not Found", { status: 404 }); + }, + port: 3000, + hostname: "0.0.0.0", +}); + +console.log(`Server running at ${server.url}`); +console.log(`Access from your network at http://192.168.1.41:${server.port}`); diff --git a/src/ewe.gleam b/src/ewe.gleam index b492267..be4dc06 100644 --- a/src/ewe.gleam +++ b/src/ewe.gleam @@ -84,12 +84,12 @@ //// ] //// } //// ] -//// +//// //// const callback = () => { //// const list = document.querySelector(".sidebar > ul:last-of-type") //// const sortedLists = document.createDocumentFragment() //// const sortedMembers = document.createDocumentFragment() -//// +//// //// for (const section of docs) { //// sortedLists.append((() => { //// const node = document.createElement("h3") @@ -101,12 +101,12 @@ //// node.append(section.header) //// return node //// })()) -//// +//// //// const sortedList = document.createElement("ul") //// sortedLists.append(sortedList) -//// +//// //// const sortedFunctions = [...section.functions].sort() -//// +//// //// for (const funcName of sortedFunctions) { //// const href = `#${funcName}` //// const member = document.querySelector( @@ -117,7 +117,7 @@ //// sortedMembers.append(member) //// } //// } -//// +//// //// document.querySelector(".sidebar").insertBefore(sortedLists, list) //// document //// .querySelector(".module-members:has(#module-values)") @@ -126,7 +126,7 @@ //// document.querySelector("#module-values").nextSibling //// ) //// } -//// +//// //// document.readyState !== "loading" //// ? callback() //// : document.addEventListener( @@ -136,10 +136,6 @@ //// ) //// -// ----------------------------------------------------------------------------- -// IMPORTS -// ----------------------------------------------------------------------------- - import ewe/internal/file import ewe/internal/handler import ewe/internal/http1 as ewe_http @@ -192,7 +188,7 @@ pub type IpAddress { /// pub fn ip_address_to_string(address address: IpAddress) -> String { ewe_to_glisten_ip(address) - |> glisten.ip_address_to_string() + |> glisten.ip_address_to_string } fn glisten_to_ewe_ip(ip: glisten.IpAddress) -> IpAddress { @@ -618,8 +614,10 @@ pub fn supervised( /// pub type BodyError { /// Body is larger than the provided limit. + /// BodyTooLarge /// Body is malformed. + /// InvalidBody } @@ -691,9 +689,9 @@ fn consumer_adapter( // CHUNKED RESPONSE // ----------------------------------------------------------------------------- -/// Represents a chunked response body. This type is used to send a chunked +/// Represents a chunked response body. This type is used to send a chunked /// response to the client. -/// +/// pub type ChunkedBody = chunked.ChunkedBody @@ -710,7 +708,7 @@ pub opaque type ChunkedNext(user_state) { } /// Instructs chunked response to continue processing. -/// +/// pub fn chunked_continue(user_state: user_state) -> ChunkedNext(user_state) { ChunkedContinue(user_state) } @@ -738,9 +736,9 @@ fn to_internal_chunked_next( } /// Sets up the connection for chunked response. -/// -/// `on_init` function is called once the chunked response process is -/// initialized. The argument is subject that can be used to send chunks to the +/// +/// `on_init` function is called once the chunked response process is +/// initialized. The argument is subject that can be used to send chunks to the /// client. It must return initial state. /// /// `handler` function is called for every message received. It must return @@ -992,47 +990,47 @@ pub fn send_text_frame( ) } -/// WebSocket close codes that can be sent when closing a connection. The `data` -/// parameter allows you to include payload up to 123 bytes in size. -/// +/// WebSocket close codes that can be sent when closing a connection. The `data` +/// parameter allows you to include payload up to 123 bytes in size. +/// pub type CloseCode { - /// Standard graceful shutdown (1000). Use when connection completed + /// Standard graceful shutdown (1000). Use when connection completed /// successfully. - /// + /// NormalClosure(data: String) - /// Invalid message format (1007). Received payload that doesn't match what + /// Invalid message format (1007). Received payload that doesn't match what /// you expected. - /// + /// InvalidPayloadData(data: String) /// Application policy violation (1008).Client broke your rules - failed /// authentication, hit rate limits, or violated business logic. - /// + /// PolicyViolation(data: String) - /// Message exceeds size limits (1009). Client sent something bigger than + /// Message exceeds size limits (1009). Client sent something bigger than /// your application allows. - /// + /// MessageTooBig(data: String) - /// Server encountered unexpected error (1011). Something went wrong on your + /// Server encountered unexpected error (1011). Something went wrong on your /// side that prevents handling the connection. - /// + /// InternalError(data: String) - /// Server is restarting (1012). Planned restart - clients can reconnect + /// Server is restarting (1012). Planned restart - clients can reconnect /// after a bit. - /// + /// ServiceRestart(data: String) - /// Temporary server overload (1013). Use when server is temporarily + /// Temporary server overload (1013). Use when server is temporarily /// unavailable, client should retry. - /// + /// TryAgainLater(data: String) - /// Gateway/proxy received invalid response (1014). You're acting as a proxy + /// Gateway/proxy received invalid response (1014). You're acting as a proxy /// and the upstream server gave you garbage. - /// + /// BadGateway(data: String) /// Custom close codes 3000-4999 for application-specific use. - /// + /// CustomCloseCode(code: Int, data: String) /// Close without a specific reason. - /// + /// NoCloseReason } @@ -1120,7 +1118,7 @@ fn to_internal_sse_next(next: SSENext(user_state)) -> sse.SSENext(user_state) { /// - `id`: event ID. /// - `retry`: The reconnection time. If the connection to the server is lost, /// the browser will wait for the specified time before attempting to reconnect. -/// +/// /// Can be created using `ewe.event` and modified with `ewe.event_name`, /// `ewe.event_id`, and `ewe.event_retry`. /// @@ -1160,7 +1158,7 @@ pub fn event_retry(event: SSEEvent, retry: Int) -> SSEEvent { /// /// `handler` function is called for every subject's message received. It must /// return instruction on how SSE connection should proceed. -/// +/// /// `on_close` function is called when SSE process is going to be stopped. /// pub fn sse( diff --git a/src/ewe/internal/http2.gleam b/src/ewe/internal/http2.gleam deleted file mode 100644 index 8b13789..0000000 --- a/src/ewe/internal/http2.gleam +++ /dev/null @@ -1 +0,0 @@ - diff --git a/src/ewe/internal/stream/websocket.gleam b/src/ewe/internal/stream/websocket.gleam index df1823f..7897f0e 100644 --- a/src/ewe/internal/stream/websocket.gleam +++ b/src/ewe/internal/stream/websocket.gleam @@ -1,7 +1,7 @@ import exception import gleam/bit_array import gleam/bytes_tree -import gleam/dynamic/decode +import gleam/dynamic import gleam/erlang/atom import gleam/erlang/process.{type Selector, type Subject} import gleam/option.{type Option, None, Some} @@ -39,19 +39,19 @@ pub type WebsocketNext(user_state, user_message) { AbnormalStop(reason: String) } -/// Internal state maintained by the WebSocket actor. -/// +// Internal state maintained by the WebSocket actor. +// type WebsocketState(user_state) { WebsocketState(user_state: user_state, context: websocks.Context) } -/// Type alias for actor next steps. -/// +// Type alias for actor next steps. +// type ActorNext(user_state, user_message) = actor.Next(WebsocketState(user_state), InternalMessage(user_message)) -/// Internal messages used by the WebSocket actor. -/// +// Internal messages used by the WebSocket actor. +// type InternalMessage(user_message) { Packet(BitArray) Close @@ -60,41 +60,41 @@ type InternalMessage(user_message) { Invalid } -/// Function called when the WebSocket connection is initialized. -/// +// Function called when the WebSocket connection is initialized. +// type OnInit(user_state, user_message) = fn(WebsocketConnection, Selector(user_message)) -> #(user_state, Selector(user_message)) -/// Function called to handle incoming WebSocket messages. -/// +// Function called to handle incoming WebSocket messages. +// type Handler(user_state, user_message) = fn(WebsocketConnection, user_state, WebsocketMessage(user_message)) -> WebsocketNext(user_state, user_message) -/// Function called when the WebSocket connection is closed. -/// +// Function called when the WebSocket connection is closed. +// type OnClose(user_state) = fn(WebsocketConnection, user_state) -> Nil -/// Error message for malformed messages. -/// +// Error message for malformed messages. +// const malformed = "Received malformed message" -/// Error message for crashed WebSocket handler. -/// +// Error message for crashed WebSocket handler. +// const crashed = "Crash in websocket handler" -/// Error message for failed PONG frame. -/// +// Error message for failed PONG frame. +// const failed_pong = "Failed to send PONG frame" -/// Error message for sending WebSocket message from non-owning process. -/// +// Error message for sending WebSocket message from non-owning process. +// const non_owning_process = "Sending WebSocket message from non-owning process" -/// Active count for socket. -/// +// Active count for socket. +// const socket_active_count = 100 /// Starts a new WebSocket connection. @@ -164,34 +164,26 @@ pub fn start( |> result.map(after_start(_, transport, socket)) } -/// Creates a selector for valid TCP/SSL records. -/// -fn select_valid_record( - selector: Selector(InternalMessage(user_message)), - binary_atom: String, -) -> Selector(InternalMessage(user_message)) { - process.select_record(selector, atom.create(binary_atom), 2, fn(record) { - decode.run(record, { - use data <- decode.field(2, decode.bit_array) - decode.success(Packet(data)) - }) - |> result.unwrap(Invalid) - }) -} - -/// Creates selector for glisten socket events. -/// +// Creates selector for glisten socket events. +// fn create_socket_selector() -> Selector(InternalMessage(user_message)) { process.new_selector() - |> select_valid_record("tcp") - |> select_valid_record("ssl") + |> process.select_record(atom.create("tcp"), 2, fn(record) { + Packet(coerce_tcp_message(record)) + }) + |> process.select_record(atom.create("tcp"), 2, fn(record) { + Packet(coerce_tcp_message(record)) + }) |> process.select_record(atom.create("tcp_closed"), 1, fn(_) { Close }) |> process.select_record(atom.create("ssl_closed"), 1, fn(_) { Close }) |> process.select_record(atom.create("tcp_passive"), 1, fn(_) { TcpPassive }) } -/// Handles incoming packet data, decoding frames and processing them. -/// +@external(erlang, "ewe_ffi", "coerce_tcp_message") +fn coerce_tcp_message(record: dynamic.Dynamic) -> BitArray + +// Handles incoming packet data, decoding frames and processing them. +// fn handle_valid_packet( transport: Transport, socket: Socket, @@ -201,7 +193,7 @@ fn handle_valid_packet( on_close: OnClose(user_state), ) -> ActorNext(user_state, user_message) { let conn = WebsocketConnection(transport, socket, state.context) - let result = + let processed = websocks.process_incoming_frames( data, state.context, @@ -214,7 +206,7 @@ fn handle_valid_packet( handle_frame, ) - case result { + case processed { Ok(#(resolved_state, context)) -> { case resolved_state.next { Continue(user_state, selector) -> { @@ -230,15 +222,12 @@ fn handle_valid_packet( handle_close(on_close, state, conn, Some(reason)) } } - Error(_violation) -> { - // echo violation as "violation during frame resolving" - handle_close(on_close, state, conn, Some(malformed)) - } + Error(_violation) -> handle_close(on_close, state, conn, Some(malformed)) } } -/// Represents the state of the WebSocket connection when resolving frames. -/// +// Represents the state of the WebSocket connection when resolving frames. +// type ResolveState(user_state, user_message) { ResolveState( socket: Socket, @@ -286,6 +275,7 @@ fn handle_frame( } } } + websocks.Control(websocks.Close(reason)) -> { let _ = transport.send( @@ -309,7 +299,7 @@ fn handle_frame( case call { Ok(Continue(user_state, new_selector)) -> { let next_selector = - user_selector(new_selector) + option.map(new_selector, process.map_selector(_, User)) |> option.or(selector) |> option.map(process.merge_selector(create_socket_selector(), _)) @@ -327,8 +317,8 @@ fn handle_frame( } } -/// Handles user messages sent to the WebSocket. -/// +// Handles user messages sent to the WebSocket. +// fn handle_user_message( transport: Transport, socket: Socket, @@ -346,7 +336,7 @@ fn handle_user_message( case call { Ok(Continue(new_user_state, new_selector)) -> { let next_selector = - user_selector(new_selector) + option.map(new_selector, process.map_selector(_, User)) |> option.map(process.merge_selector(create_socket_selector(), _)) let next = @@ -364,16 +354,8 @@ fn handle_user_message( } } -/// Maps user selector to internal message. -/// -fn user_selector( - selector: Option(Selector(user_message)), -) -> Option(Selector(InternalMessage(user_message))) { - option.map(selector, fn(selector) { process.map_selector(selector, User) }) -} - -/// Handles WebSocket connection closure. -/// +// Handles WebSocket connection closure. +// fn handle_close( on_close: OnClose(user_state), state: WebsocketState(user_state), @@ -395,8 +377,8 @@ fn handle_close( } } -/// Maps actor's starting value to Nil. -/// +// Maps actor's starting value to Nil. +// fn after_start( started: actor.Started(Subject(InternalMessage(user_message))), transport: Transport, diff --git a/src/ewe_ffi.erl b/src/ewe_ffi.erl index 0381273..212b6aa 100644 --- a/src/ewe_ffi.erl +++ b/src/ewe_ffi.erl @@ -1,7 +1,13 @@ -module(ewe_ffi). -export([close_file/1, decode_packet/3, init_clock_storage/0, lookup_http_date/0, now/0, - now_microseconds/0, open_file/1, rescue/1, set_http_date/1, validate_field_value/1]). + now_microseconds/0, open_file/1, set_http_date/1, validate_field_value/1, coerce_tcp_message/1]). + +% Socket +% ----------------------------------------------------------------------------- + +coerce_tcp_message({tcp, _Socket, Data}) -> Data; +coerce_tcp_message({ssl, _Socket, Data}) -> Data. % HTTP % ----------------------------------------------------------------------------- @@ -110,18 +116,3 @@ close_file(File) -> {error, _} -> {error, eunknown} end. - -% RESCUING -% ----------------------------------------------------------------------------- - -rescue(Callable) -> - try - {ok, Callable()} - catch - error:Error -> - {error, {errored, Error}}; - Error -> - {error, {thrown, Error}}; - exit:Error -> - {error, {exited, Error}} - end. diff --git a/test/autobahn.gleam b/test/autobahn.gleam index 865d3f9..e5a77c7 100644 --- a/test/autobahn.gleam +++ b/test/autobahn.gleam @@ -30,7 +30,7 @@ pub fn main() -> Nil { }) |> ewe.enable_ipv6() |> ewe.bind_all() - |> ewe.listening(port: 8080) + |> ewe.listening(port: 8081) |> ewe.supervised() let assert Ok(_) = diff --git a/test/bench.gleam b/test/bench.gleam new file mode 100644 index 0000000..ab0afd0 --- /dev/null +++ b/test/bench.gleam @@ -0,0 +1,50 @@ +import ewe +import gleam/erlang/process +import gleam/http +import gleam/http/request.{type Request} +import gleam/http/response +import gleam/result + +pub fn main() { + let empty_response = + response.new(200) + |> response.set_body(ewe.Empty) + + let not_found = + response.new(404) + |> response.set_body(ewe.Empty) + + let too_large = + response.new(413) + |> response.set_body(ewe.Empty) + + let assert Ok(_) = + fn(req: Request(ewe.Connection)) { + case req.method, request.path_segments(req) { + http.Get, [] -> empty_response + http.Get, ["user", id] -> + response.new(200) |> response.set_body(ewe.TextData(id)) + http.Post, ["user"] -> { + case ewe.read_body(req, 40_000_000) { + Ok(req) -> { + let content_type = + req + |> request.get_header("content-type") + |> result.unwrap("application/octet-stream") + + response.new(200) + |> response.set_body(ewe.BitsData(req.body)) + |> response.prepend_header("content-type", content_type) + } + Error(_) -> too_large + } + } + _, _ -> not_found + } + } + |> ewe.new + |> ewe.listening(port: 8080) + |> ewe.start() + + process.sleep_forever() +}