fix outbox

This commit is contained in:
anlicheng 2026-04-22 16:13:30 +08:00
parent 30171d0e42
commit 7a1d341835

View File

@ -40,7 +40,6 @@
}).
-record(outbox, {
endpoint_id :: integer() | binary() | atom(),
dir :: file:filename_all(),
metadata_path :: file:filename_all(),
max_records = 500000 :: pos_integer(),
@ -54,18 +53,29 @@
-type outbox() :: #outbox{}.
-define(DEFAULT_ROOT_DIR, "var/outbox").
-define(METADATA_FILE, "metadata.json").
-define(SEGMENT_EXT, ".log").
-define(MAX_PAYLOAD_BYTES, 32 * 1024).
-define(RECORD_MAGIC, <<"EPOBOX01">>).
-define(RECORD_VERSION, 1).
-define(RECORD_HEADER_BYTES, 21).
-define(RECORD_FOOTER_BYTES, 4).
-define(RECORD_OVERHEAD_BYTES, ?RECORD_HEADER_BYTES + ?RECORD_FOOTER_BYTES).
-define(RESYNC_READ_BYTES, 4096).
%%%===================================================================
%%% API
%%%===================================================================
-spec open(integer() | binary() | atom(), map()) -> {ok, outbox()} | {error, term()}.
open(EndpointId, Opts) when is_map(Opts) ->
RootDir = maps:get(root_dir, Opts, ?DEFAULT_ROOT_DIR),
Dir = filename:join(RootDir, endpoint_dir_name(EndpointId)),
%% Opts
%% - max_records => pos_integer()
%% segment 500000
%% - max_bytes => pos_integer()
%% segment 134217728
%% - max_segments => pos_integer()
%% segment 10
-spec open(file:filename_all(), map()) -> {ok, outbox()} | {error, term()}.
open(Dir, Opts) when is_map(Opts) ->
ok = ensure_dir(Dir),
MetadataPath = filename:join(Dir, ?METADATA_FILE),
Segments = load_segments(Dir),
@ -75,7 +85,6 @@ open(EndpointId, Opts) when is_map(Opts) ->
AckedSeq = min(maps:get(acked_seq, Meta, 0), LastSeq),
AckedSegment = segment_id_for_seq(Segments, AckedSeq, maps:get(acked_segment, Meta, 1)),
Outbox0 = #outbox{
endpoint_id = EndpointId,
dir = Dir,
metadata_path = MetadataPath,
max_records = maps:get(max_records, Opts, 500000),
@ -98,10 +107,11 @@ open(EndpointId, Opts) when is_map(Opts) ->
offset = ReadOffset
}
},
case open_fds(Outbox1) of
{ok, Outbox} ->
ok = persist_metadata(Outbox),
{ok, Outbox};
maybe
{ok, Outbox} ?= open_fds(Outbox1),
ok ?= persist_metadata(Outbox),
{ok, Outbox}
else
{error, Reason} ->
{error, Reason}
end.
@ -114,18 +124,19 @@ close(Outbox = #outbox{}) ->
-spec append(binary(), outbox()) ->
{ok, Seq :: pos_integer(), outbox()} | {dropped, capacity_reached, outbox()} | {error, term()}.
append(Payload, Outbox = #outbox{}) when is_binary(Payload) ->
case validate_payload_size(Payload) of
ok ->
RecordBytes = encoded_record_size(Payload),
case ensure_writable_segment(RecordBytes, Outbox) of
{drop, NOutbox} ->
{dropped, capacity_reached, NOutbox};
{ok, SegmentId, NOutbox0} ->
case ensure_write_fd(SegmentId, NOutbox0) of
{ok, NOutbox1} ->
maybe
{ok, NOutbox1} ?= ensure_write_fd(SegmentId, NOutbox0),
Writer1 = NOutbox1#outbox.writer,
Seq = Writer1#outbox_writer.next_seq,
SegmentPath = segment_path(NOutbox1#outbox.dir, SegmentId),
case append_record(Writer1#outbox_writer.fd, Seq, Payload) of
ok ->
ok ?= append_record(Writer1#outbox_writer.fd, Seq, Payload),
Segments = update_written_segment(
NOutbox1#outbox.segments, SegmentId, SegmentPath, Seq, RecordBytes),
NOutbox2 = NOutbox1#outbox{
@ -135,19 +146,16 @@ append(Payload, Outbox = #outbox{}) when is_binary(Payload) ->
},
segments = Segments
},
case maybe_reset_reader_after_append(Outbox, NOutbox2) of
{ok, NOutbox} ->
ok = persist_metadata(NOutbox),
{ok, Seq, NOutbox};
{error, Reason} ->
{error, Reason}
end;
{error, Reason} ->
{error, Reason}
end;
{ok, NOutbox} ?= maybe_reset_reader_after_append(Outbox, NOutbox2),
ok ?= persist_metadata(NOutbox),
{ok, Seq, NOutbox}
else
{error, Reason} ->
{error, Reason}
end
end;
{error, Reason} ->
{error, Reason}
end.
-spec ack(pos_integer(), outbox()) -> {ok, outbox()} | {error, term()}.
@ -170,13 +178,13 @@ ack(Seq, Outbox = #outbox{checkpoint = Checkpoint}) when is_integer(Seq), Seq >
pending_acks = PendingAcks
}
},
{ok, Outbox1} = prune_acked_segments(Outbox0),
case NackedSeq =:= AckedSeq of
true ->
{ok, Outbox1};
false ->
ok = persist_metadata(Outbox1),
maybe
{ok, Outbox1} ?= prune_acked_segments(Outbox0),
ok ?= maybe_persist_metadata(NackedSeq =/= AckedSeq, Outbox1),
{ok, Outbox1}
else
{error, Reason} ->
{error, Reason}
end
end.
@ -193,7 +201,7 @@ next(Outbox = #outbox{
case ensure_read_fd(Outbox) of
{ok, NOutbox0} ->
Reader0 = NOutbox0#outbox.reader,
case read_record(Reader0#outbox_reader.fd) of
case read_next_record(Reader0#outbox_reader.fd, Offset) of
eof ->
case next_segment_id(Segments, SegmentId) of
undefined ->
@ -206,9 +214,9 @@ next(Outbox = #outbox{
{error, Reason}
end
end;
{ok, Seq, Payload, RecordBytes} ->
{ok, _RecordOffset, Seq, Payload, NextOffset} ->
{ok, Seq, Payload, NOutbox0#outbox{
reader = Reader0#outbox_reader{offset = Offset + RecordBytes}
reader = Reader0#outbox_reader{offset = NextOffset}
}};
{error, Reason} ->
{error, Reason}
@ -249,25 +257,6 @@ stat(#outbox{
%%% Internal functions
%%%===================================================================
-spec endpoint_dir_name(integer() | binary() | atom()) -> string().
endpoint_dir_name(EndpointId) when is_integer(EndpointId) ->
"endpoint-" ++ integer_to_list(EndpointId);
endpoint_dir_name(EndpointId) when is_atom(EndpointId) ->
endpoint_dir_name(atom_to_binary(EndpointId, utf8));
endpoint_dir_name(EndpointId) when is_binary(EndpointId) ->
"endpoint-" ++ sanitize_binary(EndpointId).
-spec sanitize_binary(binary()) -> string().
sanitize_binary(Bin) ->
lists:map(fun
(C) when C >= $a, C =< $z -> C;
(C) when C >= $A, C =< $Z -> C;
(C) when C >= $0, C =< $9 -> C;
($-) -> $-;
($_) -> $_;
(_) -> $_
end, binary_to_list(Bin)).
-spec ensure_dir(file:filename_all()) -> ok.
ensure_dir(Dir) ->
ok = filelib:ensure_dir(filename:join(Dir, "dummy")),
@ -376,7 +365,7 @@ scan_segment(Path, Id) ->
case file:open(Path, [read, raw, binary]) of
{ok, Fd} ->
try
scan_segment_records(Fd, #segment{id = Id, path = Path})
scan_segment_records(Fd, 0, #segment{id = Id, path = Path})
after
ok = file:close(Fd)
end;
@ -384,12 +373,13 @@ scan_segment(Path, Id) ->
{error, Reason}
end.
-spec scan_segment_records(file:fd(), #segment{}) -> {ok, #segment{}} | {error, term()}.
scan_segment_records(Fd, Segment = #segment{}) ->
case read_record(Fd) of
-spec scan_segment_records(file:fd(), non_neg_integer(), #segment{}) -> {ok, #segment{}} | {error, term()}.
scan_segment_records(Fd, Offset, Segment = #segment{}) ->
case read_next_record(Fd, Offset) of
eof ->
{ok, Segment};
{ok, Seq, _Payload, RecordBytes} ->
{ok, RecordOffset, Seq, _Payload, NextOffset} ->
RecordBytes = NextOffset - RecordOffset,
NextSegment = Segment#segment{
start_seq = case Segment#segment.records of
0 -> Seq;
@ -399,14 +389,20 @@ scan_segment_records(Fd, Segment = #segment{}) ->
records = Segment#segment.records + 1,
bytes = Segment#segment.bytes + RecordBytes
},
scan_segment_records(Fd, NextSegment);
scan_segment_records(Fd, NextOffset, NextSegment);
{error, Reason} ->
{error, Reason}
end.
-spec encoded_record_size(binary()) -> pos_integer().
encoded_record_size(Payload) ->
12 + byte_size(Payload).
?RECORD_OVERHEAD_BYTES + byte_size(Payload).
-spec validate_payload_size(binary()) -> ok | {error, payload_too_large}.
validate_payload_size(Payload) when byte_size(Payload) =< ?MAX_PAYLOAD_BYTES ->
ok;
validate_payload_size(_Payload) ->
{error, payload_too_large}.
-spec ensure_writable_segment(pos_integer(), outbox()) ->
{ok, pos_integer(), outbox()} | {drop, outbox()}.
@ -442,7 +438,16 @@ ensure_writable_segment(RecordBytes, Outbox = #outbox{
-spec append_record(file:fd(), pos_integer(), binary()) -> ok | {error, term()}.
append_record(Fd, Seq, Payload) when Fd =/= undefined ->
Bin = <<Seq:64/unsigned-big-integer, (byte_size(Payload)):32/unsigned-big-integer, Payload/binary>>,
PayloadSize = byte_size(Payload),
Crc = record_crc(?RECORD_VERSION, Seq, PayloadSize, Payload),
Bin = <<
?RECORD_MAGIC/binary,
?RECORD_VERSION:8,
Seq:64/unsigned-big-integer,
PayloadSize:32/unsigned-big-integer,
Payload/binary,
Crc:32/unsigned-big-integer
>>,
case file:write(Fd, Bin) of
ok ->
ok;
@ -612,37 +617,162 @@ locate_offset_for_seq(#segment{path = Path}, TargetSeq) ->
-spec locate_offset_loop(file:fd(), pos_integer(), non_neg_integer()) ->
{ok, non_neg_integer()} | {error, term()}.
locate_offset_loop(Fd, TargetSeq, Offset) ->
case read_record(Fd) of
case read_next_record(Fd, Offset) of
eof ->
{ok, Offset};
{ok, Seq, _Payload, _RecordBytes} when Seq >= TargetSeq ->
{ok, Offset};
{ok, _Seq, _Payload, RecordBytes} ->
locate_offset_loop(Fd, TargetSeq, Offset + RecordBytes);
{ok, RecordOffset, Seq, _Payload, _NextOffset} when Seq >= TargetSeq ->
{ok, RecordOffset};
{ok, _RecordOffset, _Seq, _Payload, NextOffset} ->
locate_offset_loop(Fd, TargetSeq, NextOffset);
{error, Reason} ->
{error, Reason}
end.
-spec read_record(file:fd()) -> eof | {ok, pos_integer(), binary(), pos_integer()} | {error, term()}.
read_record(Fd) ->
case file:read(Fd, 12) of
-spec read_next_record(file:fd(), non_neg_integer()) ->
eof | {ok, non_neg_integer(), pos_integer(), binary(), non_neg_integer()} | {error, term()}.
read_next_record(Fd, Offset) ->
case read_record_at(Fd, Offset) of
eof ->
eof;
{ok, <<Seq:64/unsigned-big-integer, PayloadSize:32/unsigned-big-integer>>} ->
case file:read(Fd, PayloadSize) of
{ok, Payload} when byte_size(Payload) =:= PayloadSize ->
{ok, Seq, Payload, 12 + PayloadSize};
{ok, Seq, Payload, NextOffset} ->
{ok, Offset, Seq, Payload, NextOffset};
{error, _Reason} ->
resync_next_record(Fd, Offset + 1)
end.
-spec read_record_at(file:fd(), non_neg_integer()) ->
eof | {ok, pos_integer(), binary(), non_neg_integer()} | {error, term()}.
read_record_at(Fd, Offset) ->
maybe
{ok, _} ?= file:position(Fd, Offset),
read_record(Fd, Offset)
else
{error, Reason} ->
{error, Reason}
end.
-spec read_record(file:fd(), non_neg_integer()) ->
eof | {ok, pos_integer(), binary(), non_neg_integer()} | {error, term()}.
read_record(Fd, Offset) ->
case file:read(Fd, ?RECORD_HEADER_BYTES) of
eof ->
{error, truncated_payload};
eof;
{ok, <<
Magic:8/binary,
Version:8,
Seq:64/unsigned-big-integer,
PayloadSize:32/unsigned-big-integer
>>} when Magic =:= ?RECORD_MAGIC, Version =:= ?RECORD_VERSION ->
maybe
ok ?= validate_read_payload_size(PayloadSize),
read_record_body(Fd, Offset, Seq, PayloadSize)
else
{error, Reason} ->
{error, Reason}
end;
{ok, _} ->
{ok, Header} when byte_size(Header) < ?RECORD_HEADER_BYTES ->
{error, truncated_header};
{ok, _} ->
{error, invalid_record_header};
{error, Reason} ->
{error, Reason}
end.
-spec validate_read_payload_size(non_neg_integer()) -> ok | {error, payload_too_large}.
validate_read_payload_size(PayloadSize) when PayloadSize =< ?MAX_PAYLOAD_BYTES ->
ok;
validate_read_payload_size(_PayloadSize) ->
{error, payload_too_large}.
-spec read_record_body(file:fd(), non_neg_integer(), pos_integer(), non_neg_integer()) ->
{ok, pos_integer(), binary(), non_neg_integer()} | {error, term()}.
read_record_body(Fd, Offset, Seq, PayloadSize) ->
case file:read(Fd, PayloadSize + ?RECORD_FOOTER_BYTES) of
{ok, <<Payload:PayloadSize/binary, StoredCrc:32/unsigned-big-integer>>} ->
ExpectedCrc = record_crc(?RECORD_VERSION, Seq, PayloadSize, Payload),
case StoredCrc =:= ExpectedCrc of
true ->
{ok, Seq, Payload, Offset + ?RECORD_OVERHEAD_BYTES + PayloadSize};
false ->
{error, invalid_record_crc}
end;
eof ->
{error, truncated_payload};
{ok, _} ->
{error, truncated_payload};
{error, Reason} ->
{error, Reason}
end.
-spec record_crc(byte(), pos_integer(), non_neg_integer(), binary()) -> non_neg_integer().
record_crc(Version, Seq, PayloadSize, Payload) ->
erlang:crc32(<<
Version:8,
Seq:64/unsigned-big-integer,
PayloadSize:32/unsigned-big-integer,
Payload/binary
>>).
-spec resync_next_record(file:fd(), non_neg_integer()) ->
eof | {ok, non_neg_integer(), pos_integer(), binary(), non_neg_integer()} | {error, term()}.
resync_next_record(Fd, SearchOffset) ->
case find_next_magic(Fd, SearchOffset) of
eof ->
eof;
{ok, RecordOffset} ->
case read_record_at(Fd, RecordOffset) of
eof ->
eof;
{ok, Seq, Payload, NextOffset} ->
{ok, RecordOffset, Seq, Payload, NextOffset};
{error, _Reason} ->
resync_next_record(Fd, RecordOffset + 1)
end;
{error, Reason} ->
{error, Reason}
end.
-spec find_next_magic(file:fd(), non_neg_integer()) ->
eof | {ok, non_neg_integer()} | {error, term()}.
find_next_magic(Fd, Offset) ->
maybe
{ok, _} ?= file:position(Fd, Offset),
find_next_magic_loop(Fd, Offset, <<>>)
else
{error, Reason} ->
{error, Reason}
end.
-spec find_next_magic_loop(file:fd(), non_neg_integer(), binary()) ->
eof | {ok, non_neg_integer()} | {error, term()}.
find_next_magic_loop(Fd, Offset, Tail) ->
case file:read(Fd, ?RESYNC_READ_BYTES) of
eof ->
eof;
{ok, Bin} ->
SearchBin = <<Tail/binary, Bin/binary>>,
BaseOffset = Offset - byte_size(Tail),
case binary:match(SearchBin, ?RECORD_MAGIC) of
{Pos, _Len} ->
{ok, BaseOffset + Pos};
nomatch ->
NextTail = trailing_bytes(SearchBin, byte_size(?RECORD_MAGIC) - 1),
find_next_magic_loop(Fd, Offset + byte_size(Bin), NextTail)
end;
{error, Reason} ->
{error, Reason}
end.
-spec trailing_bytes(binary(), non_neg_integer()) -> binary().
trailing_bytes(Bin, KeepBytes) ->
Size = byte_size(Bin),
case Size =< KeepBytes of
true ->
Bin;
false ->
binary:part(Bin, Size - KeepBytes, KeepBytes)
end.
-spec next_segment_id([#segment{}], pos_integer()) -> undefined | pos_integer().
next_segment_id(Segments, SegmentId) ->
case lists:dropwhile(fun(#segment{id = Id}) -> Id =< SegmentId end, Segments) of
@ -660,9 +790,10 @@ segment_path(Dir, SegmentId) ->
open_fds(Outbox) ->
case maybe_open_write_fd(Outbox) of
{ok, Outbox0} ->
case maybe_open_read_fd(Outbox0) of
{ok, NOutbox} ->
{ok, NOutbox};
maybe
{ok, NOutbox} ?= maybe_open_read_fd(Outbox0),
{ok, NOutbox}
else
{error, Reason} ->
_ = close_open_fds(Outbox0),
{error, Reason}
@ -729,17 +860,11 @@ ensure_read_fd(Outbox = #outbox{
SegmentId = Reader#outbox_reader.segment,
Offset = Reader#outbox_reader.offset,
Path = segment_path(Dir, SegmentId),
case file:open(Path, [read, raw, binary]) of
case open_read_fd(Path, Offset) of
{ok, Fd} ->
case file:position(Fd, Offset) of
{ok, _} ->
{ok, Outbox#outbox{
reader = Reader#outbox_reader{fd = Fd}
}};
{error, Reason} ->
ok = file:close(Fd),
{error, Reason}
end;
{error, Reason} ->
{error, Reason}
end.
@ -814,3 +939,25 @@ maybe_close_fd(Fd) ->
{error, _Reason} ->
ok
end.
-spec maybe_persist_metadata(boolean(), outbox()) -> ok | {error, term()}.
maybe_persist_metadata(true, Outbox) ->
persist_metadata(Outbox);
maybe_persist_metadata(false, _Outbox) ->
ok.
-spec open_read_fd(file:filename_all(), non_neg_integer()) -> {ok, file:fd()} | {error, term()}.
open_read_fd(Path, Offset) ->
case file:open(Path, [read, raw, binary]) of
{ok, Fd} ->
maybe
{ok, _} ?= file:position(Fd, Offset),
{ok, Fd}
else
{error, Reason} ->
ok = file:close(Fd),
{error, Reason}
end;
{error, Reason} ->
{error, Reason}
end.