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