修复本地的数据缓存

This commit is contained in:
anlicheng 2026-05-12 15:29:09 +08:00
parent ae37cbe3d8
commit d7106f5ccd
4 changed files with 562 additions and 127 deletions

View File

@ -40,7 +40,6 @@ WebSocket server
- `efka_logger`:部署日志落盘。
- `efka_service_sup`:动态管理每个已注册微服务对应的 `efka_service` 进程。
- `cache_model`DETS 离线缓存。
- `efka_service_model`DETS 服务状态表。
- `efka_subscription`:本地 topic 订阅中心。
- `efka_iot_client`:连接上游 TLS server 的状态机。
@ -92,7 +91,7 @@ WebSocket server
1. 微服务发送 `ServiceCast.MetricData`
2. `efka_service_channel` 调用 `efka_service:metric_data(ServicePid, RouteKey, Metric)`
3. `efka_service` 转发给 `efka_iot_client:metric_data(RouteKey, Metric)`
4. `efka_iot_client` 如果处于 activated 状态,直接发给上游;否则写入 `cache_model` 离线缓存
4. `efka_iot_client` 如果处于 activated 状态,直接发给上游;否则写入持久化 outbox
### 3.2 EFKA 到上游TLS 长连接
@ -111,7 +110,7 @@ WebSocket server
2. 读取 `efka.tls_server_address`,用 `ssl:connect/4` 建立 TLS 连接。
3. TLS socket 使用 `{packet, 4}``{active, true}`
4. 连接成功后发送 `AuthRequest`
5. 收到鉴权成功 reply 后进入 `activated`,并触发 `flush_cache`
5. 收到鉴权成功 reply 后进入 `activated`,并触发 outbox 刷出
6. 连接或鉴权失败时关闭 socket5 秒后重连。
上游协议同样用第一个字节区分帧类型:
@ -124,7 +123,7 @@ WebSocket server
`efka_iot_client` 上报的内容:
- `metric_data`业务指标数据activated 时实时发送,否则进入 DETS 缓存
- `metric_data`业务指标数据activated 时实时发送,否则进入持久化 outbox
- `task_event_stream`Docker 部署任务流式日志,只在 activated 时发送。
- `close_task_event_stream`:任务结束事件,只在 activated 时发送。
@ -240,16 +239,15 @@ topic 匹配规则:
- channel 关闭时把服务状态改成 stopped。
- 支持查询所有服务、运行中服务和单个服务状态。
### 6.2 `cache_model`
### 6.2 `efka_iot_outbox`
使用 DETS 表 `cache`
使用 append-only log 文件和 metadata 文件
作用:
- 当 `efka_iot_client` 不在 activated 状态时,把待上报 packet 缓存下来。
- `efka_iot_client` 激活后循环 `fetch_next -> send -> delete` 刷缓存。
缓存 id 使用 `os:system_time(microsecond)` 生成。
- `efka_iot_client` 激活后循环 `next -> send -> ack` 刷出。
- 所有记录 ack 后会截断 log下一轮从 seq 1 重新开始。
## 7. 日志
@ -272,7 +270,7 @@ topic 匹配规则:
注意:
- `efka_service_model` `cache_model` 打开 DETS 前假设 `dets_dir` 已存在,代码里没有显式创建目录。
- `efka_service_model` 打开 DETS 前假设 `dets_dir` 已存在,代码里没有显式创建目录。
- `docker_client` 固定使用 `/var/run/docker.sock`,每次请求内部由短生命周期普通进程执行 open/request/close。
## 9. 当前看到的几个注意点
@ -282,7 +280,6 @@ topic 匹配规则:
- `efka_iot_client:send_result_reply/3``send_error_reply/3` 编码 `ReplyFrame` 后没有加 `FRAME_REPLY` 前缀;接收侧是否期望裸 protobuf 需要确认。
- `docker_deployer:ensure_container_absent/2` 当前没有真正确保旧容器不存在,只是上报日志。
- `efka_subscription` 计算了 topic `order`,但匹配广播时没有使用优先级排序。
- `cache_model` 使用 DETS bag但 id 由微秒时间生成,理论上极端并发下可能碰撞。
- `docker_commands` 部分错误响应解码没有统一使用 `[return_maps]`,有些分支可能匹配不到 map。
- `docker_events` 存在但未启动,且使用 shell 命令 `docker events`,与其他 Docker API 访问方式不同。

