diff --git a/fluxer_api/src/api/infrastructure/GatewayRpcError.ts b/fluxer_api/src/api/infrastructure/GatewayRpcError.ts index 673d851db..dbaf279ae 100644 --- a/fluxer_api/src/api/infrastructure/GatewayRpcError.ts +++ b/fluxer_api/src/api/infrastructure/GatewayRpcError.ts @@ -2,6 +2,7 @@ export const GatewayRpcMethodErrorCodes = { OVERLOADED: 'overloaded', + GUILD_OVERLOADED: 'guild_overloaded', INTERNAL_ERROR: 'internal_error', TIMEOUT: 'timeout', NO_RESPONDERS: 'no_responders', diff --git a/fluxer_api/src/api/infrastructure/GatewayService.ts b/fluxer_api/src/api/infrastructure/GatewayService.ts index 705153232..fb2951425 100644 --- a/fluxer_api/src/api/infrastructure/GatewayService.ts +++ b/fluxer_api/src/api/infrastructure/GatewayService.ts @@ -296,6 +296,9 @@ export class GatewayService { if (error.code === GatewayRpcMethodErrorCodes.TIMEOUT) { return new GatewayTimeoutError(); } + if (error.code === GatewayRpcMethodErrorCodes.GUILD_OVERLOADED) { + return new ServiceUnavailableError({headers: {'Retry-After': '1'}}); + } if (error.code === GatewayRpcMethodErrorCodes.OVERLOADED) { return new ServiceUnavailableError(); } diff --git a/fluxer_api/src/api/infrastructure/tests/GatewayServiceErrorMapping.test.ts b/fluxer_api/src/api/infrastructure/tests/GatewayServiceErrorMapping.test.ts index 546936460..870138613 100644 --- a/fluxer_api/src/api/infrastructure/tests/GatewayServiceErrorMapping.test.ts +++ b/fluxer_api/src/api/infrastructure/tests/GatewayServiceErrorMapping.test.ts @@ -8,6 +8,7 @@ import type {IGatewayRpcTransport} from '@app/api/infrastructure/IGatewayRpcTran import {APIErrorCodes} from '@fluxer/constants/src/ApiErrorCodes'; import {BadGatewayError} from '@fluxer/errors/src/domains/core/BadGatewayError'; import {BadRequestError} from '@fluxer/errors/src/domains/core/BadRequestError'; +import {ServiceUnavailableError} from '@fluxer/errors/src/domains/core/ServiceUnavailableError'; import {UnknownGuildError} from '@fluxer/errors/src/domains/guild/UnknownGuildError'; import {afterEach, describe, expect, it} from 'vitest'; @@ -70,4 +71,30 @@ describe('GatewayService gateway error mapping', () => { expect((error as BadRequestError).status).toBe(400); expect((error as BadRequestError).code).toBe(APIErrorCodes.INVALID_FORM_BODY); }); + + it('returns 503 with Retry-After for an overloaded guild', async () => { + const service = serviceRaising(GatewayRpcMethodErrorCodes.GUILD_OVERLOADED); + const error = await service + .getUserPermissions({guildId: createGuildID(1n), userId: createUserID(2n)}) + .catch((raised: unknown) => raised); + expect(error).toBeInstanceOf(ServiceUnavailableError); + const response = (error as ServiceUnavailableError).getResponse(); + expect(response.status).toBe(503); + expect(response.headers.get('Retry-After')).toBe('1'); + }); + + it('does not retry a guild overload response', async () => { + let calls = 0; + const client = GatewayRpcClient.createForTests({ + async call(): Promise { + calls += 1; + throw new GatewayRpcMethodError(GatewayRpcMethodErrorCodes.GUILD_OVERLOADED); + }, + async destroy(): Promise {}, + }); + await expect(client.call('guild.dispatch', {guild_id: '1'})).rejects.toMatchObject({ + code: GatewayRpcMethodErrorCodes.GUILD_OVERLOADED, + }); + expect(calls).toBe(1); + }); }); diff --git a/fluxer_app/src/features/gateway/events/EventRouter.ts b/fluxer_app/src/features/gateway/events/EventRouter.ts index 115defcfa..2d83041dd 100644 --- a/fluxer_app/src/features/gateway/events/EventRouter.ts +++ b/fluxer_app/src/features/gateway/events/EventRouter.ts @@ -20,6 +20,7 @@ import {handleGuildCountsUpdate} from '@app/features/guild/events/GuildCountsUpd import {handleGuildCreate} from '@app/features/guild/events/GuildCreate'; import {handleGuildDelete} from '@app/features/guild/events/GuildDelete'; import {handleGuildEmojisUpdate} from '@app/features/guild/events/GuildEmojisUpdate'; +import {handleGuildHealthUpdate} from '@app/features/guild/events/GuildHealthUpdate'; import {handleGuildMemberAdd} from '@app/features/guild/events/GuildMemberAdd'; import {handleGuildMemberListUpdate} from '@app/features/guild/events/GuildMemberListUpdate'; import {handleGuildMemberRemove} from '@app/features/guild/events/GuildMemberRemove'; @@ -112,6 +113,7 @@ export function createHandlerRegistry(): GatewayHandlerRegistry { registry.set('GUILD_MEMBERS_CHUNK', handleGuildMembersChunk as GatewayEventHandler); registry.set('GUILD_MEMBER_LIST_UPDATE', handleGuildMemberListUpdate as GatewayEventHandler); registry.set('GUILD_COUNTS_UPDATE', handleGuildCountsUpdate as GatewayEventHandler); + registry.set('GUILD_HEALTH_UPDATE', handleGuildHealthUpdate as GatewayEventHandler); registry.set('CHANNEL_MEMBER_COUNTS_UPDATE', handleChannelMemberCountsUpdate as GatewayEventHandler); registry.set('GUILD_ROLE_CREATE', handleGuildRoleCreate as GatewayEventHandler); registry.set('GUILD_ROLE_UPDATE', handleGuildRoleUpdate as GatewayEventHandler); diff --git a/fluxer_app/src/features/gateway/types/GatewayGuildTypes.ts b/fluxer_app/src/features/gateway/types/GatewayGuildTypes.ts index fcfd430f8..7e6acf7c4 100644 --- a/fluxer_app/src/features/gateway/types/GatewayGuildTypes.ts +++ b/fluxer_app/src/features/gateway/types/GatewayGuildTypes.ts @@ -23,4 +23,5 @@ export type GuildReadyData = Readonly<{ joined_at: string; unavailable?: boolean; unavailable_hidden?: boolean; + degraded?: boolean; }>; diff --git a/fluxer_app/src/features/guild/events/GuildCreate.ts b/fluxer_app/src/features/guild/events/GuildCreate.ts index 1ddb83662..a12be1b76 100644 --- a/fluxer_app/src/features/guild/events/GuildCreate.ts +++ b/fluxer_app/src/features/guild/events/GuildCreate.ts @@ -48,6 +48,9 @@ export function handleGuildCreate(data: GuildReadyData, _context: GatewayHandler return; } GuildAvailability.setGuildAvailable(data.id); + if (!isSync || data.degraded !== undefined) { + GuildAvailability.setGuildDegraded(data.id, data.degraded === true); + } Guilds.handleGuildCreate(data); GuildCount.handleGuildCreate(data); if (!isSync) { diff --git a/fluxer_app/src/features/guild/events/GuildDelete.ts b/fluxer_app/src/features/guild/events/GuildDelete.ts index 8b9e0254a..abee6351a 100644 --- a/fluxer_app/src/features/guild/events/GuildDelete.ts +++ b/fluxer_app/src/features/guild/events/GuildDelete.ts @@ -32,6 +32,7 @@ interface GuildDeletePayload { } export function handleGuildDelete(data: GuildDeletePayload, _context: GatewayHandlerContext): void { + GuildAvailability.setGuildDegraded(data.id, false); GuildAvailability.handleGuildAvailability(data.id, data.unavailable, data.unavailable_hidden); Guilds.handleGuildDelete({guildId: data.id, unavailable: data.unavailable}); GuildList.handleGuildDelete(data.id, data.unavailable); diff --git a/fluxer_app/src/features/guild/events/GuildHealthUpdate.ts b/fluxer_app/src/features/guild/events/GuildHealthUpdate.ts new file mode 100644 index 000000000..a8be9f40e --- /dev/null +++ b/fluxer_app/src/features/guild/events/GuildHealthUpdate.ts @@ -0,0 +1,10 @@ +// SPDX-License-Identifier: AGPL-3.0-or-later + +import GuildAvailability from '@app/features/guild/state/GuildAvailability'; +import Guilds from '@app/features/guild/state/Guilds'; + +export function handleGuildHealthUpdate(data: {guild_id: string; degraded: boolean}): void { + if (Guilds.getGuild(data.guild_id)) { + GuildAvailability.setGuildDegraded(data.guild_id, data.degraded); + } +} diff --git a/fluxer_app/src/features/guild/state/GuildAvailability.ts b/fluxer_app/src/features/guild/state/GuildAvailability.ts index 969fd4eb6..f92c7ab86 100644 --- a/fluxer_app/src/features/guild/state/GuildAvailability.ts +++ b/fluxer_app/src/features/guild/state/GuildAvailability.ts @@ -5,12 +5,14 @@ import {makeAutoObservable, observable} from 'mobx'; class GuildAvailability { unavailableGuilds: Set = observable.set(); + degradedGuilds: Set = observable.set(); constructor() { makeAutoObservable( this, { unavailableGuilds: false, + degradedGuilds: false, }, {autoBind: true}, ); @@ -23,11 +25,20 @@ class GuildAvailability { } setGuildUnavailable(guildId: string): void { + this.degradedGuilds.delete(guildId); if (!this.unavailableGuilds.has(guildId)) { this.unavailableGuilds.add(guildId); } } + setGuildDegraded(guildId: string, degraded: boolean): void { + if (degraded) { + this.degradedGuilds.add(guildId); + } else { + this.degradedGuilds.delete(guildId); + } + } + handleGuildAvailability(guildId: string, unavailable = false, unavailableHidden = false): void { if (unavailable && !unavailableHidden) { this.setGuildUnavailable(guildId); @@ -38,7 +49,11 @@ class GuildAvailability { loadUnavailableGuilds(guilds: ReadonlyArray): void { this.unavailableGuilds.clear(); + this.degradedGuilds.clear(); for (const guild of guilds) { + if (guild.degraded && !guild.unavailable) { + this.degradedGuilds.add(guild.id); + } if (guild.unavailable && !guild.unavailable_hidden) { this.unavailableGuilds.add(guild.id); } diff --git a/fluxer_gateway/src/gateway/fluxer_gateway_sup.erl b/fluxer_gateway/src/gateway/fluxer_gateway_sup.erl index 1e55bc545..71c695813 100644 --- a/fluxer_gateway/src/gateway/fluxer_gateway_sup.erl +++ b/fluxer_gateway/src/gateway/fluxer_gateway_sup.erl @@ -90,6 +90,8 @@ role_specs(presence, _Role) -> ]; role_specs(guilds, _Role) -> [ + child_spec(gateway_clock_offset, gateway_clock_offset), + child_spec(guild_health, guild_health), child_spec(guild_counts_cache, guild_counts_cache), child_spec(guild_manager, guild_manager), child_spec(voice_state_counts_sync, voice_state_counts_sync) diff --git a/fluxer_gateway/src/gateway/gateway_clock_offset.erl b/fluxer_gateway/src/gateway/gateway_clock_offset.erl new file mode 100644 index 000000000..e1386622b --- /dev/null +++ b/fluxer_gateway/src/gateway/gateway_clock_offset.erl @@ -0,0 +1,142 @@ +%% SPDX-License-Identifier: AGPL-3.0-or-later + +-module(gateway_clock_offset). +-typing([eqwalizer]). +-behaviour(gen_server). + +-export([start_link/0, start_link/1, offset/1, sample/3]). +-export([init/1, handle_call/3, handle_cast/2, handle_info/2, terminate/2, code_change/3]). + +-define(TABLE, gateway_clock_offset). +-define(PROBE_INTERVAL_MS, 10_000). +-define(PROBE_TIMEOUT_MS, 1_000). +-define(MAX_PROBE_RTT_MS, 100). + +-type read_fun() :: fun((node()) -> integer()). +-type options() :: #{ + read_fun => read_fun(), + nodes_fun => fun(() -> [node()]), + interval_ms => pos_integer() +}. +-type state() :: #{ + read_fun := read_fun(), + nodes_fun := fun(() -> [node()]), + interval_ms := pos_integer(), + probes := #{node() => reference()} +}. + +-spec start_link() -> gen_server:start_ret(). +start_link() -> + start_link(#{}). + +-spec start_link(options()) -> gen_server:start_ret(). +start_link(Options) when is_map(Options) -> + gen_server:start_link({local, ?MODULE}, ?MODULE, Options, []). + +-spec offset(node()) -> integer() | undefined. +offset(Node) when Node =:= node() -> + 0; +offset(Node) -> + try ets:lookup(?TABLE, Node) of + [{Node, Offset}] when is_integer(Offset) -> Offset; + _ -> undefined + catch + error:badarg -> undefined + end. + +-spec sample(integer(), integer(), integer()) -> {ok, integer()} | discard. +sample(RemoteMs, LocalAfterMs, RttMs) when RttMs >= 0, RttMs =< ?MAX_PROBE_RTT_MS -> + {ok, RemoteMs - LocalAfterMs}; +sample(_RemoteMs, _LocalAfterMs, _RttMs) -> + discard. + +-spec init(options()) -> {ok, state()}. +init(Options) -> + _ = ets:new(?TABLE, [named_table, protected, set, {read_concurrency, true}]), + _ = net_kernel:monitor_nodes(true), + self() ! probe_all, + {ok, #{ + read_fun => maps:get(read_fun, Options, fun read_remote_clock/1), + nodes_fun => maps:get(nodes_fun, Options, fun erlang:nodes/0), + interval_ms => maps:get(interval_ms, Options, ?PROBE_INTERVAL_MS), + probes => #{} + }}. + +-spec handle_call(term(), gen_server:from(), state()) -> {reply, ok, state()}. +handle_call(_Request, _From, State) -> + {reply, ok, State}. + +-spec handle_cast(term(), state()) -> {noreply, state()}. +handle_cast(_Msg, State) -> + {noreply, State}. + +-spec handle_info(term(), state()) -> {noreply, state()}. +handle_info(probe_all, #{nodes_fun := NodesFun, interval_ms := IntervalMs} = State) -> + erlang:send_after(IntervalMs, self(), probe_all), + {noreply, lists:foldl(fun spawn_probe/2, State, NodesFun())}; +handle_info({nodeup, Node}, State) when is_atom(Node) -> + {noreply, spawn_probe(Node, State)}; +handle_info({nodedown, Node}, #{probes := Probes} = State) when is_atom(Node) -> + ets:delete(?TABLE, Node), + {noreply, State#{probes := maps:remove(Node, Probes)}}; +handle_info({clock_offset, Node, Ref, Offset}, #{probes := Probes} = State) when + is_atom(Node), is_reference(Ref), is_integer(Offset) +-> + case maps:get(Node, Probes, undefined) of + Ref -> ets:insert(?TABLE, {Node, Offset}); + _ -> ok + end, + {noreply, State}; +handle_info(_Info, State) -> + {noreply, State}. + +-spec terminate(term(), state()) -> ok. +terminate(_Reason, _State) -> + ok. + +-spec code_change(term(), state(), term()) -> {ok, state()}. +code_change(_OldVsn, State, _Extra) -> + {ok, State}. + +-spec spawn_probe(node(), state()) -> state(). +spawn_probe(Node, #{read_fun := ReadFun, probes := Probes} = State) -> + Server = self(), + Ref = make_ref(), + _ = spawn(fun() -> probe(Server, Node, Ref, ReadFun) end), + State#{probes := Probes#{Node => Ref}}. + +-spec probe(pid(), node(), reference(), read_fun()) -> ok. +probe(Server, Node, Ref, ReadFun) -> + Start = erlang:monotonic_time(millisecond), + try ReadFun(Node) of + RemoteMs when is_integer(RemoteMs) -> + End = erlang:monotonic_time(millisecond), + report(Server, Node, Ref, sample(RemoteMs, End, End - Start)) + catch + _:_ -> ok + end. + +-spec report(pid(), node(), reference(), {ok, integer()} | discard) -> ok. +report(Server, Node, Ref, {ok, Offset}) -> + Server ! {clock_offset, Node, Ref, Offset}, + ok; +report(_Server, _Node, _Ref, discard) -> + ok. + +-spec read_remote_clock(node()) -> integer(). +read_remote_clock(Node) -> + erpc:call(Node, erlang, monotonic_time, [millisecond], ?PROBE_TIMEOUT_MS). + +-ifdef(TEST). +-include_lib("eunit/include/eunit.hrl"). + +sample_retains_the_full_return_trip_uncertainty_test() -> + ?assertEqual({ok, 4990}, sample(15000, 10010, 20)), + ?assertEqual({ok, -5000}, sample(5000, 10000, 0)), + ?assertEqual(discard, sample(15000, 10000, 101)), + ?assertEqual(discard, sample(15000, 10000, -1)). + +local_clock_does_not_need_a_probe_test() -> + ?assertEqual(0, offset(node())). + +-endif. diff --git a/fluxer_gateway/src/gateway/gateway_rpc_guild_infra.erl b/fluxer_gateway/src/gateway/gateway_rpc_guild_infra.erl index 4d917822e..0cc928e2d 100644 --- a/fluxer_gateway/src/gateway/gateway_rpc_guild_infra.erl +++ b/fluxer_gateway/src/gateway/gateway_rpc_guild_infra.erl @@ -6,6 +6,7 @@ -export([ with_guild/2, with_guild/3, + with_guild_unchecked/2, with_voice_server/2, ensure_guild_pid/1, get_guild_pid/1, @@ -47,7 +48,29 @@ with_guild(GuildId, Fun, NotFoundError) -> -spec run_with_guild_pid(integer(), pid(), fun((pid()) -> T), binary()) -> T when T :: term(). run_with_guild_pid(GuildId, Pid, Fun, NotFoundError) -> - run_with_guild_pid_guard(GuildId, Pid, fun() -> Fun(Pid) end, Fun, NotFoundError). + Checked = fun(CurrentPid) -> + ensure_responsive(CurrentPid), + Fun(CurrentPid) + end, + run_with_guild_pid_guard(GuildId, Pid, fun() -> Checked(Pid) end, Checked, NotFoundError). + +-spec with_guild_unchecked(integer(), fun((pid()) -> T)) -> T when T :: term(). +with_guild_unchecked(GuildId, Fun) -> + case ensure_guild_pid(GuildId) of + {ok, Pid} -> + run_with_guild_pid_guard( + GuildId, Pid, fun() -> Fun(Pid) end, Fun, <<"guild_not_found">> + ); + error -> + gateway_rpc_error:raise(<<"guild_not_found">>) + end. + +-spec ensure_responsive(pid()) -> ok. +ensure_responsive(Pid) -> + case guild_health:is_degraded(Pid) of + true -> gateway_rpc_error:raise(<<"guild_overloaded">>); + false -> ok + end. -spec with_voice_server(integer(), fun((pid(), pid()) -> T)) -> T when T :: term(). with_voice_server(GuildId, Fun) -> @@ -352,6 +375,7 @@ safe_gen_server_call(Pid, Request, Timeout) -> -spec safe_guild_call(integer(), pid(), {atom(), map()}, pos_integer()) -> {ok, term()} | error. safe_guild_call(GuildId, Pid, Request, Timeout) -> + ensure_responsive(Pid), try guild_query_handler:call(Pid, Request, Timeout) of Reply -> {ok, Reply} catch @@ -385,6 +409,7 @@ retry_guild_call(GuildId, Request, Timeout) -> -spec retry_guild_call_pid(integer(), pid(), {atom(), map()}, pos_integer()) -> {ok, term()} | error. retry_guild_call_pid(GuildId, NewPid, Request, Timeout) -> + ensure_responsive(NewPid), try guild_query_handler:call(NewPid, Request, Timeout) of Reply -> {ok, Reply} catch @@ -401,6 +426,29 @@ retry_guild_call_pid(GuildId, NewPid, Request, Timeout) -> -ifdef(TEST). +overloaded_guild_is_rejected_before_enqueuing_work_test() -> + ok = guild_ets_owner:ensure_table(guild_health_status, [named_table, public, set]), + Pid = spawn(fun() -> + receive + stop -> ok + end + end), + true = ets:insert(guild_health_status, {Pid, 42, true, undefined, undefined}), + try + ?assertError( + {gateway_rpc_error, <<"guild_overloaded">>}, + safe_guild_call(42, Pid, {get_viewer_counts, #{user_id => 200}}, 4000) + ), + ?assertError( + {gateway_rpc_error, <<"guild_overloaded">>}, + retry_guild_call_pid(42, Pid, {get_viewer_counts, #{user_id => 200}}, 4000) + ), + ?assertEqual({message_queue_len, 0}, process_info(Pid, message_queue_len)) + after + ets:delete(guild_health_status, Pid), + Pid ! stop + end. + guild_start_backoff_delay_exponential_test() -> Delay1 = guild_start_backoff_delay(1), ?assert(Delay1 >= ?GUILD_START_BASE_MS), diff --git a/fluxer_gateway/src/gateway/gateway_rpc_guild_lifecycle.erl b/fluxer_gateway/src/gateway/gateway_rpc_guild_lifecycle.erl index f9617a9a2..fd7cd4c77 100644 --- a/fluxer_gateway/src/gateway/gateway_rpc_guild_lifecycle.erl +++ b/fluxer_gateway/src/gateway/gateway_rpc_guild_lifecycle.erl @@ -25,7 +25,7 @@ handle(<<"guild.shutdown">>, P) -> handle_shutdown(P). -spec handle_dispatch(map()) -> term(). handle_dispatch(#{<<"guild_id">> := GuildIdBin, <<"event">> := Event, <<"data">> := Data}) -> GuildId = validation:snowflake_or_throw(<<"guild_id">>, GuildIdBin), - gateway_rpc_guild_infra:with_guild(GuildId, fun(Pid) -> + gateway_rpc_guild_infra:with_guild_unchecked(GuildId, fun(Pid) -> EventAtom = constants:dispatch_event_atom(Event), IsAlive = gateway_rpc_guild_infra:is_cached_guild_pid_alive(Pid), logger:debug( @@ -40,26 +40,14 @@ handle_dispatch(#{<<"guild_id">> := GuildIdBin, <<"event">> := Event, <<"data">> handle_get_data(#{<<"guild_id">> := GuildIdBin, <<"user_id">> := UserIdBin}) -> GuildId = validation:snowflake_or_throw(<<"guild_id">>, GuildIdBin), UserId = optional_user_id(UserIdBin), - gateway_rpc_guild_infra:with_guild( - GuildId, - fun(Pid) -> - get_data_from_guild(Pid, UserId) - end, - <<"guild_not_found">> - ). + get_data_from_guild(GuildId, UserId). -spec handle_get_auth_context(map()) -> term(). handle_get_auth_context(#{<<"guild_id">> := GuildIdBin, <<"user_id">> := UserIdBin} = Params) -> GuildId = validation:snowflake_or_throw(<<"guild_id">>, GuildIdBin), UserId = optional_user_id(UserIdBin), ChannelId = optional_channel_id(maps:get(<<"channel_id">>, Params, null)), - gateway_rpc_guild_infra:with_guild( - GuildId, - fun(Pid) -> - get_auth_context_from_guild(Pid, UserId, ChannelId) - end, - <<"guild_not_found">> - ). + get_auth_context_from_guild(GuildId, UserId, ChannelId). -spec optional_channel_id(term()) -> integer() | null. optional_channel_id(Value) -> @@ -69,10 +57,10 @@ optional_channel_id(Value) -> _ -> gateway_rpc_error:raise(validation_invalid_params) end. --spec get_auth_context_from_guild(pid(), integer() | null, integer() | null) -> term(). -get_auth_context_from_guild(Pid, UserId, ChannelId) -> +-spec get_auth_context_from_guild(integer(), integer() | null, integer() | null) -> term(). +get_auth_context_from_guild(GuildId, UserId, ChannelId) -> Request = {get_guild_auth_context, #{user_id => UserId, channel_id => ChannelId}}, - case guild_query_handler:call(Pid, Request, ?GUILD_CALL_TIMEOUT) of + case read_guild(GuildId, Request) of #{auth_context := null} -> gateway_rpc_error:raise(<<"forbidden">>); #{auth_context := AuthContext} -> @@ -89,10 +77,10 @@ optional_user_id(Value) -> _ -> gateway_rpc_error:raise(validation_invalid_params) end. --spec get_data_from_guild(pid(), integer() | null) -> term(). -get_data_from_guild(Pid, UserId) -> +-spec get_data_from_guild(integer(), integer() | null) -> term(). +get_data_from_guild(GuildId, UserId) -> Request = {get_guild_data, #{user_id => UserId}}, - case guild_query_handler:call(Pid, Request, ?GUILD_CALL_TIMEOUT) of + case read_guild(GuildId, Request) of #{guild_data := null, error_reason := <<"forbidden">>} -> gateway_rpc_error:raise(<<"forbidden">>); #{guild_data := null} -> @@ -103,6 +91,17 @@ get_data_from_guild(Pid, UserId) -> gateway_rpc_error:raise(<<"guild_data_error">>) end. +-spec read_guild(integer(), {atom(), map()}) -> term(). +read_guild(GuildId, Request) -> + case guild_read_model:query(GuildId, Request) of + {ok, Reply} -> + Reply; + miss -> + gateway_rpc_guild_infra:with_guild(GuildId, fun(Pid) -> + guild_query_handler:call(Pid, Request, ?GUILD_CALL_TIMEOUT) + end) + end. + -spec handle_start(map()) -> true. handle_start(#{<<"guild_id">> := GuildIdBin}) -> GuildId = validation:snowflake_or_throw(<<"guild_id">>, GuildIdBin), diff --git a/fluxer_gateway/src/gateway/gateway_rpc_guild_members.erl b/fluxer_gateway/src/gateway/gateway_rpc_guild_members.erl index e6b34532a..d1b19afb8 100644 --- a/fluxer_gateway/src/gateway/gateway_rpc_guild_members.erl +++ b/fluxer_gateway/src/gateway/gateway_rpc_guild_members.erl @@ -157,22 +157,10 @@ has_member_from_guild(GuildId, Pid, Msg) -> integer(), integer() ) -> {ok, map() | undefined} | {error, guild_not_found} | error. get_member_cached_or_rpc(GuildId, UserId) -> - case guild_permission_cache:get_member(GuildId, UserId) of - {ok, MemberData} when is_map(MemberData) -> - maybe_refresh_member(GuildId, UserId, MemberData); - {ok, MemberOrUndefined} -> - {ok, MemberOrUndefined}; - {error, not_found} -> - get_member_via_rpc(GuildId, UserId) - end. - --spec maybe_refresh_member(integer(), integer(), map()) -> - {ok, map() | undefined} | {error, guild_not_found} | error. -maybe_refresh_member(GuildId, UserId, MemberData) -> - Key = <<"communication_disabled_until">>, - case maps:is_key(Key, MemberData) of - true -> {ok, MemberData}; - false -> get_member_via_rpc(GuildId, UserId) + case guild_read_model:query(GuildId, {get_guild_member, #{user_id => UserId}}) of + {ok, #{success := true, member_data := Member}} -> {ok, Member}; + {ok, #{success := false}} -> {ok, undefined}; + miss -> get_member_via_rpc(GuildId, UserId) end. -spec get_member_via_rpc( diff --git a/fluxer_gateway/src/guild/guild.erl b/fluxer_gateway/src/guild/guild.erl index 52bba3ea8..7f3fbe7c7 100644 --- a/fluxer_gateway/src/guild/guild.erl +++ b/fluxer_gateway/src/guild/guild.erl @@ -43,30 +43,62 @@ init(GuildState) -> State3 = guild_init:init_caches_and_timers(State2), State4 = guild_init:init_voice_server(State3), erlang:garbage_collect(), - {ok, State4, ?HIBERNATE_TIMEOUT}. + State5 = guild_health:register_guild(State4), + ok = guild_read_model:put_state(State5), + {ok, State5, ?HIBERNATE_TIMEOUT}. -spec handle_call(term(), gen_server:from(), guild_state()) -> call_reply(). -handle_call({session_connect, Request}, {CallerPid, _}, State) -> +handle_call(Msg, From, State) -> + Result = handle_call_internal(Msg, From, State), + ok = publish_read_model(Result, State), + Result. + +-spec handle_cast(term(), guild_state()) -> cast_reply(). +handle_cast(Msg, State) -> + Result = handle_cast_internal(Msg, State), + ok = publish_read_model(Result, State), + Result. + +-spec handle_info(term(), guild_state()) -> info_reply(). +handle_info(Msg, State) -> + Result = handle_info_internal(Msg, State), + ok = publish_read_model(Result, State), + Result. + +-spec publish_read_model(tuple(), guild_state()) -> ok. +publish_read_model({reply, _Reply, NewState}, OldState) -> + guild_read_model:update(OldState, NewState); +publish_read_model({noreply, NewState}, OldState) -> + guild_read_model:update(OldState, NewState); +publish_read_model({noreply, NewState, _Timeout}, OldState) -> + guild_read_model:update(OldState, NewState); +publish_read_model(_Result, _OldState) -> + ok. + +-spec handle_call_internal(term(), gen_server:from(), guild_state()) -> call_reply(). +handle_call_internal({session_connect, Request}, {CallerPid, _}, State) -> handle_session_connect_call(Request, CallerPid, State); -handle_call(export_handoff_state, _From, State) -> +handle_call_internal(export_handoff_state, _From, State) -> {reply, {ok, guild_handoff:export_handoff_state(State)}, State}; -handle_call({get_guild_id}, _From, State) -> +handle_call_internal({get_guild_id}, _From, State) -> {reply, maps:get(id, State, undefined), State}; -handle_call({get_voice_guild_state}, _From, State) -> +handle_call_internal({get_voice_guild_state}, _From, State) -> {reply, voice_guild_state(State), State}; -handle_call({dispatch, Request}, _From, State) -> +handle_call_internal({dispatch, Request}, _From, State) -> handle_dispatch_call(Request, State); -handle_call({reload, NewData}, _From, State) -> +handle_call_internal({reload, NewData}, _From, State) -> handle_reload_call(NewData, State); -handle_call(get_voice_server_pid, _From, State) -> +handle_call_internal(get_voice_server_pid, _From, State) -> guild_voice_lifecycle:reply_voice_server_pid(State); -handle_call({released_push_holds, SessionIds}, _From, State) when is_list(SessionIds) -> +handle_call_internal({released_push_holds, SessionIds}, _From, State) when + is_list(SessionIds) +-> {reply, guild_sessions:released_push_holds(SessionIds, State), State}; -handle_call({terminate}, _From, State) -> +handle_call_internal({terminate}, _From, State) -> {stop, normal, ok, State}; -handle_call(Msg, From, State) when is_tuple(Msg) -> +handle_call_internal(Msg, From, State) when is_tuple(Msg) -> route_call(element(1, Msg), Msg, From, State); -handle_call(_, _From, State) -> +handle_call_internal(_, _From, State) -> {reply, ok, State}. -spec route_call(atom(), term(), gen_server:from(), guild_state()) -> call_reply(). @@ -143,38 +175,40 @@ voice_call_handler(Tag) -> subscription_call_handler(Tag). subscription_call_handler(lazy_subscribe) -> subscription; subscription_call_handler(_) -> undefined. --spec handle_cast(term(), guild_state()) -> cast_reply(). -handle_cast({dispatch, Request}, State) -> +-spec handle_cast_internal(term(), guild_state()) -> cast_reply(). +handle_cast_internal({dispatch, Request}, State) -> handle_dispatch_cast(Request, State); -handle_cast( +handle_cast_internal( {session_connect_async, #{guild_id := GuildId, attempt := Attempt, request := Request} = Msg}, State ) -> handle_session_connect_async_cast(GuildId, Attempt, Request, Msg, State); -handle_cast({session_connect_worker_done, SessionId, Attempt, Result0, Computed}, State) -> +handle_cast_internal( + {session_connect_worker_done, SessionId, Attempt, Result0, Computed}, State +) -> handle_session_connect_worker_done_cast(SessionId, Attempt, Result0, Computed, State); -handle_cast({set_session_active, SessionId}, State) -> +handle_cast_internal({set_session_active, SessionId}, State) -> handle_set_session_active_cast(SessionId, State); -handle_cast({set_session_passive, SessionId}, State) -> +handle_cast_internal({set_session_passive, SessionId}, State) -> handle_set_session_passive_cast(SessionId, State); -handle_cast({drop_session_member_lists, SessionId}, State) when is_binary(SessionId) -> +handle_cast_internal({drop_session_member_lists, SessionId}, State) when is_binary(SessionId) -> {noreply, guild_member_list:unsubscribe_session(SessionId, State)}; -handle_cast({set_session_typing_override, SessionId, TypingFlag}, State) -> +handle_cast_internal({set_session_typing_override, SessionId, TypingFlag}, State) -> handle_set_session_typing_override_cast(SessionId, TypingFlag, State); -handle_cast({set_session_push_hold, SessionId, Hold}, State) when +handle_cast_internal({set_session_push_hold, SessionId, Hold}, State) when is_binary(SessionId), is_boolean(Hold) -> {noreply, guild_sessions:set_session_push_hold(SessionId, Hold, State)}; -handle_cast({send_guild_sync, SessionId}, State) -> +handle_cast_internal({send_guild_sync, SessionId}, State) -> handle_send_guild_sync_cast(SessionId, State); -handle_cast({send_members_chunk, SessionId, ChunkData}, State) -> +handle_cast_internal({send_members_chunk, SessionId, ChunkData}, State) -> handle_send_members_chunk_cast(SessionId, ChunkData, State); -handle_cast({patch_everyone_perms, Bit}, State) when is_integer(Bit), Bit > 0 -> +handle_cast_internal({patch_everyone_perms, Bit}, State) when is_integer(Bit), Bit > 0 -> {noreply, guild_maintenance:apply_everyone_perm_bit(Bit, State)}; -handle_cast(Msg, State) when is_tuple(Msg) -> +handle_cast_internal(Msg, State) when is_tuple(Msg) -> route_cast(element(1, Msg), Msg, State); -handle_cast(_, State) -> +handle_cast_internal(_, State) -> {noreply, State}. -spec route_cast(atom(), term(), guild_state()) -> cast_reply(). @@ -197,47 +231,52 @@ cast_handler(update_member_subscriptions) -> subscription; cast_handler(update_dm_partners) -> dm_partners; cast_handler(_) -> undefined. --spec handle_info(term(), guild_state()) -> info_reply(). -handle_info({presence, UserId, Payload}, State) -> +-spec handle_info_internal(term(), guild_state()) -> info_reply(). +handle_info_internal({guild_health_probe, Ref}, State) when is_reference(Ref) -> + gen_server:cast(guild_health, {pong, self(), Ref, erlang:monotonic_time(millisecond)}), + {noreply, State}; +handle_info_internal(guild_health_register, State) -> + {noreply, guild_health:register_guild(State)}; +handle_info_internal({presence, UserId, Payload}, State) -> handle_presence_info(UserId, Payload, State); -handle_info({'EXIT', Pid, Reason}, State) -> +handle_info_internal({'EXIT', Pid, Reason}, State) -> handle_exit_info(Pid, Reason, State); -handle_info({'DOWN', Ref, process, _Pid, Reason}, State) -> +handle_info_internal({'DOWN', Ref, process, _Pid, Reason}, State) -> handle_down_info(Ref, Reason, State); -handle_info(count_cache_refresh, State) -> +handle_info_internal(count_cache_refresh, State) -> State1 = update_counts(State), _ = guild_maintenance:schedule_count_cache_refresh(State1), {noreply, State1}; -handle_info(availability_recheck, State) -> +handle_info_internal(availability_recheck, State) -> {noreply, guild_availability:handle_availability_recheck(State)}; -handle_info(passive_sync, State) -> +handle_info_internal(passive_sync, State) -> guild_passive_sync:handle_passive_sync(State); -handle_info(presence_reconcile, State) -> +handle_info_internal(presence_reconcile, State) -> guild_presence_reconcile:start_async(State), _ = guild_presence_reconcile:schedule(), {noreply, State}; -handle_info({presence_reconcile_apply, Mismatches}, State) when is_list(Mismatches) -> +handle_info_internal({presence_reconcile_apply, Mismatches}, State) when is_list(Mismatches) -> {noreply, guild_presence_reconcile:apply_mismatches(Mismatches, State)}; -handle_info({reconcile_user_presence, UserId}, State) -> +handle_info_internal({reconcile_user_presence, UserId}, State) -> {noreply, guild_presence_reconcile:reconcile_user(UserId, State)}; -handle_info({clear_stale_cached_voice_states, ConnectionIds}, State) -> +handle_info_internal({clear_stale_cached_voice_states, ConnectionIds}, State) -> handle_clear_stale_cached_voice_states_info(ConnectionIds, State); -handle_info(flush_lazy_subscribe_buffer, State) -> +handle_info_internal(flush_lazy_subscribe_buffer, State) -> guild_subscription_handler:handle_info(flush_lazy_subscribe_buffer, State); -handle_info(flush_member_list_sync_batch, State) -> +handle_info_internal(flush_member_list_sync_batch, State) -> {noreply, guild_member_list:flush_pending_member_list_syncs(State)}; -handle_info({check_auto_stop_empty, Token}, State) -> +handle_info_internal({check_auto_stop_empty, Token}, State) -> handle_auto_stop_info(Token, State); -handle_info(check_auto_stop_empty, State) -> +handle_info_internal(check_auto_stop_empty, State) -> {noreply, State}; -handle_info({timeout, TimerRef, member_list_sync_item_cache_rotate}, State) when +handle_info_internal({timeout, TimerRef, member_list_sync_item_cache_rotate}, State) when is_reference(TimerRef) -> ok = guild_member_list_subscribe:handle_sync_item_cache_timeout(TimerRef), {noreply, State}; -handle_info(timeout, State) -> +handle_info_internal(timeout, State) -> {noreply, State, hibernate}; -handle_info(_, State) -> +handle_info_internal(_, State) -> {noreply, State}. -spec handle_session_connect_call(term(), pid(), guild_state()) -> call_reply(). @@ -403,6 +442,10 @@ handle_auto_stop(Token, State) -> -spec terminate(term(), guild_state() | term()) -> ok. terminate(Reason, State) when is_map(State) -> + safe_cleanup( + fun() -> guild_read_model:delete(maps:get(id, State, undefined)) end, + "read_model_delete" + ), safe_cleanup( fun() -> PresenceSubs = presence_subscriptions(State), diff --git a/fluxer_gateway/src/guild/guild_data.erl b/fluxer_gateway/src/guild/guild_data.erl index f28fe1fa4..225fade72 100644 --- a/fluxer_gateway/src/guild/guild_data.erl +++ b/fluxer_gateway/src/guild/guild_data.erl @@ -23,7 +23,11 @@ -type guild_id() :: integer(). -define(CONNECT_SNAPSHOT_HEAVY_MEMBER_KEYS, [ - <<"members">>, members_normalized, <<"member_role_index">>, members_sorted_ids + <<"members">>, + members_normalized, + <<"member_role_index">>, + members_sorted_ids, + member_list_revision ]). -define(CONNECT_SNAPSHOT_HEAVY_SESSION_KEYS, [active_guilds, user_roles, viewable_channels]). diff --git a/fluxer_gateway/src/guild/guild_data_index.erl b/fluxer_gateway/src/guild/guild_data_index.erl index 7fb974980..0a0c7cda0 100644 --- a/fluxer_gateway/src/guild/guild_data_index.erl +++ b/fluxer_gateway/src/guild/guild_data_index.erl @@ -56,6 +56,7 @@ normalize_map(Data) -> Data0#{ <<"members">> => MemberMap, members_normalized => MemberMap, + member_list_revision => make_ref(), members_sorted_ids => lists:sort(maps:keys(MemberMap)), <<"roles">> => Roles, <<"channels">> => Channels, @@ -215,4 +216,19 @@ ensure_list_test() -> ?assertEqual([1, 2], ensure_list([1, 2])), ?assertEqual([], ensure_list(not_a_list)). +member_list_revision_rotates_on_every_member_writer_test() -> + Member = #{<<"user">> => #{<<"id">> => 1}}, + Data0 = normalize_map(#{<<"members">> => [Member]}), + Data1 = normalize_map(Data0), + Data2 = put_member(Member#{<<"nick">> => <<"renamed">>}, Data1), + Data3 = put_member_map(member_map(Data2), Data2), + Data4 = put_member_list(member_list(Data3), Data3), + Data5 = remove_member(1, Data4), + Revisions = [ + maps:get(member_list_revision, D) + || D <- [Data0, Data1, Data2, Data3, Data4, Data5] + ], + ?assert(lists:all(fun is_reference/1, Revisions)), + ?assertEqual(length(Revisions), length(lists:usort(Revisions))). + -endif. diff --git a/fluxer_gateway/src/guild/guild_data_index_members.erl b/fluxer_gateway/src/guild/guild_data_index_members.erl index e43d02739..2a8e83635 100644 --- a/fluxer_gateway/src/guild/guild_data_index_members.erl +++ b/fluxer_gateway/src/guild/guild_data_index_members.erl @@ -147,6 +147,7 @@ put_member(Member, Data) when is_map(Member), is_map(Data) -> Data1 = Data#{ <<"members">> => NewMemberMap, members_normalized => NewMemberMap, + member_list_revision => make_ref(), <<"member_role_index">> => RoleIndex2 }, invalidate_sorted_member_ids(Data1) @@ -160,6 +161,7 @@ put_member_map(MemberMap, Data) when is_map(MemberMap), is_map(Data) -> Data#{ <<"members">> => NormalizedMemberMap, members_normalized => NormalizedMemberMap, + member_list_revision => make_ref(), members_sorted_ids => lists:sort(maps:keys(NormalizedMemberMap)), <<"member_role_index">> => build_member_role_index(NormalizedMemberMap) }; @@ -184,6 +186,7 @@ remove_member(UserId, Data) when is_integer(UserId), is_map(Data) -> Data1 = Data#{ <<"members">> => RemovedMap, members_normalized => RemovedMap, + member_list_revision => make_ref(), <<"member_role_index">> => RoleIndex1 }, invalidate_sorted_member_ids(Data1); diff --git a/fluxer_gateway/src/guild/guild_data_wire.erl b/fluxer_gateway/src/guild/guild_data_wire.erl index 36ff1819b..6e9315516 100644 --- a/fluxer_gateway/src/guild/guild_data_wire.erl +++ b/fluxer_gateway/src/guild/guild_data_wire.erl @@ -158,6 +158,7 @@ named_key_kind(<<"recipient_ids">>) -> drop; named_key_kind(<<"role_index">>) -> drop; named_key_kind(<<"channel_index">>) -> drop; named_key_kind(<<"member_role_index">>) -> drop; +named_key_kind(<<"member_list_revision">>) -> drop; named_key_kind(<<"role_perms_cache">>) -> drop; named_key_kind(<<"overwrite_perms_cache">>) -> drop; named_key_kind(_) -> unknown. @@ -267,6 +268,7 @@ sample_keys() -> <<"role_index">>, <<"channel_index">>, <<"member_role_index">>, + <<"member_list_revision">>, <<"role_perms_cache">>, <<"overwrite_perms_cache">>, <<"guild_folders">>, @@ -347,6 +349,7 @@ reference_keep_payload_field(Key) -> <<"role_index">>, <<"channel_index">>, <<"member_role_index">>, + <<"member_list_revision">>, <<"role_perms_cache">>, <<"overwrite_perms_cache">> ]). @@ -520,4 +523,11 @@ reference_path_push(Key, Path) -> reference_has_any_path(Keys, Path) -> lists:any(fun(Key) -> lists:member(Key, Path) end, Keys). +member_list_revision_is_internal_at_every_depth_test() -> + Data = #{ + member_list_revision => make_ref(), + <<"nested">> => [#{<<"member_list_revision">> => make_ref(), <<"id">> => 9}] + }, + ?assertEqual(#{<<"nested">> => [#{<<"id">> => <<"9">>}]}, payload(Data)). + -endif. diff --git a/fluxer_gateway/src/guild/guild_dispatch_push.erl b/fluxer_gateway/src/guild/guild_dispatch_push.erl index 9b45c2f28..71d3add0f 100644 --- a/fluxer_gateway/src/guild/guild_dispatch_push.erl +++ b/fluxer_gateway/src/guild/guild_dispatch_push.erl @@ -590,7 +590,8 @@ send_compact_scanned_push(MessageData, GuildId, Scan, State) -> -spec compact_format_data(map(), map()) -> map(). compact_format_data(Data, FormatMembers) -> - Data#{ + WithoutRevision = maps:remove(member_list_revision, Data), + WithoutRevision#{ <<"members">> => FormatMembers, members_normalized => FormatMembers, members_sorted_ids => lists:sort(maps:keys(FormatMembers)), @@ -1904,9 +1905,14 @@ compact_push_carries_large_guild_metadata_test() -> end. compact_format_data_restricts_member_map_test() -> - Data = #{<<"guild">> => #{}, <<"members">> => #{1 => #{}, 2 => #{}, 3 => #{}}}, + Data = #{ + <<"guild">> => #{}, + <<"members">> => #{1 => #{}, 2 => #{}, 3 => #{}}, + member_list_revision => make_ref() + }, FormatMembers = #{2 => #{<<"roles">> => [<<"200">>]}}, Result = compact_format_data(Data, FormatMembers), + ?assertNot(maps:is_key(member_list_revision, Result)), ?assertEqual(FormatMembers, maps:get(<<"members">>, Result)), ?assertEqual(FormatMembers, maps:get(members_normalized, Result)), ?assertEqual([2], maps:get(members_sorted_ids, Result)), diff --git a/fluxer_gateway/src/guild/guild_ets_owner.erl b/fluxer_gateway/src/guild/guild_ets_owner.erl index 4c0913a21..b5cc225db 100644 --- a/fluxer_gateway/src/guild/guild_ets_owner.erl +++ b/fluxer_gateway/src/guild/guild_ets_owner.erl @@ -100,6 +100,8 @@ ensure_core_tables() -> {guild_unavailability_cache, RO}, {guild_circuit_breaker, RW}, {guild_permission_cache, RO}, + {guild_read_model, RO}, + {guild_health_status, RO}, {voice_update_queue, RW}, {voice_update_rate_limit, RW} ]). diff --git a/fluxer_gateway/src/guild/guild_health.erl b/fluxer_gateway/src/guild/guild_health.erl new file mode 100644 index 000000000..ec540d0ec --- /dev/null +++ b/fluxer_gateway/src/guild/guild_health.erl @@ -0,0 +1,384 @@ +%% SPDX-License-Identifier: AGPL-3.0-or-later + +-module(guild_health). +-typing([eqwalizer]). +-behaviour(gen_server). + +-export([ + start_link/0, + is_degraded/1, + register_guild/1, + put_session/2, + remove_session/2, + send_current/2 +]). +-export([init/1, handle_call/3, handle_cast/2, handle_info/2, terminate/2, code_change/3]). + +-define(TABLE, guild_health_status). +-define(INTERVAL_MS, 250). +-define(DEGRADED_MS, 2000). +-define(RECOVERED_MS, 500). +-define(REMOTE_LOOKUP_TIMEOUT_MS, 250). + +-spec start_link() -> gen_server:start_ret(). +start_link() -> + gen_server:start_link({local, ?MODULE}, ?MODULE, [], []). + +-spec is_degraded(pid()) -> boolean(). +is_degraded(Pid) when node(Pid) =:= node() -> + case lookup(Pid) of + {Pid, _GuildId, Degraded, _Targets, _Pending} -> Degraded; + undefined -> false + end; +is_degraded(Pid) -> + try erpc:call(node(Pid), ?MODULE, is_degraded, [Pid], ?REMOTE_LOOKUP_TIMEOUT_MS) of + true -> true; + _ -> false + catch + _:_ -> false + end. + +-spec register_guild(map()) -> map(). +register_guild(State) -> + case maps:get(id, State, undefined) of + GuildId when is_integer(GuildId), GuildId > 0 -> + State1 = ensure_targets(State), + Targets = maps:get(guild_health_sessions, State1), + try ets:insert_new(?TABLE, {self(), GuildId, false, Targets, undefined}) of + _ -> gen_server:cast(?MODULE, {register, self()}) + catch + error:badarg -> ok + end, + State1; + _ -> + State + end. + +-spec ensure_targets(map()) -> map(). +ensure_targets(#{guild_health_sessions := Tab} = State) -> + case table_owner(Tab) of + Pid when Pid =:= self() -> State; + _ -> ensure_targets(maps:remove(guild_health_sessions, State)) + end; +ensure_targets(State) -> + Tab = ets:new(guild_health_sessions, [set, protected, {read_concurrency, true}]), + State1 = State#{guild_health_sessions => Tab}, + maps:foreach( + fun(SessionId, _) -> put_session(SessionId, State1) end, maps:get(sessions, State, #{}) + ), + State1. + +-spec table_owner(ets:table()) -> pid() | undefined. +table_owner(Tab) -> + try ets:info(Tab, owner) of + Owner -> Owner + catch + error:badarg -> undefined + end. + +-spec put_session(term(), map()) -> ok. +put_session(SessionId, #{guild_health_sessions := Tab} = State) -> + case maps:get(SessionId, maps:get(sessions, State, #{}), undefined) of + #{pid := Pid} when is_pid(Pid) -> + ets:insert(Tab, {SessionId, Pid}), + ok; + _ -> + ok + end; +put_session(_, _) -> + ok. + +-spec remove_session(term(), map()) -> ok. +remove_session(SessionId, #{guild_health_sessions := Tab}) -> + ets:delete(Tab, SessionId), + ok; +remove_session(_, _) -> + ok. + +-spec send_current(pid(), pid()) -> ok. +send_current(GuildPid, SessionPid) -> + gen_server:cast({?MODULE, node(GuildPid)}, {current, GuildPid, SessionPid}). + +-spec init([]) -> {ok, map()}. +init([]) -> + ok = guild_ets_owner:ensure_table(?TABLE, [ + named_table, public, set, {read_concurrency, true} + ]), + State = lists:foldl( + fun({Pid, _, _, _, _}, Acc) -> track_guild(Pid, Acc) end, #{}, ets:tab2list(?TABLE) + ), + discover_guilds(), + schedule_tick(), + {ok, State}. + +-spec handle_call(term(), gen_server:from(), map()) -> {reply, ok, map()}. +handle_call(_, _, State) -> + {reply, ok, State}. + +-spec handle_cast(term(), map()) -> {noreply, map()}. +handle_cast({register, Pid}, State) when is_pid(Pid), node(Pid) =:= node() -> + {noreply, track_guild(Pid, State)}; +handle_cast({pong, Pid, Ref, HandledAt}, State) when is_pid(Pid), is_integer(HandledAt) -> + case lookup(Pid) of + {Pid, GuildId, Degraded, Targets, {Ref, SentAt}} -> + NewDegraded = next_degraded(Degraded, max(0, HandledAt - SentAt)), + update({Pid, GuildId, Degraded, Targets, undefined}, NewDegraded), + {noreply, State}; + _ -> + {noreply, State} + end; +handle_cast({current, Pid, SessionPid}, State) when is_pid(Pid), is_pid(SessionPid) -> + case lookup(Pid) of + {Pid, GuildId, Degraded, _Targets, _Pending} -> + notify_session(SessionPid, GuildId, Pid, Degraded); + _ -> + ok + end, + {noreply, State}; +handle_cast(_, State) -> + {noreply, State}. + +-spec handle_info(term(), map()) -> {noreply, map()}. +handle_info(tick, State) -> + Now = erlang:monotonic_time(millisecond), + maps:foreach(fun(Pid, _) -> check_guild(Pid, Now) end, State), + schedule_tick(), + {noreply, State}; +handle_info({'DOWN', Ref, process, Pid, _}, State) -> + case maps:get(Pid, State, undefined) of + Ref -> + case lookup(Pid) of + {Pid, GuildId, _, _, _} -> guild_read_model:delete(GuildId, Pid); + _ -> ok + end, + ets:delete(?TABLE, Pid), + {noreply, maps:remove(Pid, State)}; + _ -> + {noreply, State} + end; +handle_info(_, State) -> + {noreply, State}. + +-spec terminate(term(), map()) -> ok. +terminate(_, _) -> ok. + +-spec code_change(term(), map(), term()) -> {ok, map()}. +code_change(_, State, _) -> {ok, State}. + +-spec lookup(pid()) -> tuple() | undefined. +lookup(Pid) -> + try ets:lookup(?TABLE, Pid) of + [Entry] -> Entry; + [] -> undefined + catch + error:badarg -> undefined + end. + +-spec track_guild(pid(), map()) -> map(). +track_guild(Pid, State) -> + case maps:is_key(Pid, State) orelse lookup(Pid) =:= undefined of + true -> State; + false -> State#{Pid => erlang:monitor(process, Pid)} + end. + +-spec check_guild(pid(), integer()) -> ok. +check_guild(Pid, Now) -> + case lookup(Pid) of + {Pid, GuildId, Degraded, Targets, undefined} -> + case should_probe(Pid, Degraded) of + true -> + Ref = make_ref(), + ets:insert(?TABLE, {Pid, GuildId, Degraded, Targets, {Ref, Now}}), + Pid ! {guild_health_probe, Ref}, + ok; + false -> + ok + end; + {Pid, _, Degraded, _, {_Ref, SentAt}} = Entry -> + update(Entry, Degraded orelse Now - SentAt > ?DEGRADED_MS); + _ -> + ok + end. + +-spec should_probe(pid(), boolean()) -> boolean(). +should_probe(_Pid, true) -> + true; +should_probe(Pid, false) -> + case process_info(Pid, [message_queue_len, status]) of + [{message_queue_len, Len}, {status, Status}] -> Len > 0 orelse Status =:= running; + undefined -> false + end. + +-spec next_degraded(boolean(), non_neg_integer()) -> boolean(). +next_degraded(_Degraded, Delay) when Delay > ?DEGRADED_MS -> true; +next_degraded(_Degraded, Delay) when Delay < ?RECOVERED_MS -> false; +next_degraded(Degraded, _Delay) -> Degraded. + +-spec update(tuple(), boolean()) -> ok. +update({Pid, GuildId, Previous, Targets, Pending}, Degraded) -> + ets:insert(?TABLE, {Pid, GuildId, Degraded, Targets, Pending}), + case Previous =/= Degraded of + true -> notify_targets(Targets, GuildId, Pid, Degraded); + false -> ok + end. + +-spec notify_targets(ets:table(), integer(), pid(), boolean()) -> ok. +notify_targets(Targets, GuildId, Pid, Degraded) -> + try ets:tab2list(Targets) of + Rows -> + lists:foreach( + fun({_SessionId, SessionPid}) -> + notify_session(SessionPid, GuildId, Pid, Degraded) + end, + Rows + ) + catch + error:badarg -> ok + end. + +-spec notify_session(pid(), integer(), pid(), boolean()) -> ok. +notify_session(SessionPid, GuildId, GuildPid, Degraded) -> + gen_server:cast(SessionPid, {guild_health, GuildId, GuildPid, Degraded}). + +-spec discover_guilds() -> ok. +discover_guilds() -> + try ets:tab2list(guild_pid_cache) of + Rows -> + lists:foreach( + fun + ({_GuildId, Pid}) when is_pid(Pid), node(Pid) =:= node() -> + Pid ! guild_health_register; + (_) -> + ok + end, + Rows + ) + catch + error:badarg -> ok + end. + +-spec schedule_tick() -> reference(). +schedule_tick() -> erlang:send_after(?INTERVAL_MS, self(), tick). + +-ifdef(TEST). +-include_lib("eunit/include/eunit.hrl"). + +hysteresis_test() -> + ?assert(next_degraded(false, 2001)), + ?assertNot(next_degraded(false, 2000)), + ?assert(next_degraded(true, 500)), + ?assertNot(next_degraded(true, 499)), + ?assert(next_degraded(true, 3000)). + +unreachable_remote_pid_is_unknown_test() -> + Pid = binary_to_term(<<131, 88, 119, 12, "fake@nowhere", 1:32, 0:32, 1:32>>), + ?assertNotEqual(node(), node(Pid)), + {ElapsedUs, Result} = timer:tc(?MODULE, is_degraded, [Pid]), + ?assertNot(Result), + ?assert(ElapsedUs < 1000000). + +pending_probe_is_bounded_and_recovers_only_after_fresh_reply_test() -> + ok = guild_ets_owner:ensure_table(?TABLE, [named_table, public, set]), + Targets = ets:new(health_test_targets, [set]), + ets:insert(Targets, {<<"session">>, self()}), + Pid = spawn(fun() -> + receive + stop -> ok + end + end), + Entry = {Pid, 42, false, Targets, undefined}, + ets:insert(?TABLE, Entry), + try + Pid ! queued_work, + check_guild(Pid, 100), + {Pid, 42, false, Targets, {Ref, 100}} = lookup(Pid), + check_guild(Pid, 2201), + ?assert(is_degraded(Pid)), + receive + {'$gen_cast', {guild_health, 42, Pid, true}} -> ok + after 100 -> ?assert(false) + end, + check_guild(Pid, 5000), + ?assertEqual({message_queue_len, 2}, process_info(Pid, message_queue_len)), + {noreply, #{}} = handle_cast({pong, Pid, Ref, 5000}, #{}), + ?assert(is_degraded(Pid)), + check_guild(Pid, 5100), + {Pid, 42, true, Targets, {Ref2, 5100}} = lookup(Pid), + {noreply, #{}} = handle_cast({pong, Pid, Ref, 5110}, #{}), + ?assert(is_degraded(Pid)), + {noreply, #{}} = handle_cast({pong, Pid, Ref2, 5110}, #{}), + ?assertNot(is_degraded(Pid)), + receive + {'$gen_cast', {guild_health, 42, Pid, false}} -> ok + after 100 -> ?assert(false) + end + after + Pid ! stop, + ets:delete(?TABLE, Pid), + ets:delete(Targets) + end. + +session_targets_follow_replacement_and_removal_test() -> + Tab = ets:new(health_test_targets, [set]), + State = #{guild_health_sessions => Tab, sessions => #{<<"s">> => #{pid => self()}}}, + ok = put_session(<<"s">>, State), + ?assertEqual([{<<"s">>, self()}], ets:tab2list(Tab)), + Other = spawn(fun() -> + receive + stop -> ok + end + end), + ok = put_session(<<"s">>, State#{sessions => #{<<"s">> => #{pid => Other}}}), + ?assertEqual([{<<"s">>, Other}], ets:tab2list(Tab)), + ok = remove_session(<<"s">>, State), + ?assertEqual([], ets:tab2list(Tab)), + Other ! stop, + ets:delete(Tab). + +idle_guild_is_not_probed_test() -> + ok = guild_ets_owner:ensure_table(?TABLE, [named_table, public, set]), + Targets = ets:new(health_idle_targets, [set]), + Parent = self(), + Pid = spawn(fun() -> + Parent ! {idle, self()}, + receive + stop -> ok + end + end), + receive + {idle, Pid} -> ok + end, + ets:insert(?TABLE, {Pid, 43, false, Targets, undefined}), + try + check_guild(Pid, 100), + check_guild(Pid, 1000), + ?assertEqual({message_queue_len, 0}, process_info(Pid, message_queue_len)) + after + Pid ! stop, + ets:delete(?TABLE, Pid), + ets:delete(Targets) + end. + +foreign_target_table_is_rebuilt_test() -> + Parent = self(), + Owner = spawn(fun() -> + Tab = ets:new(health_foreign_targets, [set, public]), + Parent ! {foreign, Tab}, + receive + stop -> ok + end + end), + receive + {foreign, Foreign} -> + State = ensure_targets(#{ + guild_health_sessions => Foreign, sessions => #{<<"s">> => #{pid => self()}} + }), + Local = maps:get(guild_health_sessions, State), + ?assertNotEqual(Foreign, Local), + ?assertEqual(self(), ets:info(Local, owner)), + ?assertEqual([{<<"s">>, self()}], ets:tab2list(Local)), + ets:delete(Local) + end, + Owner ! stop. + +-endif. diff --git a/fluxer_gateway/src/guild/guild_member_list_engine.erl b/fluxer_gateway/src/guild/guild_member_list_engine.erl index a86f2599a..b0b04ca2e 100644 --- a/fluxer_gateway/src/guild/guild_member_list_engine.erl +++ b/fluxer_gateway/src/guild/guild_member_list_engine.erl @@ -18,7 +18,8 @@ get_all_item_keys/1, get_sorted_user_ids/1, index_of/2, - is_member_online/2 + is_member_online/2, + version/1 ]). -export([info/1]). @@ -42,6 +43,7 @@ new() -> {hoisted_role_ids, []}, {total_count, 0}, {online_count, 0}, + {version, 0}, {{section_count, ?ONLINE_IDX}, 0}, {{section_count, ?OFFLINE_IDX}, 0} ]), @@ -86,7 +88,7 @@ bulk_load(Ref, Members, HoistedRoleIds) -> guild_member_list_oset:from_sorted(OSet, Keys), store_counts(Ref, Total, Online), store_section_counts(Ref, SC), - ok. + bump_version(Ref). -spec bulk_load_member( {integer(), binary(), [integer()], boolean()}, @@ -130,7 +132,7 @@ add_member(Ref, UserId, SortKey, RoleIds, IsOnline) -> _ = ets:update_counter(Ref, {section_count, SIdx}, 1), _ = ets:update_counter(Ref, total_count, 1), ok = adjust_online_counter(Ref, IsOnline, 1), - ok. + bump_version(Ref). -spec remove_member(ets:table(), integer()) -> ok. remove_member(Ref, UserId) -> @@ -276,7 +278,7 @@ do_remove(Ref, UserId) -> _ = ets:update_counter(Ref, {section_count, SIdx}, -1), _ = ets:update_counter(Ref, total_count, -1), ok = adjust_online_counter(Ref, IsOnline, -1), - ok; + bump_version(Ref); [] -> ok end. @@ -326,6 +328,20 @@ move_section(Ref, OSet, ITab, UserId, SortKey, OldSIdx, RoleIds, IsOnline) -> false -> -1 end, _ = ets:update_counter(Ref, online_count, OnlineDelta), + bump_version(Ref). + +-spec version(ets:table()) -> non_neg_integer() | undefined. +version(Ref) -> + try ets:lookup_element(Ref, version, 2) of + Version when is_integer(Version) -> Version; + _ -> undefined + catch + error:badarg -> undefined + end. + +-spec bump_version(ets:table()) -> ok. +bump_version(Ref) -> + _ = ets:update_counter(Ref, version, 1, {version, 0}), ok. -spec adjust_online_counter(ets:table(), boolean(), integer()) -> ok. @@ -410,7 +426,7 @@ do_set_hoisted_roles(Ref, NewHoistedRoleIds) -> {hoisted_role_ids, NewHoistedRoleIds} ]), store_section_counts(Ref, SC), - ok. + bump_version(Ref). -spec rebuild_members_into( ets:table(), guild_member_list_oset:oset(), #{integer() => non_neg_integer()} diff --git a/fluxer_gateway/src/guild/guild_member_list_read.erl b/fluxer_gateway/src/guild/guild_member_list_read.erl index d934daf46..031f0e4e8 100644 --- a/fluxer_gateway/src/guild/guild_member_list_read.erl +++ b/fluxer_gateway/src/guild/guild_member_list_read.erl @@ -14,6 +14,8 @@ member_list_snapshot/2, snapshot/2, hydrate_engine_items/2, + with_member_item_memo/1, + note_presence_write/1, get_members_cursor/2 ]). @@ -27,12 +29,73 @@ list_id := list_id(), hide_offline := boolean(), member_map := map(), + member_revision := term(), presence_context := map(), state := guild_state() }. -export_type([guild_state/0, list_id/0, range/0, group_item/0, list_item/0, store_item/0]). +-define(MEMBER_ITEM_MEMO_KEY, guild_member_list_member_item_memo). +-define(PRESENCE_WRITES_KEY, guild_member_list_presence_writes). +-define(PRESENCE_WRITES_KEPT, 256). + +-type range_stamp() :: {non_neg_integer(), boolean(), reference(), ets:tid()}. +-type range_entry() :: {range_stamp(), non_neg_integer(), map(), [term()], map(), map()}. + +-spec note_presence_write(integer()) -> ok. +note_presence_write(UserId) -> + {Seq, Kept, Log} = presence_writes(), + {Kept1, Log1} = + case Kept >= 2 * ?PRESENCE_WRITES_KEPT of + true -> {?PRESENCE_WRITES_KEPT, lists:sublist(Log, ?PRESENCE_WRITES_KEPT)}; + false -> {Kept, Log} + end, + _ = erlang:put(?PRESENCE_WRITES_KEY, {Seq + 1, Kept1 + 1, [{Seq + 1, UserId} | Log1]}), + ok. + +-spec presence_writes() -> {non_neg_integer(), non_neg_integer(), [{pos_integer(), integer()}]}. +presence_writes() -> + case erlang:get(?PRESENCE_WRITES_KEY) of + {Seq, Kept, Log} when is_integer(Seq), is_integer(Kept), is_list(Log) -> + {Seq, Kept, Log}; + _ -> + {0, 0, []} + end. + +-spec presence_written_since(non_neg_integer()) -> {ok, [integer()]} | stale. +presence_written_since(Since) -> + {Seq, _Kept, Log} = presence_writes(), + case Seq - Since of + 0 -> + {ok, []}; + Count when Count > 0 -> + Written = [ + UserId + || {WriteSeq, UserId} <- lists:sublist(Log, Count), WriteSeq > Since + ], + case length(Written) =:= Count of + true -> {ok, Written}; + false -> stale + end; + _ -> + stale + end. + +-spec with_member_item_memo(fun(() -> T)) -> T. +with_member_item_memo(Fun) -> + case erlang:get(?MEMBER_ITEM_MEMO_KEY) of + undefined -> + _ = erlang:put(?MEMBER_ITEM_MEMO_KEY, #{}), + try + Fun() + after + erlang:erase(?MEMBER_ITEM_MEMO_KEY) + end; + _ -> + Fun() + end. + -spec get_member_groups(list_id(), guild_state()) -> [group_item()]. get_member_groups(ListId, State) -> case store_ref_for(ListId, State) of @@ -131,6 +194,7 @@ read_context(ListId, StoreGroups, State) -> list_id => ListId, hide_offline => offline_hidden(StoreGroups), member_map => guild_data_index:member_map(Data), + member_revision => maps:get(member_list_revision, Data, undefined), presence_context => guild_member_list_connected:presence_context(State), state => State }. @@ -142,8 +206,156 @@ read_context(ListId, StoreGroups, State) -> read_context() ) -> map(). range_sync_op_from_context(Ref, {Start, End} = Range, StoreGroups, ReadCtx) -> + case range_stamp(Ref, ReadCtx) of + undefined -> + {_Keys, Op} = build_range_sync_op(Ref, Range, StoreGroups, ReadCtx), + Op; + Stamp -> + Key = {sync_range, Ref, Start, End}, + Entry = cached_range_entry(Key, Stamp, Ref, Range, StoreGroups, ReadCtx), + ok = guild_member_list_subscribe:sync_cache_store(Key, Entry), + element(5, Entry) + end. + +-spec range_stamp(guild_member_list_store:store_ref(), read_context()) -> + range_stamp() | undefined. +range_stamp(Ref, #{hide_offline := HideOffline, member_revision := Revision} = ReadCtx) -> + #{presence_context := PresenceCtx} = ReadCtx, + Presence = maps:get(member_presence, PresenceCtx, undefined), + case is_reference(Revision) andalso presence_writes_seen_here(Presence) of + true -> + case guild_member_list_engine:version(Ref) of + Version when is_integer(Version) -> {Version, HideOffline, Revision, Presence}; + _ -> undefined + end; + false -> + undefined + end. + +-spec presence_writes_seen_here(term()) -> boolean(). +presence_writes_seen_here(Presence) when is_reference(Presence) -> + try ets:info(eqwalizer:dynamic_cast(Presence), owner) of + Owner -> Owner =:= self() + catch + _:_ -> false + end; +presence_writes_seen_here(_) -> + false. + +-spec cached_range_entry( + term(), + range_stamp(), + guild_member_list_store:store_ref(), + range(), + [{binary(), non_neg_integer()}], + read_context() +) -> range_entry(). +cached_range_entry(Key, Stamp, Ref, Range, StoreGroups, ReadCtx) -> + case guild_member_list_subscribe:sync_cache_find(Key) of + {ok, {Stamp0, Seq0, Users, Keys, Op, Connected0}} -> + Connected = range_connected_users(Users, ReadCtx), + case changed_range_users(Stamp0, Stamp, Seq0, Users, Connected0, Connected) of + {ok, Changed} -> + patch_range_entry(Stamp, Users, Keys, Op, Changed, Connected, ReadCtx); + rebuild -> + new_range_entry(Stamp, Ref, Range, StoreGroups, ReadCtx) + end; + _ -> + new_range_entry(Stamp, Ref, Range, StoreGroups, ReadCtx) + end. + +-spec changed_range_users(term(), range_stamp(), term(), term(), map(), map()) -> + {ok, #{integer() => true}} | rebuild. +changed_range_users(Stamp, Stamp, Seq0, Users, Connected0, Connected) when + is_integer(Seq0), is_map(Users) +-> + case presence_written_since(Seq0) of + {ok, Written} -> + Rewritten = maps:from_list([{U, true} || U <- Written, maps:is_key(U, Users)]), + Reconnected = + case Connected0 =:= Connected of + true -> + #{}; + false -> + maps:filter( + fun(UserId, _) -> + maps:is_key(UserId, Connected0) =/= + maps:is_key(UserId, Connected) + end, + Users + ) + end, + {ok, maps:merge(Rewritten, Reconnected)}; + stale -> + rebuild + end; +changed_range_users(_Stamp0, _Stamp, _Seq0, _Users, _Connected0, _Connected) -> + rebuild. + +-spec range_connected_users(map(), read_context()) -> map(). +range_connected_users(Users, #{presence_context := PresenceCtx}) -> + Connected = maps:get(connected_user_ids, PresenceCtx, undefined), + maps:filter(fun(UserId, _) -> is_connected(UserId, Connected) end, Users). + +-spec is_connected(integer(), term()) -> boolean(). +is_connected(UserId, Connected) -> + try + sets:is_element(UserId, eqwalizer:dynamic_cast(Connected)) + catch + _:_ -> false + end. + +-spec patch_range_entry( + range_stamp(), map(), [term()], map(), #{integer() => true}, map(), read_context() +) -> range_entry(). +patch_range_entry(Stamp, Users, Keys, Op, Changed, Connected, _ReadCtx) when + map_size(Changed) =:= 0 +-> + {Stamp, presence_seq(), Users, Keys, Op, Connected}; +patch_range_entry( + Stamp, Users, Keys, #{<<"items">> := Items} = Op, Changed, Connected, ReadCtx +) -> + #{member_map := MemberMap, presence_context := PresenceCtx} = ReadCtx, + Patched = lists:zipwith( + fun + (UserId, _Item) when is_map_key(UserId, Changed) -> + {true, New} = hydrate_member_ref(UserId, MemberMap, PresenceCtx), + New; + (_Key, Item) -> + Item + end, + Keys, + Items + ), + {Stamp, presence_seq(), Users, Keys, Op#{<<"items">> => Patched}, Connected}. + +-spec new_range_entry( + range_stamp(), + guild_member_list_store:store_ref(), + range(), + [{binary(), non_neg_integer()}], + read_context() +) -> range_entry(). +new_range_entry(Stamp, Ref, Range, StoreGroups, ReadCtx) -> + Seq = presence_seq(), + {Keys, Op} = build_range_sync_op(Ref, Range, StoreGroups, ReadCtx), + Users = maps:from_list([{UserId, true} || UserId <- Keys, is_integer(UserId)]), + {Stamp, Seq, Users, Keys, Op, range_connected_users(Users, ReadCtx)}. + +-spec presence_seq() -> non_neg_integer(). +presence_seq() -> + element(1, presence_writes()). + +-spec build_range_sync_op( + guild_member_list_store:store_ref(), + range(), + [{binary(), non_neg_integer()}], + read_context() +) -> {[term()], map()}. +build_range_sync_op(Ref, {Start, End} = Range, StoreGroups, ReadCtx) -> StoreItems = store_range_items_from_context(Ref, StoreGroups, Start, End, ReadCtx), - sync_op(Range, hydrate_engine_items_from_context(StoreItems, ReadCtx)). + Keyed = hydrate_keyed_items_from_context(StoreItems, ReadCtx), + {[K || {K, _} <- Keyed], sync_op(Range, [Item || {_, Item} <- Keyed])}. -spec store_range_items_from_context( guild_member_list_store:store_ref(), @@ -194,17 +406,34 @@ hydrate_engine_items(ListId, StoreItems, StoreGroups, State) -> HideOffline = offline_hidden(StoreGroups), hydrate_engine_items_with_context(StoreItems, ListId, HideOffline, MemberMap, State). --spec hydrate_engine_items_from_context([store_item()], read_context()) -> [list_item()]. -hydrate_engine_items_from_context(StoreItems, #{ +-spec hydrate_keyed_items_from_context([store_item()], read_context()) -> + [{integer() | group, list_item()}]. +hydrate_keyed_items_from_context(StoreItems, #{ list_id := ListId, hide_offline := HideOffline, member_map := MemberMap, presence_context := PresenceCtx, state := State }) -> - hydrate_engine_items_with_context( - StoreItems, ListId, HideOffline, MemberMap, PresenceCtx, State - ). + {Keyed, _Section} = lists:foldl( + fun(StoreItem, {Acc, Section}) -> + case + hydrate_store_item( + StoreItem, ListId, HideOffline, MemberMap, PresenceCtx, State, {[], Section} + ) + of + {[Item], Section1} -> {[{store_item_key(StoreItem), Item} | Acc], Section1}; + {[], Section1} -> {Acc, Section1} + end + end, + {[], undefined}, + StoreItems + ), + lists:reverse(Keyed). + +-spec store_item_key(store_item()) -> integer() | group. +store_item_key({member, UserId}) -> UserId; +store_item_key({group, _Id, _Count}) -> group. -spec hydrate_engine_items_with_context( [store_item()], list_id(), boolean(), map(), guild_state() @@ -265,14 +494,39 @@ hydrate_member_ref(UserId, MemberMap, PresenceCtx) -> error -> false; {ok, Member} -> - {true, #{ - <<"member">> => - guild_member_list_connected:add_presence_to_member( - Member, UserId, PresenceCtx - ) - }} + {true, member_item(UserId, Member, PresenceCtx)} end. +-spec member_item(integer(), map(), map()) -> list_item(). +member_item(UserId, Member, PresenceCtx) -> + Presence = maps:get(member_presence, PresenceCtx, undefined), + case {erlang:get(?MEMBER_ITEM_MEMO_KEY), is_reference(Presence)} of + {Items, true} when is_map(Items) -> + Connected = is_connected( + UserId, maps:get(connected_user_ids, PresenceCtx, undefined) + ), + Stamp = {Presence, Connected, presence_seq()}, + case maps:find(UserId, Items) of + {ok, {Member, Stamp, Item}} -> + Item; + _ -> + Item = build_member_item(UserId, Member, PresenceCtx), + _ = erlang:put(?MEMBER_ITEM_MEMO_KEY, Items#{ + UserId => {Member, Stamp, Item} + }), + Item + end; + _ -> + build_member_item(UserId, Member, PresenceCtx) + end. + +-spec build_member_item(integer(), map(), map()) -> list_item(). +build_member_item(UserId, Member, PresenceCtx) -> + #{ + <<"member">> => + guild_member_list_connected:add_presence_to_member(Member, UserId, PresenceCtx) + }. + -spec get_members_cursor(map(), guild_state()) -> {reply, map(), guild_state()}. get_members_cursor(Request, State) -> guild_member_list_read_cursor:get_members_cursor(Request, State). @@ -396,4 +650,365 @@ build_sync_response_without_store_still_emits_sync_windows_test() -> maps:get(<<"ops">>, Resp) ). +memo_member(UserId, Username) -> + #{<<"user">> => #{<<"id">> => integer_to_binary(UserId), <<"username">> => Username}}. + +memo_presence_ctx(Tab, Connected) -> + #{member_presence => Tab, connected_user_ids => sets:from_list(Connected)}. + +with_memo_presence_tab(Fun) -> + Tab = ets:new(memo_presence, [set, public]), + try + true = ets:insert(Tab, [ + {1, #{<<"status">> => <<"online">>, <<"mobile">> => true, <<"afk">> => false}}, + {2, #{<<"status">> => <<"idle">>, <<"mobile">> => false, <<"afk">> => true}} + ]), + Fun(Tab) + after + ets:delete(Tab) + end. + +member_item_memo_shares_items_within_scope_test() -> + with_memo_presence_tab(fun(Tab) -> + MemberMap = #{1 => memo_member(1, <<"one">>), 2 => memo_member(2, <<"two">>)}, + Ctx = memo_presence_ctx(Tab, [1, 2]), + Unscoped1 = hydrate_member_ref(1, MemberMap, Ctx), + Unscoped2 = hydrate_member_ref(1, MemberMap, Ctx), + ?assertNot(erts_debug:same(element(2, Unscoped1), element(2, Unscoped2))), + {Scoped1, Scoped2, Other} = with_member_item_memo(fun() -> + { + hydrate_member_ref(1, MemberMap, Ctx), + hydrate_member_ref(1, MemberMap, memo_presence_ctx(Tab, [1, 2])), + hydrate_member_ref(2, MemberMap, Ctx) + } + end), + ?assertEqual(Unscoped1, Scoped1), + ?assertEqual(Scoped1, Scoped2), + ?assert(erts_debug:same(element(2, Scoped1), element(2, Scoped2))), + ?assertEqual(hydrate_member_ref(2, MemberMap, Ctx), Other), + ?assertEqual( + false, with_member_item_memo(fun() -> hydrate_member_ref(3, MemberMap, Ctx) end) + ), + ?assertEqual(undefined, erlang:get(?MEMBER_ITEM_MEMO_KEY)) + end). + +member_item_memo_rebuilds_changed_member_test() -> + with_memo_presence_tab(fun(Tab) -> + Ctx = memo_presence_ctx(Tab, [1]), + Before = #{1 => memo_member(1, <<"one">>)}, + After = #{1 => memo_member(1, <<"renamed">>)}, + {Old, New} = with_member_item_memo(fun() -> + {hydrate_member_ref(1, Before, Ctx), hydrate_member_ref(1, After, Ctx)} + end), + ?assertEqual(hydrate_member_ref(1, Before, Ctx), Old), + ?assertEqual(hydrate_member_ref(1, After, Ctx), New), + ?assertNotEqual(Old, New) + end). + +member_item_memo_rebuilds_on_presence_context_change_test() -> + with_memo_presence_tab(fun(Tab) -> + MemberMap = #{1 => memo_member(1, <<"one">>)}, + Connected = memo_presence_ctx(Tab, [1]), + Disconnected = memo_presence_ctx(Tab, []), + {Online, Offline, OnlineAgain} = with_member_item_memo(fun() -> + { + hydrate_member_ref(1, MemberMap, Connected), + hydrate_member_ref(1, MemberMap, Disconnected), + hydrate_member_ref(1, MemberMap, Connected) + } + end), + {true, #{<<"member">> := #{<<"presence">> := OnlinePresence}}} = Online, + {true, #{<<"member">> := #{<<"presence">> := OfflinePresence}}} = Offline, + ?assertEqual(<<"online">>, maps:get(<<"status">>, OnlinePresence)), + ?assertEqual(true, maps:get(<<"mobile">>, OnlinePresence)), + ?assertEqual(guild_member_list_connected:default_presence(), OfflinePresence), + ?assertEqual(Online, OnlineAgain) + end). + +member_item_memo_scope_nests_and_clears_on_error_test() -> + with_memo_presence_tab(fun(Tab) -> + MemberMap = #{1 => memo_member(1, <<"one">>)}, + Ctx = memo_presence_ctx(Tab, [1]), + with_member_item_memo(fun() -> + Outer = hydrate_member_ref(1, MemberMap, Ctx), + Inner = with_member_item_memo(fun() -> hydrate_member_ref(1, MemberMap, Ctx) end), + ?assert(erts_debug:same(element(2, Outer), element(2, Inner))), + ?assertNotEqual(undefined, erlang:get(?MEMBER_ITEM_MEMO_KEY)) + end), + ?assertEqual(undefined, erlang:get(?MEMBER_ITEM_MEMO_KEY)), + ?assertError(boom, with_member_item_memo(fun() -> erlang:error(boom) end)), + ?assertEqual(undefined, erlang:get(?MEMBER_ITEM_MEMO_KEY)) + end). + +member_item_memo_sync_responses_match_unscoped_test() -> + with_memo_presence_tab(fun(Tab) -> + Ref = guild_member_list_engine:new(), + Other = guild_member_list_engine:new(), + try + ok = guild_member_list_engine:bulk_load( + Ref, [{1, <<"one">>, [], true}, {2, <<"two">>, [], false}], [] + ), + ok = guild_member_list_engine:bulk_load(Other, [{1, <<"one">>, [], true}], []), + State = #{ + id => 9, + data => #{ + <<"members">> => #{ + 1 => memo_member(1, <<"one">>), 2 => memo_member(2, <<"two">>) + } + }, + member_presence => Tab, + connected_user_ids => sets:from_list([1]), + channel_member_list_engines => #{<<"500">> => Ref, <<"600">> => Other} + }, + Build = fun() -> + [ + build_sync_response(9, ListId, [{0, 99}], State) + || ListId <- [<<"500">>, <<"600">>, <<"500">>] + ] + end, + ?assertEqual(Build(), with_member_item_memo(Build)) + after + guild_member_list_engine:destroy(Ref), + guild_member_list_engine:destroy(Other) + end + end). + +range_presence(Status) -> + #{ + <<"status">> => Status, + <<"mobile">> => false, + <<"afk">> => false, + <<"custom_status">> => null + }. + +with_range_world(Fun) -> + Tab = ets:new(range_presence, [set, public]), + Ref = guild_member_list_engine:new(), + _ = erlang:erase(guild_member_list_sync_item_cache), + _ = erlang:erase(?PRESENCE_WRITES_KEY), + try + Users = lists:seq(1, 6), + true = ets:insert(Tab, [{U, range_presence(<<"online">>)} || U <- Users]), + ok = guild_member_list_engine:bulk_load( + Ref, + [{U, <<"u", (integer_to_binary(U))/binary>>, [], U =< 4} || U <- Users], + [] + ), + State = #{ + id => 9, + data => #{ + member_list_revision => make_ref(), + <<"members">> => maps:from_list([ + {U, memo_member(U, <<"u", (integer_to_binary(U))/binary>>)} + || U <- Users + ]) + }, + member_presence => Tab, + connected_user_ids => sets:from_list([1, 2, 3, 4]), + channel_member_list_engines => #{<<"500">> => Ref} + }, + Fun(State, Tab, Ref) + after + _ = erlang:erase(guild_member_list_sync_item_cache), + _ = erlang:erase(?PRESENCE_WRITES_KEY), + guild_member_list_engine:destroy(Ref), + ets:delete(Tab), + flush_range_timers() + end. + +flush_range_timers() -> + receive + {timeout, _, member_list_sync_item_cache_rotate} -> flush_range_timers() + after 0 -> + ok + end. + +range_op(State) -> + [Op] = maps:get(<<"ops">>, build_sync_response(9, <<"500">>, [{0, 99}], State)), + Op. + +uncached_range_op(State) -> + Parent = self(), + Pid = spawn(fun() -> Parent ! {range_op, self(), range_op(State)} end), + receive + {range_op, Pid, Op} -> Op + after 5000 -> + erlang:error(uncached_range_op_timeout) + end. + +range_items_by_user(Op) -> + maps:from_list([ + {maps:get(<<"id">>, maps:get(<<"user">>, M)), I} + || #{<<"member">> := M} = I <- maps:get(<<"items">>, Op) + ]). + +range_cache_reuses_unchanged_range_test() -> + with_range_world(fun(State, _Tab, _Ref) -> + First = range_op(State), + Second = range_op(State), + ?assert(erts_debug:same(First, Second)), + ?assertEqual(uncached_range_op(State), Second), + ok = note_presence_write(77), + ?assert(erts_debug:same(First, range_op(State))) + end). + +range_cache_patches_only_the_written_member_test() -> + with_range_world(fun(State, Tab, _Ref) -> + Before = range_items_by_user(range_op(State)), + true = ets:insert(Tab, {2, range_presence(<<"dnd">>)}), + ok = note_presence_write(2), + Op = range_op(State), + ?assertEqual(uncached_range_op(State), Op), + After = range_items_by_user(Op), + ?assertNotEqual(maps:get(<<"2">>, Before), maps:get(<<"2">>, After)), + [ + ?assert(erts_debug:same(maps:get(Id, Before), maps:get(Id, After))) + || Id <- maps:keys(Before), Id =/= <<"2">> + ] + end). + +range_cache_follows_presence_deletes_test() -> + with_range_world(fun(State, Tab, _Ref) -> + _ = range_op(State), + true = ets:delete(Tab, 3), + ok = note_presence_write(3), + ?assertEqual(uncached_range_op(State), range_op(State)) + end). + +range_cache_rebuilds_after_engine_change_test() -> + with_range_world(fun(State, _Tab, Ref) -> + First = range_op(State), + ok = guild_member_list_engine:set_online(Ref, 2, false), + Op = range_op(State), + ?assertNotEqual(First, Op), + ?assertEqual(uncached_range_op(State), Op) + end). + +range_cache_patches_connectivity_change_test() -> + with_range_world(fun(State, _Tab, _Ref) -> + _ = range_op(State), + Disconnected = State#{connected_user_ids => sets:from_list([1, 3, 4])}, + Op = range_op(Disconnected), + ?assertEqual(uncached_range_op(Disconnected), Op), + ?assertEqual(uncached_range_op(State), range_op(State)) + end). + +range_cache_rebuilds_on_member_change_test() -> + with_range_world(fun(State, _Tab, _Ref) -> + _ = range_op(State), + #{data := Data} = State, + Renamed = State#{ + data => guild_data_index:put_member(memo_member(5, <<"renamed">>), Data) + }, + ?assertEqual(uncached_range_op(Renamed), range_op(Renamed)) + end). + +range_cache_rebuilds_when_write_log_overflows_test() -> + with_range_world(fun(State, Tab, _Ref) -> + _ = range_op(State), + true = ets:insert(Tab, {1, range_presence(<<"idle">>)}), + ok = note_presence_write(1), + [ok = note_presence_write(1000 + N) || N <- lists:seq(1, 3 * ?PRESENCE_WRITES_KEPT)], + ?assertEqual(stale, presence_written_since(0)), + ?assertEqual(uncached_range_op(State), range_op(State)) + end). + +range_cache_is_off_when_disabled_test() -> + with_range_world(fun(State, _Tab, _Ref) -> + application:set_env(fluxer_gateway, member_list_sync_item_cache_enabled, false), + try + First = range_op(State), + ?assertNot(erts_debug:same(First, range_op(State))), + ?assertEqual(undefined, erlang:get(guild_member_list_sync_item_cache)) + after + application:unset_env(fluxer_gateway, member_list_sync_item_cache_enabled) + end + end). + +presence_written_since_reports_writes_in_order_test() -> + _ = erlang:erase(?PRESENCE_WRITES_KEY), + try + ?assertEqual({ok, []}, presence_written_since(0)), + ok = note_presence_write(5), + ok = note_presence_write(6), + ?assertEqual({ok, [6, 5]}, presence_written_since(0)), + ?assertEqual({ok, [6]}, presence_written_since(1)), + ?assertEqual({ok, []}, presence_written_since(2)), + ?assertEqual(stale, presence_written_since(3)) + after + erlang:erase(?PRESENCE_WRITES_KEY) + end. + +range_cache_requires_a_member_revision_test() -> + with_range_world(fun(State, _Tab, _Ref) -> + #{data := Data} = State, + lists:foreach( + fun(LegacyData) -> + Legacy = State#{data => LegacyData}, + First = range_op(Legacy), + ?assertEqual(uncached_range_op(Legacy), First), + ?assertNot(erts_debug:same(First, range_op(Legacy))), + ?assertEqual(undefined, erlang:get(guild_member_list_sync_item_cache)) + end, + [maps:remove(member_list_revision, Data), Data#{member_list_revision => invalid}] + ) + end). + +range_cache_does_not_retain_unrendered_members_or_connections_test() -> + with_range_world(fun(State, _Tab, _Ref) -> + #{data := Data} = State, + Members = guild_data_index:member_map(Data), + SmallData = guild_data_index:put_member_map(Members, Data), + Small = State#{data => SmallData}, + SmallOp = range_op(Small), + SmallSize = erts_debug:flat_size(erlang:get(guild_member_list_sync_item_cache)), + Extra = maps:from_list([{U, memo_member(U, <<"extra">>)} || U <- lists:seq(7, 5000)]), + LargeData = guild_data_index:put_member_map(maps:merge(Members, Extra), SmallData), + Large = Small#{ + data => LargeData, + connected_user_ids => sets:from_list([1, 2, 3, 4] ++ lists:seq(7, 5000)) + }, + _ = erlang:erase(guild_member_list_sync_item_cache), + ?assertEqual(SmallOp, range_op(Large)), + ?assertEqual( + SmallSize, erts_debug:flat_size(erlang:get(guild_member_list_sync_item_cache)) + ) + end). + +range_cache_does_not_retain_map_backed_presence_test() -> + with_range_world(fun(State, Tab, _Ref) -> + MapState = State#{member_presence => maps:from_list(ets:tab2list(Tab))}, + First = range_op(MapState), + ?assertEqual(uncached_range_op(MapState), First), + ?assertNot(erts_debug:same(First, range_op(MapState))), + ?assertEqual(undefined, erlang:get(guild_member_list_sync_item_cache)) + end). + +member_item_memo_does_not_retain_unrelated_connections_test() -> + with_memo_presence_tab(fun(Tab) -> + MemberMap = #{1 => memo_member(1, <<"one">>)}, + Small = memo_presence_ctx(Tab, [1]), + Large = memo_presence_ctx(Tab, lists:seq(1, 5000)), + with_member_item_memo(fun() -> + {true, First} = hydrate_member_ref(1, MemberMap, Small), + SmallSize = erts_debug:flat_size(erlang:get(?MEMBER_ITEM_MEMO_KEY)), + {true, Second} = hydrate_member_ref(1, MemberMap, Large), + ?assert(erts_debug:same(First, Second)), + ?assertEqual(SmallSize, erts_debug:flat_size(erlang:get(?MEMBER_ITEM_MEMO_KEY))) + end) + end). + +member_item_memo_follows_presence_writes_in_scope_test() -> + with_memo_presence_tab(fun(Tab) -> + MemberMap = #{1 => memo_member(1, <<"one">>)}, + Ctx = memo_presence_ctx(Tab, [1]), + with_member_item_memo(fun() -> + Before = hydrate_member_ref(1, MemberMap, Ctx), + true = ets:insert(Tab, {1, range_presence(<<"dnd">>)}), + ok = note_presence_write(1), + After = hydrate_member_ref(1, MemberMap, Ctx), + ?assertNotEqual(Before, After) + end) + end). + -endif. diff --git a/fluxer_gateway/src/guild/guild_member_list_subscribe.erl b/fluxer_gateway/src/guild/guild_member_list_subscribe.erl index c47beb840..63d60d9d2 100644 --- a/fluxer_gateway/src/guild/guild_member_list_subscribe.erl +++ b/fluxer_gateway/src/guild/guild_member_list_subscribe.erl @@ -9,7 +9,9 @@ send_member_list_update_to_sessions/5, dispatch_sync_to_subscribed_list/7, dispatch_sync_to_subscribed_sessions/6, - handle_sync_item_cache_timeout/1 + handle_sync_item_cache_timeout/1, + sync_cache_find/1, + sync_cache_store/2 ]). -type guild_state() :: map(). @@ -270,13 +272,37 @@ encode_sync_payload(Payload) -> -spec encode_sync_payload_cached(map(), list()) -> {pre_encoded, binary()}. encode_sync_payload_cached(Payload, Ops) -> {TimerRef, Current, Previous} = sync_item_cache(), - {FragmentOps, {Current1, Previous1}} = - lists:mapfoldl(fun fragment_op/2, {Current, Previous}, Ops), + Key = sync_payload_key(Payload, Ops), + {Encoded, {Current1, Previous1}} = + case Current of + #{Key := {Payload, Cached}} -> + {Cached, {Current, Previous}}; + _ -> + encode_sync_payload_fragments(Key, Payload, Ops, {Current, Previous}) + end, _ = erlang:put(?SYNC_ITEM_CACHE_KEY, {TimerRef, Current1, Previous1}), + Encoded. + +-spec encode_sync_payload_fragments(term(), map(), list(), {map(), map()}) -> + {{pre_encoded, binary()}, {map(), map()}}. +encode_sync_payload_fragments(Key, Payload, Ops, Generations) -> + {FragmentOps, Generations1} = lists:mapfoldl(fun fragment_op/2, Generations, Ops), WirePayload = eqwalizer:dynamic_cast( guild_data_wire:payload(Payload#{<<"ops">> => FragmentOps}) ), - {pre_encoded, iolist_to_binary(json:encode(WirePayload, fun encode_fragment_value/2))}. + Encoded = + {pre_encoded, iolist_to_binary(json:encode(WirePayload, fun encode_fragment_value/2))}, + {Encoded, cache_fragment(Key, {Payload, Encoded}, Generations1)}. + +-spec sync_payload_key(map(), list()) -> term(). +sync_payload_key(Payload, Ops) -> + {sync_payload, maps:get(<<"id">>, Payload, undefined), [op_range(Op) || Op <- Ops]}. + +-spec op_range(term()) -> term(). +op_range(#{<<"range">> := Range}) -> + Range; +op_range(_Op) -> + undefined. -spec fragment_op(term(), {map(), map()}) -> {term(), {map(), map()}}. fragment_op(#{<<"items">> := Items} = Op, Generations) when is_list(Items) -> @@ -309,8 +335,7 @@ previous_or_encoded_fragment(Id, Item, Previous) -> {json_fragment, iolist_to_binary(json:encode(guild_data_wire:payload(Item)))} end. --spec cache_fragment(term(), {term(), {json_fragment, binary()}}, {map(), map()}) -> - {map(), map()}. +-spec cache_fragment(term(), {term(), term()}, {map(), map()}) -> {map(), map()}. cache_fragment(Id, Entry, {Current, _Previous}) when map_size(Current) >= ?SYNC_ITEM_CACHE_MAX_ENTRIES -> @@ -324,6 +349,26 @@ encode_fragment_value({json_fragment, Encoded}, _Encode) -> encode_fragment_value(Value, Encode) -> json:encode_value(Value, Encode). +-spec sync_cache_find(term()) -> {ok, term()} | error. +sync_cache_find(Key) -> + case {sync_item_cache_enabled(), erlang:get(?SYNC_ITEM_CACHE_KEY)} of + {true, {_TimerRef, #{Key := Value}, _Previous}} -> {ok, Value}; + {true, {_TimerRef, _Current, #{Key := Value}}} -> {ok, Value}; + _ -> error + end. + +-spec sync_cache_store(term(), term()) -> ok. +sync_cache_store(Key, Value) -> + case sync_item_cache_enabled() of + true -> + {TimerRef, Current, Previous} = sync_item_cache(), + {Current1, Previous1} = cache_fragment(Key, Value, {Current, Previous}), + _ = erlang:put(?SYNC_ITEM_CACHE_KEY, {TimerRef, Current1, Previous1}), + ok; + false -> + ok = erase_sync_item_cache() + end. + -spec sync_item_cache() -> {reference(), map(), map()}. sync_item_cache() -> case erlang:get(?SYNC_ITEM_CACHE_KEY) of @@ -528,6 +573,7 @@ sync_item_cache_keeps_entries_touched_each_generation_test() -> ?assertEqual(8, map_size(previous_entries())), assert_matches_uncached(Payload), ?assertEqual(8, map_size(current_entries())), + ?assertEqual(1, map_size(payload_entries())), rotate_now(), ?assertEqual(8, map_size(previous_entries())), rotate_now(), @@ -570,6 +616,52 @@ dispatch_sync_group_sends_uncached_bytes_test() -> ?assertEqual([Expected, Expected], received_dispatches()) end). +sync_payload_equal_content_reuses_encoded_bytes_test() -> + with_clean_cache(fun() -> + Members = [test_member(N) || N <- lists:seq(1, 30)], + Payload = sync_payload(<<"500">>, [{0, 99}], Members), + assert_matches_uncached(Payload), + [Key] = maps:keys(payload_entries()), + Sentinel = {pre_encoded, <<"sentinel">>}, + {TimerRef, Current, Previous} = erlang:get(?SYNC_ITEM_CACHE_KEY), + _ = erlang:put( + ?SYNC_ITEM_CACHE_KEY, {TimerRef, Current#{Key => {Payload, Sentinel}}, Previous} + ), + Copy = binary_to_term(term_to_binary(Payload)), + ?assertEqual(Sentinel, encode_sync_payload(Copy)), + [First | Rest] = Members, + Afk = put_in(First, [<<"member">>, <<"presence">>, <<"afk">>], true), + assert_matches_uncached(sync_payload(<<"500">>, [{0, 99}], [Afk | Rest])), + ?assertEqual([Key], maps:keys(payload_entries())) + end). + +sync_payload_cache_is_keyed_by_list_and_ranges_test() -> + with_clean_cache(fun() -> + Members = [test_member(N) || N <- lists:seq(1, 40)], + Payloads = [ + sync_payload(<<"500">>, [{0, 99}], Members), + sync_payload(<<"600">>, [{0, 99}], Members), + sync_payload(<<"500">>, [{0, 19}], Members), + sync_payload(<<"500">>, [{0, 19}, {20, 39}], Members) + ], + [assert_matches_uncached(P) || P <- Payloads ++ lists:reverse(Payloads)], + ?assertEqual(4, map_size(payload_entries())), + Header = (hd(Payloads))#{<<"online_count">> => 1}, + assert_matches_uncached(Header), + ?assertEqual(4, map_size(payload_entries())) + end). + +sync_payload_flag_off_ignores_cached_payload_test() -> + with_clean_cache(fun() -> + Payload = sync_payload(<<"500">>, [{0, 99}], [test_member(N) || N <- lists:seq(1, 5)]), + assert_matches_uncached(Payload), + ?assertEqual(1, map_size(payload_entries())), + with_cache_flag(false, fun() -> + assert_matches_uncached(Payload), + ?assertEqual(#{}, payload_entries()) + end) + end). + assert_matches_uncached(Payload) -> Expected = uncached(Payload), Actual = encode_sync_payload(Payload), @@ -686,13 +778,22 @@ flush_rotation_timers() -> current_entries() -> case erlang:get(?SYNC_ITEM_CACHE_KEY) of - {_, Current, _} -> Current; + {_, Current, _} -> fragment_entries(Current); _ -> #{} end. previous_entries() -> case erlang:get(?SYNC_ITEM_CACHE_KEY) of - {_, _, Previous} -> Previous; + {_, _, Previous} -> fragment_entries(Previous); + _ -> #{} + end. + +fragment_entries(Generation) -> + maps:filter(fun(Key, _) -> is_binary(Key) end, Generation). + +payload_entries() -> + case erlang:get(?SYNC_ITEM_CACHE_KEY) of + {_, Current, _} -> maps:filter(fun(Key, _) -> not is_binary(Key) end, Current); _ -> #{} end. diff --git a/fluxer_gateway/src/guild/guild_member_list_sync_batch.erl b/fluxer_gateway/src/guild/guild_member_list_sync_batch.erl index 86285bc84..2c614d703 100644 --- a/fluxer_gateway/src/guild/guild_member_list_sync_batch.erl +++ b/fluxer_gateway/src/guild/guild_member_list_sync_batch.erl @@ -94,9 +94,11 @@ dispatch_pending_syncs(ListIds, State) -> dispatch_pending_syncs(ListIds, State, SubsTab) -> _ = guild_member_list_write_context:with_guild_id(State, fun(GuildId) -> Sessions = maps:get(sessions, State, #{}), - dispatch_pending_syncs_for_guild( - GuildId, lists:sort(ListIds), Sessions, State, SubsTab - ), + guild_member_list_read:with_member_item_memo(fun() -> + dispatch_pending_syncs_for_guild( + GuildId, lists:sort(ListIds), Sessions, State, SubsTab + ) + end), {ok, State} end), State. diff --git a/fluxer_gateway/src/guild/guild_member_list_write.erl b/fluxer_gateway/src/guild/guild_member_list_write.erl index 8507fc85b..ebbec28a3 100644 --- a/fluxer_gateway/src/guild/guild_member_list_write.erl +++ b/fluxer_gateway/src/guild/guild_member_list_write.erl @@ -221,15 +221,17 @@ fold_connection_change_lists(GuildId, UserId, Mark, State, SubsTab) -> dispatch_user_change_to_subscribed_lists( UserId, OldMember, NewMember, SubsTab, State ) -> - lists:foldl( - fun(ListId, AccState) -> - dispatch_user_change_to_subscribed_list( - UserId, OldMember, NewMember, ListId, AccState - ) - end, - State, - guild_member_list_subs:list_ids(SubsTab) - ). + guild_member_list_read:with_member_item_memo(fun() -> + lists:foldl( + fun(ListId, AccState) -> + dispatch_user_change_to_subscribed_list( + UserId, OldMember, NewMember, ListId, AccState + ) + end, + State, + guild_member_list_subs:list_ids(SubsTab) + ) + end). -spec dispatch_user_change_to_subscribed_list( user_id(), diff --git a/fluxer_gateway/src/guild/guild_presence.erl b/fluxer_gateway/src/guild/guild_presence.erl index 7d70c61c3..17cc3e0ba 100644 --- a/fluxer_gateway/src/guild/guild_presence.erl +++ b/fluxer_gateway/src/guild/guild_presence.erl @@ -380,6 +380,7 @@ find_member_by_user_id(UserId, State) -> store_member_presence(UserId, PresenceMap, State) -> Tab = maps:get(member_presence, State), ets:insert(Tab, {UserId, PresenceMap}), + ok = guild_member_list_read:note_presence_write(UserId), State. -ifdef(TEST). diff --git a/fluxer_gateway/src/guild/guild_query_handler.erl b/fluxer_gateway/src/guild/guild_query_handler.erl index 9faf84a40..833529d85 100644 --- a/fluxer_gateway/src/guild/guild_query_handler.erl +++ b/fluxer_gateway/src/guild/guild_query_handler.erl @@ -19,21 +19,31 @@ -spec call(pid(), {atom(), map()}, pos_integer()) -> term(). call(GuildPid, {Tag, Request}, Timeout) -> Deadline = os:system_time(millisecond) + Timeout, - gen_server:call(GuildPid, {Tag, Request#{deadline => Deadline}}, Timeout). + MonotonicDeadline = erlang:monotonic_time(millisecond) + Timeout, + gen_server:call( + GuildPid, + {Tag, Request#{deadline => Deadline, deadline_monotonic => MonotonicDeadline}}, + Timeout + ). -spec handle_call(term(), gen_server:from(), guild_state()) -> {reply, term(), guild_state()} | {noreply, guild_state()}. handle_call(Msg, From, State) -> - case is_expired(Msg) of + case is_expired(Msg, From) of true -> {noreply, State}; false -> handle_query(Msg, From, State) end. --spec is_expired(term()) -> boolean(). -is_expired({_Tag, #{deadline := Deadline}}) when is_integer(Deadline) -> - os:system_time(millisecond) > Deadline; -is_expired(_Msg) -> +-spec is_expired(term(), gen_server:from()) -> boolean(). +is_expired({_Tag, #{deadline_monotonic := Deadline}}, {Caller, _ReplyTag}) when + is_integer(Deadline) +-> + case gateway_clock_offset:offset(node(Caller)) of + undefined -> false; + Offset -> erlang:monotonic_time(millisecond) > Deadline - Offset + end; +is_expired(_Msg, _From) -> false. -spec handle_query(term(), gen_server:from(), guild_state()) -> @@ -455,3 +465,36 @@ resolve_data_payload(#{<<"members">> := _} = State) -> State; resolve_data_payload(_State) -> undefined. + +-ifdef(TEST). +-include_lib("eunit/include/eunit.hrl"). + +monotonic_deadline_ignores_the_legacy_wall_clock_test() -> + From = {self(), make_ref()}, + Now = erlang:monotonic_time(millisecond), + ?assertNot( + is_expired({get_data, #{deadline => 0, deadline_monotonic => Now + 5000}}, From) + ), + ?assert( + is_expired( + {get_data, #{deadline => 9999999999999, deadline_monotonic => Now - 1}}, From + ) + ), + ?assertNot(is_expired({get_data, #{deadline => 0}}, From)). + +call_keeps_the_legacy_deadline_and_adds_a_monotonic_deadline_test() -> + Guild = spawn(fun() -> + receive + {'$gen_call', From, {get_data, Request}} -> gen_server:reply(From, Request) + end + end), + WallBefore = os:system_time(millisecond), + MonotonicBefore = erlang:monotonic_time(millisecond), + #{deadline := WallDeadline, deadline_monotonic := MonotonicDeadline} = + call(Guild, {get_data, #{}}, 2000), + ?assert(WallDeadline >= WallBefore + 2000), + ?assert(WallDeadline =< os:system_time(millisecond) + 2000), + ?assert(MonotonicDeadline >= MonotonicBefore + 2000), + ?assert(MonotonicDeadline =< erlang:monotonic_time(millisecond) + 2000). + +-endif. diff --git a/fluxer_gateway/src/guild/guild_read_model.erl b/fluxer_gateway/src/guild/guild_read_model.erl new file mode 100644 index 000000000..3f53625a5 --- /dev/null +++ b/fluxer_gateway/src/guild/guild_read_model.erl @@ -0,0 +1,123 @@ +%% SPDX-License-Identifier: AGPL-3.0-or-later + +-module(guild_read_model). +-typing([eqwalizer]). + +-export([put_state/1, update/2, delete/1, delete/2, query/2]). + +-define(TABLE, guild_read_model). +-define(STATE_KEYS, [id, member_count, member_list_engine, virtual_channel_access]). +-define(DATA_KEYS, [ + <<"guild">>, + <<"roles">>, + <<"channels">>, + <<"channel_index">>, + <<"emojis">>, + <<"stickers">>, + role_perms_cache, + overwrite_perms_cache, + members_ets +]). + +-spec put_state(map()) -> ok. +put_state(#{id := GuildId, data := Data} = State) when is_integer(GuildId), is_map(Data) -> + case maps:get(disable_permission_cache_updates, State, false) of + true -> ok; + false -> store(GuildId, State, Data) + end; +put_state(_) -> + ok. + +-spec store(integer(), map(), map()) -> ok. +store(GuildId, State, Data) -> + ok = guild_ets_utils:ensure_table(?TABLE, [ + named_table, public, set, {read_concurrency, true} + ]), + ReadData = maps:with(?DATA_KEYS, Data), + Members = + case maps:get(members_ets, Data, undefined) of + Tab when is_reference(Tab) -> #{}; + _ -> guild_data_index:member_map(Data) + end, + Snapshot = (maps:with(?STATE_KEYS, State))#{ + data => ReadData#{<<"members">> => Members, members_normalized => Members} + }, + true = ets:insert(?TABLE, {GuildId, self(), Snapshot}), + ok. + +-spec update(map(), map()) -> ok. +update(OldState, NewState) -> + Keys = [data | ?STATE_KEYS], + case + lists:all( + fun(Key) -> + erts_debug:same( + maps:get(Key, OldState, undefined), maps:get(Key, NewState, undefined) + ) + end, + Keys + ) + of + true -> ok; + false -> put_state(NewState) + end. + +-spec delete(term()) -> ok. +delete(GuildId) -> + delete(GuildId, self()). + +-spec delete(term(), pid()) -> ok. +delete(GuildId, Owner) -> + try ets:match_delete(?TABLE, {GuildId, Owner, '_'}) of + true -> ok + catch + error:badarg -> ok + end. + +-spec query(integer(), {atom(), map()}) -> {ok, map()} | miss. +query(GuildId, Request) -> + try + case ets:lookup(?TABLE, GuildId) of + [{GuildId, Owner, Snapshot}] -> + case is_process_alive(Owner) of + true -> read(Request, Snapshot); + false -> miss + end; + [] -> + miss + end + catch + error:badarg -> miss; + error:{badmatch, false} -> miss + end. + +-spec read({atom(), map()}, map()) -> {ok, map()} | miss. +read({Tag, Request}, Snapshot) -> + State = materialize_member(maps:get(user_id, Request, null), Snapshot), + case Tag of + get_guild_data -> reply(guild_data:get_guild_data(Request, State)); + get_guild_auth_context -> reply(guild_data:get_auth_context(Request, State)); + get_guild_member -> reply(guild_data:get_guild_member(Request, State)); + has_member -> reply(guild_data:has_member(Request, State)); + _ -> miss + end. + +-spec materialize_member(integer() | null, map()) -> map(). +materialize_member(UserId, #{data := #{members_ets := Tab} = Data} = Snapshot) -> + MemberCount = ets:info(Tab, size), + true = is_integer(MemberCount), + Members = + case ets:lookup(Tab, UserId) of + [{UserId, Member}] -> #{UserId => Member}; + [] -> #{} + end, + Snapshot#{ + member_count => MemberCount, + data => Data#{<<"members">> => Members, members_normalized => Members} + }; +materialize_member(_UserId, Snapshot) -> + Snapshot. + +-spec reply({reply, map(), map()}) -> {ok, map()}. +reply({reply, Reply, _State}) -> + {ok, Reply}. diff --git a/fluxer_gateway/src/guild/guild_request_counts.erl b/fluxer_gateway/src/guild/guild_request_counts.erl index 3de4d2cbb..f5f625b3a 100644 --- a/fluxer_gateway/src/guild/guild_request_counts.erl +++ b/fluxer_gateway/src/guild/guild_request_counts.erl @@ -100,19 +100,34 @@ spawn_fetch_worker(Self, Tag, GuildId, GuildPid, UserId) -> -spec worker(pid(), reference(), integer(), pid(), integer()) -> ok. worker(Parent, Tag, GuildId, GuildPid, UserId) -> - Request = {get_viewer_counts, #{user_id => UserId}}, - Result = - try guild_query_handler:call(GuildPid, Request, ?GUILD_CALL_TIMEOUT_MS) of - #{member_count := MemberCount, online_count := OnlineCount} -> - {ok, MemberCount, OnlineCount}; - _ -> - error - catch - _:_ -> error - end, - Parent ! {Tag, GuildId, Result}, + Parent ! {Tag, GuildId, fetch_counts(GuildPid, UserId)}, ok. +-spec fetch_counts(pid(), integer()) -> {ok, non_neg_integer(), non_neg_integer()} | error. +fetch_counts(GuildPid, UserId) -> + Request = {get_viewer_counts, #{user_id => UserId}}, + try guild_query_handler:call(GuildPid, Request, ?GUILD_CALL_TIMEOUT_MS) of + ok -> fetch_legacy_counts(GuildPid, UserId); + Reply -> counts_result(Reply) + catch + _:_ -> error + end. + +-spec fetch_legacy_counts(pid(), integer()) -> + {ok, non_neg_integer(), non_neg_integer()} | error. +fetch_legacy_counts(GuildPid, UserId) -> + try gen_server:call(GuildPid, {get_user_counts, UserId}, ?GUILD_CALL_TIMEOUT_MS) of + Reply -> counts_result(Reply) + catch + _:_ -> error + end. + +-spec counts_result(term()) -> {ok, non_neg_integer(), non_neg_integer()} | error. +counts_result(#{member_count := MemberCount, online_count := OnlineCount}) -> + {ok, MemberCount, OnlineCount}; +counts_result(_) -> + error. + -spec collect_responses(non_neg_integer(), reference(), integer(), [map()]) -> [map()]. collect_responses(0, _Tag, _Deadline, Acc) -> lists:reverse(Acc); @@ -254,6 +269,91 @@ handle_request_fetches_viewer_counts_with_deadline_test() -> ?assert(false) end. +legacy_guild(Replies) -> + spawn(fun() -> legacy_guild_loop(Replies) end). + +legacy_guild_loop(Replies) -> + receive + {'$gen_call', From, {get_viewer_counts, #{user_id := 100, deadline := D}}} when + is_integer(D) + -> + gen_server:reply(From, ok), + legacy_guild_loop(Replies); + {'$gen_call', From, {get_user_counts, 100}} -> + [Reply | Rest] = Replies, + gen_server:reply(From, Reply), + legacy_guild_loop(Rest) + after 5000 -> + ok + end. + +request_counts_payload(Guilds) -> + Self = self(), + SessionState = #{session_pid => Self, user_id => <<"100">>, guilds => Guilds}, + GuildIds = [integer_to_binary(Id) || Id <- maps:keys(Guilds)], + ok = handle_request(#{<<"guild_ids">> => GuildIds}, Self, SessionState), + receive + {'$gen_cast', {dispatch, guild_counts_update, Payload}} -> Payload + after 1000 -> + error(no_dispatch) + end. + +handle_request_falls_back_to_user_counts_on_legacy_guild_test() -> + Guild = legacy_guild([#{member_count => 50, online_count => 10}]), + Payload = request_counts_payload(#{7 => {Guild, make_ref()}}), + ?assertEqual([build_entry(7, 50, 10)], maps:get(<<"counts">>, Payload)). + +handle_request_omits_guild_when_legacy_fallback_fails_test() -> + Guild = legacy_guild([ok]), + Payload = request_counts_payload(#{7 => {Guild, make_ref()}}), + ?assertEqual([], maps:get(<<"counts">>, Payload)). + +handle_request_mixes_legacy_and_current_guilds_test() -> + Legacy = legacy_guild([#{member_count => 50, online_count => 10}]), + Current = spawn(fun() -> + receive + {'$gen_call', From, {get_viewer_counts, #{user_id := 100}}} -> + gen_server:reply(From, #{member_count => 80, online_count => 20}) + end, + receive + {'$gen_call', From2, _} -> gen_server:reply(From2, unexpected) + after 500 -> ok + end + end), + Payload = request_counts_payload(#{7 => {Legacy, make_ref()}, 9 => {Current, make_ref()}}), + ?assertEqual( + [build_entry(7, 50, 10), build_entry(9, 80, 20)], + lists:sort(maps:get(<<"counts">>, Payload)) + ). + +fetch_counts_does_not_fall_back_on_current_guild_test() -> + Self = self(), + Guild = spawn(fun() -> + receive + {'$gen_call', From, {get_viewer_counts, #{user_id := 100}}} -> + gen_server:reply(From, #{member_count => 3, online_count => 1}) + end, + receive + {'$gen_call', From2, Msg} -> + Self ! {unexpected_call, Msg}, + gen_server:reply(From2, ok) + after 300 -> ok + end + end), + ?assertEqual({ok, 3, 1}, fetch_counts(Guild, 100)), + receive + {unexpected_call, Msg} -> ?assertEqual(none, Msg) + after 400 -> ok + end. + +fetch_counts_errors_on_malformed_reply_test() -> + Guild = spawn(fun() -> + receive + {'$gen_call', From, {get_viewer_counts, _}} -> gen_server:reply(From, #{}) + end + end), + ?assertEqual(error, fetch_counts(Guild, 100)). + parse_nonce_test() -> ?assertEqual(<<"x">>, parse_nonce(<<"x">>)), ?assertEqual(<<"abc">>, parse_nonce(<<"abc">>)), diff --git a/fluxer_gateway/src/guild/guild_sessions_connect.erl b/fluxer_gateway/src/guild/guild_sessions_connect.erl index 85bb8e231..46f6957d9 100644 --- a/fluxer_gateway/src/guild/guild_sessions_connect.erl +++ b/fluxer_gateway/src/guild/guild_sessions_connect.erl @@ -351,6 +351,7 @@ remove_session_id(undefined, Sessions) -> -spec put_session_ref(session_id(), term(), guild_state()) -> guild_state(). put_session_ref(SessionId, Ref, State) when is_reference(Ref) -> + ok = guild_health:put_session(SessionId, State), Refs = maps:get(guild_session_refs, State, #{}), State#{guild_session_refs => Refs#{Ref => SessionId}}; put_session_ref(_SessionId, _Ref, State) -> @@ -359,6 +360,7 @@ put_session_ref(_SessionId, _Ref, State) -> -spec remove_session_ref(term(), guild_state()) -> guild_state(). remove_session_ref(Ref, State) when is_reference(Ref) -> Refs = maps:get(guild_session_refs, State, #{}), + ok = guild_health:remove_session(maps:get(Ref, Refs, undefined), State), State#{guild_session_refs => maps:remove(Ref, Refs)}; remove_session_ref(_Ref, State) -> State. diff --git a/fluxer_gateway/src/guild/guild_sessions_presence.erl b/fluxer_gateway/src/guild/guild_sessions_presence.erl index 498f5ebb3..b74958e2b 100644 --- a/fluxer_gateway/src/guild/guild_sessions_presence.erl +++ b/fluxer_gateway/src/guild/guild_sessions_presence.erl @@ -98,6 +98,7 @@ handle_user_offline(UserId, State) -> remove_member_presence(UserId, State) -> Tab = maps:get(member_presence, State), ets:delete(Tab, UserId), + ok = guild_member_list_read:note_presence_write(UserId), State. -spec maybe_send_cached_presence(user_id(), guild_state()) -> guild_state(). diff --git a/fluxer_gateway/src/session/session.erl b/fluxer_gateway/src/session/session.erl index 868b995b8..6cb35bd99 100644 --- a/fluxer_gateway/src/session/session.erl +++ b/fluxer_gateway/src/session/session.erl @@ -48,6 +48,7 @@ offline_timer => {reference(), reference()} | undefined, guilds => #{guild_id() => guild_ref()}, active_guilds => sets:set(guild_id()), + guild_health => #{guild_id() => {pid(), boolean()}}, calls => #{channel_id() => call_ref()}, channels => #{channel_id() => map()}, ready => map() | undefined, @@ -161,6 +162,10 @@ handle_cast_presences(Presences, State) -> -spec handle_cast_guild_or_lifecycle(term(), session_state()) -> {noreply, session_state()} | {stop, normal, session_state()}. +handle_cast_guild_or_lifecycle({guild_health, GuildId, GuildPid, Degraded}, State) when + is_integer(GuildId), is_pid(GuildPid), is_boolean(Degraded) +-> + session_guild_health:handle_update(GuildId, GuildPid, Degraded, State); handle_cast_guild_or_lifecycle(handoff_fence, State) -> session_lifecycle:handle_handoff_fence(State); handle_cast_guild_or_lifecycle({reconnect_drain, SocketPid}, State) when diff --git a/fluxer_gateway/src/session/session_connection_guild.erl b/fluxer_gateway/src/session/session_connection_guild.erl index 92b491c5e..bcd2462c1 100644 --- a/fluxer_gateway/src/session/session_connection_guild.erl +++ b/fluxer_gateway/src/session/session_connection_guild.erl @@ -459,8 +459,10 @@ finalize_guild_monitor(GuildId, GuildPid, Guilds0, State, ReadyFun) -> session_state() ) -> session_result(). apply_ready_fun(GuildId, GuildPid, ReadyFun, State) -> - case ReadyFun(State) of + State1 = session_guild_health:forget(GuildId, State), + case ReadyFun(State1) of {noreply, ReadyState} -> + ok = guild_health:send_current(GuildPid, self()), ReplayedState = maybe_replay_guild_subscriptions(GuildId, GuildPid, ReadyState), {noreply, session_dm_partners:register_guild(GuildId, GuildPid, ReplayedState)}; {stop, normal, ReadyState} -> diff --git a/fluxer_gateway/src/session/session_guild_health.erl b/fluxer_gateway/src/session/session_guild_health.erl new file mode 100644 index 000000000..8ecae50b3 --- /dev/null +++ b/fluxer_gateway/src/session/session_guild_health.erl @@ -0,0 +1,119 @@ +%% SPDX-License-Identifier: AGPL-3.0-or-later + +-module(session_guild_health). +-typing([eqwalizer]). + +-export([handle_update/4, dispatch_ready/1, forget/2]). + +-spec handle_update(integer(), pid(), boolean(), map()) -> {noreply, map()}. +handle_update(GuildId, GuildPid, Degraded, State) -> + case maps:get(GuildId, maps:get(guilds, State, #{}), undefined) of + {GuildPid, _Ref} -> apply_update(GuildId, GuildPid, Degraded, State); + _ -> {noreply, State} + end. + +-spec apply_update(integer(), pid(), boolean(), map()) -> {noreply, map()}. +apply_update(GuildId, GuildPid, Degraded, State) -> + Health = maps:get(guild_health, State, #{}), + Previous = + case maps:get(GuildId, Health, undefined) of + {_OldPid, Value} -> Value; + _ -> false + end, + State1 = State#{guild_health => Health#{GuildId => {GuildPid, Degraded}}}, + case Previous =/= Degraded andalso maps:get(ready, State, undefined) =:= undefined of + true -> {noreply, dispatch(GuildId, Degraded, State1)}; + false -> {noreply, State1} + end. + +-spec dispatch_ready(map()) -> map(). +dispatch_ready(State) -> + Guilds = maps:get(guilds, State, #{}), + maps:fold( + fun + (GuildId, {Pid, true}, Acc) -> + case maps:get(GuildId, Guilds, undefined) of + {Pid, _Ref} -> dispatch(GuildId, true, Acc); + _ -> Acc + end; + (_, _, Acc) -> + Acc + end, + State, + maps:get(guild_health, State, #{}) + ). + +-spec forget(integer(), map()) -> map(). +forget(GuildId, State) -> + State#{guild_health => maps:remove(GuildId, maps:get(guild_health, State, #{}))}. + +-spec dispatch(integer(), boolean(), map()) -> map(). +dispatch(GuildId, Degraded, State) -> + Payload = #{<<"guild_id">> => integer_to_binary(GuildId), <<"degraded">> => Degraded}, + {noreply, NewState} = session_dispatch:handle_dispatch(guild_health_update, Payload, State), + NewState. + +-ifdef(TEST). +-include_lib("eunit/include/eunit.hrl"). + +health_state() -> + #{ + id => <<"health">>, + user_id => 100, + seq => 0, + buffer => [], + guilds => #{42 => {self(), make_ref()}}, + socket_pid => undefined, + channels => #{}, + relationships => #{}, + ready => undefined, + suppress_presence_updates => false, + pending_presences => [], + presence_pid => undefined, + ignored_events => #{}, + debounce_reactions => false + }. + +health_transitions_replay_without_reconnect_test() -> + State = health_state(), + {noreply, Degraded} = handle_update(42, self(), true, State), + {noreply, Duplicate} = handle_update(42, self(), true, Degraded), + ?assertEqual(1, maps:get(seq, Duplicate)), + {noreply, Recovered} = handle_update(42, self(), false, Duplicate), + ?assertEqual(maps:get(guilds, State), maps:get(guilds, Recovered)), + Events = limited_deque:to_list(maps:get(buffer, Recovered)), + ?assertEqual([guild_health_update, guild_health_update], [maps:get(event, E) || E <- Events]), + ?assertEqual([true, false], [maps:get(<<"degraded">>, maps:get(data, E)) || E <- Events]). + +health_waits_for_ready_test() -> + State = (health_state())#{ready => #{}}, + {noreply, Degraded} = handle_update(42, self(), true, State), + ?assertEqual(0, maps:get(seq, Degraded)), + Ready = dispatch_ready(Degraded#{ready => undefined}), + ?assertEqual(1, maps:get(seq, Ready)). + +stale_guild_owner_is_ignored_test() -> + Other = spawn(fun() -> + receive + stop -> ok + end + end), + State = health_state(), + ?assertEqual({noreply, State}, handle_update(42, Other, true, State)), + Other ! stop. + +healthy_reconnect_clears_previous_owner_state_test() -> + State = (health_state())#{guild_health => #{42 => {old_owner, true}}}, + {noreply, Recovered} = handle_update(42, self(), false, State), + ?assertEqual(1, maps:get(seq, Recovered)), + ?assertEqual({self(), false}, maps:get(42, maps:get(guild_health, Recovered))). + +same_owner_reconnect_resends_degraded_state_test() -> + {noreply, Degraded} = handle_update(42, self(), true, health_state()), + Reconnected = forget(42, Degraded), + {noreply, Resynced} = handle_update(42, self(), true, Reconnected), + ?assertEqual(2, maps:get(seq, Resynced)), + Events = limited_deque:to_list(maps:get(buffer, Resynced)), + ?assertEqual([true, true], [maps:get(<<"degraded">>, maps:get(data, E)) || E <- Events]). + +-endif. diff --git a/fluxer_gateway/src/session/session_guilds.erl b/fluxer_gateway/src/session/session_guilds.erl index 285da4ad1..2952798c0 100644 --- a/fluxer_gateway/src/session/session_guilds.erl +++ b/fluxer_gateway/src/session/session_guilds.erl @@ -24,7 +24,9 @@ handle_guild_leave(GuildId, #{guilds := Guilds} = State) -> DeleteData = #{<<"id">> => integer_to_binary(GuildId)}, {noreply, DispatchedState} = session_dispatch:handle_dispatch(guild_delete, DeleteData, State), - State1 = DispatchedState#{guilds => maps:remove(GuildId, Guilds)}, + State1 = session_guild_health:forget(GuildId, DispatchedState#{ + guilds => maps:remove(GuildId, Guilds) + }), {noreply, remove_guild_subscription_state(GuildId, State1)}; _ -> {noreply, State} @@ -39,7 +41,10 @@ handle_forced_unavailable_guild_leave(GuildId, UnavailableHidden, #{guilds := Gu {noreply, State1} = session_dispatch:handle_dispatch(guild_delete, GuildDeleteData, State), self() ! {guild_connect, GuildId, 0}, - {noreply, State1#{guilds => Guilds#{GuildId => cached_unavailable}}}. + {noreply, + session_guild_health:forget(GuildId, State1#{ + guilds => Guilds#{GuildId => cached_unavailable} + })}. -spec demonitor_guild_if_connected(term()) -> ok. demonitor_guild_if_connected({Pid, Ref}) when is_pid(Pid), is_reference(Ref) -> diff --git a/fluxer_gateway/src/session/session_init.erl b/fluxer_gateway/src/session/session_init.erl index 087a05f7c..cc9137ae3 100644 --- a/fluxer_gateway/src/session/session_init.erl +++ b/fluxer_gateway/src/session/session_init.erl @@ -259,6 +259,7 @@ extract_extra_fields(D, Ready) -> collected_sessions => maps:get(collected_sessions, D, []), collected_presences => maps:get(collected_presences, D, []), guild_subscription_state => maps:get(guild_subscription_state, D, #{}), + guild_health => maps:get(guild_health, D, #{}), relationships => maps:get(relationships, D, load_relationships(Ready)), suppress_presence_updates => true, pending_presences => [], diff --git a/fluxer_gateway/src/session/session_lifecycle.erl b/fluxer_gateway/src/session/session_lifecycle.erl index b0ad4c7c7..0ef9de6f4 100644 --- a/fluxer_gateway/src/session/session_lifecycle.erl +++ b/fluxer_gateway/src/session/session_lifecycle.erl @@ -608,6 +608,7 @@ serialize_state(State) -> collected_guild_states => maps:get(collected_guild_states, State), collected_sessions => maps:get(collected_sessions, State), collected_presences => maps:get(collected_presences, State, []), + guild_health => maps:get(guild_health, State, #{}), guild_subscription_state => maps:get(guild_subscription_state, State, #{}) }. @@ -657,6 +658,7 @@ serialize_transfer_runtime(State) -> collected_guild_states => maps:get(collected_guild_states, State, []), collected_sessions => maps:get(collected_sessions, State, []), collected_presences => maps:get(collected_presences, State, []), + guild_health => maps:get(guild_health, State, #{}), guild_subscription_state => maps:get(guild_subscription_state, State, #{}) }. diff --git a/fluxer_gateway/src/session/session_ready_dispatch.erl b/fluxer_gateway/src/session/session_ready_dispatch.erl index ef8ee73fc..a0027d1c9 100644 --- a/fluxer_gateway/src/session/session_ready_dispatch.erl +++ b/fluxer_gateway/src/session/session_ready_dispatch.erl @@ -69,8 +69,9 @@ dispatch_ready_to_socket(State) -> StateAfterGuilds = dispatch_bot_guild_creates( IsBot, CollectedGuilds, Guilds, StateAfterReady ), - schedule_call_creates(StateAfterGuilds, SessionId), - StateAfterPresences = release_pending_presences(StateAfterGuilds, CollectedPresences), + StateAfterHealth = session_guild_health:dispatch_ready(StateAfterGuilds), + schedule_call_creates(StateAfterHealth, SessionId), + StateAfterPresences = release_pending_presences(StateAfterHealth, CollectedPresences), FinalState = StateAfterPresences#{ ready => undefined, collected_guild_states => [], diff --git a/fluxer_gateway/src/utils/event_atoms.erl b/fluxer_gateway/src/utils/event_atoms.erl index d66269ab0..bf8beb169 100644 --- a/fluxer_gateway/src/utils/event_atoms.erl +++ b/fluxer_gateway/src/utils/event_atoms.erl @@ -87,6 +87,7 @@ guild_event_map() -> <<"GUILD_CREATE">> => guild_create, <<"GUILD_DELETE">> => guild_delete, <<"GUILD_EMOJIS_UPDATE">> => guild_emojis_update, + <<"GUILD_HEALTH_UPDATE">> => guild_health_update, <<"GUILD_MEMBER_ADD">> => guild_member_add, <<"GUILD_MEMBER_LIST_UPDATE">> => guild_member_list_update, <<"GUILD_MEMBER_REMOVE">> => guild_member_remove, diff --git a/fluxer_gateway/test/fluxer_gateway_sup_tests.erl b/fluxer_gateway/test/fluxer_gateway_sup_tests.erl index fdff882af..fc8ab19ae 100644 --- a/fluxer_gateway/test/fluxer_gateway_sup_tests.erl +++ b/fluxer_gateway/test/fluxer_gateway_sup_tests.erl @@ -118,6 +118,7 @@ init_guilds_role_includes_only_guild_state_children_test() -> ?assert(lists:member(presence_bus, Ids)), ?assert(lists:member(guild_manager, Ids)), ?assert(lists:member(guild_counts_cache, Ids)), + ?assert(lists:member(gateway_clock_offset, Ids)), ?assert(lists:member(voice_state_counts_sync, Ids)), ?assert(lists:member(gateway_cluster_handoff, Ids)), ?assertNot(lists:member(session_manager, Ids)), @@ -137,6 +138,7 @@ init_sessions_role_includes_cluster_handoff_test() -> ?assert(lists:member(session_state_transfer, Ids)), ?assert(lists:member(gateway_cluster_handoff, Ids)), ?assertNot(lists:member(guild_manager, Ids)), + ?assertNot(lists:member(gateway_clock_offset, Ids)), ?assertNot(lists:member(presence_manager, Ids)), persistent_term:erase({fluxer_gateway, runtime_config}). diff --git a/fluxer_gateway/test/guild_data_tests.erl b/fluxer_gateway/test/guild_data_tests.erl index 0f463ebdc..f0d78fbd6 100644 --- a/fluxer_gateway/test/guild_data_tests.erl +++ b/fluxer_gateway/test/guild_data_tests.erl @@ -5,6 +5,124 @@ -include_lib("eunit/include/eunit.hrl"). +read_model_preserves_query_results_test() -> + State = read_model_state(), + try + ok = guild_read_model:put_state(State), + lists:foreach( + fun({Tag, Request, Handler}) -> + {reply, Expected, _} = Handler(Request, State), + ?assertEqual({ok, Expected}, guild_read_model:query(100, {Tag, Request})) + end, + [ + {get_guild_data, #{user_id => 200}, fun guild_data:get_guild_data/2}, + {get_guild_data, #{user_id => 999}, fun guild_data:get_guild_data/2}, + {get_guild_data, #{user_id => null}, fun guild_data:get_guild_data/2}, + {get_guild_auth_context, #{user_id => 200, channel_id => 500}, + fun guild_data:get_auth_context/2}, + {get_guild_auth_context, #{user_id => 999, channel_id => null}, + fun guild_data:get_auth_context/2}, + {get_guild_member, #{user_id => 200}, fun guild_data:get_guild_member/2}, + {get_guild_member, #{user_id => 999}, fun guild_data:get_guild_member/2}, + {has_member, #{user_id => 200}, fun guild_data:has_member/2} + ] + ) + after + cleanup_read_model(State) + end. + +read_model_uses_live_member_rows_without_copying_members_test() -> + State = read_model_state(), + #{data := #{members_ets := Tab}} = State, + try + ok = guild_read_model:put_state(State), + [{100, _, #{data := SnapshotData}}] = ets:lookup(guild_read_model, 100), + ?assertEqual(#{}, maps:get(<<"members">>, SnapshotData)), + [{200, Member}] = ets:lookup(Tab, 200), + Updated = Member#{ + <<"nick">> => <<"new nickname">>, <<"communication_disabled_until">> => null + }, + ets:insert(Tab, {200, Updated}), + ?assertEqual( + {ok, #{success => true, member_data => Updated}}, + guild_read_model:query(100, {get_guild_member, #{user_id => 200}}) + ), + ets:delete(Tab, 200), + ?assertEqual( + {ok, #{has_member => false}}, + guild_read_model:query(100, {has_member, #{user_id => 200}}) + ), + ?assertMatch( + {ok, #{guild_data := null}}, + guild_read_model:query(100, {get_guild_data, #{user_id => 200}}) + ) + after + cleanup_read_model(State) + end. + +read_model_observes_role_and_collection_changes_test() -> + State = read_model_state(), + try + ok = guild_read_model:put_state(State), + Data = maps:get(data, State), + UpdatedData = guild_data_index:normalize_map(Data#{ + <<"roles">> => [#{<<"id">> => 100, <<"permissions">> => 0}], + <<"emojis">> => [#{<<"id">> => 700, <<"name">> => <<"new">>}] + }), + Updated = State#{data => UpdatedData}, + ok = guild_read_model:update(State, Updated), + {reply, Expected, _} = guild_data:get_guild_data(#{user_id => 200}, Updated), + ?assertEqual( + {ok, Expected}, guild_read_model:query(100, {get_guild_data, #{user_id => 200}}) + ), + #{guild_data := GuildData} = Expected, + ?assertEqual([], maps:get(<<"channels">>, GuildData)) + after + cleanup_read_model(State) + end. + +read_model_survives_blocked_owner_and_rejects_dead_tables_test() -> + State = read_model_state(), + Self = self(), + Owner = spawn(fun() -> + ok = guild_read_model:put_state(State), + Self ! published, + receive + stop -> ok + end + end), + try + receive + published -> ok + after 1000 -> error(publish_timeout) + end, + ?assertMatch( + {ok, #{auth_context := #{}}}, + guild_read_model:query( + 100, {get_guild_auth_context, #{user_id => 200, channel_id => 500}} + ) + ), + ?assertEqual({message_queue_len, 0}, process_info(Owner, message_queue_len)), + ets:delete(maps:get(members_ets, maps:get(data, State))), + ?assertEqual(miss, guild_read_model:query(100, {has_member, #{user_id => 200}})) + after + Owner ! stop, + cleanup_read_model(State) + end. + +read_model_state() -> + State = test_state(), + Data = guild_data_index:normalize_map(maps:get(data, State)), + Members = guild_data_index:member_map(Data), + Tab = ets:new(read_model_members, [set, public]), + ets:insert(Tab, maps:to_list(Members)), + State#{member_count => map_size(Members), data => Data#{members_ets => Tab}}. + +cleanup_read_model(#{data := Data} = State) -> + guild_read_model:delete(100), + catch ets:delete(maps:get(members_ets, Data)), + ets:delete(maps:get(member_presence, State)). + get_guild_data_membership_gate_test() -> State = test_state(), {reply, Reply1, _} = guild_data:get_guild_data(#{user_id => 999}, State), diff --git a/fluxer_gateway/test/guild_member_list_engine_tests.erl b/fluxer_gateway/test/guild_member_list_engine_tests.erl index f6b9c8c16..c93e3663a 100644 --- a/fluxer_gateway/test/guild_member_list_engine_tests.erl +++ b/fluxer_gateway/test/guild_member_list_engine_tests.erl @@ -551,3 +551,53 @@ engine_roles_for_test_member(UserId) when UserId rem 5 =:= 0 -> [20]; engine_roles_for_test_member(_UserId) -> []. + +version_bumps_on_every_mutation_test() -> + Ref = guild_member_list_engine:new(), + try + V0 = guild_member_list_engine:version(Ref), + ok = guild_member_list_engine:bulk_load(Ref, [{1, <<"a">>, [], true}], []), + V1 = guild_member_list_engine:version(Ref), + ok = guild_member_list_engine:add_member(Ref, 2, <<"b">>, [], false), + V2 = guild_member_list_engine:version(Ref), + ok = guild_member_list_engine:update_member(Ref, 2, <<"c">>, [], false), + V3 = guild_member_list_engine:version(Ref), + ok = guild_member_list_engine:set_online(Ref, 2, true), + V4 = guild_member_list_engine:version(Ref), + changed = guild_member_list_engine:set_hoisted_roles(Ref, [7]), + V5 = guild_member_list_engine:version(Ref), + ok = guild_member_list_engine:remove_member(Ref, 1), + V6 = guild_member_list_engine:version(Ref), + Versions = [V0, V1, V2, V3, V4, V5, V6], + ?assertEqual(Versions, lists:usort(Versions)), + ?assertEqual(7, length(lists:usort(Versions))) + after + guild_member_list_engine:destroy(Ref) + end. + +version_is_stable_for_noop_mutations_test() -> + Ref = guild_member_list_engine:new(), + try + ok = guild_member_list_engine:add_member(Ref, 1, <<"a">>, [], true), + V = guild_member_list_engine:version(Ref), + ok = guild_member_list_engine:set_online(Ref, 1, true), + ok = guild_member_list_engine:set_online(Ref, 99, false), + ok = guild_member_list_engine:remove_member(Ref, 99), + unchanged = guild_member_list_engine:set_hoisted_roles(Ref, []), + _ = guild_member_list_engine:get_items(Ref, 0, 10), + ?assertEqual(V, guild_member_list_engine:version(Ref)) + after + guild_member_list_engine:destroy(Ref) + end. + +version_survives_tables_without_a_version_row_test() -> + Ref = guild_member_list_engine:new(), + try + true = ets:delete(Ref, version), + ?assertEqual(undefined, guild_member_list_engine:version(Ref)), + ok = guild_member_list_engine:add_member(Ref, 1, <<"a">>, [], true), + ?assertEqual(1, guild_member_list_engine:version(Ref)) + after + guild_member_list_engine:destroy(Ref) + end, + ?assertEqual(undefined, guild_member_list_engine:version(Ref)).