diff --git a/apps/efka/src/iot/efka_iot_outbox.erl b/apps/efka/src/iot/efka_iot_outbox.erl index 8e74f1e..1b44b93 100644 --- a/apps/efka/src/iot/efka_iot_outbox.erl +++ b/apps/efka/src/iot/efka_iot_outbox.erl @@ -21,8 +21,16 @@ %%% - max_record_bytes:允许读取的单条 record 最大 payload 字节数。 %%% %%% record 格式: -%%% - 4 字节 unsigned big-endian payload size。 -%%% - term_to_binary({Seq, Packet}) payload。 +%%% - 8 字节 magic bytes,当前为 <<"EFKAOBX1">>。 +%%% - 1 字节 version,当前为 1。 +%%% - 1 字节 header size,当前为 22。 +%%% - 4 字节 unsigned big-endian packet size。 +%%% - 8 字节 unsigned big-endian Seq。 +%%% - Packet 原始二进制内容。 +%%% - 4 字节 unsigned big-endian crc32,校验范围为 header + Packet。 +%%% - 读取时会先搜索 magic,再校验 version、header size、packet size 和 +%%% crc32。遇到损坏 record 时,会继续向后寻找下一个合法 magic,尽量恢复 +%%% 后续可读数据。 %%% %%% 写入逻辑: %%% - `append/2' 将一条 record 写入当前 writer segment,随后 fsync 当前 @@ -32,15 +40,16 @@ %%% - outbox 最多保留多少个 live segment 由 open/1 的 max_segments 参数 %%% 决定。当当前 segment 已满,且 live segment 数已经达到上限时,新 %%% packet 会被丢弃,已有磁盘数据保持不变。 -%%% - 当当前 segment 已满,但 live segment 数量少于 5 个时,writer 会关闭 -%%% 旧文件,并开始写入 `segment-(N + 1).log'。 +%%% - 当当前 segment 已满,但 live segment 数量尚未达到上限时,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 之后的数据会重新发送。 +%%% - `ack/2' 接受当前读取到的 Seq。正常情况下 Seq 连续递增;如果某条 +%%% record 损坏并被读取端跳过,ack 允许推进到后续合法 record 的 Seq, +%%% 避免单条坏数据阻断整个 segment 后续数据的消费。 %%% %%% segment 清理逻辑: %%% - 当 `ack/2' 消费完某个 segment 的最后一条 record,并且后续 segment @@ -83,6 +92,11 @@ -define(METADATA_FILE, "metadata.term"). -define(SEGMENT_PREFIX, "segment-"). -define(SEGMENT_EXT, ".log"). +-define(RECORD_MAGIC, <<"EFKAOBX1">>). +-define(RECORD_MAGIC_SIZE, 8). +-define(RECORD_VERSION, 1). +-define(RECORD_HEADER_SIZE, 22). +-define(RECORD_CRC_SIZE, 4). -type open_options() :: #{ dir := file:filename_all(), @@ -176,7 +190,7 @@ next(#outbox{segments = Segments, acked_seq = AckedSeq, max_record_bytes = MaxRe -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 -> +ack(Seq, Outbox = #outbox{acked_seq = AckedSeq}) when is_integer(Seq), Seq > AckedSeq -> AckedOutbox = Outbox#outbox{acked_seq = Seq}, maybe ok ?= persist_metadata(AckedOutbox), @@ -188,7 +202,7 @@ ack(Seq, Outbox = #outbox{acked_seq = AckedSeq}) when is_integer(Seq), Seq =:= A {error, Reason} end; ack(Seq, #outbox{acked_seq = AckedSeq}) when is_integer(Seq) -> - {error, {non_contiguous_ack, Seq, AckedSeq}}. + {error, {invalid_ack, Seq, AckedSeq}}. -spec ensure_dir(file:filename_all()) -> ok | {error, term()}. ensure_dir(Dir) -> @@ -491,9 +505,12 @@ scan_segment(Fd, Id, Path, MaxRecordBytes, StartSeq, EndSeq, Records) -> -spec encode_record(pos_integer(), binary()) -> binary(). encode_record(Seq, Packet) -> - Payload = term_to_binary({Seq, Packet}), - Size = byte_size(Payload), - <>. + PacketSize = byte_size(Packet), + Magic = ?RECORD_MAGIC, + Header = <>, + Crc32 = erlang:crc32(<
>), + <
>. -spec read_next_unacked(file:fd(), non_neg_integer(), pos_integer()) -> eof | {ok, pos_integer(), binary()} | {error, term()}. @@ -511,41 +528,121 @@ read_next_unacked(Fd, AckedSeq, MaxRecordBytes) -> -spec read_record(file:fd(), pos_integer()) -> eof | {ok, pos_integer(), binary()} | {error, term()}. read_record(Fd, MaxRecordBytes) -> - case file:read(Fd, 4) of + case find_next_magic(Fd) of eof -> eof; - {ok, <>} when Size > 0, Size =< MaxRecordBytes -> - read_record_payload(Fd, Size); - {ok, <>} -> - {error, {invalid_record_size, Size}}; - {ok, Partial} -> - {error, {truncated_record_header, Partial}}; + {ok, StartPos} -> + case read_record_after_magic(Fd, MaxRecordBytes) of + {ok, Seq, Packet} -> + {ok, Seq, Packet}; + {resync, _Reason} -> + case file:position(Fd, {bof, StartPos + 1}) of + {ok, _} -> + read_record(Fd, MaxRecordBytes); + {error, Reason} -> + {error, Reason} + end + end; {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)}}; +-spec find_next_magic(file:fd()) -> eof | {ok, non_neg_integer()} | {error, term()}. +find_next_magic(Fd) -> + case file:position(Fd, cur) of + {ok, StartPos} -> + case file:read(Fd, ?RECORD_MAGIC_SIZE) of + eof -> + eof; + {ok, Magic} when byte_size(Magic) < ?RECORD_MAGIC_SIZE -> + eof; + {ok, ?RECORD_MAGIC} -> + {ok, StartPos}; + {ok, Window} -> + find_next_magic(Fd, Window); + {error, Reason} -> + {error, Reason} + end; + {error, Reason} -> + {error, Reason} + end. + +-spec find_next_magic(file:fd(), binary()) -> eof | {ok, non_neg_integer()} | {error, term()}. +find_next_magic(Fd, Window) -> + case file:read(Fd, 1) of eof -> - {error, {truncated_record_payload, Size, 0}}; + eof; + {ok, Byte} -> + Tail = binary:part(Window, 1, ?RECORD_MAGIC_SIZE - 1), + NWindow = <>, + case NWindow of + ?RECORD_MAGIC -> + case file:position(Fd, cur) of + {ok, Pos} -> + {ok, Pos - ?RECORD_MAGIC_SIZE}; + {error, Reason} -> + {error, Reason} + end; + _ -> + find_next_magic(Fd, NWindow) + end; {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}} +-spec read_record_after_magic(file:fd(), pos_integer()) -> + {ok, pos_integer(), binary()} | {resync, term()}. +read_record_after_magic(Fd, MaxRecordBytes) -> + HeaderRestSize = ?RECORD_HEADER_SIZE - ?RECORD_MAGIC_SIZE, + case file:read(Fd, HeaderRestSize) of + {ok, <>} -> + read_record_payload(Fd, Version, HeaderSize, PacketSize, Seq, MaxRecordBytes); + {ok, Partial} -> + {resync, {truncated_record_header, Partial}}; + eof -> + {resync, truncated_record_header}; + {error, Reason} -> + {resync, Reason} + end. + +-spec read_record_payload(file:fd(), non_neg_integer(), non_neg_integer(), + non_neg_integer(), non_neg_integer(), pos_integer()) -> + {ok, pos_integer(), binary()} | {resync, term()}. +read_record_payload(_Fd, Version, _HeaderSize, _PacketSize, _Seq, _MaxRecordBytes) + when Version =/= ?RECORD_VERSION -> + {resync, {invalid_record_version, Version}}; +read_record_payload(_Fd, _Version, HeaderSize, _PacketSize, _Seq, _MaxRecordBytes) + when HeaderSize =/= ?RECORD_HEADER_SIZE -> + {resync, {invalid_record_header_size, HeaderSize}}; +read_record_payload(_Fd, _Version, _HeaderSize, PacketSize, _Seq, MaxRecordBytes) + when PacketSize =:= 0; PacketSize > MaxRecordBytes -> + {resync, {invalid_record_packet_size, PacketSize}}; +read_record_payload(_Fd, _Version, _HeaderSize, _PacketSize, Seq, _MaxRecordBytes) + when Seq =:= 0 -> + {resync, {invalid_record_seq, Seq}}; +read_record_payload(Fd, Version, HeaderSize, PacketSize, Seq, _MaxRecordBytes) -> + case file:read(Fd, PacketSize + ?RECORD_CRC_SIZE) of + {ok, <>} -> + validate_record_crc(Version, HeaderSize, PacketSize, Seq, Packet, Crc32); + {ok, Partial} -> + {resync, {truncated_record_payload, PacketSize, byte_size(Partial)}}; + eof -> + {resync, {truncated_record_payload, PacketSize, 0}}; + {error, Reason} -> + {resync, Reason} + end. + +-spec validate_record_crc(non_neg_integer(), non_neg_integer(), non_neg_integer(), + pos_integer(), binary(), non_neg_integer()) -> + {ok, pos_integer(), binary()} | {resync, term()}. +validate_record_crc(Version, HeaderSize, PacketSize, Seq, Packet, Crc32) -> + Magic = ?RECORD_MAGIC, + Header = <>, + case erlang:crc32(<
>) of + Crc32 -> + {ok, Seq, Packet}; + Expected -> + {resync, {crc32_mismatch, Seq, Expected, Crc32}} end.