View File

@ -1,103 +0,0 @@
%%%-------------------------------------------------------------------
%%% @author anlicheng
%%% @copyright (C) 2026, <COMPANY>
%%% @doc
%%% DETS backed outbound packet cache for efka_iot_client.
%%% @end
%%%-------------------------------------------------------------------
-module(efka_iot_cache).
-author("anlicheng").
-define(CACHE_TAB, cache).
-define(MAX_CACHE_ITEMS, 5000000).
-define(MAX_CACHE_FILE_SIZE, 1073741824).
-export([open/0, close/0, insert/1, fetch_next/0, delete/1]).
-spec open() -> ok | {error, term()}.
open() ->
{ok, DetsDir} = application:get_env(efka, dets_dir),
File = DetsDir ++ "cache.dets",
case dets:open_file(?CACHE_TAB, [{file, File}, {type, bag}, {keypos, 1}]) of
{ok, ?CACHE_TAB} ->
ok;
{error, Reason} ->
{error, Reason}
end.
-spec close() -> ok.
close() ->
case dets:close(?CACHE_TAB) of
ok ->
ok;
{error, not_owner} ->
ok
end.
-spec insert(binary()) -> {ok, non_neg_integer()} | {error, term()}.
insert(Data) when is_binary(Data) ->
case dets:insert(?CACHE_TAB, {generate_cache_id(), Data}) of
ok ->
trim_limits(0);
{error, Reason} ->
{error, Reason}
end.
-spec fetch_next() -> error | {ok, {integer(), binary()}}.
fetch_next() ->
case dets:first(?CACHE_TAB) of
'$end_of_table' ->
error;
Key ->
case dets:lookup(?CACHE_TAB, Key) of
[Entry | _] ->
{ok, Entry};
[] ->
fetch_next()
end
end.
-spec delete(integer()) -> ok | {error, term()}.
delete(Id) when is_integer(Id) ->
dets:delete(?CACHE_TAB, Id).
-spec generate_cache_id() -> integer().
generate_cache_id() ->
erlang:unique_integer([monotonic, positive]).
-spec over_limit() -> boolean().
over_limit() ->
item_count() > ?MAX_CACHE_ITEMS orelse file_size() > ?MAX_CACHE_FILE_SIZE.
-spec item_count() -> non_neg_integer().
item_count() ->
dets:info(?CACHE_TAB, size).
-spec file_size() -> non_neg_integer().
file_size() ->
dets:info(?CACHE_TAB, file_size).
-spec trim_limits(non_neg_integer()) -> {ok, non_neg_integer()} | {error, term()}.
trim_limits(DroppedCount) ->
case over_limit() of
true ->
case delete_oldest_entry() of
ok ->
trim_limits(DroppedCount + 1);
error ->
{ok, DroppedCount};
{error, Reason} ->
{error, Reason}
end;
false ->
{ok, DroppedCount}
end.
-spec delete_oldest_entry() -> ok | error | {error, term()}.
delete_oldest_entry() ->
case dets:first(?CACHE_TAB) of
'$end_of_table' ->
error;
Key ->
dets:delete(?CACHE_TAB, Key)
end.

View File

