diff --git a/src/outbox/endpoint_outbox.erl b/src/outbox/endpoint_outbox.erl index 2d77853..3d9c746 100644 --- a/src/outbox/endpoint_outbox.erl +++ b/src/outbox/endpoint_outbox.erl @@ -9,7 +9,7 @@ %%%------------------------------------------------------------------- -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]). -record(segment, { @@ -21,6 +21,24 @@ 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, { endpoint_id :: integer() | binary() | atom(), dir :: file:filename_all(), @@ -28,14 +46,10 @@ 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() + writer = #outbox_writer{} :: #outbox_writer{}, + reader = #outbox_reader{} :: #outbox_reader{}, + checkpoint = #outbox_checkpoint{} :: #outbox_checkpoint{} }). -type outbox() :: #outbox{}. @@ -67,16 +81,35 @@ open(EndpointId, Opts) when is_map(Opts) -> 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, + 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), - Outbox = Outbox0#outbox{read_segment = ReadSegment, read_offset = ReadOffset}, - ok = persist_metadata(Outbox), - {ok, Outbox}. + Outbox1 = Outbox0#outbox{ + reader = #outbox_reader{ + segment = ReadSegment, + offset = ReadOffset + } + }, + case open_fds(Outbox1) of + {ok, Outbox} -> + ok = persist_metadata(Outbox), + {ok, Outbox}; + {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()}. @@ -86,37 +119,56 @@ append(Payload, Outbox = #outbox{}) when is_binary(Payload) -> {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}; + case ensure_write_fd(SegmentId, NOutbox0) of + {ok, NOutbox1} -> + 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 -> + 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 + }, + 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; {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 -> +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, Outbox#outbox.pending_acks), + 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, Outbox#outbox.acked_segment), + NackedSegment = segment_id_for_seq( + Outbox#outbox.segments, + NackedSeq, + Checkpoint#outbox_checkpoint.acked_segment), Outbox0 = Outbox#outbox{ - acked_seq = NackedSeq, - acked_segment = NackedSegment, - pending_acks = PendingAcks + checkpoint = Checkpoint#outbox_checkpoint{ + acked_seq = NackedSeq, + acked_segment = NackedSegment, + pending_acks = PendingAcks + } }, {ok, Outbox1} = prune_acked_segments(Outbox0), case NackedSeq =:= AckedSeq of @@ -130,20 +182,37 @@ ack(Seq, Outbox = #outbox{acked_seq = AckedSeq}) when is_integer(Seq), Seq > 0 - -spec next(outbox()) -> eof | {ok, Seq :: pos_integer(), Payload :: binary(), outbox()} | {error, term()}. -next(#outbox{read_segment = undefined}) -> +next(#outbox{reader = #outbox_reader{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}) +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_record(Reader0#outbox_reader.fd) 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, Seq, Payload, RecordBytes} -> + {ok, Seq, Payload, NOutbox0#outbox{ + reader = Reader0#outbox_reader{offset = Offset + RecordBytes} + }}; + {error, Reason} -> + {error, Reason} end; - {ok, Seq, Payload, NextOffset} -> - {ok, Seq, Payload, Outbox#outbox{read_offset = NextOffset}}; {error, Reason} -> {error, Reason} end. @@ -151,15 +220,19 @@ next(Outbox = #outbox{read_segment = SegmentId, read_offset = Offset, dir = Dir, -spec reset_reader(outbox()) -> {ok, outbox()}. reset_reader(Outbox = #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(). stat(#outbox{ dir = Dir, - write_segment = WriteSegment, - next_seq = NextSeq, - acked_segment = AckedSegment, - acked_seq = AckedSeq, + writer = #outbox_writer{ + segment = WriteSegment, + next_seq = NextSeq + }, + checkpoint = #outbox_checkpoint{ + acked_segment = AckedSegment, + acked_seq = AckedSeq + }, segments = Segments }) -> #{ @@ -191,7 +264,7 @@ sanitize_binary(Bin) -> (C) when C >= $A, C =< $Z -> C; (C) when C >= $0, C =< $9 -> C; ($-) -> $-; - ($_ ) -> $_; + ($_) -> $_; (_) -> $_ end, binary_to_list(Bin)). @@ -238,10 +311,14 @@ default_metadata() -> -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 + writer = #outbox_writer{ + segment = WriteSegment, + next_seq = NextSeq + }, + checkpoint = #outbox_checkpoint{ + acked_segment = AckedSegment, + acked_seq = AckedSeq + } }) -> TmpPath = MetadataPath ++ ".tmp", Metadata = #{ @@ -334,12 +411,12 @@ encoded_record_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 + writer = #outbox_writer{segment = WriteSegment}, + checkpoint = #outbox_checkpoint{acked_seq = AckedSeq} }) -> case find_segment(Segments, WriteSegment) of undefined -> @@ -354,21 +431,21 @@ ensure_writable_segment(RecordBytes, Outbox = #outbox{ true -> {drop, Outbox}; 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. --spec append_record(file:filename_all(), pos_integer(), binary()) -> ok | {error, term()}. -append_record(Path, Seq, Payload) -> +-spec append_record(file:fd(), pos_integer(), binary()) -> ok | {error, term()}. +append_record(Fd, Seq, Payload) when Fd =/= undefined -> Bin = <>, - case file:open(Path, [append, raw, binary]) of - {ok, Fd} -> - try - ok = file:write(Fd, Bin) - after - ok = file:close(Fd) - end; + case file:write(Fd, Bin) of + ok -> + ok; {error, Reason} -> {error, Reason} end. @@ -453,10 +530,14 @@ advance_acked_seq(AckedSeq, PendingAcks) -> end. -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}) -> 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 -> @@ -467,20 +548,36 @@ prune_acked_segments(Outbox = #outbox{segments = Segments, acked_seq = AckedSeq} throw({delete_segment_failed, Path, Reason}) end 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(). -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 true -> Outbox; false -> - {NReadSegment, NReadOffset} = locate_reader(Outbox), - Outbox#outbox{read_segment = NReadSegment, read_offset = NReadOffset} + 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, acked_seq = AckedSeq}) -> +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 [] -> @@ -526,30 +623,6 @@ locate_offset_loop(Fd, TargetSeq, Offset) -> {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 @@ -582,3 +655,162 @@ next_segment_id(Segments, SegmentId) -> -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} -> + 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.