增加外部存储缓存策略
This commit is contained in:
parent
45a2fe5300
commit
7dd67330fb
584
src/outbox/endpoint_outbox.erl
Normal file
584
src/outbox/endpoint_outbox.erl
Normal file
@ -0,0 +1,584 @@
|
||||
%%%-------------------------------------------------------------------
|
||||
%%% @author Codex
|
||||
%%% @doc
|
||||
%%% endpoint 持久化 outbox 的第一阶段实现:
|
||||
%%% - 一个 endpoint 一个目录
|
||||
%%% - segment 按递增编号命名
|
||||
%%% - metadata 只记录连续 ack 指针
|
||||
%%% - inflight 不落盘,崩溃后从 acked_seq + 1 重放
|
||||
%%%-------------------------------------------------------------------
|
||||
-module(endpoint_outbox).
|
||||
|
||||
-export([open/2, 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, {
|
||||
endpoint_id :: integer() | binary() | atom(),
|
||||
dir :: file:filename_all(),
|
||||
metadata_path :: file:filename_all(),
|
||||
max_records = 500000 :: pos_integer(),
|
||||
max_bytes = 134217728 :: pos_integer(),
|
||||
max_segments = 10 :: pos_integer(),
|
||||
write_segment = 1 :: pos_integer(),
|
||||
next_seq = 1 :: pos_integer(),
|
||||
acked_segment = 1 :: pos_integer(),
|
||||
acked_seq = 0 :: non_neg_integer(),
|
||||
segments = [] :: [#segment{}],
|
||||
pending_acks = [] :: [pos_integer()],
|
||||
read_segment :: undefined | pos_integer(),
|
||||
read_offset = 0 :: non_neg_integer()
|
||||
}).
|
||||
|
||||
-type outbox() :: #outbox{}.
|
||||
|
||||
-define(DEFAULT_ROOT_DIR, "var/outbox").
|
||||
-define(METADATA_FILE, "metadata.json").
|
||||
-define(SEGMENT_EXT, ".log").
|
||||
|
||||
%%%===================================================================
|
||||
%%% 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)),
|
||||
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{
|
||||
endpoint_id = EndpointId,
|
||||
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),
|
||||
write_segment = LastSegmentId,
|
||||
next_seq = LastSeq + 1,
|
||||
acked_segment = AckedSegment,
|
||||
acked_seq = AckedSeq,
|
||||
segments = Segments
|
||||
},
|
||||
{ReadSegment, ReadOffset} = locate_reader(Outbox0),
|
||||
Outbox = Outbox0#outbox{read_segment = ReadSegment, read_offset = ReadOffset},
|
||||
ok = persist_metadata(Outbox),
|
||||
{ok, Outbox}.
|
||||
|
||||
-spec append(binary(), outbox()) ->
|
||||
{ok, Seq :: pos_integer(), outbox()} | {dropped, capacity_reached, outbox()} | {error, term()}.
|
||||
append(Payload, Outbox = #outbox{}) when is_binary(Payload) ->
|
||||
RecordBytes = encoded_record_size(Payload),
|
||||
case ensure_writable_segment(RecordBytes, Outbox) of
|
||||
{drop, NOutbox} ->
|
||||
{dropped, capacity_reached, NOutbox};
|
||||
{ok, SegmentId, NOutbox0} ->
|
||||
Seq = NOutbox0#outbox.next_seq,
|
||||
SegmentPath = segment_path(NOutbox0#outbox.dir, SegmentId),
|
||||
case append_record(SegmentPath, Seq, Payload) of
|
||||
ok ->
|
||||
Segments = update_written_segment(
|
||||
NOutbox0#outbox.segments, SegmentId, SegmentPath, Seq, RecordBytes),
|
||||
NOutbox = NOutbox0#outbox{
|
||||
write_segment = SegmentId,
|
||||
next_seq = Seq + 1,
|
||||
segments = Segments
|
||||
},
|
||||
ok = persist_metadata(NOutbox),
|
||||
{ok, Seq, NOutbox};
|
||||
{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 > 0 ->
|
||||
case Seq =< AckedSeq of
|
||||
true ->
|
||||
{ok, Outbox};
|
||||
false ->
|
||||
PendingAcks0 = ordsets:add_element(Seq, Outbox#outbox.pending_acks),
|
||||
{NackedSeq, PendingAcks} = advance_acked_seq(AckedSeq, PendingAcks0),
|
||||
NackedSegment = segment_id_for_seq(Outbox#outbox.segments, NackedSeq, Outbox#outbox.acked_segment),
|
||||
Outbox0 = Outbox#outbox{
|
||||
acked_seq = NackedSeq,
|
||||
acked_segment = NackedSegment,
|
||||
pending_acks = PendingAcks
|
||||
},
|
||||
{ok, Outbox1} = prune_acked_segments(Outbox0),
|
||||
case NackedSeq =:= AckedSeq of
|
||||
true ->
|
||||
{ok, Outbox1};
|
||||
false ->
|
||||
ok = persist_metadata(Outbox1),
|
||||
{ok, Outbox1}
|
||||
end
|
||||
end.
|
||||
|
||||
-spec next(outbox()) ->
|
||||
eof | {ok, Seq :: pos_integer(), Payload :: binary(), outbox()} | {error, term()}.
|
||||
next(#outbox{read_segment = undefined}) ->
|
||||
eof;
|
||||
next(Outbox = #outbox{read_segment = SegmentId, read_offset = Offset, dir = Dir, segments = Segments}) ->
|
||||
SegmentPath = segment_path(Dir, SegmentId),
|
||||
case read_record_at(SegmentPath, Offset) of
|
||||
eof ->
|
||||
case next_segment_id(Segments, SegmentId) of
|
||||
undefined ->
|
||||
eof;
|
||||
NextSegmentId ->
|
||||
next(Outbox#outbox{read_segment = NextSegmentId, read_offset = 0})
|
||||
end;
|
||||
{ok, Seq, Payload, NextOffset} ->
|
||||
{ok, Seq, Payload, Outbox#outbox{read_offset = NextOffset}};
|
||||
{error, Reason} ->
|
||||
{error, Reason}
|
||||
end.
|
||||
|
||||
-spec reset_reader(outbox()) -> {ok, outbox()}.
|
||||
reset_reader(Outbox = #outbox{}) ->
|
||||
{ReadSegment, ReadOffset} = locate_reader(Outbox),
|
||||
{ok, Outbox#outbox{read_segment = ReadSegment, read_offset = ReadOffset}}.
|
||||
|
||||
-spec stat(outbox()) -> map().
|
||||
stat(#outbox{
|
||||
dir = Dir,
|
||||
write_segment = WriteSegment,
|
||||
next_seq = NextSeq,
|
||||
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 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")),
|
||||
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,
|
||||
write_segment = WriteSegment,
|
||||
next_seq = NextSeq,
|
||||
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, #segment{id = Id, path = Path})
|
||||
after
|
||||
ok = file:close(Fd)
|
||||
end;
|
||||
{error, Reason} ->
|
||||
{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
|
||||
eof ->
|
||||
{ok, Segment};
|
||||
{ok, Seq, _Payload, RecordBytes} ->
|
||||
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, NextSegment);
|
||||
{error, Reason} ->
|
||||
{error, Reason}
|
||||
end.
|
||||
|
||||
-spec encoded_record_size(binary()) -> pos_integer().
|
||||
encoded_record_size(Payload) ->
|
||||
12 + byte_size(Payload).
|
||||
|
||||
-spec ensure_writable_segment(pos_integer(), outbox()) ->
|
||||
{ok, pos_integer(), outbox()} | {drop, outbox()}.
|
||||
ensure_writable_segment(RecordBytes, Outbox = #outbox{
|
||||
write_segment = WriteSegment,
|
||||
max_records = MaxRecords,
|
||||
max_bytes = MaxBytes,
|
||||
max_segments = MaxSegments,
|
||||
segments = Segments,
|
||||
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{write_segment = WriteSegment + 1}}
|
||||
end
|
||||
end
|
||||
end.
|
||||
|
||||
-spec append_record(file:filename_all(), pos_integer(), binary()) -> ok | {error, term()}.
|
||||
append_record(Path, Seq, Payload) ->
|
||||
Bin = <<Seq:64/unsigned-big-integer, (byte_size(Payload)):32/unsigned-big-integer, Payload/binary>>,
|
||||
case file:open(Path, [append, raw, binary]) of
|
||||
{ok, Fd} ->
|
||||
try
|
||||
ok = file:write(Fd, Bin)
|
||||
after
|
||||
ok = file:close(Fd)
|
||||
end;
|
||||
{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, acked_seq = AckedSeq}) ->
|
||||
{DeleteSegments, KeepSegments} = lists:partition(fun(#segment{end_seq = EndSeq}) ->
|
||||
EndSeq =< AckedSeq
|
||||
end, Segments),
|
||||
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(Outbox#outbox{segments = KeepSegments})}.
|
||||
|
||||
-spec reset_reader_after_prune(outbox()) -> outbox().
|
||||
reset_reader_after_prune(Outbox = #outbox{read_segment = ReadSegment, segments = Segments}) ->
|
||||
case ReadSegment =:= undefined orelse lists:keymember(ReadSegment, #segment.id, Segments) of
|
||||
true ->
|
||||
Outbox;
|
||||
false ->
|
||||
{NReadSegment, NReadOffset} = locate_reader(Outbox),
|
||||
Outbox#outbox{read_segment = NReadSegment, read_offset = NReadOffset}
|
||||
end.
|
||||
|
||||
-spec locate_reader(outbox()) -> {undefined | pos_integer(), non_neg_integer()}.
|
||||
locate_reader(#outbox{segments = Segments, 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_record(Fd) of
|
||||
eof ->
|
||||
{ok, Offset};
|
||||
{ok, Seq, _Payload, _RecordBytes} when Seq >= TargetSeq ->
|
||||
{ok, Offset};
|
||||
{ok, _Seq, _Payload, RecordBytes} ->
|
||||
locate_offset_loop(Fd, TargetSeq, Offset + RecordBytes);
|
||||
{error, Reason} ->
|
||||
{error, Reason}
|
||||
end.
|
||||
|
||||
-spec read_record_at(file:filename_all(), non_neg_integer()) ->
|
||||
eof | {ok, pos_integer(), binary(), non_neg_integer()} | {error, term()}.
|
||||
read_record_at(Path, Offset) ->
|
||||
case file:open(Path, [read, raw, binary]) of
|
||||
{ok, Fd} ->
|
||||
try
|
||||
{ok, _} = file:position(Fd, Offset),
|
||||
case read_record(Fd) of
|
||||
eof ->
|
||||
eof;
|
||||
{ok, Seq, Payload, RecordBytes} ->
|
||||
{ok, Seq, Payload, Offset + RecordBytes};
|
||||
{error, Reason} ->
|
||||
{error, Reason}
|
||||
end
|
||||
after
|
||||
ok = file:close(Fd)
|
||||
end;
|
||||
{error, enoent} ->
|
||||
eof;
|
||||
{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
|
||||
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};
|
||||
eof ->
|
||||
{error, truncated_payload};
|
||||
{error, Reason} ->
|
||||
{error, Reason}
|
||||
end;
|
||||
{ok, _} ->
|
||||
{error, truncated_header};
|
||||
{error, Reason} ->
|
||||
{error, Reason}
|
||||
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])).
|
||||
Loading…
x
Reference in New Issue
Block a user