diff --git a/src/osiris.hrl b/src/osiris.hrl index 0fae757..39200bf 100644 --- a/src/osiris.hrl +++ b/src/osiris.hrl @@ -37,7 +37,7 @@ -define(IS_STRING(S), is_list(S) orelse is_binary(S)). --define(C_NUM_LOG_FIELDS, 5). +-define(C_NUM_LOG_FIELDS, 6). -define(MAGIC, 5). %% chunk format version diff --git a/src/osiris_log.erl b/src/osiris_log.erl index 3f59c38..ad6f89b 100644 --- a/src/osiris_log.erl +++ b/src/osiris_log.erl @@ -88,6 +88,7 @@ -define(C_FIRST_TIMESTAMP, 3). -define(C_CHUNKS, 4). -define(C_SEGMENTS, 5). +-define(C_SEGMENT_SIZE_BYTES, 6). -define(COUNTER_FIELDS, [ {offset, ?C_OFFSET, counter, @@ -95,7 +96,8 @@ {first_offset, ?C_FIRST_OFFSET, counter, "First offset, not updated for readers"}, {first_timestamp, ?C_FIRST_TIMESTAMP, counter, "First timestamp, not updated for readers"}, {chunks, ?C_CHUNKS, counter, "Number of chunks read or written, incremented even if a reader only reads the header"}, - {segments, ?C_SEGMENTS, counter, "Number of segments"} + {segments, ?C_SEGMENTS, counter, "Number of segments"}, + {segment_size_bytes, ?C_SEGMENT_SIZE_BYTES, counter, "Total size of all segment files in bytes"} ] ). @@ -561,6 +563,8 @@ init(#{dir := Dir, shared = Shared, filter_size = FilterSize}, ok = maybe_fix_corrupted_files(Config), + IndexFiles = sorted_index_files(Config), + Config0 = Config#{index_files => IndexFiles}, %% for an acceptor this is derived from the writer's offset range, for a %% writer it is the offset the stream was configured to start at DefaultNextOffset = maps:get(initial_offset, Config, 0), @@ -616,6 +620,8 @@ init(#{dir := Dir, %% at a valid chunk we can now truncate the segment to size in %% case there is trailing data ok = file:truncate(SegFd), + InitSegBytes = sum_log_sizes(IndexFiles), + counters:put(Cnt, ?C_SEGMENT_SIZE_BYTES, InitSegBytes), {ok, IdxFd} = open(IdxFilename, ?FILE_OPTS_WRITE), {ok, IdxEof} = file:position(IdxFd, eof), NumChunks = (IdxEof - ?IDX_HEADER_SIZE) div ?INDEX_RECORD_SIZE_B, @@ -636,6 +642,7 @@ init(#{dir := Dir, {ok, IdxFd} = open(IdxFilename, ?FILE_OPTS_WRITE), {ok, _} = file:position(SegFd, ?LOG_HEADER_SIZE), counters:put(Cnt, ?C_SEGMENTS, 1), + counters:put(Cnt, ?C_SEGMENT_SIZE_BYTES, ?LOG_HEADER_SIZE), %% the segment could potentially have trailing data here so we'll %% do a truncate just in case. The index would have been truncated %% earlier @@ -943,7 +950,7 @@ truncate_to(_Name, _Range, _EpochOffsets, []) -> []; truncate_to(_Name, _Range, [], IdxFiles) -> %% ????? this means the entire log is out - [begin ok = delete_segment_from_index(I) end || I <- IdxFiles], + _ = [delete_segment_from_index(I) || I <- IdxFiles], []; truncate_to(Name, RemoteRange, [{E, ChId} | NextEOs], IdxFiles) -> case find_segment_for_offset(ChId, IdxFiles) of @@ -968,8 +975,7 @@ truncate_to(Name, RemoteRange, [{E, ChId} | NextEOs], IdxFiles) -> true -> %% there is no overlap, need to delete all %% local segments - [begin ok = delete_segment_from_index(I) end - || I <- IdxFiles], + _ = [delete_segment_from_index(I) || I <- IdxFiles], []; false -> %% there is overlap @@ -1011,7 +1017,7 @@ truncate_to(Name, RemoteRange, [{E, ChId} | NextEOs], IdxFiles) -> fun (I) -> case index_file_first_offset(I) > ChId of true -> - ok = delete_segment_from_index(I), + _ = delete_segment_from_index(I), false; false -> true @@ -2007,9 +2013,7 @@ orphaned_segments([_Unexpected | Rem], Acc) -> orphaned_segments(Rem, Acc). first_and_last_seginfos(#{index_files := IdxFiles}) -> - first_and_last_seginfos0(IdxFiles); -first_and_last_seginfos(#{dir := Dir}) -> - first_and_last_seginfos0(sorted_index_files(Dir)). + first_and_last_seginfos0(IdxFiles). first_and_last_seginfos0([]) -> none; @@ -2261,38 +2265,49 @@ update_retention(Retention, State = State0#?MODULE{cfg = Cfg#cfg{retention = Retention}}, trigger_retention_eval(State). +-type retention_result() :: + #{range := range(), + first_timestamp := osiris:timestamp(), + num_remaining_segments := non_neg_integer(), + deleted_segment_bytes := non_neg_integer()}. -spec evaluate_retention(file:filename_all(), [retention_spec()]) -> - {range(), FirstTimestamp :: osiris:timestamp(), - NumRemainingFiles :: non_neg_integer()}. + retention_result(). evaluate_retention(Dir, Specs) when is_list(Dir) -> % convert to binary for faster operations later % mostly in segment_from_index_file/1 evaluate_retention(unicode:characters_to_binary(Dir), Specs); evaluate_retention(Dir, Specs) when is_binary(Dir) -> - {Time, Result} = timer:tc( fun() -> IdxFiles0 = sorted_index_files(Dir), - IdxFiles = evaluate_retention0(IdxFiles0, Specs), - OffsetRange = offset_range_from_idx_files(IdxFiles), - FirstTs = first_timestamp_from_index_files(IdxFiles), - {OffsetRange, FirstTs, length(IdxFiles)} + {IdxFiles, DelSeg} = + evaluate_retention0(IdxFiles0, Specs), + #{range => offset_range_from_idx_files(IdxFiles), + first_timestamp => first_timestamp_from_index_files(IdxFiles), + num_remaining_segments => length(IdxFiles), + deleted_segment_bytes => DelSeg} end), ?DEBUG_(<<>>," (~w) completed in ~fms", [Specs, Time/1_000]), Result. -evaluate_retention0(IdxFiles, []) -> - IdxFiles; -evaluate_retention0(IdxFiles, [{max_bytes, MaxSize} | Specs]) -> - RemIdxFiles = eval_max_bytes(IdxFiles, MaxSize), - evaluate_retention0(RemIdxFiles, Specs); -evaluate_retention0(IdxFiles, [{max_age, Age} | Specs]) -> - RemIdxFiles = eval_age(IdxFiles, Age), - evaluate_retention0(RemIdxFiles, Specs). - -eval_age([_] = IdxFiles, _Age) -> - IdxFiles; -eval_age([IdxFile | IdxFiles] = AllIdxFiles, Age) -> +evaluate_retention0(IdxFiles, Specs) -> + evaluate_retention0(IdxFiles, Specs, 0). + +evaluate_retention0(IdxFiles, [], DeletedBytes) -> + {IdxFiles, DeletedBytes}; +evaluate_retention0(IdxFiles, [{max_bytes, MaxSize} | Specs], Acc) -> + {RemIdxFiles, DelSeg} = eval_max_bytes(IdxFiles, MaxSize), + evaluate_retention0(RemIdxFiles, Specs, Acc + DelSeg); +evaluate_retention0(IdxFiles, [{max_age, Age} | Specs], Acc) -> + {RemIdxFiles, DelSeg} = eval_age(IdxFiles, Age), + evaluate_retention0(RemIdxFiles, Specs, Acc + DelSeg). + +eval_age(IdxFiles, Age) -> + eval_age(IdxFiles, Age, 0). + +eval_age([_] = IdxFiles, _Age, Acc) -> + {IdxFiles, Acc}; +eval_age([IdxFile | IdxFiles] = AllIdxFiles, Age, Acc) -> case last_timestamp_in_index_file(IdxFile) of {ok, Ts} -> Now = erlang:system_time(millisecond), @@ -2301,41 +2316,59 @@ eval_age([IdxFile | IdxFiles] = AllIdxFiles, Age) -> %% the oldest timestamp is older than retention %% and there are other segments available %% we can delete + SegSize = file_size_or_zero(segment_from_index_file(IdxFile)), ok = delete_segment_from_index(IdxFile), - eval_age(IdxFiles, Age); + eval_age(IdxFiles, Age, Acc + SegSize); false -> - AllIdxFiles + {AllIdxFiles, Acc} end; _Err -> - AllIdxFiles + {AllIdxFiles, Acc} end; -eval_age([], _Age) -> +eval_age([], _Age, Acc) -> %% this could happen if retention is evaluated whilst %% a stream is being deleted - []. + {[], Acc}. -eval_max_bytes([], _) -> []; +eval_max_bytes([], _) -> {[], 0}; eval_max_bytes(IdxFiles, MaxSize) -> - [Latest|Older] = lists:reverse(IdxFiles), + [Latest | Older] = lists:reverse(IdxFiles), eval_max_bytes(Older, %% for retention eval it is ok to use a file size function %% that implicitly return 0 when file is not found MaxSize - file_size_or_zero( segment_from_index_file(Latest)), - [Latest]). + [Latest], + 0). -eval_max_bytes([], _, Acc) -> - Acc; -eval_max_bytes([IdxFile | Rest], Limit, Acc) -> +eval_max_bytes([], _, Acc, DeletedBytes) -> + {Acc, DeletedBytes}; +eval_max_bytes([IdxFile | Rest], Limit, Acc, AccSeg) -> SegFile = segment_from_index_file(IdxFile), Size = file_size(SegFile), case Size =< Limit of true -> - eval_max_bytes(Rest, Limit - Size, [IdxFile | Acc]); + eval_max_bytes(Rest, Limit - Size, [IdxFile | Acc], AccSeg); false -> - [ok = delete_segment_from_index(Seg) || Seg <- [IdxFile | Rest]], - Acc - end. + Total = lists:foldl(fun(Seg, S) -> + SegSize = file_size_or_zero( + segment_from_index_file(Seg)), + _ = delete_segment_from_index(Seg), + S + SegSize + end, + AccSeg, + [IdxFile | Rest]), + {Acc, Total} + end. + +sum_log_sizes(IdxFiles) -> + lists:foldl( + fun(IdxFile, SegAcc) -> + SegFile = segment_from_index_file(IdxFile), + SegAcc + file_size_or_zero(SegFile) + end, + 0, + IdxFiles). file_size(Path) -> case prim_file:read_file_info(Path) of @@ -2633,6 +2666,7 @@ write_chunk(Chunk, %% update counters counters:put(CntRef, ?C_OFFSET, NextOffset - 1), counters:add(CntRef, ?C_CHUNKS, 1), + counters:add(CntRef, ?C_SEGMENT_SIZE_BYTES, Size), maybe_set_first_offset(LastChunk, Next, Timestamp, Cfg), State#?MODULE{mode = Write#write{tail_info = {NextOffset, @@ -2873,6 +2907,7 @@ open_new_segment(#?MODULE{cfg = #cfg{name = Name, {ok, _} = file:position(Fd, eof), {ok, _} = file:position(IdxFd, eof), counters:add(Cnt, ?C_SEGMENTS, 1), + counters:add(Cnt, ?C_SEGMENT_SIZE_BYTES, ?LOG_HEADER_SIZE), State0#?MODULE{current_file = Filename, fd = Fd, @@ -3243,13 +3278,17 @@ trigger_retention_eval(#?MODULE{cfg = %% updates first offset and first timestamp %% after retention has been evaluated - EvalFun = fun ({{FstOff, _}, FstTs, NumSegLeft}) + EvalFun = fun (#{range := {FstOff, _}, + first_timestamp := FstTs, + num_remaining_segments := NumSegLeft, + deleted_segment_bytes := DelSegBytes}) when is_integer(FstOff), is_integer(FstTs) -> osiris_log_shared:set_first_chunk_id(Shared, FstOff), counters:put(Cnt, ?C_FIRST_OFFSET, FstOff), counters:put(Cnt, ?C_FIRST_TIMESTAMP, FstTs), - counters:put(Cnt, ?C_SEGMENTS, NumSegLeft); + counters:put(Cnt, ?C_SEGMENTS, NumSegLeft), + counters:sub(Cnt, ?C_SEGMENT_SIZE_BYTES, DelSegBytes); (_) -> ok end, diff --git a/src/osiris_retention.erl b/src/osiris_retention.erl index eca37f6..ae87eea 100644 --- a/src/osiris_retention.erl +++ b/src/osiris_retention.erl @@ -97,7 +97,7 @@ evaluate_retention({eval, Pid, Name, Dir, Specs, Fun} = Eval, State) -> end. schedule({eval, _Pid, Name, _Dir, Specs, _Fun} = Eval, - {_, _, NumSegmentRemaining}, + #{num_remaining_segments := NumSegmentRemaining}, #state{scheduled = Scheduled0} = State) -> %% we need to check the scheduled map even if the current specs do not %% include max_age as the retention config could have changed diff --git a/test/osiris_SUITE.erl b/test/osiris_SUITE.erl index 0e6834f..12d7fc4 100644 --- a/test/osiris_SUITE.erl +++ b/test/osiris_SUITE.erl @@ -62,6 +62,7 @@ all_tests() -> replica_unknown_command, diverged_replica, retention, + retention_size_counters, retention_max_age_eventually, retention_max_age_update_retention, retention_max_age_noproc, @@ -1308,6 +1309,31 @@ retention(Config) -> osiris:stop_cluster(Conf1), ok. +retention_size_counters(Config) -> + Name = ?config(cluster_name, Config), + SegSize = 100 * 1000, + Conf0 = + #{name => Name, + epoch => 1, + leader_node => node(), + retention => [{max_bytes, SegSize * 100}], + max_segment_size_bytes => SegSize, + replica_nodes => []}, + {ok, #{leader_pid := Leader} = Conf1} = osiris:start_cluster(Conf0), + write_n(Leader, 1000, 0, 500 * 8, #{}), + Key = {osiris_writer, Name}, + #{segment_size_bytes := SegBefore} = + osiris_counters:counters(Key, [segment_size_bytes]), + ok = osiris:update_retention(Leader, [{max_bytes, SegSize}]), + await_condition( + fun() -> + #{segment_size_bytes := S} = + osiris_counters:counters(Key, [segment_size_bytes]), + S < SegBefore + end, 200, 50), + osiris:stop_cluster(Conf1), + ok. + retention_max_age_eventually(Config) -> DataDir = ?config(data_dir, Config), Num = 150000, diff --git a/test/osiris_log_SUITE.erl b/test/osiris_log_SUITE.erl index 2e3c538..0a745d5 100644 --- a/test/osiris_log_SUITE.erl +++ b/test/osiris_log_SUITE.erl @@ -36,6 +36,11 @@ all_tests() -> init_recover_with_initial_offset, init_recover_with_writers, init_with_lower_epoch, + size_counters_fresh_log, + size_counters_after_write, + size_counters_restart, + size_counters_after_retention_max_bytes, + size_counters_after_retention_max_age, write_batch, write_first_chunk_updates_first_timestamp, write_first_chunk_with_initial_offset, @@ -265,6 +270,88 @@ init_with_lower_epoch(Config) -> osiris_log:init(Conf#{epoch => 0})), ok. +size_counters_fresh_log(Config) -> + Conf = ?config(osiris_conf, Config), + S0 = osiris_log:init(Conf), + Cnt = osiris_log:counters_ref(S0), + ?assertEqual(?LOG_HEADER_SIZE, counter_get(Cnt, segment_size_bytes)), + ok = osiris_log:close(S0), + ok. + +size_counters_after_write(Config) -> + Conf = ?config(osiris_conf, Config), + S0 = osiris_log:init(Conf), + Cnt = osiris_log:counters_ref(S0), + SegBefore = counter_get(Cnt, segment_size_bytes), + S1 = osiris_log:write([<<"hello">>], S0), + ?assert(counter_get(Cnt, segment_size_bytes) > SegBefore), + ok = osiris_log:close(S1), + ok. + +size_counters_restart(Config) -> + Conf = ?config(osiris_conf, Config), + LDir = ?config(leader_dir, Config), + EpochChunks = [{1, [crypto:strong_rand_bytes(512) || _ <- lists:seq(1, 10)]}], + Log0 = seed_log(Conf#{dir => LDir}, EpochChunks, Config), + Cnt0 = osiris_log:counters_ref(Log0), + SegBytes = counter_get(Cnt0, segment_size_bytes), + ok = osiris_log:close(Log0), + Log1 = osiris_log:init(Conf#{dir => LDir}), + Cnt1 = osiris_log:counters_ref(Log1), + ?assertEqual(SegBytes, counter_get(Cnt1, segment_size_bytes)), + ok = osiris_log:close(Log1), + ok. + +size_counters_after_retention_max_bytes(Config) -> + Conf = ?config(osiris_conf, Config), + LDir = ?config(leader_dir, Config), + Data = crypto:strong_rand_bytes(1500), + EpochChunks = + [begin {1, [Data || _ <- lists:seq(1, 50)]} end + || _ <- lists:seq(1, 20)], + Log0 = seed_log(Conf#{dir => LDir, max_segment_size_bytes => 1000 * 1000}, + EpochChunks, Config), + Cnt = osiris_log:counters_ref(Log0), + SegBefore = counter_get(Cnt, segment_size_bytes), + ok = osiris_log:close(Log0), + Log1 = osiris_log:init(Conf#{dir => LDir}), + Cnt1 = osiris_log:counters_ref(Log1), + SegInit = counter_get(Cnt1, segment_size_bytes), + ?assertEqual(SegBefore, SegInit), + ok = osiris_log:close(Log1), + #{deleted_segment_bytes := DelSeg} = + osiris_log:evaluate_retention(LDir, [{max_bytes, 1500 * 100}]), + ?assert(DelSeg > 0), + Log2 = osiris_log:init(Conf#{dir => LDir}), + Cnt2 = osiris_log:counters_ref(Log2), + ?assertEqual(SegInit - DelSeg, counter_get(Cnt2, segment_size_bytes)), + ok = osiris_log:close(Log2), + ok. + +size_counters_after_retention_max_age(Config) -> + Conf = ?config(osiris_conf, Config), + LDir = ?config(leader_dir, Config), + Data = crypto:strong_rand_bytes(1500), + Ts = now_ms() - 2000, + EpochChunks = + [begin {1, Ts, [Data || _ <- lists:seq(1, 50)]} end + || _ <- lists:seq(1, 20)], + Log0 = seed_log(Conf#{dir => LDir, max_segment_size_bytes => 1000 * 1000}, + EpochChunks, Config), + ok = osiris_log:close(Log0), + Log1 = osiris_log:init(Conf#{dir => LDir}), + Cnt1 = osiris_log:counters_ref(Log1), + SegInit = counter_get(Cnt1, segment_size_bytes), + ok = osiris_log:close(Log1), + #{deleted_segment_bytes := DelSeg} = + osiris_log:evaluate_retention(LDir, [{max_bytes, 100000000}, {max_age, 1000}]), + ?assert(DelSeg > 0), + Log2 = osiris_log:init(Conf#{dir => LDir}), + Cnt2 = osiris_log:counters_ref(Log2), + ?assertEqual(SegInit - DelSeg, counter_get(Cnt2, segment_size_bytes)), + ok = osiris_log:close(Log2), + ok. + write_batch(Config) -> Conf = ?config(osiris_conf, Config), S0 = osiris_log:init(Conf), @@ -1846,9 +1933,15 @@ evaluate_retention_max_bytes(Config) -> osiris_log:close(Log), %% this should delete at least one segment Spec = {max_bytes, 1500 * 100}, - Range = osiris_log:evaluate_retention(LDir, [Spec]), + #{range := OffRange, + first_timestamp := FstTs, + num_remaining_segments := NumSeg} = + osiris_log:evaluate_retention(LDir, [Spec]), %% idempotency check - Range = osiris_log:evaluate_retention(LDir, [Spec]), + #{range := OffRange, + first_timestamp := FstTs, + num_remaining_segments := NumSeg} = + osiris_log:evaluate_retention(LDir, [Spec]), SegFiles = filelib:wildcard( filename:join(LDir, "*.segment")), @@ -1890,9 +1983,15 @@ evaluate_retention_max_age(Config) -> %% this should delete at least one segment as all chunks should be older %% than the retention of 1000ms; max_bytes shouldn't affect the result Spec = [{max_bytes, 100000000}, {max_age, 1000}], - Range = osiris_log:evaluate_retention(LDir, Spec), - %% idempotency - Range = osiris_log:evaluate_retention(LDir, Spec), + #{range := OffRange, + first_timestamp := FstTs, + num_remaining_segments := NumSeg} = + osiris_log:evaluate_retention(LDir, Spec), + %% idempotency check: range and segment count must be stable; deleted bytes differ + #{range := OffRange, + first_timestamp := FstTs, + num_remaining_segments := NumSeg} = + osiris_log:evaluate_retention(LDir, Spec), SegFiles = filelib:wildcard( filename:join(LDir, "*.segment")), @@ -3323,3 +3422,8 @@ recv(Socket, Expected, Acc) -> Other -> Other end. + +counter_get(Cnt, Name) -> + Fields = osiris_log:counter_fields(), + {Name, Pos, _, _} = lists:keyfind(Name, 1, Fields), + counters:get(Cnt, Pos).