diff --git a/src/outbox/endpoint_outbox.erl b/src/outbox/endpoint_outbox.erl new file mode 100644 index 0000000..2d77853 --- /dev/null +++ b/src/outbox/endpoint_outbox.erl @@ -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 = <>, + 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, <>} -> + 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])).