fix outbox
This commit is contained in:
parent
ee0d9fcc08
commit
50864cf7bd
@ -21,8 +21,16 @@
|
|||||||
%%% - max_record_bytes:允许读取的单条 record 最大 payload 字节数。
|
%%% - max_record_bytes:允许读取的单条 record 最大 payload 字节数。
|
||||||
%%%
|
%%%
|
||||||
%%% record 格式:
|
%%% record 格式:
|
||||||
%%% - 4 字节 unsigned big-endian payload size。
|
%%% - 8 字节 magic bytes,当前为 <<"EFKAOBX1">>。
|
||||||
%%% - term_to_binary({Seq, Packet}) payload。
|
%%% - 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 当前
|
%%% - `append/2' 将一条 record 写入当前 writer segment,随后 fsync 当前
|
||||||
@ -32,15 +40,16 @@
|
|||||||
%%% - outbox 最多保留多少个 live segment 由 open/1 的 max_segments 参数
|
%%% - outbox 最多保留多少个 live segment 由 open/1 的 max_segments 参数
|
||||||
%%% 决定。当当前 segment 已满,且 live segment 数已经达到上限时,新
|
%%% 决定。当当前 segment 已满,且 live segment 数已经达到上限时,新
|
||||||
%%% packet 会被丢弃,已有磁盘数据保持不变。
|
%%% packet 会被丢弃,已有磁盘数据保持不变。
|
||||||
%%% - 当当前 segment 已满,但 live segment 数量少于 5 个时,writer 会关闭
|
%%% - 当当前 segment 已满,但 live segment 数量尚未达到上限时,writer 会
|
||||||
%%% 旧文件,并开始写入 `segment-(N + 1).log'。
|
%%% 关闭旧文件,并开始写入 `segment-(N + 1).log'。
|
||||||
%%%
|
%%%
|
||||||
%%% 读取逻辑:
|
%%% 读取逻辑:
|
||||||
%%% - `next/1' 找到第一个 end_seq 大于 acked_seq 的 live segment,打开该
|
%%% - `next/1' 找到第一个 end_seq 大于 acked_seq 的 live segment,打开该
|
||||||
%%% segment,并只返回第一条 Seq 大于 acked_seq 的 record。
|
%%% segment,并只返回第一条 Seq 大于 acked_seq 的 record。
|
||||||
%%% - 本模块不做批量读取;efka_iot_client 每次只发送并 ack 一条 packet。
|
%%% - 本模块不做批量读取;efka_iot_client 每次只发送并 ack 一条 packet。
|
||||||
%%% - `ack/2' 只接受连续的下一个 Seq。这样可以保持重放逻辑简单:进程崩溃后,
|
%%% - `ack/2' 接受当前读取到的 Seq。正常情况下 Seq 连续递增;如果某条
|
||||||
%%% acked_seq 之后的数据会重新发送。
|
%%% record 损坏并被读取端跳过,ack 允许推进到后续合法 record 的 Seq,
|
||||||
|
%%% 避免单条坏数据阻断整个 segment 后续数据的消费。
|
||||||
%%%
|
%%%
|
||||||
%%% segment 清理逻辑:
|
%%% segment 清理逻辑:
|
||||||
%%% - 当 `ack/2' 消费完某个 segment 的最后一条 record,并且后续 segment
|
%%% - 当 `ack/2' 消费完某个 segment 的最后一条 record,并且后续 segment
|
||||||
@ -83,6 +92,11 @@
|
|||||||
-define(METADATA_FILE, "metadata.term").
|
-define(METADATA_FILE, "metadata.term").
|
||||||
-define(SEGMENT_PREFIX, "segment-").
|
-define(SEGMENT_PREFIX, "segment-").
|
||||||
-define(SEGMENT_EXT, ".log").
|
-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() :: #{
|
-type open_options() :: #{
|
||||||
dir := file:filename_all(),
|
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()}.
|
-spec ack(pos_integer(), outbox()) -> {ok, outbox()} | {error, term()}.
|
||||||
ack(Seq, Outbox = #outbox{acked_seq = AckedSeq}) when is_integer(Seq), Seq =< AckedSeq ->
|
ack(Seq, Outbox = #outbox{acked_seq = AckedSeq}) when is_integer(Seq), Seq =< AckedSeq ->
|
||||||
{ok, Outbox};
|
{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},
|
AckedOutbox = Outbox#outbox{acked_seq = Seq},
|
||||||
maybe
|
maybe
|
||||||
ok ?= persist_metadata(AckedOutbox),
|
ok ?= persist_metadata(AckedOutbox),
|
||||||
@ -188,7 +202,7 @@ ack(Seq, Outbox = #outbox{acked_seq = AckedSeq}) when is_integer(Seq), Seq =:= A
|
|||||||
{error, Reason}
|
{error, Reason}
|
||||||
end;
|
end;
|
||||||
ack(Seq, #outbox{acked_seq = AckedSeq}) when is_integer(Seq) ->
|
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()}.
|
-spec ensure_dir(file:filename_all()) -> ok | {error, term()}.
|
||||||
ensure_dir(Dir) ->
|
ensure_dir(Dir) ->
|
||||||
@ -491,9 +505,12 @@ scan_segment(Fd, Id, Path, MaxRecordBytes, StartSeq, EndSeq, Records) ->
|
|||||||
|
|
||||||
-spec encode_record(pos_integer(), binary()) -> binary().
|
-spec encode_record(pos_integer(), binary()) -> binary().
|
||||||
encode_record(Seq, Packet) ->
|
encode_record(Seq, Packet) ->
|
||||||
Payload = term_to_binary({Seq, Packet}),
|
PacketSize = byte_size(Packet),
|
||||||
Size = byte_size(Payload),
|
Magic = ?RECORD_MAGIC,
|
||||||
<<Size:32/unsigned-big, Payload/binary>>.
|
Header = <<Magic/binary, ?RECORD_VERSION:8/unsigned, ?RECORD_HEADER_SIZE:8/unsigned,
|
||||||
|
PacketSize:32/unsigned-big, Seq:64/unsigned-big>>,
|
||||||
|
Crc32 = erlang:crc32(<<Header/binary, Packet/binary>>),
|
||||||
|
<<Header/binary, Packet/binary, Crc32:32/unsigned-big>>.
|
||||||
|
|
||||||
-spec read_next_unacked(file:fd(), non_neg_integer(), pos_integer()) ->
|
-spec read_next_unacked(file:fd(), non_neg_integer(), pos_integer()) ->
|
||||||
eof | {ok, pos_integer(), binary()} | {error, term()}.
|
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()}.
|
-spec read_record(file:fd(), pos_integer()) -> eof | {ok, pos_integer(), binary()} | {error, term()}.
|
||||||
read_record(Fd, MaxRecordBytes) ->
|
read_record(Fd, MaxRecordBytes) ->
|
||||||
case file:read(Fd, 4) of
|
case find_next_magic(Fd) of
|
||||||
eof ->
|
eof ->
|
||||||
eof;
|
eof;
|
||||||
{ok, <<Size:32/unsigned-big>>} when Size > 0, Size =< MaxRecordBytes ->
|
{ok, StartPos} ->
|
||||||
read_record_payload(Fd, Size);
|
case read_record_after_magic(Fd, MaxRecordBytes) of
|
||||||
{ok, <<Size:32/unsigned-big>>} ->
|
{ok, Seq, Packet} ->
|
||||||
{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};
|
{ok, Seq, Packet};
|
||||||
Other ->
|
{resync, _Reason} ->
|
||||||
{error, {invalid_record, Other}}
|
case file:position(Fd, {bof, StartPos + 1}) of
|
||||||
catch
|
{ok, _} ->
|
||||||
error:Reason ->
|
read_record(Fd, MaxRecordBytes);
|
||||||
{error, {invalid_record, Reason}}
|
{error, Reason} ->
|
||||||
|
{error, Reason}
|
||||||
|
end
|
||||||
|
end;
|
||||||
|
{error, Reason} ->
|
||||||
|
{error, Reason}
|
||||||
|
end.
|
||||||
|
|
||||||
|
-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 ->
|
||||||
|
eof;
|
||||||
|
{ok, Byte} ->
|
||||||
|
Tail = binary:part(Window, 1, ?RECORD_MAGIC_SIZE - 1),
|
||||||
|
NWindow = <<Tail/binary, Byte/binary>>,
|
||||||
|
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 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, <<Version:8/unsigned, HeaderSize:8/unsigned, PacketSize:32/unsigned-big,
|
||||||
|
Seq:64/unsigned-big>>} ->
|
||||||
|
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, <<Packet:PacketSize/binary, Crc32:32/unsigned-big>>} ->
|
||||||
|
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 = <<Magic/binary, Version:8/unsigned, HeaderSize:8/unsigned,
|
||||||
|
PacketSize:32/unsigned-big, Seq:64/unsigned-big>>,
|
||||||
|
case erlang:crc32(<<Header/binary, Packet/binary>>) of
|
||||||
|
Crc32 ->
|
||||||
|
{ok, Seq, Packet};
|
||||||
|
Expected ->
|
||||||
|
{resync, {crc32_mismatch, Seq, Expected, Crc32}}
|
||||||
end.
|
end.
|
||||||
|
|||||||
Loading…
x
Reference in New Issue
Block a user