fix retry

This commit is contained in:
anlicheng 2026-04-22 17:23:17 +08:00
parent 3b1e4500d7
commit 143a88ba01
5 changed files with 89 additions and 153 deletions

View File

@ -13,15 +13,17 @@
%% %%
-define(RETRY_INTERVAL, 5000). -define(RETRY_INTERVAL, 5000).
%%
-define(MAX_RETRY_TIMES, 3).
-export([new/2, append/2, trigger_next/1, trigger_n/1, cleanup/1, recover_inflight/1, ack/2, stat/1, resize/2]). -export([new/2, append/2, trigger_next/1, trigger_n/1, handle_timeout/4, cleanup/1, recover_inflight/1, ack/2, stat/1, resize/2]).
-export_type([buffer/0]). -export_type([buffer/0]).
-record(buffer, { -record(buffer, {
endpoint :: #endpoint{}, endpoint :: #endpoint{},
outbox :: endpoint_outbox:outbox(), outbox :: endpoint_outbox:outbox(),
%% %% ack #{Id => {TimerRef, Payload, RetryTimes}}
timer_pid :: pid(), timer_map = #{} :: #{integer() => {reference(), binary(), non_neg_integer()}},
%% %%
window_size = 10, window_size = 10,
%% %%
@ -40,9 +42,7 @@ new(Endpoint = #endpoint{id = Id}, WindowSize) when is_integer(WindowSize), Wind
OutboxDir = filename:join(RootDir, integer_to_list(Id)), OutboxDir = filename:join(RootDir, integer_to_list(Id)),
{ok, Outbox} = endpoint_outbox:open(OutboxDir, #{}), {ok, Outbox} = endpoint_outbox:open(OutboxDir, #{}),
%% #buffer{outbox = Outbox, endpoint = Endpoint, window_size = WindowSize}.
{ok, TimerPid} = endpoint_timer:start_link(?RETRY_INTERVAL),
#buffer{outbox = Outbox, timer_pid = TimerPid, endpoint = Endpoint, window_size = WindowSize}.
-spec append(Payload :: binary(), Buffer :: #buffer{}) -> NBuffer :: #buffer{}. -spec append(Payload :: binary(), Buffer :: #buffer{}) -> NBuffer :: #buffer{}.
append(Payload, Buffer = #buffer{outbox = Outbox, window_size = WindowSize, flight_num = FlightNum}) when is_binary(Payload) -> append(Payload, Buffer = #buffer{outbox = Outbox, window_size = WindowSize, flight_num = FlightNum}) when is_binary(Payload) ->
@ -71,7 +71,7 @@ trigger_n(Buffer = #buffer{window_size = WindowSize}) ->
%% %%
-spec trigger_next(Buffer :: #buffer{}) -> NBuffer :: #buffer{}. -spec trigger_next(Buffer :: #buffer{}) -> NBuffer :: #buffer{}.
trigger_next(Buffer = #buffer{outbox = Outbox, timer_pid = TimerPid, flight_num = FlightNum, window_size = WindowSize}) -> trigger_next(Buffer = #buffer{outbox = Outbox, flight_num = FlightNum, window_size = WindowSize}) ->
case FlightNum < WindowSize of case FlightNum < WindowSize of
false -> false ->
Buffer; Buffer;
@ -82,8 +82,7 @@ trigger_next(Buffer = #buffer{outbox = Outbox, timer_pid = TimerPid, flight_num
{ok, Id, Payload, NOutbox} -> {ok, Id, Payload, NOutbox} ->
ReceiverPid = self(), ReceiverPid = self(),
ReceiverPid ! {next_data, Id, Payload}, ReceiverPid ! {next_data, Id, Payload},
endpoint_timer:task(TimerPid, Id, fun() -> ReceiverPid ! {next_data, Id, Payload} end), schedule_retry(Id, Payload, 0, Buffer#buffer{outbox = NOutbox, flight_num = FlightNum + 1});
Buffer#buffer{outbox = NOutbox, flight_num = FlightNum + 1};
{error, Reason} -> {error, Reason} ->
logger:warning("[endpoint_buffer] read next outbox failed, endpoint_id: ~p, reason: ~p", logger:warning("[endpoint_buffer] read next outbox failed, endpoint_id: ~p, reason: ~p",
[buffer_endpoint_id(Buffer), Reason]), [buffer_endpoint_id(Buffer), Reason]),
@ -91,13 +90,35 @@ trigger_next(Buffer = #buffer{outbox = Outbox, timer_pid = TimerPid, flight_num
end end
end. end.
-spec handle_timeout(reference(), integer(), binary(), #buffer{}) -> #buffer{}.
handle_timeout(TimerRef, Id, _Payload, Buffer = #buffer{timer_map = TimerMap})
when is_reference(TimerRef), is_integer(Id) ->
case maps:take(Id, TimerMap) of
{{TimerRef, Payload, RetryTimes}, NTimerMap} when RetryTimes < ?MAX_RETRY_TIMES ->
logger:warning("[endpoint_buffer] retry message, endpoint_id: ~p, id: ~p, retry: ~p/~p",
[buffer_endpoint_id(Buffer), Id, RetryTimes + 1, ?MAX_RETRY_TIMES]),
ReceiverPid = self(),
ReceiverPid ! {next_data, Id, Payload},
schedule_retry(Id, Payload, RetryTimes + 1, Buffer#buffer{timer_map = NTimerMap});
{{TimerRef, Payload, RetryTimes}, NTimerMap} ->
logger:warning("[endpoint_buffer] drop message after retries exhausted, endpoint_id: ~p, id: ~p, retries: ~p",
[buffer_endpoint_id(Buffer), Id, RetryTimes]),
drop_message(Id, Payload, RetryTimes, Buffer#buffer{timer_map = NTimerMap});
{{_OtherTimerRef, _OtherPayload, _RetryTimes}, _NTimerMap} ->
Buffer;
error ->
Buffer
end.
-spec ack(Id :: integer(), Buffer :: #buffer{}) -> NBuffer :: #buffer{}. -spec ack(Id :: integer(), Buffer :: #buffer{}) -> NBuffer :: #buffer{}.
ack(Id, Buffer = #buffer{timer_pid = TimerPid, outbox = Outbox, acc_num = AccNum, flight_num = FlightNum}) when is_integer(Id) -> ack(Id, Buffer = #buffer{timer_map = TimerMap, outbox = Outbox, acc_num = AccNum, flight_num = FlightNum}) when is_integer(Id) ->
case endpoint_timer:ack(TimerPid, Id) of case maps:take(Id, TimerMap) of
true -> {{TimerRef, Payload, RetryTimes}, NTimerMap} ->
_ = erlang:cancel_timer(TimerRef),
case endpoint_outbox:ack(Id, Outbox) of case endpoint_outbox:ack(Id, Outbox) of
{ok, NOutbox} -> {ok, NOutbox} ->
NBuffer = Buffer#buffer{ NBuffer = Buffer#buffer{
timer_map = NTimerMap,
outbox = NOutbox, outbox = NOutbox,
acc_num = AccNum + 1, acc_num = AccNum + 1,
flight_num = max(FlightNum - 1, 0) flight_num = max(FlightNum - 1, 0)
@ -106,9 +127,9 @@ ack(Id, Buffer = #buffer{timer_pid = TimerPid, outbox = Outbox, acc_num = AccNum
{error, Reason} -> {error, Reason} ->
logger:warning("[endpoint_buffer] ack outbox failed, endpoint_id: ~p, id: ~p, reason: ~p", logger:warning("[endpoint_buffer] ack outbox failed, endpoint_id: ~p, id: ~p, reason: ~p",
[buffer_endpoint_id(Buffer), Id, Reason]), [buffer_endpoint_id(Buffer), Id, Reason]),
Buffer schedule_retry(Id, Payload, RetryTimes, Buffer#buffer{timer_map = NTimerMap})
end; end;
false -> error ->
Buffer Buffer
end. end.
@ -126,20 +147,20 @@ stat(#buffer{acc_num = AccNum, outbox = Outbox, flight_num = FlightNum}) ->
}. }.
-spec cleanup(Buffer :: #buffer{}) -> #buffer{}. -spec cleanup(Buffer :: #buffer{}) -> #buffer{}.
cleanup(Buffer = #buffer{timer_pid = TimerPid}) -> cleanup(Buffer = #buffer{timer_map = TimerMap}) ->
endpoint_timer:cleanup(TimerPid), cancel_timers(TimerMap),
Buffer. Buffer#buffer{timer_map = #{}}.
-spec recover_inflight(Buffer :: #buffer{}) -> #buffer{}. -spec recover_inflight(Buffer :: #buffer{}) -> #buffer{}.
recover_inflight(Buffer = #buffer{outbox = Outbox, timer_pid = TimerPid}) -> recover_inflight(Buffer = #buffer{outbox = Outbox, timer_map = TimerMap}) ->
endpoint_timer:cleanup(TimerPid), cancel_timers(TimerMap),
case endpoint_outbox:reset_reader(Outbox) of case endpoint_outbox:reset_reader(Outbox) of
{ok, NOutbox} -> {ok, NOutbox} ->
Buffer#buffer{outbox = NOutbox, flight_num = 0}; Buffer#buffer{outbox = NOutbox, timer_map = #{}, flight_num = 0};
{error, Reason} -> {error, Reason} ->
logger:warning("[endpoint_buffer] recover inflight failed, endpoint_id: ~p, reason: ~p", logger:warning("[endpoint_buffer] recover inflight failed, endpoint_id: ~p, reason: ~p",
[buffer_endpoint_id(Buffer), Reason]), [buffer_endpoint_id(Buffer), Reason]),
Buffer#buffer{flight_num = 0} Buffer#buffer{timer_map = #{}, flight_num = 0}
end. end.
-spec resize(Buffer :: #buffer{}, WindowSize :: integer()) -> #buffer{}. -spec resize(Buffer :: #buffer{}, WindowSize :: integer()) -> #buffer{}.
@ -153,3 +174,32 @@ resize(Buffer = #buffer{}, WindowSize) when is_integer(WindowSize), WindowSize >
-spec buffer_endpoint_id(buffer()) -> integer(). -spec buffer_endpoint_id(buffer()) -> integer().
buffer_endpoint_id(#buffer{endpoint = #endpoint{id = Id}}) -> buffer_endpoint_id(#buffer{endpoint = #endpoint{id = Id}}) ->
Id. Id.
-spec schedule_retry(integer(), binary(), non_neg_integer(), buffer()) -> buffer().
schedule_retry(Id, Payload, RetryTimes, Buffer = #buffer{timer_map = TimerMap})
when is_integer(Id), is_binary(Payload), is_integer(RetryTimes), RetryTimes >= 0 ->
TimerRef = erlang:start_timer(?RETRY_INTERVAL, self(), {endpoint_buffer_retry, Id, Payload}),
Buffer#buffer{timer_map = maps:put(Id, {TimerRef, Payload, RetryTimes}, TimerMap)}.
-spec cancel_timers(#{integer() => {reference(), binary(), non_neg_integer()}}) -> ok.
cancel_timers(TimerMap) ->
lists:foreach(fun({_Id, {TimerRef, _Payload, _RetryTimes}}) ->
_ = erlang:cancel_timer(TimerRef),
ok
end, maps:to_list(TimerMap)),
ok.
-spec drop_message(integer(), binary(), non_neg_integer(), buffer()) -> buffer().
drop_message(Id, Payload, RetryTimes, Buffer = #buffer{outbox = Outbox, acc_num = AccNum, flight_num = FlightNum}) ->
case endpoint_outbox:ack(Id, Outbox) of
{ok, NOutbox} ->
trigger_next(Buffer#buffer{
outbox = NOutbox,
acc_num = AccNum + 1,
flight_num = max(FlightNum - 1, 0)
});
{error, Reason} ->
logger:warning("[endpoint_buffer] ack dropped message failed, endpoint_id: ~p, id: ~p, reason: ~p",
[buffer_endpoint_id(Buffer), Id, Reason]),
schedule_retry(Id, Payload, RetryTimes, Buffer)
end.

View File

@ -70,7 +70,7 @@ handle_call(get_stat, _From, State = #state{buffer = Buffer}) ->
{noreply, NewState :: #state{}} | {noreply, NewState :: #state{}} |
{noreply, NewState :: #state{}, timeout() | hibernate} | {noreply, NewState :: #state{}, timeout() | hibernate} |
{stop, Reason :: term(), NewState :: #state{}}). {stop, Reason :: term(), NewState :: #state{}}).
handle_cast({forward, Metric}, State = #state{buffer = Buffer, endpoint = #endpoint{config = #http_endpoint{token = Token}}}) -> handle_cast({forward, Metric}, State = #state{buffer = Buffer}) ->
NBuffer = endpoint_buffer:append(Metric, Buffer), NBuffer = endpoint_buffer:append(Metric, Buffer),
{noreply, State#state{buffer = NBuffer}}; {noreply, State#state{buffer = NBuffer}};
handle_cast({reload, NEndpoint = #endpoint{matcher = NMatcher, config = #http_endpoint{pool_size = PoolSize}}}, State = #state{endpoint = #endpoint{matcher = Matcher}, buffer = Buffer}) -> handle_cast({reload, NEndpoint = #endpoint{matcher = NMatcher, config = #http_endpoint{pool_size = PoolSize}}}, State = #state{endpoint = #endpoint{matcher = Matcher}, buffer = Buffer}) ->
@ -103,7 +103,10 @@ handle_info({next_data, Id, Metric}, State = #state{buffer = Buffer, endpoint =
{error, Reason} -> {error, Reason} ->
logger:warning("[endpoint_http] url: ~p, get error: ~p", [Url, Reason]), logger:warning("[endpoint_http] url: ~p, get error: ~p", [Url, Reason]),
{noreply, State} {noreply, State}
end. end;
handle_info({timeout, TimerRef, {endpoint_buffer_retry, Id, Payload}}, State = #state{buffer = Buffer}) ->
NBuffer = endpoint_buffer:handle_timeout(TimerRef, Id, Payload, Buffer),
{noreply, State#state{buffer = NBuffer}}.
%% @private %% @private
%% @doc This function is called by a gen_server when it is about to %% @doc This function is called by a gen_server when it is about to

View File

@ -70,6 +70,9 @@ disconnected(cast, {reload, NEndpoint = #endpoint{matcher = NMatcher}},
reload_endpoint(Matcher, NMatcher, ClientId, NEndpoint, State); reload_endpoint(Matcher, NMatcher, ClientId, NEndpoint, State);
disconnected(state_timeout, connect, State) -> disconnected(state_timeout, connect, State) ->
try_connect(State); try_connect(State);
disconnected(info, {timeout, TimerRef, {endpoint_buffer_retry, Id, Payload}}, State = #state{buffer = Buffer}) ->
NBuffer = endpoint_buffer:handle_timeout(TimerRef, Id, Payload, Buffer),
{keep_state, State#state{buffer = NBuffer}};
disconnected(info, {next_data, _Id, _Tuple}, State) -> disconnected(info, {next_data, _Id, _Tuple}, State) ->
{keep_state, State}; {keep_state, State};
disconnected(info, {ack, Id}, State = #state{buffer = Buffer}) -> disconnected(info, {ack, Id}, State = #state{buffer = Buffer}) ->
@ -120,6 +123,9 @@ connected(info, {next_data, Id, Metric},
{next_state, disconnected, State#state{client_pid = undefined, buffer = NBuffer}, {next_state, disconnected, State#state{client_pid = undefined, buffer = NBuffer},
[{state_timeout, ?RETRY_INTERVAL, connect}]} [{state_timeout, ?RETRY_INTERVAL, connect}]}
end; end;
connected(info, {timeout, TimerRef, {endpoint_buffer_retry, Id, Payload}}, State = #state{buffer = Buffer}) ->
NBuffer = endpoint_buffer:handle_timeout(TimerRef, Id, Payload, Buffer),
{keep_state, State#state{buffer = NBuffer}};
connected(info, {ack, Id}, State = #state{buffer = Buffer}) -> connected(info, {ack, Id}, State = #state{buffer = Buffer}) ->
ack_buffer(Id, Buffer, State); ack_buffer(Id, Buffer, State);
connected(info, {'EXIT', ClientPid, Reason}, connected(info, {'EXIT', ClientPid, Reason},

View File

@ -110,6 +110,9 @@ disconnected(internal, do_connect,
{keep_state, State#state{conn_pid = undefined, inflight = #{}}, {keep_state, State#state{conn_pid = undefined, inflight = #{}},
[{state_timeout, ?RETRY_INTERVAL, connect}]} [{state_timeout, ?RETRY_INTERVAL, connect}]}
end; end;
disconnected(info, {timeout, TimerRef, {endpoint_buffer_retry, Id, Payload}}, State = #state{buffer = Buffer}) ->
NBuffer = endpoint_buffer:handle_timeout(TimerRef, Id, Payload, Buffer),
{keep_state, State#state{buffer = NBuffer}};
disconnected(info, {next_data, _Id, _Tuple}, State) -> disconnected(info, {next_data, _Id, _Tuple}, State) ->
{keep_state, State}; {keep_state, State};
disconnected(info, {'EXIT', ConnPid, Reason}, disconnected(info, {'EXIT', ConnPid, Reason},
@ -156,6 +159,9 @@ connected(info, {next_data, Id, Metric},
{next_state, disconnected, State#state{conn_pid = undefined, inflight = #{}, buffer = NBuffer}, {next_state, disconnected, State#state{conn_pid = undefined, inflight = #{}, buffer = NBuffer},
[{state_timeout, ?RETRY_INTERVAL, connect}]} [{state_timeout, ?RETRY_INTERVAL, connect}]}
end; end;
connected(info, {timeout, TimerRef, {endpoint_buffer_retry, Id, Payload}}, State = #state{buffer = Buffer}) ->
NBuffer = endpoint_buffer:handle_timeout(TimerRef, Id, Payload, Buffer),
{keep_state, State#state{buffer = NBuffer}};
connected(info, {disconnected, ReasonCode, Properties}, connected(info, {disconnected, ReasonCode, Properties},
State = #state{conn_pid = ConnPid, buffer = Buffer}) -> State = #state{conn_pid = ConnPid, buffer = Buffer}) ->
logger:debug("[endpoint_mqtt] Recv a DISONNECT packet - ReasonCode: ~p, Properties: ~p", [ReasonCode, Properties]), logger:debug("[endpoint_mqtt] Recv a DISONNECT packet - ReasonCode: ~p, Properties: ~p", [ReasonCode, Properties]),

View File

@ -1,129 +0,0 @@
%%%-------------------------------------------------------------------
%%% @author anlicheng
%%% @copyright (C) 2024, <COMPANY>
%%% @doc
%%%
%%% @end
%%% Created : 07. 5 2024 10:30
%%%-------------------------------------------------------------------
-module(endpoint_timer).
-author("anlicheng").
-include("endpoint.hrl").
-behaviour(gen_server).
%% API
-export([start_link/1]).
-export([task/3, ack/2, cleanup/1]).
%% gen_server callbacks
-export([init/1, handle_call/3, handle_cast/2, handle_info/2, terminate/2, code_change/3]).
-record(state, {
retry_interval = 0,
%%
timer_map = #{}
}).
%%%===================================================================
%%% API
%%%===================================================================
task(Pid, Id, Task) when is_pid(Pid), is_integer(Id), is_function(Task, 0) ->
gen_server:cast(Pid, {task, Id, Task}).
-spec ack(pid(), integer()) -> boolean().
ack(Pid, Id) when is_pid(Pid), is_integer(Id) ->
gen_server:call(Pid, {ack, Id}).
cleanup(Pid) when is_pid(Pid) ->
gen_server:cast(Pid, cleanup).
%% @doc Spawns the server and registers the local name (unique)
-spec(start_link(RetryInterval :: integer()) ->
{ok, Pid :: pid()} | ignore | {error, Reason :: term()}).
start_link(RetryInterval) when is_integer(RetryInterval) ->
gen_server:start_link(?MODULE, [RetryInterval], []).
%%%===================================================================
%%% gen_server callbacks
%%%===================================================================
%% @private
%% @doc Initializes the server
-spec(init(Args :: term()) ->
{ok, State :: #state{}} | {ok, State :: #state{}, timeout() | hibernate} |
{stop, Reason :: term()} | ignore).
init([RetryInterval]) ->
ok = iot_log:set_metadata(),
{ok, #state{retry_interval = RetryInterval}}.
%% @private
%% @doc Handling call messages
-spec(handle_call(Request :: term(), From :: {pid(), Tag :: term()},
State :: #state{}) ->
{reply, Reply :: term(), NewState :: #state{}} |
{reply, Reply :: term(), NewState :: #state{}, timeout() | hibernate} |
{noreply, NewState :: #state{}} |
{noreply, NewState :: #state{}, timeout() | hibernate} |
{stop, Reason :: term(), Reply :: term(), NewState :: #state{}} |
{stop, Reason :: term(), NewState :: #state{}}).
handle_call({ack, Id}, _From, State = #state{timer_map = TimerMap}) ->
case maps:take(Id, TimerMap) of
error ->
{reply, false, State};
{TimerRef, NTimerMap} ->
is_reference(TimerRef) andalso erlang:cancel_timer(TimerRef),
{reply, true, State#state{timer_map = NTimerMap}}
end;
handle_call(_Request, _From, State = #state{}) ->
{reply, ok, State}.
%% @private
%% @doc Handling cast messages
-spec(handle_cast(Request :: term(), State :: #state{}) ->
{noreply, NewState :: #state{}} |
{noreply, NewState :: #state{}, timeout() | hibernate} |
{stop, Reason :: term(), NewState :: #state{}}).
handle_cast({task, Id, Task}, State = #state{retry_interval = RetryInterval, timer_map = TimerMap}) ->
TimerRef = erlang:start_timer(RetryInterval, self(), {repost_ticker, {Id, Task}}),
{noreply, State#state{timer_map = maps:put(Id, TimerRef, TimerMap)}};
%%
handle_cast(cleanup, State = #state{timer_map = TimerMap}) ->
lists:foreach(fun({_, TimerRef}) -> catch erlang:cancel_timer(TimerRef) end, maps:to_list(TimerMap)),
{noreply, State#state{timer_map = #{}}}.
%% @private
%% @doc Handling all non call/cast messages
-spec(handle_info(Info :: timeout() | term(), State :: #state{}) ->
{noreply, NewState :: #state{}} |
{noreply, NewState :: #state{}, timeout() | hibernate} |
{stop, Reason :: term(), NewState :: #state{}}).
handle_info({timeout, _, {repost_ticker, {Id, Task}}}, State = #state{retry_interval = RetryInterval, timer_map = TimerMap}) ->
Task(),
TimerRef = erlang:start_timer(RetryInterval, self(), {repost_ticker, {Id, Task}}),
{noreply, State#state{timer_map = maps:put(Id, TimerRef, TimerMap)}}.
%% @private
%% @doc This function is called by a gen_server when it is about to
%% terminate. It should be the opposite of Module:init/1 and do any
%% necessary cleaning up. When it returns, the gen_server terminates
%% with Reason. The return value is ignored.
-spec(terminate(Reason :: (normal | shutdown | {shutdown, term()} | term()),
State :: #state{}) -> term()).
terminate(_Reason, _State = #state{}) ->
ok.
%% @private
%% @doc Convert process state when code is changed
-spec(code_change(OldVsn :: term() | {down, term()}, State :: #state{},
Extra :: term()) ->
{ok, NewState :: #state{}} | {error, Reason :: term()}).
code_change(_OldVsn, State = #state{}, _Extra) ->
{ok, State}.
%%%===================================================================
%%% Internal functions
%%%===================================================================