Rheo.Query (rheo v1.0.0)

Copy Markdown View Source

Portable query against a Rheo event stream.

Backends translate this struct into native filters. Prefer Rheo.query/1, Rheo.query/2, Rheo.query_page/2, or Rheo.stream_query/2 rather than constructing backend-specific documents (Mongo filters, SQL, …).

Part of the SemVer-frozen surface (ADR 029 / ADR 030). Field meanings and page-cursor shape (%{partition => after_sequence}, legacy %{after_sequence: n} still accepted) will not change without a major version.

Fields

  • :stream — stream name (required)
  • :where — keyword filters (type, key, partition, payload/metadata fields)
  • :from / :to — optional DateTime bounds on timestamp
  • :after_sequence — exclusive lower bound on sequence (like Rheo.read/2 :after)
  • :after_sequences — per-partition exclusive lower bounds (%{partition => seq}), used by page cursors on multi-partition streams
  • :until_sequence — inclusive upper bound on sequence
  • :order_by — [{field, :asc | :desc}] (default [sequence: :asc]); 1 / -1 are accepted and normalized

  • :limit — max events (default 100)
  • :cursor — opaque page cursor from %Rheo.Page{next_cursor} (sequence-asc pages)

Building queries

Keyword options at the top level are merged into :where (except recognized control keys). These are equivalent:

iex> a = Rheo.Query.new("market-events", type: "curve_update", currency: "EUR")
iex> b = Rheo.Query.new("market-events", where: [type: "curve_update", currency: "EUR"])
iex> a.where == b.where
true

Explicit :where plus flat filters:

iex> q = Rheo.Query.new("market-events", where: [type: "curve_update"], currency: "USD", limit: 5)
iex> {q.where[:type], q.where[:currency], q.limit}
{"curve_update", "USD", 5}

Common filters

Event fields:

iex> q = Rheo.Query.new("market-events", type: "curve_update", key: "EUR-EURIBOR-6M", partition: 0)
iex> {q.where[:type], q.where[:key], q.where[:partition]}
{"curve_update", "EUR-EURIBOR-6M", 0}

Payload fields (any atom other than the reserved ones becomes a payload match, e.g. currency, curve, price):

iex> q = Rheo.Query.new("s", currency: "EUR", curve: "EUR-EURIBOR-6M")
iex> {q.where[:currency], q.where[:curve]}
{"EUR", "EUR-EURIBOR-6M"}

Lineage / metadata (see Rheo.Event.Lineage):

iex> q = Rheo.Query.new("s", correlation_id: "trade-42", producer: "pricing-v3", schema: "curve_update")
iex> {q.where[:correlation_id], q.where[:producer], q.where[:schema]}
{"trade-42", "pricing-v3", "curve_update"}

Sequence and time ranges

iex> q = Rheo.Query.new("s", after_sequence: 100, until_sequence: 200, limit: 50)
iex> {q.after_sequence, q.until_sequence, q.limit}
{100, 200, 50}

iex> from = ~U[2026-01-01 00:00:00.000Z]
iex> to = ~U[2026-01-31 23:59:59.000Z]
iex> q = Rheo.Query.new("s", from: from, to: to, order_by: [sequence: :desc])
iex> {q.from, q.to, q.order_by}
{~U[2026-01-01 00:00:00.000Z], ~U[2026-01-31 23:59:59.000Z], [sequence: :desc]}

Ordering and limits

iex> Rheo.Query.new("s").order_by
[sequence: :asc]

iex> Rheo.Query.new("s", order_by: [timestamp: :desc, sequence: :desc], limit: 10).limit
10

Pagination cursors

Rheo.query_page/2 returns %Rheo.Page{next_cursor: …}. Pass that map back as :cursor (usually with ascending sequence order):

iex> q = Rheo.Query.new("s", limit: 100, cursor: %{0 => 50})
iex> q.cursor
%{0 => 50}

iex> q = Rheo.Query.apply_cursor(Rheo.Query.new("s", after_sequence: 10, cursor: %{after_sequence: 50}))
iex> {q.after_sequence, q.cursor}
{50, nil}

iex> q = Rheo.Query.apply_cursor(Rheo.Query.new("s", cursor: %{0 => 3, 1 => 2}))
iex> {q.after_sequences, q.cursor}
{%{0 => 3, 1 => 2}, nil}

Running queries (facade)

These call the configured backend (not doctested here):

Rheo.query("market-events", type: "curve_update", currency: "EUR")

