diff --git a/benchmark/ewe@4/.gitignore b/benchmark/ewe@4/.gitignore new file mode 100644 index 0000000..c795b05 --- /dev/null +++ b/benchmark/ewe@4/.gitignore @@ -0,0 +1 @@ +build \ No newline at end of file diff --git a/benchmark/ewe@4/gleam.toml b/benchmark/ewe@4/gleam.toml new file mode 100644 index 0000000..21b8a6b --- /dev/null +++ b/benchmark/ewe@4/gleam.toml @@ -0,0 +1,12 @@ +name = "app" +version = "1.0.0" + +[dependencies] +gleam_stdlib = ">= 1.0.0 and < 2.0.0" +ewe = ">= 4.0.1 and < 5.0.0" +gleam_http = ">= 4.3.0 and < 5.0.0" +gleam_erlang = ">= 1.3.0 and < 2.0.0" +logging = ">= 1.5.0 and < 2.0.0" + +[dev_dependencies] +gleeunit = ">= 1.0.0 and < 2.0.0" diff --git a/benchmark/ewe@4/manifest.toml b/benchmark/ewe@4/manifest.toml new file mode 100644 index 0000000..0889d38 --- /dev/null +++ b/benchmark/ewe@4/manifest.toml @@ -0,0 +1,31 @@ +# 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.1", build_tools = ["gleam"], requirements = ["compresso", "exception", "gleam_erlang", "gleam_http", "gleam_otp", "gleam_stdlib", "glisten", "logging", "websocks"], otp_app = "ewe", source = "hex", outer_checksum = "1A6CBE2E07B7E784F9B9503C94C0F23CDD644AC5EA4E1B957F3A5550541699DB" }, + { 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" }, + { 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.3", build_tools = ["gleam"], requirements = [], otp_app = "gleam_stdlib", source = "hex", outer_checksum = "1F543AFBA5D33DA493E6087F4E4C4F20D899411343512686C98A8ABB2963CF22" }, + { 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.11.0", build_tools = ["gleam"], requirements = ["gleam_stdlib"], otp_app = "gleeunit", source = "hex", outer_checksum = "EC31ABA74256AEA531EDF8169931D775BBB384FED0A8A1BDC4DD9354E3E21826" }, + { 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 = "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" }, +] + +[requirements] +ewe = { version = ">= 4.0.1 and < 5.0.0" } +gleam_erlang = { version = ">= 1.3.0 and < 2.0.0" } +gleam_http = { version = ">= 4.3.0 and < 5.0.0" } +gleam_stdlib = { version = ">= 1.0.0 and < 2.0.0" } +gleeunit = { version = ">= 1.0.0 and < 2.0.0" } +logging = { version = ">= 1.5.0 and < 2.0.0" } diff --git a/benchmark/ewe@4/src/app.gleam b/benchmark/ewe@4/src/app.gleam new file mode 100644 index 0000000..4283d55 --- /dev/null +++ b/benchmark/ewe@4/src/app.gleam @@ -0,0 +1,68 @@ +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 assert Ok(_started) = + ewe.new(handle_request) + |> ewe.listening(port: 3001) + |> ewe.start + + process.sleep_forever() +} + +fn handle_request( + request: request.Request(ewe.Connection), +) -> response.Response(ewe.ResponseBody) { + case request.path { + "/hello" -> + response.new(200) + |> response.set_body(ewe.TextData("Hello, Joe!")) + "/echo" -> + case ewe.read_body(request, 10_000_000) { + Ok(req) -> + response.new(200) + |> response.set_body(ewe.BitsData(req.body)) + Error(_error) -> response.new(400) |> response.set_body(ewe.Empty) + } + "/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.Response( + status: 200, + headers: [#("content-type", "application/octet-stream")], + body: file, + ) + } + _ -> + response.new(404) + |> response.set_body(ewe.Empty) + } +} diff --git a/benchmark/ewe@5/.gitignore b/benchmark/ewe@5/.gitignore new file mode 100644 index 0000000..c795b05 --- /dev/null +++ b/benchmark/ewe@5/.gitignore @@ -0,0 +1 @@ +build \ No newline at end of file diff --git a/benchmark/ewe@5/gleam.toml b/benchmark/ewe@5/gleam.toml new file mode 100644 index 0000000..d4b95c6 --- /dev/null +++ b/benchmark/ewe@5/gleam.toml @@ -0,0 +1,12 @@ +name = "app" +version = "1.0.0" + +[dependencies] +gleam_stdlib = ">= 1.0.0 and < 2.0.0" +gleam_erlang = ">= 1.3.0 and < 2.0.0" +gleam_http = ">= 4.3.0 and < 5.0.0" +logging = ">= 1.5.0 and < 2.0.0" +ewe = { path = "../../" } + +[dev_dependencies] +gleeunit = ">= 1.0.0 and < 2.0.0" diff --git a/benchmark/ewe@5/manifest.toml b/benchmark/ewe@5/manifest.toml new file mode 100644 index 0000000..ec6f173 --- /dev/null +++ b/benchmark/ewe@5/manifest.toml @@ -0,0 +1,28 @@ +# 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 = "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.3", build_tools = ["gleam"], requirements = [], otp_app = "gleam_stdlib", source = "hex", outer_checksum = "1F543AFBA5D33DA493E6087F4E4C4F20D899411343512686C98A8ABB2963CF22" }, + { name = "gleeunit", version = "1.11.0", build_tools = ["gleam"], requirements = ["gleam_stdlib"], otp_app = "gleeunit", source = "hex", outer_checksum = "EC31ABA74256AEA531EDF8169931D775BBB384FED0A8A1BDC4DD9354E3E21826" }, + { 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" }, +] + +[requirements] +ewe = { path = "../../" } +gleam_erlang = { version = ">= 1.3.0 and < 2.0.0" } +gleam_http = { version = ">= 4.3.0 and < 5.0.0" } +gleam_stdlib = { version = ">= 1.0.0 and < 2.0.0" } +gleeunit = { version = ">= 1.0.0 and < 2.0.0" } +logging = { version = ">= 1.5.0 and < 2.0.0" } diff --git a/benchmark/ewe@5/src/app.gleam b/benchmark/ewe@5/src/app.gleam new file mode 100644 index 0000000..9f3e0db --- /dev/null +++ b/benchmark/ewe@5/src/app.gleam @@ -0,0 +1,71 @@ +import ewe +import gleam/bytes_tree +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") + + let assert Ok(_started) = + ewe.new(listener_name:, connection_factory_name:, handler: handle_request) + |> ewe.start + + process.sleep_forever() +} + +fn handle_request( + request: request.Request(ewe.Connection), +) -> response.Response(ewe.Body) { + case request.path { + "/hello" -> + response.new(200) + |> response.set_body(ewe.Text("Hello, Joe!")) + "/echo" -> + case ewe.read_body(request, 10_000_000) { + Ok(req) -> + response.new(200) + |> response.set_body(ewe.Bytes(bytes_tree.from_bit_array(req.body))) + Error(_error) -> response.new(400) |> response.set_body(ewe.Empty) + } + "/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.Response( + status: 200, + headers: [#("content-type", "application/octet-stream")], + body: file, + ) + } + _ -> + response.new(404) + |> response.set_body(ewe.Empty) + } +} diff --git a/benchmark/mist/.gitignore b/benchmark/mist/.gitignore new file mode 100644 index 0000000..599be4e --- /dev/null +++ b/benchmark/mist/.gitignore @@ -0,0 +1,4 @@ +*.beam +*.ez +/build +erl_crash.dump diff --git a/benchmark/mist/gleam.toml b/benchmark/mist/gleam.toml new file mode 100644 index 0000000..115f3e9 --- /dev/null +++ b/benchmark/mist/gleam.toml @@ -0,0 +1,23 @@ +name = "app" +version = "1.0.0" + +# Fill out these fields if you intend to generate HTML documentation or publish +# your project to the Hex package manager. +# +# description = "" +# licences = ["Apache-2.0"] +# repository = { type = "github", user = "", repo = "" } +# links = [{ title = "Website", href = "" }] +# +# For a full reference of all the available options, you can have a look at +# https://gleam.run/writing-gleam/gleam-toml/. + +[dependencies] +gleam_stdlib = ">= 1.0.0 and < 2.0.0" +mist = ">= 6.0.3 and < 7.0.0" +logging = ">= 1.5.0 and < 2.0.0" +gleam_erlang = ">= 1.3.0 and < 2.0.0" +gleam_http = ">= 4.3.0 and < 5.0.0" + +[dev_dependencies] +gleeunit = ">= 1.0.0 and < 2.0.0" diff --git a/benchmark/mist/manifest.toml b/benchmark/mist/manifest.toml new file mode 100644 index 0000000..da6fcc4 --- /dev/null +++ b/benchmark/mist/manifest.toml @@ -0,0 +1,30 @@ +# 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 = "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" }, + { 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.3", build_tools = ["gleam"], requirements = [], otp_app = "gleam_stdlib", source = "hex", outer_checksum = "1F543AFBA5D33DA493E6087F4E4C4F20D899411343512686C98A8ABB2963CF22" }, + { name = "gleeunit", version = "1.11.0", build_tools = ["gleam"], requirements = ["gleam_stdlib"], otp_app = "gleeunit", source = "hex", outer_checksum = "EC31ABA74256AEA531EDF8169931D775BBB384FED0A8A1BDC4DD9354E3E21826" }, + { 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 = "gramps", version = "6.0.1", build_tools = ["gleam"], requirements = ["gleam_crypto", "gleam_erlang", "gleam_http", "gleam_stdlib"], otp_app = "gramps", source = "hex", outer_checksum = "D55636072DEE173F6586A5679D3C02EC7A0DE3F8646B78C351B72908FF223DF7" }, + { name = "hpack_erl", version = "0.3.0", build_tools = ["rebar3"], requirements = [], otp_app = "hpack", source = "hex", outer_checksum = "D6137D7079169D8C485C6962DFE261AF5B9EF60FBC557344511C1E65E3D95FB0" }, + { name = "logging", version = "1.5.0", build_tools = ["gleam"], requirements = ["gleam_stdlib"], otp_app = "logging", source = "hex", outer_checksum = "BC5F18CE5DD9686100229FE5409BDC3DD5C46D5A7DF2F804AD2D8F0DD6C5060E" }, + { name = "mist", version = "6.0.3", build_tools = ["gleam"], requirements = ["exception", "gleam_erlang", "gleam_http", "gleam_otp", "gleam_stdlib", "glisten", "gramps", "hpack_erl", "logging"], otp_app = "mist", source = "hex", outer_checksum = "1B07F321D5FA0CB162D81496F2DE96AEB6EF8980F4F38230A4CC3F849497E020" }, +] + +[requirements] +gleam_erlang = { version = ">= 1.3.0 and < 2.0.0" } +gleam_http = { version = ">= 4.3.0 and < 5.0.0" } +gleam_stdlib = { version = ">= 1.0.0 and < 2.0.0" } +gleeunit = { version = ">= 1.0.0 and < 2.0.0" } +logging = { version = ">= 1.5.0 and < 2.0.0" } +mist = { version = ">= 6.0.3 and < 7.0.0" } diff --git a/benchmark/mist/src/app.gleam b/benchmark/mist/src/app.gleam new file mode 100644 index 0000000..18f8f06 --- /dev/null +++ b/benchmark/mist/src/app.gleam @@ -0,0 +1,80 @@ +import gleam/bytes_tree +import gleam/erlang/process +import gleam/http +import gleam/http/request +import gleam/http/response +import gleam/option +import gleam/string +import logging +import mist + +pub fn main() -> Nil { + logging.configure() + logging.set_level(logging.Debug) + + let assert Ok(_started) = + mist.new(handle_request) + |> mist.port(3002) + |> mist.start + + process.sleep_forever() +} + +fn handle_request( + request: request.Request(mist.Connection), +) -> response.Response(mist.ResponseData) { + case request.path { + "/hello" -> + response.new(200) + |> response.set_body(mist.Bytes(bytes_tree.from_string("Hello, Joe!"))) + "/whoami" -> + response.new(200) + |> response.set_body(mist.Bytes(bytes_tree.from_string( + "method=" <> method_to_string(request.method), + ))) + "/echo" -> + case mist.read_body(request, 10_000_000) { + Ok(req) -> + response.new(200) + |> response.set_body(mist.Bytes(bytes_tree.from_bit_array(req.body))) + Error(_error) -> + response.new(400) |> response.set_body(mist.Bytes(bytes_tree.new())) + } + "/file/small" -> { + // head -c 100K /dev/urandom > file_100kb.bin + let assert Ok(file) = + mist.send_file( + "./dev/priv/file_100kb.bin", + offset: 0, + 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) = + mist.send_file("./dev/priv/file_1gb.bin", offset: 0, limit: option.None) + + response.Response( + status: 200, + headers: [#("content-type", "application/octet-stream")], + body: file, + ) + } + _ -> + response.new(404) + |> response.set_body(mist.Bytes(bytes_tree.new())) + } +} + +fn method_to_string(method: http.Method) -> String { + case method { + http.Other(name) -> "Other(" <> name <> ")" + _other -> string.inspect(method) + } +} diff --git a/dev/priv/.gitkeep b/dev/priv/.gitkeep deleted file mode 100644 index e69de29..0000000 diff --git a/dev/priv/chunked.lua b/dev/priv/chunked.lua new file mode 100644 index 0000000..6941185 --- /dev/null +++ b/dev/priv/chunked.lua @@ -0,0 +1,22 @@ +wrk.method = "POST" + +local chunks = { + string.rep("a", 256), + string.rep("b", 256), + string.rep("c", 256), +} + +local encoded = "" +for _, chunk in ipairs(chunks) do + encoded = encoded .. string.format("%x", #chunk) .. "\r\n" .. chunk .. "\r\n" +end +encoded = encoded .. "0\r\n\r\n" + +request = function() + return "POST /echo HTTP/1.1\r\n" .. + "Host: 127.0.0.1\r\n" .. + "Transfer-Encoding: chunked\r\n" .. + "Content-Type: application/octet-stream\r\n" .. + "Connection: keep-alive\r\n\r\n" .. + encoded +end \ No newline at end of file diff --git a/dev/serve.gleam b/dev/serve.gleam index d527148..9f3e0db 100644 --- a/dev/serve.gleam +++ b/dev/serve.gleam @@ -1,4 +1,5 @@ import ewe +import gleam/bytes_tree import gleam/erlang/process import gleam/http/request import gleam/http/response @@ -24,7 +25,15 @@ fn handle_request( ) -> response.Response(ewe.Body) { case request.path { "/hello" -> - response.Response(status: 200, headers: [], body: ewe.Text("Hello, Joe!")) + response.new(200) + |> response.set_body(ewe.Text("Hello, Joe!")) + "/echo" -> + case ewe.read_body(request, 10_000_000) { + Ok(req) -> + response.new(200) + |> response.set_body(ewe.Bytes(bytes_tree.from_bit_array(req.body))) + Error(_error) -> response.new(400) |> response.set_body(ewe.Empty) + } "/file/small" -> { // head -c 100K /dev/urandom > file_100kb.bin let assert Ok(file) = diff --git a/src/ewe.gleam b/src/ewe.gleam index 628bfb9..db04422 100644 --- a/src/ewe.gleam +++ b/src/ewe.gleam @@ -1,6 +1,7 @@ import ewe/internal/connection import ewe/internal/file import ewe/internal/handler as handler_ +import ewe/internal/http1 import gleam/bytes_tree import gleam/erlang/process import gleam/http @@ -8,6 +9,7 @@ import gleam/http/request import gleam/http/response import gleam/int import gleam/io +import gleam/list import gleam/option.{type Option, None, Some} import gleam/otp/actor import gleam/otp/factory_supervisor as factory @@ -76,14 +78,19 @@ fn convert_socket_address(address: glisten.SocketAddress) -> SocketAddress { /// Retrieves the client's socket address from the connection. Returns error if /// the socket information is unavailable. pub fn get_client_info(connection: Connection) -> Result(SocketAddress, Nil) { - let peername = transport.peername(connection.transport, connection.socket) - use info <- result.map(over: peername) - - case info { - socket.TcpSockName(ip_address:, port:) -> - unsafe_from_internal_options_ip_address(ip_address) - |> TcpSocketAddress(port:) - socket.UnixSockName(path:) -> UnixSocketAddress(path:) + case connection { + connection.Http1(transport:, socket:, ..) -> { + let peername = transport.peername(transport, socket) + use info <- result.map(over: peername) + + case info { + socket.TcpSockName(ip_address:, port:) -> + unsafe_from_internal_options_ip_address(ip_address) + |> TcpSocketAddress(port:) + socket.UnixSockName(path:) -> UnixSocketAddress(path:) + } + } + connection.Http2 -> todo as "HTTP/2 is not implemented yet!" } } @@ -372,3 +379,36 @@ pub fn file( Error(error) -> Error(unsafe_from_internal_file_error(error)) } } + +pub type BodyError { + /// The declared body is bigger than the `limit` passed to `read_body`. + BodyTooLarge + /// The body couldn't be fully read: the connection dropped, timed out, or + /// the chunked framing was malformed. + InvalidBody +} + +// BodyError and http1.BodyError are structurally identical. +@external(erlang, "gleam_stdlib", "identity") +fn unsafe_from_internal_body_error(error: http1.BodyError) -> BodyError + +/// Reads the entire request body into memory, up to `limit` bytes. For a +/// chunked request, any trailer fields are appended to the returned request's +/// `headers`. +pub fn read_body( + req: request.Request(Connection), + limit limit: Int, +) -> Result(request.Request(BitArray), BodyError) { + case req.body { + connection.Http1(..) -> { + use #(body, trailers) <- result.try( + http1.read_body(http1.unsafe_to_http1_connection(req.body), limit) + |> result.map_error(unsafe_from_internal_body_error), + ) + + request.Request(..req, headers: list.append(req.headers, trailers), body:) + |> Ok + } + connection.Http2 -> todo as "HTTP/2 is not implemented yet!" + } +} diff --git a/src/ewe/internal/connection.gleam b/src/ewe/internal/connection.gleam index 35fcd8f..385a290 100644 --- a/src/ewe/internal/connection.gleam +++ b/src/ewe/internal/connection.gleam @@ -10,7 +10,17 @@ pub type Connection { socket: socket.Socket, self: process.Subject(handler.Message(Message)), buffer: BitArray, + framing: Framing, ) + Http2 +} + +/// How to find the end of a request body on the wire, derived once from +/// `Content-Length`/`Transfer-Encoding` at parse time. +pub type Framing { + Fixed(length: Int) + Chunked + NoBody } pub type Body { @@ -26,4 +36,10 @@ pub type File { pub type Message { Timeout + /// Sent by `http1.read_body` to itself once the request body has been + /// fully consumed, carrying whatever bytes came after it. + BodyDrained(leftover: BitArray) + /// Sent by `http1.read_body` to itself when it gave up on the body + /// part-way through, leaving the connection in an unknown position. + BodyAbandoned } diff --git a/src/ewe/internal/handler.gleam b/src/ewe/internal/handler.gleam index 9909e3c..f361479 100644 --- a/src/ewe/internal/handler.gleam +++ b/src/ewe/internal/handler.gleam @@ -60,6 +60,16 @@ pub fn loop( logging.log(logging.Debug, "Connection idled for too long, closing.") glisten.stop() } + glisten.User(connection.BodyDrained(..)) + | glisten.User(connection.BodyAbandoned) -> { + // http1.handle_message always drains these itself before returning + logging.log( + logging.Critical, + "Web server loop received a message that should not be reached to the handler! That means there is a bug somewhere within the implementation. Please open an issue addressing the alert, thank you! https://github.com/vshakitskiy/ewe/issues/new", + ) + + glisten.continue(state) + } glisten.Packet(data) -> case state { Initialised(handler:, buffer:, idle_timer:) -> { @@ -94,7 +104,16 @@ pub fn loop( } } } - Http1(..) -> todo as "HTTP/1.x connection loop not implemented yet" + Http1(state) -> { + let next = + http1.State(..state, buffer: <>) + |> http1.handle_message(connection) + + case next { + http1.Continue(state) -> glisten.continue(Http1(state)) + http1.Close -> glisten.stop() + } + } 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 3f42a90..4a5e36b 100644 --- a/src/ewe/internal/http1.gleam +++ b/src/ewe/internal/http1.gleam @@ -14,6 +14,7 @@ import gleam/result import gleam/string import glisten import glisten/internal/handler +import glisten/socket import glisten/transport import logging @@ -49,16 +50,20 @@ pub fn handle_message( transport.Ssl -> http.Https } + let body_connection = + connection.Http1( + transport: connection.transport, + socket: connection.socket, + self: connection.subject, + buffer: remaining, + framing: metadata.framing, + ) + let request = request.Request( method: head.method, headers: head.headers, - body: connection.Http1( - transport: connection.transport, - socket: connection.socket, - self: connection.subject, - buffer: remaining, - ), + body: body_connection, scheme:, host: head.host, port: head.port, @@ -68,6 +73,11 @@ pub fn handle_message( let response = state.handler(request) + let #(buffer, body_drained) = + resolve_body(unsafe_to_http1_connection(body_connection)) + let metadata = + Metadata(..metadata, keep_alive: metadata.keep_alive && body_drained) + let sent = case encode_response(response, head.method, metadata) { Ok(Encoded(bytes:, keep_alive:, file: file_body)) -> { use Nil <- result.try(transport.send( @@ -110,7 +120,7 @@ pub fn handle_message( handler.User(connection.Timeout), ) - State(..state, buffer: remaining, idle_timer: option.Some(timer)) + State(..state, buffer:, idle_timer: option.Some(timer)) |> Continue } Ok(False) -> Close @@ -138,6 +148,251 @@ pub fn handle_message( } } +/// Errors from consuming a request body via `read_body`. +pub type BodyError { + /// A `Fixed` body's declared length alone exceeds `limit`, so nothing was + /// read; or a `Chunked` body's running total went over `limit` mid-stream, + /// so the connection is no longer reusable. + BodyTooLarge + /// The body couldn't be fully read off the socket: closed, timed out, or + /// malformed chunked framing. + InvalidBody +} + +pub type Connection { + Http1( + transport: transport.Transport, + socket: socket.Socket, + self: process.Subject(handler.Message(connection.Message)), + buffer: BitArray, + framing: connection.Framing, + ) +} + +@external(erlang, "gleam_stdlib", "identity") +pub fn unsafe_to_http1_connection(conn: connection.Connection) -> Connection + +const body_read_timeout = 10_000 + +// Cap for auto draining a body the handler never read, so the connection can +// still be reused for the next request. Bodies bigger than this force a +// close instead of an unbounded drain. +const auto_drain_limit = 1_048_576 + +/// Reads the entire request body into memory, up to `limit` bytes, along with +/// any chunked trailer fields. Blocks until the full body has arrived. +/// `Fixed` and `NoBody` requests never have trailers. +pub fn read_body( + conn: Connection, + limit: Int, +) -> Result(#(BitArray, List(#(String, String))), BodyError) { + let Http1(transport:, socket:, self:, buffer:, framing:) = conn + + case framing, consume_body(transport, socket, buffer, framing, limit) { + _framing, Ok(#(body, trailers, leftover)) -> { + process.send(self, handler.User(connection.BodyDrained(leftover:))) + Ok(#(body, trailers)) + } + connection.Fixed(_length), Error(BodyTooLarge) -> Error(BodyTooLarge) + _framing, Error(error) -> { + process.send(self, handler.User(connection.BodyAbandoned)) + Error(error) + } + } +} + +// Reconciles what the handler did with the request body so the connection can +// be safely reused. If the handler called `read_body`, picks up the leftover +// bytes it reported via a message to `conn.self`. Otherwise drains the declared +// body itself, up to `auto_drain_limit`, so a well behaved client isn't +// punished with a dropped connection just because the handler didn't care about +// its body. +fn resolve_body(conn: Connection) -> #(BitArray, Bool) { + let Http1(transport:, socket:, self:, buffer:, framing:) = conn + + case process.receive(self, 0) { + Ok(handler.User(connection.BodyDrained(leftover))) -> #(leftover, True) + Ok(handler.User(connection.BodyAbandoned)) -> #(<<>>, False) + _nothing_pending -> + case consume_body(transport, socket, buffer, framing, auto_drain_limit) { + Ok(#(_discarded, _trailers, leftover)) -> #(leftover, True) + Error(_reason) -> #(<<>>, False) + } + } +} + +// Consumes exactly the declared body from `buffer`, pulling more from the +// socket as needed, and splits off whatever comes right after it, along with +// any chunked trailer fields. +fn consume_body( + transport: transport.Transport, + socket: socket.Socket, + buffer: BitArray, + framing: connection.Framing, + limit: Int, +) -> Result(#(BitArray, List(#(String, String)), BitArray), BodyError) { + case framing { + connection.NoBody -> Ok(#(<<>>, [], buffer)) + connection.Fixed(length) if length > limit -> Error(BodyTooLarge) + connection.Fixed(length) -> + read_fixed(transport, socket, buffer, length) |> to_body_result + connection.Chunked -> + read_chunked(transport, socket, buffer, limit, bytes_tree.new(), 0) + |> to_body_result + } +} + +fn to_body_result( + result: Result(#(BitArray, List(#(String, String)), BitArray), ParseError), +) -> Result(#(BitArray, List(#(String, String)), BitArray), BodyError) { + case result { + Ok(value) -> Ok(value) + Error(ChunkTooLarge) -> Error(BodyTooLarge) + Error(_other) -> Error(InvalidBody) + } +} + +// A `Content-Length` body. Whatever's missing beyond `buffer` is read in a +// single exact size call, since the socket blocks until exactly that many +// bytes arrive or the connection drops. +fn read_fixed( + transport: transport.Transport, + socket: socket.Socket, + buffer: BitArray, + length: Int, +) -> Result(#(BitArray, List(#(String, String)), BitArray), ParseError) { + case buffer { + <> -> Ok(#(body, [], leftover)) + _ -> { + case + transport.receive_timeout( + transport, + socket, + length - bit_array.byte_size(buffer), + body_read_timeout, + ) + { + Ok(more) -> Ok(#(<>, [], <<>>)) + Error(_reason) -> Error(BodyReadFailed) + } + } + } +} + +// A `Transfer-Encoding: chunked` body. Each call advances by one chunk (or +// trailers), pulling more from the socket only for whichever piece is +// currently short. +fn read_chunked( + transport: transport.Transport, + socket: socket.Socket, + buffer: BitArray, + limit: Int, + acc: bytes_tree.BytesTree, + total: Int, +) -> Result(#(BitArray, List(#(String, String)), BitArray), ParseError) { + use #(size, remaining) <- result.try(pull_until( + transport, + socket, + buffer, + parse_chunk_line, + )) + + case size { + 0 -> { + use #(trailers, _state, remaining) <- result.try( + pull_until(transport, socket, remaining, fn(buffer) { + parse_headers(buffer, [], 0, initial_header_state()) + }), + ) + Ok(#(bytes_tree.to_bit_array(acc), trailers, remaining)) + } + size -> { + let total = total + size + case total > limit { + True -> Error(ChunkTooLarge) + False -> { + use #(data, remaining) <- result.try( + pull_until(transport, socket, remaining, take_chunk_data(_, size)), + ) + read_chunked( + transport, + socket, + remaining, + limit, + bytes_tree.append(acc, data), + total, + ) + } + } + } + } +} + +// Runs `step` against `buffer`, pulling more bytes from the socket only when +// `step` itself reports it's short. +fn pull_until( + transport: transport.Transport, + socket: socket.Socket, + buffer: BitArray, + step: fn(BitArray) -> Step(a), +) -> Result(a, ParseError) { + case step(buffer) { + Done(value) -> Ok(value) + ParseError(error) -> Error(error) + More -> + case transport.receive_timeout(transport, socket, 0, body_read_timeout) { + Ok(more) -> + pull_until(transport, socket, <>, step) + Error(_reason) -> Error(BodyReadFailed) + } + } +} + +fn parse_chunk_line(buffer: BitArray) -> Step(#(Int, BitArray)) { + use #(line, remaining) <- try_step(extract_line( + buffer, + max_chunk_size_line, + ChunkSizeLineTooLong, + BadChunkSize, + )) + use size <- try_step(parse_chunk_size(line)) + Done(#(size, remaining)) +} + +// Chunk-size lines may carry `;extensions`. Only the hex size before them +// matters here +fn parse_chunk_size(line: BitArray) -> Step(Int) { + case parse_hex_digits(line, 0, False) { + Ok(size) -> Done(size) + Error(Nil) -> ParseError(BadChunkSize) + } +} + +fn parse_hex_digits(bits: BitArray, acc: Int, any: Bool) -> Result(Int, Nil) { + case bits { + <> if byte >= 48 && byte <= 57 -> + parse_hex_digits(remaining, acc * 16 + { byte - 48 }, True) + <> if byte >= 97 && byte <= 102 -> + parse_hex_digits(remaining, acc * 16 + { byte - 87 }, True) + <> if byte >= 65 && byte <= 70 -> + parse_hex_digits(remaining, acc * 16 + { byte - 55 }, True) + _bits if any -> Ok(acc) + _bits -> Error(Nil) + } +} + +fn take_chunk_data(buffer: BitArray, size: Int) -> Step(#(BitArray, BitArray)) { + case buffer { + <> -> + Done(#(data, remaining)) + _buffer -> + case bit_array.byte_size(buffer) < size + 2 { + True -> More + False -> ParseError(BadChunkFraming) + } + } +} + /// Errors that can occur while turning a handler's `response.Response` into /// wire bytes. pub type EncodeError { @@ -368,8 +623,7 @@ pub type Head { /// Cheap conclusions drawn from `Head.headers` in one pass. pub type Metadata { Metadata( - content_length: option.Option(Int), - chunked: Bool, + framing: connection.Framing, keep_alive: Bool, upgrade: option.Option(String), ) @@ -390,6 +644,11 @@ pub type ParseError { BadHost MissingHost AmbiguousFraming + ChunkSizeLineTooLong + BadChunkSize + BadChunkFraming + ChunkTooLarge + BodyReadFailed } pub fn error_to_string(error: ParseError) -> String { @@ -412,6 +671,14 @@ pub fn error_to_string(error: ParseError) -> String { MissingHost -> "missing required Host header" AmbiguousFraming -> "conflicting Content-Length and Transfer-Encoding headers" + ChunkSizeLineTooLong -> + "chunk size line exceeds " + <> int.to_string(max_chunk_size_line) + <> " bytes" + BadChunkSize -> "malformed chunk size" + BadChunkFraming -> "malformed chunk data framing" + ChunkTooLarge -> "chunked body exceeds size limit" + BodyReadFailed -> "failed to read request body from the socket" } } @@ -426,6 +693,8 @@ const max_header_line = 8192 const max_headers = 100 +const max_chunk_size_line = 128 + /// Parses as much of a request head as `buffer` contains. pub fn parse(buffer: BitArray) -> Result(Parsed, ParseError) { let step = { @@ -743,6 +1012,13 @@ fn resolve_metadata(state: HeaderState, version: Version) -> Step(Metadata) { case state.content_length, state.chunked { option.Some(_length), True -> ParseError(AmbiguousFraming) content_length, chunked -> { + let framing = case content_length, chunked { + option.Some(length), False -> connection.Fixed(length) + option.None, True -> connection.Chunked + option.None, False -> connection.NoBody + option.Some(_length), True -> + panic as "AmbiguousFraming already rejected above" + } let keep_alive = case state.connection, version { option.Some(keep_alive), _version -> keep_alive option.None, Http11 -> True @@ -752,7 +1028,7 @@ fn resolve_metadata(state: HeaderState, version: Version) -> Step(Metadata) { True -> state.upgrade False -> option.None } - Done(Metadata(content_length:, chunked:, keep_alive:, upgrade:)) + Done(Metadata(framing:, keep_alive:, upgrade:)) } } } diff --git a/test/ewe/internal/http1_test.gleam b/test/ewe/internal/http1_test.gleam index 21b3fb0..8bd3e09 100644 --- a/test/ewe/internal/http1_test.gleam +++ b/test/ewe/internal/http1_test.gleam @@ -1,3 +1,4 @@ +import ewe/internal/connection import ewe/internal/http1 import gleam/http import gleam/option.{None, Some} @@ -18,8 +19,7 @@ pub fn simple_get_test() { assert head.headers == [#("host", "example.com"), #("connection", "keep-alive")] assert metadata.keep_alive - assert !metadata.chunked - assert metadata.content_length == None + assert metadata.framing == connection.NoBody assert remaining == <<>> } @@ -99,8 +99,7 @@ pub fn mixed_case_and_ows_headers_test() { "POST /submit HTTP/1.1\r\nHost: example.com\r\nContent-Length: 13 \r\nConnection: close\r\n\r\nHELLO WORLD!!":utf8, >> - let assert Ok(http1.Complete(head, metadata, remaining)) = - http1.parse(buffer) + let assert Ok(http1.Complete(head, metadata, remaining)) = http1.parse(buffer) assert head.headers == [ @@ -108,7 +107,7 @@ pub fn mixed_case_and_ows_headers_test() { #("content-length", "13"), #("connection", "close"), ] - assert metadata.content_length == Some(13) + assert metadata.framing == connection.Fixed(13) assert !metadata.keep_alive assert remaining == <<"HELLO WORLD!!":utf8>> } @@ -132,7 +131,7 @@ pub fn chunked_transfer_encoding_test() { let assert Ok(http1.Complete(_head, metadata, _remaining)) = http1.parse(buffer) - assert metadata.chunked + assert metadata.framing == connection.Chunked } pub fn transfer_encoding_chunked_among_multiple_tokens_test() { @@ -143,7 +142,7 @@ pub fn transfer_encoding_chunked_among_multiple_tokens_test() { let assert Ok(http1.Complete(_head, metadata, _remaining)) = http1.parse(buffer) - assert metadata.chunked + assert metadata.framing == connection.Chunked } pub fn conflicting_content_length_and_chunked_rejected_test() { @@ -288,24 +287,25 @@ pub fn no_query_string_test() { let buffer = <<"GET /plain HTTP/1.1\r\nHost: example.com\r\n\r\n":utf8>> assert http1.parse(buffer) - == Ok(http1.Complete( - http1.Head( - method: http.Get, - host: "example.com", - port: None, - path: "/plain", - query: None, - version: http1.Http11, - headers: [#("host", "example.com")], - ), - http1.Metadata( - content_length: None, - chunked: False, - keep_alive: True, - upgrade: None, + == Ok( + http1.Complete( + http1.Head( + method: http.Get, + host: "example.com", + port: None, + path: "/plain", + query: None, + version: http1.Http11, + headers: [#("host", "example.com")], + ), + http1.Metadata( + framing: connection.NoBody, + keep_alive: True, + upgrade: None, + ), + <<>>, ), - <<>>, - )) + ) } pub fn websocket_upgrade_requested_test() {