Current section

Files

Jump to
reckon_db src reckon_db_streams.erl
Raw

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_gater_stream_id for the format rules (moved out
%% of reckon-db in 3.0.0 — protocol contract belongs in the
%% gateway layer, shared with reckon-evoq).
case reckon_gater_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).