@ -32,6 +32,7 @@
-record(state, {
socket :: undefined | ssl:sslsocket(),
outbox :: efka_iot_outbox:outbox(),
%% auth请求的refauth请求和响应的对应关系
auth_ref = undefined :: undefined | binary(),
ping_timer_ref = undefined :: undefined | reference(),
@ -77,10 +78,10 @@ start_link() ->
-spec init(list()) -> {ok, atom(), #state{}}.
init([]) ->
case efka_iot_cache:open() of
ok ->
case efka_iot_outbox:open() of
{ok, Outbox} ->
erlang:start_timer(0, self(), create_transport),
{ok, ?STATE_DISCONNECTED, #state{socket = undefined}};
{ok, ?STATE_DISCONNECTED, #state{socket = undefined, outbox = Outbox}};
{error, Reason} ->
{stop, Reason}
end.
@ -89,7 +90,7 @@ init([]) ->
callback_mode() ->
handle_event_function.
%% , DETS
%% outbox
-spec handle_event(term(), term(), atom(), #state{}) -> term().
handle_event(cast, {metric_data, RouteKey, Metric}, StateName, State = #state{socket = Socket}) ->
Packet = term_to_binary({<<"message">>, {<<"data">>, #{<<"route_key">> => RouteKey, <<"metric">> => Metric}}}),
@ -98,8 +99,19 @@ handle_event(cast, {metric_data, RouteKey, Metric}, StateName, State = #state{so
ok = ssl:send(Socket, Packet),
{keep_state, State};
_ ->
{ok, DroppedCount} = efka_iot_cache:insert(Packet),
{keep_state, State#state{dropped_message_count = State#state.dropped_message_count + DroppedCount}}
case efka_iot_outbox:append(Packet, State#state.outbox) of
{ok, Outbox} ->
{keep_state, State#state{outbox = Outbox}};
{dropped, capacity_reached, Outbox} ->
logger:warning("[efka_iot_client] outbox capacity reached, drop offline metric"),
{keep_state, State#state{
outbox = Outbox,
dropped_message_count = State#state.dropped_message_count + 1
}};
{error, Reason} ->
logger:warning("[efka_iot_client] append outbox failed, reason: ~p", [Reason]),
{keep_state, State#state{dropped_message_count = State#state.dropped_message_count + 1}}
end
end;
%% Task的stream流
@ -161,12 +173,20 @@ handle_event(info, {timeout, _TimerRef, ssl_ping}, _StateName, State) ->
%%
handle_event(info, flush_cache, ?STATE_ACTIVATED, State = #state{socket = Socket}) ->
case efka_iot_cache:fetch_next() of
{ok, {Id, Packet}} ->
case efka_iot_outbox:next(State#state.outbox) of
{ok, Seq, Packet} ->
ok = ssl:send(Socket, Packet),
ok = efka_iot_cache:delete(Id),
{keep_state, State, [{next_event, info, flush_cache}]};
error ->
case efka_iot_outbox:ack(Seq, State#state.outbox) of
{ok, Outbox} ->
{keep_state, State#state{outbox = Outbox}, [{next_event, info, flush_cache}]};
{error, Reason} ->
logger:warning("[efka_iot_client] ack outbox failed, seq: ~p, reason: ~p", [Seq, Reason]),
{keep_state, State}
end;
eof ->
{keep_state, State};
{error, Reason} ->
logger:warning("[efka_iot_client] read outbox failed, reason: ~p", [Reason]),
{keep_state, State}
end;
handle_event(info, flush_cache, _, State) ->
@ -276,10 +296,10 @@ handle_container_command(Ref, Request, Socket) ->
ok.
-spec terminate(term(), atom(), #state{}) -> ok.
terminate(Reason, _StateName, State = #state{socket = Socket}) ->
terminate(Reason, _StateName, State = #state{socket = Socket, outbox = Outbox}) ->
cancel_ssl_ping(State),
disconnect(Socket),
efka_iot_cache:close(),
efka_iot_outbox:close(Outbox),
logger:notice("[efka_iot_client] terminate with reason: ~p", [Reason]),
ok.

View File

@ -0,0 +1,521 @@
%%%-------------------------------------------------------------------
%%% @doc
%%% efka_iot_client
%%%
%%% efka_iot_client 线 packet使
%%% segment
%%%
%%%
%%%
%%% - `segment-N.log' segment record
%%% - `metadata.term' #{write_seq => N, acked_seq => N,
%%% writer_segment => N} segment segment
%%% 使 metadata
%%%
%%%
%%% record
%%% - 4 unsigned big-endian payload size
%%% - term_to_binary({Seq, Packet}) payload
%%%
%%%
%%% - `append/2' record writer segment fsync
%%% segment write_seq metadata
%%% - segment 200000
%%% - outbox 5 live segment segment
%%% 5 live segment packet
%%% - segment live segment 5 writer
%%% `segment-(N + 1).log'
%%%
%%%
%%% - `next/1' end_seq acked_seq live segment
%%% segment Seq acked_seq record
%%% - efka_iot_client ack packet
%%% - `ack/2' Seq
%%% acked_seq
%%%
%%% segment
%%% - `ack/2' segment record segment
%%% segment
%%% segment
%%% - acked_seq write_seq record
%%% segment 0 append
%%% `segment-1.log' Seq 1
%%% @end
%%%-------------------------------------------------------------------
-module(efka_iot_outbox).
-export([open/0, close/1, append/2, next/1, ack/2]).
-export_type([outbox/0]).
-record(segment, {
id :: pos_integer(),
path :: file:filename_all(),
start_seq = 0 :: non_neg_integer(),
end_seq = 0 :: non_neg_integer(),
records = 0 :: non_neg_integer()
}).
-record(outbox, {
dir :: file:filename_all(),
metadata_path :: file:filename_all(),
writer_segment = 1 :: pos_integer(),
fd :: file:fd(),
segments = [] :: [#segment{}],
next_seq = 1 :: pos_integer(),
write_seq = 0 :: non_neg_integer(),
acked_seq = 0 :: non_neg_integer()
}).
-type outbox() :: #outbox{}.
-define(METADATA_FILE, "metadata.term").
-define(SEGMENT_PREFIX, "segment-").
-define(SEGMENT_EXT, ".log").
-define(SEGMENT_RECORD_LIMIT, 200000).
-define(MAX_SEGMENTS, 5).
-define(MAX_RECORD_BYTES, 16 * 1024 * 1024).
-spec open() -> {ok, outbox()} | {error, term()}.
open() ->
{ok, DetsDir} = application:get_env(efka, dets_dir),
Dir = filename:join(DetsDir, "iot_outbox"),
MetadataPath = filename:join(Dir, ?METADATA_FILE),
maybe
ok ?= ensure_dir(Dir),
{_MetaWriteSeq, MetaAckedSeq, MetaWriterSegment} = read_metadata(MetadataPath),
{ok, Segments} ?= load_segments(Dir),
ScannedWriteSeq = max_segment_end(Segments),
WriteSeq = ScannedWriteSeq,
AckedSeq = min(MetaAckedSeq, WriteSeq),
WriterSegment = writer_segment_id(Segments, MetaWriterSegment),
{ok, Fd} ?= open_writer(segment_path(Dir, WriterSegment)),
Outbox0 = #outbox{
dir = Dir,
metadata_path = MetadataPath,
writer_segment = WriterSegment,
fd = Fd,
segments = Segments,
next_seq = WriteSeq + 1,
write_seq = WriteSeq,
acked_seq = AckedSeq
},
{ok, Outbox1} ?= normalize_open_outbox(Outbox0),
ok ?= persist_metadata(Outbox1),
{ok, Outbox1}
else
{error, Reason} ->
{error, Reason}
end.
-spec close(outbox()) -> ok.
close(#outbox{fd = Fd}) ->
_ = file:close(Fd),
ok.
-spec append(binary(), outbox()) ->
{ok, outbox()} | {dropped, capacity_reached, outbox()} | {error, term()}.
append(Packet, Outbox0) when is_binary(Packet) ->
case ensure_writable_segment(Outbox0) of
{ok, Outbox1 = #outbox{fd = Fd, next_seq = Seq}} ->
Record = encode_record(Seq, Packet),
maybe
ok ?= file:write(Fd, Record),
ok ?= file:sync(Fd),
Outbox2 = mark_written(Outbox1, Seq),
ok ?= persist_metadata(Outbox2),
{ok, Outbox2}
else
{error, Reason} ->
{error, Reason}
end;
{dropped, capacity_reached, Outbox1} ->
{dropped, capacity_reached, Outbox1};
{error, Reason} ->
{error, Reason}
end.
-spec next(outbox()) -> eof | {ok, pos_integer(), binary()} | {error, term()}.
next(#outbox{segments = Segments, acked_seq = AckedSeq}) ->
case next_segment(Segments, AckedSeq) of
undefined ->
eof;
#segment{path = Path} ->
case file:open(Path, [read, binary]) of
{ok, Fd} ->
Result = read_next_unacked(Fd, AckedSeq),
_ = file:close(Fd),
Result;
{error, enoent} ->
eof;
{error, Reason} ->
{error, Reason}
end
end.
-spec ack(pos_integer(), outbox()) -> {ok, outbox()} | {error, term()}.
ack(Seq, Outbox = #outbox{acked_seq = AckedSeq}) when is_integer(Seq), Seq =< AckedSeq ->
{ok, Outbox};
ack(Seq, Outbox = #outbox{acked_seq = AckedSeq}) when is_integer(Seq), Seq =:= AckedSeq + 1 ->
AckedOutbox = Outbox#outbox{acked_seq = Seq},
maybe
ok ?= persist_metadata(AckedOutbox),
{ok, NOutbox} ?= maybe_reset_or_prune(AckedOutbox),
ok ?= persist_metadata(NOutbox),
{ok, NOutbox}
else
{error, Reason} ->
{error, Reason}
end;
ack(Seq, #outbox{acked_seq = AckedSeq}) when is_integer(Seq) ->
{error, {non_contiguous_ack, Seq, AckedSeq}}.
-spec ensure_dir(file:filename_all()) -> ok | {error, term()}.
ensure_dir(Dir) ->
filelib:ensure_dir(filename:join(Dir, "dummy")).
-spec read_metadata(file:filename_all()) ->
{non_neg_integer(), non_neg_integer(), pos_integer()}.
read_metadata(MetadataPath) ->
case file:read_file(MetadataPath) of
{ok, Bin} ->
safe_metadata(binary_to_term(Bin, [safe]));
{error, _} ->
{0, 0, 1}
end.
-spec safe_metadata(term()) -> {non_neg_integer(), non_neg_integer(), pos_integer()}.
safe_metadata(#{write_seq := WriteSeq, acked_seq := AckedSeq, writer_segment := WriterSegment})
when is_integer(WriteSeq), WriteSeq >= 0,
is_integer(AckedSeq), AckedSeq >= 0,
is_integer(WriterSegment), WriterSegment > 0 ->
{WriteSeq, AckedSeq, WriterSegment};
safe_metadata(#{write_seq := WriteSeq, acked_seq := AckedSeq})
when is_integer(WriteSeq), WriteSeq >= 0, is_integer(AckedSeq), AckedSeq >= 0 ->
{WriteSeq, AckedSeq, 1};
safe_metadata(_) ->
{0, 0, 1}.
-spec persist_metadata(outbox()) -> ok | {error, term()}.
persist_metadata(#outbox{
metadata_path = MetadataPath,
writer_segment = WriterSegment,
write_seq = WriteSeq,
acked_seq = AckedSeq
}) ->
Metadata = term_to_binary(#{
write_seq => WriteSeq,
acked_seq => AckedSeq,
writer_segment => WriterSegment
}),
TmpPath = MetadataPath ++ ".tmp",
case file:write_file(TmpPath, Metadata, [write, binary]) of
ok ->
file:rename(TmpPath, MetadataPath);
{error, Reason} ->
{error, Reason}
end.
-spec load_segments(file:filename_all()) -> {ok, [#segment{}]} | {error, term()}.
load_segments(Dir) ->
case file:list_dir(Dir) of
{ok, Names} ->
Ids = lists:sort([Id || Name <- Names, {ok, Id} <- [parse_segment_id(Name)]]),
load_segments(Dir, Ids, []);
{error, enoent} ->
{ok, []};
{error, Reason} ->
{error, Reason}
end.
-spec load_segments(file:filename_all(), [pos_integer()], [#segment{}]) ->
{ok, [#segment{}]} | {error, term()}.
load_segments(_Dir, [], Acc) ->
{ok, lists:reverse(Acc)};
load_segments(Dir, [Id | Rest], Acc) ->
Path = segment_path(Dir, Id),
case scan_segment(Path, Id) of
{ok, undefined} ->
load_segments(Dir, Rest, Acc);
{ok, Segment} ->
load_segments(Dir, Rest, [Segment | Acc]);
{error, Reason} ->
{error, Reason}
end.
-spec parse_segment_id(string()) -> {ok, pos_integer()} | error.
parse_segment_id(Name) ->
case re:run(Name, "^" ++ ?SEGMENT_PREFIX ++ "([0-9]+)\\" ++ ?SEGMENT_EXT ++ "$",
[{capture, [1], list}])
of
{match, [Digits]} ->
case list_to_integer(Digits) of
Id when Id > 0 ->
{ok, Id};
_ ->
error
end;
nomatch ->
error
end.
-spec segment_path(file:filename_all(), pos_integer()) -> file:filename_all().
segment_path(Dir, Id) ->
filename:join(Dir, ?SEGMENT_PREFIX ++ integer_to_list(Id) ++ ?SEGMENT_EXT).
-spec open_writer(file:filename_all()) -> {ok, file:fd()} | {error, term()}.
open_writer(Path) ->
file:open(Path, [append, binary]).
-spec writer_segment_id([#segment{}], pos_integer()) -> pos_integer().
writer_segment_id([], MetaWriterSegment) ->
max(1, MetaWriterSegment);
writer_segment_id(Segments, _MetaWriterSegment) ->
(lists:last(Segments))#segment.id.
-spec max_segment_end([#segment{}]) -> non_neg_integer().
max_segment_end([]) ->
0;
max_segment_end(Segments) ->
lists:max([Segment#segment.end_seq || Segment <- Segments]).
-spec normalize_open_outbox(outbox()) -> {ok, outbox()} | {error, term()}.
normalize_open_outbox(Outbox = #outbox{acked_seq = Seq, write_seq = Seq}) when Seq > 0 ->
reset_empty_outbox(Outbox);
normalize_open_outbox(Outbox) ->
prune_acked_segments(Outbox).
-spec ensure_writable_segment(outbox()) ->
{ok, outbox()} | {dropped, capacity_reached, outbox()} | {error, term()}.
ensure_writable_segment(Outbox = #outbox{segments = Segments, writer_segment = WriterSegment}) ->
case find_segment(WriterSegment, Segments) of
undefined ->
{ok, Outbox};
#segment{records = Records} when Records < ?SEGMENT_RECORD_LIMIT ->
{ok, Outbox};
#segment{} when length(Segments) >= ?MAX_SEGMENTS ->
{dropped, capacity_reached, Outbox};
#segment{} ->
rotate_writer(Outbox, WriterSegment + 1)
end.
-spec rotate_writer(outbox(), pos_integer()) -> {ok, outbox()} | {error, term()}.
rotate_writer(Outbox = #outbox{dir = Dir, fd = Fd}, NewSegment) ->
maybe
ok ?= file:sync(Fd),
{ok, NFd} ?= open_writer(segment_path(Dir, NewSegment)),
_ = file:close(Fd),
{ok, Outbox#outbox{writer_segment = NewSegment, fd = NFd}}
else
{error, Reason} ->
{error, Reason}
end.
-spec mark_written(outbox(), pos_integer()) -> outbox().
mark_written(Outbox = #outbox{
dir = Dir,
writer_segment = WriterSegment,
segments = Segments
}, Seq) ->
Path = segment_path(Dir, WriterSegment),
Segment0 = case find_segment(WriterSegment, Segments) of
undefined ->
#segment{id = WriterSegment, path = Path, start_seq = Seq};
Segment ->
Segment
end,
Records = Segment0#segment.records + 1,
Segment1 = Segment0#segment{end_seq = Seq, records = Records},
Outbox#outbox{
segments = replace_segment(Segment1, Segments),
next_seq = Seq + 1,
write_seq = Seq
}.
-spec replace_segment(#segment{}, [#segment{}]) -> [#segment{}].
replace_segment(Segment, Segments) ->
lists:sort(
fun(A, B) -> A#segment.id < B#segment.id end,
[Segment | [S || S <- Segments, S#segment.id =/= Segment#segment.id]]
).
-spec find_segment(pos_integer(), [#segment{}]) -> #segment{} | undefined.
find_segment(Id, Segments) ->
case [Segment || Segment <- Segments, Segment#segment.id =:= Id] of
[Segment] ->
Segment;
[] ->
undefined
end.
-spec next_segment([#segment{}], non_neg_integer()) -> #segment{} | undefined.
next_segment([], _AckedSeq) ->
undefined;
next_segment([Segment = #segment{end_seq = EndSeq} | _Rest], AckedSeq) when EndSeq > AckedSeq ->
Segment;
next_segment([_Segment | Rest], AckedSeq) ->
next_segment(Rest, AckedSeq).
-spec maybe_reset_or_prune(outbox()) -> {ok, outbox()} | {error, term()}.
maybe_reset_or_prune(Outbox = #outbox{acked_seq = Seq, write_seq = Seq}) ->
reset_empty_outbox(Outbox);
maybe_reset_or_prune(Outbox) ->
prune_acked_segments(Outbox).
-spec prune_acked_segments(outbox()) -> {ok, outbox()} | {error, term()}.
prune_acked_segments(Outbox = #outbox{segments = Segments, acked_seq = AckedSeq}) ->
{DeleteSegments, KeepSegments} = take_acked_segments(Segments, AckedSeq, []),
maybe
ok ?= delete_files([Segment#segment.path || Segment <- DeleteSegments]),
{ok, Outbox#outbox{segments = KeepSegments}}
else
{error, Reason} ->
{error, Reason}
end.
-spec take_acked_segments([#segment{}], non_neg_integer(), [#segment{}]) ->
{[#segment{}], [#segment{}]}.
take_acked_segments([Segment = #segment{end_seq = EndSeq} | Rest], AckedSeq, Acc)
when EndSeq =< AckedSeq ->
take_acked_segments(Rest, AckedSeq, [Segment | Acc]);
take_acked_segments(Segments, _AckedSeq, Acc) ->
{lists:reverse(Acc), Segments}.
-spec reset_empty_outbox(outbox()) -> {ok, outbox()} | {error, term()}.
reset_empty_outbox(Outbox = #outbox{
dir = Dir,
fd = Fd,
writer_segment = WriterSegment,
segments = Segments
}) ->
CurrentWriterPath = segment_path(Dir, WriterSegment),
SegmentPaths = [Segment#segment.path || Segment <- Segments],
Paths = lists:usort([CurrentWriterPath | SegmentPaths]),
maybe
ok ?= file:close(Fd),
ok ?= delete_files(Paths),
{ok, NFd} ?= open_writer(segment_path(Dir, 1)),
{ok, Outbox#outbox{
writer_segment = 1,
fd = NFd,
segments = [],
next_seq = 1,
write_seq = 0,
acked_seq = 0
}}
else
{error, Reason} ->
{error, Reason}
end.
-spec delete_files([file:filename_all()]) -> ok | {error, term()}.
delete_files([]) ->
ok;
delete_files([Path | Rest]) ->
case file:delete(Path) of
ok ->
delete_files(Rest);
{error, enoent} ->
delete_files(Rest);
{error, Reason} ->
{error, {delete_failed, Path, Reason}}
end.
-spec scan_segment(file:filename_all(), pos_integer()) ->
{ok, #segment{} | undefined} | {error, term()}.
scan_segment(Path, Id) ->
case file:open(Path, [read, binary]) of
{ok, Fd} ->
Result = scan_segment(Fd, Id, Path, 0, 0, 0),
_ = file:close(Fd),
Result;
{error, enoent} ->
{ok, undefined};
{error, Reason} ->
{error, Reason}
end.
-spec scan_segment(file:fd(), pos_integer(), file:filename_all(),
non_neg_integer(), non_neg_integer(), non_neg_integer()) ->
{ok, #segment{} | undefined} | {error, term()}.
scan_segment(Fd, Id, Path, StartSeq, EndSeq, Records) ->
case read_record(Fd) of
eof when Records =:= 0 ->
{ok, undefined};
eof ->
{ok, #segment{
id = Id,
path = Path,
start_seq = StartSeq,
end_seq = EndSeq,
records = Records
}};
{ok, Seq, _Packet} ->
NStartSeq = case StartSeq of
0 -> Seq;
_ -> StartSeq
end,
scan_segment(Fd, Id, Path, NStartSeq, Seq, Records + 1);
{error, Reason} ->
{error, Reason}
end.
-spec encode_record(pos_integer(), binary()) -> binary().
encode_record(Seq, Packet) ->
Payload = term_to_binary({Seq, Packet}),
Size = byte_size(Payload),
<<Size:32/unsigned-big, Payload/binary>>.
-spec read_next_unacked(file:fd(), non_neg_integer()) ->
eof | {ok, pos_integer(), binary()} | {error, term()}.
read_next_unacked(Fd, AckedSeq) ->
case read_record(Fd) of
eof ->
eof;
{ok, Seq, Packet} when Seq > AckedSeq ->
{ok, Seq, Packet};
{ok, _Seq, _Packet} ->
read_next_unacked(Fd, AckedSeq);
{error, Reason} ->
{error, Reason}
end.
-spec read_record(file:fd()) -> eof | {ok, pos_integer(), binary()} | {error, term()}.
read_record(Fd) ->
case file:read(Fd, 4) of
eof ->
eof;
{ok, <<Size:32/unsigned-big>>} when Size > 0, Size =< ?MAX_RECORD_BYTES ->
read_record_payload(Fd, Size);
{ok, <<Size:32/unsigned-big>>} ->
{error, {invalid_record_size, Size}};
{ok, Partial} ->
{error, {truncated_record_header, Partial}};
{error, Reason} ->
{error, Reason}
end.
-spec read_record_payload(file:fd(), pos_integer()) ->
{ok, pos_integer(), binary()} | {error, term()}.
read_record_payload(Fd, Size) ->
case file:read(Fd, Size) of
{ok, Payload} when byte_size(Payload) =:= Size ->
safe_record(Payload);
{ok, Payload} ->
{error, {truncated_record_payload, Size, byte_size(Payload)}};
eof ->
{error, {truncated_record_payload, Size, 0}};
{error, Reason} ->
{error, Reason}
end.
-spec safe_record(binary()) -> {ok, pos_integer(), binary()} | {error, term()}.
safe_record(Payload) ->
try binary_to_term(Payload, [safe]) of
{Seq, Packet} when is_integer(Seq), Seq > 0, is_binary(Packet) ->
{ok, Seq, Packet};
Other ->
{error, {invalid_record, Other}}
catch
error:Reason ->
{error, {invalid_record, Reason}}
end.