defmodule NLdoc.Stream.Message do @moduledoc """ Stream Consumer """ alias NLdoc.Stream.Message.Header @type t :: map() @doc """ Retrieves the stream offset from the provided headers. ## Examples iex> headers = [{"x-stream-offset", :long, 42}] iex> NLdoc.Stream.Message.get_stream_offset(%{headers: headers}) 42 iex> NLdoc.Stream.Message.get_stream_offset(%{headers: []}) nil """ @spec get_stream_offset(t() | [Header.t()]) :: integer() | nil def get_stream_offset(%{headers: headers}), do: headers |> Header.find_value("x-stream-offset", :long) @doc """ Retrieves the stream filter value from the provided headers. ## Examples iex> headers = [{"x-stream-filter-value", :longstr, "value"}] iex> NLdoc.Stream.Message.get_stream_filter_value(%{headers: headers}) "value" iex> NLdoc.Stream.Message.get_stream_filter_value(%{headers: []}) nil """ @spec get_stream_filter_value(t() | [Header.t()]) :: String.t() | nil def get_stream_filter_value(%{headers: headers}), do: headers |> Header.find_value("x-stream-filter-value", :longstr) @doc """ Even though we set `x-stream-filter` to `id`, we still need to check the `x-stream-filter-value` header to be sure, because: ... server-side filtering is probabilistic: messages that do not match the filter value can still be sent to the consumer. The server uses a Bloom filter, a space-efficient probabilistic data structure, where false positives are possible. See: [RabbitMQ Docs on Streams Filtering](https://www.rabbitmq.com/docs/streams#filtering) Returns `true` if the stream filter value in the message equals the given value. ## Examples iex> msg = %{headers: [{"x-stream-filter-value", :longstr,"uuid"}]} iex> NLdoc.Stream.Message.matches_stream_filter?(msg, "uuid") true iex> NLdoc.Stream.Message.matches_stream_filter?(msg, "wrong-id") false """ @spec matches_stream_filter?(t(), String.t()) :: boolean() def matches_stream_filter?(msg, value) when is_map(msg) and is_binary(value), do: value == msg |> get_stream_filter_value() end