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