Skip to content

Commit efccc2d

Browse files
authored
Merge pull request #195 from emqx/dev/fix-autoclean
Make autoclean configuration dynamic
2 parents 4354af0 + f3cc2c1 commit efccc2d

12 files changed

Lines changed: 200 additions & 162 deletions

.github/workflows/run_test_case.yaml

Lines changed: 68 additions & 53 deletions
Original file line numberDiff line numberDiff line change
@@ -3,57 +3,72 @@ name: Run test case
33
on: [push, pull_request]
44

55
jobs:
6-
76
run_test_case:
8-
runs-on: ubuntu-latest
9-
10-
container: ghcr.io/emqx/emqx-builder/5.3-5:1.15.7-26.2.1-2-ubuntu24.04
11-
12-
steps:
13-
- uses: actions/checkout@9bb56186c3b09b4f86b1c65136769dd318469633 # v4.1.2
14-
15-
- name: Install prerequisites
16-
run: |
17-
apt update
18-
apt install -y cmake
19-
20-
- name: Configure git
21-
run: |
22-
git config --global --add safe.directory "*"
23-
24-
- name: Compile
25-
run: |
26-
make
27-
28-
- name: Concuerror tests
29-
run : |
30-
make concuerror_test
31-
32-
- name: Smoke test
33-
run: |
34-
make smoke-test
35-
36-
- name: Fault-tolerance tests
37-
run: |
38-
make ct-fault-tolerance
39-
40-
- name: Consistency tests
41-
run: |
42-
make ct-consistency
43-
44-
- name: Coveralls
45-
env:
46-
GITHUB_TOKEN: ${{ secrets.GITHUB_TOKEN }}
47-
run: |
48-
make coveralls
49-
50-
- uses: actions/upload-artifact@5d5d22a31266ced268874388b861e4b58bb5c2f3 # v4.3.1
51-
if: always()
52-
with:
53-
name: logs
54-
path: _build/test/logs
55-
56-
- uses: actions/upload-artifact@5d5d22a31266ced268874388b861e4b58bb5c2f3 # v4.3.1
57-
with:
58-
name: cover
59-
path: _build/test/cover
7+
runs-on: ubuntu-latest
8+
strategy:
9+
fail-fast: false
10+
matrix:
11+
os:
12+
- ubuntu24.04
13+
arch:
14+
- amd64
15+
otp:
16+
- "28.2-2"
17+
builder:
18+
- '6.0-9'
19+
elixir:
20+
- '1.19.1'
21+
22+
container: "ghcr.io/emqx/emqx-builder/${{ matrix.builder }}:${{ matrix.elixir }}-${{ matrix.otp }}-${{ matrix.os }}"
23+
24+
steps:
25+
- uses: actions/checkout@8e8c483db84b4bee98b60c0593521ed34d9990e8 # v6.0.1
26+
with:
27+
# Needed for backward-compatibility test:
28+
fetch-depth: 0
29+
30+
- name: Install prerequisites
31+
run: |
32+
apt update
33+
apt install -y cmake
34+
35+
- name: Configure git
36+
run: |
37+
git config --global --add safe.directory "*"
38+
39+
- name: Compile
40+
run: |
41+
make
42+
43+
- name: Concuerror tests
44+
run : |
45+
make concuerror_test
46+
47+
- name: Smoke test
48+
run: |
49+
make smoke-test
50+
51+
- name: Fault-tolerance tests
52+
run: |
53+
make ct-fault-tolerance
54+
55+
- name: Consistency tests
56+
run: |
57+
make ct-consistency
58+
59+
- name: Coveralls
60+
env:
61+
GITHUB_TOKEN: ${{ secrets.GITHUB_TOKEN }}
62+
run: |
63+
make coveralls
64+
65+
- uses: actions/upload-artifact@330a01c490aca151604b8cf639adc76d48f6c5d4 # v5.0.0
66+
if: always()
67+
with:
68+
name: logs
69+
path: _build/test/logs
70+
71+
- uses: actions/upload-artifact@330a01c490aca151604b8cf639adc76d48f6c5d4 # v5.0.0
72+
with:
73+
name: cover
74+
path: _build/test/cover

