fix outbox
This commit is contained in:
parent
7dd67330fb
commit
30171d0e42
@ -9,7 +9,7 @@
|
|||||||
%%%-------------------------------------------------------------------
|
%%%-------------------------------------------------------------------
|
||||||
-module(endpoint_outbox).
|
-module(endpoint_outbox).
|
||||||
|
|
||||||
-export([open/2, append/2, ack/2, next/1, reset_reader/1, stat/1]).
|
-export([open/2, close/1, append/2, ack/2, next/1, reset_reader/1, stat/1]).
|
||||||
-export_type([outbox/0]).
|
-export_type([outbox/0]).
|
||||||
|
|
||||||
-record(segment, {
|
-record(segment, {
|
||||||
@ -21,6 +21,24 @@
|
|||||||
bytes = 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, {
|
-record(outbox, {
|
||||||
endpoint_id :: integer() | binary() | atom(),
|
endpoint_id :: integer() | binary() | atom(),
|
||||||
dir :: file:filename_all(),
|
dir :: file:filename_all(),
|
||||||
@ -28,14 +46,10 @@
|
|||||||
max_records = 500000 :: pos_integer(),
|
max_records = 500000 :: pos_integer(),
|
||||||
max_bytes = 134217728 :: pos_integer(),
|
max_bytes = 134217728 :: pos_integer(),
|
||||||
max_segments = 10 :: 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{}],
|
segments = [] :: [#segment{}],
|
||||||
pending_acks = [] :: [pos_integer()],
|
writer = #outbox_writer{} :: #outbox_writer{},
|
||||||
read_segment :: undefined | pos_integer(),
|
reader = #outbox_reader{} :: #outbox_reader{},
|
||||||
read_offset = 0 :: non_neg_integer()
|
checkpoint = #outbox_checkpoint{} :: #outbox_checkpoint{}
|
||||||
}).
|
}).
|
||||||
|
|
||||||
-type outbox() :: #outbox{}.
|
-type outbox() :: #outbox{}.
|
||||||
@ -67,16 +81,35 @@ open(EndpointId, Opts) when is_map(Opts) ->
|
|||||||
max_records = maps:get(max_records, Opts, 500000),
|
max_records = maps:get(max_records, Opts, 500000),
|
||||||
max_bytes = maps:get(max_bytes, Opts, 134217728),
|
max_bytes = maps:get(max_bytes, Opts, 134217728),
|
||||||
max_segments = maps:get(max_segments, Opts, 10),
|
max_segments = maps:get(max_segments, Opts, 10),
|
||||||
write_segment = LastSegmentId,
|
writer = #outbox_writer{
|
||||||
next_seq = LastSeq + 1,
|
segment = LastSegmentId,
|
||||||
|
next_seq = LastSeq + 1
|
||||||
|
},
|
||||||
|
checkpoint = #outbox_checkpoint{
|
||||||
acked_segment = AckedSegment,
|
acked_segment = AckedSegment,
|
||||||
acked_seq = AckedSeq,
|
acked_seq = AckedSeq
|
||||||
|
},
|
||||||
segments = Segments
|
segments = Segments
|
||||||
},
|
},
|
||||||
{ReadSegment, ReadOffset} = locate_reader(Outbox0),
|
{ReadSegment, ReadOffset} = locate_reader(Outbox0),
|
||||||
Outbox = Outbox0#outbox{read_segment = ReadSegment, read_offset = ReadOffset},
|
Outbox1 = Outbox0#outbox{
|
||||||
|
reader = #outbox_reader{
|
||||||
|
segment = ReadSegment,
|
||||||
|
offset = ReadOffset
|
||||||
|
}
|
||||||
|
},
|
||||||
|
case open_fds(Outbox1) of
|
||||||
|
{ok, Outbox} ->
|
||||||
ok = persist_metadata(Outbox),
|
ok = persist_metadata(Outbox),
|
||||||
{ok, Outbox}.
|
{ok, Outbox};
|
||||||
|
{error, Reason} ->
|
||||||
|
{error, Reason}
|
||||||
|
end.
|
||||||
|
|
||||||
|
-spec close(outbox()) -> ok.
|
||||||
|
close(Outbox = #outbox{}) ->
|
||||||
|
_ = close_open_fds(Outbox),
|
||||||
|
ok.
|
||||||
|
|
||||||
-spec append(binary(), outbox()) ->
|
-spec append(binary(), outbox()) ->
|
||||||
{ok, Seq :: pos_integer(), outbox()} | {dropped, capacity_reached, outbox()} | {error, term()}.
|
{ok, Seq :: pos_integer(), outbox()} | {dropped, capacity_reached, outbox()} | {error, term()}.
|
||||||
@ -86,37 +119,56 @@ append(Payload, Outbox = #outbox{}) when is_binary(Payload) ->
|
|||||||
{drop, NOutbox} ->
|
{drop, NOutbox} ->
|
||||||
{dropped, capacity_reached, NOutbox};
|
{dropped, capacity_reached, NOutbox};
|
||||||
{ok, SegmentId, NOutbox0} ->
|
{ok, SegmentId, NOutbox0} ->
|
||||||
Seq = NOutbox0#outbox.next_seq,
|
case ensure_write_fd(SegmentId, NOutbox0) of
|
||||||
SegmentPath = segment_path(NOutbox0#outbox.dir, SegmentId),
|
{ok, NOutbox1} ->
|
||||||
case append_record(SegmentPath, Seq, Payload) of
|
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 ->
|
||||||
Segments = update_written_segment(
|
Segments = update_written_segment(
|
||||||
NOutbox0#outbox.segments, SegmentId, SegmentPath, Seq, RecordBytes),
|
NOutbox1#outbox.segments, SegmentId, SegmentPath, Seq, RecordBytes),
|
||||||
NOutbox = NOutbox0#outbox{
|
NOutbox2 = NOutbox1#outbox{
|
||||||
write_segment = SegmentId,
|
writer = Writer1#outbox_writer{
|
||||||
next_seq = Seq + 1,
|
segment = SegmentId,
|
||||||
|
next_seq = Seq + 1
|
||||||
|
},
|
||||||
segments = Segments
|
segments = Segments
|
||||||
},
|
},
|
||||||
|
case maybe_reset_reader_after_append(Outbox, NOutbox2) of
|
||||||
|
{ok, NOutbox} ->
|
||||||
ok = persist_metadata(NOutbox),
|
ok = persist_metadata(NOutbox),
|
||||||
{ok, Seq, NOutbox};
|
{ok, Seq, NOutbox};
|
||||||
{error, Reason} ->
|
{error, Reason} ->
|
||||||
{error, Reason}
|
{error, Reason}
|
||||||
|
end;
|
||||||
|
{error, Reason} ->
|
||||||
|
{error, Reason}
|
||||||
|
end;
|
||||||
|
{error, Reason} ->
|
||||||
|
{error, Reason}
|
||||||
end
|
end
|
||||||
end.
|
end.
|
||||||
|
|
||||||
-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 > 0 ->
|
ack(Seq, Outbox = #outbox{checkpoint = Checkpoint}) when is_integer(Seq), Seq > 0 ->
|
||||||
|
AckedSeq = Checkpoint#outbox_checkpoint.acked_seq,
|
||||||
case Seq =< AckedSeq of
|
case Seq =< AckedSeq of
|
||||||
true ->
|
true ->
|
||||||
{ok, Outbox};
|
{ok, Outbox};
|
||||||
false ->
|
false ->
|
||||||
PendingAcks0 = ordsets:add_element(Seq, Outbox#outbox.pending_acks),
|
PendingAcks0 = ordsets:add_element(Seq, Checkpoint#outbox_checkpoint.pending_acks),
|
||||||
{NackedSeq, PendingAcks} = advance_acked_seq(AckedSeq, PendingAcks0),
|
{NackedSeq, PendingAcks} = advance_acked_seq(AckedSeq, PendingAcks0),
|
||||||
NackedSegment = segment_id_for_seq(Outbox#outbox.segments, NackedSeq, Outbox#outbox.acked_segment),
|
NackedSegment = segment_id_for_seq(
|
||||||
|
Outbox#outbox.segments,
|
||||||
|
NackedSeq,
|
||||||
|
Checkpoint#outbox_checkpoint.acked_segment),
|
||||||
Outbox0 = Outbox#outbox{
|
Outbox0 = Outbox#outbox{
|
||||||
|
checkpoint = Checkpoint#outbox_checkpoint{
|
||||||
acked_seq = NackedSeq,
|
acked_seq = NackedSeq,
|
||||||
acked_segment = NackedSegment,
|
acked_segment = NackedSegment,
|
||||||
pending_acks = PendingAcks
|
pending_acks = PendingAcks
|
||||||
|
}
|
||||||
},
|
},
|
||||||
{ok, Outbox1} = prune_acked_segments(Outbox0),
|
{ok, Outbox1} = prune_acked_segments(Outbox0),
|
||||||
case NackedSeq =:= AckedSeq of
|
case NackedSeq =:= AckedSeq of
|
||||||
@ -130,20 +182,37 @@ ack(Seq, Outbox = #outbox{acked_seq = AckedSeq}) when is_integer(Seq), Seq > 0 -
|
|||||||
|
|
||||||
-spec next(outbox()) ->
|
-spec next(outbox()) ->
|
||||||
eof | {ok, Seq :: pos_integer(), Payload :: binary(), outbox()} | {error, term()}.
|
eof | {ok, Seq :: pos_integer(), Payload :: binary(), outbox()} | {error, term()}.
|
||||||
next(#outbox{read_segment = undefined}) ->
|
next(#outbox{reader = #outbox_reader{segment = undefined}}) ->
|
||||||
eof;
|
eof;
|
||||||
next(Outbox = #outbox{read_segment = SegmentId, read_offset = Offset, dir = Dir, segments = Segments}) ->
|
next(Outbox = #outbox{
|
||||||
SegmentPath = segment_path(Dir, SegmentId),
|
reader = Reader,
|
||||||
case read_record_at(SegmentPath, Offset) of
|
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_record(Reader0#outbox_reader.fd) of
|
||||||
eof ->
|
eof ->
|
||||||
case next_segment_id(Segments, SegmentId) of
|
case next_segment_id(Segments, SegmentId) of
|
||||||
undefined ->
|
undefined ->
|
||||||
eof;
|
eof;
|
||||||
NextSegmentId ->
|
NextSegmentId ->
|
||||||
next(Outbox#outbox{read_segment = NextSegmentId, read_offset = 0})
|
case switch_read_segment(NextSegmentId, 0, NOutbox0) of
|
||||||
|
{ok, NOutbox1} ->
|
||||||
|
next(NOutbox1);
|
||||||
|
{error, Reason} ->
|
||||||
|
{error, Reason}
|
||||||
|
end
|
||||||
|
end;
|
||||||
|
{ok, Seq, Payload, RecordBytes} ->
|
||||||
|
{ok, Seq, Payload, NOutbox0#outbox{
|
||||||
|
reader = Reader0#outbox_reader{offset = Offset + RecordBytes}
|
||||||
|
}};
|
||||||
|
{error, Reason} ->
|
||||||
|
{error, Reason}
|
||||||
end;
|
end;
|
||||||
{ok, Seq, Payload, NextOffset} ->
|
|
||||||
{ok, Seq, Payload, Outbox#outbox{read_offset = NextOffset}};
|
|
||||||
{error, Reason} ->
|
{error, Reason} ->
|
||||||
{error, Reason}
|
{error, Reason}
|
||||||
end.
|
end.
|
||||||
@ -151,15 +220,19 @@ next(Outbox = #outbox{read_segment = SegmentId, read_offset = Offset, dir = Dir,
|
|||||||
-spec reset_reader(outbox()) -> {ok, outbox()}.
|
-spec reset_reader(outbox()) -> {ok, outbox()}.
|
||||||
reset_reader(Outbox = #outbox{}) ->
|
reset_reader(Outbox = #outbox{}) ->
|
||||||
{ReadSegment, ReadOffset} = locate_reader(Outbox),
|
{ReadSegment, ReadOffset} = locate_reader(Outbox),
|
||||||
{ok, Outbox#outbox{read_segment = ReadSegment, read_offset = ReadOffset}}.
|
switch_read_segment(ReadSegment, ReadOffset, Outbox).
|
||||||
|
|
||||||
-spec stat(outbox()) -> map().
|
-spec stat(outbox()) -> map().
|
||||||
stat(#outbox{
|
stat(#outbox{
|
||||||
dir = Dir,
|
dir = Dir,
|
||||||
write_segment = WriteSegment,
|
writer = #outbox_writer{
|
||||||
next_seq = NextSeq,
|
segment = WriteSegment,
|
||||||
|
next_seq = NextSeq
|
||||||
|
},
|
||||||
|
checkpoint = #outbox_checkpoint{
|
||||||
acked_segment = AckedSegment,
|
acked_segment = AckedSegment,
|
||||||
acked_seq = AckedSeq,
|
acked_seq = AckedSeq
|
||||||
|
},
|
||||||
segments = Segments
|
segments = Segments
|
||||||
}) ->
|
}) ->
|
||||||
#{
|
#{
|
||||||
@ -238,10 +311,14 @@ default_metadata() ->
|
|||||||
-spec persist_metadata(outbox()) -> ok | {error, term()}.
|
-spec persist_metadata(outbox()) -> ok | {error, term()}.
|
||||||
persist_metadata(#outbox{
|
persist_metadata(#outbox{
|
||||||
metadata_path = MetadataPath,
|
metadata_path = MetadataPath,
|
||||||
write_segment = WriteSegment,
|
writer = #outbox_writer{
|
||||||
next_seq = NextSeq,
|
segment = WriteSegment,
|
||||||
|
next_seq = NextSeq
|
||||||
|
},
|
||||||
|
checkpoint = #outbox_checkpoint{
|
||||||
acked_segment = AckedSegment,
|
acked_segment = AckedSegment,
|
||||||
acked_seq = AckedSeq
|
acked_seq = AckedSeq
|
||||||
|
}
|
||||||
}) ->
|
}) ->
|
||||||
TmpPath = MetadataPath ++ ".tmp",
|
TmpPath = MetadataPath ++ ".tmp",
|
||||||
Metadata = #{
|
Metadata = #{
|
||||||
@ -334,12 +411,12 @@ encoded_record_size(Payload) ->
|
|||||||
-spec ensure_writable_segment(pos_integer(), outbox()) ->
|
-spec ensure_writable_segment(pos_integer(), outbox()) ->
|
||||||
{ok, pos_integer(), outbox()} | {drop, outbox()}.
|
{ok, pos_integer(), outbox()} | {drop, outbox()}.
|
||||||
ensure_writable_segment(RecordBytes, Outbox = #outbox{
|
ensure_writable_segment(RecordBytes, Outbox = #outbox{
|
||||||
write_segment = WriteSegment,
|
|
||||||
max_records = MaxRecords,
|
max_records = MaxRecords,
|
||||||
max_bytes = MaxBytes,
|
max_bytes = MaxBytes,
|
||||||
max_segments = MaxSegments,
|
max_segments = MaxSegments,
|
||||||
segments = Segments,
|
segments = Segments,
|
||||||
acked_seq = AckedSeq
|
writer = #outbox_writer{segment = WriteSegment},
|
||||||
|
checkpoint = #outbox_checkpoint{acked_seq = AckedSeq}
|
||||||
}) ->
|
}) ->
|
||||||
case find_segment(Segments, WriteSegment) of
|
case find_segment(Segments, WriteSegment) of
|
||||||
undefined ->
|
undefined ->
|
||||||
@ -354,21 +431,21 @@ ensure_writable_segment(RecordBytes, Outbox = #outbox{
|
|||||||
true ->
|
true ->
|
||||||
{drop, Outbox};
|
{drop, Outbox};
|
||||||
false ->
|
false ->
|
||||||
{ok, WriteSegment + 1, Outbox#outbox{write_segment = WriteSegment + 1}}
|
{ok, WriteSegment + 1, Outbox#outbox{
|
||||||
|
writer = ((Outbox#outbox.writer)#outbox_writer{
|
||||||
|
segment = WriteSegment + 1
|
||||||
|
})
|
||||||
|
}}
|
||||||
end
|
end
|
||||||
end
|
end
|
||||||
end.
|
end.
|
||||||
|
|
||||||
-spec append_record(file:filename_all(), pos_integer(), binary()) -> ok | {error, term()}.
|
-spec append_record(file:fd(), pos_integer(), binary()) -> ok | {error, term()}.
|
||||||
append_record(Path, Seq, Payload) ->
|
append_record(Fd, Seq, Payload) when Fd =/= undefined ->
|
||||||
Bin = <<Seq:64/unsigned-big-integer, (byte_size(Payload)):32/unsigned-big-integer, Payload/binary>>,
|
Bin = <<Seq:64/unsigned-big-integer, (byte_size(Payload)):32/unsigned-big-integer, Payload/binary>>,
|
||||||
case file:open(Path, [append, raw, binary]) of
|
case file:write(Fd, Bin) of
|
||||||
{ok, Fd} ->
|
ok ->
|
||||||
try
|
ok;
|
||||||
ok = file:write(Fd, Bin)
|
|
||||||
after
|
|
||||||
ok = file:close(Fd)
|
|
||||||
end;
|
|
||||||
{error, Reason} ->
|
{error, Reason} ->
|
||||||
{error, Reason}
|
{error, Reason}
|
||||||
end.
|
end.
|
||||||
@ -453,10 +530,14 @@ advance_acked_seq(AckedSeq, PendingAcks) ->
|
|||||||
end.
|
end.
|
||||||
|
|
||||||
-spec prune_acked_segments(outbox()) -> {ok, outbox()} | {error, term()}.
|
-spec prune_acked_segments(outbox()) -> {ok, outbox()} | {error, term()}.
|
||||||
prune_acked_segments(Outbox = #outbox{segments = Segments, acked_seq = AckedSeq}) ->
|
prune_acked_segments(Outbox = #outbox{
|
||||||
|
segments = Segments,
|
||||||
|
checkpoint = #outbox_checkpoint{acked_seq = AckedSeq}
|
||||||
|
}) ->
|
||||||
{DeleteSegments, KeepSegments} = lists:partition(fun(#segment{end_seq = EndSeq}) ->
|
{DeleteSegments, KeepSegments} = lists:partition(fun(#segment{end_seq = EndSeq}) ->
|
||||||
EndSeq =< AckedSeq
|
EndSeq =< AckedSeq
|
||||||
end, Segments),
|
end, Segments),
|
||||||
|
Outbox0 = detach_deleted_fds(DeleteSegments, Outbox#outbox{segments = KeepSegments}),
|
||||||
ok = lists:foldl(fun(#segment{path = Path}, ok) ->
|
ok = lists:foldl(fun(#segment{path = Path}, ok) ->
|
||||||
case file:delete(Path) of
|
case file:delete(Path) of
|
||||||
ok ->
|
ok ->
|
||||||
@ -467,20 +548,36 @@ prune_acked_segments(Outbox = #outbox{segments = Segments, acked_seq = AckedSeq}
|
|||||||
throw({delete_segment_failed, Path, Reason})
|
throw({delete_segment_failed, Path, Reason})
|
||||||
end
|
end
|
||||||
end, ok, DeleteSegments),
|
end, ok, DeleteSegments),
|
||||||
{ok, reset_reader_after_prune(Outbox#outbox{segments = KeepSegments})}.
|
{ok, reset_reader_after_prune(Outbox0)}.
|
||||||
|
|
||||||
-spec reset_reader_after_prune(outbox()) -> outbox().
|
-spec reset_reader_after_prune(outbox()) -> outbox().
|
||||||
reset_reader_after_prune(Outbox = #outbox{read_segment = ReadSegment, segments = Segments}) ->
|
reset_reader_after_prune(Outbox = #outbox{
|
||||||
|
reader = #outbox_reader{segment = ReadSegment},
|
||||||
|
segments = Segments
|
||||||
|
}) ->
|
||||||
case ReadSegment =:= undefined orelse lists:keymember(ReadSegment, #segment.id, Segments) of
|
case ReadSegment =:= undefined orelse lists:keymember(ReadSegment, #segment.id, Segments) of
|
||||||
true ->
|
true ->
|
||||||
Outbox;
|
Outbox;
|
||||||
false ->
|
false ->
|
||||||
{NReadSegment, NReadOffset} = locate_reader(Outbox),
|
case reset_reader(Outbox) of
|
||||||
Outbox#outbox{read_segment = NReadSegment, read_offset = NReadOffset}
|
{ok, NOutbox} ->
|
||||||
|
NOutbox;
|
||||||
|
{error, _} ->
|
||||||
|
Outbox#outbox{
|
||||||
|
reader = #outbox_reader{
|
||||||
|
segment = undefined,
|
||||||
|
offset = 0,
|
||||||
|
fd = undefined
|
||||||
|
}
|
||||||
|
}
|
||||||
|
end
|
||||||
end.
|
end.
|
||||||
|
|
||||||
-spec locate_reader(outbox()) -> {undefined | pos_integer(), non_neg_integer()}.
|
-spec locate_reader(outbox()) -> {undefined | pos_integer(), non_neg_integer()}.
|
||||||
locate_reader(#outbox{segments = Segments, acked_seq = AckedSeq}) ->
|
locate_reader(#outbox{
|
||||||
|
segments = Segments,
|
||||||
|
checkpoint = #outbox_checkpoint{acked_seq = AckedSeq}
|
||||||
|
}) ->
|
||||||
NextSeq = AckedSeq + 1,
|
NextSeq = AckedSeq + 1,
|
||||||
case lists:dropwhile(fun(#segment{end_seq = EndSeq}) -> EndSeq < NextSeq end, Segments) of
|
case lists:dropwhile(fun(#segment{end_seq = EndSeq}) -> EndSeq < NextSeq end, Segments) of
|
||||||
[] ->
|
[] ->
|
||||||
@ -526,30 +623,6 @@ locate_offset_loop(Fd, TargetSeq, Offset) ->
|
|||||||
{error, Reason}
|
{error, Reason}
|
||||||
end.
|
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()}.
|
-spec read_record(file:fd()) -> eof | {ok, pos_integer(), binary(), pos_integer()} | {error, term()}.
|
||||||
read_record(Fd) ->
|
read_record(Fd) ->
|
||||||
case file:read(Fd, 12) of
|
case file:read(Fd, 12) of
|
||||||
@ -582,3 +655,162 @@ next_segment_id(Segments, SegmentId) ->
|
|||||||
-spec segment_path(file:filename_all(), pos_integer()) -> file:filename_all().
|
-spec segment_path(file:filename_all(), pos_integer()) -> file:filename_all().
|
||||||
segment_path(Dir, SegmentId) ->
|
segment_path(Dir, SegmentId) ->
|
||||||
filename:join(Dir, io_lib:format("~6..0B~s", [SegmentId, ?SEGMENT_EXT])).
|
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} ->
|
||||||
|
case maybe_open_read_fd(Outbox0) of
|
||||||
|
{ok, NOutbox} ->
|
||||||
|
{ok, NOutbox};
|
||||||
|
{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 file:open(Path, [read, raw, binary]) 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.
|
||||||
|
|
||||||
|
-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.
|
||||||
|
|||||||
Loading…
x
Reference in New Issue
Block a user