diff --git a/server/gleam.toml b/server/gleam.toml index a98d545..25cffde 100644 --- a/server/gleam.toml +++ b/server/gleam.toml @@ -18,8 +18,8 @@ gtfs = { path = "../gtfs" } gleam_stdlib = ">= 1.0.3 and < 2.0.0" wisp = ">= 2.0.0 and < 3.0.0" +ewe = ">= 7.0.0 and < 8.0.0" gleam_erlang = ">= 1.3.0 and < 2.0.0" -mist = ">= 6.0.2 and < 7.0.0" gtfs_rt_nyct = "0.2.0" # gtfs_rt_nyct = { path = "../gtfs-rt-nyct-gleam/" } gleam_httpc = ">= 5.0.0 and < 6.0.0" @@ -30,12 +30,14 @@ lustre = ">= 5.3.5 and < 6.0.0" gsv = ">= 5.0.0 and < 6.0.0" simplifile = ">= 2.3.0 and < 3.0.0" repeatedly = ">= 2.1.2 and < 3.0.0" -gleam_otp = ">= 1.2.0 and < 2.0.0" gleam_json = ">= 3.1.0 and < 4.0.0" tzif = ">= 1.1.2 and < 2.0.0" booklet = ">= 1.1.0 and < 2.0.0" envoy = ">= 1.1.0 and < 2.0.0" logging = ">= 1.3.0 and < 2.0.0" +# This is only for the wisp_ewe adapter. +# See +exception = ">= 2.1.1 and < 3.0.0" [dev-dependencies] gleeunit = ">= 1.0.0 and < 2.0.0" diff --git a/server/manifest.toml b/server/manifest.toml index 1972dbc..b06b2de 100644 --- a/server/manifest.toml +++ b/server/manifest.toml @@ -7,9 +7,11 @@ # You should check this file into your source control repository. packages = [ + { name = "alpacki", version = "3.0.1", build_tools = ["gleam"], requirements = ["gleam_stdlib"], otp_app = "alpacki", source = "hex", outer_checksum = "5BCB617A9606E56018790D4F24AB1E06B6D05BC17ACD6686B99E8D4D18B1F5FD" }, { name = "booklet", version = "1.1.0", build_tools = ["gleam"], requirements = [], otp_app = "booklet", source = "hex", outer_checksum = "08E0FDB78DC4D8A5D3C80295B021505C7D2A2E7B6C6D5EAB7286C36F4A53C851" }, { name = "directories", version = "1.2.0", build_tools = ["gleam"], requirements = ["envoy", "gleam_stdlib", "platform", "simplifile"], otp_app = "directories", source = "hex", outer_checksum = "D13090CFCDF6759B87217E8DDD73A75903A700148A82C1D33799F333E249BF9E" }, { name = "envoy", version = "1.2.0", build_tools = ["gleam"], requirements = ["gleam_stdlib"], otp_app = "envoy", source = "hex", outer_checksum = "9C6FBB6BFA02A52798BEEC5977A738CAD6E4A057F4B67FD0C8061AD2502C191A" }, + { name = "ewe", version = "7.0.0", build_tools = ["gleam"], requirements = ["alpacki", "gleam_erlang", "gleam_http", "gleam_otp", "gleam_stdlib", "logging", "websocks"], otp_app = "ewe", source = "hex", outer_checksum = "13ACE8FB49EBA631ECA8BA2B62D44B159BDA2F5D69144E509FE2273C1812AB58" }, { name = "exception", version = "2.1.1", build_tools = ["gleam"], requirements = ["gleam_stdlib"], otp_app = "exception", source = "hex", outer_checksum = "6BDEA95248093599391C3B5DF1835C5C6A86C353C2F99CE539B450E3432FE117" }, { name = "filepath", version = "1.1.2", build_tools = ["gleam"], requirements = ["gleam_stdlib"], otp_app = "filepath", source = "hex", outer_checksum = "B06A9AF0BF10E51401D64B98E4B627F1D2E48C154967DA7AF4D0914780A6D40A" }, { name = "gleam_community_maths", version = "2.0.2", build_tools = ["gleam"], requirements = ["gleam_stdlib", "gleam_yielder"], otp_app = "gleam_community_maths", source = "hex", outer_checksum = "8B0DA56738D4387666A9FA63554B159E2EDF269CAD4B847AB055364E288A71EC" }, @@ -44,17 +46,19 @@ packages = [ { name = "simplifile", version = "2.5.0", build_tools = ["gleam"], requirements = ["filepath", "gleam_stdlib"], otp_app = "simplifile", source = "hex", outer_checksum = "6C72DCCDF25C38A5931740B30E823969F33106831FD1637719B5EDBCA30027A4" }, { name = "splitter", version = "1.2.0", build_tools = ["gleam"], requirements = ["gleam_stdlib"], otp_app = "splitter", source = "hex", outer_checksum = "3DFD6B6C49E61EDAF6F7B27A42054A17CFF6CA2135FF553D0CB61C234D281DD0" }, { name = "tzif", version = "1.1.2", build_tools = ["gleam"], requirements = ["filepath", "gleam_stdlib", "gleam_time", "simplifile"], otp_app = "tzif", source = "hex", outer_checksum = "04C8EBAFB7F6A6E61B8B0C16E06F96E7C7F76C31172EF07204CE9813B2E54AD2" }, + { name = "websocks", version = "4.0.1", build_tools = ["gleam"], requirements = ["gleam_crypto", "gleam_erlang", "gleam_stdlib"], otp_app = "websocks", source = "hex", outer_checksum = "89B0C31A032CBE28D4C5FB5CC25A8A690669FEBA501C63F99A4E8F5748750B2A" }, { name = "wisp", version = "2.2.2", build_tools = ["gleam"], requirements = ["directories", "exception", "filepath", "gleam_crypto", "gleam_erlang", "gleam_http", "gleam_json", "gleam_stdlib", "houdini", "logging", "marceau", "mist", "simplifile"], otp_app = "wisp", source = "hex", outer_checksum = "5FF5F1E288C3437252ABB93D8F9CF42FF652CE7AD54480CFE736038DC09C4F22" }, ] [requirements] booklet = { version = ">= 1.1.0 and < 2.0.0" } envoy = { version = ">= 1.1.0 and < 2.0.0" } +ewe = { version = ">= 7.0.0 and < 8.0.0" } +exception = { version = ">= 2.1.1 and < 3.0.0" } gleam_erlang = { version = ">= 1.3.0 and < 2.0.0" } gleam_http = { version = ">= 4.3.0 and < 5.0.0" } gleam_httpc = { version = ">= 5.0.0 and < 6.0.0" } gleam_json = { version = ">= 3.1.0 and < 4.0.0" } -gleam_otp = { version = ">= 1.2.0 and < 2.0.0" } gleam_stdlib = { version = ">= 1.0.3 and < 2.0.0" } gleam_time = { version = ">= 1.4.0 and < 2.0.0" } gleeunit = { version = ">= 1.0.0 and < 2.0.0" } @@ -63,7 +67,6 @@ gtfs = { path = "../gtfs" } gtfs_rt_nyct = { version = "0.2.0" } logging = { version = ">= 1.3.0 and < 2.0.0" } lustre = { version = ">= 5.3.5 and < 6.0.0" } -mist = { version = ">= 6.0.2 and < 7.0.0" } protobin = { version = ">= 2.0.0 and < 3.0.0" } repeatedly = { version = ">= 2.1.2 and < 3.0.0" } shared = { path = "../shared" } diff --git a/server/src/subway_gleam/server.gleam b/server/src/subway_gleam/server.gleam index ddf9460..4c3e399 100644 --- a/server/src/subway_gleam/server.gleam +++ b/server/src/subway_gleam/server.gleam @@ -1,4 +1,5 @@ import eflame +import ewe import gleam/erlang/process import gleam/float import gleam/http/request @@ -9,11 +10,10 @@ import gleam/option import gleam/result import gleam/string import logging -import mist import repeatedly import tzif/database as tzif import wisp -import wisp/wisp_mist +import wisp/wisp_ewe import subway_gleam/gtfs/env as gtfs_env import subway_gleam/gtfs/st @@ -87,7 +87,7 @@ pub fn start(sleeping_after sleep_after_ms: Result(Int, Nil)) -> Nil { }) let secret_key_base = wisp.random_string(64) - let wisp_handler = wisp_mist.handler(handler(state, _), secret_key_base) + let wisp_handler = wisp_ewe.handler(handler(state, _), secret_key_base) let host = env.host() let http_port = env.http_port() @@ -101,20 +101,25 @@ pub fn start(sleeping_after sleep_after_ms: Result(Int, Nil)) -> Nil { ]), ) + let listener_name = process.new_name("sbwy_listener") + let connection_factory_name = process.new_name("sbwy_connection_factory") + let handler = ewe_handler(_, state, wisp_handler) + let server = + ewe.new(listener_name:, connection_factory_name:, handler:) + |> ewe.bind(to: host) + |> ewe.with_http2( + ewe.Http2Options(..ewe.default_http2_options(), websocket: True), + ) let assert Ok(_service) = case env.certfile(), env.keyfile() { Ok(certfile), Ok(keyfile) -> - mist_handler(_, state, wisp_handler) - |> mist.new - |> mist.bind(host) - |> mist.port(https_port) - |> mist.with_tls(certfile:, keyfile:) - |> mist.start + server + |> ewe.listening(on: https_port) + |> ewe.with_tls(ewe.Disk(cert: certfile, key: keyfile)) + |> ewe.start _, _ -> - mist_handler(_, state, wisp_handler) - |> mist.new - |> mist.bind(host) - |> mist.port(http_port) - |> mist.start + server + |> ewe.listening(on: http_port) + |> ewe.start } case sleep_after_ms { @@ -136,12 +141,12 @@ pub fn start(sleeping_after sleep_after_ms: Result(Int, Nil)) -> Nil { // This is done b/c wisp doesn't support some features (e.g. websockets, // server-sent events). So for routes that use the features wisp does support, // they go in the wisp handler. Otherwise, they go here. -fn mist_handler( - req: request.Request(mist.Connection), +fn ewe_handler( + req: request.Request(ewe.Connection), state_ref: state.Ref, - wisp_handler: fn(request.Request(mist.Connection)) -> - response.Response(mist.ResponseData), -) -> response.Response(mist.ResponseData) { + wisp_handler: fn(request.Request(ewe.Connection)) -> + response.Response(ewe.Body), +) -> response.Response(ewe.Body) { use <- log.time("mist_handler") let state = state.get(state_ref) case request.path_segments(req) { diff --git a/server/src/subway_gleam/server/sse_gtfs.gleam b/server/src/subway_gleam/server/sse_gtfs.gleam index 87d90f2..7864e42 100644 --- a/server/src/subway_gleam/server/sse_gtfs.gleam +++ b/server/src/subway_gleam/server/sse_gtfs.gleam @@ -1,62 +1,64 @@ +import ewe import gleam/erlang/process import gleam/http/request import gleam/http/response import gleam/json -import gleam/otp/actor import gleam/result -import mist import subway_gleam/server/log import subway_gleam/server/state import subway_gleam/server/state/gtfs_store pub fn sse_gtfs( - req: request.Request(mist.Connection), + req: request.Request(ewe.Connection), state: state.Ref, model: fn() -> Result(json.Json, anyerror), -) -> response.Response(mist.ResponseData) { - use req, context <- log.request(req) - - mist.server_sent_events( - request: req, - initial_response: response.new(200) - // Prevent nginx from buffering SSE responses - |> response.set_header("X-Accel-Buffering", "no"), - init: fn(self) { +) -> response.Response(ewe.Body) { + use _req, context <- log.request(req) + response.new(200) + // Prevent nginx from buffering SSE responses + |> response.set_header("X-Accel-Buffering", "no") + |> ewe.sse( + on_init: fn(_conn, selector) { + let self = process.new_subject() let state = state.get(state) gtfs_store.subscribe_watcher(self, to: state.gtfs_store) - log.debug("Subscribed to gtfs store.", with: context) - self + #(self, process.select(selector, for: self)) }, - loop: fn(self: process.Subject(Nil), _msg: Nil, conn: mist.SSEConnection) -> actor.Next( - process.Subject(Nil), - Nil, - ) { + handler: fn(conn, self, _message: Nil) { log.debug("Notified of gtfs update; updating model...", with: context) let model = model() |> result.replace_error(Nil) - |> result.map(json.to_string_tree) + |> result.map(json.to_string) log.debug("Finished updating model.", with: context) - let event = model |> result.map(mist.event) - case result.try(event, mist.send_event(conn, _)) { - Ok(Nil) -> { + let event = model |> result.map(ewe.event) + case result.map(event, ewe.send_event(conn, _)) { + Ok(Ok(Nil)) -> { log.debug("Sent model.", with: context) - actor.continue(self) + ewe.continue(self) } - Error(Nil) -> { + _ -> { log.debug( "Failed to send model: connection closed; unsubscribing from gtfs store.", with: context, ) let state = state.get(state) gtfs_store.unsubscribe_watcher(self, from: state.gtfs_store) - actor.stop() + ewe.stop() } } }, + on_close: fn(_conn, self) { + log.debug( + "Connection closed; unsubscribing from gtfs store.", + with: context, + ) + let state = state.get(state) + gtfs_store.unsubscribe_watcher(self, from: state.gtfs_store) + }, ) } diff --git a/server/src/wisp/wisp_ewe.gleam b/server/src/wisp/wisp_ewe.gleam new file mode 100644 index 0000000..398cf23 --- /dev/null +++ b/server/src/wisp/wisp_ewe.gleam @@ -0,0 +1,99 @@ +//// See . + +import ewe +import exception +import gleam/http/request +import gleam/http/response +import gleam/option +import gleam/string +import wisp +import wisp/internal + +const max_body_size = 8_000_000 + +/// Convert a Wisp request handler into a function that can be run with the Ewe +/// web server. +/// +/// # Examples +/// +/// ```gleam +/// pub fn main() { +/// let secret_key_base = "..." +/// let listener_name = process.new_name("ewe_listener") +/// let connection_factory_name = process.new_name("ewe_connection_factory") +/// +/// let assert Ok(_) = +/// handle_request +/// |> wisp_ewe.handler(secret_key_base) +/// |> ewe.new(listener_name:, connection_factory_name:, handler: _) +/// |> ewe.listening(on: 8000) +/// |> ewe.start +/// +/// process.sleep_forever() +/// } +/// ``` +/// +/// The secret key base is used for signing and encryption. To be able to +/// verify and decrypt messages you will need to use the same key each time +/// your program is run. Keep this value secret! Malicious people with this +/// value will likely be able to hack your application. +/// +pub fn handler( + handler: fn(wisp.Request) -> wisp.Response, + secret_key_base: String, +) -> fn(request.Request(ewe.Connection)) -> response.Response(ewe.Body) { + fn(req: request.Request(ewe.Connection)) { + let connection = req.body + let wisp_req = + internal.make_connection(ewe_body_reader(req), secret_key_base) + |> request.set_body(req, _) + + use <- exception.defer(fn() { + let assert Ok(_) = wisp.delete_temporary_files(wisp_req) + }) + + handler(wisp_req) + |> ewe_response(connection) + } +} + +fn ewe_body_reader(req: request.Request(ewe.Connection)) -> internal.Reader { + fn(size) { + case ewe.read_body_chunk(req, max_chunk_bytes: size, limit: max_body_size) { + Ok(ewe.Chunk(data:, request:)) -> + Ok(internal.Chunk(data, ewe_body_reader(request))) + Ok(ewe.Done(..)) -> Ok(internal.ReadingFinished) + Error(_) -> Error(Nil) + } + } +} + +fn ewe_response( + resp: response.Response(wisp.Body), + connection: ewe.Connection, +) -> response.Response(ewe.Body) { + case resp.body { + wisp.Text(text) -> response.set_body(resp, ewe.Text(text)) + wisp.Bytes(bytes) -> response.set_body(resp, ewe.Bytes(bytes)) + wisp.File(path:, offset:, limit:) -> + ewe_send_file(resp, connection, path, offset, limit) + } +} + +fn ewe_send_file( + resp: response.Response(wisp.Body), + connection: ewe.Connection, + path: String, + offset: Int, + limit: option.Option(Int), +) -> response.Response(ewe.Body) { + case ewe.file(connection, path, offset: option.Some(offset), limit:) { + Ok(file) -> response.set_body(resp, file) + Error(error) -> { + string.inspect(error) + |> wisp.log_error + + response.new(500) |> response.set_body(ewe.Empty) + } + } +}