Packages
reckon_db
2.3.5
5.11.0
5.10.4
5.10.3
5.10.1
5.10.0
5.9.1
5.9.0
5.8.3
5.8.2
5.8.1
5.8.0
5.7.0
5.6.1
5.6.0
5.5.5
5.5.4
5.5.3
5.5.2
5.5.1
5.5.0
5.4.0
5.2.2
5.2.1
5.2.0
5.1.0
5.0.0
4.0.0
3.1.2
3.1.1
3.0.0
2.3.7
2.3.6
2.3.5
2.3.4
2.3.3
2.3.2
2.3.1
2.3.0
2.2.2
2.2.0
2.1.4
2.1.3
2.1.2
2.1.1
2.1.0
2.0.0
1.7.7
1.7.6
1.7.5
1.7.4
1.7.3
1.7.2
1.7.1
1.7.0
1.6.3
1.6.2
1.6.1
1.6.0
1.5.1
1.5.0
1.4.5
1.4.4
1.4.3
1.4.2
1.4.1
1.4.0
1.3.3
1.3.2
1.3.1
1.3.0
1.2.7
1.2.6
1.2.5
1.2.4
1.2.3
1.2.2
1.2.1
1.2.0
1.1.1
1.1.0
1.0.3
1.0.2
1.0.1
1.0.0
BEAM-native Event Store built on Khepri/Ra with Raft consensus. Event sourcing, persistent subscriptions, snapshots, and automatic cluster formation via UDP multicast discovery. Ships embedded Rust NIFs for 3-15x acceleration of crypto, hashing, compression, aggregation, filter matching, and grap...
Current section
Files
Jump to
Current section
Files
src/reckon_db_streams.erl
%% @doc Streams API facade for reckon-db
%%
%% Provides the public API for stream operations:
%% - append: Write events to a stream with optimistic concurrency
%% - read: Read events from a stream
%% - get_version: Get current stream version
%% - exists: Check if stream exists
%% - list_streams: List all streams in the store
%%
%% @author rgfaber
-module(reckon_db_streams).
-include("reckon_db.hrl").
-include("reckon_db_telemetry.hrl").
-include_lib("khepri/include/khepri.hrl").
%% API
-export([
append/4,
append/5,
read/5,
read/6,
read_all/4,
read_all_global/3,
read_by_event_types/3,
read_by_tags/4,
get_version/2,
exists/2,
has_events/1,
list_streams/1,
delete/2
]).
%% Internal exports for workers
-export([
do_append/4,
do_read/5
]).
%%====================================================================
%% Types
%%====================================================================
-type new_event() :: #{
event_type := binary(),
data := map() | binary(),
metadata => map(),
tags => [binary()],
event_id => binary()
}.
-type direction() :: forward | backward.
-export_type([new_event/0, direction/0]).
%%====================================================================
%% API
%%====================================================================
%% @doc Append events to a stream with expected version check
%%
%% Expected version semantics:
%% -1 (NO_STREAM) - Stream must not exist (first write)
%% -2 (ANY_VERSION) - No version check, always append
%% N >= 0 - Stream version must equal N
%%
%% Returns {ok, NewVersion} on success or {error, Reason} on failure.
-spec append(atom(), binary(), integer(), [new_event()]) ->
{ok, non_neg_integer()} | {error, term()}.
append(StoreId, StreamId, ExpectedVersion, Events) ->
append(StoreId, StreamId, ExpectedVersion, Events, #{}).
-spec append(atom(), binary(), integer(), [new_event()], map()) ->
{ok, non_neg_integer()} | {error, term()}.
append(StoreId, StreamId, ExpectedVersion, Events, _Opts) ->
%% Stream-id format gate. Rejecting at the head of append/4
%% means no malformed id reaches Khepri, so the store can't
%% accumulate polluted paths from misbehaving tests / clients.
%% See reckon_db_stream_id for the format rules.
case reckon_db_stream_id:validate(StreamId) of
ok ->
do_append_with_telemetry(StoreId, StreamId, ExpectedVersion, Events);
{error, Reason} ->
{error, {invalid_stream_id, Reason, StreamId}}
end.
%% @private validation-gated wrapper around do_append that also
%% emits the write-lifecycle telemetry.
do_append_with_telemetry(StoreId, StreamId, ExpectedVersion, Events) ->
StartTime = erlang:monotonic_time(),
%% Emit start telemetry
telemetry:execute(
?STREAM_WRITE_START,
#{system_time => erlang:system_time(millisecond)},
#{store_id => StoreId, stream_id => StreamId,
event_count => length(Events), expected_version => ExpectedVersion}
),
Result = do_append(StoreId, StreamId, ExpectedVersion, Events),
Duration = erlang:monotonic_time() - StartTime,
%% Emit stop/error telemetry
case Result of
{ok, NewVersion} ->
telemetry:execute(
?STREAM_WRITE_STOP,
#{duration => Duration, event_count => length(Events)},
#{store_id => StoreId, stream_id => StreamId, new_version => NewVersion}
),
Result;
{error, Reason} ->
telemetry:execute(
?STREAM_WRITE_ERROR,
#{duration => Duration},
#{store_id => StoreId, stream_id => StreamId, reason => Reason}
),
Result
end.
%% @doc Read events from a stream
%%
%% Parameters:
%% StoreId - The store identifier
%% StreamId - The stream identifier
%% StartVersion - Starting version (0-based)
%% Count - Maximum number of events to read
%% Direction - forward or backward
%%
%% Returns {ok, [Event]} or {error, Reason}
-spec read(atom(), binary(), non_neg_integer(), pos_integer(), direction()) ->
{ok, [event()]} | {error, term()}.
read(StoreId, StreamId, StartVersion, Count, Direction) ->
read(StoreId, StreamId, StartVersion, Count, Direction, #{}).
%% @doc Read events from a stream with explicit options.
%%
%% Currently supported options:
%%
%% `verify' :: `skip_legacy' | `strict' | `skip_all'
%% Tamper-resistance enforcement mode. Default: `skip_legacy'.
%% - `skip_legacy' (default): events with version below the
%% per-stream chain_start watermark are returned untouched
%% (legacy data); events at or above the watermark are
%% verified strictly and an integrity_violation is returned
%% on any failure.
%% - `strict': every event must carry integrity fields and
%% verify; legacy events surface as missing_integrity.
%% - `skip_all': no verification (dangerous; intended for
%% migration tooling only).
%%
%% Backward-direction reads always bypass chain verification in 2.1.0;
%% the MAC alone could still be checked but is not in this release.
%% Forward reads receive full chain + MAC verification.
-type verify_mode() :: skip_legacy | strict | skip_all.
-type read_opts() :: #{verify => verify_mode()}.
-spec read(
atom(), binary(), non_neg_integer(), pos_integer(), direction(), read_opts()
) -> {ok, [event()]} | {error, term()}.
read(StoreId, StreamId, StartVersion, Count, Direction, Opts) ->
StartTime = erlang:monotonic_time(),
%% Emit start telemetry
telemetry:execute(
?STREAM_READ_START,
#{system_time => erlang:system_time(millisecond)},
#{store_id => StoreId, stream_id => StreamId,
start_version => StartVersion, count => Count, direction => Direction}
),
Result = do_read_with_verify(
StoreId, StreamId, StartVersion, Count, Direction, Opts),
Duration = erlang:monotonic_time() - StartTime,
%% Emit stop telemetry
case Result of
{ok, Events} ->
telemetry:execute(
?STREAM_READ_STOP,
#{duration => Duration, event_count => length(Events)},
#{store_id => StoreId, stream_id => StreamId}
),
Result;
Error ->
Error
end.
%% @doc Read all events from a stream
-spec read_all(atom(), binary(), pos_integer(), direction()) ->
{ok, [event()]} | {error, term()}.
read_all(StoreId, StreamId, BatchSize, Direction) ->
read(StoreId, StreamId, 0, BatchSize, Direction).
%% @doc Read all events across all streams in global epoch_us order.
%%
%% Returns events sorted by epoch_us, skipping `Offset' events and
%% returning up to `BatchSize' events. Used by catch-up subscriptions
%% to replay historical events to a subscriber.
%%
%% Parameters:
%% StoreId - The store identifier
%% Offset - Number of events to skip (0-based)
%% BatchSize - Maximum number of events to return
%%
%% Returns events sorted by epoch_us (global ordering).
-spec read_all_global(atom(), non_neg_integer(), pos_integer()) ->
{ok, [event()]} | {error, term()}.
read_all_global(StoreId, Offset, BatchSize) ->
Path = [streams,
?KHEPRI_WILDCARD_STAR,
#if_all{conditions = [
?KHEPRI_WILDCARD_STAR,
#if_has_data{has_data = true}
]}],
case khepri:get_many(StoreId, Path) of
{ok, Results} when is_map(Results) ->
Events = [convert_result_to_event(PathKey, Value)
|| {PathKey, Value} <- maps:to_list(Results)],
ValidEvents = [E || E <- Events, E =/= undefined],
SortedEvents = lists:sort(
fun(#event{epoch_us = E1}, #event{epoch_us = E2}) -> E1 =< E2 end,
ValidEvents),
Skipped = safe_nthtail(Offset, SortedEvents),
{ok, lists:sublist(Skipped, BatchSize)};
{ok, _} ->
{ok, []};
{error, _} = Error ->
Error
end.
%% @doc Read all events of specific types from all streams using Khepri native filtering.
%%
%% This function uses Khepri's built-in #if_data_matches condition to filter
%% events by type at the database level, avoiding loading all events into memory.
%%
%% Parameters:
%% StoreId - The store identifier
%% EventTypes - List of event type binaries to match
%% BatchSize - Maximum number of events to return (for pagination)
%%
%% Returns events sorted by epoch_us (global ordering).
-spec read_by_event_types(atom(), [binary()], pos_integer()) ->
{ok, [event()]} | {error, term()}.
read_by_event_types(StoreId, EventTypes, BatchSize) when is_list(EventTypes) ->
%% Build Khepri path pattern with native data matching
%% Path structure: [streams, StreamId, PaddedVersion]
%% We use #if_any to match any of the specified event types
TypeConditions = [
#if_data_matches{pattern = #event{event_type = ET, _ = '_'}}
|| ET <- EventTypes
],
%% Combine conditions: match any event type AND has data
DataCondition = case TypeConditions of
[] ->
#if_has_data{has_data = true};
[SingleCondition] ->
#if_all{conditions = [
#if_has_data{has_data = true},
SingleCondition
]};
_ ->
#if_all{conditions = [
#if_has_data{has_data = true},
#if_any{conditions = TypeConditions}
]}
end,
%% Build path pattern matching all events across all streams
Path = [streams,
?KHEPRI_WILDCARD_STAR, %% Any stream ID
#if_all{conditions = [
?KHEPRI_WILDCARD_STAR, %% Any version
DataCondition
]}],
case khepri:get_many(StoreId, Path) of
{ok, Results} when is_map(Results) ->
%% Convert results to events and sort by epoch_us
Events = [convert_result_to_event(PathKey, Value)
|| {PathKey, Value} <- maps:to_list(Results)],
ValidEvents = [E || E <- Events, E =/= undefined],
%% Sort by epoch_us for global ordering
SortedEvents = lists:sort(
fun(#event{epoch_us = E1}, #event{epoch_us = E2}) -> E1 =< E2 end,
ValidEvents
),
%% Apply batch size limit
LimitedEvents = lists:sublist(SortedEvents, BatchSize),
{ok, LimitedEvents};
{ok, _} ->
{ok, []};
{error, _} = Error ->
Error
end.
%% @doc Read all events matching tags from all streams.
%%
%% Tags provide a mechanism for cross-stream querying without affecting
%% stream-based concurrency control. This is useful for the process-centric
%% model where you want to find all events related to specific participants.
%%
%% == Match Modes ==
%%
%% `any' (default): Returns events containing ANY of the specified tags (union).
%% Example: `read_by_tags(Store, [<<"student:456">>, <<"student:789">>], any, 100)'
%% Returns events for either student.
%%
%% `all': Returns events containing ALL of the specified tags (intersection).
%% Example: `read_by_tags(Store, [<<"student:456">>, <<"course:CS101">>], all, 100)'
%% Returns only events tagged with both student 456 AND course CS101.
%%
%% == Parameters ==
%%
%% StoreId - The store identifier
%% Tags - List of tag binaries to match
%% Match - `any' | `all' (matching strategy)
%% BatchSize - Maximum number of events to return
%%
%% == Returns ==
%%
%% Events sorted by epoch_us (global ordering).
-spec read_by_tags(atom(), [binary()], any | all, pos_integer()) ->
{ok, [event()]} | {error, term()}.
read_by_tags(StoreId, Tags, Match, BatchSize) when is_list(Tags), is_atom(Match) ->
%% Query all events from all streams and filter by tags in Erlang.
%% Khepri's pattern matching doesn't easily support list membership checks,
%% so we fetch events with data and filter client-side.
%%
%% For large stores, consider maintaining a separate tag index.
Path = [streams,
?KHEPRI_WILDCARD_STAR, %% Any stream ID
#if_all{conditions = [
?KHEPRI_WILDCARD_STAR, %% Any version
#if_has_data{has_data = true}
]}],
case khepri:get_many(StoreId, Path) of
{ok, Results} when is_map(Results) ->
%% Convert results to events
AllEvents = [convert_result_to_event(PathKey, Value)
|| {PathKey, Value} <- maps:to_list(Results)],
ValidEvents = [E || E <- AllEvents, E =/= undefined],
%% Filter by tag match mode (only events with tags list)
FilteredEvents = filter_events_by_tags(ValidEvents, Tags, Match),
%% Sort by epoch_us for global ordering
SortedEvents = lists:sort(
fun(#event{epoch_us = E1}, #event{epoch_us = E2}) -> E1 =< E2 end,
FilteredEvents
),
%% Apply batch size limit
LimitedEvents = lists:sublist(SortedEvents, BatchSize),
{ok, LimitedEvents};
{ok, _} ->
{ok, []};
{error, _} = Error ->
Error
end.
%% @private Filter events by tags according to match mode
%% Empty search tags always returns empty (no criteria = no results)
-spec filter_events_by_tags([event()], [binary()], any | all) -> [event()].
filter_events_by_tags(_Events, [], _Match) ->
[];
filter_events_by_tags(Events, Tags, any) ->
%% Return events that have ANY of the tags (union)
TagSet = sets:from_list(Tags),
lists:filter(
fun(#event{tags = EventTags}) when is_list(EventTags), EventTags =/= [] ->
EventTagSet = sets:from_list(EventTags),
sets:size(sets:intersection(TagSet, EventTagSet)) > 0;
(_) ->
false
end,
Events
);
filter_events_by_tags(Events, Tags, all) ->
%% Return events that have ALL of the tags (intersection)
TagSet = sets:from_list(Tags),
lists:filter(
fun(#event{tags = EventTags}) when is_list(EventTags), EventTags =/= [] ->
EventTagSet = sets:from_list(EventTags),
sets:is_subset(TagSet, EventTagSet);
(_) ->
false
end,
Events
).
%% @private Convert a Khepri result to an event, adding stream_id from path
-spec convert_result_to_event([atom() | binary()], term()) -> event() | undefined.
convert_result_to_event([streams, StreamId | _], Event) when is_record(Event, event) ->
Event#event{stream_id = StreamId};
convert_result_to_event([streams, StreamId | _], EventMap) when is_map(EventMap) ->
Event = map_to_event(EventMap),
Event#event{stream_id = StreamId};
convert_result_to_event(_, _) ->
undefined.
%% @doc Get current version of a stream
%%
%% Returns:
%% -1 - if stream doesn't exist or is empty
%% N >= 0 - representing the version of the latest event
-spec get_version(atom(), binary()) -> integer().
get_version(StoreId, StreamId) ->
Path = ?STREAMS_PATH ++ [StreamId],
case khepri:count(StoreId, Path ++ [?KHEPRI_WILDCARD_STAR]) of
{ok, 0} -> ?NO_STREAM;
{ok, Count} -> Count - 1;
{error, _} -> ?NO_STREAM
end.
%% @doc Check if a stream exists
-spec exists(atom(), binary()) -> boolean().
exists(StoreId, StreamId) ->
Path = ?STREAMS_PATH ++ [StreamId],
case khepri:exists(StoreId, Path) of
true -> true;
false -> false;
{error, _} -> false
end.
%% @doc Check if a store contains at least one event.
%% Cannot rely on stream existence alone — streams can survive
%% after all their events are deleted (truncation, GDPR erasure).
%% Checks for actual event data by reading 1 event globally.
-spec has_events(atom()) -> boolean().
has_events(StoreId) ->
case read_all_global(StoreId, 0, 1) of
{ok, [_ | _]} -> true;
_ -> false
end.
%% @doc List all streams in the store
-spec list_streams(atom()) -> {ok, [binary()]} | {error, term()}.
list_streams(StoreId) ->
%% Use get_many with wildcard to find all stream IDs
Path = ?STREAMS_PATH ++ [?KHEPRI_WILDCARD_STAR],
case khepri:get_many(StoreId, Path) of
{ok, Results} when is_map(Results) ->
%% Extract stream IDs from paths: {[streams, StreamId, ...], _}
StreamIds = lists:usort([
extract_stream_id(P) || P <- maps:keys(Results)
]),
{ok, StreamIds};
{ok, _} ->
{ok, []};
{error, _} = Error ->
Error
end.
%% @private Extract stream ID from path [streams, StreamId, ...]
-spec extract_stream_id([atom() | binary()]) -> binary().
extract_stream_id([streams, StreamId | _]) ->
StreamId;
extract_stream_id(_) ->
<<>>.
%% @doc Delete a stream and all its events
-spec delete(atom(), binary()) -> ok | {error, term()}.
delete(StoreId, StreamId) ->
Path = ?STREAMS_PATH ++ [StreamId],
case khepri:delete(StoreId, Path) of
ok -> ok;
{error, _} = Error -> Error
end.
%%====================================================================
%% Internal functions
%%====================================================================
%% @private
-spec do_append(atom(), binary(), integer(), [new_event()]) ->
{ok, non_neg_integer()} | {error, term()}.
do_append(StoreId, StreamId, ExpectedVersion, Events) ->
CurrentVersion = get_version(StoreId, StreamId),
%% Check expected version
case check_expected_version(ExpectedVersion, CurrentVersion) of
ok ->
append_events_to_stream(StoreId, StreamId, CurrentVersion, Events);
{error, _} = Error ->
Error
end.
%% @private
-spec check_expected_version(integer(), integer()) -> ok | {error, term()}.
check_expected_version(?ANY_VERSION, _CurrentVersion) ->
ok;
check_expected_version(?NO_STREAM, CurrentVersion) when CurrentVersion =:= ?NO_STREAM ->
ok;
check_expected_version(?NO_STREAM, CurrentVersion) ->
{error, {wrong_expected_version, ?NO_STREAM, CurrentVersion}};
check_expected_version(ExpectedVersion, CurrentVersion) when ExpectedVersion =:= CurrentVersion ->
ok;
check_expected_version(ExpectedVersion, CurrentVersion) ->
{error, {wrong_expected_version, ExpectedVersion, CurrentVersion}}.
%% @private
-spec append_events_to_stream(atom(), binary(), integer(), [new_event()]) ->
{ok, non_neg_integer()}.
append_events_to_stream(StoreId, StreamId, CurrentVersion, Events) ->
Now = erlang:system_time(millisecond),
EpochUs = erlang:system_time(microsecond),
%% Resolve the tamper-resistance context for this batch. For
%% stores with integrity disabled (the default and only mode
%% pre-2.1) this is a constant `disabled` pass-through. For
%% integrity-enabled stores we load the HMAC key, ensure the
%% per-stream watermark exists, and resolve the initial chain
%% tip — all once, before the per-event loop.
IntegrityCtx = setup_integrity(StoreId, StreamId, CurrentVersion),
InitialTip = resolve_initial_tip(StoreId, StreamId, CurrentVersion + 1, IntegrityCtx),
{FinalVersion, _FinalTip} = lists:foldl(
fun(Event, {AccVersion, AccTip}) ->
NewVersion = AccVersion + 1,
RecordedEvent0 = create_event_record(
Event, StreamId, NewVersion, Now, EpochUs),
{RecordedEvent, NextTip} = apply_integrity_if_enabled(
RecordedEvent0, AccTip, IntegrityCtx),
PaddedVersion = pad_version(NewVersion, ?VERSION_PADDING),
Path = ?STREAMS_PATH ++ [StreamId, PaddedVersion],
ok = khepri:put(StoreId, Path, RecordedEvent),
{NewVersion, NextTip}
end,
{CurrentVersion, InitialTip},
Events
),
{ok, FinalVersion}.
%% @private Tamper-resistance context for a single append batch.
%%
%% Either `disabled` (no integrity work) or
%% `{enabled, Key, ChainStart}` carrying the HMAC key and the
%% chain-start watermark for this stream.
-type integrity_ctx() ::
disabled |
{enabled, Key :: binary(), ChainStart :: non_neg_integer()}.
-spec setup_integrity(atom(), binary(), integer()) -> integrity_ctx().
setup_integrity(StoreId, StreamId, CurrentVersion) ->
case reckon_db_integrity_key:is_enabled(StoreId) of
false ->
disabled;
true ->
NextVersion = CurrentVersion + 1,
{ok, ChainStart} = reckon_db_chain_watermark:set_if_absent(
StoreId, StreamId, NextVersion),
Key = reckon_db_integrity_key:get(StoreId),
{enabled, Key, ChainStart}
end.
%% @private Resolve the chain-tip value that the FIRST event in this
%% batch must reference as its `prev_event_hash`.
%%
%% Disabled context: tip is `undefined` (unused).
%% Enabled, first integrity event in stream (NextVersion =:= ChainStart):
%% tip is the genesis 32-zero-byte value.
%% Enabled, later batch (NextVersion > ChainStart):
%% tip is computed from the predecessor event on disk.
-spec resolve_initial_tip(
atom(), binary(), non_neg_integer(), integrity_ctx()
) -> binary() | undefined.
resolve_initial_tip(_StoreId, _StreamId, _NextVersion, disabled) ->
undefined;
resolve_initial_tip(_StoreId, _StreamId, NextVersion,
{enabled, _Key, ChainStart})
when NextVersion =:= ChainStart ->
reckon_gater_integrity:genesis_prev_hash();
resolve_initial_tip(StoreId, StreamId, NextVersion,
{enabled, _Key, ChainStart})
when NextVersion > ChainStart ->
PrevVersion = NextVersion - 1,
PaddedVersion = pad_version(PrevVersion, ?VERSION_PADDING),
Path = ?STREAMS_PATH ++ [StreamId, PaddedVersion],
case khepri:get(StoreId, Path) of
{ok, #event{prev_event_hash = PrevPrevHash} = PrevEvent}
when is_binary(PrevPrevHash) ->
reckon_gater_integrity:compute_chain_hash(PrevEvent, PrevPrevHash);
Other ->
%% Invariant violation: stream's watermark says the
%% predecessor should be integrity-bearing, but we cannot
%% find a usable predecessor. Surface fast — replay
%% won't be able to verify anyway.
erlang:error({integrity_setup_failed,
#{stream_id => StreamId,
looking_for_version => PrevVersion,
chain_start => ChainStart,
got => Other}})
end.
%% @private Compute and attach the integrity fields for one event.
%%
%% Disabled context: pass-through.
%% Enabled: set prev_event_hash to the running tip, compute MAC, then
%% compute the next tip (= chain hash of the just-built event) for the
%% next iteration of the fold.
-spec apply_integrity_if_enabled(
event(),
Tip :: binary() | undefined,
integrity_ctx()
) -> {event(), NextTip :: binary() | undefined}.
apply_integrity_if_enabled(Event, _Tip, disabled) ->
{Event, undefined};
apply_integrity_if_enabled(#event{} = Event, Tip, {enabled, Key, _ChainStart})
when is_binary(Tip) ->
Event1 = Event#event{prev_event_hash = Tip},
Mac = reckon_gater_integrity:compute_event_mac(Event1, Key),
Event2 = Event1#event{mac = Mac},
NextTip = reckon_gater_integrity:compute_chain_hash(Event2, Tip),
{Event2, NextTip}.
%% @private
-spec create_event_record(new_event(), binary(), non_neg_integer(), integer(), integer()) -> event().
create_event_record(Event, StreamId, Version, Timestamp, EpochUs) ->
EventId = maps:get(event_id, Event, generate_event_id()),
EventType = maps:get(event_type, Event),
Data = maps:get(data, Event),
Metadata = maps:get(metadata, Event, #{}),
Tags = maps:get(tags, Event, undefined),
DataContentType = maps:get(data_content_type, Event, ?CONTENT_TYPE_JSON),
MetadataContentType = maps:get(metadata_content_type, Event, ?CONTENT_TYPE_JSON),
#event{
event_id = EventId,
event_type = EventType,
stream_id = StreamId,
version = Version,
data = Data,
metadata = Metadata,
tags = Tags,
timestamp = Timestamp,
epoch_us = EpochUs,
data_content_type = DataContentType,
metadata_content_type = MetadataContentType
}.
%% @private
-spec generate_event_id() -> binary().
generate_event_id() ->
reckon_gater_uuid:to_string(reckon_gater_uuid:v7()).
%% @private
-spec pad_version(non_neg_integer(), pos_integer()) -> binary().
pad_version(Version, Length) ->
VersionStr = integer_to_list(Version),
Padding = Length - length(VersionStr),
PaddedStr = lists:duplicate(Padding, $0) ++ VersionStr,
list_to_binary(PaddedStr).
%% @private Backward-compatible: no verification (for internal callers
%% that have not yet adopted Opts).
-spec do_read(atom(), binary(), non_neg_integer(), pos_integer(), direction()) ->
{ok, [event()]} | {error, term()}.
do_read(StoreId, StreamId, StartVersion, Count, Direction) ->
do_read_with_verify(
StoreId, StreamId, StartVersion, Count, Direction, #{verify => skip_all}).
%% @private Read with verification applied per the resolved options.
-spec do_read_with_verify(
atom(), binary(), non_neg_integer(), pos_integer(), direction(), read_opts()
) -> {ok, [event()]} | {error, term()}.
do_read_with_verify(StoreId, StreamId, StartVersion, Count, Direction, Opts) ->
case exists(StoreId, StreamId) of
false ->
{error, {stream_not_found, StreamId}};
true ->
case read_events(StoreId, StreamId, StartVersion, Count, Direction) of
{ok, Events} ->
case maybe_verify_events(
StoreId, StreamId, StartVersion, Direction,
Events, Opts) of
{ok, _} = Ok ->
Ok;
{integrity_violation, _} = Violation ->
%% Wrap at the public API boundary so
%% callers see {error, _}. Internal
%% helpers and the gater module return
%% the bare tuple.
{error, Violation}
end;
Other ->
Other
end
end.
%% @private Apply the configured verification mode to a fresh result.
%%
%% Backward reads verify the same chain as forward reads by walking
%% the result in forward order (smallest version first) and then
%% reversing the returned list to preserve the caller's requested
%% ordering. The chain semantics are direction-independent: every
%% event's `prev_event_hash` must equal the chain hash of its
%% predecessor, regardless of the order the caller chose to receive
%% them in.
-spec maybe_verify_events(
atom(), binary(), non_neg_integer(), direction(), [event()], read_opts()
) -> {ok, [event()]} | {error, term()}.
maybe_verify_events(StoreId, StreamId, StartVersion, Direction, Events, Opts) ->
Mode = maps:get(verify, Opts, skip_legacy),
case verify_required(StoreId, Direction, Mode) of
false ->
{ok, Events};
true ->
verify_in_direction(
StoreId, StreamId, StartVersion, Direction, Events, Mode)
end.
%% @private Verification runs on integrity-enabled stores regardless
%% of read direction. The verification surface is the same in both
%% directions — the only difference is the result-ordering of the
%% returned events.
-spec verify_required(atom(), direction(), verify_mode()) -> boolean().
verify_required(_StoreId, _Direction, skip_all) -> false;
verify_required(StoreId, _Direction, _Mode) ->
reckon_db_integrity_key:is_enabled(StoreId).
%% @private Dispatch on direction: forward verifies in place; backward
%% reverses to forward order, verifies, and reverses the result back.
-spec verify_in_direction(
atom(), binary(), non_neg_integer(), direction(), [event()], verify_mode()
) -> {ok, [event()]} | {error, term()}.
verify_in_direction(StoreId, StreamId, StartVersion, forward, Events, Mode) ->
verify_events_forward(StoreId, StreamId, StartVersion, Events, Mode);
verify_in_direction(StoreId, StreamId, StartVersion, backward, Events, Mode) ->
%% Backward reads return events highest-version-first. To verify
%% the chain we re-order them lowest-version-first, walk the
%% forward verifier, then reverse the verified list before
%% returning so the caller gets the ordering they asked for.
%% The forward verifier's `StartVersion` is the LOWEST version in
%% the batch, which for a backward read is the LAST event.
ForwardEvents = lists:reverse(Events),
ForwardStart = case ForwardEvents of
[] -> StartVersion;
[#event{version = V} | _] -> V
end,
case verify_events_forward(
StoreId, StreamId, ForwardStart, ForwardEvents, Mode) of
{ok, Verified} ->
{ok, lists:reverse(Verified)};
{integrity_violation, _} = Violation ->
Violation
end.
%% @private Walk events in forward order, verifying each against the
%% running chain tip. Short-circuits on the first integrity_violation.
-spec verify_events_forward(
atom(), binary(), non_neg_integer(), [event()], verify_mode()
) -> {ok, [event()]} | {error, term()}.
verify_events_forward(StoreId, StreamId, StartVersion, Events, Mode) ->
{ok, ChainStart} = reckon_db_chain_watermark:lookup(StoreId, StreamId),
Key = reckon_db_integrity_key:get(StoreId),
InitialTip = resolve_read_initial_tip(
StoreId, StreamId, StartVersion, ChainStart),
verify_events_loop(Events, InitialTip, ChainStart, Key, Mode, StreamId, []).
%% @private
verify_events_loop([], _Tip, _ChainStart, _Key, _Mode, _StreamId, Acc) ->
{ok, lists:reverse(Acc)};
verify_events_loop([Event | Rest], Tip, ChainStart, Key, Mode, StreamId, Acc) ->
case is_legacy_event(Event, ChainStart) of
true ->
handle_legacy(Event, Rest, Tip, ChainStart, Key, Mode, StreamId, Acc);
false ->
handle_integrity(Event, Rest, Tip, ChainStart, Key, Mode, StreamId, Acc)
end.
%% @private An event is legacy if it predates the watermark (or the
%% watermark is absent, meaning the stream never had integrity events).
is_legacy_event(_Event, undefined) -> true;
is_legacy_event(#event{version = V}, ChainStart) when is_integer(ChainStart) ->
V < ChainStart.
handle_legacy(Event, Rest, Tip, ChainStart, Key, strict, StreamId, _Acc) ->
%% Strict mode refuses to return legacy events at all.
_ = Tip, _ = ChainStart, _ = Key, _ = Rest,
{integrity_violation, #{
layer => storage,
stream_id => StreamId,
version => Event#event.version,
kind => missing_integrity,
context => #{detail => legacy_event_under_strict_mode}
}};
handle_legacy(Event, Rest, Tip, ChainStart, Key, Mode, StreamId, Acc) ->
%% skip_legacy: return the legacy event untouched. Emit telemetry
%% so operators can monitor remediation progress.
telemetry:execute(
[reckon, db, read, legacy_event_returned],
#{system_time => erlang:system_time(millisecond)},
#{store_id_hint => StreamId, version => Event#event.version}
),
verify_events_loop(Rest, Tip, ChainStart, Key, Mode, StreamId, [Event | Acc]).
handle_integrity(Event, Rest, undefined, ChainStart, Key, Mode, StreamId, Acc)
when is_integer(ChainStart),
is_record(Event, event),
Event#event.version =:= ChainStart ->
%% Mid-read transition from legacy region to integrity region: the
%% first integrity-bearing event in a stream was written with
%% prev_event_hash = genesis. Seed the running tip accordingly so
%% the verifier has a concrete value to compare against.
Genesis = reckon_gater_integrity:genesis_prev_hash(),
handle_integrity(Event, Rest, Genesis, ChainStart, Key, Mode, StreamId, Acc);
handle_integrity(Event, Rest, Tip, ChainStart, Key, Mode, StreamId, Acc)
when is_binary(Tip) ->
case reckon_gater_integrity:verify_event(Event, Tip, Key) of
ok ->
NextTip = reckon_gater_integrity:compute_chain_hash(Event, Tip),
verify_events_loop(
Rest, NextTip, ChainStart, Key, Mode, StreamId,
[Event | Acc]);
{integrity_violation, _} = Violation ->
Violation
end.
%% @private Resolve the chain tip that the FIRST event in the returned
%% batch should reference as its `prev_event_hash`.
%%
%% - undefined watermark or StartVersion < watermark: legacy-only;
%% the running tip is irrelevant (verification won't be called).
%% - StartVersion == watermark: tip = genesis (this is the first
%% integrity-bearing event in the stream).
%% - StartVersion > watermark: tip = chain_hash of the event at
%% (StartVersion - 1), which must itself be integrity-bearing.
-spec resolve_read_initial_tip(
atom(), binary(), non_neg_integer(), non_neg_integer() | undefined
) -> binary() | undefined.
resolve_read_initial_tip(_StoreId, _StreamId, _StartVersion, undefined) ->
undefined;
resolve_read_initial_tip(_StoreId, _StreamId, StartVersion, ChainStart)
when StartVersion < ChainStart ->
undefined;
resolve_read_initial_tip(_StoreId, _StreamId, StartVersion, ChainStart)
when StartVersion =:= ChainStart ->
reckon_gater_integrity:genesis_prev_hash();
resolve_read_initial_tip(StoreId, StreamId, StartVersion, _ChainStart)
when StartVersion > 0 ->
PrevVersion = StartVersion - 1,
PaddedVersion = pad_version(PrevVersion, ?VERSION_PADDING),
Path = ?STREAMS_PATH ++ [StreamId, PaddedVersion],
case khepri:get(StoreId, Path) of
{ok, #event{prev_event_hash = PrevPrevHash} = PrevEvent}
when is_binary(PrevPrevHash) ->
reckon_gater_integrity:compute_chain_hash(PrevEvent, PrevPrevHash);
_ ->
%% No usable predecessor; legacy regime in practice.
undefined
end.
%% @private
-spec read_events(atom(), binary(), non_neg_integer(), pos_integer(), direction()) ->
{ok, [event()]}.
read_events(StoreId, StreamId, StartVersion, Count, Direction) ->
Versions = calculate_versions(StartVersion, Count, Direction),
Events = lists:filtermap(
fun(Version) ->
PaddedVersion = pad_version(Version, ?VERSION_PADDING),
Path = ?STREAMS_PATH ++ [StreamId, PaddedVersion],
case khepri:get(StoreId, Path) of
{ok, Event} when is_record(Event, event) ->
{true, Event};
{ok, EventMap} when is_map(EventMap) ->
%% Convert map to record if needed
{true, map_to_event(EventMap)};
_ ->
false
end
end,
Versions
),
{ok, Events}.
%% @private
-spec calculate_versions(non_neg_integer(), pos_integer(), direction()) -> [non_neg_integer()].
calculate_versions(StartVersion, Count, forward) ->
lists:seq(StartVersion, StartVersion + Count - 1);
calculate_versions(StartVersion, Count, backward) ->
EndVersion = max(0, StartVersion - Count + 1),
lists:reverse(lists:seq(EndVersion, StartVersion)).
%% @private
-spec map_to_event(map()) -> event().
map_to_event(Map) ->
#event{
event_id = maps:get(event_id, Map, undefined),
event_type = maps:get(event_type, Map, undefined),
stream_id = maps:get(stream_id, Map, undefined),
version = maps:get(version, Map, 0),
data = maps:get(data, Map, #{}),
metadata = maps:get(metadata, Map, #{}),
tags = maps:get(tags, Map, undefined),
timestamp = maps:get(timestamp, Map, 0),
epoch_us = maps:get(epoch_us, Map, 0),
data_content_type = maps:get(data_content_type, Map, ?CONTENT_TYPE_JSON),
metadata_content_type = maps:get(metadata_content_type, Map, ?CONTENT_TYPE_JSON)
}.
%% @private Safe version of lists:nthtail that returns [] when Offset >= length.
-spec safe_nthtail(non_neg_integer(), list()) -> list().
safe_nthtail(0, List) -> List;
safe_nthtail(_, []) -> [];
safe_nthtail(N, [_ | Rest]) -> safe_nthtail(N - 1, Rest).