From d7106f5ccd61de74ab8ed707dabf5636cfc6a4b6 Mon Sep 17 00:00:00 2001 From: anlicheng <244108715@qq.com> Date: Tue, 12 May 2026 15:29:09 +0800 Subject: [PATCH] =?UTF-8?q?=E4=BF=AE=E5=A4=8D=E6=9C=AC=E5=9C=B0=E7=9A=84?= =?UTF-8?q?=E6=95=B0=E6=8D=AE=E7=BC=93=E5=AD=98?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- CODE_LOGIC_OVERVIEW.md | 19 +- apps/efka/src/iot/efka_iot_cache.erl | 103 ----- apps/efka/src/iot/efka_iot_client.erl | 46 ++- apps/efka/src/iot/efka_iot_outbox.erl | 521 ++++++++++++++++++++++++++ 4 files changed, 562 insertions(+), 127 deletions(-) delete mode 100644 apps/efka/src/iot/efka_iot_cache.erl create mode 100644 apps/efka/src/iot/efka_iot_outbox.erl diff --git a/CODE_LOGIC_OVERVIEW.md b/CODE_LOGIC_OVERVIEW.md index 39cceb7..c8b88bc 100644 --- a/CODE_LOGIC_OVERVIEW.md +++ b/CODE_LOGIC_OVERVIEW.md @@ -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. 连接或鉴权失败时关闭 socket,5 秒后重连。 上游协议同样用第一个字节区分帧类型: @@ -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 访问方式不同。 diff --git a/apps/efka/src/iot/efka_iot_cache.erl b/apps/efka/src/iot/efka_iot_cache.erl deleted file mode 100644 index b21e1b9..0000000 --- a/apps/efka/src/iot/efka_iot_cache.erl +++ /dev/null @@ -1,103 +0,0 @@ -%%%------------------------------------------------------------------- -%%% @author anlicheng -%%% @copyright (C) 2026, -%%% @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. diff --git a/apps/efka/src/iot/efka_iot_client.erl b/apps/efka/src/iot/efka_iot_client.erl index 96fc000..ce0ff8f 100644 --- a/apps/efka/src/iot/efka_iot_client.erl +++ b/apps/efka/src/iot/efka_iot_client.erl @@ -32,6 +32,7 @@ -record(state, { socket :: undefined | ssl:sslsocket(), + outbox :: efka_iot_outbox:outbox(), %% 保存当前auth请求的ref,用来建立auth请求和响应的对应关系 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. diff --git a/apps/efka/src/iot/efka_iot_outbox.erl b/apps/efka/src/iot/efka_iot_outbox.erl new file mode 100644 index 0000000..dc31ee0 --- /dev/null +++ b/apps/efka/src/iot/efka_iot_outbox.erl @@ -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), + <>. + +-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, <>} when Size > 0, Size =< ?MAX_RECORD_BYTES -> + read_record_payload(Fd, Size); + {ok, <>} -> + {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.