Current section
Files
Jump to
Current section
Files
lib/uma_db_client.ex
defmodule UmaDbClient do
@moduledoc """
Elixir gRPC client for UmaDB's DCB (Dynamic Consistency Boundary) service.
## Usage
{:ok, channel} = UmaDbClient.connect("localhost:50051")
# Append events
event = UmaDbClient.Builder.event("OrderPlaced", ["order:1"], Jason.encode!(%{id: 1}))
{:ok, position} = UmaDbClient.append(channel, [event])
# Read events as a lazy stream
{:ok, stream} = UmaDbClient.read(channel)
Enum.each(stream, fn %UmaDb.V1.SequencedEvent{position: pos, event: e} ->
IO.inspect({pos, e.event_type})
end)
## TLS
Pass gRPC credentials via `opts` to `connect/2`:
cred = GRPC.Credential.new(ssl: [cacertfile: "/path/to/ca.pem"])
{:ok, channel} = UmaDbClient.connect("myserver:443", cred: cred)
"""
alias UmaDb.V1.DCB.Stub
alias UmaDb.V1
@type channel :: GRPC.Channel.t()
@type position :: non_neg_integer()
@doc """
Opens a gRPC channel to a UmaDB server.
`target` is a `"host:port"` string. `opts` are passed directly to
`GRPC.Stub.connect/2` (e.g. `cred:` for TLS, `interceptors:` etc.).
"""
@spec connect(String.t(), keyword()) :: {:ok, channel()} | {:error, term()}
def connect(target, opts \\ []) do
GRPC.Stub.connect(target, opts)
end
@doc """
Returns the current head position of the event log.
Returns `{:ok, nil}` when the log is empty.
"""
@spec head(channel()) :: {:ok, position() | nil} | {:error, term()}
def head(channel) do
case Stub.head(channel, %V1.HeadRequest{}) do
{:ok, %V1.HeadResponse{position: pos}} -> {:ok, pos}
{:error, _} = err -> err
end
end
@doc """
Returns the last tracked position for the given source identifier.
Returns `{:ok, nil}` when no tracking record exists for `source`.
"""
@spec get_tracking_info(channel(), String.t()) :: {:ok, position() | nil} | {:error, term()}
def get_tracking_info(channel, source) do
case Stub.get_tracking_info(channel, %V1.TrackingRequest{source: source}) do
{:ok, %V1.TrackingResponse{position: pos}} -> {:ok, pos}
{:error, _} = err -> err
end
end
@doc """
Appends one or more events to the log.
Returns `{:ok, position}` where `position` is the log position after the append.
## Options
- `:condition` — an `UmaDb.V1.AppendCondition` struct for optimistic concurrency
- `:tracking_info` — an `UmaDb.V1.TrackingInfo` struct to update a tracking cursor
"""
@spec append(channel(), [V1.Event.t()], keyword()) :: {:ok, position()} | {:error, term()}
def append(channel, events, opts \\ []) do
request = %V1.AppendRequest{
events: events,
condition: Keyword.get(opts, :condition),
tracking_info: Keyword.get(opts, :tracking_info)
}
case Stub.append(channel, request) do
{:ok, %V1.AppendResponse{position: pos}} -> {:ok, pos}
{:error, _} = err -> err
end
end
@doc """
Returns a lazy `Enumerable` of `UmaDb.V1.SequencedEvent` structs from the log.
The stream is server-driven: each element is received from the server as it is
produced. Iterate with `Enum.to_list/1`, `Stream.each/2`, etc.
Mid-stream gRPC errors raise `GRPC.RPCError`. Wrap iteration in `try/rescue`
if you need to recover from them.
## Options
- `:query` — `UmaDb.V1.Query` to filter by event type and/or tags
- `:start` — starting position (inclusive)
- `:backwards` — if `true`, read in reverse order
- `:limit` — maximum number of events to return
- `:subscribe` — if `true`, keep the stream open for new events after catching up
- `:batch_size` — number of events per server response batch
"""
@spec read(channel(), keyword()) :: {:ok, Enumerable.t()} | {:error, term()}
def read(channel, opts \\ []) do
request = %V1.ReadRequest{
query: Keyword.get(opts, :query),
start: Keyword.get(opts, :start),
backwards: Keyword.get(opts, :backwards),
limit: Keyword.get(opts, :limit),
subscribe: Keyword.get(opts, :subscribe),
batch_size: Keyword.get(opts, :batch_size)
}
case Stub.read(channel, request) do
{:ok, grpc_stream} ->
{:ok, flatten_sequenced_events(grpc_stream, V1.ReadResponse)}
{:error, _} = err ->
err
end
end
@doc """
Returns a lazy `Enumerable` of `UmaDb.V1.SequencedEvent` structs, starting
from the given position and continuing indefinitely as new events arrive.
This is a long-running server-push stream. Blocking iteration (e.g. via
`Enum.each/2`) will not return until the server closes the stream or an
error occurs. Consider running it in a dedicated process.
## Options
- `:query` — `UmaDb.V1.Query` to filter by event type and/or tags
- `:after` — only receive events after this position
- `:batch_size` — number of events per server response batch
"""
@spec subscribe(channel(), keyword()) :: {:ok, Enumerable.t()} | {:error, term()}
def subscribe(channel, opts \\ []) do
request = %V1.SubscribeRequest{
query: Keyword.get(opts, :query),
after: Keyword.get(opts, :after),
batch_size: Keyword.get(opts, :batch_size)
}
case Stub.subscribe(channel, request) do
{:ok, grpc_stream} ->
{:ok, flatten_sequenced_events(grpc_stream, V1.SubscribeResponse)}
{:error, _} = err ->
err
end
end
defp flatten_sequenced_events(grpc_stream, response_module) do
Stream.flat_map(grpc_stream, fn
{:ok, response} ->
Map.fetch!(response, :events)
{:error, reason} ->
raise GRPC.RPCError,
status: GRPC.Status.internal(),
message: "stream error from #{inspect(response_module)}: #{inspect(reason)}"
end)
end
end