964 lines
34 KiB
Erlang
964 lines
34 KiB
Erlang
%%%-------------------------------------------------------------------
|
|
%%% @author Codex
|
|
%%% @doc
|
|
%%% endpoint 持久化 outbox 的第一阶段实现:
|
|
%%% - 一个 endpoint 一个目录
|
|
%%% - segment 按递增编号命名
|
|
%%% - metadata 只记录连续 ack 指针
|
|
%%% - inflight 不落盘,崩溃后从 acked_seq + 1 重放
|
|
%%%-------------------------------------------------------------------
|
|
-module(endpoint_outbox).
|
|
|
|
-export([open/2, close/1, append/2, ack/2, next/1, reset_reader/1, stat/1]).
|
|
-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(),
|
|
bytes = 0 :: non_neg_integer()
|
|
}).
|
|
|
|
-record(outbox_writer, {
|
|
segment = 1 :: pos_integer(),
|
|
next_seq = 1 :: pos_integer(),
|
|
fd = undefined :: undefined | file:fd()
|
|
}).
|
|
|
|
-record(outbox_reader, {
|
|
segment :: undefined | pos_integer(),
|
|
offset = 0 :: non_neg_integer(),
|
|
fd = undefined :: undefined | file:fd()
|
|
}).
|
|
|
|
-record(outbox_checkpoint, {
|
|
acked_segment = 1 :: pos_integer(),
|
|
acked_seq = 0 :: non_neg_integer(),
|
|
pending_acks = [] :: [pos_integer()]
|
|
}).
|
|
|
|
-record(outbox, {
|
|
dir :: file:filename_all(),
|
|
metadata_path :: file:filename_all(),
|
|
max_records = 500000 :: pos_integer(),
|
|
max_bytes = 134217728 :: pos_integer(),
|
|
max_segments = 10 :: pos_integer(),
|
|
segments = [] :: [#segment{}],
|
|
writer = #outbox_writer{} :: #outbox_writer{},
|
|
reader = #outbox_reader{} :: #outbox_reader{},
|
|
checkpoint = #outbox_checkpoint{} :: #outbox_checkpoint{}
|
|
}).
|
|
|
|
-type outbox() :: #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
|
|
%%%===================================================================
|
|
|
|
%% 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),
|
|
Meta = read_metadata(MetadataPath),
|
|
LastSegmentId = last_segment_id(Segments, maps:get(write_segment, Meta, 1)),
|
|
LastSeq = last_seq(Segments, maps:get(write_seq, Meta, 0)),
|
|
AckedSeq = min(maps:get(acked_seq, Meta, 0), LastSeq),
|
|
AckedSegment = segment_id_for_seq(Segments, AckedSeq, maps:get(acked_segment, Meta, 1)),
|
|
Outbox0 = #outbox{
|
|
dir = Dir,
|
|
metadata_path = MetadataPath,
|
|
max_records = maps:get(max_records, Opts, 500000),
|
|
max_bytes = maps:get(max_bytes, Opts, 134217728),
|
|
max_segments = maps:get(max_segments, Opts, 10),
|
|
writer = #outbox_writer{
|
|
segment = LastSegmentId,
|
|
next_seq = LastSeq + 1
|
|
},
|
|
checkpoint = #outbox_checkpoint{
|
|
acked_segment = AckedSegment,
|
|
acked_seq = AckedSeq
|
|
},
|
|
segments = Segments
|
|
},
|
|
{ReadSegment, ReadOffset} = locate_reader(Outbox0),
|
|
Outbox1 = Outbox0#outbox{
|
|
reader = #outbox_reader{
|
|
segment = ReadSegment,
|
|
offset = ReadOffset
|
|
}
|
|
},
|
|
maybe
|
|
{ok, Outbox} ?= open_fds(Outbox1),
|
|
ok ?= persist_metadata(Outbox),
|
|
{ok, Outbox}
|
|
else
|
|
{error, Reason} ->
|
|
{error, Reason}
|
|
end.
|
|
|
|
-spec close(outbox()) -> ok.
|
|
close(Outbox = #outbox{}) ->
|
|
_ = close_open_fds(Outbox),
|
|
ok.
|
|
|
|
-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} ->
|
|
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),
|
|
ok ?= append_record(Writer1#outbox_writer.fd, Seq, Payload),
|
|
Segments = update_written_segment(
|
|
NOutbox1#outbox.segments, SegmentId, SegmentPath, Seq, RecordBytes),
|
|
NOutbox2 = NOutbox1#outbox{
|
|
writer = Writer1#outbox_writer{
|
|
segment = SegmentId,
|
|
next_seq = Seq + 1
|
|
},
|
|
segments = Segments
|
|
},
|
|
{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()}.
|
|
ack(Seq, Outbox = #outbox{checkpoint = Checkpoint}) when is_integer(Seq), Seq > 0 ->
|
|
AckedSeq = Checkpoint#outbox_checkpoint.acked_seq,
|
|
case Seq =< AckedSeq of
|
|
true ->
|
|
{ok, Outbox};
|
|
false ->
|
|
PendingAcks0 = ordsets:add_element(Seq, Checkpoint#outbox_checkpoint.pending_acks),
|
|
{NackedSeq, PendingAcks} = advance_acked_seq(AckedSeq, PendingAcks0),
|
|
NackedSegment = segment_id_for_seq(
|
|
Outbox#outbox.segments,
|
|
NackedSeq,
|
|
Checkpoint#outbox_checkpoint.acked_segment),
|
|
Outbox0 = Outbox#outbox{
|
|
checkpoint = Checkpoint#outbox_checkpoint{
|
|
acked_seq = NackedSeq,
|
|
acked_segment = NackedSegment,
|
|
pending_acks = PendingAcks
|
|
}
|
|
},
|
|
maybe
|
|
{ok, Outbox1} ?= prune_acked_segments(Outbox0),
|
|
ok ?= maybe_persist_metadata(NackedSeq =/= AckedSeq, Outbox1),
|
|
{ok, Outbox1}
|
|
else
|
|
{error, Reason} ->
|
|
{error, Reason}
|
|
end
|
|
end.
|
|
|
|
-spec next(outbox()) ->
|
|
eof | {ok, Seq :: pos_integer(), Payload :: binary(), outbox()} | {error, term()}.
|
|
next(#outbox{reader = #outbox_reader{segment = undefined}}) ->
|
|
eof;
|
|
next(Outbox = #outbox{
|
|
reader = Reader,
|
|
segments = Segments
|
|
}) ->
|
|
SegmentId = Reader#outbox_reader.segment,
|
|
Offset = Reader#outbox_reader.offset,
|
|
case ensure_read_fd(Outbox) of
|
|
{ok, NOutbox0} ->
|
|
Reader0 = NOutbox0#outbox.reader,
|
|
case read_next_record(Reader0#outbox_reader.fd, Offset) of
|
|
eof ->
|
|
case next_segment_id(Segments, SegmentId) of
|
|
undefined ->
|
|
eof;
|
|
NextSegmentId ->
|
|
case switch_read_segment(NextSegmentId, 0, NOutbox0) of
|
|
{ok, NOutbox1} ->
|
|
next(NOutbox1);
|
|
{error, Reason} ->
|
|
{error, Reason}
|
|
end
|
|
end;
|
|
{ok, _RecordOffset, Seq, Payload, NextOffset} ->
|
|
{ok, Seq, Payload, NOutbox0#outbox{
|
|
reader = Reader0#outbox_reader{offset = NextOffset}
|
|
}};
|
|
{error, Reason} ->
|
|
{error, Reason}
|
|
end;
|
|
{error, Reason} ->
|
|
{error, Reason}
|
|
end.
|
|
|
|
-spec reset_reader(outbox()) -> {ok, outbox()}.
|
|
reset_reader(Outbox = #outbox{}) ->
|
|
{ReadSegment, ReadOffset} = locate_reader(Outbox),
|
|
switch_read_segment(ReadSegment, ReadOffset, Outbox).
|
|
|
|
-spec stat(outbox()) -> map().
|
|
stat(#outbox{
|
|
dir = Dir,
|
|
writer = #outbox_writer{
|
|
segment = WriteSegment,
|
|
next_seq = NextSeq
|
|
},
|
|
checkpoint = #outbox_checkpoint{
|
|
acked_segment = AckedSegment,
|
|
acked_seq = AckedSeq
|
|
},
|
|
segments = Segments
|
|
}) ->
|
|
#{
|
|
dir => Dir,
|
|
write_segment => WriteSegment,
|
|
write_seq => max(NextSeq - 1, 0),
|
|
acked_segment => AckedSegment,
|
|
acked_seq => AckedSeq,
|
|
segment_num => length(Segments),
|
|
pending_segment_num => pending_segment_count(Segments, AckedSeq)
|
|
}.
|
|
|
|
%%%===================================================================
|
|
%%% Internal functions
|
|
%%%===================================================================
|
|
|
|
-spec ensure_dir(file:filename_all()) -> ok.
|
|
ensure_dir(Dir) ->
|
|
ok = filelib:ensure_dir(filename:join(Dir, "dummy")),
|
|
ok.
|
|
|
|
-spec read_metadata(file:filename_all()) -> map().
|
|
read_metadata(Path) ->
|
|
case file:read_file(Path) of
|
|
{ok, Bin} ->
|
|
case catch jiffy:decode(Bin, [return_maps]) of
|
|
#{<<"write_segment">> := WriteSegment,
|
|
<<"write_seq">> := WriteSeq,
|
|
<<"acked_segment">> := AckedSegment,
|
|
<<"acked_seq">> := AckedSeq}
|
|
when is_integer(WriteSegment), is_integer(WriteSeq),
|
|
is_integer(AckedSegment), is_integer(AckedSeq) ->
|
|
#{
|
|
write_segment => WriteSegment,
|
|
write_seq => WriteSeq,
|
|
acked_segment => AckedSegment,
|
|
acked_seq => AckedSeq
|
|
};
|
|
_ ->
|
|
default_metadata()
|
|
end;
|
|
{error, enoent} ->
|
|
default_metadata();
|
|
{error, _} ->
|
|
default_metadata()
|
|
end.
|
|
|
|
-spec default_metadata() -> map().
|
|
default_metadata() ->
|
|
#{
|
|
write_segment => 1,
|
|
write_seq => 0,
|
|
acked_segment => 1,
|
|
acked_seq => 0
|
|
}.
|
|
|
|
-spec persist_metadata(outbox()) -> ok | {error, term()}.
|
|
persist_metadata(#outbox{
|
|
metadata_path = MetadataPath,
|
|
writer = #outbox_writer{
|
|
segment = WriteSegment,
|
|
next_seq = NextSeq
|
|
},
|
|
checkpoint = #outbox_checkpoint{
|
|
acked_segment = AckedSegment,
|
|
acked_seq = AckedSeq
|
|
}
|
|
}) ->
|
|
TmpPath = MetadataPath ++ ".tmp",
|
|
Metadata = #{
|
|
<<"write_segment">> => WriteSegment,
|
|
<<"write_seq">> => max(NextSeq - 1, 0),
|
|
<<"acked_segment">> => AckedSegment,
|
|
<<"acked_seq">> => AckedSeq
|
|
},
|
|
Bin = iolist_to_binary(jiffy:encode(Metadata, [force_utf8])),
|
|
ok = file:write_file(TmpPath, Bin),
|
|
ok = file:rename(TmpPath, MetadataPath),
|
|
ok.
|
|
|
|
-spec load_segments(file:filename_all()) -> [#segment{}].
|
|
load_segments(Dir) ->
|
|
Pattern = filename:join(Dir, "*" ++ ?SEGMENT_EXT),
|
|
SegmentFiles = lists:sort(filelib:wildcard(Pattern)),
|
|
lists:filtermap(fun(Path) ->
|
|
case parse_segment_id(filename:basename(Path)) of
|
|
{ok, Id} ->
|
|
case scan_segment(Path, Id) of
|
|
{ok, Segment = #segment{records = Records}} when Records > 0 ->
|
|
{true, Segment};
|
|
_ ->
|
|
false
|
|
end;
|
|
error ->
|
|
false
|
|
end
|
|
end, SegmentFiles).
|
|
|
|
-spec parse_segment_id(string()) -> {ok, pos_integer()} | error.
|
|
parse_segment_id(Filename) ->
|
|
case filename:extension(Filename) of
|
|
?SEGMENT_EXT ->
|
|
Base = filename:rootname(Filename, ?SEGMENT_EXT),
|
|
try
|
|
Id = list_to_integer(Base),
|
|
case Id > 0 of
|
|
true ->
|
|
{ok, Id};
|
|
false ->
|
|
error
|
|
end
|
|
catch
|
|
_:_ ->
|
|
error
|
|
end;
|
|
_ ->
|
|
error
|
|
end.
|
|
|
|
-spec scan_segment(file:filename_all(), pos_integer()) -> {ok, #segment{}} | {error, term()}.
|
|
scan_segment(Path, Id) ->
|
|
case file:open(Path, [read, raw, binary]) of
|
|
{ok, Fd} ->
|
|
try
|
|
scan_segment_records(Fd, 0, #segment{id = Id, path = Path})
|
|
after
|
|
ok = file:close(Fd)
|
|
end;
|
|
{error, Reason} ->
|
|
{error, Reason}
|
|
end.
|
|
|
|
-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, RecordOffset, Seq, _Payload, NextOffset} ->
|
|
RecordBytes = NextOffset - RecordOffset,
|
|
NextSegment = Segment#segment{
|
|
start_seq = case Segment#segment.records of
|
|
0 -> Seq;
|
|
_ -> Segment#segment.start_seq
|
|
end,
|
|
end_seq = Seq,
|
|
records = Segment#segment.records + 1,
|
|
bytes = Segment#segment.bytes + RecordBytes
|
|
},
|
|
scan_segment_records(Fd, NextOffset, NextSegment);
|
|
{error, Reason} ->
|
|
{error, Reason}
|
|
end.
|
|
|
|
-spec encoded_record_size(binary()) -> pos_integer().
|
|
encoded_record_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()}.
|
|
ensure_writable_segment(RecordBytes, Outbox = #outbox{
|
|
max_records = MaxRecords,
|
|
max_bytes = MaxBytes,
|
|
max_segments = MaxSegments,
|
|
segments = Segments,
|
|
writer = #outbox_writer{segment = WriteSegment},
|
|
checkpoint = #outbox_checkpoint{acked_seq = AckedSeq}
|
|
}) ->
|
|
case find_segment(Segments, WriteSegment) of
|
|
undefined ->
|
|
{ok, WriteSegment, Outbox};
|
|
#segment{records = Records, bytes = Bytes} ->
|
|
NeedRoll = Records > 0 andalso (Records >= MaxRecords orelse Bytes + RecordBytes > MaxBytes),
|
|
case NeedRoll of
|
|
false ->
|
|
{ok, WriteSegment, Outbox};
|
|
true ->
|
|
case pending_segment_count(Segments, AckedSeq) >= MaxSegments of
|
|
true ->
|
|
{drop, Outbox};
|
|
false ->
|
|
{ok, WriteSegment + 1, Outbox#outbox{
|
|
writer = ((Outbox#outbox.writer)#outbox_writer{
|
|
segment = WriteSegment + 1
|
|
})
|
|
}}
|
|
end
|
|
end
|
|
end.
|
|
|
|
-spec append_record(file:fd(), pos_integer(), binary()) -> ok | {error, term()}.
|
|
append_record(Fd, Seq, Payload) when Fd =/= undefined ->
|
|
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;
|
|
{error, Reason} ->
|
|
{error, Reason}
|
|
end.
|
|
|
|
-spec update_written_segment([#segment{}], pos_integer(), file:filename_all(), pos_integer(), pos_integer()) -> [#segment{}].
|
|
update_written_segment(Segments, SegmentId, Path, Seq, RecordBytes) ->
|
|
case take_segment(SegmentId, Segments) of
|
|
{undefined, Rest} ->
|
|
lists:sort(fun compare_segment/2, [
|
|
#segment{
|
|
id = SegmentId,
|
|
path = Path,
|
|
start_seq = Seq,
|
|
end_seq = Seq,
|
|
records = 1,
|
|
bytes = RecordBytes
|
|
} | Rest
|
|
]);
|
|
{Segment, Rest} ->
|
|
lists:sort(fun compare_segment/2, [
|
|
Segment#segment{
|
|
end_seq = Seq,
|
|
records = Segment#segment.records + 1,
|
|
bytes = Segment#segment.bytes + RecordBytes
|
|
} | Rest
|
|
])
|
|
end.
|
|
|
|
-spec compare_segment(#segment{}, #segment{}) -> boolean().
|
|
compare_segment(#segment{id = Id0}, #segment{id = Id1}) ->
|
|
Id0 =< Id1.
|
|
|
|
-spec find_segment([#segment{}], pos_integer()) -> undefined | #segment{}.
|
|
find_segment(Segments, SegmentId) ->
|
|
lists:keyfind(SegmentId, #segment.id, Segments).
|
|
|
|
-spec take_segment(pos_integer(), [#segment{}]) -> {undefined | #segment{}, [#segment{}]}.
|
|
take_segment(SegmentId, Segments) ->
|
|
case lists:partition(fun(#segment{id = Id}) -> Id =:= SegmentId end, Segments) of
|
|
{[Segment], Rest} ->
|
|
{Segment, Rest};
|
|
{[], Rest} ->
|
|
{undefined, Rest}
|
|
end.
|
|
|
|
-spec last_segment_id([#segment{}], pos_integer()) -> pos_integer().
|
|
last_segment_id([], Default) ->
|
|
Default;
|
|
last_segment_id(Segments, _Default) ->
|
|
(lists:last(Segments))#segment.id.
|
|
|
|
-spec last_seq([#segment{}], non_neg_integer()) -> non_neg_integer().
|
|
last_seq([], Default) ->
|
|
Default;
|
|
last_seq(Segments, _Default) ->
|
|
(lists:last(Segments))#segment.end_seq.
|
|
|
|
-spec segment_id_for_seq([#segment{}], non_neg_integer(), pos_integer()) -> pos_integer().
|
|
segment_id_for_seq(_Segments, 0, Default) ->
|
|
Default;
|
|
segment_id_for_seq(Segments, Seq, Default) ->
|
|
case lists:dropwhile(fun(#segment{end_seq = EndSeq}) -> EndSeq < Seq end, Segments) of
|
|
[#segment{id = SegmentId, start_seq = StartSeq, end_seq = EndSeq} | _]
|
|
when Seq >= StartSeq, Seq =< EndSeq ->
|
|
SegmentId;
|
|
_ ->
|
|
Default
|
|
end.
|
|
|
|
-spec pending_segment_count([#segment{}], non_neg_integer()) -> non_neg_integer().
|
|
pending_segment_count(Segments, AckedSeq) ->
|
|
length([Segment || Segment = #segment{end_seq = EndSeq} <- Segments, EndSeq > AckedSeq]).
|
|
|
|
-spec advance_acked_seq(non_neg_integer(), [pos_integer()]) -> {non_neg_integer(), [pos_integer()]}.
|
|
advance_acked_seq(AckedSeq, PendingAcks) ->
|
|
NextSeq = AckedSeq + 1,
|
|
case ordsets:is_element(NextSeq, PendingAcks) of
|
|
true ->
|
|
advance_acked_seq(NextSeq, ordsets:del_element(NextSeq, PendingAcks));
|
|
false ->
|
|
{AckedSeq, PendingAcks}
|
|
end.
|
|
|
|
-spec prune_acked_segments(outbox()) -> {ok, outbox()} | {error, term()}.
|
|
prune_acked_segments(Outbox = #outbox{
|
|
segments = Segments,
|
|
checkpoint = #outbox_checkpoint{acked_seq = AckedSeq}
|
|
}) ->
|
|
{DeleteSegments, KeepSegments} = lists:partition(fun(#segment{end_seq = EndSeq}) ->
|
|
EndSeq =< AckedSeq
|
|
end, Segments),
|
|
Outbox0 = detach_deleted_fds(DeleteSegments, Outbox#outbox{segments = KeepSegments}),
|
|
ok = lists:foldl(fun(#segment{path = Path}, ok) ->
|
|
case file:delete(Path) of
|
|
ok ->
|
|
ok;
|
|
{error, enoent} ->
|
|
ok;
|
|
{error, Reason} ->
|
|
throw({delete_segment_failed, Path, Reason})
|
|
end
|
|
end, ok, DeleteSegments),
|
|
{ok, reset_reader_after_prune(Outbox0)}.
|
|
|
|
-spec reset_reader_after_prune(outbox()) -> outbox().
|
|
reset_reader_after_prune(Outbox = #outbox{
|
|
reader = #outbox_reader{segment = ReadSegment},
|
|
segments = Segments
|
|
}) ->
|
|
case ReadSegment =:= undefined orelse lists:keymember(ReadSegment, #segment.id, Segments) of
|
|
true ->
|
|
Outbox;
|
|
false ->
|
|
case reset_reader(Outbox) of
|
|
{ok, NOutbox} ->
|
|
NOutbox;
|
|
{error, _} ->
|
|
Outbox#outbox{
|
|
reader = #outbox_reader{
|
|
segment = undefined,
|
|
offset = 0,
|
|
fd = undefined
|
|
}
|
|
}
|
|
end
|
|
end.
|
|
|
|
-spec locate_reader(outbox()) -> {undefined | pos_integer(), non_neg_integer()}.
|
|
locate_reader(#outbox{
|
|
segments = Segments,
|
|
checkpoint = #outbox_checkpoint{acked_seq = AckedSeq}
|
|
}) ->
|
|
NextSeq = AckedSeq + 1,
|
|
case lists:dropwhile(fun(#segment{end_seq = EndSeq}) -> EndSeq < NextSeq end, Segments) of
|
|
[] ->
|
|
{undefined, 0};
|
|
[#segment{id = SegmentId, start_seq = StartSeq} = Segment | _] ->
|
|
case NextSeq =< StartSeq of
|
|
true ->
|
|
{SegmentId, 0};
|
|
false ->
|
|
case locate_offset_for_seq(Segment, NextSeq) of
|
|
{ok, Offset} ->
|
|
{SegmentId, Offset};
|
|
{error, _} ->
|
|
{SegmentId, 0}
|
|
end
|
|
end
|
|
end.
|
|
|
|
-spec locate_offset_for_seq(#segment{}, pos_integer()) -> {ok, non_neg_integer()} | {error, term()}.
|
|
locate_offset_for_seq(#segment{path = Path}, TargetSeq) ->
|
|
case file:open(Path, [read, raw, binary]) of
|
|
{ok, Fd} ->
|
|
try
|
|
locate_offset_loop(Fd, TargetSeq, 0)
|
|
after
|
|
ok = file:close(Fd)
|
|
end;
|
|
{error, Reason} ->
|
|
{error, Reason}
|
|
end.
|
|
|
|
-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_next_record(Fd, Offset) of
|
|
eof ->
|
|
{ok, Offset};
|
|
{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_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, 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 ->
|
|
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, 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
|
|
[#segment{id = NextSegmentId} | _] ->
|
|
NextSegmentId;
|
|
[] ->
|
|
undefined
|
|
end.
|
|
|
|
-spec segment_path(file:filename_all(), pos_integer()) -> file:filename_all().
|
|
segment_path(Dir, SegmentId) ->
|
|
filename:join(Dir, io_lib:format("~6..0B~s", [SegmentId, ?SEGMENT_EXT])).
|
|
|
|
-spec open_fds(outbox()) -> {ok, outbox()} | {error, term()}.
|
|
open_fds(Outbox) ->
|
|
case maybe_open_write_fd(Outbox) of
|
|
{ok, Outbox0} ->
|
|
maybe
|
|
{ok, NOutbox} ?= maybe_open_read_fd(Outbox0),
|
|
{ok, NOutbox}
|
|
else
|
|
{error, Reason} ->
|
|
_ = close_open_fds(Outbox0),
|
|
{error, Reason}
|
|
end;
|
|
{error, Reason} ->
|
|
{error, Reason}
|
|
end.
|
|
|
|
-spec maybe_open_write_fd(outbox()) -> {ok, outbox()} | {error, term()}.
|
|
maybe_open_write_fd(Outbox = #outbox{
|
|
segments = Segments,
|
|
writer = #outbox_writer{
|
|
segment = WriteSegment,
|
|
fd = undefined
|
|
}
|
|
}) ->
|
|
case find_segment(Segments, WriteSegment) of
|
|
undefined ->
|
|
{ok, Outbox};
|
|
_ ->
|
|
ensure_write_fd(WriteSegment, Outbox)
|
|
end;
|
|
maybe_open_write_fd(Outbox = #outbox{}) ->
|
|
{ok, Outbox}.
|
|
|
|
-spec maybe_open_read_fd(outbox()) -> {ok, outbox()} | {error, term()}.
|
|
maybe_open_read_fd(Outbox = #outbox{reader = #outbox_reader{segment = undefined}}) ->
|
|
{ok, Outbox};
|
|
maybe_open_read_fd(Outbox = #outbox{reader = #outbox_reader{fd = undefined}}) ->
|
|
ensure_read_fd(Outbox);
|
|
maybe_open_read_fd(Outbox = #outbox{}) ->
|
|
{ok, Outbox}.
|
|
|
|
-spec ensure_write_fd(pos_integer(), outbox()) -> {ok, outbox()} | {error, term()}.
|
|
ensure_write_fd(SegmentId, Outbox = #outbox{
|
|
writer = #outbox_writer{
|
|
segment = SegmentId,
|
|
fd = Fd
|
|
}
|
|
})
|
|
when Fd =/= undefined ->
|
|
{ok, Outbox};
|
|
ensure_write_fd(SegmentId, Outbox = #outbox{dir = Dir}) ->
|
|
Writer0 = (Outbox#outbox.writer)#outbox_writer{segment = SegmentId},
|
|
Outbox0 = close_write_fd(Outbox#outbox{writer = Writer0}),
|
|
Path = segment_path(Dir, SegmentId),
|
|
case file:open(Path, [append, raw, binary]) of
|
|
{ok, Fd} ->
|
|
Writer1 = (Outbox0#outbox.writer)#outbox_writer{fd = Fd},
|
|
{ok, Outbox0#outbox{writer = Writer1}};
|
|
{error, Reason} ->
|
|
{error, Reason}
|
|
end.
|
|
|
|
-spec ensure_read_fd(outbox()) -> {ok, outbox()} | {error, term()}.
|
|
ensure_read_fd(Outbox = #outbox{reader = #outbox_reader{segment = undefined}}) ->
|
|
{ok, Outbox};
|
|
ensure_read_fd(Outbox = #outbox{reader = #outbox_reader{fd = Fd}}) when Fd =/= undefined ->
|
|
{ok, Outbox};
|
|
ensure_read_fd(Outbox = #outbox{
|
|
reader = Reader,
|
|
dir = Dir
|
|
}) ->
|
|
SegmentId = Reader#outbox_reader.segment,
|
|
Offset = Reader#outbox_reader.offset,
|
|
Path = segment_path(Dir, SegmentId),
|
|
case open_read_fd(Path, Offset) of
|
|
{ok, Fd} ->
|
|
{ok, Outbox#outbox{
|
|
reader = Reader#outbox_reader{fd = Fd}
|
|
}};
|
|
{error, Reason} ->
|
|
{error, Reason}
|
|
end.
|
|
|
|
-spec switch_read_segment(undefined | pos_integer(), non_neg_integer(), outbox()) ->
|
|
{ok, outbox()} | {error, term()}.
|
|
switch_read_segment(undefined, _Offset, Outbox) ->
|
|
Reader0 = #outbox_reader{segment = undefined, offset = 0},
|
|
{ok, close_read_fd(Outbox#outbox{reader = Reader0})};
|
|
switch_read_segment(SegmentId, Offset, Outbox = #outbox{reader = Reader}) ->
|
|
ensure_read_fd(close_read_fd(Outbox#outbox{
|
|
reader = Reader#outbox_reader{
|
|
segment = SegmentId,
|
|
offset = Offset
|
|
}
|
|
})).
|
|
|
|
-spec maybe_reset_reader_after_append(outbox(), outbox()) -> {ok, outbox()} | {error, term()}.
|
|
maybe_reset_reader_after_append(#outbox{reader = #outbox_reader{segment = undefined}}, Outbox) ->
|
|
reset_reader(Outbox);
|
|
maybe_reset_reader_after_append(_PrevOutbox, Outbox) ->
|
|
{ok, Outbox}.
|
|
|
|
-spec detach_deleted_fds([#segment{}], outbox()) -> outbox().
|
|
detach_deleted_fds(DeleteSegments, Outbox = #outbox{
|
|
reader = #outbox_reader{segment = ReadSegment},
|
|
writer = #outbox_writer{segment = WriteSegment}
|
|
}) ->
|
|
DeleteIds = [SegmentId || #segment{id = SegmentId} <- DeleteSegments],
|
|
Outbox0 = case lists:member(ReadSegment, DeleteIds) of
|
|
true ->
|
|
close_read_fd(Outbox);
|
|
false ->
|
|
Outbox
|
|
end,
|
|
case lists:member(WriteSegment, DeleteIds) of
|
|
true ->
|
|
close_write_fd(Outbox0);
|
|
false ->
|
|
Outbox0
|
|
end.
|
|
|
|
-spec close_open_fds(outbox()) -> outbox().
|
|
close_open_fds(Outbox) ->
|
|
close_write_fd(close_read_fd(Outbox)).
|
|
|
|
-spec close_read_fd(outbox()) -> outbox().
|
|
close_read_fd(Outbox = #outbox{reader = #outbox_reader{fd = undefined}}) ->
|
|
Outbox;
|
|
close_read_fd(Outbox = #outbox{reader = Reader}) ->
|
|
Fd = Reader#outbox_reader.fd,
|
|
_ = maybe_close_fd(Fd),
|
|
Outbox#outbox{reader = Reader#outbox_reader{fd = undefined}}.
|
|
|
|
-spec close_write_fd(outbox()) -> outbox().
|
|
close_write_fd(Outbox = #outbox{writer = #outbox_writer{fd = undefined}}) ->
|
|
Outbox;
|
|
close_write_fd(Outbox = #outbox{writer = Writer}) ->
|
|
Fd = Writer#outbox_writer.fd,
|
|
_ = maybe_close_fd(Fd),
|
|
Outbox#outbox{writer = Writer#outbox_writer{fd = undefined}}.
|
|
|
|
-spec maybe_close_fd(file:fd()) -> ok.
|
|
maybe_close_fd(Fd) ->
|
|
case file:close(Fd) of
|
|
ok ->
|
|
ok;
|
|
{error, badarg} ->
|
|
ok;
|
|
{error, terminated} ->
|
|
ok;
|
|
{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.
|