diff --git a/.gitignore b/.gitignore index 9c8e0fa..672388c 100644 --- a/.gitignore +++ b/.gitignore @@ -7,3 +7,5 @@ CLAUDE.md autobahn/server # examples /bench +/dev/priv/file_1gb.bin +/dev/priv/file_100kb.bin \ No newline at end of file diff --git a/dev/priv/.gitkeep b/dev/priv/.gitkeep new file mode 100644 index 0000000..e69de29 diff --git a/dev/serve.gleam b/dev/serve.gleam index d375331..d527148 100644 --- a/dev/serve.gleam +++ b/dev/serve.gleam @@ -2,8 +2,13 @@ import ewe import gleam/erlang/process import gleam/http/request import gleam/http/response +import gleam/option +import logging pub fn main() -> Nil { + logging.configure() + logging.set_level(logging.Debug) + let listener_name = process.new_name("listener_name") let connection_factory_name = process.new_name("connection_factory_name") @@ -17,8 +22,41 @@ pub fn main() -> Nil { fn handle_request( request: request.Request(ewe.Connection), ) -> response.Response(ewe.Body) { - echo request + case request.path { + "/hello" -> + response.Response(status: 200, headers: [], body: ewe.Text("Hello, Joe!")) + "/file/small" -> { + // head -c 100K /dev/urandom > file_100kb.bin + let assert Ok(file) = + ewe.file( + "./dev/priv/file_100kb.bin", + offset: option.None, + limit: option.None, + ) + + response.Response( + status: 200, + headers: [#("content-type", "application/octet-stream")], + body: file, + ) + } + "/file/big" -> { + // head -c 1G /dev/urandom > file_1gb.bin + let assert Ok(file) = + ewe.file( + "./dev/priv/file_1gb.bin", + offset: option.None, + limit: option.None, + ) - response.new(200) - |> response.set_body(ewe.Text("Hello, Joe!")) + response.Response( + status: 200, + headers: [#("content-type", "application/octet-stream")], + body: file, + ) + } + _ -> + response.new(404) + |> response.set_body(ewe.Empty) + } } diff --git a/examples/manifest.toml b/examples/manifest.toml index 27c2db5..e195c62 100644 --- a/examples/manifest.toml +++ b/examples/manifest.toml @@ -1,18 +1,20 @@ -# This file was generated by Gleam -# You typically do not need to edit this file +# Do not manually edit this file, it is managed by Gleam. +# +# This file locks the dependency versions used, to make your build +# deterministic and to prevent unexpected versions from being included +# in your application. +# +# You should check this file into your source control repository. packages = [ - { name = "compresso", version = "0.1.0", build_tools = ["gleam"], requirements = ["exception", "gleam_erlang", "gleam_stdlib", "gleam_yielder", "logging"], otp_app = "compresso", source = "hex", outer_checksum = "8BE29A1EDA42F70826ED148EAE40C46BB3FC18E78FE472663DB01DD4A38172D4" }, - { name = "ewe", version = "4.0.0", build_tools = ["gleam"], requirements = ["compresso", "exception", "gleam_erlang", "gleam_http", "gleam_otp", "gleam_stdlib", "glisten", "logging", "websocks"], source = "local", path = ".." }, - { name = "exception", version = "2.1.0", build_tools = ["gleam"], requirements = ["gleam_stdlib"], otp_app = "exception", source = "hex", outer_checksum = "329D269D5C2A314F7364BD2711372B6F2C58FA6F39981572E5CA68624D291F8C" }, + { name = "ewe", version = "5.0.0", build_tools = ["gleam"], requirements = ["gleam_erlang", "gleam_http", "gleam_otp", "gleam_stdlib", "glisten", "logging", "websocks"], source = "local", path = ".." }, { name = "gleam_crypto", version = "1.6.0", build_tools = ["gleam"], requirements = ["gleam_stdlib"], otp_app = "gleam_crypto", source = "hex", outer_checksum = "2DE9E4EF53CF6FEE049D4F765731F7178F7A11AEFAE00EEE63BF7536B354AD3F" }, { name = "gleam_erlang", version = "1.3.0", build_tools = ["gleam"], requirements = ["gleam_stdlib"], otp_app = "gleam_erlang", source = "hex", outer_checksum = "1124AD3AA21143E5AF0FC5CF3D9529F6DB8CA03E43A55711B60B6B7B3874375C" }, { name = "gleam_http", version = "4.3.0", build_tools = ["gleam"], requirements = ["gleam_stdlib"], otp_app = "gleam_http", source = "hex", outer_checksum = "82EA6A717C842456188C190AFB372665EA56CE13D8559BF3B1DD9E40F619EE0C" }, { name = "gleam_otp", version = "1.2.0", build_tools = ["gleam"], requirements = ["gleam_erlang", "gleam_stdlib"], otp_app = "gleam_otp", source = "hex", outer_checksum = "BA6A294E295E428EC1562DC1C11EA7530DCB981E8359134BEABC8493B7B2258E" }, { name = "gleam_stdlib", version = "1.0.2", build_tools = ["gleam"], requirements = [], otp_app = "gleam_stdlib", source = "hex", outer_checksum = "0C5506589DF4C63DF5D6FFBB834562D6865C6C2AEE0019D7B37886BD6D128141" }, - { name = "gleam_yielder", version = "1.1.0", build_tools = ["gleam"], requirements = ["gleam_stdlib"], otp_app = "gleam_yielder", source = "hex", outer_checksum = "8E4E4ECFA7982859F430C57F549200C7749823C106759F4A19A78AEA6687717A" }, { name = "gleeunit", version = "1.10.0", build_tools = ["gleam"], requirements = ["gleam_stdlib"], otp_app = "gleeunit", source = "hex", outer_checksum = "254B697FE72EEAD7BF82E941723918E421317813AC49923EE76A18C788C61E72" }, - { name = "glisten", version = "9.0.1", build_tools = ["gleam"], requirements = ["gleam_erlang", "gleam_otp", "gleam_stdlib", "logging"], otp_app = "glisten", source = "hex", outer_checksum = "7795AA50830656F3A0316A6B26595F893C83272DA901B3405E31339CAA31A10B" }, + { name = "glisten", version = "9.0.1", build_tools = ["gleam"], requirements = ["gleam_erlang", "gleam_otp", "gleam_stdlib", "logging"], source = "git", repo = "https://github.com/vshakitskiy/glisten.git", commit = "63e9a39f7dc35526c2f62128f44f871dab245d15" }, { name = "logging", version = "1.5.0", build_tools = ["gleam"], requirements = ["gleam_stdlib"], otp_app = "logging", source = "hex", outer_checksum = "BC5F18CE5DD9686100229FE5409BDC3DD5C46D5A7DF2F804AD2D8F0DD6C5060E" }, { name = "websocks", version = "3.0.1", build_tools = ["gleam"], requirements = ["gleam_crypto", "gleam_erlang", "gleam_stdlib"], otp_app = "websocks", source = "hex", outer_checksum = "C70340E5B6C3390383ADA17029DCA6F8903863A7AD8CD8E1520EDCC4FE70D6FD" }, ] diff --git a/gleam.toml b/gleam.toml index 6fdb55a..3b99791 100644 --- a/gleam.toml +++ b/gleam.toml @@ -16,7 +16,6 @@ gleam_http = ">= 4.3.0 and < 5.0.0" logging = ">= 1.3.0 and < 2.0.0" gleam_erlang = ">= 1.3.0 and < 2.0.0" websocks = ">= 3.0.0 and < 4.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 df1cb9f..7e5d050 100644 --- a/manifest.toml +++ b/manifest.toml @@ -7,7 +7,6 @@ # You should check this file into your source control repository. packages = [ - { name = "exception", version = "2.1.1", build_tools = ["gleam"], requirements = ["gleam_stdlib"], otp_app = "exception", source = "hex", outer_checksum = "6BDEA95248093599391C3B5DF1835C5C6A86C353C2F99CE539B450E3432FE117" }, { name = "gleam_crypto", version = "1.6.0", build_tools = ["gleam"], requirements = ["gleam_stdlib"], otp_app = "gleam_crypto", source = "hex", outer_checksum = "2DE9E4EF53CF6FEE049D4F765731F7178F7A11AEFAE00EEE63BF7536B354AD3F" }, { name = "gleam_erlang", version = "1.3.0", build_tools = ["gleam"], requirements = ["gleam_stdlib"], otp_app = "gleam_erlang", source = "hex", outer_checksum = "1124AD3AA21143E5AF0FC5CF3D9529F6DB8CA03E43A55711B60B6B7B3874375C" }, { name = "gleam_http", version = "4.3.0", build_tools = ["gleam"], requirements = ["gleam_stdlib"], otp_app = "gleam_http", source = "hex", outer_checksum = "82EA6A717C842456188C190AFB372665EA56CE13D8559BF3B1DD9E40F619EE0C" }, @@ -21,7 +20,6 @@ packages = [ ] [requirements] -exception = { version = ">= 2.1.0 and < 3.0.0" } gleam_erlang = { version = ">= 1.3.0 and < 2.0.0" } gleam_http = { version = ">= 4.3.0 and < 5.0.0" } gleam_httpc = { version = ">= 5.0.0 and < 6.0.0" } diff --git a/src/ewe.gleam b/src/ewe.gleam index d076899..628bfb9 100644 --- a/src/ewe.gleam +++ b/src/ewe.gleam @@ -1,5 +1,7 @@ import ewe/internal/connection +import ewe/internal/file import ewe/internal/handler as handler_ +import gleam/bytes_tree import gleam/erlang/process import gleam/http import gleam/http/request @@ -23,8 +25,10 @@ pub type Connection = connection.Connection pub type Body { - Bytes(BitArray) + Bytes(bytes_tree.BytesTree) Text(String) + Empty + File(connection.File) } pub type IpAddress { @@ -270,15 +274,26 @@ pub fn quiet(builder: Builder) -> Builder { Builder(..builder, on_start: fn(_scheme, _address) { Nil }) } +// Body and connection.Body are structurally identical. +@external(erlang, "gleam_stdlib", "identity") +fn unsafe_to_internal_response( + response: response.Response(Body), +) -> response.Response(connection.Body) + /// Starts the server with the provided configuration. pub fn start( builder: Builder, ) -> Result(actor.Started(supervisor.Supervisor), actor.StartError) { + let handler = fn(request) { + builder.handler(request) + |> unsafe_to_internal_response + } + let pool = glisten.new( listener_name: builder.listener_name, connection_factory_name: builder.connection_factory_name, - on_init: handler_.on_init, + on_init: handler_.on_init(handler), loop: handler_.loop, ) @@ -311,7 +326,7 @@ pub fn start( }) let scheme = case builder.tls { - Some(_) -> http.Https + Some(_config) -> http.Https None -> http.Http } let address = @@ -330,3 +345,30 @@ pub fn supervised( fn() { start(builder) } |> supervision.supervisor } + +pub type FileError { + NotFound + IsDirectory + AccessDenied + UnknownError + InvalidOffset + InvalidLimit +} + +// FileError and file.FileError are structurally identical. +@external(erlang, "gleam_stdlib", "identity") +fn unsafe_from_internal_file_error(error: file.FileError) -> FileError + +/// Prepares a file to be streamed as a response body. `offset` and `limit` in +/// bytes let you serve a byte range from the file. leave either as `None` to +/// serve from the start or through the end. +pub fn file( + path: String, + offset offset: Option(Int), + limit limit: Option(Int), +) -> Result(Body, FileError) { + case file.resolve(path, offset, limit) { + Ok(file) -> Ok(File(file)) + Error(error) -> Error(unsafe_from_internal_file_error(error)) + } +} diff --git a/src/ewe/internal/connection.gleam b/src/ewe/internal/connection.gleam index 75503a8..35fcd8f 100644 --- a/src/ewe/internal/connection.gleam +++ b/src/ewe/internal/connection.gleam @@ -1,3 +1,4 @@ +import gleam/bytes_tree import gleam/erlang/process import glisten/internal/handler import glisten/socket @@ -12,6 +13,17 @@ pub type Connection { ) } +pub type Body { + Bytes(bytes_tree.BytesTree) + Text(String) + Empty + File(File) +} + +pub type File { + FileMetadata(path: String, offset: Int, length: Int) +} + pub type Message { Timeout } diff --git a/src/ewe/internal/file.gleam b/src/ewe/internal/file.gleam new file mode 100644 index 0000000..83235ff --- /dev/null +++ b/src/ewe/internal/file.gleam @@ -0,0 +1,113 @@ +import ewe/internal/connection +import gleam/bytes_tree +import gleam/int +import gleam/option +import gleam/result +import glisten/socket +import glisten/transport + +pub type FileError { + NotFound + IsDirectory + AccessDenied + UnknownError + InvalidOffset + InvalidLimit +} + +/// Stats `path` and validates `offset` and `limit` against its size, producing +/// wire-agnostic metadata. +pub fn resolve( + path: String, + offset: option.Option(Int), + limit: option.Option(Int), +) -> Result(connection.File, FileError) { + use size <- result.try(stat(path)) + let offset = option.unwrap(offset, 0) + let available = size - offset + + case offset >= 0 && offset <= size, limit { + True, option.None -> { + Ok(connection.FileMetadata(path:, offset:, length: available)) + } + True, option.Some(limit) -> { + let length = int.min(limit, available) + Ok(connection.FileMetadata(path:, offset:, length:)) + } + _, option.Some(limit) if limit < 0 -> Error(InvalidLimit) + False, _ -> Error(InvalidOffset) + } +} + +/// Streams a resolved file straight to the wire. TLS payloads must pass +/// through userspace to be encrypted, so they're read and sent in fixed-size +/// chunks instead. +pub fn send( + transport: transport.Transport, + socket: socket.Socket, + file: connection.File, +) -> Result(Nil, socket.SocketReason) { + case file.length { + 0 -> Ok(Nil) + length -> { + use fd <- result.try(open(file.path)) + let result = case transport { + transport.Tcp -> do_sendfile(fd, socket, file.offset, length) + transport.Ssl -> send_chunks(transport, socket, fd, file.offset, length) + } + close(fd) + result + } + } +} + +const chunk_size = 65_536 + +fn send_chunks( + transport: transport.Transport, + socket: socket.Socket, + fd: FileDescriptor, + offset: Int, + remaining: Int, +) -> Result(Nil, socket.SocketReason) { + case remaining { + 0 -> Ok(Nil) + _remaining -> { + let amount = int.min(remaining, chunk_size) + use data <- result.try(pread(fd, offset, amount)) + use Nil <- result.try(transport.send( + transport, + socket, + bytes_tree.from_bit_array(data), + )) + + send_chunks(transport, socket, fd, offset + amount, remaining - amount) + } + } +} + +pub type FileDescriptor + +@external(erlang, "file_ffi", "stat") +fn stat(path: String) -> Result(Int, FileError) + +@external(erlang, "file_ffi", "sendfile") +fn do_sendfile( + fd: FileDescriptor, + socket: socket.Socket, + offset: Int, + bytes: Int, +) -> Result(Nil, socket.SocketReason) + +@external(erlang, "file_ffi", "open") +fn open(path: String) -> Result(FileDescriptor, socket.SocketReason) + +@external(erlang, "file_ffi", "pread") +fn pread( + fd: FileDescriptor, + offset: Int, + length: Int, +) -> Result(BitArray, socket.SocketReason) + +@external(erlang, "file_ffi", "close") +fn close(fd: FileDescriptor) -> Nil diff --git a/src/ewe/internal/file_ffi.erl b/src/ewe/internal/file_ffi.erl new file mode 100644 index 0000000..44446a1 --- /dev/null +++ b/src/ewe/internal/file_ffi.erl @@ -0,0 +1,36 @@ +-module(file_ffi). + +-include_lib("kernel/include/file.hrl"). + +-export([stat/1, sendfile/4, open/1, pread/3, close/1]). + +stat(Path) -> + case file:read_file_info(Path, [raw, {time, posix}]) of + {ok, #file_info{type = directory}} -> {error, is_directory}; + {ok, #file_info{size = Size}} -> {ok, Size}; + {error, enoent} -> {error, not_found}; + {error, eacces} -> {error, access_denied}; + {error, _Reason} -> {error, unknown_error} + end. + +%% the kernel streams the file straight to the socket without passing through +%% userspace memory. +sendfile(Fd, Socket, Offset, Bytes) -> + case file:sendfile(Fd, Socket, Offset, Bytes, []) of + {ok, _Sent} -> {ok, nil}; + {error, Reason} -> {error, Reason} + end. + +open(Path) -> + file:open(Path, [raw, binary, read]). + +pread(Fd, Offset, Length) -> + case file:pread(Fd, Offset, Length) of + {ok, Data} -> {ok, Data}; + eof -> {ok, <<>>}; + {error, Reason} -> {error, Reason} + end. + +close(Fd) -> + file:close(Fd), + nil. diff --git a/src/ewe/internal/handler.gleam b/src/ewe/internal/handler.gleam index 8c0aea3..9909e3c 100644 --- a/src/ewe/internal/handler.gleam +++ b/src/ewe/internal/handler.gleam @@ -2,10 +2,12 @@ import ewe/internal/connection import ewe/internal/http1 import gleam/bit_array import gleam/erlang/process -import gleam/http +import gleam/http/request +import gleam/http/response import gleam/option import glisten -import glisten/transport +import glisten/internal/handler +import logging /// The state of a connection for its entire lifetime. It starts at /// `Initialised`, is classified into `Http1` or `Http2` and then stays in that @@ -13,16 +15,39 @@ import glisten/transport pub type State { /// Accumulates bytes in `buffer` until `sniff_preface` can tell HTTP/1.x /// apart from an HTTP/2 prior-knowledge preface. - Initialised(buffer: BitArray) + Initialised( + handler: fn(request.Request(connection.Connection)) -> + response.Response(connection.Body), + buffer: BitArray, + idle_timer: option.Option(process.Timer), + ) Http1(http1.State) /// No HTTP/2 connection handling exists yet! Http2 } +pub const idle_timeout = 10_000 + pub fn on_init( - _connection: glisten.Connection(connection.Message), -) -> #(State, option.Option(process.Selector(connection.Message))) { - #(Initialised(buffer: <<>>), option.None) + handler: fn(request.Request(connection.Connection)) -> + response.Response(connection.Body), +) { + fn(connection: glisten.Connection(connection.Message)) -> #( + State, + option.Option(process.Selector(connection.Message)), + ) { + let timer = + process.send_after( + connection.subject, + idle_timeout, + handler.User(connection.Timeout), + ) + + #( + Initialised(handler:, buffer: <<>>, idle_timer: option.Some(timer)), + option.None, + ) + } } pub fn loop( @@ -31,35 +56,47 @@ pub fn loop( connection: glisten.Connection(connection.Message), ) -> glisten.Next(State, glisten.Message(connection.Message)) { case message { - glisten.User(connection.Timeout) -> todo as "Timeout not implemented yet" + glisten.User(connection.Timeout) -> { + logging.log(logging.Debug, "Connection idled for too long, closing.") + glisten.stop() + } glisten.Packet(data) -> case state { - Initialised(buffer:) -> classify(<>, connection) - Http1(..) -> todo as "HTTP/1.x connection loop not implemented yet" - Http2(..) -> todo as "HTTP/2 connection handling not implemented yet" - } - } -} + Initialised(handler:, buffer:, idle_timer:) -> { + case idle_timer { + option.Some(timer) -> process.cancel_timer(timer) + option.None -> process.TimerNotFound + } -// Runs the preface sniff exactly once, on the accumulated bytes from -// `Initialised`, and hands off to whichever protocol state it resolves to. -fn classify( - buffer: BitArray, - connection: glisten.Connection(connection.Message), -) -> glisten.Next(State, glisten.Message(connection.Message)) { - case sniff_preface(buffer) { - NeedMoreData -> glisten.continue(Initialised(buffer:)) - Http2Preface(_remaining) -> glisten.continue(Http2) - NotHttp2(buffer:) -> { - let next = - http1.State(buffer:, idle_timer: option.None) - |> http1.handle_message(connection) + let buffer = <> + case sniff_preface(buffer) { + NeedMoreData -> { + let timer = + process.send_after( + connection.subject, + idle_timeout, + handler.User(connection.Timeout), + ) + + Initialised(handler:, buffer:, idle_timer: option.Some(timer)) + |> glisten.continue + } + Http2Preface(_remaining) -> glisten.continue(Http2) + NotHttp2(buffer:) -> { + let next = + http1.State(handler:, buffer:, idle_timer: option.None) + |> http1.handle_message(connection) - case next { - http1.Continue(state) -> glisten.continue(Http1(state)) - http1.Close -> glisten.stop() + case next { + http1.Continue(state) -> glisten.continue(Http1(state)) + http1.Close -> glisten.stop() + } + } + } + } + Http1(..) -> todo as "HTTP/1.x connection loop not implemented yet" + Http2(..) -> todo as "HTTP/2 connection handling not implemented yet" } - } } } diff --git a/src/ewe/internal/http1.gleam b/src/ewe/internal/http1.gleam index 60885c2..3f42a90 100644 --- a/src/ewe/internal/http1.gleam +++ b/src/ewe/internal/http1.gleam @@ -1,19 +1,33 @@ +import ewe/internal/clock import ewe/internal/connection +import ewe/internal/file import gleam/bit_array +import gleam/bytes_tree import gleam/erlang/process import gleam/http import gleam/http/request +import gleam/http/response import gleam/int import gleam/list import gleam/option +import gleam/result +import gleam/string import glisten +import glisten/internal/handler import glisten/transport import logging pub type State { - State(buffer: BitArray, idle_timer: option.Option(process.Timer)) + State( + handler: fn(request.Request(connection.Connection)) -> + response.Response(connection.Body), + buffer: BitArray, + idle_timer: option.Option(process.Timer), + ) } +pub const idle_timeout = 10_000 + pub type Next { Continue(State) Close @@ -23,26 +37,28 @@ pub fn handle_message( state: State, connection: glisten.Connection(connection.Message), ) -> Next { + case state.idle_timer { + option.Some(timer) -> process.cancel_timer(timer) + option.None -> process.TimerNotFound + } + case parse(state.buffer) { - Ok(Complete(head, _metadata, remaining)) -> { + Ok(Complete(head, metadata, remaining)) -> { let scheme = case connection.transport { transport.Tcp -> http.Http transport.Ssl -> http.Https } - let connection = - connection.Http1( - transport: connection.transport, - socket: connection.socket, - self: connection.subject, - buffer: remaining, - ) - let request = request.Request( method: head.method, headers: head.headers, - body: connection, + body: connection.Http1( + transport: connection.transport, + socket: connection.socket, + self: connection.subject, + buffer: remaining, + ), scheme:, host: head.host, port: head.port, @@ -50,21 +66,285 @@ pub fn handle_message( query: head.query, ) - echo request + let response = state.handler(request) + + let sent = case encode_response(response, head.method, metadata) { + Ok(Encoded(bytes:, keep_alive:, file: file_body)) -> { + use Nil <- result.try(transport.send( + connection.transport, + connection.socket, + bytes, + )) + + case file_body { + option.None -> Ok(keep_alive) + option.Some(data) -> + case file.send(connection.transport, connection.socket, data) { + Ok(Nil) -> Ok(keep_alive) + Error(reason) -> Error(reason) + } + } + } + Error(UnsafeHeader(name)) -> { + logging.log( + logging.Error, + "Handler produced an unsafe response header: " <> name, + ) + + use Nil <- result.try(transport.send( + connection.transport, + connection.socket, + internal_server_error(), + )) + + Ok(False) + } + } + + case sent { + Ok(True) -> { + let timer = + process.send_after( + connection.subject, + idle_timeout, + handler.User(connection.Timeout), + ) + + State(..state, buffer: remaining, idle_timer: option.Some(timer)) + |> Continue + } + Ok(False) -> Close + Error(_reason) -> Close + } + } + Ok(Incomplete) -> { + let timer = + process.send_after( + connection.subject, + idle_timeout, + handler.User(connection.Timeout), + ) - Continue(State(..state, buffer: remaining)) + Continue(State(..state, idle_timer: option.Some(timer))) } - Ok(Incomplete) -> Continue(state) Error(error) -> { logging.log( logging.Error, "Failed to parse HTTP/1.x request: " <> error_to_string(error), ) + Close } } } +/// Errors that can occur while turning a handler's `response.Response` into +/// wire bytes. +pub type EncodeError { + /// A header name or value contained a CR, LF, or NUL byte, which would let + /// it inject a new header or terminate the header block early. + UnsafeHeader(name: String) +} + +// Threaded through the pass over `response.headers` with the tree built so far +// and whether the handler's own `connection` header asked to close. +type EncodeState { + EncodeState(tree: bytes_tree.BytesTree, force_close: Bool) +} + +fn initial_encode_state() -> EncodeState { + EncodeState(tree: bytes_tree.new(), force_close: False) +} + +/// The result of encoding a response: `bytes` is ready for a single +/// `transport.send`, and `file`, when present, still needs to be streamed +/// separately since it was never loaded into `bytes`. +pub type Encoded { + Encoded( + bytes: bytes_tree.BytesTree, + keep_alive: Bool, + file: option.Option(connection.File), + ) +} + +/// Builds the response as a `BytesTree`. `Bytes`, `Text` and `Empty` bodies +/// are folded straight into it, so the whole response goes out in a single +/// `transport.send`. +pub fn encode_response( + response: response.Response(connection.Body), + method: http.Method, + metadata: Metadata, +) -> Result(Encoded, EncodeError) { + use state <- result.try(encode_headers(response.headers)) + let length = body_length(response.body) + let keep_alive = metadata.keep_alive && !state.force_close + + let head = + state.tree + |> append_date() + |> append_connection(keep_alive) + |> bytes_tree.prepend(status_line(response.status)) + |> bytes_tree.append_string("content-length: " <> int.to_string(length)) + |> bytes_tree.append(<<"\r\n\r\n":utf8>>) + + case method, response.body { + http.Head, _body -> Ok(Encoded(head, keep_alive, option.None)) + _method, connection.File(data) -> + Ok(Encoded(head, keep_alive, option.Some(data))) + _method, connection.Bytes(tree) -> + Ok(Encoded(bytes_tree.append_tree(head, tree), keep_alive, option.None)) + _method, connection.Text(text) -> + Ok(Encoded(bytes_tree.append_string(head, text), keep_alive, option.None)) + _method, connection.Empty -> Ok(Encoded(head, keep_alive, option.None)) + } +} + +// Builds the header block. +fn encode_headers( + headers: List(#(String, String)), +) -> Result(EncodeState, EncodeError) { + use state, #(name, value) <- list.try_fold(headers, initial_encode_state()) + case name { + "content-length" | "transfer-encoding" | "date" -> Ok(state) + "connection" -> + case find_unsafe_header_byte(value) { + Error(Nil) -> { + let lowered = value |> bit_array.from_string |> lowercase_ascii + let force_close = + state.force_close || has_token(lowered, <<"close":utf8>>) + Ok(EncodeState(..state, force_close:)) + } + Ok(_position) -> Error(UnsafeHeader(name)) + } + _name -> + case find_unsafe_header_byte(name), find_unsafe_header_byte(value) { + Error(Nil), Error(Nil) -> { + let tree = + bytes_tree.append_string(state.tree, name) + |> bytes_tree.append(<<": ":utf8>>) + |> bytes_tree.append_string(value) + |> bytes_tree.append(<<"\r\n":utf8>>) + + Ok(EncodeState(..state, tree:)) + } + _other, _other -> Error(UnsafeHeader(name)) + } + } +} + +// Fills in `Date` from the shared clock (RFC 9110 ยง6.6.1). The clock actor +// refreshes the cached value once a second. +fn append_date(tree: bytes_tree.BytesTree) -> bytes_tree.BytesTree { + bytes_tree.append_string(tree, "date: ") + |> bytes_tree.append(clock.get()) + |> bytes_tree.append(<<"\r\n":utf8>>) +} + +fn append_connection( + tree: bytes_tree.BytesTree, + keep_alive: Bool, +) -> bytes_tree.BytesTree { + let value = case keep_alive { + True -> <<"keep-alive":utf8>> + False -> <<"close":utf8>> + } + + bytes_tree.append(tree, <<"connection: ":utf8>>) + |> bytes_tree.append(value) + |> bytes_tree.append(<<"\r\n":utf8>>) +} + +// Precomputed `HTTP/1.1 NNN Reason\r\n` literals for every standard status +// code. Codes with no registered meaning in this range fall through to the +// general form with an empty reason phrase. +fn status_line(status: Int) -> BitArray { + case status { + 100 -> <<"HTTP/1.1 100 Continue\r\n":utf8>> + 101 -> <<"HTTP/1.1 101 Switching Protocols\r\n":utf8>> + 102 -> <<"HTTP/1.1 102 Processing\r\n":utf8>> + 103 -> <<"HTTP/1.1 103 Early Hints\r\n":utf8>> + 200 -> <<"HTTP/1.1 200 OK\r\n":utf8>> + 201 -> <<"HTTP/1.1 201 Created\r\n":utf8>> + 202 -> <<"HTTP/1.1 202 Accepted\r\n":utf8>> + 203 -> <<"HTTP/1.1 203 Non-Authoritative Information\r\n":utf8>> + 204 -> <<"HTTP/1.1 204 No Content\r\n":utf8>> + 205 -> <<"HTTP/1.1 205 Reset Content\r\n":utf8>> + 206 -> <<"HTTP/1.1 206 Partial Content\r\n":utf8>> + 207 -> <<"HTTP/1.1 207 Multi-Status\r\n":utf8>> + 208 -> <<"HTTP/1.1 208 Already Reported\r\n":utf8>> + 226 -> <<"HTTP/1.1 226 IM Used\r\n":utf8>> + 300 -> <<"HTTP/1.1 300 Multiple Choices\r\n":utf8>> + 301 -> <<"HTTP/1.1 301 Moved Permanently\r\n":utf8>> + 302 -> <<"HTTP/1.1 302 Found\r\n":utf8>> + 303 -> <<"HTTP/1.1 303 See Other\r\n":utf8>> + 304 -> <<"HTTP/1.1 304 Not Modified\r\n":utf8>> + 305 -> <<"HTTP/1.1 305 Use Proxy\r\n":utf8>> + 307 -> <<"HTTP/1.1 307 Temporary Redirect\r\n":utf8>> + 308 -> <<"HTTP/1.1 308 Permanent Redirect\r\n":utf8>> + 400 -> <<"HTTP/1.1 400 Bad Request\r\n":utf8>> + 401 -> <<"HTTP/1.1 401 Unauthorized\r\n":utf8>> + 402 -> <<"HTTP/1.1 402 Payment Required\r\n":utf8>> + 403 -> <<"HTTP/1.1 403 Forbidden\r\n":utf8>> + 404 -> <<"HTTP/1.1 404 Not Found\r\n":utf8>> + 405 -> <<"HTTP/1.1 405 Method Not Allowed\r\n":utf8>> + 406 -> <<"HTTP/1.1 406 Not Acceptable\r\n":utf8>> + 407 -> <<"HTTP/1.1 407 Proxy Authentication Required\r\n":utf8>> + 408 -> <<"HTTP/1.1 408 Request Timeout\r\n":utf8>> + 409 -> <<"HTTP/1.1 409 Conflict\r\n":utf8>> + 410 -> <<"HTTP/1.1 410 Gone\r\n":utf8>> + 411 -> <<"HTTP/1.1 411 Length Required\r\n":utf8>> + 412 -> <<"HTTP/1.1 412 Precondition Failed\r\n":utf8>> + 413 -> <<"HTTP/1.1 413 Content Too Large\r\n":utf8>> + 414 -> <<"HTTP/1.1 414 URI Too Long\r\n":utf8>> + 415 -> <<"HTTP/1.1 415 Unsupported Media Type\r\n":utf8>> + 416 -> <<"HTTP/1.1 416 Range Not Satisfiable\r\n":utf8>> + 417 -> <<"HTTP/1.1 417 Expectation Failed\r\n":utf8>> + 418 -> <<"HTTP/1.1 418 I'm a Teapot\r\n":utf8>> + 421 -> <<"HTTP/1.1 421 Misdirected Request\r\n":utf8>> + 422 -> <<"HTTP/1.1 422 Unprocessable Content\r\n":utf8>> + 423 -> <<"HTTP/1.1 423 Locked\r\n":utf8>> + 424 -> <<"HTTP/1.1 424 Failed Dependency\r\n":utf8>> + 425 -> <<"HTTP/1.1 425 Too Early\r\n":utf8>> + 426 -> <<"HTTP/1.1 426 Upgrade Required\r\n":utf8>> + 428 -> <<"HTTP/1.1 428 Precondition Required\r\n":utf8>> + 429 -> <<"HTTP/1.1 429 Too Many Requests\r\n":utf8>> + 431 -> <<"HTTP/1.1 431 Request Header Fields Too Large\r\n":utf8>> + 451 -> <<"HTTP/1.1 451 Unavailable For Legal Reasons\r\n":utf8>> + 500 -> <<"HTTP/1.1 500 Internal Server Error\r\n":utf8>> + 501 -> <<"HTTP/1.1 501 Not Implemented\r\n":utf8>> + 502 -> <<"HTTP/1.1 502 Bad Gateway\r\n":utf8>> + 503 -> <<"HTTP/1.1 503 Service Unavailable\r\n":utf8>> + 504 -> <<"HTTP/1.1 504 Gateway Timeout\r\n":utf8>> + 505 -> <<"HTTP/1.1 505 HTTP Version Not Supported\r\n":utf8>> + 506 -> <<"HTTP/1.1 506 Variant Also Negotiates\r\n":utf8>> + 507 -> <<"HTTP/1.1 507 Insufficient Storage\r\n":utf8>> + 508 -> <<"HTTP/1.1 508 Loop Detected\r\n":utf8>> + 510 -> <<"HTTP/1.1 510 Not Extended\r\n":utf8>> + 511 -> <<"HTTP/1.1 511 Network Authentication Required\r\n":utf8>> + _other -> <<"HTTP/1.1 ":utf8, int.to_string(status):utf8, " \r\n":utf8>> + } +} + +fn body_length(body: connection.Body) -> Int { + case body { + connection.Bytes(tree) -> bytes_tree.byte_size(tree) + connection.Text(text) -> string.byte_size(text) + connection.Empty -> 0 + connection.File(data) -> data.length + } +} + +fn internal_server_error() -> bytes_tree.BytesTree { + bytes_tree.new() + |> bytes_tree.append(<< + "HTTP/1.1 500 Internal Server Error\r\ndate: ":utf8, + >>) + |> bytes_tree.append(clock.get()) + |> bytes_tree.append(<< + "\r\nconnection: close\r\ncontent-length: 0\r\n\r\n":utf8, + >>) +} + pub type Version { Http10 Http11 @@ -717,6 +997,9 @@ fn find_question(bits: BitArray) -> Result(Int, Nil) @external(erlang, "http1_ffi", "find_close_bracket") fn find_close_bracket(bits: BitArray) -> Result(Int, Nil) +@external(erlang, "http1_ffi", "find_unsafe_header_byte") +fn find_unsafe_header_byte(bits: String) -> Result(Int, Nil) + @external(erlang, "http1_ffi", "split_comma") fn split_comma(bits: BitArray) -> List(BitArray) diff --git a/src/ewe/internal/http1_ffi.erl b/src/ewe/internal/http1_ffi.erl index a924110..c4f7ea5 100644 --- a/src/ewe/internal/http1_ffi.erl +++ b/src/ewe/internal/http1_ffi.erl @@ -8,6 +8,7 @@ find_space/1, find_question/1, find_close_bracket/1, + find_unsafe_header_byte/1, split_comma/1, list_to_bit_array/1, bit_array_to_string/1 @@ -21,6 +22,10 @@ init() -> persistent_term:put({?MODULE, question}, binary:compile_pattern(<<"?">>)), persistent_term:put({?MODULE, comma}, binary:compile_pattern(<<",">>)), persistent_term:put({?MODULE, close_bracket}, binary:compile_pattern(<<"]">>)), + persistent_term:put( + {?MODULE, unsafe_header}, + binary:compile_pattern([<<"\r">>, <<"\n">>, <<0>>]) + ), ok. find_lf(Bin) -> find(Bin, lf). @@ -28,6 +33,7 @@ find_colon(Bin) -> find(Bin, colon). find_space(Bin) -> find(Bin, space). find_question(Bin) -> find(Bin, question). find_close_bracket(Bin) -> find(Bin, close_bracket). +find_unsafe_header_byte(Bin) -> find(Bin, unsafe_header). find(Bin, Key) -> case binary:match(Bin, persistent_term:get({?MODULE, Key})) of