%%%------------------------------------------------------------------- %%% @author aresei %%% @copyright (C) 2023, %%% @doc %%% %%% @end %%% Created : 06. 7月 2023 12:02 %%%------------------------------------------------------------------- -module(endpoint_buffer). -include("endpoint.hrl"). %% 消息重发间隔 -define(RETRY_INTERVAL, 5000). %% 最大重试次数,不包含首次发送 -define(MAX_RETRY_TIMES, 3). %% 与 endpoint_outbox 的单条记录限制保持一致;fast path 也不能绕过该限制。 -define(MAX_PAYLOAD_BYTES, 32 * 1024). -export([new/2, append/2, append_only/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]). -type flight_source() :: outbox | memory. -type timer_entry() :: {reference(), binary(), non_neg_integer(), flight_source()}. -record(buffer, { endpoint :: #endpoint{}, outbox :: endpoint_outbox:outbox(), %% 当前待 ack 的数据及其重试定时器 #{Id => {TimerRef, Payload, RetryTimes, Source}} timer_map = #{} :: #{integer() => timer_entry()}, %% 内存 fast path 使用负数 id,避免与 outbox 的正整数 seq 冲突。 next_memory_id = -1 :: integer(), %% 窗口大小,允许最大的未确认消息数 window_size = 10, %% 未确认的消息数 flight_num = 0, %% 记录成功处理的消息数 acc_num = 0 }). -type buffer() :: #buffer{}. -spec new(Endpoint :: #endpoint{}, WindowSize :: integer()) -> Buffer :: #buffer{}. new(Endpoint = #endpoint{id = Id}, WindowSize) when is_integer(WindowSize), WindowSize > 0 -> %% 读取配置文件 {ok, Endpoints} = application:get_env(iot, endpoints), RootDir = proplists:get_value(root_dir, Endpoints), OutboxDir = filename:join(RootDir, integer_to_list(Id)), {ok, Outbox} = endpoint_outbox:open(OutboxDir, #{}), #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) -> case validate_payload_size(Payload, Buffer) of ok -> case FlightNum < WindowSize andalso outbox_empty(Outbox) of true -> dispatch_memory(Payload, Buffer); false -> case append_to_outbox(Payload, Buffer) of {ok, NBuffer} -> trigger_next(NBuffer); {dropped, NBuffer} -> trigger_next(NBuffer); {error, NBuffer} -> NBuffer end end; error -> Buffer end. -spec append_only(Payload :: binary(), Buffer :: #buffer{}) -> NBuffer :: #buffer{}. append_only(Payload, Buffer = #buffer{}) when is_binary(Payload) -> case validate_payload_size(Payload, Buffer) of ok -> case append_to_outbox(Payload, Buffer) of {ok, NBuffer} -> NBuffer; {dropped, NBuffer} -> NBuffer; {error, NBuffer} -> NBuffer end; error -> Buffer end. -spec trigger_n(Buffer :: #buffer{}) -> NBuffer :: #buffer{}. trigger_n(Buffer = #buffer{window_size = WindowSize}) -> %% 最多允许window_size lists:foldl(fun(_, Buffer0) -> trigger_next(Buffer0) end, Buffer, lists:seq(1, WindowSize)). %% 触发读取下一条数据 -spec trigger_next(Buffer :: #buffer{}) -> NBuffer :: #buffer{}. trigger_next(Buffer = #buffer{outbox = Outbox, flight_num = FlightNum, window_size = WindowSize}) -> case FlightNum < WindowSize of false -> Buffer; true -> case endpoint_outbox:next(Outbox) of eof -> Buffer; {ok, Id, Payload, NOutbox} -> ReceiverPid = self(), ReceiverPid ! {next_data, Id, Payload}, schedule_retry(Id, Payload, 0, outbox, 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]), Buffer 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, Source}, 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, Source, Buffer#buffer{timer_map = NTimerMap}); {{TimerRef, Payload, RetryTimes, outbox}, NTimerMap} -> logger:warning("[endpoint_buffer] drop message after retries exhausted, endpoint_id: ~p, id: ~p, retries: ~p", [buffer_endpoint_id(Buffer), Id, RetryTimes]), drop_outbox_message(Id, Payload, RetryTimes, Buffer#buffer{timer_map = NTimerMap}); {{TimerRef, _, RetryTimes, memory}, NTimerMap} -> logger:warning("[endpoint_buffer] drop memory message after retries exhausted, endpoint_id: ~p, id: ~p, retries: ~p", [buffer_endpoint_id(Buffer), Id, RetryTimes]), drop_memory_message(Buffer#buffer{timer_map = NTimerMap}); {{_OtherTimerRef, _OtherPayload, _RetryTimes, _Source}, _NTimerMap} -> Buffer; error -> Buffer end. -spec ack(Id :: integer(), Buffer :: #buffer{}) -> NBuffer :: #buffer{}. 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, memory}, NTimerMap} -> _ = erlang:cancel_timer(TimerRef), NBuffer = Buffer#buffer{ timer_map = NTimerMap, acc_num = AccNum + 1, flight_num = max(FlightNum - 1, 0) }, trigger_next(NBuffer); {{TimerRef, Payload, RetryTimes, outbox}, 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) }, trigger_next(NBuffer); {error, Reason} -> logger:warning("[endpoint_buffer] ack outbox failed, endpoint_id: ~p, id: ~p, reason: ~p", [buffer_endpoint_id(Buffer), Id, Reason]), schedule_retry(Id, Payload, RetryTimes, outbox, Buffer#buffer{timer_map = NTimerMap}) end; error -> Buffer end. %% 获取当前统计信息 -spec stat(Buffer :: #buffer{}) -> map(). stat(#buffer{acc_num = AccNum, outbox = Outbox, flight_num = FlightNum, timer_map = TimerMap}) -> OutboxStat = endpoint_outbox:stat(Outbox), WriteSeq = maps:get(write_seq, OutboxStat, 0), AckedSeq = maps:get(acked_seq, OutboxStat, 0), OutboxFlightNum = count_inflight(outbox, TimerMap), MemoryFlightNum = count_inflight(memory, TimerMap), QueueNum = max(WriteSeq - AckedSeq - OutboxFlightNum, 0), OutboxStat#{ <<"acc_num">> => AccNum, <<"queue_num">> => QueueNum, <<"inflight_num">> => FlightNum, <<"outbox_inflight_num">> => OutboxFlightNum, <<"memory_inflight_num">> => MemoryFlightNum }. -spec cleanup(Buffer :: #buffer{}) -> #buffer{}. cleanup(Buffer = #buffer{timer_map = TimerMap}) -> cancel_timers(TimerMap), NBuffer0 = persist_memory_inflight(Buffer), reset_reader_after_recover(NBuffer0#buffer{timer_map = #{}, flight_num = 0}). -spec recover_inflight(Buffer :: #buffer{}) -> #buffer{}. recover_inflight(Buffer = #buffer{timer_map = TimerMap}) -> cancel_timers(TimerMap), NBuffer0 = persist_memory_inflight(Buffer), case endpoint_outbox:reset_reader(NBuffer0#buffer.outbox) of {ok, NOutbox} -> NBuffer0#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]), NBuffer0#buffer{timer_map = #{}, flight_num = 0} end. -spec resize(Buffer :: #buffer{}, WindowSize :: integer()) -> #buffer{}. resize(Buffer = #buffer{}, WindowSize) when is_integer(WindowSize), WindowSize > 0 -> trigger_n(Buffer#buffer{window_size = WindowSize}). %%%=================================================================== %%% Internal functions %%%=================================================================== -spec buffer_endpoint_id(buffer()) -> integer(). buffer_endpoint_id(#buffer{endpoint = #endpoint{id = Id}}) -> Id. -spec validate_payload_size(binary(), buffer()) -> ok | error. validate_payload_size(Payload, Buffer) -> case byte_size(Payload) =< ?MAX_PAYLOAD_BYTES of true -> ok; false -> logger:warning("[endpoint_buffer] payload too large, endpoint_id: ~p, size: ~p, max_size: ~p", [buffer_endpoint_id(Buffer), byte_size(Payload), ?MAX_PAYLOAD_BYTES]), error end. -spec outbox_empty(endpoint_outbox:outbox()) -> boolean(). outbox_empty(Outbox) -> Stat = endpoint_outbox:stat(Outbox), maps:get(write_seq, Stat, 0) =:= maps:get(acked_seq, Stat, 0). -spec append_to_outbox(binary(), buffer()) -> {ok | dropped | error, buffer()}. append_to_outbox(Payload, Buffer = #buffer{outbox = Outbox}) -> case endpoint_outbox:append(Payload, Outbox) of {ok, _Seq, NOutbox} -> {ok, Buffer#buffer{outbox = NOutbox}}; {dropped, capacity_reached, NOutbox} -> logger:warning("[endpoint_buffer] outbox capacity reached, endpoint_id: ~p", [buffer_endpoint_id(Buffer)]), {dropped, Buffer#buffer{outbox = NOutbox}}; {error, Reason} -> logger:warning("[endpoint_buffer] append outbox failed, endpoint_id: ~p, reason: ~p", [buffer_endpoint_id(Buffer), Reason]), {error, Buffer} end. -spec dispatch_memory(binary(), buffer()) -> buffer(). dispatch_memory(Payload, Buffer = #buffer{next_memory_id = Id, flight_num = FlightNum}) -> ReceiverPid = self(), ReceiverPid ! {next_data, Id, Payload}, schedule_retry(Id, Payload, 0, memory, Buffer#buffer{ next_memory_id = Id - 1, flight_num = FlightNum + 1 }). -spec schedule_retry(integer(), binary(), non_neg_integer(), flight_source(), buffer()) -> buffer(). schedule_retry(Id, Payload, RetryTimes, Source, 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, Source}, TimerMap)}. -spec cancel_timers(#{integer() => timer_entry()}) -> ok. cancel_timers(TimerMap) -> lists:foreach(fun({_Id, {TimerRef, _Payload, _RetryTimes, _Source}}) -> _ = erlang:cancel_timer(TimerRef), ok end, maps:to_list(TimerMap)), ok. -spec drop_outbox_message(integer(), binary(), non_neg_integer(), buffer()) -> buffer(). drop_outbox_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, outbox, Buffer) end. -spec drop_memory_message(buffer()) -> buffer(). drop_memory_message(Buffer = #buffer{acc_num = AccNum, flight_num = FlightNum}) -> trigger_next(Buffer#buffer{ acc_num = AccNum + 1, flight_num = max(FlightNum - 1, 0) }). -spec persist_memory_inflight(buffer()) -> buffer(). persist_memory_inflight(Buffer = #buffer{timer_map = TimerMap}) -> MemoryInflight = lists:sort( fun({IdA, _PayloadA}, {IdB, _PayloadB}) -> IdA > IdB end, [{Id, Payload} || {Id, {_TimerRef, Payload, _RetryTimes, memory}} <- maps:to_list(TimerMap)] ), lists:foldl(fun({_Id, Payload}, AccBuffer) -> case append_to_outbox(Payload, AccBuffer) of {ok, NBuffer} -> NBuffer; {dropped, NBuffer} -> NBuffer; {error, NBuffer} -> NBuffer end end, Buffer, MemoryInflight). -spec reset_reader_after_recover(buffer()) -> buffer(). reset_reader_after_recover(Buffer = #buffer{outbox = Outbox}) -> case endpoint_outbox:reset_reader(Outbox) of {ok, NOutbox} -> Buffer#buffer{outbox = NOutbox}; {error, Reason} -> logger:warning("[endpoint_buffer] reset reader failed, endpoint_id: ~p, reason: ~p", [buffer_endpoint_id(Buffer), Reason]), Buffer end. -spec count_inflight(flight_source(), #{integer() => timer_entry()}) -> non_neg_integer(). count_inflight(Source, TimerMap) -> maps:fold(fun(_Id, {_TimerRef, _Payload, _RetryTimes, EntrySource}, Acc) -> case EntrySource =:= Source of true -> Acc + 1; false -> Acc end end, 0, TimerMap).