From 143a88ba0161044868973b1d588f35d5265dc790 Mon Sep 17 00:00:00 2001 From: anlicheng <244108715@qq.com> Date: Wed, 22 Apr 2026 17:23:17 +0800 Subject: [PATCH] fix retry --- src/endpoint/endpoint_buffer.erl | 92 +++++++++++++++++----- src/endpoint/endpoint_http.erl | 9 ++- src/endpoint/endpoint_kafka.erl | 6 ++ src/endpoint/endpoint_mqtt.erl | 6 ++ src/endpoint/endpoint_timer.erl | 129 ------------------------------- 5 files changed, 89 insertions(+), 153 deletions(-) delete mode 100644 src/endpoint/endpoint_timer.erl diff --git a/src/endpoint/endpoint_buffer.erl b/src/endpoint/endpoint_buffer.erl index 9a29485..93e902c 100644 --- a/src/endpoint/endpoint_buffer.erl +++ b/src/endpoint/endpoint_buffer.erl @@ -13,15 +13,17 @@ %% 消息重发间隔 -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]). -record(buffer, { endpoint :: #endpoint{}, outbox :: endpoint_outbox:outbox(), - %% 定时器 - timer_pid :: pid(), + %% 当前待 ack 的数据及其重试定时器 #{Id => {TimerRef, Payload, RetryTimes}} + timer_map = #{} :: #{integer() => {reference(), binary(), non_neg_integer()}}, %% 窗口大小,允许最大的未确认消息数 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)), {ok, Outbox} = endpoint_outbox:open(OutboxDir, #{}), - %% 定义重发器 - {ok, TimerPid} = endpoint_timer:start_link(?RETRY_INTERVAL), - #buffer{outbox = Outbox, timer_pid = TimerPid, endpoint = Endpoint, window_size = WindowSize}. + #buffer{outbox = Outbox, endpoint = Endpoint, window_size = WindowSize}. -spec append(Payload :: binary(), Buffer :: #buffer{}) -> NBuffer :: #buffer{}. 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{}. -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 false -> Buffer; @@ -82,8 +82,7 @@ trigger_next(Buffer = #buffer{outbox = Outbox, timer_pid = TimerPid, flight_num {ok, Id, Payload, NOutbox} -> ReceiverPid = self(), ReceiverPid ! {next_data, Id, Payload}, - endpoint_timer:task(TimerPid, Id, fun() -> ReceiverPid ! {next_data, Id, Payload} end), - Buffer#buffer{outbox = NOutbox, flight_num = FlightNum + 1}; + schedule_retry(Id, Payload, 0, Buffer#buffer{outbox = NOutbox, flight_num = FlightNum + 1}); {error, Reason} -> logger:warning("[endpoint_buffer] read next outbox failed, endpoint_id: ~p, reason: ~p", [buffer_endpoint_id(Buffer), Reason]), @@ -91,13 +90,35 @@ trigger_next(Buffer = #buffer{outbox = Outbox, timer_pid = TimerPid, flight_num 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{}. -ack(Id, Buffer = #buffer{timer_pid = TimerPid, outbox = Outbox, acc_num = AccNum, flight_num = FlightNum}) when is_integer(Id) -> - case endpoint_timer:ack(TimerPid, Id) of - true -> +ack(Id, Buffer = #buffer{timer_map = TimerMap, outbox = Outbox, acc_num = AccNum, flight_num = FlightNum}) when is_integer(Id) -> + case maps:take(Id, TimerMap) of + {{TimerRef, Payload, RetryTimes}, NTimerMap} -> + _ = erlang:cancel_timer(TimerRef), case endpoint_outbox:ack(Id, Outbox) of {ok, NOutbox} -> NBuffer = Buffer#buffer{ + timer_map = NTimerMap, outbox = NOutbox, acc_num = AccNum + 1, flight_num = max(FlightNum - 1, 0) @@ -106,9 +127,9 @@ ack(Id, Buffer = #buffer{timer_pid = TimerPid, outbox = Outbox, acc_num = AccNum {error, Reason} -> logger:warning("[endpoint_buffer] ack outbox failed, endpoint_id: ~p, id: ~p, reason: ~p", [buffer_endpoint_id(Buffer), Id, Reason]), - Buffer + schedule_retry(Id, Payload, RetryTimes, Buffer#buffer{timer_map = NTimerMap}) end; - false -> + error -> Buffer end. @@ -126,20 +147,20 @@ stat(#buffer{acc_num = AccNum, outbox = Outbox, flight_num = FlightNum}) -> }. -spec cleanup(Buffer :: #buffer{}) -> #buffer{}. -cleanup(Buffer = #buffer{timer_pid = TimerPid}) -> - endpoint_timer:cleanup(TimerPid), - Buffer. +cleanup(Buffer = #buffer{timer_map = TimerMap}) -> + cancel_timers(TimerMap), + Buffer#buffer{timer_map = #{}}. -spec recover_inflight(Buffer :: #buffer{}) -> #buffer{}. -recover_inflight(Buffer = #buffer{outbox = Outbox, timer_pid = TimerPid}) -> - endpoint_timer:cleanup(TimerPid), +recover_inflight(Buffer = #buffer{outbox = Outbox, timer_map = TimerMap}) -> + cancel_timers(TimerMap), case endpoint_outbox:reset_reader(Outbox) of {ok, NOutbox} -> - Buffer#buffer{outbox = NOutbox, flight_num = 0}; + Buffer#buffer{outbox = NOutbox, timer_map = #{}, flight_num = 0}; {error, Reason} -> logger:warning("[endpoint_buffer] recover inflight failed, endpoint_id: ~p, reason: ~p", [buffer_endpoint_id(Buffer), Reason]), - Buffer#buffer{flight_num = 0} + Buffer#buffer{timer_map = #{}, flight_num = 0} end. -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(). buffer_endpoint_id(#buffer{endpoint = #endpoint{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. diff --git a/src/endpoint/endpoint_http.erl b/src/endpoint/endpoint_http.erl index 61a8049..2377e9e 100644 --- a/src/endpoint/endpoint_http.erl +++ b/src/endpoint/endpoint_http.erl @@ -70,7 +70,7 @@ handle_call(get_stat, _From, State = #state{buffer = Buffer}) -> {noreply, NewState :: #state{}} | {noreply, NewState :: #state{}, timeout() | hibernate} | {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), {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}) -> @@ -103,7 +103,10 @@ handle_info({next_data, Id, Metric}, State = #state{buffer = Buffer, endpoint = {error, Reason} -> logger:warning("[endpoint_http] url: ~p, get error: ~p", [Url, Reason]), {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 %% @doc This function is called by a gen_server when it is about to @@ -151,4 +154,4 @@ patch_headers(Metric, Token) when is_binary(Token), Token /= <<>> -> Sign = iot_util:sha256(erlang:iolist_to_binary([Token, Metric, Token])), [{<<"X-Signature">>, Sign}]; patch_headers(_, _) -> - []. \ No newline at end of file + []. diff --git a/src/endpoint/endpoint_kafka.erl b/src/endpoint/endpoint_kafka.erl index 08c5c9d..777b97a 100644 --- a/src/endpoint/endpoint_kafka.erl +++ b/src/endpoint/endpoint_kafka.erl @@ -70,6 +70,9 @@ disconnected(cast, {reload, NEndpoint = #endpoint{matcher = NMatcher}}, reload_endpoint(Matcher, NMatcher, ClientId, NEndpoint, State); disconnected(state_timeout, 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) -> {keep_state, State}; 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}, [{state_timeout, ?RETRY_INTERVAL, connect}]} 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}) -> ack_buffer(Id, Buffer, State); connected(info, {'EXIT', ClientPid, Reason}, diff --git a/src/endpoint/endpoint_mqtt.erl b/src/endpoint/endpoint_mqtt.erl index d5ee42b..66e564c 100644 --- a/src/endpoint/endpoint_mqtt.erl +++ b/src/endpoint/endpoint_mqtt.erl @@ -110,6 +110,9 @@ disconnected(internal, do_connect, {keep_state, State#state{conn_pid = undefined, inflight = #{}}, [{state_timeout, ?RETRY_INTERVAL, connect}]} 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) -> {keep_state, State}; 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}, [{state_timeout, ?RETRY_INTERVAL, connect}]} 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}, State = #state{conn_pid = ConnPid, buffer = Buffer}) -> logger:debug("[endpoint_mqtt] Recv a DISONNECT packet - ReasonCode: ~p, Properties: ~p", [ReasonCode, Properties]), diff --git a/src/endpoint/endpoint_timer.erl b/src/endpoint/endpoint_timer.erl deleted file mode 100644 index d919be8..0000000 --- a/src/endpoint/endpoint_timer.erl +++ /dev/null @@ -1,129 +0,0 @@ -%%%------------------------------------------------------------------- -%%% @author anlicheng -%%% @copyright (C) 2024, -%%% @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 -%%%===================================================================