perf(gateway): speed up presence and harden guild queries (#3072)

This commit is contained in:
Hampus
2026-09-30 21:12:13 +02:00
committed by GitHub
parent 6e2f90b03c
commit ab0b483fbe
44 changed files with 2166 additions and 144 deletions
@@ -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)
@@ -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.
@@ -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),
@@ -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),
@@ -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(
+87 -44
View File
@@ -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),
+5 -1
View File
@@ -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]).
@@ -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.
@@ -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);
@@ -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.
@@ -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)),
@@ -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}
]).
+384
View File
@@ -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.
@@ -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()}
@@ -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.
@@ -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.
@@ -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.
@@ -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(),
@@ -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).
@@ -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.
@@ -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}.
+111 -11
View File
@@ -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">>)),
@@ -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.
@@ -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().
+5
View File
@@ -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
@@ -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} ->
@@ -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.
@@ -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) ->
@@ -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 => [],
@@ -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, #{})
}.
@@ -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 => [],
+1
View File
@@ -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,
@@ -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}).
+118
View File
@@ -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),
@@ -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)).