Rheo.query(%Rheo.Query{
  stream: "market-events",
  where: [type: "curve_update", currency: "EUR"],
  after_sequence: 1_000,
  order_by: [sequence: :asc],
  limit: 50
})

{:ok, page} = Rheo.query_page("market-events", type: "curve_update", limit: 100)
{:ok, page2} = Rheo.query_page("market-events", type: "curve_update", limit: 100, cursor: page.next_cursor)

Rheo.stream_query("market-events", type: "curve_update", limit: 100)
|> Enum.take(250)

Named instance:

Rheo.query("market-events", type: "curve_update", rheo: MyRheo)

Non-goals

Query never ACKs, leases, or deletes events. Consumer progress is Rheo.fetch/3 / replay — see Rheo.replay/3 and ADR 015.

Summary

Types

Sort direction for order_by.

Per-partition exclusive lower bounds (%{partition => after_sequence}).

t()

Portable query against a stream.

Functions

Exclusive sequence lower bound for one partition.

Applies an opaque page :cursor and clears :cursor.

Builds a query from a stream name and keyword options.

Whether event is strictly after the query's exclusive sequence cursor.

Types

order_dir()

@type order_dir() :: :asc | :desc

Sort direction for order_by.

partition_cursors()

@type partition_cursors() :: %{optional(non_neg_integer()) => non_neg_integer()}

Per-partition exclusive lower bounds (%{partition => after_sequence}).

t()

@type t() :: %Rheo.Query{
  after_sequence: non_neg_integer() | nil,
  after_sequences: partition_cursors() | nil,
  cursor: map() | nil,
  from: DateTime.t() | nil,
  limit: pos_integer(),
  order_by: [{atom(), order_dir()}],
  stream: String.t(),
  to: DateTime.t() | nil,
  until_sequence: pos_integer() | nil,
  where: keyword()
}

Portable query against a stream.

See the module documentation for field meanings and examples.

Functions

after_sequence_for(query, partition)

@spec after_sequence_for(t(), non_neg_integer()) :: non_neg_integer() | nil

Exclusive sequence lower bound for one partition.

Prefers :after_sequences when present, otherwise :after_sequence.

apply_cursor(query)

@spec apply_cursor(t()) :: t()

Applies an opaque page :cursor and clears :cursor.

Composite cursors (%{partition => after_sequence}) populate :after_sequences. The legacy %{after_sequence: n} shape still sets a global :after_sequence lower bound.

Used by backends and Rheo.query_page/2. Prefer passing cursor: into Rheo.Query.new/2 or Rheo.query_page/2 rather than calling this directly.

Examples

iex> q = Rheo.Query.new("s", after_sequence: 10, cursor: %{after_sequence: 50})
iex> applied = Rheo.Query.apply_cursor(q)
iex> {applied.after_sequence, applied.cursor}
{50, nil}

iex> q = Rheo.Query.apply_cursor(Rheo.Query.new("s", cursor: %{0 => 3, "1" => 2}))
iex> q.after_sequences
%{0 => 3, 1 => 2}

iex> Rheo.Query.apply_cursor(Rheo.Query.new("s")).cursor
nil

new(stream, opts \\ [])

@spec new(String.t(), keyword()) :: t()

Builds a query from a stream name and keyword options.

Recognized options: :where, :from, :to, :after_sequence, :after_sequences, :until_sequence, :order_by, :limit, :cursor, plus flat filters (:type, :key, :currency, :correlation_id, …) merged into :where. Options :rheo and :sort are ignored (use :order_by; pass :rheo to Rheo.query/2 instead).

Examples

iex> q = Rheo.Query.new("market-events", type: "curve_update", currency: "EUR", limit: 10)
iex> {q.stream, q.where[:type], q.where[:currency], q.limit}
{"market-events", "curve_update", "EUR", 10}

iex> q = Rheo.Query.new("orders", after_sequence: 5, until_sequence: 9, order_by: [sequence: :desc])
iex> {q.after_sequence, q.until_sequence, q.order_by}
{5, 9, [sequence: :desc]}

iex> q = Rheo.Query.new("s", where: [type: "x"], key: "k1", partition: 0)
iex> q.where
[type: "x", key: "k1", partition: 0]

iex> Rheo.Query.new("s", rheo: MyRheo, sort: %{"sequence" => -1}).order_by
[sequence: :asc]

past_after?(map, query)

@spec past_after?(map(), t()) :: boolean()

Whether event is strictly after the query's exclusive sequence cursor.