iot_cloud/apps/endpoint/src/endpoint_buffer.erl
2026-05-30 15:25:40 +08:00

305 lines
13 KiB
Erlang

%%%-------------------------------------------------------------------
%%% @author aresei
%%% @copyright (C) 2023, <COMPANY>
%%% @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 的单条记录限制保持一致。
-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()},
%% 窗口大小,允许最大的未确认消息数
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(endpoint, 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{}) when is_binary(Payload) ->
case validate_payload_size(Payload, Buffer) of
ok ->
%% Persist before dispatch so a VM crash does not lose in-flight data.
case append_to_outbox(Payload, Buffer) of
{ok, NBuffer} ->
trigger_next(NBuffer);
{dropped, NBuffer} ->
trigger_next(NBuffer);
{error, NBuffer} ->
NBuffer
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}
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 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 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}
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).