From 50f157e8f0e69e41e6d0a0678dabfd9fa52bdf12 Mon Sep 17 00:00:00 2001 From: vshakitskiy Date: Tue, 14 Jul 2026 14:36:05 +0300 Subject: [PATCH] clear files --- dev/autobahn.gleam | 42 - dev/benchmark.gleam | 45 - dev/https.gleam | 28 - dev/preview.gleam | 384 -------- gleam.toml | 10 +- manifest.toml | 25 +- src/ewe.gleam | 1074 +---------------------- src/ewe/internal/clock.gleam | 133 --- src/ewe/internal/decoder.gleam | 132 --- src/ewe/internal/encoder.gleam | 124 --- src/ewe/internal/file.gleam | 79 -- src/ewe/internal/handler.gleam | 88 -- src/ewe/internal/http1.gleam | 886 ------------------- src/ewe/internal/http1/buffer.gleam | 26 - src/ewe/internal/http1/handler.gleam | 332 ------- src/ewe/internal/stream/chunked.gleam | 130 --- src/ewe/internal/stream/sse.gleam | 158 ---- src/ewe/internal/stream/websocket.gleam | 442 ---------- src/ewe_ffi.erl | 199 ----- test/client/http.gleam | 67 -- test/client/tcp.gleam | 34 - test/h1spec_test.gleam | 233 ----- test/http_test.gleam | 150 ---- test/server.gleam | 43 - 24 files changed, 19 insertions(+), 4845 deletions(-) delete mode 100644 dev/autobahn.gleam delete mode 100644 dev/benchmark.gleam delete mode 100644 dev/https.gleam delete mode 100644 dev/preview.gleam delete mode 100644 src/ewe/internal/clock.gleam delete mode 100644 src/ewe/internal/decoder.gleam delete mode 100644 src/ewe/internal/encoder.gleam delete mode 100644 src/ewe/internal/file.gleam delete mode 100644 src/ewe/internal/handler.gleam delete mode 100644 src/ewe/internal/http1.gleam delete mode 100644 src/ewe/internal/http1/buffer.gleam delete mode 100644 src/ewe/internal/http1/handler.gleam delete mode 100644 src/ewe/internal/stream/chunked.gleam delete mode 100644 src/ewe/internal/stream/sse.gleam delete mode 100644 src/ewe/internal/stream/websocket.gleam delete mode 100644 src/ewe_ffi.erl delete mode 100644 test/client/http.gleam delete mode 100644 test/client/tcp.gleam delete mode 100644 test/h1spec_test.gleam delete mode 100644 test/http_test.gleam delete mode 100644 test/server.gleam diff --git a/dev/autobahn.gleam b/dev/autobahn.gleam deleted file mode 100644 index e067ca6..0000000 --- a/dev/autobahn.gleam +++ /dev/null @@ -1,42 +0,0 @@ -import ewe -import gleam/erlang/process -import gleam/otp/static_supervisor as supervisor -import logging - -pub fn main() -> Nil { - logging.configure() - logging.set_level(logging.Info) - - let ewe_server = - ewe.new(fn(req) { - ewe.upgrade_websocket( - req, - on_init: fn(_conn, selector) { #(Nil, selector) }, - handler: fn(conn, state, msg) { - case msg { - ewe.Text(text_frame) -> { - let _ = ewe.send_text_frame(conn, text_frame) - ewe.websocket_continue(state) - } - ewe.Binary(binary_frame) -> { - let _ = ewe.send_binary_frame(conn, binary_frame) - ewe.websocket_continue(state) - } - _ -> ewe.websocket_continue(state) - } - }, - on_close: fn(_conn, _state) { Nil }, - ) - }) - |> ewe.bind("0.0.0.0") - |> ewe.listening(port: 8080) - |> ewe.supervised() - - let assert Ok(_) = - supervisor.new(supervisor.OneForAll) - |> supervisor.add(ewe_server) - |> supervisor.restart_tolerance(intensity: 1_000_000, period: 1_000_000) - |> supervisor.start() - - process.sleep_forever() -} diff --git a/dev/benchmark.gleam b/dev/benchmark.gleam deleted file mode 100644 index 31ecddc..0000000 --- a/dev/benchmark.gleam +++ /dev/null @@ -1,45 +0,0 @@ -import ewe.{type Request, type Response} -import gleam/erlang/process -import gleam/http -import gleam/http/request -import gleam/http/response -import gleam/result - -pub fn main() { - let assert Ok(_) = - ewe.new(handle_request) - |> ewe.listening(port: 8080) - |> ewe.start - - process.sleep_forever() -} - -fn handle_request(req: Request) -> Response { - case req.method, request.path_segments(req) { - http.Get, [] -> - response.new(200) - |> response.set_body(ewe.Empty) - http.Get, ["user", id] -> - response.new(200) |> response.set_body(ewe.TextData(id)) - http.Post, ["user"] -> { - case ewe.read_body(req, 40_000_000) { - Ok(req) -> { - let content_type = - req - |> request.get_header("content-type") - |> result.unwrap("application/octet-stream") - - response.new(200) - |> response.set_body(ewe.BitsData(req.body)) - |> response.prepend_header("content-type", content_type) - } - Error(_) -> - response.new(413) - |> response.set_body(ewe.Empty) - } - } - _, _ -> - response.new(404) - |> response.set_body(ewe.Empty) - } -} diff --git a/dev/https.gleam b/dev/https.gleam deleted file mode 100644 index fdbd72d..0000000 --- a/dev/https.gleam +++ /dev/null @@ -1,28 +0,0 @@ -import ewe.{type Request, type Response} -import gleam/erlang/process -import gleam/http/response -import logging - -pub fn main() { - logging.configure() - logging.set_level(logging.Debug) - - // Start the server that has TLS enabled. - let assert Ok(_) = - ewe.new(handler) - |> ewe.bind("0.0.0.0") - |> ewe.listening(port: 8080) - |> ewe.enable_tls( - certificate_file: "examples/priv/localhost.crt", - key_file: "examples/priv/localhost.key", - ) - |> ewe.start - - process.sleep_forever() -} - -fn handler(_req: Request) -> Response { - response.new(200) - |> response.set_header("content-type", "text/plain; charset=utf-8") - |> response.set_body(ewe.TextData("Hello, World!")) -} diff --git a/dev/preview.gleam b/dev/preview.gleam deleted file mode 100644 index efb4a03..0000000 --- a/dev/preview.gleam +++ /dev/null @@ -1,384 +0,0 @@ -import gleam/bit_array -import gleam/crypto -import gleam/dict -import gleam/erlang/charlist.{type Charlist} -import gleam/erlang/process.{type Name, type Pid, type Subject} -import gleam/http/request -import gleam/http/response -import gleam/int -import gleam/list -import gleam/option.{None, Some} -import gleam/otp/actor -import gleam/otp/static_supervisor as supervisor -import gleam/otp/supervision.{type ChildSpecification} -import gleam/result -import logging - -import ewe.{type Request, type Response} - -pub fn main() { - logging.configure() - logging.set_level(logging.Debug) - - // Create a named subject for the pubsub worker - let pubsub_name = process.new_name("pubsub") - let pubsub = process.named_subject(pubsub_name) - - // Configure and start the supervision tree with pubsub worker and the ewe - // server, that listens on port 8080 - let assert Ok(_) = - supervisor.new(supervisor.OneForAll) - |> supervisor.add(pubsub_worker(pubsub_name)) - |> supervisor.add( - ewe.new(handler(_, pubsub)) - |> ewe.bind("0.0.0.0") - |> ewe.listening(port: 8080) - |> ewe.supervised(), - ) - |> supervisor.start() - - process.sleep_forever() -} - -// Define the messages that can be sent to the pubsub worker -type PubSubMessage { - Subscribe(topic: String, client: Subject(Broadcast)) - Publish(topic: String, message: Broadcast) - Unsubscribe(topic: String, client: Subject(Broadcast)) -} - -// Define the messages that could be received by websocket and SSE clients -type Broadcast { - Text(String) - Bytes(BitArray) -} - -// Define the state of the websocket connection -type WebsocketState { - WebsocketState( - pubsub: Subject(PubSubMessage), - topic: String, - client: Subject(Broadcast), - ) -} - -// Main logic of the pubsub worker, that handles the messages and keeps track of -// the clients on topics. Its implementation is not really important -fn pubsub_worker( - named: Name(PubSubMessage), -) -> ChildSpecification(Subject(PubSubMessage)) { - let pubsub = - actor.new(dict.new()) - |> actor.on_message(fn(state, msg) { - case msg { - Subscribe(topic:, client:) -> { - let new_state = - dict.upsert(in: state, update: topic, with: fn(clients) { - case clients { - Some(clients) -> [client, ..clients] - None -> { - logging.log(logging.Info, "Creating topic " <> topic) - [client] - } - } - }) - - let assert Ok(pid) = process.subject_owner(client) - logging.log( - logging.Info, - "Subscribing client " <> pid_to_string(pid) <> " to topic " <> topic, - ) - - actor.continue(new_state) - } - Publish(topic:, message:) -> { - case message { - Text(text) -> - logging.log( - logging.Info, - "Publishing text message `" <> text <> "` to topic " <> topic, - ) - Bytes(_binary) -> - logging.log( - logging.Info, - "Publishing binary message to topic " <> topic, - ) - } - - case dict.get(state, topic) { - Ok(clients) -> list.each(clients, actor.send(_, message)) - Error(_) -> Nil - } - - actor.continue(state) - } - Unsubscribe(topic:, client:) -> { - let assert Ok(pid) = process.subject_owner(client) - logging.log( - logging.Info, - "Unsubscribing client " - <> pid_to_string(pid) - <> " from topic " - <> topic, - ) - - let new_state = case dict.get(state, topic) { - Ok([_]) | Ok([]) -> { - logging.log(logging.Info, "Dropping topic " <> topic) - dict.drop(state, [topic]) - } - Ok(clients) -> { - list.filter(clients, fn(c) { c != client }) - |> dict.insert(state, topic, _) - } - Error(_) -> state - } - - actor.continue(new_state) - } - } - }) - |> actor.named(named) - - supervision.worker(fn() { - logging.log(logging.Info, "Starting pubsub worker") - actor.start(pubsub) - }) -} - -// Main HTTP request handler that routes requests to different endpoints -fn handler(req: Request, pubsub: Subject(PubSubMessage)) -> Response { - case request.path_segments(req) { - // GET /hello/:name - Simple greeting endpoint - ["hello", name] -> { - response.new(200) - |> response.set_header("content-type", "text/plain; charset=utf-8") - |> response.set_body(ewe.TextData("Hello, " <> name <> "!")) - } - // GET /bytes/:amount - Generate random N bytes - ["bytes", amount] -> { - let random_bytes = - int.parse(amount) - |> result.unwrap(0) - |> crypto.strong_random_bytes() - - response.new(200) - |> response.set_header("content-type", "application/octet-stream") - |> response.set_body(ewe.BitsData(random_bytes)) - } - - // POST /echo - Echo back the request body - ["echo"] -> handle_echo(req) - - // POST /stream/:chunk_size - Stream and echo back the request body in chunks - ["stream", chunk_size] -> - handle_stream(req, int.parse(chunk_size) |> result.unwrap(16)) - // GET /file/:path - Serve a file from the public directory - ["file", path] -> serve_file(path) - - // POST /topic/:topic/ws - Upgrade to WebSocket connection - ["topic", topic, "ws"] -> - ewe.upgrade_websocket( - req, - on_init: fn(_conn, selector) { - logging.log( - logging.Info, - "WebSocket connection opened: " <> pid_to_string(process.self()), - ) - - let client = process.new_subject() - process.send(pubsub, Subscribe(topic:, client:)) - - let state = WebsocketState(pubsub:, topic:, client:) - let selector = process.select(selector, client) - - #(state, selector) - }, - handler: handle_websocket, - on_close: fn(_conn, state) { - let assert Ok(pid) = process.subject_owner(state.client) - logging.log( - logging.Info, - "WebSocket connection closed: " <> pid_to_string(pid), - ) - - process.send(pubsub, Unsubscribe(state.topic, state.client)) - }, - ) - - // POST /topic/:topic/sse - Switch to Server-Sent Events connection - ["topic", topic, "sse"] -> - ewe.sse( - req, - on_init: fn(client) { - logging.log( - logging.Info, - "SSE connection opened: " <> pid_to_string(process.self()), - ) - - process.send(pubsub, Subscribe(topic:, client:)) - client - }, - handler: fn(conn, client, message) { - let assert Ok(_) = case message { - Text(text) -> ewe.send_event(conn, ewe.event(text)) - _ -> Ok(Nil) - } - - ewe.sse_continue(client) - }, - on_close: fn(_conn, client) { - logging.log( - logging.Info, - "SSE connection closed: " <> pid_to_string(process.self()), - ) - - process.send(pubsub, Unsubscribe(topic:, client:)) - }, - ) - - // All other routes return 404 - _ -> - response.new(404) - |> response.set_body(ewe.Empty) - } -} - -fn handle_echo(req: Request) -> Response { - let content_type = - request.get_header(req, "content-type") - |> result.unwrap("application/octet-stream") - - case ewe.read_body(req, 1024) { - Ok(req) -> - response.new(200) - |> response.set_header("content-type", content_type) - |> response.set_body(ewe.BitsData(req.body)) - Error(ewe.BodyTooLarge) -> - response.new(413) - |> response.set_header("content-type", "text/plain; charset=utf-8") - |> response.set_body(ewe.TextData("Body too large")) - Error(ewe.InvalidBody) -> - response.new(400) - |> response.set_header("content-type", "text/plain; charset=utf-8") - |> response.set_body(ewe.TextData("Invalid request")) - } -} - -pub type StreamMessage { - Chunk(BitArray) - Done - BodyError(ewe.BodyError) -} - -fn stream_resource( - consumer: ewe.Consumer, - subject: Subject(StreamMessage), - chunk_size: Int, -) -> Nil { - process.sleep(int.random(250)) - case consumer(chunk_size) { - Ok(ewe.Consumed(data, next)) -> { - logging.log(logging.Info, { - "Consumed " <> int.to_string(bit_array.byte_size(data)) <> " bytes" - }) - - process.send(subject, Chunk(data)) - stream_resource(next, subject, chunk_size) - } - Ok(ewe.Done) -> { - process.send(subject, Done) - } - Error(body_error) -> { - process.send(subject, BodyError(body_error)) - } - } -} - -fn handle_stream(req: Request, chunk_size: Int) -> Response { - let content_type = - request.get_header(req, "content-type") - |> result.unwrap("application/octet-stream") - - case ewe.stream_body(req) { - Ok(consumer) -> { - ewe.chunked_body( - req, - response.new(200) |> response.set_header("content-type", content_type), - on_init: fn(subject) { - process.spawn(fn() { stream_resource(consumer, subject, chunk_size) }) - - Nil - }, - handler: fn(chunked_body, state, message) { - case message { - Chunk(data) -> - case ewe.send_chunk(chunked_body, data) { - Ok(Nil) -> ewe.chunked_continue(state) - Error(_) -> ewe.chunked_stop_abnormal("Failed to send chunk") - } - Done -> ewe.chunked_stop() - BodyError(_body_error) -> - ewe.chunked_stop_abnormal("failed to read body") - } - }, - on_close: fn(_conn, _state) { - logging.log(logging.Info, "Stream closed") - }, - ) - } - Error(_) -> - response.new(400) - |> response.set_header("content-type", "text/plain; charset=utf-8") - |> response.set_body(ewe.TextData("Invalid request")) - } -} - -fn serve_file(path: String) -> Response { - case ewe.file("public/" <> path, offset: None, limit: None) { - Ok(file) -> { - response.new(200) - |> response.set_header("content-type", "application/octet-stream") - |> response.set_body(file) - } - Error(_) -> { - response.new(404) - |> response.set_header("content-type", "text/plain; charset=utf-8") - |> response.set_body(ewe.TextData("File not found")) - } - } -} - -fn handle_websocket( - conn: ewe.WebsocketConnection, - state: WebsocketState, - msg: ewe.WebsocketMessage(Broadcast), -) -> ewe.WebsocketNext(WebsocketState, Broadcast) { - case msg { - ewe.Text(text) -> { - process.send(state.pubsub, Publish(state.topic, Text(text))) - ewe.websocket_continue(state) - } - - ewe.Binary(binary) -> { - process.send(state.pubsub, Publish(state.topic, Bytes(binary))) - ewe.websocket_continue(state) - } - - ewe.User(message) -> { - let assert Ok(_) = case message { - Text(text) -> ewe.send_text_frame(conn, text) - Bytes(binary) -> ewe.send_binary_frame(conn, binary) - } - - ewe.websocket_continue(state) - } - } -} - -fn pid_to_string(pid: Pid) -> String { - charlist.to_string(pid_to_list(pid)) -} - -@external(erlang, "erlang", "pid_to_list") -fn pid_to_list(pid: Pid) -> Charlist diff --git a/gleam.toml b/gleam.toml index ca69559..abcd83e 100644 --- a/gleam.toml +++ b/gleam.toml @@ -1,5 +1,5 @@ name = "ewe" -version = "4.0.1" +version = "5.0.0" description = "🐑 a fluffy Gleam web server" target = "erlang" @@ -9,20 +9,18 @@ links = [{ title = "Gleam", href = "https://gleam.run" }] [dependencies] gleam_stdlib = ">= 1.0.0 and < 2.0.0" -glisten = ">= 9.0.0 and < 10.0.0" +# glisten = ">= 9.0.0 and < 10.0.0" +glisten = { git = "https://github.com/vshakitskiy/glisten.git", ref = "ab2422b" } gleam_otp = ">= 1.1.0 and < 2.0.0" 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" -compresso = "0.1.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" gleam_httpc = ">= 5.0.0 and < 6.0.0" -gleam_json = ">= 3.0.2 and < 4.0.0" -gleam_crypto = ">= 1.5.1 and < 2.0.0" [erlang] -application_start_module = "ewe@internal@clock" +# application_start_module = "ewe@internal@clock" diff --git a/manifest.toml b/manifest.toml index ba45267..37511ab 100644 --- a/manifest.toml +++ b/manifest.toml @@ -1,34 +1,33 @@ -# 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 = "exception", version = "2.1.0", build_tools = ["gleam"], requirements = ["gleam_stdlib"], otp_app = "exception", source = "hex", outer_checksum = "329D269D5C2A314F7364BD2711372B6F2C58FA6F39981572E5CA68624D291F8C" }, + { 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_httpc", version = "5.0.0", build_tools = ["gleam"], requirements = ["gleam_erlang", "gleam_http", "gleam_stdlib"], otp_app = "gleam_httpc", source = "hex", outer_checksum = "C545172618D07811494E97AAA4A0FB34DA6F6D0061FDC8041C2F8E3BE2B2E48F" }, - { name = "gleam_json", version = "3.1.0", build_tools = ["gleam"], requirements = ["gleam_stdlib"], otp_app = "gleam_json", source = "hex", outer_checksum = "44FDAA8847BE8FC48CA7A1C089706BD54BADCC4C45B237A992EDDF9F2CDB2836" }, { 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.0", build_tools = ["gleam"], requirements = [], otp_app = "gleam_stdlib", source = "hex", outer_checksum = "960090C2FB391784BB34267B099DC9315CC1B1F6013E7415BC763CEF1905D7D3" }, - { 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 = "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 = "ab2422bb27a9acfe8cadc1d755901ff61b0f96d8" }, { 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] -compresso = { version = "0.1.0" } exception = { version = ">= 2.1.0 and < 3.0.0" } -gleam_crypto = { version = ">= 1.5.1 and < 2.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" } -gleam_json = { version = ">= 3.0.2 and < 4.0.0" } gleam_otp = { version = ">= 1.1.0 and < 2.0.0" } gleam_stdlib = { version = ">= 1.0.0 and < 2.0.0" } gleeunit = { version = ">= 1.0.0 and < 2.0.0" } -glisten = { version = ">= 9.0.0 and < 10.0.0" } +glisten = { git = "https://github.com/vshakitskiy/glisten.git", ref = "ab2422b" } logging = { version = ">= 1.3.0 and < 2.0.0" } websocks = { version = ">= 3.0.0 and < 4.0.0" } diff --git a/src/ewe.gleam b/src/ewe.gleam index b9609f9..8e58a24 100644 --- a/src/ewe.gleam +++ b/src/ewe.gleam @@ -1,1071 +1,3 @@ -//// - -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 -import gleam/erlang/process.{type Selector, type Subject} -import gleam/http -import gleam/http/request.{type Request as HttpRequest} -import gleam/http/response.{type Response as HttpResponse} -import gleam/int -import gleam/io -import gleam/option.{type Option, None, Some} -import gleam/otp/actor -import gleam/otp/factory_supervisor as factory -import gleam/otp/static_supervisor.{type Supervisor} as supervisor -import gleam/otp/supervision -import gleam/result -import gleam/string_tree.{type StringTree} -import glisten -import glisten/internal/listener -import glisten/socket/options as glisten_options -import glisten/transport -import websocks - -// CONNECTION -// ----------------------------------------------------------------------------- - -/// Represents the request body and connection metadata. Access the body using -/// `ewe.read_body`, or retrieve client information with `ewe.get_client_info`. -pub type Connection = - ewe_http.Connection - -// IP ADDRESS -// ----------------------------------------------------------------------------- - -/// Represents an IP address. Can be either IPv4 or IPv6. -pub type IpAddress { - IpV4(Int, Int, Int, Int) - IpV6(Int, Int, Int, Int, Int, Int, Int, Int) -} - -/// Converts an `IpAddress` to its string representation. -pub fn ip_address_to_string(address address: IpAddress) -> String { - ewe_to_glisten_ip(address) - |> glisten.ip_address_to_string -} - -fn glisten_to_ewe_ip(ip: glisten.IpAddress) -> IpAddress { - case ip { - glisten.IpV4(n1, n2, n3, n4) -> IpV4(n1, n2, n3, n4) - glisten.IpV6(n1, n2, n3, n4, n5, n6, n7, n8) -> - IpV6(n1, n2, n3, n4, n5, n6, n7, n8) - } -} - -fn glisten_options_to_ewe_ip(ip: glisten_options.IpAddress) -> IpAddress { - case ip { - glisten_options.IpV4(n1, n2, n3, n4) -> IpV4(n1, n2, n3, n4) - glisten_options.IpV6(n1, n2, n3, n4, n5, n6, n7, n8) -> - IpV6(n1, n2, n3, n4, n5, n6, n7, n8) - } -} - -fn ewe_to_glisten_ip(ip: IpAddress) -> glisten.IpAddress { - case ip { - IpV4(n1, n2, n3, n4) -> glisten.IpV4(n1, n2, n3, n4) - IpV6(n1, n2, n3, n4, n5, n6, n7, n8) -> - glisten.IpV6(n1, n2, n3, n4, n5, n6, n7, n8) - } -} - -// INFORMATION -// ----------------------------------------------------------------------------- - -/// Represents a socket address with IP and port. Use `ewe.get_client_info` to -/// get the client's address from a connection, or `ewe.get_server_info` to get -/// the server's bound address. -pub type SocketAddress { - SocketAddress(ip: IpAddress, port: Int) -} - -/// Retrieves the client's socket address from the connection. Returns error if -/// the socket information is unavailable. -pub fn get_client_info( - connection connection: Connection, -) -> Result(SocketAddress, Nil) { - transport.peername(connection.transport, connection.socket) - |> result.map(fn(server_info) { - SocketAddress(glisten_options_to_ewe_ip(server_info.0), server_info.1) - }) -} - -/// Gets the server's bound address and port. Requires the server to be running -/// and the listener name to match the one set in `ewe.with_name`. -pub fn get_server_info( - listener_name name: process.Name(listener.Message), -) -> SocketAddress { - let server_info = glisten.get_server_info(name, 10_000) - let ip_address = glisten_to_ewe_ip(server_info.ip_address) - - SocketAddress(ip: ip_address, port: server_info.port) -} - -// RESPONSE -// ----------------------------------------------------------------------------- - -/// Represents possible response body options. -/// -/// Types for direct usage: -/// - Regular data: `TextData`, `BytesData`, `BitsData`, `StringTreeData`, -/// `Empty`. -/// -/// Types that should never be used directly: -/// - `File`: see `ewe.file` to construct it. -/// - `Chunked`: indicates that response body is being sent in chunks with -/// `chunked` transfer encoding. -/// - `Websocket`: indicates that request is being upgraded to a WebSocket -/// connection. -/// - `SSE`: indicates that request is being upgraded to a Server-Sent Events -/// connection. -pub type ResponseBody { - /// Allows to set response body from a string. - TextData(String) - /// Allows to set response body from bytes. - BytesData(BytesTree) - /// Allows to set response body from bits. - BitsData(BitArray) - /// Allows to set response body from a string tree. - StringTreeData(StringTree) - /// Allows to set empty response body. - Empty - - /// Allows to set response body from a file more efficiently rather than - /// sending contents in regular data types. - File(descriptor: file.IoDevice, offset: Int, size: Int) - - /// Indicates that response body is being sent in chunks with `chunked` - /// transfer encoding. - Chunked - /// Indicates that request is being upgraded to a WebSocket connection. - Websocket - /// Indicates that request is being upgraded to a Server-Sent Events - /// connection. - SSE -} - -/// A convenient alias for a HTTP response with a `ResponseBody` as the body. -/// -pub type Response = - HttpResponse(ResponseBody) - -fn transform_response_body( - resp: Response, -) -> HttpResponse(ewe_http.ResponseBody) { - response.set_body(resp, case resp.body { - TextData(text) -> ewe_http.TextData(text) - BytesData(bytes) -> ewe_http.BytesData(bytes) - BitsData(bits) -> ewe_http.BitsData(bits) - StringTreeData(string_tree) -> ewe_http.StringTreeData(string_tree) - - Chunked -> ewe_http.Chunked - File(descriptor, offset, size) -> ewe_http.File(descriptor, offset, size) - - Websocket -> ewe_http.Websocket - SSE -> ewe_http.SSE - - Empty -> ewe_http.Empty - }) -} - -/// Error type returned by `ewe.file` when opening a file for the response body. -pub type FileError { - /// File does not exist. - NoEntry - /// Missing permission for reading the file, or for searching one of the - /// parents directories. - NoAccess - /// The named file is a directory. - IsDirectory - /// Untypical file error. - UnknownFileError(dynamic.Dynamic) -} - -fn internal_to_file_error(error: file.FileError) -> FileError { - case error { - file.Enoent -> NoEntry - file.Eacces -> NoAccess - file.Eisdir -> IsDirectory - file.Eunknown(error) -> UnknownFileError(error) - } -} - -/// Creates a file response body. Use `offset` to skip bytes from the start, and -/// `limit` to send only a portion of the file. -pub fn file( - path: String, - offset offset: Option(Int), - limit limit: Option(Int), -) -> Result(ResponseBody, FileError) { - // TODO: handle invalid offset + limit? - case file.open(path) { - Ok(file) -> - Ok(File( - file.descriptor, - offset: option.unwrap(offset, 0), - size: option.unwrap(limit, file.size), - )) - Error(error) -> Error(internal_to_file_error(error)) - } -} - -// BUILDER -// ----------------------------------------------------------------------------- - -/// Contains all server configurations, can be adjusted by different builder -/// functions. -pub opaque type Builder { - Builder( - handler: fn(Request) -> Response, - port: Int, - interface: String, - ipv6: Bool, - tls: Option(#(String, String)), - on_start: fn(http.Scheme, SocketAddress) -> Nil, - on_crash: Response, - listener_name: process.Name(listener.Message), - idle_timeout: Int, - ) -} - -/// Creates new server builder with handler provided. -pub fn new(handler: fn(Request) -> Response) -> Builder { - Builder( - handler:, - port: 8080, - interface: "127.0.0.1", - ipv6: False, - tls: None, - on_start: fn(scheme, server) { - let address = case server.ip { - IpV6(..) -> "[" <> ip_address_to_string(server.ip) <> "]" - IpV4(..) -> ip_address_to_string(server.ip) - } - - let url = - http.scheme_to_string(scheme) - <> "://" - <> address - <> ":" - <> int.to_string(server.port) - - io.println("Listening on " <> url) - }, - on_crash: response.new(500) |> response.set_body(Empty), - listener_name: process.new_name("glisten_listener"), - idle_timeout: 10_000, - ) -} - -/// Binds server to a specific network interface; e.g., "0.0.0.0" for all IPv4 -/// interfaces, "127.0.0.1" for localhost, "::" for all IPv6 interfaces, or -/// "::1" for IPv6 loopback. Crashes the program if the interface is invalid. -pub fn bind(builder: Builder, interface interface: String) -> Builder { - Builder(..builder, interface:) -} - -/// Sets the listening port for server. -pub fn listening(builder: Builder, port port: Int) -> Builder { - Builder(..builder, port:) -} - -/// Sets the listening port to 0, which causes the OS to assign a random -/// available port. -pub fn listening_random(builder: Builder) -> Builder { - Builder(..builder, port: 0) -} - -/// Forces the underlying socket to use IPv6. On IPv4 provided in `ewe.bind` or -/// if the system does not support IPv6, the server crashes. The exceptions are -/// `localhost`, `127.0.0.1` or `0.0.0.0`, they are automatically bound to work -/// with either address family. -pub fn force_ipv6(builder: Builder) -> Builder { - Builder(..builder, ipv6: True) -} - -/// Enables TLS (HTTPS) support, with provided certificate and key files. -/// Crashes the program if the files don't exist or are invalid. -pub fn enable_tls( - builder: Builder, - certificate_file certificate_file: String, - key_file key_file: String, -) -> Builder { - Builder(..builder, tls: Some(#(certificate_file, key_file))) -} - -/// Sets a custom listener process name. This name is required when calling -/// `ewe.get_server_info` to retrieve the server's bound address and port. -pub fn with_name( - builder: Builder, - name: process.Name(listener.Message), -) -> Builder { - Builder(..builder, listener_name: name) -} - -/// Sets a callback function called after the server starts. Receives the scheme -/// and server's socket address. -pub fn on_start( - builder: Builder, - on_start: fn(http.Scheme, SocketAddress) -> Nil, -) -> Builder { - Builder(..builder, on_start:) -} - -/// Sets an empty `on_start` function. -pub fn quiet(builder: Builder) -> Builder { - Builder(..builder, on_start: fn(_, _) { Nil }) -} - -/// Sets a custom response that will be sent when server crashes. -pub fn on_crash(builder: Builder, on_crash: Response) -> Builder { - Builder(..builder, on_crash:) -} - -/// Sets the idle timeout in milliseconds. Connections are closed after this -/// period of inactivity. Defaults to 10_000ms if the value is negative. -pub fn idle_timeout(builder: Builder, idle_timeout: Int) -> Builder { - case idle_timeout { - idle_timeout if idle_timeout >= 0 -> Builder(..builder, idle_timeout:) - _ -> Builder(..builder, idle_timeout: 10_000) - } -} - -// SERVER -// ----------------------------------------------------------------------------- - -/// Starts the server with the provided configuration. -pub fn start( - builder: Builder, -) -> Result(actor.Started(Supervisor), actor.StartError) { - let handler = fn(req) { transform_response_body(builder.handler(req)) } - let on_crash = transform_response_body(builder.on_crash) - - let factory_name = process.new_name("ewe_streams") - - let factory_child = - factory.worker_child(fn(start) { start() }) - |> factory.restart_strategy(supervision.Temporary) - |> factory.named(factory_name) - |> factory.supervised() - - let glisten = - glisten.new( - handler.init, - handler.loop(handler, on_crash, factory_name, builder.idle_timeout), - ) - |> glisten.with_listener_name(builder.listener_name) - - let glisten = case builder.ipv6 { - True -> glisten.with_ipv6(glisten) - False -> glisten - } - - let glisten = glisten.bind(glisten, builder.interface) - - let glisten = case builder.tls { - Some(#(cert, key)) -> glisten.with_tls(glisten, cert, key) - // Uncomment once http2 will be implemented! - // |> glisten.with_http2 - None -> glisten - } - - supervisor.new(supervisor.OneForAll) - |> supervisor.add(glisten.supervised(glisten, builder.port)) - |> supervisor.add(factory_child) - |> supervisor.start() - |> result.map(fn(started) { - let scheme = case builder.tls { - Some(#(_, _)) -> http.Https - None -> http.Http - } - - let server_info = glisten.get_server_info(builder.listener_name, 10_000) - let ip_address = glisten_to_ewe_ip(server_info.ip_address) - - let server = SocketAddress(ip: ip_address, port: server_info.port) - - builder.on_start(scheme, server) - - started - }) -} - -/// Returns a child specification for use in a supervision tree. -pub fn supervised( - builder: Builder, -) -> supervision.ChildSpecification(supervisor.Supervisor) { - supervision.supervisor(fn() { start(builder) }) -} - -// REQUEST -// ----------------------------------------------------------------------------- - -/// Possible errors that can occur when reading a body. -pub type BodyError { - /// Body is larger than the provided limit. - BodyTooLarge - /// Body is malformed. - InvalidBody -} - -/// A convenient alias for a HTTP request with a `Connection` as the body. -pub type Request = - HttpRequest(Connection) - -/// Reads body from the request. Returns `BodyTooLarge` if body exceeds -/// `bytes_limit`, or `InvalidBody` if malformed. Supports both chunked and -/// content-length bodies. -pub fn read_body( - req: Request, - bytes_limit bytes_limit: Int, -) -> Result(HttpRequest(BitArray), BodyError) { - case ewe_http.read_body(req, bytes_limit) { - Ok(req) -> Ok(req) - Error(ewe_http.BodyTooLarge) -> Error(BodyTooLarge) - Error(_) -> Error(InvalidBody) - } -} - -/// A convenient alias for a consumer that reads `N` amount of bytes from the -/// request body stream. -pub type Consumer = - fn(Int) -> Result(Stream, BodyError) - -/// The progress of reading the request body stream. -pub type Stream { - /// Chunk of data has been consumed. - Consumed(data: BitArray, next: Consumer) - /// Signifies that the request body stream has been fully consumed. - Done -} - -/// Returns a consumer for streaming the request body in chunks. -pub fn stream_body(req: Request) -> Result(Consumer, BodyError) { - case ewe_http.stream_body(req) { - Ok(consumer) -> Ok(consumer_adapter(consumer)) - Error(_) -> Error(InvalidBody) - } -} - -fn consumer_adapter( - internal_consumer: fn(Int) -> Result(ewe_http.Stream, ewe_http.ParseError), -) -> Consumer { - fn(size) { - case internal_consumer(size) { - Ok(ewe_http.Done) -> Ok(Done) - Ok(ewe_http.Consumed(data, next)) -> { - Ok(Consumed(data, consumer_adapter(next))) - } - Error(_) -> Error(InvalidBody) - } - } -} - -// CHUNKED RESPONSE -// ----------------------------------------------------------------------------- - -/// Represents a chunked response body. This type is used to send a chunked -/// response to the client. -pub type ChunkedBody = - chunked.ChunkedBody - -/// Represents an instruction on how chunked response should be processed. -/// -/// - continue processing the chunked response. -/// - stop the chunked response normally. -/// - stop the chunked response with abnormal reason. -pub opaque type ChunkedNext(user_state) { - ChunkedContinue(user_state) - ChunkedStop - ChunkedAbnormalStop(reason: String) -} - -/// Instructs chunked response to continue processing. -pub fn chunked_continue(user_state: user_state) -> ChunkedNext(user_state) { - ChunkedContinue(user_state) -} - -/// Instructs chunked response to stop normally. -pub fn chunked_stop() -> ChunkedNext(user_state) { - ChunkedStop -} - -/// Instructs chunked response to stop with abnormal reason. -pub fn chunked_stop_abnormal(reason: String) -> ChunkedNext(user_state) { - ChunkedAbnormalStop(reason) -} - -fn to_internal_chunked_next( - next: ChunkedNext(user_state), -) -> chunked.ChunkedNext(user_state) { - case next { - ChunkedContinue(user_state) -> chunked.Continue(user_state) - ChunkedStop -> chunked.NormalStop - ChunkedAbnormalStop(reason) -> chunked.AbnormalStop(reason) - } -} - -/// Sets up the connection for chunked response. -/// -/// `on_init` function is called once the chunked response process is -/// initialized. The argument is subject that can be used to send chunks to the -/// client. It must return initial state. -/// -/// `handler` function is called for every message received. It must return -/// instruction on how chunked response should proceed. -/// -/// `on_close` function is called when the chunked response process is going to be stopped. -pub fn chunked_body( - req: Request, - resp: HttpResponse(a), - on_init on_init: fn(Subject(user_message)) -> user_state, - handler handler: fn(ChunkedBody, user_state, user_message) -> - ChunkedNext(user_state), - on_close on_close: fn(ChunkedBody, user_state) -> Nil, -) -> Response { - let handler = fn(conn, state, msg) { - handler(conn, state, msg) - |> to_internal_chunked_next() - } - - let transport = req.body.transport - let socket = req.body.socket - - case chunked.send_response(resp, transport, socket) { - Ok(Nil) -> { - let started = - factory.get_by_name(req.body.factory_name) - |> factory.start_child(fn() { - chunked.start(transport, socket, on_init, handler, on_close) - }) - - case started { - Ok(started) -> { - let _ = transport.controlling_process(transport, socket, started.pid) - response.new(200) |> response.set_body(Chunked) - } - Error(_) -> response.new(400) |> response.set_body(Empty) - } - } - Error(Nil) -> response.new(400) |> response.set_body(Empty) - } -} - -/// Sends a chunk to the client. -pub fn send_chunk( - body: ChunkedBody, - chunk: BitArray, -) -> Result(Nil, glisten.SocketReason) { - chunked.send_chunk(body.transport, body.socket, chunk) -} - -// WEBSOCKET -// ----------------------------------------------------------------------------- - -/// Represents a WebSocket connection between a client and a server. -pub type WebsocketConnection = - websocket.WebsocketConnection - -/// Represents an instruction on how WebSocket connection should proceed. -/// -/// - continue processing the WebSocket connection. -/// - continue processing the WebSocket connection with selector for custom -/// messages. -/// - stop the WebSocket connection. -/// - stop the WebSocket connection with abnormal reason. -pub opaque type WebsocketNext(user_state, user_message) { - WebsocketContinue(user_state, Option(Selector(user_message))) - WebsocketNormalStop - WebsocketAbnormalStop(reason: String) -} - -/// Instructs WebSocket connection to continue processing. -pub fn websocket_continue( - user_state: user_state, -) -> WebsocketNext(user_state, user_message) { - WebsocketContinue(user_state, None) -} - -/// Instructs WebSocket connection to continue processing, including selector -/// for custom messages. -pub fn websocket_continue_with_selector( - user_state: user_state, - selector: Selector(user_message), -) -> WebsocketNext(user_state, user_message) { - WebsocketContinue(user_state, Some(selector)) -} - -/// Instructs WebSocket connection to stop. -pub fn websocket_stop() -> WebsocketNext(user_state, user_message) { - WebsocketNormalStop -} - -/// Instructs WebSocket connection to stop with abnormal reason. -pub fn websocket_stop_abnormal( - reason: String, -) -> WebsocketNext(user_state, user_message) { - WebsocketAbnormalStop(reason) -} - -fn to_websocket_next( - next: websocket.WebsocketNext(user_state, user_message), -) -> WebsocketNext(user_state, user_message) { - case next { - websocket.Continue(user_state, selector) -> - WebsocketContinue(user_state, selector) - websocket.NormalStop -> WebsocketNormalStop - websocket.AbnormalStop(reason) -> WebsocketAbnormalStop(reason) - } -} - -fn to_internal_websocket_next( - next: WebsocketNext(user_state, user_message), -) -> websocket.WebsocketNext(user_state, user_message) { - case next { - WebsocketContinue(user_state, selector) -> - websocket.Continue(user_state, selector) - WebsocketNormalStop -> websocket.NormalStop - WebsocketAbnormalStop(reason) -> websocket.AbnormalStop(reason) - } -} - -/// Represents a WebSocket message received from the client. -pub type WebsocketMessage(user_message) { - /// Indicate that text frame has been received. - Text(String) - /// Indicate that binary frame has been received. - Binary(BitArray) - /// Indicate that user message has been received from WebSocket selector. - User(user_message) -} - -fn transform_websocket_message( - message: websocket.WebsocketMessage(user_message), -) -> Result(WebsocketMessage(user_message), Nil) { - case message { - websocket.Frame(websocks.Text(payload)) -> - Ok(Text(unsafe_to_string(payload))) - websocket.Frame(websocks.Binary(payload)) -> Ok(Binary(payload)) - websocket.UserMessage(user_message) -> Ok(User(user_message)) - _ -> Error(Nil) - } -} - -@external(erlang, "gleam_stdlib", "identity") -fn unsafe_to_string(a: BitArray) -> String - -/// Upgrade request to a WebSocket connection. If the initial request is not -/// valid for WebSocket upgrade, 400 response is sent. -/// -/// `on_init` function is called once process that handles WebSocket connection -/// is initialized. It must return a tuple with initial state and selector for -/// custom messages. If there is no custom messages, user can pass the same -/// selector from the argument -/// -/// `handler` function is called for every WebSocket message received. It must -/// return instruction on how WebSocket connection should proceed. -/// -/// `on_close` function is called when WebSocket process is going to be stopped. -pub fn upgrade_websocket( - req: Request, - on_init on_init: fn(WebsocketConnection, Selector(user_message)) -> - #(user_state, Selector(user_message)), - handler handler: fn( - WebsocketConnection, - user_state, - WebsocketMessage(user_message), - ) -> WebsocketNext(user_state, user_message), - on_close on_close: fn(WebsocketConnection, user_state) -> Nil, -) -> Response { - let handler = fn(conn, state, msg) { - transform_websocket_message(msg) - |> result.map(handler(conn, state, _)) - |> result.unwrap(websocket_continue(state)) - |> to_internal_websocket_next() - } - - let transport = req.body.transport - let socket = req.body.socket - - case ewe_http.upgrade_websocket(req, transport, socket) { - Ok(#(extensions, per_message_deflate)) -> { - let started = - factory.get_by_name(req.body.factory_name) - |> factory.start_child(fn() { - websocket.start( - transport, - socket, - on_init, - handler, - on_close, - extensions, - per_message_deflate, - ) - }) - - case started { - Ok(actor.Started(pid:, ..)) -> { - let _ = transport.controlling_process(transport, socket, pid) - response.new(200) |> response.set_body(Websocket) - } - Error(_) -> response.new(500) |> response.set_body(Empty) - } - } - Error(_) -> response.new(400) |> response.set_body(Empty) - } -} - -/// Sends a binary frame to the websocket client. -pub fn send_binary_frame( - conn: WebsocketConnection, - bits: BitArray, -) -> Result(Nil, glisten.SocketReason) { - websocks.encode_binary_frame - |> websocket.send_frame(conn.transport, conn.socket, conn.context, bits) -} - -/// Sends a text frame to the websocket client. -pub fn send_text_frame( - conn: WebsocketConnection, - text: String, -) -> Result(Nil, glisten.SocketReason) { - websocket.send_frame( - websocks.encode_text_frame, - conn.transport, - conn.socket, - conn.context, - bit_array.from_string(text), - ) -} - -/// WebSocket close codes that can be sent when closing a connection. The `data` -/// parameter allows you to include payload up to 123 bytes in size. -pub type CloseCode { - /// Standard graceful shutdown (1000). Use when connection completed - /// successfully. - NormalClosure(data: String) - /// Invalid message format (1007). Received payload that doesn't match what - /// you expected. - InvalidPayloadData(data: String) - /// Application policy violation (1008).Client broke your rules - failed - /// authentication, hit rate limits, or violated business logic. - PolicyViolation(data: String) - /// Message exceeds size limits (1009). Client sent something bigger than - /// your application allows. - MessageTooBig(data: String) - /// Server encountered unexpected error (1011). Something went wrong on your - /// side that prevents handling the connection. - InternalError(data: String) - /// Server is restarting (1012). Planned restart - clients can reconnect - /// after a bit. - ServiceRestart(data: String) - /// Temporary server overload (1013). Use when server is temporarily - /// unavailable, client should retry. - TryAgainLater(data: String) - /// Gateway/proxy received invalid response (1014). You're acting as a proxy - /// and the upstream server gave you garbage. - BadGateway(data: String) - /// Custom close codes 3000-4999 for application-specific use. - CustomCloseCode(code: Int, data: String) - /// Close without a specific reason. - NoCloseReason -} - -fn to_internal_close_code(code: CloseCode) -> websocks.CloseReason { - case code { - NormalClosure(data) -> websocks.NormalClosure(bit_array.from_string(data)) - InvalidPayloadData(data) -> - websocks.InvalidPayloadData(bit_array.from_string(data)) - PolicyViolation(data) -> - websocks.PolicyViolation(bit_array.from_string(data)) - MessageTooBig(data) -> websocks.MessageTooBig(bit_array.from_string(data)) - InternalError(data) -> websocks.InternalError(bit_array.from_string(data)) - ServiceRestart(data) -> websocks.ServiceRestart(bit_array.from_string(data)) - TryAgainLater(data) -> websocks.TryAgainLater(bit_array.from_string(data)) - BadGateway(data) -> websocks.BadGateway(bit_array.from_string(data)) - CustomCloseCode(code, data) -> - websocks.CustomCloseCode(code, bit_array.from_string(data)) - NoCloseReason -> websocks.NoCloseReason - } -} - -/// Sends a close frame to the websocket client. Once this function is called, -/// no other frames can be sent on this connection. Returns how the WebSocket -/// connection should proceed - make sure your handler returns this value. -pub fn send_close_frame( - conn: WebsocketConnection, - code: CloseCode, -) -> WebsocketNext(user_state, user_message) { - to_internal_close_code(code) - |> websocket.send_close_frame(conn.transport, conn.socket, _) - |> to_websocket_next() -} - -// SERVER-SENT EVENT -// ----------------------------------------------------------------------------- - -/// Represents a Server-Sent Events connection between a client and a server. -pub type SSEConnection = - sse.SSEConnection - -/// Represents an instruction on how Server-Sent Events connection should -/// proceed. -/// -/// - continue processing the Server-Sent Events connection. -/// - stop the Server-Sent Events connection. -/// - stop the Server-Sent Events connection with abnormal reason. -pub opaque type SSENext(user_state) { - SSEContinue(user_state) - SSENormalStop - SSEAbnormalStop(reason: String) -} - -/// Instructs Server-Sent Events connection to continue processing. -pub fn sse_continue(user_state: user_state) -> SSENext(user_state) { - SSEContinue(user_state) -} - -/// Instructs Server-Sent Events connection to stop. -pub fn sse_stop() -> SSENext(user_state) { - SSENormalStop -} - -/// Instructs Server-Sent Events connection to stop with abnormal reason. -pub fn sse_stop_abnormal(reason: String) -> SSENext(user_state) { - SSEAbnormalStop(reason) -} - -fn to_internal_sse_next(next: SSENext(user_state)) -> sse.SSENext(user_state) { - case next { - SSEContinue(user_state) -> sse.Continue(user_state) - SSENormalStop -> sse.NormalStop - SSEAbnormalStop(reason) -> sse.AbnormalStop(reason) - } -} - -/// Represents a Server-Sent Events event. The event fields are: -/// - `event`: a string identifying the type of event described. -/// - `data`: the data field for the message. -/// - `id`: event ID. -/// - `retry`: The reconnection time. If the connection to the server is lost, -/// the browser will wait for the specified time before attempting to reconnect. -/// -/// Can be created using `ewe.event` and modified with `ewe.event_name`, -/// `ewe.event_id`, and `ewe.event_retry`. -pub type SSEEvent = - sse.SSEEvent - -/// Creates a new SSE event with the given data. Use `ewe.event_name`, -/// `ewe.event_id`, and `ewe.event_retry` to modify other fields of the event. -pub fn event(data: String) -> SSEEvent { - sse.SSEEvent(event: None, data:, id: None, retry: None) -} - -/// Sets the name of the event. -pub fn event_name(event: SSEEvent, name: String) -> SSEEvent { - sse.SSEEvent(..event, event: Some(name)) -} - -/// Sets the ID of the event. -pub fn event_id(event: SSEEvent, id: String) -> SSEEvent { - sse.SSEEvent(..event, id: Some(id)) -} - -/// Sets the retry time of the event. -pub fn event_retry(event: SSEEvent, retry: Int) -> SSEEvent { - sse.SSEEvent(..event, retry: Some(retry)) -} - -/// Sets up the connection for Server-Sent Events. -/// -/// `on_init` function is called once process that handles SSE connection -/// is initialized. The argument is subject that can be used to send messages -/// to the client. It must return initial state. -/// -/// `handler` function is called for every subject's message received. It must -/// return instruction on how SSE connection should proceed. -/// -/// `on_close` function is called when SSE process is going to be stopped. -pub fn sse( - req: Request, - on_init on_init: fn(Subject(user_message)) -> user_state, - handler handler: fn(SSEConnection, user_state, user_message) -> - SSENext(user_state), - on_close on_close: fn(SSEConnection, user_state) -> Nil, -) { - let handler = fn(conn, state, msg) { - handler(conn, state, msg) - |> to_internal_sse_next() - } - - let transport = req.body.transport - let socket = req.body.socket - - case sse.send_response(transport, socket) { - Ok(Nil) -> { - let started = - factory.get_by_name(req.body.factory_name) - |> factory.start_child(fn() { - sse.start(transport, socket, on_init, handler, on_close) - }) - - case started { - Ok(actor.Started(pid:, ..)) -> { - let _ = transport.controlling_process(transport, socket, pid) - response.new(200) |> response.set_body(SSE) - } - Error(_) -> response.new(400) |> response.set_body(Empty) - } - } - Error(Nil) -> response.new(400) |> response.set_body(Empty) - } -} - -/// Sends a Server-Sent Events event to the client. -pub fn send_event( - conn: SSEConnection, - event: SSEEvent, -) -> Result(Nil, glisten.SocketReason) { - sse.send_event(conn.transport, conn.socket, event) -} +pub fn main() { + echo "hi" +} \ No newline at end of file diff --git a/src/ewe/internal/clock.gleam b/src/ewe/internal/clock.gleam deleted file mode 100644 index 74a40e9..0000000 --- a/src/ewe/internal/clock.gleam +++ /dev/null @@ -1,133 +0,0 @@ -import gleam/erlang/atom -import gleam/erlang/process -import gleam/int -import gleam/otp/actor -import gleam/result -import gleam/string -import gleam/string_tree -import logging - -/// Message that can be sent to the clock actor. -/// -type Message { - Tick -} - -/// Starts the clock application. -/// -pub fn start(_type, _args) -> Result(process.Pid, actor.StartError) { - actor.new_with_initialiser(1000, fn(subject) { - init_clock_storage() - set_http_date(calculate_http_date()) - process.send_after(subject, 1000, Tick) - - actor.initialised(subject) - |> actor.returning(subject) - |> Ok - }) - |> actor.on_message(fn(subject, _msg) { - process.send_after(subject, 1000, Tick) - - set_http_date(calculate_http_date()) - - actor.continue(subject) - }) - |> actor.start() - |> result.map(fn(started) { - let assert Ok(pid) = process.subject_owner(started.data) - pid - }) -} - -/// 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. -/// -pub fn get_http_date() -> String { - case lookup_http_date() { - Ok(date) -> date - Error(Nil) -> { - logging.log( - logging.Warning, - "Failed to lookup HTTP date, calculating new one", - ) - calculate_http_date() - } - } -} - -/// 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() - |> string_tree.append(weekday_to_string(weekday)) - |> string_tree.append(", ") - |> string_tree.append(int.to_string(day) |> string.pad_start(2, "0")) - |> string_tree.append(" ") - |> string_tree.append(month_to_string(month)) - |> string_tree.append(" ") - |> string_tree.append(int.to_string(year) |> string.pad_start(4, "0")) - |> string_tree.append(" ") - |> string_tree.append(int.to_string(hour) |> string.pad_start(2, "0")) - |> string_tree.append(":") - |> string_tree.append(int.to_string(minute) |> string.pad_start(2, "0")) - |> string_tree.append(":") - |> string_tree.append(int.to_string(second) |> string.pad_start(2, "0")) - |> string_tree.append(" GMT") - |> string_tree.to_string() -} - -/// Converts a weekday number to a string. -/// -fn weekday_to_string(weekday: Int) -> String { - case weekday { - 1 -> "Mon" - 2 -> "Tue" - 3 -> "Wed" - 4 -> "Thu" - 5 -> "Fri" - 6 -> "Sat" - 7 -> "Sun" - _ -> - panic as "erlang is breaking the fourth wall: erlang weekday outside of 1-7 range" - } -} - -/// Converts a month number to a string. -/// -fn month_to_string(month: Int) -> String { - case month { - 1 -> "Jan" - 2 -> "Feb" - 3 -> "Mar" - 4 -> "Apr" - 5 -> "May" - 6 -> "Jun" - 7 -> "Jul" - 8 -> "Aug" - 9 -> "Sep" - 10 -> "Oct" - 11 -> "Nov" - 12 -> "Dec" - _ -> - panic as "erlang is breaking the fourth wall: erlang month outside of 1-12 range" - } -} - -@external(erlang, "ewe_ffi", "now") -fn now() -> #(Int, #(Int, Int, Int), #(Int, Int, Int)) - -@external(erlang, "ewe_ffi", "init_clock_storage") -fn init_clock_storage() -> Nil - -@external(erlang, "ewe_ffi", "set_http_date") -fn set_http_date(date: String) -> Nil - -@external(erlang, "ewe_ffi", "lookup_http_date") -fn lookup_http_date() -> Result(String, Nil) diff --git a/src/ewe/internal/decoder.gleam b/src/ewe/internal/decoder.gleam deleted file mode 100644 index 837fdae..0000000 --- a/src/ewe/internal/decoder.gleam +++ /dev/null @@ -1,132 +0,0 @@ -import ewe/internal/http1/buffer -import gleam/dynamic -import gleam/http -import gleam/option - -/// Type of HTTP packet being decoded. -/// -pub type PacketType { - HttpBin - HttphBin -} - -/// Absolute path in HTTP request. -/// -pub type AbsPath { - AbsPath(BitArray) -} - -/// HTTP version as major and minor numbers. -/// -pub type Version = - #(Int, Int) - -/// HTTP packet structure. -/// -pub type HttpPacket { - HttpRequest(method: BitArray, path: AbsPath, version: Version) - HttpHeader(idx: Int, field: BitArray, value: BitArray) - HttpEoh - Http2Upgrade -} - -/// Complete packet with data and remaining bytes. -/// -pub type Packet { - Packet(HttpPacket, rest: BitArray) - More(length: option.Option(Int)) -} - -/// Decodes HTTP packets using external FFI implementation. -/// -pub fn decode_packet( - type_ type_: PacketType, - buffer buffer: buffer.Buffer, -) -> Result(Packet, dynamic.Dynamic) { - decode_packet_ffi(type_, buffer.data, []) -} - -@external(erlang, "ewe_ffi", "decode_packet") -fn decode_packet_ffi( - type_ type_: PacketType, - packet packet: BitArray, - options options: List(a), -) -> Result(Packet, dynamic.Dynamic) - -/// Decodes HTTP method from binary data. -/// -pub fn decode_method(method: BitArray) -> Result(http.Method, Nil) { - case method { - <<"GET">> -> Ok(http.Get) - <<"POST">> -> Ok(http.Post) - <<"HEAD">> -> Ok(http.Head) - <<"PUT">> -> Ok(http.Put) - <<"DELETE">> -> Ok(http.Delete) - <<"TRACE">> -> Ok(http.Trace) - <<"CONNECT">> -> Ok(http.Connect) - <<"OPTIONS">> -> Ok(http.Options) - <<"PATCH">> -> Ok(http.Patch) - _ -> Error(Nil) - } -} - -/// Maps header field indices to their string names. -/// -pub fn formatted_field_by_idx(idx: Int) -> Result(String, Nil) { - case idx { - 0 -> Error(Nil) - 1 -> Ok("cache-control") - 2 -> Ok("connection") - 3 -> Ok("date") - 4 -> Ok("pragma") - 5 -> Ok("transfer-encoding") - 6 -> Ok("upgrade") - 7 -> Ok("via") - 8 -> Ok("accept") - 9 -> Ok("accept-charset") - 10 -> Ok("accept-encoding") - 11 -> Ok("accept-language") - 12 -> Ok("authorization") - 13 -> Ok("from") - 14 -> Ok("host") - 15 -> Ok("if-modified-since") - 16 -> Ok("if-match") - 17 -> Ok("if-none-match") - 18 -> Ok("if-range") - 19 -> Ok("if-unmodified-since") - 20 -> Ok("max-forwards") - 21 -> Ok("proxy-authorization") - 22 -> Ok("range") - 23 -> Ok("referer") - 24 -> Ok("user-agent") - 25 -> Ok("age") - 26 -> Ok("location") - 27 -> Ok("proxy-authenticate") - 28 -> Ok("public") - 29 -> Ok("retry-after") - 30 -> Ok("server") - 31 -> Ok("vary") - 32 -> Ok("warning") - 33 -> Ok("www-authenticate") - 34 -> Ok("allow") - 35 -> Ok("content-base") - 36 -> Ok("content-encoding") - 37 -> Ok("content-language") - 38 -> Ok("content-length") - 39 -> Ok("content-location") - 40 -> Ok("content-md5") - 41 -> Ok("content-range") - 42 -> Ok("content-type") - 43 -> Ok("etag") - 44 -> Ok("expires") - 45 -> Ok("last-modified") - 46 -> Ok("accept-ranges") - 47 -> Ok("set-cookie") - 48 -> Ok("set-cookie2") - 49 -> Ok("x-forwarded-for") - 50 -> Ok("cookie") - 51 -> Ok("keep-alive") - 52 -> Ok("proxy-connection") - _ -> Error(Nil) - } -} diff --git a/src/ewe/internal/encoder.gleam b/src/ewe/internal/encoder.gleam deleted file mode 100644 index a69b3da..0000000 --- a/src/ewe/internal/encoder.gleam +++ /dev/null @@ -1,124 +0,0 @@ -import gleam/bytes_tree -import gleam/http/response -import gleam/int -import gleam/list - -/// Encodes an HTTP response into bytes. -/// -pub fn encode_response( - response: response.Response(BitArray), -) -> bytes_tree.BytesTree { - bytes_tree.new() - |> bytes_tree.append(encode_status_line(response.status)) - |> bytes_tree.append(encode_headers(response.headers)) - |> bytes_tree.append(response.body) -} - -/// 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() - |> bytes_tree.append(encode_status_line(response.status)) - |> bytes_tree.append(encode_headers(response.headers)) -} - -/// 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) - - <<"HTTP/1.1 ", status:utf8, " ", status_name:bits, "\r\n">> -} - -/// Encodes HTTP headers into bytes. -/// -fn encode_headers(headers: List(#(String, String))) -> BitArray { - let headers = - list.fold(headers, <<>>, fn(acc, headers) { - let #(key, value) = headers - - << - acc:bits, - sanitize_header_value(key):utf8, - ": ", - sanitize_header_value(value):utf8, - "\r\n", - >> - }) - - <> -} - -@external(erlang, "ewe_ffi", "sanitize_header_value") -fn sanitize_header_value(value: String) -> String - -/// Maps HTTP status codes to their text descriptions. -/// -fn status_to_bit_array(status: Int) -> BitArray { - case status { - 100 -> <<"Continue">> - 101 -> <<"Switching Protocols">> - 102 -> <<"Processing">> - 103 -> <<"Early Hints">> - 200 -> <<"OK">> - 201 -> <<"Created">> - 202 -> <<"Accepted">> - 203 -> <<"Non-Authoritative Information">> - 204 -> <<"No Content">> - 205 -> <<"Reset Content">> - 206 -> <<"Partial Content">> - 207 -> <<"Multi-Status">> - 208 -> <<"Already Reported">> - 226 -> <<"IM Used">> - 300 -> <<"Multiple Choices">> - 301 -> <<"Moved Permanently">> - 302 -> <<"Found">> - 303 -> <<"See Other">> - 304 -> <<"Not Modified">> - 307 -> <<"Temporary Redirect">> - 308 -> <<"Permanent Redirect">> - 400 -> <<"Bad Request">> - 401 -> <<"Unauthorized">> - 403 -> <<"Forbidden">> - 404 -> <<"Not Found">> - 405 -> <<"Method Not Allowed">> - 406 -> <<"Not Acceptable">> - 407 -> <<"Proxy Authentication Required">> - 408 -> <<"Request Timeout">> - 409 -> <<"Conflict">> - 410 -> <<"Gone">> - 411 -> <<"Length Required">> - 412 -> <<"Precondition Failed">> - 413 -> <<"Payload Too Large">> - 414 -> <<"URI Too Long">> - 415 -> <<"Unsupported Media Type">> - 416 -> <<"Range Not Satisfiable">> - 417 -> <<"Expectation Failed">> - 418 -> <<"I'm a teapot">> - 421 -> <<"Misdirected Request">> - 422 -> <<"Unprocessable Entity">> - 423 -> <<"Locked">> - 424 -> <<"Failed Dependency">> - 425 -> <<"Too Early">> - 426 -> <<"Upgrade Required">> - 428 -> <<"Precondition Required">> - 429 -> <<"Too Many Requests">> - 431 -> <<"Request Header Fields Too Large">> - 451 -> <<"Unavailable For Legal Reasons">> - 500 -> <<"Internal Server Error">> - 501 -> <<"Not Implemented">> - 502 -> <<"Bad Gateway">> - 503 -> <<"Service Unavailable">> - 504 -> <<"Gateway Timeout">> - 505 -> <<"HTTP Version Not Supported">> - 506 -> <<"Variant Also Negotiates">> - 507 -> <<"Insufficient Storage">> - 508 -> <<"Loop Detected">> - 510 -> <<"Not Extended">> - 511 -> <<"Network Authentication Required">> - _ -> <<"Unknown HTTP Status">> - } -} diff --git a/src/ewe/internal/file.gleam b/src/ewe/internal/file.gleam deleted file mode 100644 index be7c9a9..0000000 --- a/src/ewe/internal/file.gleam +++ /dev/null @@ -1,79 +0,0 @@ -import gleam/bytes_tree -import gleam/dynamic -import gleam/result -import glisten -import glisten/socket.{type Socket} -import glisten/transport.{type Transport} - -/// Reference to a file. -/// -pub type IoDevice - -/// Errors that can occur when opening a file. -/// -pub type FileError { - Enoent - Eacces - Eisdir - Eunknown(dynamic.Dynamic) -} - -/// A file. -/// -pub type File { - File(descriptor: IoDevice, size: Int) -} - -/// Errors that can occur when sending a file. -/// -pub type SendError { - FileIssue(FileError) - SocketIssue(glisten.SocketReason) -} - -/// Sends a file to the client. -/// -pub fn send( - transport: Transport, - socket: Socket, - descriptor: IoDevice, - offset: Int, - size: Int, -) -> Result(Nil, SendError) { - case transport { - transport.Tcp(..) -> { - send_file(descriptor, socket, offset, size, []) - |> result.map_error(SocketIssue) - } - transport.Ssl(..) -> { - pread(descriptor, offset, size) - |> result.map_error(FileIssue) - |> result.try(fn(bits) { - transport.send(transport, socket, bytes_tree.from_bit_array(bits)) - |> result.map_error(SocketIssue) - }) - } - } -} - -@external(erlang, "ewe_ffi", "open_file") -pub fn open(path: String) -> Result(File, FileError) - -@external(erlang, "ewe_ffi", "close_file") -pub fn close(file: IoDevice) -> Result(Nil, FileError) - -@external(erlang, "file", "sendfile") -fn send_file( - descriptor descriptor: IoDevice, - socket socket: Socket, - offset offset: Int, - bytes bytes: Int, - options options: List(a), -) -> Result(Nil, glisten.SocketReason) - -@external(erlang, "file", "pread") -fn pread( - descriptor: IoDevice, - location: Int, - number: Int, -) -> Result(BitArray, FileError) diff --git a/src/ewe/internal/handler.gleam b/src/ewe/internal/handler.gleam deleted file mode 100644 index e1c1ace..0000000 --- a/src/ewe/internal/handler.gleam +++ /dev/null @@ -1,88 +0,0 @@ -import ewe/internal/http1.{type Connection, type ResponseBody} as ewe_http -import ewe/internal/http1/handler as http1_handler -import gleam/bytes_tree -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 glisten -import glisten/transport -import logging - -/// State of the request handler. -/// -pub type Handler { - Http1(state: http1_handler.Http1Handler, self: process.Subject(Nil)) -} - -/// Initializes the request handler state. -/// -pub fn init(_) -> #(Handler, Option(process.Selector(Nil))) { - let subject = process.new_subject() - let selector = - process.new_selector() - |> process.select(subject) - - #(Http1(http1_handler.init(), self: subject), Some(selector)) -} - -/// Main loop that processes incoming messages. -/// -pub fn loop( - handler: fn(Request(Connection)) -> Response(ResponseBody), - on_crash: Response(ResponseBody), - factory_name: process.Name( - factory.Message(fn() -> Result(actor.Started(Nil), actor.StartError), Nil), - ), - idle_timeout: Int, -) -> glisten.Loop(Handler, Nil) { - fn( - 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, message { - Http1(state, self), glisten.Packet(message) -> { - let result = - http1_handler.handle_packet( - state, - conn, - message, - sender, - handler, - on_crash, - idle_timeout, - ) - - case result { - http1_handler.Continue(state) -> glisten.continue(Http1(state, self)) - http1_handler.Http2Upgrade(http1_handler.Direct(_data)) -> { - logging.log(logging.Notice, "HTTP/2 upgrade; sending goaway") - let goaway = << - 0:size(24), 4, 0, 0:size(32), 8:size(24), 7, 0, 0:size(32), - 0:size(32), 13:size(32), - >> - - let _ = - transport.send( - conn.transport, - conn.socket, - bytes_tree.from_bit_array(goaway), - ) - - glisten.stop() - } - http1_handler.Stop - | http1_handler.Http2Upgrade(http1_handler.Upgrade(..)) -> - glisten.stop() - } - } - _, _ -> glisten.stop() - } - } -} diff --git a/src/ewe/internal/http1.gleam b/src/ewe/internal/http1.gleam deleted file mode 100644 index c723f5b..0000000 --- a/src/ewe/internal/http1.gleam +++ /dev/null @@ -1,886 +0,0 @@ -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 ewe/internal/http1/buffer.{type Buffer, Buffer} -import gleam/bit_array -import gleam/bool -import gleam/bytes_tree.{type BytesTree} -import gleam/dict.{type Dict} -import gleam/erlang/atom -import gleam/erlang/process -import gleam/http -import gleam/http/request.{type Request, Request} -import gleam/http/response.{type Response} -import gleam/int -import gleam/list -import gleam/option.{None, Some} -import gleam/otp/actor -import gleam/otp/factory_supervisor as factory -import gleam/result.{replace_error, try} -import gleam/set.{type Set} -import gleam/string -import gleam/string_tree.{type StringTree} -import glisten -import glisten/socket.{type Socket} -import glisten/transport.{type Transport} -import websocks - -// 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), - ), - ) -} - -/// 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(<<>>, 0), - factory_name:, - ) -} - -/// 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.pending, max_reading_size) - - use data <- try( - transport.receive_timeout(transport, socket, read_size, 5000) - |> replace_error(on_error), - ) - let new_buffer = buffer.append(buffer, data) - - case new_buffer.pending { - 0 -> Ok(new_buffer) - _ -> read_from_socket(transport:, socket:, buffer: new_buffer, on_error:) - } -} - -// HTTP/1.1 -// ----------------------------------------------------------------------------- - -/// Errors that can occur when parsing a request. -/// -pub type ParseError { - // request line - InvalidMethod - InvalidPath - InvalidVersion - // headers - InvalidHeaders - MissingHost - DuplicateHost - InvalidContentLength - // body - InvalidBody - BodyTooLarge - // anomalies - MalformedRequest - PacketDiscard -} - -/// HTTP version enumeration. -/// -pub type HttpVersion { - Http10 - Http11 -} - -/// Result of parsing a request. -/// -pub type ParsedRequest { - Http1Request(req: Request(Connection), version: HttpVersion) - Http2Upgrade(upgrade: Http2Upgrade) -} - -/// HTTP/2 upgrade options. -/// -pub type Http2Upgrade { - Upgrade(req: Request(Connection), settings: String) - Direct(data: BitArray) -} - -/// Parses an HTTP request from the given buffer. -/// -pub fn parse_request( - conn: Connection, - buffer: Buffer, -) -> Result(ParsedRequest, ParseError) { - let transport = conn.transport - let socket = conn.socket - - case decoder.decode_packet(HttpBin, buffer) { - Ok(Packet(HttpRequest(atom_method, AbsPath(target), version), rest)) -> { - // Request Line - use method <- try( - decoder.decode_method(atom_method) - |> replace_error(InvalidMethod), - ) - - use #(path, query) <- try( - bit_array.to_string(target) - |> try(parse_path) - |> replace_error(InvalidPath), - ) - - // Headers - use #(headers, rest) <- try(parse_headers( - transport, - socket, - buffer: Buffer(rest, 0), - headers: dict.new(), - )) - - // Forming the request - let scheme = case transport { - transport.Tcp(..) -> http.Http - transport.Ssl(..) -> http.Https - } - - use host <- try( - dict.get(headers, "host") - |> result.replace_error(MissingHost), - ) - - let #(host, port) = case string.split_once(host, ":") { - Ok(#(host, port)) -> #(host, Some(port)) - Error(_) -> #(host, None) - } - - let port = - option.map(port, fn(port) { - int.parse(port) - |> result.unwrap(case scheme { - http.Http -> 80 - http.Https -> 443 - }) - }) - - let req = - Request( - method:, - headers: dict.to_list(headers), - body: Connection(..conn, buffer: Buffer(rest, 0)), - scheme:, - host:, - port:, - path:, - query:, - ) - - case version { - #(1, 0) -> Ok(Http1Request(req:, version: Http10)) - #(1, 1) -> { - let connection = dict.get(headers, "connection") - let upgrade = dict.get(headers, "upgrade") - let settings = dict.get(headers, "http2-settings") - - case connection, upgrade, settings { - Ok(connection), Ok("h2c"), Ok(settings) -> { - case string.contains(connection, "upgrade") { - True -> Ok(Http2Upgrade(Upgrade(req:, settings:))) - False -> Ok(Http1Request(req:, version: Http11)) - } - } - _, _, _ -> Ok(Http1Request(req:, version: Http11)) - } - } - _ -> Error(InvalidVersion) - } - } - Ok(Packet(decoder.Http2Upgrade, <<"\r\nSM\r\n\r\n":utf8, data:bits>>)) -> - Ok(Http2Upgrade(Direct(data:))) - Ok(More(size)) -> { - use new_buffer <- try(read_from_socket( - transport, - socket, - buffer: Buffer(buffer.data, option.unwrap(size, 0)), - on_error: MalformedRequest, - )) - - parse_request(conn, new_buffer) - } - _ -> Error(PacketDiscard) - } -} - -@external(erlang, "ewe_ffi", "parse_path") -fn parse_path(string: String) -> Result(#(String, option.Option(String)), Nil) - -/// 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) -> validate_lowercase_field(field) - }) - - use value <- try(case field { - "transfer-encoding" | "connection" | "upgrade" | "expect" | "trailer" -> - validate_lowercase_field_value(value) - _ -> validate_field_value(value) - }) - - let new_buffer = Buffer(rest, 0) - - 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(buffer.data, 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_lowercase_field") -fn validate_lowercase_field(field: BitArray) -> Result(String, ParseError) - -@external(erlang, "ewe_ffi", "validate_field_value") -fn validate_field_value(value: BitArray) -> Result(String, ParseError) - -@external(erlang, "ewe_ffi", "validate_lowercase_field_value") -fn validate_lowercase_field_value(value: BitArray) -> Result(String, ParseError) - -/// 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, -) -> Result(Request(BitArray), ParseError) { - use _ <- try(handle_continue(req)) - - let transport = req.body.transport - let socket = req.body.socket - - case request.get_header(req, "transfer-encoding") { - Ok("chunked") -> { - use #(body, rest_buffer) <- try(read_chunked_body( - transport, - socket, - req.body.buffer, - <<>>, - size_limit, - 0, - )) - - let req = request.set_body(req, body) - - case list.key_find(req.headers, "trailer") { - Ok(trailer) -> { - let set = - trailer - |> string.split(",") - |> list.fold(set.new(), fn(set, field) { - set.insert(set, string.trim(field)) - }) - - Ok(handle_trailers(req, set, rest_buffer)) - } - Error(Nil) -> Ok(req) - } - } - _ -> { - let content_length = - request.get_header(req, "content-length") - |> try(int.parse) - |> result.unwrap(0) - - use <- bool.guard(content_length > size_limit, Error(BodyTooLarge)) - - let left = content_length - bit_array.byte_size(req.body.buffer.data) - - case content_length, left { - 0, 0 -> Ok(<<>>) - 0, _l | _cl, 0 -> Ok(req.body.buffer.data) - _cl, _l -> - read_from_socket( - transport, - socket, - buffer: Buffer(req.body.buffer.data, left), - on_error: InvalidBody, - ) - |> result.map(fn(buffer) { buffer.data }) - } - |> result.map(request.set_body(req, _)) - } - } -} - -/// 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 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, - )) - - read_chunked_body( - transport:, - socket:, - buffer: new_buffer, - accumulated_body:, - body_size_limit:, - body_current_size:, - ) - } - 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) - } -} - -/// 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(rest, 0))) - [chunk_size, rest] -> { - use size <- try( - bit_array.to_string(chunk_size) - |> try(int.base_parse(_, 16)) - |> replace_error(InvalidBody), - ) - - case split(rest, <<"\r\n">>, []) { - [chunk, rest] -> { - case bit_array.byte_size(chunk) == size { - True -> Ok(Chunk(chunk, size, Buffer(rest, 0))) - False -> Error(InvalidBody) - } - } - _ -> Ok(Incomplete) - } - } - _ -> Ok(Incomplete) - } -} - -@external(erlang, "binary", "split") -fn split( - subject: BitArray, - pattern: BitArray, - options: List(atom.Atom), -) -> List(BitArray) - -/// 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), headers)) -> { - let field_name = case decoder.formatted_field_by_idx(idx) { - Ok(field_name) -> Ok(field_name) - Error(Nil) -> validate_lowercase_field(field) - } - - case field_name { - Ok(field_name) -> { - case set.contains(set, field_name) && is_allowed_trailer(field_name) { - True -> { - case validate_field_value(value) { - Ok(value) -> { - request.set_header(req, field_name, value) - |> handle_trailers(set, Buffer(headers, 0)) - } - Error(_) -> handle_trailers(req, set, Buffer(headers, 0)) - } - } - False -> handle_trailers(req, set, Buffer(headers, 0)) - } - } - Error(_) -> handle_trailers(req, set, Buffer(headers, 0)) - } - } - _ -> req - } -} - -/// Checks if a trailer field is allowed. -fn is_allowed_trailer(field: String) -> Bool { - case field { - "server-timing" | "content-digest" | "repr-digest" -> True - _ -> False - } -} - -// Streaming Body -// ----------------------------------------------------------------------------- - -/// 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 -} - -/// 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) -} - -/// Streams the request body from the socket. -/// -pub fn stream_body(req: Request(Connection)) { - use _ <- result.try( - handle_continue(req) - |> result.replace_error(InvalidBody), - ) - - case request.get_header(req, "transfer-encoding") { - Ok("chunked") -> { - let state = ChunkedStreamState(Buffer(<<>>, 0), 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 pending = content_length - bit_array.byte_size(req.body.buffer.data) - let stream_buffer = Buffer(req.body.buffer.data, int.max(0, pending)) - - do_stream_body(req, stream_buffer) - |> Ok - } - } -} - -/// Creates a consumer function that reads `N` amount of bytes from the chunked -/// request body until it is fully consumed. -fn do_stream_body_chunked( - req: Request(Connection), - chunked_stream_state: ChunkedStreamState, -) -> fn(Int) -> Result(Stream, ParseError) { - fn(size: Int) { - let read_result = - read_from_socket_until( - transport: req.body.transport, - socket: req.body.socket, - state: chunked_stream_state, - until: size, - ) - - case read_result { - Ok(#(data, ChunkedStreamState(done: True, ..))) -> - Ok(Consumed(data, fn(_) { Ok(Done) })) - Ok(#(data, state)) -> - Ok(Consumed(data, do_stream_body_chunked(req, state))) - Error(_) -> Error(InvalidBody) - } - } -} - -/// Reads data from the socket until `N` amount of bytes are read. -fn read_from_socket_until( - transport transport: Transport, - socket socket: Socket, - state state: ChunkedStreamState, - until until: Int, -) -> Result(#(BitArray, ChunkedStreamState), ParseError) { - let size = bit_array.byte_size(state.data.data) - - case state.done, size { - // Data buffer contains enough data to consume `until` bytes - _, size if size >= until -> { - let #(data, rest) = buffer.split(state.data, until) - Ok(#(data, ChunkedStreamState(..state, data: Buffer(rest, 0)))) - } - - // Accomplished the reading - True, _ -> Ok(#(state.data.data, state)) - - // Data buffer does not contain enough data to consume `until` bytes - False, _ -> { - case parse_body_chunk(state.chunk) { - Ok(FinalChunk(_)) -> - read_from_socket_until( - transport:, - socket:, - state: ChunkedStreamState( - ..state, - chunk: Buffer(<<>>, 0), - done: True, - ), - until:, - ) - Ok(Incomplete) -> { - use new_buffer <- try(read_from_socket( - transport:, - socket:, - buffer: state.chunk, - on_error: InvalidBody, - )) - - read_from_socket_until( - transport:, - socket:, - state: ChunkedStreamState(..state, chunk: new_buffer), - until:, - ) - } - Ok(Chunk(chunk, _, rest)) -> { - read_from_socket_until( - transport:, - socket:, - state: ChunkedStreamState( - ..state, - data: buffer.append(state.data, chunk), - chunk: rest, - ), - until:, - ) - } - Error(error) -> Error(error) - } - } - } -} - -/// 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) - - case buffer.pending, 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(rest, 0)))) - } - - // 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(rest, buffer.pending) - 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(<<>>, 0), - on_error: InvalidBody, - )) - - let new_buffer = - Buffer( - <>, - int.max(0, buffer.pending - bit_array.byte_size(read_buffer.data)), - ) - - let #(data, rest) = buffer.split(new_buffer, size) - Ok(Consumed(data, do_stream_body(req, Buffer(rest, 0)))) - } - } - } -} - -// 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(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") { - 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)) -} - -// Response -// ----------------------------------------------------------------------------- - -/// 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 -} - -/// 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" && 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/buffer.gleam b/src/ewe/internal/http1/buffer.gleam deleted file mode 100644 index a3c0fb4..0000000 --- a/src/ewe/internal/http1/buffer.gleam +++ /dev/null @@ -1,26 +0,0 @@ -import gleam/bit_array -import gleam/int - -/// Buffer holding data read from a socket, along with how many more bytes are -/// expected (pending) before the current read operation is complete. -/// -pub type Buffer { - Buffer(data: BitArray, pending: Int) -} - -/// Appends data to the buffer and decrements pending bytes accordingly. -/// -pub fn append(buffer: Buffer, data: BitArray) -> Buffer { - let pending = int.max(0, buffer.pending - bit_array.byte_size(data)) - Buffer(<>, pending) -} - -/// Splits the buffer data at the given byte boundary. Returns the first part -/// and the rest. -/// -pub fn split(buffer: Buffer, bytes: Int) -> #(BitArray, BitArray) { - case buffer.data { - <> -> #(partition, rest) - _ -> #(buffer.data, <<>>) - } -} diff --git a/src/ewe/internal/http1/handler.gleam b/src/ewe/internal/http1/handler.gleam deleted file mode 100644 index e814412..0000000 --- a/src/ewe/internal/http1/handler.gleam +++ /dev/null @@ -1,332 +0,0 @@ -import compresso -import ewe/internal/encoder -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 ewe/internal/http1/buffer.{Buffer} -import exception -import gleam/bit_array -import gleam/bytes_tree -import gleam/erlang/process -import gleam/http/request.{type Request} -import gleam/http/response.{type Response} -import gleam/int -import gleam/option.{type Option, None, Some} -import gleam/result -import gleam/string -import gleam/string_tree -import glisten -import glisten/internal/handler.{Close, Internal} as glisten_handler -import glisten/socket -import glisten/transport -import logging - -/// HTTP/1.1 handler state. -/// -pub type Http1Handler { - Http1Handler(idle_timer: Option(process.Timer)) -} - -/// Initializes the HTTP/1.1 handler state. -/// -pub fn init() -> Http1Handler { - Http1Handler(idle_timer: None) -} - -/// Action to take after handling a packet. -/// -pub type Next { - Continue(state: Http1Handler) - Stop - Http2Upgrade(upgrade: Http2Upgrade) -} - -/// HTTP/2 upgrade options. -/// -pub type Http2Upgrade { - Upgrade(request: Request(Connection), settings: String) - Direct(data: BitArray) -} - -/// Handles received glisten packet. -/// -pub fn handle_packet( - 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, -) -> Next { - case state.idle_timer { - Some(timer) -> process.cancel_timer(timer) - None -> process.TimerNotFound - } - - case ewe_http.parse_request(connection, Buffer(data, 0)) { - Ok(ewe_http.Http1Request(request, version)) -> { - let call_result = - call(request, version, glisten_subject, handler, on_crash, idle_timeout) - - case call_result { - Ok(state) -> Continue(state) - Error(Nil) -> Stop - } - } - Ok(ewe_http.Http2Upgrade(ewe_http.Direct(data))) -> - Http2Upgrade(Direct(data:)) - Ok(ewe_http.Http2Upgrade(ewe_http.Upgrade(request, _settings))) -> { - logging.log(logging.Notice, "HTTP/2 upgrade; using HTTP/1.1") - let call_result = - call( - request, - ewe_http.Http11, - glisten_subject, - handler, - on_crash, - idle_timeout, - ) - - case call_result { - Ok(state) -> Continue(state) - Error(Nil) -> Stop - } - } - Error(reason) -> { - let #(status, message) = case reason { - ewe_http.InvalidMethod -> #( - 400, - "Rejected HTTP request with invalid method", - ) - ewe_http.InvalidPath -> #( - 400, - "Rejected HTTP request with invalid path", - ) - ewe_http.InvalidVersion -> #( - 505, - "Rejected HTTP request with unsupported version", - ) - ewe_http.InvalidHeaders -> #( - 400, - "Rejected HTTP request with malformed headers", - ) - ewe_http.MissingHost -> #( - 400, - "Rejected HTTP request with missing Host header", - ) - ewe_http.DuplicateHost -> #( - 400, - "Rejected HTTP request with duplicate Host header", - ) - ewe_http.InvalidContentLength -> #( - 400, - "Rejected HTTP request with invalid Content-Length", - ) - ewe_http.InvalidBody -> #( - 400, - "Rejected HTTP request with malformed body", - ) - ewe_http.BodyTooLarge -> #( - 400, - "Rejected HTTP request with body exceeding size limit", - ) - ewe_http.MalformedRequest -> #( - 400, - "Rejected HTTP request due to malformed packet", - ) - ewe_http.PacketDiscard -> #( - 400, - "Rejected HTTP request due to unrecognized packet", - ) - } - - logging.log(logging.Warning, message) - - let _ = - response.new(status) - |> response.set_body(<<>>) - |> response.set_header("connection", "close") - |> encoder.encode_response() - |> transport.send(connection.transport, connection.socket, _) - - Stop - } - } -} - -/// 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(Http1Handler, Nil) { - let response = case exception.rescue(fn() { handler(request) }) { - Ok(response) -> response - Error(_exception) -> { - logging.log(logging.Error, "Caught crash in request handler") - response.set_header(on_crash, "connection", "close") - } - } - - case response.body { - Websocket | SSE | Chunked -> Error(Nil) - File(descriptor, offset, size) -> - send_file(request, version, response, descriptor, offset, size) - |> on_sent(response, glisten_subject, idle_timeout) - - _ -> - 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) - } -} - -/// 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 response = case response.get_header(response, "content-length") { - Ok(_) -> response - Error(Nil) -> - response.set_header(response, "content-length", int.to_string(size)) - } - - let sent = - 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( - request.body.transport, - request.body.socket, - descriptor, - offset, - size, - ) - |> result.replace_error(socket.Badarg) - }) - - let _ = file.close(descriptor) - - sent -} - -/// Sends a body to the client. -/// -fn send_body( - request: Request(Connection), - version: HttpVersion, - response: Response(ResponseBody), -) -> Result(Nil, glisten.SocketReason) { - 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 - BitsData(bits) -> bits - BytesData(bytes) -> bytes_tree.to_bit_array(bytes) - Empty -> <<>> - _ -> panic - } - - let content_length = bit_array.byte_size(bits) - let response = case content_length > 1024 { - True -> - case can_encode_gzip(request, response) { - True -> { - let compressed = compresso.gzip(bits) - let content_length = bit_array.byte_size(compressed) - - remove_charset(response) - |> response.set_header("content-encoding", "gzip") - |> response.set_header("vary", "Accept-Encoding") - |> response.set_header( - "content-length", - int.to_string(content_length), - ) - |> response.set_body(compressed) - } - _ -> - response.set_body(response, bits) - |> response.set_header( - "content-length", - int.to_string(content_length), - ) - } - False -> - response.set_body(response, bits) - |> response.set_header("content-length", int.to_string(content_length)) - } - - ewe_http.append_default_headers(response, request, version) - |> encoder.encode_response() - |> transport.send(request.body.transport, request.body.socket, _) -} - -/// 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 content_encoding = response.get_header(response, "content-encoding") - - case accept_encoding, content_encoding { - Ok(True), Error(Nil) -> True - _, _ -> False - } -} - -/// 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(response, "content-type", parts.0) - }) - |> result.unwrap(response) -} - -/// 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 deleted file mode 100644 index 8b9c5d0..0000000 --- a/src/ewe/internal/stream/chunked.gleam +++ /dev/null @@ -1,130 +0,0 @@ -import ewe/internal/encoder -import gleam/bit_array -import gleam/bytes_tree -import gleam/erlang/process.{type Subject} -import gleam/http/response.{type Response} -import gleam/otp/actor -import gleam/result -import glisten -import glisten/socket.{type Socket} -import glisten/transport.{type Transport} -import logging - -/// Sends a response for a chunked transfer encoding. -/// -pub fn send_response( - resp: Response(a), - transport: Transport, - socket: Socket, -) -> Result(Nil, Nil) { - case response.get_header(resp, "transfer-encoding") { - Ok("chunked") -> resp - _ -> response.set_header(resp, "transfer-encoding", "chunked") - } - |> encoder.encode_response_partially() - |> transport.send(transport, socket, _) - |> result.replace_error(Nil) -} - -/// 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, - on_init: fn(Subject(user_message)) -> user_state, - handler: fn(ChunkedBody, user_state, user_message) -> ChunkedNext(user_state), - on_close: fn(ChunkedBody, user_state) -> Nil, -) -> Result(actor.Started(Nil), actor.StartError) { - actor.new_with_initialiser(1000, fn(_subject) { - let subject = process.new_subject() - let state = on_init(subject) - - let selector = - process.new_selector() - |> process.select(subject) - - actor.initialised(state) - |> actor.returning(Nil) - |> actor.selecting(selector) - |> Ok - }) - |> actor.on_message(fn(state, message) { - let conn = ChunkedBody(transport, socket) - - case handler(conn, state, message) { - Continue(new_state) -> actor.continue(new_state) - NormalStop -> { - case send_end(transport, socket) { - Ok(Nil) -> { - on_close(conn, state) - actor.stop() - } - Error(socket_reason) -> { - let message = - "Failed to send chunked response terminator: " - <> socket.reason_to_string(socket_reason) - - logging.log(logging.Warning, message) - on_close(conn, state) - actor.stop_abnormal(message) - } - } - } - AbnormalStop(reason) -> { - logging.log(logging.Warning, "Chunked response stopped: " <> reason) - on_close(conn, state) - actor.stop_abnormal(reason) - } - } - }) - |> actor.start() -} - -/// Sends the end marker for chunked transfer encoding. -/// -fn send_end( - transport: Transport, - socket: Socket, -) -> Result(Nil, glisten.SocketReason) { - transport.send(transport, socket, bytes_tree.from_bit_array(<<"0\r\n\r\n">>)) -} - -/// Sends a chunk to the client. -/// -pub fn send_chunk( - 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, _) -} - -/// 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. -/// -@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 deleted file mode 100644 index 00994e6..0000000 --- a/src/ewe/internal/stream/sse.gleam +++ /dev/null @@ -1,158 +0,0 @@ -import ewe/internal/encoder -import gleam/bytes_tree -import gleam/erlang/atom -import gleam/erlang/process.{type Selector, type Subject} -import gleam/http/response -import gleam/int -import gleam/list -import gleam/option.{type Option} -import gleam/otp/actor -import gleam/result -import gleam/string_tree -import glisten/socket.{type Socket} -import glisten/socket/options.{Active, ActiveMode} -import glisten/transport.{type Transport} - -/// Sends a response for a Server-Sent Events connection. -/// -pub fn send_response(transport: Transport, socket: Socket) -> Result(Nil, Nil) { - response.new(200) - |> response.set_header("content-type", "text/event-stream") - |> 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. -/// -pub type SSEConnection { - SSEConnection(transport: Transport, socket: Socket) -} - -/// Represents an instruction on how Server-Sent Events connection should proceed. -/// -pub type SSENext(user_state) { - Continue(user_state) - NormalStop - AbnormalStop(reason: String) -} - -/// Represents a message that can be sent to or received from the Server-Sent -/// Events connection. -/// -pub type SSEMessages(user_message) { - User(user_message) - Close -} - -/// Starts a new Server-Sent Events connection. -/// -pub fn start( - transport: Transport, - socket: Socket, - on_init: fn(Subject(user_message)) -> user_state, - handler: fn(SSEConnection, user_state, user_message) -> SSENext(user_state), - on_close: fn(SSEConnection, user_state) -> Nil, -) -> Result(actor.Started(Nil), actor.StartError) { - actor.new_with_initialiser(1000, fn(_self) { - let _ = transport.set_opts(transport, socket, [ActiveMode(Active)]) - - let subject = process.new_subject() - let state = on_init(subject) - let selector = create_socket_selector(subject) - - actor.initialised(state) - |> actor.returning(Nil) - |> actor.selecting(selector) - |> Ok - }) - |> actor.on_message(fn(state, message) { - case message { - User(message) -> { - let conn = SSEConnection(transport, socket) - case handler(conn, state, message) { - Continue(new_state) -> actor.continue(new_state) - NormalStop -> { - on_close(conn, state) - actor.stop() - } - AbnormalStop(reason) -> { - on_close(conn, state) - actor.stop_abnormal(reason) - } - } - } - Close -> { - on_close(SSEConnection(transport, socket), state) - actor.stop() - } - } - }) - |> actor.start() -} - -/// 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 }) -} - -/// 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, - event: SSEEvent, -) -> Result(Nil, socket.SocketReason) { - let id = - option.map(event.id, format("id", _)) - |> option.unwrap("") - - let retry = - option.map(event.retry, int.to_string) - |> option.map(format("retry", _)) - |> option.unwrap("") - - let data = - string_tree.from_string(event.data) - |> string_tree.split("\n") - |> list.map(string_tree.prepend(_, "data: ")) - |> string_tree.join("\n") - - let event = - option.map(event.event, format("event", _)) - |> option.unwrap("") - - string_tree.new() - |> string_tree.append(event) - |> string_tree.append(id) - |> string_tree.append(retry) - |> string_tree.append_tree(data) - |> string_tree.append("\n\n") - |> bytes_tree.from_string_tree() - |> transport.send(transport, socket, _) -} - -/// Formats a field and value for a Server-Sent Events event. -/// -fn format(field: String, value: String) { - field <> ": " <> value <> "\n" -} diff --git a/src/ewe/internal/stream/websocket.gleam b/src/ewe/internal/stream/websocket.gleam deleted file mode 100644 index 795c1ce..0000000 --- a/src/ewe/internal/stream/websocket.gleam +++ /dev/null @@ -1,442 +0,0 @@ -import exception -import gleam/bit_array -import gleam/bytes_tree -import gleam/dynamic -import gleam/erlang/atom -import gleam/erlang/process.{type Selector} -import gleam/option.{type Option, None, Some} -import gleam/otp/actor -import glisten/socket.{type Socket, type SocketReason} -import glisten/socket/options.{ActiveMode, Count} -import glisten/transport.{type Transport} -import logging -import websocks - -/// Represents a WebSocket connection. -/// -pub type WebsocketConnection { - WebsocketConnection( - transport: Transport, - socket: Socket, - context: websocks.Context, - ) -} - -/// 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. -/// -pub type WebsocketNext(user_state, user_message) { - Continue(user_state: user_state, selector: Option(Selector(user_message))) - NormalStop - AbnormalStop(reason: String) -} - -// 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 ActorNext(user_state, user_message) = - actor.Next(WebsocketState(user_state), InternalMessage(user_message)) - -// Internal messages used by the WebSocket actor. -// -type InternalMessage(user_message) { - Packet(BitArray) - Close - TcpPassive - User(user_message) - Invalid -} - -// 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. -// -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. -// -type OnClose(user_state) = - fn(WebsocketConnection, user_state) -> Nil - -// 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" - -// Active count for socket. -// -const socket_active_count = 100 - -/// Starts a new WebSocket connection. -/// -pub fn start( - transport: Transport, - socket: Socket, - on_init: OnInit(user_state, user_message), - handler: Handler(user_state, user_message), - on_close: OnClose(user_state), - extensions: List(String), - permessage_deflate: Bool, -) -> Result(actor.Started(Nil), actor.StartError) { - actor.new_with_initialiser(1000, fn(_self) { - let _ = - transport.set_opts(transport, socket, [ - ActiveMode(Count(socket_active_count)), - ]) - - let compression = case permessage_deflate { - True -> Some(websocks.get_compression_extensions(extensions)) - False -> None - } - - let context = websocks.create_context(compression, websocks.Server) - - let #(user_state, user_selector) = - WebsocketConnection(transport, socket, context) - |> on_init(process.new_selector()) - - let selector = - process.map_selector(user_selector, User) - |> process.merge_selector(create_socket_selector()) - - WebsocketState(user_state:, context:) - |> actor.initialised() - |> actor.selecting(selector) - |> actor.returning(Nil) - |> Ok - }) - |> actor.on_message(fn(state, msg) { - case msg { - Packet(data) -> - handle_valid_packet(transport, socket, state, data, handler, on_close) - User(user_message) -> - handle_user_message( - transport, - socket, - state, - user_message, - handler, - on_close, - ) - Close -> { - let conn = WebsocketConnection(transport, socket, state.context) - handle_close(on_close, state, conn, None) - } - Invalid -> { - let conn = WebsocketConnection(transport, socket, state.context) - handle_close(on_close, state, conn, Some(malformed)) - } - TcpPassive -> { - let _ = - transport.set_opts(transport, socket, [ - ActiveMode(Count(socket_active_count)), - ]) - actor.continue(state) - } - } - }) - |> actor.start() - // |> result.map(after_start(_, transport, socket)) -} - -// Creates selector for glisten socket events. -// -fn create_socket_selector() -> Selector(InternalMessage(user_message)) { - process.new_selector() - |> process.select_record(atom.create("tcp"), 2, fn(record) { - Packet(coerce_tcp_message(record)) - }) - |> process.select_record(atom.create("ssl"), 2, fn(record) { - Packet(coerce_tcp_message(record)) - }) - |> process.select_record(atom.create("tcp_closed"), 1, fn(_) { Close }) - |> process.select_record(atom.create("ssl_closed"), 1, fn(_) { Close }) - |> process.select_record(atom.create("tcp_passive"), 1, fn(_) { TcpPassive }) -} - -@external(erlang, "ewe_ffi", "coerce_tcp_message") -fn coerce_tcp_message(record: dynamic.Dynamic) -> BitArray - -// Handles incoming packet data, decoding frames and processing them. -// -fn handle_valid_packet( - transport: Transport, - socket: Socket, - state: WebsocketState(user_state), - data: BitArray, - handler: Handler(user_state, user_message), - on_close: OnClose(user_state), -) -> ActorNext(user_state, user_message) { - let conn = WebsocketConnection(transport, socket, state.context) - let processed = - websocks.process_incoming_frames( - data, - state.context, - ResolveState( - socket:, - transport:, - handler:, - next: Continue(state.user_state, None), - ), - handle_frame, - ) - - case processed { - Ok(#(resolved_state, context)) -> { - case resolved_state.next { - Continue(user_state, selector) -> { - let next = actor.continue(WebsocketState(user_state:, context:)) - - case selector { - Some(selector) -> actor.with_selector(next, selector) - None -> next - } - } - NormalStop -> handle_close(on_close, state, conn, None) - AbnormalStop(reason) -> - handle_close(on_close, state, conn, Some(reason)) - } - } - Error(_violation) -> handle_close(on_close, state, conn, Some(malformed)) - } -} - -// Represents the state of the WebSocket connection when resolving frames. -// -type ResolveState(user_state, user_message) { - ResolveState( - socket: Socket, - transport: Transport, - handler: Handler(user_state, user_message), - next: WebsocketNext(user_state, InternalMessage(user_message)), - ) -} - -/// Processes a list of frames sequentially. -/// -fn handle_frame( - state: ResolveState(user_state, user_message), - context: websocks.Context, - frame: websocks.Frame, -) -> websocks.ResolveNext(ResolveState(user_state, user_message)) { - case frame { - websocks.Control(websocks.Ping(payload)) -> { - case bit_array.byte_size(payload) { - size if size > 125 -> - websocks.Stop( - ResolveState( - ..state, - next: AbnormalStop( - "control frames are only allowed to have payload up to and including 125 octets", - ), - ), - ) - _ -> { - let sent = - transport.send( - state.transport, - state.socket, - websocks.encode_pong_frame(payload, None) - |> bytes_tree.from_bit_array(), - ) - - case sent { - Ok(Nil) -> websocks.Continue(state) - Error(_) -> - websocks.Stop( - ResolveState(..state, next: AbnormalStop(failed_pong)), - ) - } - } - } - } - - websocks.Control(websocks.Close(reason)) -> { - let _ = - transport.send( - state.transport, - state.socket, - websocks.encode_close_frame(reason, None) - |> bytes_tree.from_bit_array(), - ) - - websocks.Stop(ResolveState(..state, next: NormalStop)) - } - - frame -> { - let assert Continue(user_state, selector) = state.next - - let conn = WebsocketConnection(state.transport, state.socket, context) - - let call = - exception.rescue(fn() { state.handler(conn, user_state, Frame(frame)) }) - - case call { - Ok(Continue(user_state, new_selector)) -> { - let next_selector = - option.map(new_selector, process.map_selector(_, User)) - |> option.or(selector) - |> option.map(process.merge_selector(create_socket_selector(), _)) - - websocks.Continue( - ResolveState(..state, next: Continue(user_state, next_selector)), - ) - } - Ok(NormalStop) -> websocks.Stop(ResolveState(..state, next: NormalStop)) - Ok(AbnormalStop(reason)) -> - websocks.Stop(ResolveState(..state, next: AbnormalStop(reason))) - Error(_) -> - websocks.Stop(ResolveState(..state, next: AbnormalStop(crashed))) - } - } - } -} - -// Handles user messages sent to the WebSocket. -// -fn handle_user_message( - transport: Transport, - socket: Socket, - state: WebsocketState(user_state), - user_message: user_message, - handler: Handler(user_state, user_message), - on_close: OnClose(user_state), -) -> ActorNext(user_state, user_message) { - let conn = WebsocketConnection(transport, socket, state.context) - let call = - exception.rescue(fn() { - handler(conn, state.user_state, UserMessage(user_message)) - }) - - case call { - Ok(Continue(new_user_state, new_selector)) -> { - let next_selector = - option.map(new_selector, process.map_selector(_, User)) - |> option.map(process.merge_selector(create_socket_selector(), _)) - - let next = - actor.continue(WebsocketState(..state, user_state: new_user_state)) - - case next_selector { - Some(selector) -> actor.with_selector(next, selector) - None -> next - } - } - Ok(NormalStop) -> handle_close(on_close, state, conn, None) - Ok(AbnormalStop(reason)) -> - handle_close(on_close, state, conn, Some(reason)) - Error(_) -> handle_close(on_close, state, conn, Some(crashed)) - } -} - -// Handles WebSocket connection closure. -// -fn handle_close( - on_close: OnClose(user_state), - state: WebsocketState(user_state), - conn: WebsocketConnection, - abnormal_reason: Option(String), -) -> actor.Next(WebsocketState(user_state), InternalMessage(user_message)) { - websocks.close_context(state.context) - on_close(conn, state.user_state) - - case abnormal_reason { - Some(reason) -> { - let level = case reason == crashed { - True -> logging.Error - False -> logging.Warning - } - logging.log(level, "WebSocket closed: " <> reason) - actor.stop_abnormal(reason) - } - None -> actor.stop() - } -} - -/// 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(_socket_reason) -> { - logging.log( - logging.Error, - "Frame should be sent from the WebSocket connection, but was sent from different process.", - ) - 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( - "Socket error occured while trying to send close frame: " - <> socket.reason_to_string(reason), - ) - Error(_reason) -> { - logging.log( - logging.Error, - "Frame should be sent from the WebSocket connection, but was sent from different process.", - ) - - panic as non_owning_process - } - } -} diff --git a/src/ewe_ffi.erl b/src/ewe_ffi.erl deleted file mode 100644 index ff913e6..0000000 --- a/src/ewe_ffi.erl +++ /dev/null @@ -1,199 +0,0 @@ --module(ewe_ffi). - --export([close_file/1, decode_packet/3, init_clock_storage/0, lookup_http_date/0, now/0, - now_microseconds/0, open_file/1, set_http_date/1, validate_lowercase_field/1, - validate_field_value/1, validate_lowercase_field_value/1, sanitize_header_value/1, - coerce_tcp_message/1, parse_path/1]). - -% Socket -% ----------------------------------------------------------------------------- - -coerce_tcp_message({tcp, _Socket, Data}) -> - Data; -coerce_tcp_message({ssl, _Socket, Data}) -> - Data. - -% HTTP -% ----------------------------------------------------------------------------- - -decode_packet(Type, Packet, Options) -> - case erlang:decode_packet(Type, Packet, Options) of - {ok, {http_request, <<"PRI">>, '*', {2, 0}}, Rest} -> - {ok, {packet, http2_upgrade, Rest}}; - {ok, {http_request, Method, Uri, Version}, Rest} -> - MethodBin = - if is_atom(Method) -> - atom_to_binary(Method); - true -> - Method - end, - {ok, {packet, {http_request, MethodBin, Uri, Version}, Rest}}; - {ok, {http_header, Idx, _, Field, Value}, Rest} -> - {ok, {packet, {http_header, Idx, Field, Value}, Rest}}; - {ok, Bin, Rest} -> - {ok, {packet, Bin, Rest}}; - {more, undefined} -> - {ok, {more, none}}; - {more, Length} -> - {ok, {more, {some, Length}}}; - {error, Reason} -> - {error, Reason} - end. - -parse_path(Value) -> - case uri_string:parse(Value) of - {error, _, _} -> - {error, nil}; - Uri -> - Query = - try - {some, maps:get(query, Uri)} - catch - _:_ -> - none - end, - {ok, {maps:get(path, Uri), Query}} - end. - -validate_lowercase_field(<<>>) -> - {error, invalid_headers}; -validate_lowercase_field(Value) -> - validate_lowercase_field(Value, <<>>). - -validate_lowercase_field(<<>>, Acc) -> - {ok, Acc}; -validate_lowercase_field(<>, Acc) when C >= $A, C =< $Z -> - validate_lowercase_field(Rest, <>); -validate_lowercase_field(<>, Acc) - when C >= $a, C =< $z; - C >= $0, C =< $9; - C =:= $!; - C =:= $#; - C =:= $$; - C =:= $%; - C =:= $&; - C =:= $'; - C =:= $*; - C =:= $+; - C =:= $-; - C =:= $.; - C =:= $^; - C =:= $_; - C =:= $`; - C =:= $|; - C =:= $~ -> - validate_lowercase_field(Rest, <>); -validate_lowercase_field(_, _) -> - {error, invalid_headers}. - -validate_field_value(Value) -> - case do_validate_field_value(Value) of - true -> - {ok, Value}; - false -> - {error, invalid_headers} - end. - -% HTTP field values can contain: -% - VCHAR: 0x21-0x7E (visible ASCII characters) -% - WSP: 0x20 (space), 0x09 (tab) -% - obs-text: 0x80-0xFF (for backward compatibility) -% Invalid: control characters 0x00-0x08, 0x0A-0x1F, 0x7F -do_validate_field_value(Value) -> - case Value of - <<>> -> - true; - <> - when C =:= 16#09 - orelse C >= 16#20 andalso C =< 16#7E - orelse C >= 16#80 andalso C =< 16#FF -> - do_validate_field_value(Rest); - _ -> - false - end. - -validate_lowercase_field_value(Value) -> - do_validate_lowercase_field_value(Value, <<>>). - -do_validate_lowercase_field_value(<<>>, Acc) -> - {ok, Acc}; -do_validate_lowercase_field_value(<>, Acc) when C >= $A, C =< $Z -> - do_validate_lowercase_field_value(Rest, <>); -do_validate_lowercase_field_value(<>, Acc) - when C =:= 16#09 - orelse C >= 16#20 andalso C =< 16#7E - orelse C >= 16#80 andalso C =< 16#FF -> - do_validate_lowercase_field_value(Rest, <>); -do_validate_lowercase_field_value(_, _) -> - {error, invalid_headers}. - -sanitize_header_value(Value) -> - sanitize_header_value(Value, <<>>). - -sanitize_header_value(<<>>, Acc) -> - Acc; -sanitize_header_value(<>, Acc) when C =:= 16#0D; C =:= 16#0A -> - sanitize_header_value(Rest, Acc); -sanitize_header_value(<>, Acc) -> - sanitize_header_value(Rest, <>). - -% CLOCK -% ----------------------------------------------------------------------------- - -now() -> - Timestamp = os:system_time(microsecond), - {Date, Time} = calendar:system_time_to_universal_time(Timestamp, microsecond), - Weekday = calendar:day_of_the_week(Date), - {Weekday, Date, Time}. - -now_microseconds() -> - os:system_time(microsecond). - -init_clock_storage() -> - ets:new(ewe_clock, [set, protected, named_table, {read_concurrency, true}]). - -set_http_date(Value) -> - ets:insert(ewe_clock, {http_date, Value}). - -lookup_http_date() -> - try - {ok, ets:lookup_element(ewe_clock, http_date, 2)} - catch - _:badarg -> - {error, nil} - end. - -% FILES -% ----------------------------------------------------------------------------- - -open_file(Path) -> - case file:open(Path, [binary, raw]) of - {ok, IoDevice} -> - {ok, {file, IoDevice, filelib:file_size(Path)}}; - {error, enoent} -> - {error, enoent}; - {error, eacces} -> - {error, eacces}; - {error, eisdir} -> - {error, eisdir}; - {error, enotdir} -> - {error, enoent}; - {error, Err} -> - {error, {eunknown, Err}} - end. - -close_file(File) -> - case file:close(File) of - ok -> - {ok, nil}; - {error, enoent} -> - {error, enoent}; - {error, eacces} -> - {error, eacces}; - {error, eisdir} -> - {error, eisdir}; - {error, enotdir} -> - {error, enoent}; - {error, _} -> - {error, eunknown} - end. diff --git a/test/client/http.gleam b/test/client/http.gleam deleted file mode 100644 index 32aadab..0000000 --- a/test/client/http.gleam +++ /dev/null @@ -1,67 +0,0 @@ -import gleam/bit_array -import gleam/http/response.{type Response, Response} -import gleam/int -import gleam/list -import gleam/result -import gleam/string - -pub type ParseError { - InvalidStatus - InvalidHeaders - MalformedResponse -} - -pub fn parse(data: BitArray) -> Result(Response(String), ParseError) { - case bit_array.to_string(data) { - Ok(data) -> - case string.split_once(data, "\r\n\r\n") { - Ok(#(lines, body)) -> do_parse(lines, body) - Error(Nil) -> do_parse(data, "") - } - - Error(Nil) -> Error(MalformedResponse) - } -} - -fn do_parse( - lines: String, - body: String, -) -> Result(Response(String), ParseError) { - case string.split(lines, "\r\n") { - [status, ..headers] -> { - use status <- result.try(parse_status(status)) - use headers <- result.try(parse_headers(headers, [])) - - Ok(Response(status:, headers:, body:)) - } - _ -> Error(MalformedResponse) - } -} - -fn parse_status(status: String) -> Result(Int, ParseError) { - case string.split(status, " ") { - [_version, status, ..] -> - int.parse(status) - |> result.replace_error(InvalidStatus) - _ -> Error(InvalidStatus) - } -} - -fn parse_headers( - headers: List(String), - acc: List(#(String, String)), -) -> Result(List(#(String, String)), ParseError) { - case headers { - [] -> Ok(list.reverse(acc)) - [header, ..remaining] -> { - case string.split_once(header, ": "), string.split_once(header, ":") { - Ok(#(name, value)), _ -> - parse_headers(remaining, [#(string.lowercase(name), value), ..acc]) - Error(Nil), Ok(#(name, value)) -> - [#(string.lowercase(name), string.trim(value)), ..acc] - |> parse_headers(remaining, _) - Error(Nil), Error(Nil) -> Error(InvalidHeaders) - } - } - } -} diff --git a/test/client/tcp.gleam b/test/client/tcp.gleam deleted file mode 100644 index 45fb1f3..0000000 --- a/test/client/tcp.gleam +++ /dev/null @@ -1,34 +0,0 @@ -import gleam/dynamic -import gleam/erlang/atom -import gleam/erlang/charlist -import glisten/socket -import glisten/tcp - -pub fn with_socket( - port port: Int, - active active: Bool, - callback callback: fn(socket.Socket) -> a, -) -> a { - let assert Ok(socket) = - tcp_connect(charlist.from_string("localhost"), port, [ - from(atom.create("binary")), - from(#(atom.create("active"), active)), - ]) - - let result = callback(socket) - - let assert Ok(Nil) = tcp.close(socket) - - result -} - -@external(erlang, "gen_tcp", "connect") -fn tcp_connect( - host: charlist.Charlist, - port: Int, - options: List(dynamic.Dynamic), -) -> Result(socket.Socket, Nil) - -// https://github.com/rawhat/glisten/blob/master/test/tcp_client.gleam#L13C1-L14C29 -@external(erlang, "gleam@function", "identity") -fn from(value: a) -> dynamic.Dynamic diff --git a/test/h1spec_test.gleam b/test/h1spec_test.gleam deleted file mode 100644 index f714b2b..0000000 --- a/test/h1spec_test.gleam +++ /dev/null @@ -1,233 +0,0 @@ -import client/http -import client/tcp as client -import gleam/bytes_tree -import gleam/http/response.{type Response} -import gleam/list -import glisten/socket.{type Socket} -import glisten/tcp -import server - -fn expect_timeout(req: String) -> Nil { - let socket_address = server.start(server.echoer()) - use socket <- client.with_socket(socket_address.port, active: False) - - let assert Ok(Nil) = tcp.send(socket, bytes_tree.from_string(req)) - - assert tcp.receive_timeout(socket, 0, 500) == Error(socket.Timeout) -} - -fn run_request(socket: Socket, req: String) -> Response(String) { - let assert Ok(Nil) = tcp.send(socket, bytes_tree.from_string(req)) - - let assert Ok(resp) = tcp.receive_timeout(socket, 0, 1000) - let assert Ok(resp) = http.parse(resp) - - resp -} - -pub fn expect_status( - req: String, - status: List(#(Int, Int)), -) -> Response(String) { - let socket_address = server.start(server.echoer()) - use socket <- client.with_socket(socket_address.port, active: False) - - let resp = run_request(socket, req) - - assert list.fold_until(over: status, from: False, with: fn(_, status) { - let #(start, end) = status - - case resp.status { - status if status >= start && status <= end -> list.Stop(True) - _ -> list.Continue(False) - } - }) - == True - - resp -} - -pub fn fragmented_method_test() { - expect_timeout("G") -} - -pub fn fragmented_url_1_test() { - expect_timeout("GET ") -} - -pub fn fragmented_url_2_test() { - expect_timeout("GET /hello") -} - -pub fn fragmented_url_3_test() { - expect_timeout("GET /hello ") -} - -pub fn fragmented_http_version_test() { - expect_timeout("GET /hello HTTP") -} - -pub fn fragmented_request_line_test() { - expect_timeout("GET /hello HTTP/1.1") -} - -pub fn fragmented_request_line_newline_1_test() { - expect_timeout("GET /hello HTTP/1.1\r") -} - -pub fn fragmented_request_line_newline_2_test() { - expect_timeout("GET /hello HTTP/1.1\r\n") -} - -pub fn fragmented_field_name_test() { - expect_timeout("GET /hello HTTP/1.1\r\nHos") -} - -pub fn fragmented_field_value_1_test() { - expect_timeout("GET /hello HTTP/1.1\r\nHost:") -} - -pub fn fragmented_field_value_2_test() { - expect_timeout("GET /hello HTTP/1.1\r\nHost: ") -} - -pub fn fragmented_field_value_3_test() { - expect_timeout("GET /hello HTTP/1.1\r\nHost: localhost") -} - -pub fn fragmented_field_value_4_test() { - expect_timeout("GET /hello HTTP/1.1\r\nHost: localhost\r") -} - -pub fn fragmented_request_test() { - expect_timeout("GET /hello HTTP/1.1\r\nHost: localhost\r\n") -} - -pub fn fragmented_request_termination_test() { - expect_timeout("GET /hello HTTP/1.1\r\nHost: localhost\r\n\r") -} - -pub fn request_without_http_version_test() { - expect_status("GET / \r\n\r\n", [#(400, 599)]) -} - -pub fn request_with_expect_header_test() { - expect_status( - "GET / HTTP/1.1\r\nHost: example.com\r\nExpect: 100-continue\r\n\r\n", - [#(100, 100), #(200, 299)], - ) -} - -pub fn valid_get_request_test() { - expect_status("GET / HTTP/1.1\r\nHost: example.com\r\n\r\n", [#(200, 299)]) -} - -pub fn valid_get_request_with_edge_case() { - expect_status("GET / HTTP/1.1\r\nhoSt:\texample.com\r\nempty:\r\n\r\n", [ - #(200, 299), - ]) -} - -pub fn invalid_header_characters_test() { - expect_status( - "GET / HTTP/1.1\r\nHost: example.com\r\nX-Invalid[]: test\r\n\r\n", - [#(400, 499)], - ) -} - -pub fn missing_host_header_test() { - expect_status("GET / HTTP/1.1\r\nContent-Length: 5\r\n\r\n", [#(400, 499)]) -} - -pub fn multiple_host_headers_test() { - expect_status( - "GET / HTTP/1.1\r\nHost: example.com\r\nHost: example.org\r\n\r\n", - [#(400, 499)], - ) -} - -pub fn overflowing_negative_content_length_test() { - expect_status( - "GET / HTTP/1.1\r\nHost: example.com\r\nContent-Length: -123456789123456789123456789\r\n\r\n", - [#(400, 499)], - ) -} - -pub fn negative_content_length_test() { - expect_status( - "GET / HTTP/1.1\r\nHost: example.com\r\nContent-Length: -1234\r\n\r\n", - [#(400, 499)], - ) -} - -pub fn non_numeric_content_length_test() { - expect_status( - "GET / HTTP/1.1\r\nHost: example.com\r\nContent-Length: abc\r\n\r\n", - [#(400, 499)], - ) -} - -pub fn empty_header_value_test() { - expect_status( - "GET / HTTP/1.1\r\nHost: example.com\r\nX-Empty-Header: \r\n\r\n", - [#(200, 299)], - ) -} - -pub fn header_containing_invalid_control_character_test() { - expect_status( - "GET / HTTP/1.1\r\nHost: example.com\r\nX-Bad-Control-Char: test\u{0007}\r\n\r\n", - [#(400, 499)], - ) -} - -pub fn invalid_http_version_test() { - expect_status("GET / HTTP/9.9\r\nHost: example.com\r\n\r\n", [ - #(400, 499), - #(500, 599), - ]) -} - -pub fn invalid_prefix_of_request_test() { - expect_status("Extra lineGET / HTTP/1.1\r\nHost: example.com\r\n\r\n", [ - #(400, 499), - #(500, 599), - ]) -} - -pub fn invalid_line_ending_test() { - expect_status( - "GET / HTTP/1.1\r\nHost: example.com\r\n\rSome-Header: Test\r\n\r\n", - [#(400, 499)], - ) -} - -pub fn valid_post_request_with_body_test() { - let resp = - expect_status( - "POST / HTTP/1.1\r\nHost: example.com\r\nContent-Length: 5\r\n\r\nhello", - [#(200, 299)], - ) - - assert resp.body == "hello" -} - -pub fn chunked_transfer_encoding_test() { - let resp = - expect_status( - "POST / HTTP/1.1\r\nHost: example.com\r\nTransfer-Encoding: chunked\r\n\r\nc\r\nHellO world1\r\n0\r\n\r\n", - [#(200, 299)], - ) - - assert resp.body == "HellO world1" -} - -pub fn conflicting_transfer_encoding_and_content_length_test() { - let resp = - expect_status( - "POST / HTTP/1.1\r\nHost: example.com\r\ncontent-LengtH: 5\r\nTransFer-Encoding: chunked\r\n\r\nc\r\nHellO world1\r\n0\r\n\r\n", - [#(400, 499), #(200, 299)], - ) - - assert resp.body == "HellO world1" -} diff --git a/test/http_test.gleam b/test/http_test.gleam deleted file mode 100644 index 4e78d44..0000000 --- a/test/http_test.gleam +++ /dev/null @@ -1,150 +0,0 @@ -import client/tcp as client -import ewe -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/httpc -import gleam/int -import gleam/string -import glisten/socket -import glisten/tcp -import server - -const quote = "Lorem ipsum dolor sit amet, consectetur adipiscing elit, sed do eiusmod tempor incididunt ut labore et dolore magna aliqua. Ut enim ad minim veniam, quis nostrud exercitation ullamco laboris nisi ut aliquip ex ea commodo consequat. Duis aute irure dolor in reprehenderit in voluptate velit esse cillum dolore eu fugiat nulla pariatur. Excepteur sint occaecat cupidatat non proident, sunt in culpa qui officia deserunt mollit anim id est laborum." - -pub fn simple_request_test() { - let socket_address = server.start(server.hi()) - let ip = ewe.ip_address_to_string(socket_address.ip) - let port = int.to_string(socket_address.port) - - let assert Ok(req) = request.to("http://" <> ip <> ":" <> port <> "/") - - let assert Ok(resp) = httpc.send(req) - - assert resp.status == 200 - assert resp.body == "hi" - assert response.get_header(resp, "content-type") - == Ok("text/plain; charset=utf-8") - assert response.get_header(resp, "content-length") == Ok("2") - assert response.get_header(resp, "connection") == Ok("close") -} - -pub fn chunked_body_test() { - let socket_address = server.start(server.echoer()) - let ip = ewe.ip_address_to_string(socket_address.ip) - let port = int.to_string(socket_address.port) - - let assert Ok(req) = request.to("http://" <> ip <> ":" <> port <> "/") - let req = - request.set_header(req, "Transfer-Encoding", "chunked") - |> request.set_method(http.Post) - |> request.set_header("Host", "localhost:" <> port) - |> request.set_header("Trailer", "Server-Timing") - |> request.set_header("Content-Type", "text/plain; charset=utf-8") - |> request.set_body( - "8\r\n" - <> "Mozilla \r\n" - <> "11\r\n" - <> "Developer Network\r\n" - <> "0\r\n" - <> "Server-Timing: total;dur=1000\r\n" - <> "\r\n", - ) - - let assert Ok(resp) = httpc.send(req) - - assert resp.status == 200 - assert resp.body == "Mozilla Developer Network" - assert response.get_header(resp, "content-type") - == Ok("text/plain; charset=utf-8") - assert response.get_header(resp, "content-length") == Ok("25") -} - -pub fn chunked_body_partial_test() { - let socket_address = server.start(server.echoer()) - - let req = - "POST /echo HTTP/1.1\r\n" - <> "Host: localhost:" - <> int.to_string(socket_address.port) - <> "\r\n" - <> "Content-Type: text/plain; charset=utf-8\r\n" - <> "Transfer-Encoding: chunked\r\n\r\n" - - let chunk1 = "D\r\n" <> "Hello, world!\r\n" - let chunk2 = "1\r\n" - let chunk3 = "#\r\n" - let chunk4 = "0\r\n" - - use socket <- client.with_socket(socket_address.port, active: False) - - let assert Ok(Nil) = tcp.send(socket, bytes_tree.from_string(req)) - let assert Ok(Nil) = tcp.send(socket, bytes_tree.from_string(chunk1)) - let assert Ok(Nil) = tcp.send(socket, bytes_tree.from_string(chunk2)) - let assert Ok(Nil) = tcp.send(socket, bytes_tree.from_string(chunk3)) - let assert Ok(Nil) = tcp.send(socket, bytes_tree.from_string(chunk4)) - - let assert Ok(resp) = tcp.receive(socket, 0) - let assert Ok(Nil) = tcp.close(socket) - - let assert Ok(resp) = bit_array.to_string(resp) - let assert [_, body] = string.split(resp, "\r\n\r\n") - assert body == "Hello, world!#" -} - -pub fn idle_timeout_test() { - let socket_address = server.start(server.echoer() |> ewe.idle_timeout(500)) - let port = int.to_string(socket_address.port) - - use socket <- client.with_socket(socket_address.port, active: False) - - let req = - "GET /echo HTTP/1.1\r\n" - <> "Host: localhost:" - <> port - <> "\r\n" - <> "Connection: keep-alive\r\n" - <> "\r\n" - - let assert Ok(Nil) = tcp.send(socket, bytes_tree.from_string(req)) - let assert Ok(_) = tcp.receive(socket, 0) - - process.sleep(250) - - let assert Ok(Nil) = tcp.send(socket, bytes_tree.from_string(req)) - let assert Ok(_) = tcp.receive(socket, 0) - - let assert Error(socket.Timeout) = tcp.receive_timeout(socket, 0, 300) - let assert Error(socket.Closed) = tcp.receive_timeout(socket, 0, 250) - - Nil -} - -pub fn gzip_test() { - let socket_address = server.start(server.echoer()) - let ip = ewe.ip_address_to_string(socket_address.ip) - let port = int.to_string(socket_address.port) - - let assert Ok(req) = request.to("http://" <> ip <> ":" <> port <> "/") - - let quote = string.repeat(quote, 10) - let quote_length = int.to_string(string.byte_size(quote)) - - let req = - request.set_method(req, http.Post) - |> request.set_header("Host", "localhost:" <> port) - |> request.set_header("Accept-Encoding", "gzip, deflate, br") - |> request.set_header("Content-Type", "text/plain; charset=utf-8") - |> request.set_header("Content-Length", quote_length) - |> request.set_body(<>) - - let assert Ok(resp) = httpc.send_bits(req) - - assert gunzip(resp.body) == <> -} - -@external(erlang, "zlib", "gunzip") -fn gunzip(data: BitArray) -> BitArray diff --git a/test/server.gleam b/test/server.gleam deleted file mode 100644 index 180463d..0000000 --- a/test/server.gleam +++ /dev/null @@ -1,43 +0,0 @@ -import gleam/erlang/process -import gleam/http/request -import gleam/http/response -import gleam/result - -import ewe - -pub fn start(builder: ewe.Builder) -> ewe.SocketAddress { - let name = process.new_name("ewe_test_server") - - let _ = - ewe.with_name(builder, name) - |> ewe.listening_random() - |> ewe.quiet() - |> ewe.start() - - ewe.get_server_info(name) -} - -pub fn hi() -> ewe.Builder { - ewe.new(fn(_req) { - response.new(200) - |> response.set_body(ewe.TextData("hi")) - |> response.set_header("content-type", "text/plain; charset=utf-8") - }) -} - -pub fn echoer() -> ewe.Builder { - ewe.new(fn(req) { - let content_type = - request.get_header(req, "content-type") - |> result.unwrap("text/plain") - - case ewe.read_body(req, 10_240) { - Ok(req) -> { - response.new(200) - |> response.set_body(ewe.BitsData(req.body)) - |> response.set_header("content-type", content_type) - } - Error(_) -> response.new(400) |> response.set_body(ewe.Empty) - } - }) -} -- 2.51.2