diff --git a/examples/your_first_bot/src/your_first_bot.gleam b/examples/your_first_bot/src/your_first_bot.gleam index 5896546..562fff8 100644 --- a/examples/your_first_bot/src/your_first_bot.gleam +++ b/examples/your_first_bot/src/your_first_bot.gleam @@ -1,10 +1,9 @@ import dotenv_gleam import envoy -import gleam/erlang/process.{type Subject} -import gleam/option.{None, Some} +import gleam/erlang/process +import gleam/option.{Some} import gleam/string import grom -import grom/activity import grom/command import grom/gateway import grom/gateway/intent @@ -12,7 +11,7 @@ import grom/interaction.{type Interaction} import logging type State { - State(client: grom.Client, gateway: Subject(gateway.Message)) + State(client: grom.Client) } pub fn main() -> Nil { @@ -30,14 +29,7 @@ pub fn main() -> Nil { let assert Ok(data) = gateway.get_data(client) let gateway_start_result = - gateway.new_with_initializer( - fn(subject) { - let state = State(client, subject) - Ok(state) - }, - identify, - data, - ) + gateway.new(State(client), identify, data) |> gateway.on_event(do: on_event) |> gateway.start @@ -61,14 +53,14 @@ fn on_event(state: State, event: gateway.Event) { logging.log(logging.Warning, string.inspect(error)) gateway.continue(state) } - gateway.ReadyEvent(ready) -> on_ready(state, ready) + gateway.AllShardsReadyEvent(ready) -> on_ready(state, ready) gateway.InteractionCreatedEvent(interaction) -> on_interaction_created(state, interaction) _ -> gateway.continue(state) } } -fn on_ready(state: State, ready: gateway.ReadyMessage) { +fn on_ready(state: State, ready: gateway.AllShardsReadyMessage) { logging.log(logging.Info, "Ready!") let global_commands = [ @@ -100,16 +92,6 @@ fn on_ready(state: State, ready: gateway.ReadyMessage) { } } - state.gateway - |> gateway.update_presence(using: gateway.UpdatePresenceMessage( - status: gateway.Online, - since: None, - activities: [ - activity.new(named: "the gateway connection", type_: activity.Watching), - ], - is_afk: False, - )) - gateway.continue(state) } @@ -140,61 +122,10 @@ fn on_slash_command_executed( ) { case command.name { "ping" -> on_ping_command(state, interaction) - "soundboards" -> on_soundboards_command(state, interaction) - "join" -> on_join_command(state, interaction) _ -> gateway.continue(state) } } -fn on_join_command( - state: State, - interaction: Interaction, -) -> gateway.Next(State) { - state.gateway - |> gateway.update_voice_state(using: gateway.UpdateVoiceStateMessage( - "1155216444691325049", - Some("1155216445211422795"), - False, - False, - )) - - let _response_result = - state.client - |> interaction.respond( - to: interaction, - using: interaction.RespondWithChannelMessageWithSource( - interaction.ResponseMessage( - ..interaction.new_response_message(), - content: Some("Joined!"), - ), - ), - ) - - gateway.continue(state) -} - -fn on_soundboards_command( - state: State, - interaction: Interaction, -) -> gateway.Next(State) { - state.gateway - |> gateway.request_soundboard_sounds(for: ["1155216444691325049"]) - - let _response_result = - state.client - |> interaction.respond( - to: interaction, - using: interaction.RespondWithChannelMessageWithSource( - interaction.ResponseMessage( - ..interaction.new_response_message(), - content: Some("See the console!"), - ), - ), - ) - - gateway.continue(state) -} - fn on_ping_command(state: State, interaction: Interaction) { let response = interaction.RespondWithChannelMessageWithSource( diff --git a/mise.toml b/mise.toml new file mode 100644 index 0000000..f6d5bc0 --- /dev/null +++ b/mise.toml @@ -0,0 +1,3 @@ +[tools] +erlang = "latest" +gleam = "latest" diff --git a/src/grom/gateway.gleam b/src/grom/gateway.gleam index 7e4dd75..34c00cc 100644 --- a/src/grom/gateway.gleam +++ b/src/grom/gateway.gleam @@ -772,12 +772,19 @@ pub opaque type Builder(state) { ) } -pub opaque type ConnectionManagerMessage { +type ConnectionManagerMessage { UpdateWebsocket(to: Subject(stratus.InternalMessage(StratusUserMessage))) SendUserMessage(message: StratusUserMessage) } -pub opaque type Connection { +type ConnectionManagerState { + ConnectionManagerState( + websocket: Option(Subject(stratus.InternalMessage(StratusUserMessage))), + queued_messages: List(StratusUserMessage), + ) +} + +type Connection { GettingReady( gateway_url: String, manager: Subject(ConnectionManagerMessage), @@ -812,7 +819,7 @@ type StratusUserMessage { UserMessage(UserMessage) } -pub opaque type UserMessage { +type UserMessage { StartPresenceUpdate(UpdatePresenceMessage) StartVoiceStateUpdate(UpdateVoiceStateMessage) StartGuildMembersRequest(RequestGuildMembersMessage) @@ -2384,11 +2391,12 @@ pub fn start( Some(count) -> count None -> builder.data.recommended_shards } + let max_concurrency = builder.data.session_start_limits.max_identify_requests_per_5_seconds let supervisor = static_supervisor.new(static_supervisor.OneForOne) - let shard_ids = list.range(from: 0, to: shard_count) + let shard_ids = list.range(from: 0, to: shard_count - 1) let shards = shard_ids |> list.map(fn(id) { Shard(id, shard_count) }) @@ -2399,7 +2407,7 @@ pub fn start( |> list.index_map(fn(shards, bucket) { shards |> list.map(fn(shard) { - BucketedShard(duration.seconds(max_concurrency * bucket), shard) + BucketedShard(duration.seconds(5 * bucket), shard) }) }) @@ -2422,10 +2430,10 @@ pub fn start( process.new_selector() |> process.select(subject) - bucketed_shards - |> list.each(fn(shards) { - shards - |> list.each(fn(shard) { + let supervisor = + bucketed_shards + |> list.flatten + |> list.fold(supervisor, fn(supervisor, shard) { supervisor |> static_supervisor.add(supervised_shard_spawner( builder, @@ -2433,7 +2441,6 @@ pub fn start( subject, )) }) - }) use _supervisor <- result.try( supervisor @@ -2647,7 +2654,7 @@ fn start_connection( ) use connection_manager <- result.try( - actor.new(None) + actor.new(ConnectionManagerState(None, [])) |> actor.on_message(on_connection_manager_message) |> actor.start |> result.map_error(grom.CouldNotStartActor), @@ -2848,26 +2855,36 @@ fn reconnect(connection_state: Connection) -> Nil { } fn on_connection_manager_message( - current: Option(Subject(stratus.InternalMessage(StratusUserMessage))), + current: ConnectionManagerState, message: ConnectionManagerMessage, -) -> actor.Next(Option(Subject(stratus.InternalMessage(StratusUserMessage))), a) { - case current { - Some(manager) -> { - case message { - SendUserMessage(user_message) -> { - user_message - |> stratus.to_user_message - |> process.send(manager, _) +) { + case message { + UpdateWebsocket(to: new) -> { + list.each(current.queued_messages, fn(msg) { + msg + |> stratus.to_user_message + |> process.send(new, _) + }) + actor.continue( + ConnectionManagerState(websocket: Some(new), queued_messages: []), + ) + } + SendUserMessage(msg) -> + case current.websocket { + Some(ws) -> { + msg + |> stratus.to_user_message + |> process.send(ws, _) actor.continue(current) } - UpdateWebsocket(to: new) -> actor.continue(Some(new)) - } - } - None -> - case message { - SendUserMessage(_) -> actor.continue(current) - UpdateWebsocket(to: new) -> actor.continue(Some(new)) + None -> + actor.continue( + ConnectionManagerState(..current, queued_messages: [ + msg, + ..current.queued_messages + ]), + ) } } }