Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
32 changes: 19 additions & 13 deletions src/osiris.erl
Original file line number Diff line number Diff line change
Expand Up @@ -68,8 +68,15 @@
{abs, offset()} |
offset() |
{timestamp, timestamp()}.
-type retention_fun() :: fun((IdxFiles :: [file:filename_all()]) ->
{ToDelete :: [file:filename_all()], ToKeep :: [file:filename_all()]}).
-type retention_spec() ::
{max_bytes, non_neg_integer()} | {max_age, milliseconds()}.
{max_bytes, non_neg_integer()} |
{max_age, milliseconds()} |
{'fun', retention_fun()}.
-type retention_update() ::
[retention_spec()] |
fun(([retention_spec()]) -> [retention_spec()]).
-type writer_id() :: binary().
-type batch() :: {batch, NumRecords :: non_neg_integer(),
compression_type(),
Expand All @@ -83,11 +90,6 @@

%% returned when reading
-type entry() :: binary() | batch().
-type reader_options() :: #{transport => tcp | ssl,
chunk_selector => all | user_data,
filter_spec => osiris_bloom:filter_spec(),
read_ahead => boolean() | non_neg_integer()
}.

-export_type([name/0,
config/0,
Expand All @@ -97,6 +99,8 @@
tracking_id/0,
offset_spec/0,
retention_spec/0,
retention_update/0,
retention_fun/0,
timestamp/0,
writer_id/0,
data/0,
Expand Down Expand Up @@ -228,8 +232,9 @@ init_reader(Pid, OffsetSpec, CounterSpec) ->
chunk_selector => user_data}).

-spec init_reader(pid(), offset_spec(), osiris_log:counter_spec(),
reader_options()) ->
{ok, osiris_log:state()} |
osiris_log_reader:options()) ->
{ok, osiris_log_reader:state()} |
{error, osiris_log_reader:transient_error()} |
{error, {offset_out_of_range, empty | {offset(), offset()}}} |
{error, {invalid_last_offset_epoch, offset(), offset()}}.
init_reader(Pid, OffsetSpec, {_, _} = CounterSpec, Options)
Expand All @@ -238,18 +243,19 @@ init_reader(Pid, OffsetSpec, {_, _} = CounterSpec, Options)
Ctx0 = osiris_util:get_reader_context(Pid),
Ctx = Ctx0#{counter_spec => CounterSpec,
options => Options},
osiris_log:init_offset_reader(OffsetSpec, Ctx).
(osiris_log_reader:module()):init_offset_reader(OffsetSpec, Ctx).

-spec resolve_offset_spec(pid(), offset_spec(), reader_options()) ->
-spec resolve_offset_spec(pid(), offset_spec(), osiris_log_reader:options()) ->
{ok, offset()} |
{error, osiris_log_reader:transient_error()} |
{error, no_index_file} |
{error, {offset_out_of_range, osiris_log:range()}} |
{error, retries_exhausted}.
resolve_offset_spec(Pid, OffsetSpec, Options)
when is_pid(Pid) andalso node(Pid) =:= node() ->
Ctx0 = osiris_util:get_reader_context(Pid),
Ctx = Ctx0#{options => Options},
osiris_log:resolve_offset_spec(OffsetSpec, Ctx).
(osiris_log_reader:module()):resolve_offset_spec(OffsetSpec, Ctx).

-spec register_offset_listener(pid(), offset()) -> ok.
register_offset_listener(Pid, Offset) ->
Expand All @@ -273,10 +279,10 @@ register_offset_listener(Pid, Offset, EvtFormatter) ->
end,
ok.

-spec update_retention(pid(), [osiris:retention_spec()]) ->
-spec update_retention(pid(), retention_update()) ->
ok | {error, term()}.
update_retention(Pid, Retention)
when is_pid(Pid) andalso is_list(Retention) ->
when is_pid(Pid) andalso (is_list(Retention) orelse is_function(Retention, 1)) ->
Msg = {update_retention, Retention},
try
case gen:call(Pid, '$gen_call', Msg) of
Expand Down
Loading
Loading