mirror of
https://github.com/fluxerapp/fluxer
synced 2026-10-07 19:22:14 +09:00
fix(gateway): stop anti-entropy resurrecting deleted presence (#2298)
This commit is contained in:
@@ -126,12 +126,14 @@ handoff_to_target(TargetNode) ->
|
||||
-spec put_local(integer(), map(), state()) -> {ok, state()}.
|
||||
put_local(UserId, Presence, State) ->
|
||||
{_Reply, NewState} = presence_cache_shards:forward_put(UserId, Presence, State),
|
||||
{ok, presence_cache_rebalance:increment_generation(NewState)}.
|
||||
Tracked = presence_cache_rebalance:record_put_tombstone(UserId, Presence, NewState),
|
||||
{ok, presence_cache_rebalance:increment_generation(Tracked)}.
|
||||
|
||||
-spec delete_local(integer(), state()) -> {ok, state()}.
|
||||
delete_local(UserId, State) ->
|
||||
{_Reply, NewState} = presence_cache_shards:forward_delete(UserId, State),
|
||||
{ok, presence_cache_rebalance:increment_generation(NewState)}.
|
||||
Tracked = presence_cache_rebalance:record_delete_tombstone(UserId, NewState),
|
||||
{ok, presence_cache_rebalance:increment_generation(Tracked)}.
|
||||
|
||||
-spec local_snapshot(state()) -> #{integer() => map()}.
|
||||
local_snapshot(State) ->
|
||||
@@ -149,6 +151,7 @@ init([]) ->
|
||||
pending_operations => #{},
|
||||
pending_retry_timer => undefined,
|
||||
pending_nodedown_cleanups => #{},
|
||||
delete_tombstones => #{},
|
||||
generation => 0,
|
||||
anti_entropy_timer => presence_cache_rebalance:schedule_anti_entropy()
|
||||
}}.
|
||||
|
||||
@@ -9,7 +9,9 @@
|
||||
perform_anti_entropy/1,
|
||||
handle_anti_entropy_request/3,
|
||||
handle_anti_entropy_digest_request/3,
|
||||
merge_anti_entropy_entries/2
|
||||
merge_anti_entropy_entries/2,
|
||||
record_delete/2,
|
||||
record_put/3
|
||||
]).
|
||||
|
||||
-export_type([state/0]).
|
||||
@@ -17,6 +19,8 @@
|
||||
-define(ANTI_ENTROPY_INTERVAL_MS, 30000).
|
||||
-define(ANTI_ENTROPY_MSG, anti_entropy_tick).
|
||||
-define(SNAPSHOT_CHUNK_SIZE, 500).
|
||||
-define(TOMBSTONE_TTL_MS, 90000).
|
||||
-define(TOMBSTONE_LIMIT, 50000).
|
||||
|
||||
-type state() :: map().
|
||||
|
||||
@@ -38,7 +42,7 @@ cancel_anti_entropy_timer(State) ->
|
||||
perform_anti_entropy(State) ->
|
||||
case persistent_term:get(presence_noop, false) of
|
||||
true -> State;
|
||||
false -> broadcast_anti_entropy_requests(State)
|
||||
false -> broadcast_anti_entropy_requests(prune_tombstones(State))
|
||||
end.
|
||||
|
||||
-spec handle_anti_entropy_request(node(), non_neg_integer(), state()) -> {noreply, state()}.
|
||||
@@ -69,6 +73,70 @@ merge_anti_entropy_entries(Entries, State) when is_map(Entries) ->
|
||||
merge_anti_entropy_entries(_, State) ->
|
||||
State.
|
||||
|
||||
-spec record_delete(integer(), state()) -> state().
|
||||
record_delete(UserId, State) ->
|
||||
Tombstones = tombstones(State),
|
||||
Updated = Tombstones#{UserId => now_ms() + ?TOMBSTONE_TTL_MS},
|
||||
State#{delete_tombstones => cap_tombstones(Updated)}.
|
||||
|
||||
-spec record_put(integer(), map(), state()) -> state().
|
||||
record_put(UserId, Presence, State) ->
|
||||
case is_visible_presence(Presence) of
|
||||
true -> clear_tombstone(UserId, State);
|
||||
false -> record_delete(UserId, State)
|
||||
end.
|
||||
|
||||
-spec clear_tombstone(integer(), state()) -> state().
|
||||
clear_tombstone(UserId, State) ->
|
||||
Tombstones = tombstones(State),
|
||||
case maps:is_key(UserId, Tombstones) of
|
||||
true -> State#{delete_tombstones => maps:remove(UserId, Tombstones)};
|
||||
false -> State
|
||||
end.
|
||||
|
||||
-spec tombstones(state()) -> #{integer() => integer()}.
|
||||
tombstones(State) ->
|
||||
case maps:get(delete_tombstones, State, #{}) of
|
||||
Tombstones when is_map(Tombstones) -> Tombstones;
|
||||
_ -> #{}
|
||||
end.
|
||||
|
||||
-spec prune_tombstones(state()) -> state().
|
||||
prune_tombstones(State) ->
|
||||
State#{delete_tombstones => drop_expired(tombstones(State), now_ms())}.
|
||||
|
||||
-spec drop_expired(#{integer() => integer()}, integer()) -> #{integer() => integer()}.
|
||||
drop_expired(Tombstones, Now) ->
|
||||
maps:filter(fun(_UserId, ExpiresAt) -> ExpiresAt > Now end, Tombstones).
|
||||
|
||||
-spec cap_tombstones(#{integer() => integer()}) -> #{integer() => integer()}.
|
||||
cap_tombstones(Tombstones) ->
|
||||
case maps:size(Tombstones) > ?TOMBSTONE_LIMIT of
|
||||
true -> evict_oldest_tombstones(drop_expired(Tombstones, now_ms()));
|
||||
false -> Tombstones
|
||||
end.
|
||||
|
||||
-spec evict_oldest_tombstones(#{integer() => integer()}) -> #{integer() => integer()}.
|
||||
evict_oldest_tombstones(Tombstones) ->
|
||||
Size = maps:size(Tombstones),
|
||||
case Size > ?TOMBSTONE_LIMIT of
|
||||
true -> keep_newest_tombstones(Tombstones, Size);
|
||||
false -> Tombstones
|
||||
end.
|
||||
|
||||
-spec keep_newest_tombstones(#{integer() => integer()}, non_neg_integer()) ->
|
||||
#{integer() => integer()}.
|
||||
keep_newest_tombstones(Tombstones, Size) ->
|
||||
Ordered = lists:sort(
|
||||
fun({_UserIdA, ExpiresA}, {_UserIdB, ExpiresB}) -> ExpiresA =< ExpiresB end,
|
||||
maps:to_list(Tombstones)
|
||||
),
|
||||
maps:from_list(lists:nthtail(Size - ?TOMBSTONE_LIMIT, Ordered)).
|
||||
|
||||
-spec now_ms() -> integer().
|
||||
now_ms() ->
|
||||
erlang:monotonic_time(millisecond).
|
||||
|
||||
-spec broadcast_anti_entropy_requests(state()) -> state().
|
||||
broadcast_anti_entropy_requests(State) ->
|
||||
Digest = presence_cache_shards:content_digest(State),
|
||||
@@ -124,11 +192,26 @@ merge_single_entry(_UserId, _Presence, AccState) ->
|
||||
merge_if_missing(UserId, Presence, State) ->
|
||||
case presence_cache_bulk:get_local_fast(UserId) of
|
||||
not_found ->
|
||||
merge_if_visible(UserId, Presence, State);
|
||||
merge_unless_deleted(UserId, Presence, State);
|
||||
{ok, _} ->
|
||||
State
|
||||
end.
|
||||
|
||||
-spec merge_unless_deleted(integer(), map(), state()) -> state().
|
||||
merge_unless_deleted(UserId, Presence, State) ->
|
||||
Tombstones = tombstones(State),
|
||||
Now = now_ms(),
|
||||
case maps:get(UserId, Tombstones, undefined) of
|
||||
ExpiresAt when is_integer(ExpiresAt), ExpiresAt > Now ->
|
||||
State;
|
||||
undefined ->
|
||||
merge_if_visible(UserId, Presence, State);
|
||||
_ ->
|
||||
merge_if_visible(
|
||||
UserId, Presence, State#{delete_tombstones => maps:remove(UserId, Tombstones)}
|
||||
)
|
||||
end.
|
||||
|
||||
-spec merge_if_visible(integer(), map(), state()) -> state().
|
||||
merge_if_visible(UserId, Presence, State) ->
|
||||
case is_visible_presence(Presence) of
|
||||
|
||||
@@ -10,6 +10,8 @@
|
||||
handle_anti_entropy_request/3,
|
||||
handle_anti_entropy_digest_request/3,
|
||||
merge_anti_entropy_entries/2,
|
||||
record_delete_tombstone/2,
|
||||
record_put_tombstone/3,
|
||||
schedule_anti_entropy/0,
|
||||
cancel_anti_entropy_timer/1,
|
||||
start_nodedown_grace/2,
|
||||
@@ -80,6 +82,14 @@ handle_anti_entropy_digest_request(FromNode, RemoteDigest, State) ->
|
||||
merge_anti_entropy_entries(Entries, State) ->
|
||||
presence_cache_anti_entropy:merge_anti_entropy_entries(Entries, State).
|
||||
|
||||
-spec record_delete_tombstone(integer(), state()) -> state().
|
||||
record_delete_tombstone(UserId, State) ->
|
||||
presence_cache_anti_entropy:record_delete(UserId, State).
|
||||
|
||||
-spec record_put_tombstone(integer(), map(), state()) -> state().
|
||||
record_put_tombstone(UserId, Presence, State) ->
|
||||
presence_cache_anti_entropy:record_put(UserId, Presence, State).
|
||||
|
||||
-spec start_nodedown_grace(node(), state()) -> state().
|
||||
start_nodedown_grace(Node, State) ->
|
||||
PendingCleanups = maps:get(pending_nodedown_cleanups, State, #{}),
|
||||
|
||||
@@ -230,6 +230,70 @@ evict_oldest_pending_drops_oldest_inserted_test() ->
|
||||
presence_cache_rebalance:cancel_pending_retry_timer(State1)
|
||||
end.
|
||||
|
||||
anti_entropy_does_not_resurrect_deleted_presence_test() ->
|
||||
with_local_presence_cache(fun(Pid) -> assert_delete_survives_stale_peer_entry(Pid) end).
|
||||
|
||||
anti_entropy_repairs_missing_presence_test() ->
|
||||
with_local_presence_cache(fun(Pid) -> assert_merge_repairs_missing_entry(Pid) end).
|
||||
|
||||
anti_entropy_merges_once_tombstone_expires_test() ->
|
||||
with_local_presence_cache(fun(Pid) -> assert_merge_resumes_after_expiry(Pid) end).
|
||||
|
||||
anti_entropy_tombstone_cleared_when_user_returns_test() ->
|
||||
with_local_presence_cache(fun(Pid) -> assert_tombstone_cleared_on_put(Pid) end).
|
||||
|
||||
assert_delete_survives_stale_peer_entry(Pid) ->
|
||||
UserId = 77001,
|
||||
Presence = anti_entropy_presence(UserId),
|
||||
State0 = cache_state(sys:get_state(Pid)),
|
||||
?assertEqual([node()], presence_cache_bulk:resolve_owner_nodes(UserId)),
|
||||
{_PutReply, State1} = presence_cache:put_local(UserId, Presence, State0),
|
||||
{_DeleteReply, State2} = presence_cache:delete_local(UserId, State1),
|
||||
?assertMatch({not_found, _}, presence_cache_ops:get_local(UserId, State2)),
|
||||
State3 = presence_cache_rebalance:merge_anti_entropy_entries(#{UserId => Presence}, State2),
|
||||
?assertMatch({not_found, _}, presence_cache_ops:get_local(UserId, State3)).
|
||||
|
||||
assert_merge_repairs_missing_entry(Pid) ->
|
||||
UserId = 77002,
|
||||
Presence = anti_entropy_presence(UserId),
|
||||
State0 = cache_state(sys:get_state(Pid)),
|
||||
State1 = presence_cache_rebalance:merge_anti_entropy_entries(#{UserId => Presence}, State0),
|
||||
?assertMatch({{ok, Presence}, _}, presence_cache_ops:get_local(UserId, State1)).
|
||||
|
||||
assert_merge_resumes_after_expiry(Pid) ->
|
||||
UserId = 77003,
|
||||
Presence = anti_entropy_presence(UserId),
|
||||
State0 = cache_state(sys:get_state(Pid)),
|
||||
{_PutReply, State1} = presence_cache:put_local(UserId, Presence, State0),
|
||||
{_DeleteReply, State2} = presence_cache:delete_local(UserId, State1),
|
||||
Expired = State2#{
|
||||
delete_tombstones => #{UserId => erlang:monotonic_time(millisecond) - 1}
|
||||
},
|
||||
State3 = presence_cache_rebalance:merge_anti_entropy_entries(
|
||||
#{UserId => Presence}, Expired
|
||||
),
|
||||
?assertMatch({{ok, Presence}, _}, presence_cache_ops:get_local(UserId, State3)),
|
||||
?assertEqual(#{}, maps:get(delete_tombstones, State3)).
|
||||
|
||||
assert_tombstone_cleared_on_put(Pid) ->
|
||||
UserId = 77004,
|
||||
Presence = anti_entropy_presence(UserId),
|
||||
State0 = cache_state(sys:get_state(Pid)),
|
||||
{_FirstPut, State1} = presence_cache:put_local(UserId, Presence, State0),
|
||||
{_DeleteReply, State2} = presence_cache:delete_local(UserId, State1),
|
||||
?assert(maps:is_key(UserId, maps:get(delete_tombstones, State2))),
|
||||
{_SecondPut, State3} = presence_cache:put_local(UserId, Presence, State2),
|
||||
?assertNot(maps:is_key(UserId, maps:get(delete_tombstones, State3))).
|
||||
|
||||
anti_entropy_presence(UserId) ->
|
||||
#{
|
||||
<<"status">> => <<"online">>,
|
||||
<<"user">> => #{<<"id">> => integer_to_binary(UserId)}
|
||||
}.
|
||||
|
||||
with_local_presence_cache(Fun) ->
|
||||
with_presence_member_nodes([node()], fun() -> with_presence_cache(Fun) end).
|
||||
|
||||
maybe_start_for_test() ->
|
||||
case whereis(presence_cache) of
|
||||
undefined -> presence_cache:start_link();
|
||||
@@ -253,12 +317,15 @@ with_presence_cache(Fun) ->
|
||||
end.
|
||||
|
||||
with_presence_members(RemoteNode, Fun) ->
|
||||
with_presence_member_nodes([node(), RemoteNode], Fun).
|
||||
|
||||
with_presence_member_nodes(Nodes, Fun) ->
|
||||
MembersKey = {gateway_cluster_membership, members},
|
||||
RoleMembersKey = {gateway_cluster_membership, members_by_role},
|
||||
OldMembers = persistent_term:get(MembersKey, undefined),
|
||||
OldRoleMembers = persistent_term:get(RoleMembersKey, undefined),
|
||||
persistent_term:put(MembersKey, [node(), RemoteNode]),
|
||||
persistent_term:put(RoleMembersKey, #{presence => [node(), RemoteNode]}),
|
||||
persistent_term:put(MembersKey, Nodes),
|
||||
persistent_term:put(RoleMembersKey, #{presence => Nodes}),
|
||||
try
|
||||
Fun()
|
||||
after
|
||||
|
||||
Reference in New Issue
Block a user