Makefile

Lines changed: 2 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -31,15 +31,14 @@ test: smoke-test ct-consistency ct-fault-tolerance cover
3131
.PHONY: smoke-test
3232
smoke-test:
3333
$(REBAR) do eunit, ct -v --cover --readable=$(CT_READABLE)
34-
$(REBAR) ct --readable=$(CT_READABLE) -v --suite mria_compatibility_suite
3534

3635
.PHONY: ct-consistency
3736
ct-consistency:
38-
$(REBAR) ct --cover -v --readable=$(CT_READABLE) --suite mria_proper_suite,mria_proper_mixed_cluster_suite
37+
$(REBAR) ct --cover -v --readable=$(CT_READABLE) --suite mria_proper_suite_,mria_proper_mixed_cluster_suite_
3938

4039
.PHONY: ct-fault-tolerance
4140
ct-fault-tolerance:
42-
$(REBAR) ct --cover -v --readable=$(CT_READABLE) --suite mria_fault_tolerance_suite
41+
$(REBAR) ct --cover -v --readable=$(CT_READABLE) --suite mria_fault_tolerance_suite_
4342

4443
.PHONY: ct-suite
4544
ct-suite: compile

include/mria.hrl

Lines changed: 9 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -5,14 +5,15 @@
55
-type(member_address() :: {inet:ip_address(), inet:port_number()}).
66

77
-record(member, {
8-
node :: node(),
9-
addr :: undefined | member_address(),
10-
guid :: undefined | mria_guid:guid(),
11-
hash :: undefined | pos_integer(),
12-
status :: member_status(),
13-
mnesia :: undefined | running | stopped | false,
14-
ltime :: undefined | erlang:timestamp(),
15-
role :: mria_rlog:role()
8+
node :: node(),
9+
addr :: undefined | member_address(),
10+
guid :: undefined | mria_guid:guid(),
11+
hash :: undefined | pos_integer(),
12+
status :: member_status(),
13+
mnesia :: undefined | running | stopped | false,
14+
%% Timestamp of the last membership update (up, down, join, leave, ...) in seconds:
15+
last_update :: undefined | integer(),
16+
role :: mria_rlog:role()
1617
}).
1718

1819
-type(member() :: #member{}).

src/mria_autoclean.erl

Lines changed: 0 additions & 50 deletions
This file was deleted.

src/mria_membership.erl

Lines changed: 9 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,5 @@
11
%%--------------------------------------------------------------------
2-
%% Copyright (c) 2019-2023 EMQ Technologies Co., Ltd. All Rights Reserved.
2+
%% Copyright (c) 2019-2026 EMQ Technologies Co., Ltd. All Rights Reserved.
33
%%
44
%% Licensed under the Apache License, Version 2.0 (the "License");
55
%% you may not use this file except in compliance with the License.
@@ -37,6 +37,7 @@
3737
, is_member/1
3838
, oldest/1
3939
, replicants/0
40+
, now_seconds/0
4041
]).
4142

