From 1cd942d39ddb8703ba05825f7a4acd00121de9cd Mon Sep 17 00:00:00 2001 From: vshakitskiy Date: Wed, 19 Nov 2025 23:24:17 +0300 Subject: [PATCH] release v2.1.2 --- CHANGELOG.md | 5 + gleam.toml | 3 +- manifest.toml | 1 + src/ewe.gleam | 49 +- src/ewe/internal/buffer.gleam | 35 +- src/ewe/internal/clock.gleam | 42 +- src/ewe/internal/decoder.gleam | 40 +- src/ewe/internal/encoder.gleam | 31 +- src/ewe/internal/exception.gleam | 23 - src/ewe/internal/file.gleam | 27 +- src/ewe/internal/handler.gleam | 72 +- src/ewe/internal/http1.gleam | 1067 +++++++++++------------ src/ewe/internal/http1/handler.gleam | 240 ++--- src/ewe/internal/stream/chunked.gleam | 96 +- src/ewe/internal/stream/sse.gleam | 133 ++- src/ewe/internal/stream/websocket.gleam | 295 +++---- src/ewe_ffi.erl | 4 - 17 files changed, 1021 insertions(+), 1142 deletions(-) delete mode 100644 src/ewe/internal/exception.gleam diff --git a/CHANGELOG.md b/CHANGELOG.md index 96d68b1..9bd666e 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -1,5 +1,10 @@ # Changelog +# v2.1.2 - 19.11.2025 + +- Improve internal codebase +- Remove `exception` module with a package + # v2.1.1 - 14.11.2025 - Bring back active socket option for SSE diff --git a/gleam.toml b/gleam.toml index 1b84046..1d28086 100644 --- a/gleam.toml +++ b/gleam.toml @@ -1,5 +1,5 @@ name = "ewe" -version = "2.1.1" +version = "2.1.2" description = "🐑 a fluffy Gleam web server" target = "erlang" @@ -19,6 +19,7 @@ logging = ">= 1.3.0 and < 2.0.0" gleam_erlang = ">= 1.3.0 and < 2.0.0" compresso = "0.1.0" websocks = ">= 2.0.0 and < 3.0.0" +exception = ">= 2.1.0 and < 3.0.0" [dev-dependencies] gleeunit = ">= 1.0.0 and < 2.0.0" diff --git a/manifest.toml b/manifest.toml index 383780b..6d4d61e 100644 --- a/manifest.toml +++ b/manifest.toml @@ -32,3 +32,4 @@ gleeunit = { version = ">= 1.0.0 and < 2.0.0" } glisten = { version = ">= 8.0.1 and < 9.0.0" } logging = { version = ">= 1.3.0 and < 2.0.0" } websocks = { version = ">= 2.0.0 and < 3.0.0" } +exception = { version = ">= 2.1.0 and < 3.0.0" } diff --git a/src/ewe.gleam b/src/ewe.gleam index a680ad5..b492267 100644 --- a/src/ewe.gleam +++ b/src/ewe.gleam @@ -140,6 +140,12 @@ // IMPORTS // ----------------------------------------------------------------------------- +import ewe/internal/file +import ewe/internal/handler +import ewe/internal/http1 as ewe_http +import ewe/internal/stream/chunked +import ewe/internal/stream/sse +import ewe/internal/stream/websocket import gleam/bit_array import gleam/bytes_tree.{type BytesTree} import gleam/dynamic @@ -155,23 +161,13 @@ import gleam/otp/static_supervisor.{type Supervisor} as supervisor import gleam/otp/supervision import gleam/result import gleam/string_tree.{type StringTree} -import logging - -import websocks - import glisten import glisten/internal/listener import glisten/socket/options as glisten_options import glisten/transport +import logging +import websocks -import ewe/internal/file -import ewe/internal/handler -import ewe/internal/http1 as ewe_http -import ewe/internal/stream/chunked -import ewe/internal/stream/sse -import ewe/internal/stream/websocket - -// ----------------------------------------------------------------------------- // CONNECTION // ----------------------------------------------------------------------------- @@ -182,7 +178,6 @@ import ewe/internal/stream/websocket pub type Connection = ewe_http.Connection -// ----------------------------------------------------------------------------- // IP ADDRESS // ----------------------------------------------------------------------------- @@ -224,7 +219,6 @@ fn ewe_to_glisten_ip(ip: IpAddress) -> glisten.IpAddress { } } -// ----------------------------------------------------------------------------- // INFORMATION // ----------------------------------------------------------------------------- @@ -259,7 +253,6 @@ pub fn get_server_info( SocketAddress(ip: ip_address, port: server_info.port) } -// ----------------------------------------------------------------------------- // RESPONSE // ----------------------------------------------------------------------------- @@ -384,16 +377,9 @@ pub fn file( } } -// ----------------------------------------------------------------------------- // BUILDER // ----------------------------------------------------------------------------- -type Handler = - fn(Request) -> Response - -type OnStart = - fn(http.Scheme, SocketAddress) -> Nil - /// Ewe's server builder. Contains all server configurations. Can be adjusted /// with the following functions: /// - `ewe.bind` @@ -410,12 +396,12 @@ type OnStart = /// pub opaque type Builder { Builder( - handler: Handler, + handler: fn(Request) -> Response, port: Int, interface: String, ipv6: Bool, tls: Option(#(String, String)), - on_start: OnStart, + on_start: fn(http.Scheme, SocketAddress) -> Nil, on_crash: Response, listener_name: process.Name(listener.Message), idle_timeout: Int, @@ -430,11 +416,11 @@ pub opaque type Builder { /// - No ipv6 support /// - No TLS support /// - Default listener name for server information retrieval -/// - on_start: prints `Listening on ://:` +/// - on_start: logs `Listening on ://:` /// - on_crash: empty 500 response /// - idle_timeout: connection is closed after 10_000ms of inactivity /// -pub fn new(handler: Handler) -> Builder { +pub fn new(handler: fn(Request) -> Response) -> Builder { Builder( handler:, port: 8080, @@ -554,7 +540,6 @@ pub fn idle_timeout(builder: Builder, idle_timeout: Int) -> Builder { } } -// ----------------------------------------------------------------------------- // SERVER // ----------------------------------------------------------------------------- @@ -626,7 +611,6 @@ pub fn supervised( supervision.supervisor(fn() { start(builder) }) } -// ----------------------------------------------------------------------------- // REQUEST // ----------------------------------------------------------------------------- @@ -704,11 +688,12 @@ fn consumer_adapter( } } -// ----------------------------------------------------------------------------- -// Chunked Response Body +// CHUNKED RESPONSE // ----------------------------------------------------------------------------- -/// Represents a chunked response body. This type is used to send a chunked response to the client. +/// Represents a chunked response body. This type is used to send a chunked +/// response to the client. +/// pub type ChunkedBody = chunked.ChunkedBody @@ -810,7 +795,6 @@ pub fn send_chunk( chunked.send_chunk(body.transport, body.socket, chunk) } -// ----------------------------------------------------------------------------- // WEBSOCKET // ----------------------------------------------------------------------------- @@ -1083,7 +1067,6 @@ pub fn send_close_frame( |> to_websocket_next() } -// ----------------------------------------------------------------------------- // SERVER-SENT EVENT // ----------------------------------------------------------------------------- diff --git a/src/ewe/internal/buffer.gleam b/src/ewe/internal/buffer.gleam index efacd1d..1a57ba4 100644 --- a/src/ewe/internal/buffer.gleam +++ b/src/ewe/internal/buffer.gleam @@ -1,57 +1,54 @@ -// ----------------------------------------------------------------------------- -// IMPORTS -// ----------------------------------------------------------------------------- import gleam/bit_array import gleam/int -// ----------------------------------------------------------------------------- -// PUBLIC TYPES -// ----------------------------------------------------------------------------- - -// Represents a buffer of data +/// Represents a buffer of data. +/// pub type Buffer { Buffer(data: BitArray, remaining: Int) } -// ----------------------------------------------------------------------------- -// PUBLIC API -// ----------------------------------------------------------------------------- - -/// Creates a new buffer with the given initial data +/// Creates a new buffer with the given initial data. +/// pub fn new(initial: BitArray) -> Buffer { Buffer(initial, 0) } /// Creates a new buffer with the given initial data and remaining bytes to be -/// read +/// read. +/// pub fn new_sized(initial: BitArray, size: Int) { Buffer(initial, size) } -/// Creates a new empty buffer +/// Creates a new empty buffer. +/// pub fn empty() -> Buffer { Buffer(<<>>, 0) } -/// Adjusts remaining bytes to be read +/// Adjusts remaining bytes to be read. +/// pub fn sized(buffer: Buffer, size: Int) -> Buffer { Buffer(buffer.data, size) } -/// Appends the given data to the buffer +/// Appends the given data to the buffer. +/// pub fn append(buffer: Buffer, data: BitArray) -> Buffer { let remaining = int.max(0, buffer.remaining - bit_array.byte_size(data)) Buffer(<>, remaining) } -/// Appends the given data to the buffer with calculated data size +/// Appends the given data to the buffer with calculated data size. +/// pub fn append_size(buffer: Buffer, data: BitArray, size: Int) -> Buffer { let remaining = int.max(0, buffer.remaining - size) Buffer(<>, remaining) } /// Splits the buffer into two parts, the first part is the given number of -/// bytes and the second part is the remaining bytes +/// bytes and the second part is the remaining bytes. +/// pub fn split(buffer: Buffer, bytes: Int) -> #(BitArray, BitArray) { case buffer.data { <> -> #(partition, rest) diff --git a/src/ewe/internal/clock.gleam b/src/ewe/internal/clock.gleam index 9ce7251..74a40e9 100644 --- a/src/ewe/internal/clock.gleam +++ b/src/ewe/internal/clock.gleam @@ -1,6 +1,3 @@ -// ----------------------------------------------------------------------------- -// IMPORTS -// ----------------------------------------------------------------------------- import gleam/erlang/atom import gleam/erlang/process import gleam/int @@ -10,20 +7,14 @@ import gleam/string import gleam/string_tree import logging -// ----------------------------------------------------------------------------- -// INTERNAL TYPES -// ----------------------------------------------------------------------------- - -/// Message that can be sent to the clock actor +/// Message that can be sent to the clock actor. +/// type Message { Tick } -// ----------------------------------------------------------------------------- -// PUBLIC API -// ----------------------------------------------------------------------------- - -/// starts the clock application +/// Starts the clock application. +/// pub fn start(_type, _args) -> Result(process.Pid, actor.StartError) { actor.new_with_initialiser(1000, fn(subject) { init_clock_storage() @@ -48,13 +39,15 @@ pub fn start(_type, _args) -> Result(process.Pid, actor.StartError) { }) } -/// stops the clock application +/// Stops the clock application. +/// pub fn stop(_state) { atom.create("ok") } -/// looks up the HTTP date from the clock storage or calculates a new one if -/// it's not found +/// Looks up the HTTP date from the clock storage or calculates a new one if +/// it's not found. +/// pub fn get_http_date() -> String { case lookup_http_date() { Ok(date) -> date @@ -68,11 +61,8 @@ pub fn get_http_date() -> String { } } -// ----------------------------------------------------------------------------- -// INTERNAL FUNCTIONS -// ----------------------------------------------------------------------------- - -/// Calculates the HTTP date based on the current time +/// Calculates the HTTP date based on the current time. +/// fn calculate_http_date() -> String { let #(weekday, #(year, month, day), #(hour, minute, second)) = now() string_tree.new() @@ -93,7 +83,8 @@ fn calculate_http_date() -> String { |> string_tree.to_string() } -/// Converts a weekday number to a string +/// Converts a weekday number to a string. +/// fn weekday_to_string(weekday: Int) -> String { case weekday { 1 -> "Mon" @@ -108,7 +99,8 @@ fn weekday_to_string(weekday: Int) -> String { } } -/// Converts a month number to a string +/// Converts a month number to a string. +/// fn month_to_string(month: Int) -> String { case month { 1 -> "Jan" @@ -128,10 +120,6 @@ fn month_to_string(month: Int) -> String { } } -// ----------------------------------------------------------------------------- -// EXTERNAL FUNCTIONS -// ----------------------------------------------------------------------------- - @external(erlang, "ewe_ffi", "now") fn now() -> #(Int, #(Int, Int, Int), #(Int, Int, Int)) diff --git a/src/ewe/internal/decoder.gleam b/src/ewe/internal/decoder.gleam index fcde968..d84a1d7 100644 --- a/src/ewe/internal/decoder.gleam +++ b/src/ewe/internal/decoder.gleam @@ -1,31 +1,28 @@ -// ----------------------------------------------------------------------------- -// IMPORTS -// ----------------------------------------------------------------------------- import ewe/internal/buffer import gleam/dynamic import gleam/http import gleam/option -// ----------------------------------------------------------------------------- -// TYPES -// ----------------------------------------------------------------------------- - -// Type of HTTP packet being decoded +/// Type of HTTP packet being decoded. +/// pub type PacketType { HttpBin HttphBin } -// Absolute path in HTTP request +/// Absolute path in HTTP request. +/// pub type AbsPath { AbsPath(BitArray) } -// HTTP version as major and minor numbers +/// HTTP version as major and minor numbers. +/// pub type Version = #(Int, Int) -// HTTP packet structure +/// HTTP packet structure. +/// pub type HttpPacket { HttpRequest(method: BitArray, path: AbsPath, version: Version) HttpHeader(idx: Int, field: BitArray, value: BitArray) @@ -33,17 +30,15 @@ pub type HttpPacket { Http2Upgrade } -// Complete packet with data and remaining bytes +/// Complete packet with data and remaining bytes. +/// pub type Packet { Packet(HttpPacket, rest: BitArray) More(length: option.Option(Int)) } -// ----------------------------------------------------------------------------- -// DECODING -// ----------------------------------------------------------------------------- - -/// Decodes HTTP packets using external FFI implementation +/// Decodes HTTP packets using external FFI implementation. +/// pub fn decode_packet( type_ type_: PacketType, buffer buffer: buffer.Buffer, @@ -51,7 +46,6 @@ pub fn decode_packet( decode_packet_ffi(type_, buffer.data, []) } -/// Decodes HTTP packets using external FFI implementation @external(erlang, "ewe_ffi", "decode_packet") fn decode_packet_ffi( type_ type_: PacketType, @@ -59,7 +53,8 @@ fn decode_packet_ffi( options options: List(a), ) -> Result(Packet, dynamic.Dynamic) -/// Decodes HTTP method from binary data +/// Decodes HTTP method from binary data. +/// pub fn decode_method(method: BitArray) -> Result(http.Method, Nil) { case method { <<"GET">> -> Ok(http.Get) @@ -75,11 +70,8 @@ pub fn decode_method(method: BitArray) -> Result(http.Method, Nil) { } } -// ----------------------------------------------------------------------------- -// MAPPING -// ----------------------------------------------------------------------------- - -/// Maps header field indices to their string names +/// Maps header field indices to their string names. +/// pub fn formatted_field_by_idx(idx: Int) -> Result(String, Nil) { case idx { 0 -> Error(Nil) diff --git a/src/ewe/internal/encoder.gleam b/src/ewe/internal/encoder.gleam index 73f8de9..43e1824 100644 --- a/src/ewe/internal/encoder.gleam +++ b/src/ewe/internal/encoder.gleam @@ -1,16 +1,10 @@ -// ----------------------------------------------------------------------------- -// IMPORTS -// ----------------------------------------------------------------------------- import gleam/bytes_tree import gleam/http/response import gleam/int import gleam/list -// ----------------------------------------------------------------------------- -// PUBLIC API -// ----------------------------------------------------------------------------- - -/// Encodes an HTTP response into bytes +/// Encodes an HTTP response into bytes. +/// pub fn encode_response( response: response.Response(BitArray), ) -> bytes_tree.BytesTree { @@ -20,7 +14,9 @@ pub fn encode_response( |> bytes_tree.append(response.body) } -pub fn setup_encoded_response( +/// Encodes the HTTP status line and headers of an HTTP response. +/// +pub fn encode_response_partially( response: response.Response(a), ) -> bytes_tree.BytesTree { bytes_tree.new() @@ -28,11 +24,8 @@ pub fn setup_encoded_response( |> bytes_tree.append(encode_headers(response.headers)) } -// ----------------------------------------------------------------------------- -// ENCODING -// ----------------------------------------------------------------------------- - -/// Encodes the HTTP status line +/// Encodes the HTTP status line. +/// fn encode_status_line(status: Int) -> BitArray { let status_name = status_to_bit_array(status) let status = int.to_string(status) @@ -40,7 +33,8 @@ fn encode_status_line(status: Int) -> BitArray { <<"HTTP/1.1 ", status:utf8, " ", status_name:bits, "\r\n">> } -/// Encodes HTTP headers into bytes +/// Encodes HTTP headers into bytes. +/// fn encode_headers(headers: List(#(String, String))) -> BitArray { let headers = list.fold(headers, <<>>, fn(acc, headers) { @@ -52,11 +46,8 @@ fn encode_headers(headers: List(#(String, String))) -> BitArray { <> } -// ----------------------------------------------------------------------------- -// MAPPING -// ----------------------------------------------------------------------------- - -/// Maps HTTP status codes to their text descriptions +/// Maps HTTP status codes to their text descriptions. +/// fn status_to_bit_array(status: Int) -> BitArray { case status { 100 -> <<"Continue">> diff --git a/src/ewe/internal/exception.gleam b/src/ewe/internal/exception.gleam deleted file mode 100644 index 6912710..0000000 --- a/src/ewe/internal/exception.gleam +++ /dev/null @@ -1,23 +0,0 @@ -// ----------------------------------------------------------------------------- -// IMPORTS -// ----------------------------------------------------------------------------- -import gleam/dynamic - -// ----------------------------------------------------------------------------- -// TYPES -// ----------------------------------------------------------------------------- - -// Represents an exception happened during the call -pub type Exception { - Errored(dynamic.Dynamic) - Thrown(dynamic.Dynamic) - Exited(dynamic.Dynamic) -} - -// ----------------------------------------------------------------------------- -// RESCUE -// ----------------------------------------------------------------------------- - -/// Rescues an exception that happened during the call -@external(erlang, "ewe_ffi", "rescue") -pub fn rescue(callable: fn() -> a) -> Result(a, Exception) diff --git a/src/ewe/internal/file.gleam b/src/ewe/internal/file.gleam index fa452fa..be7c9a9 100644 --- a/src/ewe/internal/file.gleam +++ b/src/ewe/internal/file.gleam @@ -5,14 +5,12 @@ import glisten import glisten/socket.{type Socket} import glisten/transport.{type Transport} -// ----------------------------------------------------------------------------- -// TYPES -// ----------------------------------------------------------------------------- - -// Represents a reference to a file +/// Reference to a file. +/// pub type IoDevice -// Represents errors that can occur when opening a file +/// Errors that can occur when opening a file. +/// pub type FileError { Enoent Eacces @@ -20,22 +18,21 @@ pub type FileError { Eunknown(dynamic.Dynamic) } -// Represents a file +/// A file. +/// pub type File { File(descriptor: IoDevice, size: Int) } -// Errors that can occur when sending a file +/// Errors that can occur when sending a file. +/// pub type SendError { FileIssue(FileError) SocketIssue(glisten.SocketReason) } -// ----------------------------------------------------------------------------- -// PUBLIC API -// ----------------------------------------------------------------------------- - -/// Sends a file to the client +/// Sends a file to the client. +/// pub fn send( transport: Transport, socket: Socket, @@ -59,10 +56,6 @@ pub fn send( } } -// ----------------------------------------------------------------------------- -// FILE OPERATIONS -// ----------------------------------------------------------------------------- - @external(erlang, "ewe_ffi", "open_file") pub fn open(path: String) -> Result(File, FileError) diff --git a/src/ewe/internal/handler.gleam b/src/ewe/internal/handler.gleam index e1d3a19..b8b255c 100644 --- a/src/ewe/internal/handler.gleam +++ b/src/ewe/internal/handler.gleam @@ -1,48 +1,23 @@ -// ----------------------------------------------------------------------------- -// IMPORTS -// ----------------------------------------------------------------------------- +import ewe/internal/http1.{type Connection, type ResponseBody} as ewe_http +import ewe/internal/http1/handler as http1_handler import gleam/erlang/process import gleam/http/request.{type Request} import gleam/http/response.{type Response} import gleam/option.{type Option, Some} import gleam/otp/actor import gleam/otp/factory_supervisor as factory -import logging - import glisten +import logging -import ewe/internal/http1.{type Connection, type ResponseBody} as ewe_http - -import ewe/internal/http1/handler as http1_handler - -// ----------------------------------------------------------------------------- -// PUBLIC TYPES -// ----------------------------------------------------------------------------- - -// Custom message that can be sent to or received from the Glisten actor -pub type Message { - IdleTimeout -} - -// State of the Glisten actor -pub type GlistenState { - GlistenState(timer: Option(process.Timer), subject: process.Subject(Message)) -} - -pub type State { - Http1(state: http1_handler.State, self: process.Subject(Message)) +/// State of the request handler. +/// +pub type Handler { + Http1(state: http1_handler.Http1Handler, self: process.Subject(Nil)) } -// ----------------------------------------------------------------------------- -// INTERNAL TYPES -// ----------------------------------------------------------------------------- - -// ----------------------------------------------------------------------------- -// PUBLIC API -// ----------------------------------------------------------------------------- - -/// Initializes the Glisten actor's state and selector for custom messages -pub fn init(_) -> #(State, Option(process.Selector(Message))) { +/// Initializes the request handler state. +/// +pub fn init(_) -> #(Handler, Option(process.Selector(Nil))) { let subject = process.new_subject() let selector = process.new_selector() @@ -51,7 +26,8 @@ pub fn init(_) -> #(State, Option(process.Selector(Message))) { #(Http1(http1_handler.init(), self: subject), Some(selector)) } -/// Main handler loop that processes HTTP requests +/// Main loop that processes incoming messages. +/// pub fn loop( handler: fn(Request(Connection)) -> Response(ResponseBody), on_crash: Response(ResponseBody), @@ -59,22 +35,22 @@ pub fn loop( factory.Message(fn() -> Result(actor.Started(Nil), actor.StartError), Nil), ), idle_timeout: Int, -) -> glisten.Loop(State, Message) { +) -> glisten.Loop(Handler, Nil) { fn( - state: State, - msg: glisten.Message(Message), - conn: glisten.Connection(Message), - ) -> glisten.Next(State, glisten.Message(Message)) { + state: Handler, + message: glisten.Message(Nil), + conn: glisten.Connection(Nil), + ) -> glisten.Next(Handler, glisten.Message(Nil)) { let sender = conn.subject let conn = ewe_http.transform_connection(conn, factory_name) - case state, msg { - Http1(state, self), glisten.Packet(msg) -> { + case state, message { + Http1(state, self), glisten.Packet(message) -> { let result = http1_handler.handle_packet( state, - msg, conn, + message, sender, handler, on_crash, @@ -82,9 +58,11 @@ pub fn loop( ) case result { - http1_handler.Continue(new_state) -> - glisten.continue(Http1(new_state, self)) - http1_handler.Http2CleartextUpgrade(_req, _settings) -> { + http1_handler.Continue(state) -> glisten.continue(Http1(state, self)) + http1_handler.Http2Upgrade(http1_handler.OverCleartext( + _req, + _settings, + )) -> { logging.log( logging.Debug, "Received HTTP/2 cleartext upgrade request", diff --git a/src/ewe/internal/http1.gleam b/src/ewe/internal/http1.gleam index 4503a9a..bf11fd0 100644 --- a/src/ewe/internal/http1.gleam +++ b/src/ewe/internal/http1.gleam @@ -1,6 +1,10 @@ -// ----------------------------------------------------------------------------- -// IMPORTS -// ----------------------------------------------------------------------------- +import ewe/internal/buffer.{type Buffer} +import ewe/internal/clock +import ewe/internal/decoder.{ + AbsPath, HttpBin, HttpEoh, HttpHeader, HttpRequest, HttphBin, More, Packet, +} +import ewe/internal/encoder +import ewe/internal/file import gleam/bit_array import gleam/bool import gleam/bytes_tree.{type BytesTree} @@ -20,152 +24,110 @@ import gleam/set.{type Set} import gleam/string import gleam/string_tree.{type StringTree} import gleam/uri - -import websocks - import glisten import glisten/socket.{type Socket} import glisten/transport.{type Transport} +import websocks -import ewe/internal/buffer.{type Buffer} -import ewe/internal/clock -import ewe/internal/decoder.{ - AbsPath, HttpBin, HttpEoh, HttpHeader, HttpRequest, HttphBin, More, Packet, +// Connection +// ----------------------------------------------------------------------------- + +/// Connection to a client. +/// +pub type Connection { + Connection( + transport: Transport, + socket: Socket, + buffer: Buffer, + factory_name: process.Name( + factory.Message(fn() -> Result(actor.Started(Nil), actor.StartError), Nil), + ), + ) } -import ewe/internal/encoder -import ewe/internal/file -// ----------------------------------------------------------------------------- -// PUBLIC TYPES -// ----------------------------------------------------------------------------- +/// Transforms a glisten connection. +/// +pub fn transform_connection( + conn: glisten.Connection(a), + factory_name: process.Name(_), +) -> Connection { + Connection( + transport: conn.transport, + socket: conn.socket, + buffer: buffer.empty(), + factory_name:, + ) +} -// HTTP response body types -pub type ResponseBody { - TextData(String) - BytesData(BytesTree) - BitsData(BitArray) - StringTreeData(StringTree) +/// Reads data from the socket with timeout and size limits. +/// +fn read_from_socket( + transport transport: Transport, + socket socket: Socket, + buffer buffer: Buffer, + on_error on_error: ParseError, +) -> Result(Buffer, ParseError) { + let read_size = int.min(buffer.remaining, max_reading_size) - File(descriptor: file.IoDevice, offset: Int, size: Int) + use data <- try( + transport.receive_timeout(transport, socket, read_size, 5000) + |> replace_error(on_error), + ) - Chunked - Websocket - SSE + let new_buffer = buffer.append_size(buffer, data, read_size) - Empty + case new_buffer.remaining { + 0 -> Ok(new_buffer) + _ -> read_from_socket(transport:, socket:, buffer: new_buffer, on_error:) + } } -// HTTP parsing error types +// HTTP/1.1 +// ----------------------------------------------------------------------------- + +/// Errors that can occur when parsing a request. +/// pub type ParseError { // request line InvalidMethod InvalidTarget InvalidVersion - // headers InvalidHeaders MissingHost DuplicateHost InvalidContentLength - // body InvalidBody BodyTooLarge - // anomalies MalformedRequest PacketDiscard } -// HTTP connection -pub type Connection { - Connection( - transport: Transport, - socket: Socket, - buffer: Buffer, - factory_name: process.Name( - factory.Message(fn() -> Result(actor.Started(Nil), actor.StartError), Nil), - ), - ) -} - -// HTTP version enumeration +/// HTTP version enumeration. +/// pub type HttpVersion { Http10 Http11 } -// WebSocket upgrade error types -pub type UpgradeWebsocketError { - MethodNotGet - MissingConnectionHeader - InvalidConnectionHeader - MissingUpgradeHeader - InvalidUpgradeHeader - MissingWebsocketVersion - MissingWebsocketKey -} - -// Possible results of consuming some amount of data from the request body -pub type Stream { - Consumed(data: BitArray, next: fn(Int) -> Result(Stream, ParseError)) - Done -} - -// Result of parsing a request +/// Result of parsing a request. +/// pub type ParsedRequest { Http1Request(req: Request(Connection), version: HttpVersion) - Http2Upgrade(data: BitArray) - Http2CleartextUpgrade(req: Request(Connection), settings: String) -} - -// ----------------------------------------------------------------------------- -// INTERNAL TYPES -// ----------------------------------------------------------------------------- - -// Function that reads `N` amount of bytes from the request body -type Consumer = - fn(Int) -> Result(Stream, ParseError) - -// Chunked body parsing result -type BodyChunk { - Incomplete - Chunk(BitArray, size: Int, rest: Buffer) - FinalChunk(rest: Buffer) -} - -// State of the chunked body parsing -type ChunkedStreamState { - ChunkedStreamState(data: Buffer, chunk: Buffer, done: Bool) + Http2Upgrade(upgrade: Http2Upgrade) } -// ----------------------------------------------------------------------------- -// CONSTANTS -// ----------------------------------------------------------------------------- - -// 2MB = 2M bytes -const max_reading_size = 2_000_000 - -// ----------------------------------------------------------------------------- -// PUBLIC API -// ----------------------------------------------------------------------------- - -/// Transforms a glisten connection to `Connection` type -pub fn transform_connection( - conn: glisten.Connection(a), - factory_name: process.Name( - factory.Message(fn() -> Result(actor.Started(Nil), actor.StartError), Nil), - ), -) -> Connection { - Connection( - transport: conn.transport, - socket: conn.socket, - buffer: buffer.empty(), - factory_name:, - ) +/// HTTP/2 upgrade options. +/// +pub type Http2Upgrade { + OverCleartext(req: Request(Connection), settings: String) + OverTLS(data: BitArray) } -/// Parses an HTTP request from the given buffer +/// Parses an HTTP request from the given buffer. +/// pub fn parse_request( conn: Connection, buffer: Buffer, @@ -186,6 +148,7 @@ pub fn parse_request( |> try(uri.parse) |> replace_error(InvalidTarget), ) + // Headers use #(headers, rest) <- try(parse_headers( transport, @@ -244,7 +207,7 @@ pub fn parse_request( string.contains(string.lowercase(connection), "upgrade") case is_upgrade { - True -> Ok(Http2CleartextUpgrade(req:, settings:)) + True -> Ok(Http2Upgrade(OverCleartext(req:, settings:))) False -> Ok(Http1Request(req:, version: Http11)) } } @@ -255,7 +218,7 @@ pub fn parse_request( } } Ok(Packet(decoder.Http2Upgrade, <<"\r\nSM\r\n\r\n":utf8, data:bits>>)) -> - Ok(Http2Upgrade(data)) + Ok(Http2Upgrade(OverTLS(data:))) Ok(More(size)) -> { use new_buffer <- try(read_from_socket( transport, @@ -270,7 +233,118 @@ pub fn parse_request( } } -/// Reads the HTTP request body +/// Parses HTTP headers from the buffer. +/// +fn parse_headers( + transport transport: Transport, + socket socket: Socket, + buffer buffer: Buffer, + headers headers: Dict(String, String), +) { + case decoder.decode_packet(HttphBin, buffer) { + Ok(Packet(HttpEoh, rest)) -> Ok(#(headers, rest)) + Ok(Packet(HttpHeader(idx, field, value), rest)) -> { + use field <- try(case decoder.formatted_field_by_idx(idx) { + Ok(field) -> Ok(field) + Error(Nil) -> { + bit_array.to_string(field) + |> result.map(string.lowercase) + |> replace_error(InvalidHeaders) + } + }) + + use value <- try( + validate_field_value(value) |> replace_error(InvalidHeaders), + ) + + let new_buffer = buffer.new(rest) + + use _ <- try(case field { + "host" -> { + case dict.has_key(headers, field) { + True -> Error(DuplicateHost) + False -> Ok(Nil) + } + } + "content-length" -> { + int.parse(value) + |> result.try(fn(value) { + case value < 0 { + True -> Error(Nil) + False -> Ok(Nil) + } + }) + |> result.replace_error(InvalidContentLength) + } + _ -> Ok(Nil) + }) + + insert_header(headers, field, value) + |> parse_headers(transport:, socket:, buffer: new_buffer, headers: _) + } + Ok(More(size)) -> { + let read_size = option.unwrap(size, 0) + + let sized_buffer = buffer.sized(buffer, read_size) + + use new_buffer <- try(read_from_socket( + transport:, + socket:, + buffer: sized_buffer, + on_error: InvalidHeaders, + )) + + parse_headers(transport:, socket:, buffer: new_buffer, headers:) + } + _ -> Error(InvalidHeaders) + } +} + +@external(erlang, "ewe_ffi", "validate_field_value") +fn validate_field_value(value: BitArray) -> Result(String, Nil) + +/// Inserts a header into the headers dictionary. +/// +fn insert_header( + headers: Dict(String, String), + field: String, + value: String, +) -> Dict(String, String) { + case field != "set-cookie" { + True -> + dict.upsert(headers, field, fn(target) { + case target { + option.Some(existing) -> existing <> ", " <> value + option.None -> value + } + }) + False -> dict.insert(headers, available_cookie_key(headers, 0), value) + } +} + +/// Finds an available key for set-cookie headers. +/// +fn available_cookie_key(headers: Dict(String, String), idx: Int) -> String { + let key = case idx { + 0 -> "set-cookie" + n -> "set-cookie-" <> int.to_string(n) + } + + case dict.has_key(headers, key) { + True -> available_cookie_key(headers, idx + 1) + False -> key + } +} + +// Reading Body +// ----------------------------------------------------------------------------- + +/// 2MB (2 million bytes). +/// +const max_reading_size = 2_000_000 + +/// Reads the request body from the socket. +/// pub fn read_body( req: Request(Connection), size_limit: Int, @@ -338,417 +412,198 @@ pub fn read_body( } } -/// Streams the HTTP request body -pub fn stream_body(req: Request(Connection)) { - use _ <- result.try( - handle_continue(req) - |> result.replace_error(InvalidBody), - ) +/// Reads a chunked transfer-encoded body. +/// +fn read_chunked_body( + transport transport: Transport, + socket socket: Socket, + buffer buffer: Buffer, + accumulated_body accumulated_body: BitArray, + body_size_limit body_size_limit: Int, + body_current_size body_current_size: Int, +) -> Result(#(BitArray, Buffer), ParseError) { + use <- bool.guard(body_current_size > body_size_limit, Error(BodyTooLarge)) - case request.get_header(req, "transfer-encoding") { - Ok("chunked") -> { - let state = ChunkedStreamState(buffer.empty(), req.body.buffer, False) - Ok(do_stream_body_chunked(req, state)) - } - _ -> { - let content_length = - request.get_header(req, "content-length") - |> result.try(int.parse) - |> result.unwrap(0) - - let remaining = content_length - bit_array.byte_size(req.body.buffer.data) - let stream_buffer = buffer.sized(req.body.buffer, int.max(0, remaining)) + case parse_body_chunk(buffer) { + Ok(FinalChunk(rest)) -> Ok(#(accumulated_body, rest)) + Ok(Incomplete) -> { + use new_buffer <- try(read_from_socket( + transport:, + socket:, + buffer:, + on_error: InvalidBody, + )) - do_stream_body(req, stream_buffer) - |> Ok + read_chunked_body( + transport:, + socket:, + buffer: new_buffer, + accumulated_body:, + body_size_limit:, + body_current_size:, + ) } - } -} - -/// Upgrades an HTTP connection to WebSocket -pub fn upgrade_websocket( - req: Request(Connection), - transport: Transport, - socket: Socket, -) -> Result(#(List(String), Bool), UpgradeWebsocketError) { - use <- bool.guard(req.method != http.Get, Error(MethodNotGet)) - - let is_upgrade = - request.get_header(req, "connection") - |> result.map(fn(connection) { - string.lowercase(connection) |> string.contains("upgrade") - }) - - use _ <- try(case is_upgrade { - Ok(True) -> Ok(Nil) - Ok(False) -> Error(InvalidConnectionHeader) - Error(_) -> Error(MissingConnectionHeader) - }) - - use _ <- try( - case request.get_header(req, "upgrade") |> result.map(string.lowercase) { - Ok("websocket") -> Ok(Nil) - Ok(_) -> Error(InvalidUpgradeHeader) - Error(_) -> Error(MissingUpgradeHeader) - }, - ) - - use <- bool.guard( - request.get_header(req, "sec-websocket-version") == Error(Nil), - Error(MissingWebsocketVersion), - ) - - use key <- try( - request.get_header(req, "sec-websocket-key") - |> result.replace_error(MissingWebsocketKey), - ) - - let accept_key = websocks.compute_accept(key) - - let extensions = - request.get_header(req, "sec-websocket-extensions") - |> result.map(string.split(_, ";")) - |> result.unwrap([]) - - let permessage_deflate = websocks.has_deflate(extensions) - - let resp = - response.new(101) - |> response.set_body(<<>>) - |> response.set_header("connection", "upgrade") - |> response.set_header("upgrade", "websocket") - |> response.set_header("sec-websocket-accept", accept_key) - |> response.set_header("sec-websocket-version", "13") - - let resp = case permessage_deflate { - True -> - response.set_header( - resp, - "sec-websocket-extensions", - "permessage-deflate", + Ok(Chunk(chunk, size, rest)) -> + read_chunked_body( + transport:, + socket:, + buffer: rest, + accumulated_body: <>, + body_size_limit:, + body_current_size: body_current_size + size, ) - False -> resp + Error(error) -> Error(error) } - - let _ = - encoder.encode_response(resp) - |> transport.send(transport, socket, _) - - Ok(#(extensions, permessage_deflate)) } -/// Appends default headers to HTTP responses -pub fn append_default_headers( - resp: Response(a), - req: Request(Connection), - version: HttpVersion, -) -> Response(a) { - let set_close = request.get_header(req, "connection") == Ok("close") - - let resp = case response.get_header(resp, "date") { - Ok(_) -> resp - Error(Nil) -> response.set_header(resp, "date", clock.get_http_date()) - } +/// Parses a single chunk from the chunked body. +/// +fn parse_body_chunk(buffer: Buffer) -> Result(BodyChunk, ParseError) { + case split(buffer.data, <<"\r\n">>, []) { + [<<"0">>, rest] -> Ok(FinalChunk(buffer.new(rest))) + [chunk_size, rest] -> { + use size <- try( + bit_array.to_string(chunk_size) + |> try(int.base_parse(_, 16)) + |> replace_error(InvalidBody), + ) - case version, set_close { - Http10, _ -> response.set_header(resp, "connection", "close") - _, True -> response.set_header(resp, "connection", "close") - Http11, False -> - case response.get_header(resp, "connection") { - Ok(_) -> resp - Error(Nil) -> response.set_header(resp, "connection", "keep-alive") + case split(rest, <<"\r\n">>, []) { + [chunk, rest] -> { + case bit_array.byte_size(chunk) == size { + True -> Ok(Chunk(chunk, size, buffer.new(rest))) + False -> Error(InvalidBody) + } + } + _ -> Ok(Incomplete) } - } -} - -pub fn set_content_length(resp: Response(BitArray)) -> Response(BitArray) { - case response.get_header(resp, "content-length") { - Ok(_) -> resp - Error(Nil) -> { - let body_size = bit_array.byte_size(resp.body) |> int.to_string - response.set_header(resp, "content-length", body_size) } + _ -> Ok(Incomplete) } } -// ----------------------------------------------------------------------------- -// SOCKET OPERATIONS -// ----------------------------------------------------------------------------- - -/// Reads data from socket with timeout and size limits -fn read_from_socket( - transport transport: Transport, - socket socket: Socket, - buffer buffer: Buffer, - on_error on_error: ParseError, -) -> Result(Buffer, ParseError) { - let read_size = int.min(buffer.remaining, max_reading_size) - - use data <- try( - transport.receive_timeout(transport, socket, read_size, 5000) - |> replace_error(on_error), - ) - - let new_buffer = buffer.append_size(buffer, data, read_size) - - case new_buffer.remaining { - 0 -> Ok(new_buffer) - _ -> read_from_socket(transport:, socket:, buffer: new_buffer, on_error:) - } -} - -// ----------------------------------------------------------------------------- -// HEADERS -// ----------------------------------------------------------------------------- +@external(erlang, "binary", "split") +fn split( + subject: BitArray, + pattern: BitArray, + options: List(atom.Atom), +) -> List(BitArray) -/// Parses HTTP headers from the buffer -fn parse_headers( - transport transport: Transport, - socket socket: Socket, - buffer buffer: Buffer, - headers headers: Dict(String, String), -) { - case decoder.decode_packet(HttphBin, buffer) { - Ok(Packet(HttpEoh, rest)) -> Ok(#(headers, rest)) - Ok(Packet(HttpHeader(idx, field, value), rest)) -> { - use field <- try(case decoder.formatted_field_by_idx(idx) { - Ok(field) -> Ok(field) +/// Handles trailer headers in chunked responses. +fn handle_trailers( + req: Request(BitArray), + set: Set(String), + rest: Buffer, +) -> Request(BitArray) { + case decoder.decode_packet(HttphBin, rest) { + Ok(Packet(HttpEoh, _)) -> req + Ok(Packet(HttpHeader(idx, field, value), header_rest)) -> { + let field_name = case decoder.formatted_field_by_idx(idx) { + Ok(field_name) -> Ok(field_name) Error(Nil) -> { bit_array.to_string(field) |> result.map(string.lowercase) - |> replace_error(InvalidHeaders) } - }) - - use value <- try( - validate_field_value(value) |> replace_error(InvalidHeaders), - ) - - let new_buffer = buffer.new(rest) + } - use _ <- try(case field { - "host" -> { - case dict.has_key(headers, field) { - True -> Error(DuplicateHost) - False -> Ok(Nil) - } - } - "content-length" -> { - int.parse(value) - |> result.try(fn(value) { - case value < 0 { - True -> Error(Nil) - False -> Ok(Nil) + case field_name { + Ok(field_name) -> { + case + set.contains(set, field_name) && !is_forbidden_trailer(field_name) + { + True -> { + case bit_array.to_string(value) { + Ok(value) -> { + request.set_header(req, field_name, value) + |> handle_trailers(set, buffer.new(header_rest)) + } + Error(Nil) -> handle_trailers(req, set, rest) + } } - }) - |> result.replace_error(InvalidContentLength) + False -> handle_trailers(req, set, rest) + } } - _ -> Ok(Nil) - }) - - insert_header(headers, field, value) - |> parse_headers(transport:, socket:, buffer: new_buffer, headers: _) - } - Ok(More(size)) -> { - let read_size = option.unwrap(size, 0) - - let sized_buffer = buffer.sized(buffer, read_size) - - use new_buffer <- try(read_from_socket( - transport:, - socket:, - buffer: sized_buffer, - on_error: InvalidHeaders, - )) - - parse_headers(transport:, socket:, buffer: new_buffer, headers:) + Error(Nil) -> handle_trailers(req, set, rest) + } } - _ -> Error(InvalidHeaders) - } -} - -/// Inserts a header into the headers dictionary -fn insert_header( - headers: Dict(String, String), - field: String, - value: String, -) -> Dict(String, String) { - case field != "set-cookie" { - True -> - dict.upsert(headers, field, fn(target) { - case target { - option.Some(existing) -> existing <> ", " <> value - option.None -> value - } - }) - False -> dict.insert(headers, available_cookie_key(headers, 0), value) + _ -> req } } -/// Finds an available key for set-cookie headers -fn available_cookie_key(headers: Dict(String, String), idx: Int) -> String { - let key = case idx { - 0 -> "set-cookie" - n -> "set-cookie-" <> int.to_string(n) - } - - case dict.has_key(headers, key) { - True -> available_cookie_key(headers, idx + 1) - False -> key +/// Checks if a header field is forbidden in trailers. +fn is_forbidden_trailer(field: String) -> Bool { + case string.lowercase(field) { + "transfer-encoding" + | "content-length" + | "host" + | "cache-control" + | "expect" + | "max-forwards" + | "pragma" + | "range" + | "te" -> True + _ -> False } } +// Streaming Body // ----------------------------------------------------------------------------- -// REQUEST HANDLING -// ----------------------------------------------------------------------------- - -/// Handles 100-continue expectations -pub fn handle_continue(req: Request(Connection)) -> Result(Nil, ParseError) { - let expect = - req.headers - |> list.find(fn(tupple) { - tupple.0 == "expect" && string.lowercase(tupple.1) == "100-continue" - }) - case expect { - Ok(_) -> { - response.new(100) - |> response.set_body(<<>>) - |> encoder.encode_response() - |> transport.send(req.body.transport, req.body.socket, _) - |> result.replace_error(MalformedRequest) - } - Error(Nil) -> Ok(Nil) - } +/// Possible results of consuming some amount of data from the request body. +/// +pub type Stream { + Consumed(data: BitArray, next: fn(Int) -> Result(Stream, ParseError)) + Done } -// ----------------------------------------------------------------------------- -// BODY -// ----------------------------------------------------------------------------- - -/// Creates a consumer function that reads `N` amount of bytes from the request -/// body until it is fully consumed -fn do_stream_body(req: Request(Connection), buffer: Buffer) -> Consumer { - fn(size: Int) { - let buffer_size = bit_array.byte_size(buffer.data) - - case buffer.remaining, buffer_size { - // Request body is fully consumed - 0, 0 -> Ok(Done) - - // Request body is supposed to be fully consumed but there is more data in buffer - 0, _ -> { - let #(data, rest) = buffer.split(buffer, size) - Ok(Consumed(data, do_stream_body(req, buffer.new(rest)))) - } - - // Request body is not fully consumed and there is enough data in buffer to consume `size` bytes - _, buffer_size if buffer_size >= size -> { - let #(data, rest) = buffer.split(buffer, size) - let new_buffer = buffer.new_sized(rest, buffer.remaining) - Ok(Consumed(data, do_stream_body(req, new_buffer))) - } - - // Request body is not fully consumed and there is not enough data in buffer to consume `size` bytes - _, _ -> { - use read_buffer <- try(read_from_socket( - transport: req.body.transport, - socket: req.body.socket, - buffer: buffer.empty(), - on_error: InvalidBody, - )) - - let new_buffer = - buffer.new_sized( - <>, - int.max(0, buffer.remaining - bit_array.byte_size(read_buffer.data)), - ) - - let #(data, rest) = buffer.split(new_buffer, size) - Ok(Consumed(data, do_stream_body(req, buffer.new(rest)))) - } - } - } +/// Chunked body parsing result. +/// +type BodyChunk { + Incomplete + Chunk(BitArray, size: Int, rest: Buffer) + FinalChunk(rest: Buffer) } -// ----------------------------------------------------------------------------- -// CHUNKED BODY -// ----------------------------------------------------------------------------- - -/// Reads chunked transfer-encoded body -fn read_chunked_body( - transport transport: Transport, - socket socket: Socket, - buffer buffer: Buffer, - accumulated_body accumulated_body: BitArray, - body_size_limit body_size_limit: Int, - body_current_size body_current_size: Int, -) -> Result(#(BitArray, Buffer), ParseError) { - use <- bool.guard(body_current_size > body_size_limit, Error(BodyTooLarge)) +/// State of the chunked body parsing. +/// +type ChunkedStreamState { + ChunkedStreamState(data: Buffer, chunk: Buffer, done: Bool) +} - case parse_body_chunk(buffer) { - Ok(FinalChunk(rest)) -> Ok(#(accumulated_body, rest)) - Ok(Incomplete) -> { - use new_buffer <- try(read_from_socket( - transport:, - socket:, - buffer:, - on_error: InvalidBody, - )) +/// Streams the request body from the socket. +/// +pub fn stream_body(req: Request(Connection)) { + use _ <- result.try( + handle_continue(req) + |> result.replace_error(InvalidBody), + ) - read_chunked_body( - transport:, - socket:, - buffer: new_buffer, - accumulated_body:, - body_size_limit:, - body_current_size:, - ) + case request.get_header(req, "transfer-encoding") { + Ok("chunked") -> { + let state = ChunkedStreamState(buffer.empty(), req.body.buffer, False) + Ok(do_stream_body_chunked(req, state)) } - Ok(Chunk(chunk, size, rest)) -> - read_chunked_body( - transport:, - socket:, - buffer: rest, - accumulated_body: <>, - body_size_limit:, - body_current_size: body_current_size + size, - ) - Error(error) -> Error(error) - } -} + _ -> { + let content_length = + request.get_header(req, "content-length") + |> result.try(int.parse) + |> result.unwrap(0) -/// Parses a single chunk from the chunked body -fn parse_body_chunk(buffer: Buffer) -> Result(BodyChunk, ParseError) { - case split(buffer.data, <<"\r\n">>, []) { - [<<"0">>, rest] -> Ok(FinalChunk(buffer.new(rest))) - [chunk_size, rest] -> { - use size <- try( - bit_array.to_string(chunk_size) - |> try(int.base_parse(_, 16)) - |> replace_error(InvalidBody), - ) + let remaining = content_length - bit_array.byte_size(req.body.buffer.data) + let stream_buffer = buffer.sized(req.body.buffer, int.max(0, remaining)) - case split(rest, <<"\r\n">>, []) { - [chunk, rest] -> { - case bit_array.byte_size(chunk) == size { - True -> Ok(Chunk(chunk, size, buffer.new(rest))) - False -> Error(InvalidBody) - } - } - _ -> Ok(Incomplete) - } + do_stream_body(req, stream_buffer) + |> Ok } - _ -> Ok(Incomplete) } } /// Creates a consumer function that reads `N` amount of bytes from the chunked -/// request body until it is fully consumed +/// request body until it is fully consumed. fn do_stream_body_chunked( req: Request(Connection), chunked_stream_state: ChunkedStreamState, -) -> Consumer { +) -> fn(Int) -> Result(Stream, ParseError) { fn(size: Int) { let read_result = read_from_socket_until( @@ -768,7 +623,7 @@ fn do_stream_body_chunked( } } -/// Reads data from socket until `N` amount of bytes are read +/// Reads data from the socket until `N` amount of bytes are read. fn read_from_socket_until( transport transport: Transport, socket socket: Socket, @@ -834,77 +689,217 @@ fn read_from_socket_until( } } -// ----------------------------------------------------------------------------- -// TRAILER HEADERS -// ----------------------------------------------------------------------------- +/// Creates a consumer function that reads `N` amount of bytes from the request +/// body until it is fully consumed. +/// +fn do_stream_body( + req: Request(Connection), + buffer: Buffer, +) -> fn(Int) -> Result(Stream, ParseError) { + fn(size: Int) { + let buffer_size = bit_array.byte_size(buffer.data) -/// Handles trailer headers in chunked responses -fn handle_trailers( - req: Request(BitArray), - set: Set(String), - rest: Buffer, -) -> Request(BitArray) { - case decoder.decode_packet(HttphBin, rest) { - Ok(Packet(HttpEoh, _)) -> req - Ok(Packet(HttpHeader(idx, field, value), header_rest)) -> { - let field_name = case decoder.formatted_field_by_idx(idx) { - Ok(field_name) -> Ok(field_name) - Error(Nil) -> { - bit_array.to_string(field) - |> result.map(string.lowercase) - } + case buffer.remaining, buffer_size { + // Request body is fully consumed + 0, 0 -> Ok(Done) + + // Request body is supposed to be fully consumed but there is more data + // in buffer + 0, _ -> { + let #(data, rest) = buffer.split(buffer, size) + Ok(Consumed(data, do_stream_body(req, buffer.new(rest)))) } - case field_name { - Ok(field_name) -> { - case - set.contains(set, field_name) && !is_forbidden_trailer(field_name) - { - True -> { - case bit_array.to_string(value) { - Ok(value) -> { - request.set_header(req, field_name, value) - |> handle_trailers(set, buffer.new(header_rest)) - } - Error(Nil) -> handle_trailers(req, set, rest) - } - } - False -> handle_trailers(req, set, rest) - } - } - Error(Nil) -> handle_trailers(req, set, rest) + // Request body is not fully consumed and there is enough data in buffer + // to consume `size` bytes + _, buffer_size if buffer_size >= size -> { + let #(data, rest) = buffer.split(buffer, size) + let new_buffer = buffer.new_sized(rest, buffer.remaining) + Ok(Consumed(data, do_stream_body(req, new_buffer))) + } + + // Request body is not fully consumed and there is not enough data in + // buffer to consume `size` bytes + _, _ -> { + use read_buffer <- try(read_from_socket( + transport: req.body.transport, + socket: req.body.socket, + buffer: buffer.empty(), + on_error: InvalidBody, + )) + + let new_buffer = + buffer.new_sized( + <>, + int.max(0, buffer.remaining - bit_array.byte_size(read_buffer.data)), + ) + + let #(data, rest) = buffer.split(new_buffer, size) + Ok(Consumed(data, do_stream_body(req, buffer.new(rest)))) } } - _ -> req } } -/// Checks if a header field is forbidden in trailers -fn is_forbidden_trailer(field: String) -> Bool { - case string.lowercase(field) { - "transfer-encoding" - | "content-length" - | "host" - | "cache-control" - | "expect" - | "max-forwards" - | "pragma" - | "range" - | "te" -> True - _ -> False +// Upgrades +// ----------------------------------------------------------------------------- + +/// Errors that can occur when upgrading a WebSocket connection. +/// +pub type UpgradeWebsocketError { + MethodNotGet + MissingConnectionHeader + InvalidConnectionHeader + MissingUpgradeHeader + InvalidUpgradeHeader + MissingWebsocketVersion + MissingWebsocketKey +} + +/// Upgrades an HTTP connection to WebSocket. +/// +pub fn upgrade_websocket( + req: Request(Connection), + transport: Transport, + socket: Socket, +) -> Result(#(List(String), Bool), UpgradeWebsocketError) { + use <- bool.guard(req.method != http.Get, Error(MethodNotGet)) + + let is_upgrade = + request.get_header(req, "connection") + |> result.map(fn(connection) { + string.lowercase(connection) |> string.contains("upgrade") + }) + + use _ <- try(case is_upgrade { + Ok(True) -> Ok(Nil) + Ok(False) -> Error(InvalidConnectionHeader) + Error(_) -> Error(MissingConnectionHeader) + }) + + use _ <- try( + case request.get_header(req, "upgrade") |> result.map(string.lowercase) { + Ok("websocket") -> Ok(Nil) + Ok(_) -> Error(InvalidUpgradeHeader) + Error(_) -> Error(MissingUpgradeHeader) + }, + ) + + use <- bool.guard( + request.get_header(req, "sec-websocket-version") == Error(Nil), + Error(MissingWebsocketVersion), + ) + + use key <- try( + request.get_header(req, "sec-websocket-key") + |> result.replace_error(MissingWebsocketKey), + ) + + let accept_key = websocks.compute_accept(key) + + let extensions = + request.get_header(req, "sec-websocket-extensions") + |> result.map(string.split(_, ";")) + |> result.unwrap([]) + + let permessage_deflate = websocks.has_deflate(extensions) + + let resp = + response.new(101) + |> response.set_body(<<>>) + |> response.set_header("connection", "upgrade") + |> response.set_header("upgrade", "websocket") + |> response.set_header("sec-websocket-accept", accept_key) + |> response.set_header("sec-websocket-version", "13") + + let resp = case permessage_deflate { + True -> + response.set_header( + resp, + "sec-websocket-extensions", + "permessage-deflate", + ) + False -> resp } + + let _ = + encoder.encode_response(resp) + |> transport.send(transport, socket, _) + + Ok(#(extensions, permessage_deflate)) } -// ----------------------------------------------------------------------------- -// EXTERNAL FFI +// Response // ----------------------------------------------------------------------------- -@external(erlang, "binary", "split") -fn split( - subject: BitArray, - pattern: BitArray, - options: List(atom.Atom), -) -> List(BitArray) +/// Response body variants. +/// +pub type ResponseBody { + TextData(String) + BytesData(BytesTree) + BitsData(BitArray) + StringTreeData(StringTree) + File(descriptor: file.IoDevice, offset: Int, size: Int) + Chunked + Websocket + SSE + Empty +} -@external(erlang, "ewe_ffi", "validate_field_value") -fn validate_field_value(value: BitArray) -> Result(String, Nil) +/// Appends default headers to HTTP responses. +/// +pub fn append_default_headers( + resp: Response(a), + req: Request(Connection), + version: HttpVersion, +) -> Response(a) { + let set_close = request.get_header(req, "connection") == Ok("close") + + let resp = case response.get_header(resp, "date") { + Ok(_) -> resp + Error(Nil) -> response.set_header(resp, "date", clock.get_http_date()) + } + + case version, set_close { + Http10, _ -> response.set_header(resp, "connection", "close") + _, True -> response.set_header(resp, "connection", "close") + Http11, False -> + case response.get_header(resp, "connection") { + Ok(_) -> resp + Error(Nil) -> response.set_header(resp, "connection", "keep-alive") + } + } +} + +/// Sets the content length header if it is not already set. +/// +pub fn set_content_length(resp: Response(BitArray)) -> Response(BitArray) { + case response.get_header(resp, "content-length") { + Ok(_) -> resp + Error(Nil) -> { + let body_size = bit_array.byte_size(resp.body) |> int.to_string + response.set_header(resp, "content-length", body_size) + } + } +} + +/// Handles 100-continue expectations. +/// +pub fn handle_continue(req: Request(Connection)) -> Result(Nil, ParseError) { + let expect = + req.headers + |> list.find(fn(tupple) { + tupple.0 == "expect" && string.lowercase(tupple.1) == "100-continue" + }) + + case expect { + Ok(_) -> { + response.new(100) + |> response.set_body(<<>>) + |> encoder.encode_response() + |> transport.send(req.body.transport, req.body.socket, _) + |> result.replace_error(MalformedRequest) + } + Error(Nil) -> Ok(Nil) + } +} diff --git a/src/ewe/internal/http1/handler.gleam b/src/ewe/internal/http1/handler.gleam index 3b4ace0..4ff4493 100644 --- a/src/ewe/internal/http1/handler.gleam +++ b/src/ewe/internal/http1/handler.gleam @@ -1,8 +1,12 @@ import compresso import ewe/internal/buffer import ewe/internal/encoder -import ewe/internal/exception import ewe/internal/file +import ewe/internal/http1.{ + type Connection, type HttpVersion, type ResponseBody, BitsData, BytesData, + Chunked, Empty, File, SSE, StringTreeData, TextData, Websocket, +} as ewe_http +import exception import gleam/bit_array import gleam/bytes_tree import gleam/erlang/process @@ -13,68 +17,69 @@ import gleam/option.{type Option, None, Some} import gleam/result import gleam/string import gleam/string_tree -import glisten/internal/handler.{Close, Internal} +import glisten +import glisten/internal/handler.{Close, Internal} as glisten_handler import glisten/socket import glisten/transport import logging -import ewe/internal/http1.{ - type Connection, type ResponseBody, BitsData, BytesData, Chunked, Empty, File, - SSE, StringTreeData, TextData, Websocket, -} as ewe_http - -import glisten - -pub type State { - State(idle_timer: Option(process.Timer)) +/// HTTP/1.1 handler state. +/// +pub type Http1Handler { + Http1Handler(idle_timer: Option(process.Timer)) } -pub fn init() -> State { - State(idle_timer: None) +/// Initializes the HTTP/1.1 handler state. +/// +pub fn init() -> Http1Handler { + Http1Handler(idle_timer: None) } -pub type Message { - IdleTimeout +/// Action to take after handling a packet. +/// +pub type Next { + Continue(state: Http1Handler) + Stop + Http2Upgrade(upgrade: Http2Upgrade) } -pub type HandleResult { - Continue(new_state: State) - Stop - Http2Upgrade(data: BitArray) - Http2CleartextUpgrade(req: Request(ewe_http.Connection), settings: String) +/// HTTP/2 upgrade options. +/// +pub type Http2Upgrade { + OverCleartext(request: Request(Connection), settings: String) + OverTLS(data: BitArray) } +/// Handles received glisten packet. +/// pub fn handle_packet( - state: State, - msg: BitArray, - conn: Connection, - subject: process.Subject(handler.Message(_)), - handler: fn(Request(ewe_http.Connection)) -> Response(ewe_http.ResponseBody), - on_crash: Response(ewe_http.ResponseBody), + state: Http1Handler, + connection: Connection, + data: BitArray, + glisten_subject: process.Subject(glisten_handler.Message(_)), + handler: fn(Request(Connection)) -> Response(ResponseBody), + on_crash: Response(ResponseBody), idle_timeout: Int, -) -> HandleResult { +) -> Next { case state.idle_timer { Some(timer) -> process.cancel_timer(timer) None -> process.TimerNotFound } - let parsed = ewe_http.parse_request(conn, buffer.new(msg)) + case ewe_http.parse_request(connection, buffer.new(data)) { + Ok(ewe_http.Http1Request(request, version)) -> { + let call_result = + call(request, version, glisten_subject, handler, on_crash, idle_timeout) - case parsed { - Ok(ewe_http.Http1Request(req, version)) -> { - let called = - call_handler(req, version, subject, handler, on_crash, idle_timeout) - case called { + case call_result { Ok(state) -> Continue(state) Error(Nil) -> Stop } } - - Ok(ewe_http.Http2Upgrade(data)) -> Http2Upgrade(data) - - Ok(ewe_http.Http2CleartextUpgrade(req, settings)) -> - Http2CleartextUpgrade(req, settings) - + Ok(ewe_http.Http2Upgrade(ewe_http.OverTLS(data))) -> + Http2Upgrade(OverTLS(data:)) + Ok(ewe_http.Http2Upgrade(ewe_http.OverCleartext(request, settings))) -> + Http2Upgrade(OverCleartext(request:, settings:)) Error(reason) -> { let status = case reason { ewe_http.InvalidVersion -> 505 @@ -86,23 +91,25 @@ pub fn handle_packet( |> response.set_body(<<>>) |> response.set_header("connection", "close") |> encoder.encode_response() - |> transport.send(conn.transport, conn.socket, _) + |> transport.send(connection.transport, connection.socket, _) Stop } } } -fn call_handler( - req: Request(ewe_http.Connection), - version: ewe_http.HttpVersion, - subject: process.Subject(handler.Message(user_message)), - handler: fn(Request(ewe_http.Connection)) -> Response(ewe_http.ResponseBody), - on_crash: Response(ewe_http.ResponseBody), +/// Takes parsed HTTP request and calls the handler. +/// +fn call( + request: Request(Connection), + version: HttpVersion, + glisten_subject: process.Subject(glisten_handler.Message(_)), + handler: fn(Request(Connection)) -> Response(ResponseBody), + on_crash: Response(ResponseBody), idle_timeout: Int, -) -> Result(State, Nil) { - let resp = case exception.rescue(fn() { handler(req) }) { - Ok(resp) -> resp +) -> Result(Http1Handler, Nil) { + let response = case exception.rescue(fn() { handler(request) }) { + Ok(response) -> response Error(e) -> { logging.log(logging.Error, string.inspect(e)) @@ -110,53 +117,65 @@ fn call_handler( } } - case resp.body { + case response.body { Websocket | SSE | Chunked -> Error(Nil) - File(descriptor, offset, size) -> { - let sent = handle_resp_file(req, version, resp, descriptor, offset, size) + File(descriptor, offset, size) -> + send_file(request, version, response, descriptor, offset, size) + |> on_sent(response, glisten_subject, idle_timeout) - case sent, is_connection_close(resp) { - Ok(Nil), False -> { - let timer = process.send_after(subject, idle_timeout, Internal(Close)) - Ok(State(Some(timer))) - } - _, _ -> Error(Nil) - } - } - _ -> { - let sent = handle_resp_body(req, version, resp, resp.body) - case sent, is_connection_close(resp) { - Ok(Nil), False -> { - let timer = process.send_after(subject, idle_timeout, Internal(Close)) - Ok(State(Some(timer))) - } - _, _ -> Error(Nil) - } + _ -> + send_body(request, version, response) + |> on_sent(response, glisten_subject, idle_timeout) + } +} + +/// Actions to take after response is sent. +/// +fn on_sent( + sent: Result(Nil, glisten.SocketReason), + response: Response(ResponseBody), + glisten_subject: process.Subject(glisten_handler.Message(_)), + idle_timeout: Int, +) -> Result(Http1Handler, Nil) { + case sent, is_connection_close(response) { + Ok(Nil), False -> { + let timer = + process.send_after(glisten_subject, idle_timeout, Internal(Close)) + + Ok(Http1Handler(Some(timer))) } + _, _ -> Error(Nil) } } -/// Handles the file response body and sends it to the client -fn handle_resp_file( - req: Request(Connection), - version: ewe_http.HttpVersion, - resp: Response(ResponseBody), +/// Sends a file to the client. +/// +fn send_file( + request: Request(Connection), + version: HttpVersion, + response: Response(ResponseBody), descriptor: file.IoDevice, offset: Int, size: Int, ) -> Result(Nil, glisten.SocketReason) { - let resp = case response.get_header(resp, "content-length") { - Ok(_) -> resp + let response = case response.get_header(response, "content-length") { + Ok(_) -> response Error(Nil) -> - response.set_header(resp, "content-length", int.to_string(size)) + response.set_header(response, "content-length", int.to_string(size)) } let sent = - ewe_http.append_default_headers(resp, req, version) - |> encoder.setup_encoded_response() - |> transport.send(req.body.transport, req.body.socket, _) + ewe_http.append_default_headers(response, request, version) + |> encoder.encode_response_partially() + |> transport.send(request.body.transport, request.body.socket, _) |> result.try(fn(_) { - file.send(req.body.transport, req.body.socket, descriptor, offset, size) + file.send( + request.body.transport, + request.body.socket, + descriptor, + offset, + size, + ) |> result.replace_error(socket.Badarg) }) @@ -165,14 +184,14 @@ fn handle_resp_file( sent } -/// Handles the response body and sends it to the client -fn handle_resp_body( - req: Request(Connection), - version: ewe_http.HttpVersion, - resp: Response(ResponseBody), - body: ResponseBody, +/// Sends a body to the client. +/// +fn send_body( + request: Request(Connection), + version: HttpVersion, + response: Response(ResponseBody), ) -> Result(Nil, glisten.SocketReason) { - let bits = case body { + let bits = case response.body { TextData(text) -> bit_array.from_string(text) StringTreeData(string_tree) -> string_tree.to_string(string_tree) |> bit_array.from_string @@ -183,15 +202,14 @@ fn handle_resp_body( } let content_length = bit_array.byte_size(bits) - - let resp = case content_length > 1024 { + let response = case content_length > 1024 { True -> - case encode_gzip(req, resp) { + case can_encode_gzip(request, response) { True -> { let compressed = compresso.gzip(bits) let content_length = bit_array.byte_size(compressed) - remove_charset(resp) + remove_charset(response) |> response.set_header("content-encoding", "gzip") |> response.set_header("vary", "Accept-Encoding") |> response.set_header( @@ -201,44 +219,52 @@ fn handle_resp_body( |> response.set_body(compressed) } _ -> - response.set_body(resp, bits) + response.set_body(response, bits) |> response.set_header( "content-length", int.to_string(content_length), ) } False -> - response.set_body(resp, bits) + response.set_body(response, bits) |> response.set_header("content-length", int.to_string(content_length)) } - ewe_http.append_default_headers(resp, req, version) + ewe_http.append_default_headers(response, request, version) |> encoder.encode_response() - |> transport.send(req.body.transport, req.body.socket, _) + |> transport.send(request.body.transport, request.body.socket, _) } -fn encode_gzip(req: Request(ewe_http.Connection), resp: Response(a)) -> Bool { - let req = - request.get_header(req, "accept-encoding") +/// Can the body be encoded to gzip? +/// +fn can_encode_gzip(request: Request(Connection), response: Response(_)) -> Bool { + let accept_encoding = + request.get_header(request, "accept-encoding") |> result.map(string.contains(_, "gzip")) - let resp = response.get_header(resp, "content-encoding") + let content_encoding = response.get_header(response, "content-encoding") - case req, resp { + case accept_encoding, content_encoding { Ok(True), Error(Nil) -> True _, _ -> False } } -fn remove_charset(resp: Response(a)) -> Response(a) { - response.get_header(resp, "content-type") +/// Removes the charset from the content-type header. +/// +fn remove_charset(response: Response(_)) -> Response(_) { + response.get_header(response, "content-type") |> result.try(string.split_once(_, ";")) - |> result.map(fn(parts) { response.set_header(resp, "content-type", parts.0) }) - |> result.unwrap(resp) + |> result.map(fn(parts) { + response.set_header(response, "content-type", parts.0) + }) + |> result.unwrap(response) } -fn is_connection_close(resp: Response(a)) -> Bool { - case response.get_header(resp, "connection") { +/// Is the connection set to close? +/// +fn is_connection_close(response: Response(_)) -> Bool { + case response.get_header(response, "connection") { Ok("close") -> True _ -> False } diff --git a/src/ewe/internal/stream/chunked.gleam b/src/ewe/internal/stream/chunked.gleam index bf1579b..09ed9b0 100644 --- a/src/ewe/internal/stream/chunked.gleam +++ b/src/ewe/internal/stream/chunked.gleam @@ -1,6 +1,3 @@ -// ----------------------------------------------------------------------------- -// IMPORTS -// ----------------------------------------------------------------------------- import ewe/internal/encoder import gleam/bit_array import gleam/bytes_tree @@ -14,27 +11,8 @@ import glisten/socket.{type Socket} import glisten/transport.{type Transport} import logging -// ----------------------------------------------------------------------------- -// PUBLIC TYPES -// ----------------------------------------------------------------------------- - -// Represents a chunked response connection -pub type ChunkedBody { - ChunkedBody(transport: Transport, socket: Socket) -} - -// Represents an instruction on how chunked response should proceed -pub type ChunkedNext(user_state) { - Continue(user_state) - NormalStop - AbnormalStop(reason: String) -} - -// ----------------------------------------------------------------------------- -// PUBLIC API -// ----------------------------------------------------------------------------- - -/// Sends a response for a chunked transfer encoding +/// Sends a response for a chunked transfer encoding. +/// pub fn send_response( resp: Response(a), transport: Transport, @@ -44,12 +22,27 @@ pub fn send_response( Ok("chunked") -> resp _ -> response.set_header(resp, "transfer-encoding", "chunked") } - |> encoder.setup_encoded_response() + |> encoder.encode_response_partially() |> transport.send(transport, socket, _) |> result.replace_error(Nil) } -/// Starts a new chunked response connection +/// Represents a chunked response connection. +/// +pub type ChunkedBody { + ChunkedBody(transport: Transport, socket: Socket) +} + +/// Represents an instruction on how chunked response should proceed. +/// +pub type ChunkedNext(user_state) { + Continue(user_state) + NormalStop + AbnormalStop(reason: String) +} + +/// Starts a new chunked response connection. +/// pub fn start( transport: Transport, socket: Socket, @@ -109,25 +102,21 @@ pub fn start( |> result.map(after_start(_, transport, socket)) } -/// Sends a chunk to the client -pub fn send_chunk( +/// Maps actor's starting value to Nil. +/// +fn after_start( + started: actor.Started(Subject(user_message)), transport: Transport, socket: Socket, - chunk: BitArray, -) -> Result(Nil, glisten.SocketReason) { - bytes_tree.new() - |> bytes_tree.append_string(to_hex_string(bit_array.byte_size(chunk))) - |> bytes_tree.append(<<"\r\n">>) - |> bytes_tree.append(chunk) - |> bytes_tree.append(<<"\r\n">>) - |> transport.send(transport, socket, _) -} +) -> actor.Started(Nil) { + let assert Ok(pid) = process.subject_owner(started.data) + let _ = transport.controlling_process(transport, socket, pid) -// ----------------------------------------------------------------------------- -// INTERNAL FUNCTIONS -// ----------------------------------------------------------------------------- + actor.Started(..started, data: Nil) +} -/// Sends the end marker for chunked transfer encoding +/// Sends the end marker for chunked transfer encoding. +/// fn send_end( transport: Transport, socket: Socket, @@ -135,23 +124,28 @@ fn send_end( transport.send(transport, socket, bytes_tree.from_bit_array(<<"0\r\n\r\n">>)) } -/// Maps actor's starting value to Nil -fn after_start( - started: actor.Started(Subject(user_message)), +/// Sends a chunk to the client. +/// +pub fn send_chunk( 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) + chunk: BitArray, +) -> Result(Nil, glisten.SocketReason) { + bytes_tree.new() + |> bytes_tree.append_string(to_hex_string(bit_array.byte_size(chunk))) + |> bytes_tree.append(<<"\r\n">>) + |> bytes_tree.append(chunk) + |> bytes_tree.append(<<"\r\n">>) + |> transport.send(transport, socket, _) } -/// Converts an integer to a hexadecimal string +/// Converts an integer to a hexadecimal string. +/// fn to_hex_string(integer: Int) -> String { integer_to_list(integer, 16) } -/// Converts an integer to a string in the given base +/// Converts an integer to a string in the given base. +/// @external(erlang, "erlang", "integer_to_list") fn integer_to_list(integer: Int, base: Int) -> String diff --git a/src/ewe/internal/stream/sse.gleam b/src/ewe/internal/stream/sse.gleam index a36b710..af7be7f 100644 --- a/src/ewe/internal/stream/sse.gleam +++ b/src/ewe/internal/stream/sse.gleam @@ -1,6 +1,4 @@ -// ----------------------------------------------------------------------------- -// IMPORTS -// ----------------------------------------------------------------------------- +import ewe/internal/encoder import gleam/bytes_tree import gleam/erlang/atom import gleam/erlang/process.{type Selector, type Subject} @@ -11,62 +9,46 @@ import gleam/option.{type Option} import gleam/otp/actor import gleam/result import gleam/string_tree -import glisten/socket/options.{Active, ActiveMode} - import glisten/socket.{type Socket} +import glisten/socket/options.{Active, ActiveMode} import glisten/transport.{type Transport} -import ewe/internal/encoder - -// ----------------------------------------------------------------------------- -// PUBLIC TYPES -// ----------------------------------------------------------------------------- +/// 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") + |> response.set_header("cache-control", "no-cache") + |> response.set_header("connection", "keep-alive") + |> encoder.encode_response_partially() + |> transport.send(transport, socket, _) + |> result.replace_error(Nil) +} -// Represents a Server-Sent Events connection +/// Represents a Server-Sent Events connection. +/// pub type SSEConnection { SSEConnection(transport: Transport, socket: Socket) } -// Represents a Server-Sent Events event -pub type SSEEvent { - SSEEvent( - event: Option(String), - data: String, - id: Option(String), - retry: Option(Int), - ) -} - -// Represents an instruction on how Server-Sent Events connection should proceed +/// 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 -// Events connection +/// 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 } -// ----------------------------------------------------------------------------- -// PUBLIC API -// ----------------------------------------------------------------------------- - -/// 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") - |> response.set_header("cache-control", "no-cache") - |> response.set_header("connection", "keep-alive") - |> encoder.setup_encoded_response() - |> transport.send(transport, socket, _) - |> result.replace_error(Nil) -} - -/// Starts a new Server-Sent Events connection +/// Starts a new Server-Sent Events connection. +/// pub fn start( transport: Transport, socket: Socket, @@ -110,7 +92,45 @@ pub fn start( |> result.map(after_start(_, transport, socket)) } -/// Sends an event to the client +/// Creates a selector for the Server-Sent Events connection. +/// +fn create_socket_selector( + user_subject: Subject(user_message), +) -> Selector(SSEMessages(user_message)) { + process.new_selector() + |> process.select_map(user_subject, fn(msg) { User(msg) }) + |> process.select_record(atom.create("tcp_closed"), 1, fn(_) { Close }) + |> 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), + data: String, + id: Option(String), + retry: Option(Int), + ) +} + +/// Sends an event to the client. +/// pub fn send_event( transport: Transport, socket: Socket, @@ -145,35 +165,8 @@ pub fn send_event( |> transport.send(transport, socket, _) } -// ----------------------------------------------------------------------------- -// INTERNAL FUNCTIONS -// ----------------------------------------------------------------------------- - -/// Creates a selector for the Server-Sent Events connection -fn create_socket_selector( - user_subject: Subject(user_message), -) -> Selector(SSEMessages(user_message)) { - process.new_selector() - |> process.select_map(user_subject, fn(msg) { User(msg) }) - |> process.select_record(atom.create("tcp_closed"), 1, fn(_) { Close }) - |> process.select_record(atom.create("ssl_closed"), 1, fn(_) { Close }) -} - -/// Formats a field and value for a Server-Sent Events event +/// Formats a field and value for a Server-Sent Events event. +/// fn format(field: String, value: String) { field <> ": " <> value <> "\n" } - -/// 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) -} diff --git a/src/ewe/internal/stream/websocket.gleam b/src/ewe/internal/stream/websocket.gleam index c480267..df1823f 100644 --- a/src/ewe/internal/stream/websocket.gleam +++ b/src/ewe/internal/stream/websocket.gleam @@ -1,6 +1,4 @@ -// ----------------------------------------------------------------------------- -// IMPORTS -// ----------------------------------------------------------------------------- +import exception import gleam/bit_array import gleam/bytes_tree import gleam/dynamic/decode @@ -10,21 +8,14 @@ import gleam/option.{type Option, None, Some} import gleam/otp/actor import gleam/result import gleam/string -import logging - -import websocks - import glisten/socket.{type Socket, type SocketReason} import glisten/socket/options.{ActiveMode, Count} import glisten/transport.{type Transport} +import logging +import websocks -import ewe/internal/exception - -// ----------------------------------------------------------------------------- -// PUBLIC TYPES -// ----------------------------------------------------------------------------- - -// Represents a WebSocket connection +/// Represents a WebSocket connection. +/// pub type WebsocketConnection { WebsocketConnection( transport: Transport, @@ -33,33 +24,34 @@ pub type WebsocketConnection { ) } -// Messages that can be sent to or received from the WebSocket +/// Messages that can be sent to or received from the WebSocket. +/// pub type WebsocketMessage(user_message) { Frame(websocks.Frame) UserMessage(user_message) } -// Control flow for WebSocket message handling +/// Control flow for WebSocket message handling. +/// pub type WebsocketNext(user_state, user_message) { Continue(user_state: user_state, selector: Option(Selector(user_message))) NormalStop AbnormalStop(reason: String) } -// ----------------------------------------------------------------------------- -// INTERNAL TYPES -// ----------------------------------------------------------------------------- - -// 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 @@ -68,104 +60,45 @@ type InternalMessage(user_message) { Invalid } -// ----------------------------------------------------------------------------- -// CALLBACK TYPE ALIASES -// ----------------------------------------------------------------------------- - -// 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 -// ----------------------------------------------------------------------------- -// CONSTANTS -// ----------------------------------------------------------------------------- - +/// Error message for malformed messages. +/// const malformed = "Received malformed message" +/// Error message for crashed WebSocket handler. +/// const crashed = "Crash in websocket handler" +/// Error message for failed PONG frame. +/// const failed_pong = "Failed to send PONG frame" +/// Error message for sending WebSocket message from non-owning process. +/// const non_owning_process = "Sending WebSocket message from non-owning process" -// ----------------------------------------------------------------------------- -// COMPRESSION UTILITIES -// ----------------------------------------------------------------------------- - -// /// Gets the deflate context from the compression option -// fn get_deflate( -// compression: Option(compression.Compression), -// ) -> Option(compression.Context) { -// option.map(compression, fn(compression) { compression.deflate }) -// } - -// /// Gets the inflate context from the compression option -// fn get_inflate( -// compression: Option(compression.Compression), -// ) -> Option(compression.Context) { -// option.map(compression, fn(compression) { compression.inflate }) -// } - -// ----------------------------------------------------------------------------- -// SELECTOR UTILITIES -// ----------------------------------------------------------------------------- - -/// 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 -fn create_socket_selector() -> Selector(InternalMessage(user_message)) { - process.new_selector() - // https://github.com/rawhat/glisten/blob/master/src/glisten/internal/handler.gleam#L121 - |> select_valid_record("tcp") - // https://github.com/rawhat/glisten/blob/master/src/glisten/internal/handler.gleam#L129 - |> select_valid_record("ssl") - // https://github.com/rawhat/glisten/blob/master/src/glisten/internal/handler.gleam#L140 - |> process.select_record(atom.create("tcp_closed"), 1, fn(_) { Close }) - // https://github.com/rawhat/glisten/blob/master/src/glisten/internal/handler.gleam#L137 - |> process.select_record(atom.create("ssl_closed"), 1, fn(_) { Close }) - |> process.select_record(atom.create("tcp_passive"), 1, fn(_) { TcpPassive }) -} - -/// 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) }) -} - -// ----------------------------------------------------------------------------- -// SOCKET UTILITIES -// ----------------------------------------------------------------------------- - +/// Active count for socket. +/// const socket_active_count = 100 -// ----------------------------------------------------------------------------- -// PUBLIC API -// ----------------------------------------------------------------------------- - -/// Starts a new WebSocket connection +/// Starts a new WebSocket connection. +/// pub fn start( transport: Transport, socket: Socket, @@ -231,67 +164,34 @@ pub fn start( |> result.map(after_start(_, transport, socket)) } -/// Sends a frame to the WebSocket -pub fn send_frame( - encoder: fn(BitArray, websocks.Context, Option(BitArray)) -> BitArray, - transport: Transport, - socket: Socket, - context: websocks.Context, - payload: BitArray, -) -> Result(Nil, SocketReason) { - let frame = - exception.rescue(fn() { - encoder(payload, context, option.None) - |> bytes_tree.from_bit_array() - |> transport.send(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)) }) - - case frame { - Ok(frame) -> frame - Error(reason) -> { - logging.log( - logging.Error, - "Frame should be sent from the WebSocket connection, but was sent from different process: " - <> string.inspect(reason), - ) - panic as non_owning_process - } - } + |> result.unwrap(Invalid) + }) } -pub fn send_close_frame( - transport: Transport, - socket: Socket, - code: websocks.CloseReason, -) -> WebsocketNext(user_state, user_message) { - let frame = - exception.rescue(fn() { - websocks.encode_close_frame(code, None) - |> bytes_tree.from_bit_array() - |> transport.send(transport, socket, _) - }) - - case frame { - Ok(Ok(Nil)) -> NormalStop - Ok(Error(reason)) -> - AbnormalStop("Failed to send close frame: " <> string.inspect(reason)) - Error(reason) -> { - logging.log( - logging.Error, - "Frame should be sent from the WebSocket connection, but was sent from different process: " - <> string.inspect(reason), - ) - - panic as non_owning_process - } - } +/// 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_closed"), 1, fn(_) { Close }) + |> process.select_record(atom.create("ssl_closed"), 1, fn(_) { Close }) + |> process.select_record(atom.create("tcp_passive"), 1, fn(_) { TcpPassive }) } -// ----------------------------------------------------------------------------- -// MESSAGE HANDLING -// ----------------------------------------------------------------------------- - -/// Handles incoming packet data, decoding frames and processing them +/// Handles incoming packet data, decoding frames and processing them. +/// fn handle_valid_packet( transport: Transport, socket: Socket, @@ -337,6 +237,8 @@ fn handle_valid_packet( } } +/// Represents the state of the WebSocket connection when resolving frames. +/// type ResolveState(user_state, user_message) { ResolveState( socket: Socket, @@ -346,7 +248,8 @@ type ResolveState(user_state, user_message) { ) } -/// Processes a list of frames sequentially +/// Processes a list of frames sequentially. +/// fn handle_frame( state: ResolveState(user_state, user_message), context: websocks.Context, @@ -424,7 +327,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, @@ -460,11 +364,16 @@ fn handle_user_message( } } -// ----------------------------------------------------------------------------- -// CONNECTION MANAGEMENT -// ----------------------------------------------------------------------------- +/// 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), @@ -486,7 +395,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, @@ -499,3 +409,62 @@ fn after_start( actor.Started(..started, data: Nil) } + +/// Sends a frame to the WebSocket. +/// +pub fn send_frame( + encoder: fn(BitArray, websocks.Context, Option(BitArray)) -> BitArray, + transport: Transport, + socket: Socket, + context: websocks.Context, + payload: BitArray, +) -> Result(Nil, SocketReason) { + let frame = + exception.rescue(fn() { + encoder(payload, context, option.None) + |> bytes_tree.from_bit_array() + |> transport.send(transport, socket, _) + }) + + case frame { + Ok(frame) -> frame + Error(reason) -> { + logging.log( + logging.Error, + "Frame should be sent from the WebSocket connection, but was sent from different process: " + <> string.inspect(reason), + ) + panic as non_owning_process + } + } +} + +/// Sends a close frame to the WebSocket. +/// +pub fn send_close_frame( + transport: Transport, + socket: Socket, + code: websocks.CloseReason, +) -> WebsocketNext(user_state, user_message) { + let frame = + exception.rescue(fn() { + websocks.encode_close_frame(code, None) + |> bytes_tree.from_bit_array() + |> transport.send(transport, socket, _) + }) + + case frame { + Ok(Ok(Nil)) -> NormalStop + Ok(Error(reason)) -> + AbnormalStop("Failed to send close frame: " <> string.inspect(reason)) + Error(reason) -> { + logging.log( + logging.Error, + "Frame should be sent from the WebSocket connection, but was sent from different process: " + <> string.inspect(reason), + ) + + panic as non_owning_process + } + } +} diff --git a/src/ewe_ffi.erl b/src/ewe_ffi.erl index 9e93b89..0381273 100644 --- a/src/ewe_ffi.erl +++ b/src/ewe_ffi.erl @@ -3,7 +3,6 @@ -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]). -% ----------------------------------------------------------------------------- % HTTP % ----------------------------------------------------------------------------- @@ -51,7 +50,6 @@ do_validate_field_value(Value) -> false end. -% ----------------------------------------------------------------------------- % CLOCK % ----------------------------------------------------------------------------- @@ -78,7 +76,6 @@ lookup_http_date() -> {error, nil} end. -% ----------------------------------------------------------------------------- % FILES % ----------------------------------------------------------------------------- @@ -114,7 +111,6 @@ close_file(File) -> {error, eunknown} end. -% ----------------------------------------------------------------------------- % RESCUING % ----------------------------------------------------------------------------- -- 2.51.2