From 0a7efb254973a575d7a42d8eaca72335f97262cf Mon Sep 17 00:00:00 2001 From: Filip Hoffmann Date: Fri, 25 Jul 2025 17:58:22 +0200 Subject: [PATCH] gateway: part 1 - waiting for 3rd party update --- gleam.toml | 4 + manifest.toml | 10 ++ src/grom.gleam | 6 +- src/grom/gateway.gleam | 246 ++++++++++++++++++++++++++++++++ src/grom/gateway/handlers.gleam | 9 ++ src/grom/gateway/sequence.gleam | 47 ++++++ 6 files changed, 321 insertions(+), 1 deletion(-) create mode 100644 src/grom/gateway.gleam create mode 100644 src/grom/gateway/handlers.gleam create mode 100644 src/grom/gateway/sequence.gleam diff --git a/gleam.toml b/gleam.toml index b2784ba..dffb112 100644 --- a/gleam.toml +++ b/gleam.toml @@ -11,6 +11,10 @@ gleam_http = ">= 4.0.0 and < 5.0.0" gleam_json = ">= 3.0.1 and < 4.0.0" gleam_time = ">= 1.1.0 and < 2.0.0" multipart_form = ">= 1.1.0 and < 2.0.0" +stratus = ">= 1.0.0 and < 2.0.0" +repeatedly = ">= 2.1.2 and < 3.0.0" +gleam_erlang = ">= 1.2.0 and < 2.0.0" +gleam_otp = ">= 1.0.0 and < 2.0.0" [dev-dependencies] gleeunit = ">= 1.0.0 and < 2.0.0" diff --git a/manifest.toml b/manifest.toml index 7a617fb..ea7a891 100644 --- a/manifest.toml +++ b/manifest.toml @@ -2,21 +2,31 @@ # You typically do not need to edit this file packages = [ + { name = "gleam_crypto", version = "1.5.1", build_tools = ["gleam"], requirements = ["gleam_stdlib"], otp_app = "gleam_crypto", source = "hex", outer_checksum = "50774BAFFF1144E7872814C566C5D653D83A3EBF23ACC3156B757A1B6819086E" }, { name = "gleam_erlang", version = "1.2.0", build_tools = ["gleam"], requirements = ["gleam_stdlib"], otp_app = "gleam_erlang", source = "hex", outer_checksum = "F91CE62A2D011FA13341F3723DB7DB118541AAA5FE7311BD2716D018F01EF9E3" }, { name = "gleam_http", version = "4.1.0", build_tools = ["gleam"], requirements = ["gleam_stdlib"], otp_app = "gleam_http", source = "hex", outer_checksum = "DB25DFC8530B64B77105405B80686541A0D96F7E2D83D807D6B2155FB9A8B1B8" }, { name = "gleam_httpc", version = "4.1.1", build_tools = ["gleam"], requirements = ["gleam_erlang", "gleam_http", "gleam_stdlib"], otp_app = "gleam_httpc", source = "hex", outer_checksum = "C670EBD46FC1472AD5F1F74F1D3938D1D0AC1C7531895ED1D4DDCB6F07279F43" }, { name = "gleam_json", version = "3.0.2", build_tools = ["gleam"], requirements = ["gleam_stdlib"], otp_app = "gleam_json", source = "hex", outer_checksum = "874FA3C3BB6E22DD2BB111966BD40B3759E9094E05257899A7C08F5DE77EC049" }, + { name = "gleam_otp", version = "1.0.0", build_tools = ["gleam"], requirements = ["gleam_erlang", "gleam_stdlib"], otp_app = "gleam_otp", source = "hex", outer_checksum = "7020E652D18F9ABAC9C877270B14160519FA0856EE80126231C505D719AD68DA" }, { name = "gleam_stdlib", version = "0.62.0", build_tools = ["gleam"], requirements = [], otp_app = "gleam_stdlib", source = "hex", outer_checksum = "DC8872BC0B8550F6E22F0F698CFE7F1E4BDA7312FDEB40D6C3F44C5B706C8310" }, { name = "gleam_time", version = "1.4.0", build_tools = ["gleam"], requirements = ["gleam_stdlib"], otp_app = "gleam_time", source = "hex", outer_checksum = "DCDDC040CE97DA3D2A925CDBBA08D8A78681139745754A83998641C8A3F6587E" }, { name = "gleeunit", version = "1.6.0", build_tools = ["gleam"], requirements = ["gleam_stdlib"], otp_app = "gleeunit", source = "hex", outer_checksum = "63022D81C12C17B7F1A60E029964E830A4CBD846BBC6740004FC1F1031AE0326" }, + { name = "gramps", version = "3.0.3", build_tools = ["gleam"], requirements = ["gleam_crypto", "gleam_erlang", "gleam_http", "gleam_stdlib"], otp_app = "gramps", source = "hex", outer_checksum = "75F0F20C867A6217CBB632A7E563568D6A6366B850815041E8E0B4F179681E53" }, + { name = "logging", version = "1.3.0", build_tools = ["gleam"], requirements = ["gleam_stdlib"], otp_app = "logging", source = "hex", outer_checksum = "1098FBF10B54B44C2C7FDF0B01C1253CAFACDACABEFB4B0D027803246753E06D" }, { name = "multipart_form", version = "1.1.0", build_tools = ["gleam"], requirements = ["gleam_http", "gleam_stdlib"], otp_app = "multipart_form", source = "hex", outer_checksum = "082C77A0C3BB1128FCD55491665E9B72BC943E849B67D02B08CFA6808AD8E47C" }, + { name = "repeatedly", version = "2.1.2", build_tools = ["gleam"], requirements = [], otp_app = "repeatedly", source = "hex", outer_checksum = "93AE1938DDE0DC0F7034F32C1BF0D4E89ACEBA82198A1FE21F604E849DA5F589" }, + { name = "stratus", version = "1.0.0", build_tools = ["gleam"], requirements = ["gleam_crypto", "gleam_erlang", "gleam_http", "gleam_otp", "gleam_stdlib", "gramps", "logging"], otp_app = "stratus", source = "hex", outer_checksum = "BC2C93B8FF2BF1D18AB3A6397E4A0C74D0FD954063EEF12027F1961ED40B2DBE" }, ] [requirements] +gleam_erlang = { version = ">= 1.2.0 and < 2.0.0" } gleam_http = { version = ">= 4.0.0 and < 5.0.0" } gleam_httpc = { version = ">= 4.1.0 and < 5.0.0" } gleam_json = { version = ">= 3.0.1 and < 4.0.0" } +gleam_otp = { version = ">= 1.0.0 and < 2.0.0" } gleam_stdlib = { version = ">= 0.44.0 and < 2.0.0" } gleam_time = { version = ">= 1.1.0 and < 2.0.0" } gleeunit = { version = ">= 1.0.0 and < 2.0.0" } multipart_form = { version = ">= 1.1.0 and < 2.0.0" } +repeatedly = { version = ">= 2.1.2 and < 3.0.0" } +stratus = { version = ">= 1.0.0 and < 2.0.0" } diff --git a/src/grom.gleam b/src/grom.gleam index 67df6af..69b0587 100644 --- a/src/grom.gleam +++ b/src/grom.gleam @@ -1,6 +1,7 @@ import gleam/http/response.{type Response} import gleam/httpc import gleam/json +import gleam/otp/actor // TYPES ----------------------------------------------------------------------- @@ -13,10 +14,13 @@ pub type Error { CouldNotDecode(json.DecodeError) StatusCodeUnsuccessful(Response(String)) ResponseNotValidUtf8(BitArray) + InvalidGatewayUrl(String) + CouldNotStartActor(actor.StartError) + CouldNotSendSocketMessage } // FUNCTIONS ------------------------------------------------------------------- -pub fn version() { +pub fn version() -> String { "v0.0.0" } diff --git a/src/grom/gateway.gleam b/src/grom/gateway.gleam new file mode 100644 index 0000000..00b2755 --- /dev/null +++ b/src/grom/gateway.gleam @@ -0,0 +1,246 @@ +import gleam/dynamic/decode +import gleam/erlang/process.{type Subject} +import gleam/float +import gleam/http +import gleam/http/request +import gleam/json +import gleam/option.{type Option, None, Some} +import gleam/result +import gleam/time/duration.{type Duration} +import grom +import grom/gateway/handlers.{type Handlers} +import grom/gateway/sequence +import grom/internal/rest +import repeatedly.{type Repeater} +import stratus + +pub type GatewayResponse { + GatewayResponse( + url: String, + shards: Int, + session_start_limit: SessionStartLimit, + ) +} + +pub type ReceiveEvent { + Hello(HelloEvent) + HeartbeatAcknowledged +} + +pub type SendEvent { + Heartbeat(HeartbeatEvent) +} + +pub type HelloEvent { + HelloEvent(heartbeat_interval: Duration) +} + +pub type HeartbeatEvent { + HeartbeatEvent(sequence: Option(Int)) +} + +pub type SessionStartLimit { + SessionStartLimit( + total: Int, + remaining: Int, + reset_after: Duration, + /// Max amount of Identify requests allowed per 5 seconds + max_concurrency: Int, + ) +} + +type State { + State( + sequence_holder: Subject(sequence.Message), + heartbeat_loop: Option(Repeater(Nil)), + ) +} + +pub fn start(client: grom.Client, handlers: Handlers(a)) { + use gateway_response <- result.try( + client + |> rest.new_request(http.Get, "/gateway/bot") + |> rest.execute, + ) + + use gateway <- result.try( + gateway_response.body + |> json.parse(using: decoder()) + |> result.map_error(grom.CouldNotDecode), + ) + + let connection_url = gateway.url <> "?v=10&encoding=json" + + use connection_request <- result.try( + request.to(connection_url) + |> result.replace_error(grom.InvalidGatewayUrl(connection_url)), + ) + + use sequence_holder <- result.try(sequence.new_holder()) + let sequence_holder = sequence_holder.data + + let connection_builder = + stratus.websocket( + request: connection_request, + init: fn() { #(State(sequence_holder:, heartbeat_loop: None), None) }, + loop: fn(state, message, connection) { + case message { + stratus.Text(event) -> { + case on_receive_event(event, state, connection, handlers) { + Ok(state) -> stratus.continue(state) + Error(error) -> { + handlers.error_handler(error) + stratus.continue(state) + } + } + } + _ -> stratus.continue(state) + } + }, + ) + |> stratus.on_close(fn(state) { + case state.heartbeat_loop { + Some(loop) -> repeatedly.stop(loop) + None -> Nil + } + }) + + stratus.initialize(connection_builder) + |> result.map_error(grom.CouldNotStartActor) +} + +@internal +pub fn decoder() -> decode.Decoder(GatewayResponse) { + use url <- decode.field("url", decode.string) + use shards <- decode.field("shards", decode.int) + use session_start_limit <- decode.field( + "session_start_limit", + session_start_limit_decoder(), + ) + + decode.success(GatewayResponse(url:, shards:, session_start_limit:)) +} + +@internal +pub fn session_start_limit_decoder() -> decode.Decoder(SessionStartLimit) { + use total <- decode.field("total", decode.int) + use remaining <- decode.field("remaining", decode.int) + use reset_after <- decode.field("reset_after", { + use duration <- decode.then(decode.int) + + duration + |> duration.milliseconds + |> decode.success + }) + use max_concurrency <- decode.field("max_concurrency", decode.int) + + decode.success(SessionStartLimit( + total:, + remaining:, + reset_after:, + max_concurrency:, + )) +} + +fn on_receive_event( + event: String, + state: State, + connection: stratus.Connection, + handlers: Handlers(a), +) -> Result(State, grom.Error) { + use event <- result.try( + event + |> json.parse(using: receive_event_decoder()) + |> result.map_error(grom.CouldNotDecode), + ) + + case event { + Hello(event) -> on_hello_event(event, state, connection, handlers) + } +} + +fn on_hello_event( + event: HelloEvent, + state: State, + connection: stratus.Connection, + handlers: Handlers(a), +) -> Result(State, grom.Error) { + let interval_milliseconds = + event.heartbeat_interval + |> duration.to_seconds + |> float.multiply(1000.0) + + interval_milliseconds + |> float.multiply(get_jitter()) + |> float.round + |> process.sleep + + use _ <- result.try( + HeartbeatEvent(None) + |> heartbeat_event_to_json + |> json.to_string + |> stratus.send_text_message(connection, _) + |> result.replace_error(grom.CouldNotSendSocketMessage), + ) + + let repeater = + repeatedly.call(float.round(interval_milliseconds), Nil, fn(_state, _i) { + let sequence = + state.sequence_holder + |> sequence.get + + case + { + HeartbeatEvent(sequence:) + |> heartbeat_event_to_json + |> json.to_string + |> stratus.send_text_message(connection, _) + |> result.replace_error(grom.CouldNotSendSocketMessage) + } + { + Ok(_) -> Nil + Error(error) -> { + handlers.error_handler(error) + Nil + } + } + }) + + Ok(State(..state, heartbeat_loop: Some(repeater))) +} + +fn get_jitter() -> Float { + case float.random() { + 0.0 -> get_jitter() + x -> x + } +} + +fn receive_event_decoder() -> decode.Decoder(ReceiveEvent) { + use opcode <- decode.field("op", decode.int) + use data <- decode.optional_field("d", HeartbeatAcknowledged, case opcode { + 10 -> hello_event_decoder() + _ -> decode.failure(Hello(HelloEvent(duration.seconds(0))), "ReceiveEvent") + }) + + decode.success(data) +} + +fn hello_event_decoder() -> decode.Decoder(ReceiveEvent) { + use heartbeat_interval <- decode.field("heartbeat_interval", { + use duration <- decode.then(decode.int) + + duration + |> duration.milliseconds + |> decode.success + }) + + decode.success(Hello(HelloEvent(heartbeat_interval:))) +} + +fn heartbeat_event_to_json(event: HeartbeatEvent) -> json.Json { + json.object([ + #("op", json.int(1)), + #("d", json.nullable(event.sequence, json.int)), + ]) +} diff --git a/src/grom/gateway/handlers.gleam b/src/grom/gateway/handlers.gleam new file mode 100644 index 0000000..e0d9123 --- /dev/null +++ b/src/grom/gateway/handlers.gleam @@ -0,0 +1,9 @@ +import grom + +pub type Handlers(a) { + Handlers(error_handler: fn(grom.Error) -> a) +} + +pub fn new() { + Handlers(error_handler: fn(_error) { Nil }) +} diff --git a/src/grom/gateway/sequence.gleam b/src/grom/gateway/sequence.gleam new file mode 100644 index 0000000..49144fe --- /dev/null +++ b/src/grom/gateway/sequence.gleam @@ -0,0 +1,47 @@ +import gleam/erlang/process.{type Subject} +import gleam/option.{type Option, None} +import gleam/otp/actor +import gleam/result +import grom + +pub type Message { + Get(reply_to: Subject(Option(Int))) + Set(Option(Int)) + Shutdown +} + +pub fn new_holder() -> Result(actor.Started(Subject(Message)), grom.Error) { + actor.new(None) + |> actor.on_message(on_message) + |> actor.start + |> result.map_error(grom.CouldNotStartActor) +} + +pub fn get(from actor: Subject(Message)) -> Option(Int) { + actor + |> actor.call(10, Get) +} + +pub fn set(actor: Subject(Message), to new: Option(Int)) -> Nil { + actor + |> actor.send(Set(new)) +} + +pub fn shutdown(actor: Subject(Message)) -> Nil { + actor + |> actor.send(Shutdown) +} + +fn on_message( + current: Option(Int), + message: Message, +) -> actor.Next(Option(Int), a) { + case message { + Set(new) -> actor.continue(new) + Get(..) -> { + process.send(message.reply_to, current) + actor.continue(current) + } + Shutdown -> actor.stop() + } +} -- 2.51.2