4243
-export([ leader/0
@@ -158,6 +159,9 @@ members(Status) ->
158159
replicants() ->
159160
try select(replicant) catch error:badarg -> [] end.
160161

162+
now_seconds() ->
163+
erlang:monotonic_time(second).
164+
161165
%% Legacy API, selects a leader only from core members for compatibility reasons
162166
%% Get leader node of the members
163167
-spec(leader() -> node()).
@@ -482,9 +486,10 @@ make_new_local_member() ->
482486
true -> running;
483487
false -> stopped
484488
end,
485-
with_hash(#member{node = node(), guid = mria_guid:gen(),
489+
with_hash(#member{node = node(),
490+
guid = mria_guid:gen(),
486491
status = up, mnesia = IsMnesiaRunning,
487-
ltime = erlang:timestamp(),
492+
last_update = now_seconds(),
488493
role = mria_config:role()
489494
}).
490495

@@ -510,7 +515,7 @@ lookup(Node) ->
510515
ets:lookup(?TAB, Node).
511516

512517
insert(Member0) ->
513-
Member = Member0#member{ltime = erlang:timestamp()},
518+
Member = Member0#member{last_update = now_seconds()},
514519
?tp(mria_membership_insert, #{member => Member}),
515520
ets:insert(?TAB, Member).
516521

src/mria_node_monitor.erl

Lines changed: 31 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,5 @@
11
%%--------------------------------------------------------------------
2-
%% Copyright (c) 2019-2023 EMQ Technologies Co., Ltd. All Rights Reserved.
2+
%% Copyright (c) 2019-2026 EMQ Technologies Co., Ltd. All Rights Reserved.
33
%%
44
%% Licensed under the Apache License, Version 2.0 (the "License");
55
%% you may not use this file except in compliance with the License.
@@ -42,8 +42,7 @@
4242
-record(state, {
4343
partitions :: list(node()),
4444
heartbeat :: undefined | reference(),
45-
autoheal :: mria_autoheal:autoheal(),
46-
autoclean :: mria_autoclean:autoclean()
45+
autoheal :: mria_autoheal:autoheal()
4746
}).
4847

4948
-define(SERVER, ?MODULE).
@@ -79,8 +78,7 @@ init([]) ->
7978
{ok, _} = mnesia:subscribe(system),
8079
lists:foreach(fun(N) -> self() ! {nodeup, N, []} end, nodes() -- [node()]),
8180
State = #state{partitions = [],
82-
autoheal = mria_autoheal:init(),
83-
autoclean = mria_autoclean:init()
81+
autoheal = mria_autoheal:init()
8482
},
8583
{ok, ensure_heartbeat(State)}.
8684

@@ -199,6 +197,7 @@ handle_info(heartbeat, State) ->
199197
true -> ok
200198
end
201199
end, mria_mnesia:cluster_nodes(all)),
200+
autoclean_cores(),
202201
{noreply, ensure_heartbeat(State#state{heartbeat = undefined})};
203202

204203
handle_info(Msg = {'EXIT', Pid, _Reason}, State = #state{autoheal = Autoheal}) ->
@@ -207,10 +206,6 @@ handle_info(Msg = {'EXIT', Pid, _Reason}, State = #state{autoheal = Autoheal}) -
207206
_ -> {noreply, State}
208207
end;
209208

210-
%% Autoclean Event.
211-
handle_info(autoclean, State = #state{autoclean = AutoClean}) ->
212-
{noreply, State#state{autoclean = mria_autoclean:check(AutoClean)}};
213-
214209
handle_info(Info, State) ->
215210
logger:error("Unexpected info: ~p", [Info]),
216211
{noreply, State}.
@@ -237,3 +232,30 @@ ensure_heartbeat(State) ->
237232

238233
autoheal_handle_msg(Msg, State = #state{autoheal = Autoheal}) ->
239234
State#state{autoheal = mria_autoheal:handle_msg(Msg, Autoheal)}.
235+
236+
%%--------------------------------------------------------------------
237+
%% Autoclean of stopped core nodes
238+
%%--------------------------------------------------------------------
239+
240+
autoclean_cores() ->
241+
maybe
242+
{ok, Expiry} ?= application:get_env(mria, cluster_autoclean),
243+
[maybe_clean(Member, Expiry) || Member <- mria_membership:members(down)]
244+
end,
245+
ok.
246+
247+
maybe_clean(#member{node = Node, last_update = WentDownAt} = Member, MaxDownSecs) ->
248+
Now = mria_membership:now_seconds(),
249+
case Now - WentDownAt > MaxDownSecs of
250+
true ->
251+
?tp(notice, mria_autoclean_force_leave,
252+
#{ node => Node
253+
, limit => MaxDownSecs
254+
, status => Member
255+
, last_update => WentDownAt
256+
, now => Now
257+
}),
258+
mria:force_leave(Node);
259+
false ->
260+
ok
261+
end.

0 commit comments

Comments
 (0)