Packages

Thin client for the Diavasi data plane

Current section

Files

Jump to
Raw

README.md

# Elixir client

[![CI](https://github.com/diavasis/diavasi-elixir/actions/workflows/ci.yml/badge.svg)](https://github.com/diavasis/diavasi-elixir/actions/workflows/ci.yml)
[![Hex.pm](https://img.shields.io/hexpm/v/diavasi_client.svg)](https://hex.pm/packages/diavasi_client)
[![Hex Docs](https://img.shields.io/badge/hex-docs-6e4a7e.svg)](https://hexdocs.pm/diavasi_client)
[![license](https://img.shields.io/github/license/diavasis/diavasi-elixir)](https://github.com/diavasis/diavasi-elixir/blob/main/LICENSE)

`Diavasi.Client` is a supervised consumer of `diavasi.data.v1`. It opens a TLS stream, sends the bearer token, Hello version 1, then JoinGroup. `stream/1` yields batches. The caller acks with `ack/2`. `leave/1` is the clean stop. Dropping the process returns unacked batches to the server. The client stores no cursor and does not dedupe on `record_id`. Reconnect with the same consumer id and the server replays them.

`proto/data.proto` in this repository is the copy of `diavasi.data.v1` from [github.com/diavasis/diavasi](https://github.com/diavasis/diavasi) tag `v0.13.0`. The Hex package is `diavasi_client` version 0.1.0. The Mix app is `:diavasi_client`.

## Install

```elixir
{:diavasi_client, "~> 0.1.0"}
```

## Library

```elixir
{:ok, pid} = Diavasi.Client.start_link(
  addr: "127.0.0.1:7710",
  ca: "/tmp/diavasi-sdk/dataplane-ca.crt",
  token: "sdk-demo",
  group: "demo",
  consumer: "elixir",
  max_in_flight: 1
)

pid
|> Diavasi.Client.stream()
|> Enum.each(fn batch ->
  IO.inspect(batch.batch_id)
  Diavasi.Client.ack(pid, batch.batch_id)
end)

Diavasi.Client.leave(pid)
```

`run/1` consumes `:total` records and acks each batch. `:halt_after` closes after that many acks and does not send Leave. `disconnect/1` closes the stream the same way.

`start_link/1` returns `{:error, reason}` when the handshake fails. A bad token is a gRPC unauthorized error. `next_batch/1` returns `{:ok, batch}`, `:done`, or `{:error, reason}`. `stream/1` raises that reason. A group that is not running is `protocol error 5`. Protocol codes are 1 bad version, 2 bad state, 3 unknown ack, 4 duplicate ack, 5 group not running, 6 unsupported, 7 internal, 8 heartbeat timeout.

Each example below uses the same connection options. Recreate the group before running a second one. The synthetic group `demo` has eight records, so these snippets stop after eight.

```elixir
opts = [
  addr: "127.0.0.1:7710",
  ca: "/tmp/diavasi-sdk/dataplane-ca.crt",
  token: "sdk-demo",
  group: "demo",
  consumer: "elixir",
  max_in_flight: 1,
  total: 8
]
```

`Diavasi.Client` is already a GenServer. The modules below are the process that owns it: they pull batches, print each record, ack, and stop on an error.

### GenServer

```elixir
defmodule Demo.Consumer do
  use GenServer

  def start_link(opts), do: GenServer.start_link(__MODULE__, opts)

  def init(opts) do
    case Diavasi.Client.start_link(opts) do
      {:ok, client} ->
        send(self(), :pull)
        {:ok, %{client: client, seen: 0, total: Keyword.fetch!(opts, :total)}}

      {:error, reason} ->
        {:stop, reason}
    end
  end

  def handle_info(:pull, state) do
    case Diavasi.Client.next_batch(state.client) do
      {:ok, batch} ->
        Enum.each(batch.records, fn record ->
          IO.puts(
            "batch #{batch.batch_id} record #{record.record_id} (#{byte_size(record.payload)} bytes)"
          )
        end)

        :ok = Diavasi.Client.ack(state.client, batch.batch_id)
        seen = state.seen + length(batch.records)

        if seen >= state.total do
          :ok = Diavasi.Client.leave(state.client)
          {:stop, :normal, %{state | seen: seen}}
        else
          send(self(), :pull)
          {:noreply, %{state | seen: seen}}
        end

      :done ->
        {:stop, :normal, state}

      {:error, reason} ->
        IO.puts(:stderr, reason)
        {:stop, {:shutdown, reason}, state}
    end
  end
end

case Demo.Consumer.start_link(opts) do
  {:ok, pid} ->
    ref = Process.monitor(pid)

    receive do
      {:DOWN, ^ref, :process, ^pid, :normal} -> :ok
      {:DOWN, ^ref, :process, ^pid, reason} -> IO.puts(:stderr, inspect(reason))
    end

  {:error, reason} ->
    IO.puts(:stderr, "connect failed: #{inspect(reason)}")
end
```

### Task

`Task.async/1` runs one consume off the caller. `Task.await/2` returns the error string when `start_link/1` or `next_batch/1` fails.

```elixir
task =
  Task.async(fn ->
    case Diavasi.Client.start_link(Keyword.put(opts, :consumer, "elixir-task")) do
      {:ok, pid} ->
        try do
          pid
          |> Diavasi.Client.stream()
          |> Enum.reduce_while(0, fn batch, seen ->
            Enum.each(batch.records, fn record ->
              IO.puts(
                "batch #{batch.batch_id} record #{record.record_id} (#{byte_size(record.payload)} bytes)"
              )
            end)

            :ok = Diavasi.Client.ack(pid, batch.batch_id)
            seen = seen + length(batch.records)
            if seen >= opts[:total], do: {:halt, seen}, else: {:cont, seen}
          end)
        after
          if Process.alive?(pid), do: Diavasi.Client.leave(pid)
        end

      {:error, reason} ->
        {:error, reason}
    end
  end)

case Task.await(task, 30_000) do
  {:error, reason} -> IO.puts(:stderr, inspect(reason))
  seen -> IO.puts("acked #{seen} records")
end
```

`stream/1` raises a protocol or gRPC error from `next_batch/1`. Wrap the `Enum` call in `try/rescue` when that task should return the message instead of crashing:

```elixir
try do
  pid |> Diavasi.Client.stream() |> Enum.each(fn batch ->
    :ok = Diavasi.Client.ack(pid, batch.batch_id)
  end)
rescue
  e in RuntimeError -> IO.puts(:stderr, Exception.message(e))
end
```

### Flow

Add `{:flow, "~> 1.2"}` to the application `mix.exs`. Keep one stage and `max_demand: 1` so the client acks one batch before the next pull.

```elixir
{:ok, pid} = Diavasi.Client.start_link(Keyword.put(opts, :consumer, "elixir-flow"))

pid
|> Diavasi.Client.stream()
|> Stream.transform(0, fn batch, seen ->
  if seen >= opts[:total], do: {:halt, seen}, else: {[batch], seen + length(batch.records)}
end)
|> Flow.from_enumerable(stages: 1, max_demand: 1)
|> Flow.map(fn batch ->
  Enum.each(batch.records, fn record ->
    IO.puts(
      "batch #{batch.batch_id} record #{record.record_id} (#{byte_size(record.payload)} bytes)"
    )
  end)

  :ok = Diavasi.Client.ack(pid, batch.batch_id)
  batch
end)
|> Flow.run()

Diavasi.Client.leave(pid)
```

### GenStage

Add `{:gen_stage, "~> 1.2"}`. The producer pulls batches. The consumer processes each record and acks. Demand is 1, so the next batch waits until that ack returns.

```elixir
defmodule Diavasi.Examples.Producer do
  use GenStage

  def start_link(opts), do: GenStage.start_link(__MODULE__, opts)

  def init(opts) do
    case Diavasi.Client.start_link(opts) do
      {:ok, client} ->
        {:producer, %{client: client, seen: 0, total: Keyword.fetch!(opts, :total)}}

      {:error, reason} ->
        {:stop, reason}
    end
  end

  def handle_demand(_demand, %{seen: seen, total: total} = state) when seen >= total do
    {:noreply, [], state}
  end

  def handle_demand(demand, state) when demand > 0 do
    case Diavasi.Client.next_batch(state.client) do
      {:ok, batch} ->
        event = %{client: state.client, batch: batch}
        {:noreply, [event], %{state | seen: state.seen + length(batch.records)}}

      :done ->
        {:noreply, [], state}

      {:error, reason} ->
        {:stop, reason, state}
    end
  end
end

defmodule Diavasi.Examples.StageConsumer do
  use GenStage

  def start_link(producer), do: GenStage.start_link(__MODULE__, producer)

  def init(producer) do
    {:consumer, nil, subscribe_to: [{producer, max_demand: 1}]}
  end

  def handle_events(events, _from, state) do
    Enum.each(events, fn %{client: client, batch: batch} ->
      Enum.each(batch.records, fn record ->
        IO.puts(
          "batch #{batch.batch_id} record #{record.record_id} (#{byte_size(record.payload)} bytes)"
        )
      end)

      :ok = Diavasi.Client.ack(client, batch.batch_id)
    end)

    {:noreply, [], state}
  end
end

{:ok, producer} =
  Diavasi.Examples.Producer.start_link(Keyword.put(opts, :consumer, "elixir-stage"))

{:ok, _consumer} = Diavasi.Examples.StageConsumer.start_link(producer)
```

### Broadway

Add `{:broadway, "~> 1.0"}` and `{:gen_stage, "~> 1.2"}`. The producer is a GenStage. `handle_message/3` prints each record and acks before the batch is marked successful.

```elixir
defmodule Diavasi.Examples.BroadwayProducer do
  use GenStage

  def start_link(opts), do: GenStage.start_link(__MODULE__, opts)

  def init(opts) do
    case Diavasi.Client.start_link(opts) do
      {:ok, client} ->
        {:producer, %{client: client, seen: 0, total: Keyword.fetch!(opts, :total)}}

      {:error, reason} ->
        {:stop, reason}
    end
  end

  def handle_demand(_demand, %{seen: seen, total: total} = state) when seen >= total do
    {:noreply, [], state}
  end

  def handle_demand(_demand, state) do
    case Diavasi.Client.next_batch(state.client) do
      {:ok, batch} ->
        event = %{client: state.client, batch: batch}
        {:noreply, [event], %{state | seen: state.seen + length(batch.records)}}

      :done ->
        {:noreply, [], state}

      {:error, reason} ->
        {:stop, reason, state}
    end
  end
end

defmodule Diavasi.Examples.Pipeline do
  use Broadway

  def start_link(opts) do
    Broadway.start_link(__MODULE__,
      name: __MODULE__,
      producer: [
        module: {Diavasi.Examples.BroadwayProducer, opts},
        concurrency: 1
      ],
      processors: [default: [concurrency: 1, max_demand: 1]]
    )
  end

  @impl true
  def handle_message(_processor, message, _context) do
    %{client: client, batch: batch} = message.data

    Enum.each(batch.records, fn record ->
      IO.puts(
        "batch #{batch.batch_id} record #{record.record_id} (#{byte_size(record.payload)} bytes)"
      )
    end)

    :ok = Diavasi.Client.ack(client, batch.batch_id)
    message
  end
end

{:ok, pipeline} =
  Diavasi.Examples.Pipeline.start_link(Keyword.put(opts, :consumer, "elixir-broadway"))
```

`handle_failed/2` is the Broadway hook when `handle_message/3` raises. The ack in `handle_message/3` has already run only when the function returns the message. A raise before `ack/2` leaves the batch unacked, and the server replays it on the next join with the same consumer id.

## Run

Start the server from the repo root:

```bash
cargo build -p diavasi
export PATH="$PWD/target/debug:$PATH"
mkdir -p /tmp/diavasi-sdk
diavasi serve --bind 127.0.0.1:7700 --data-bind 127.0.0.1:7710 \
  --store /tmp/diavasi-sdk/state --token sdk-demo
```

In a second terminal, from the repo root, create the group and run one script. Each script consumes all eight records. Pause, delete, create, and start the group before the next one. Delete returns 409 while the group is running, and start on an existing group resumes its cursor, so a later client waits on heartbeats and prints nothing.

```bash
curl -fsS -X POST -H "Authorization: Bearer sdk-demo" \
  http://127.0.0.1:7700/v1/groups/demo/pause || true
curl -fsS -X DELETE -H "Authorization: Bearer sdk-demo" \
  http://127.0.0.1:7700/v1/groups/demo || true
curl -fsS -H "Authorization: Bearer sdk-demo" -H "content-type: application/json" \
  -d '{"group_id":"demo","total_records":8,"payload_size":8,"max_buffer_records":64,"max_buffer_bytes":65536,"batch_max_records":4,"batch_timeout_ms":200,"ordering_contract":"synthetic-u64"}' \
  http://127.0.0.1:7700/v1/groups
curl -fsS -X POST -H "Authorization: Bearer sdk-demo" \
  http://127.0.0.1:7700/v1/groups/demo/start

mix deps.get
mix run examples/consumer.exs
```

The other scripts are `examples/task.exs`, `examples/flow.exs`, `examples/gen_stage.exs`, and `examples/broadway.exs`. Flow, GenStage, and Broadway are dependencies of this Mix project so those scripts compile.

The flag client is still:

```bash
mix diavasi.consume --addr 127.0.0.1:7710 --ca /tmp/diavasi-sdk/dataplane-ca.crt \
  --token sdk-demo --group demo --consumer elixir --total 8
```

Flags: `--addr`, `--ca`, `--token`, `--group`, `--consumer`, `--total`, `--max-in-flight` (default 1), `--halt-after`. The task prints `record_ids` and `batch_ids`.

```bash
docker compose -f clients/docker-compose.yml --profile elixir up --abort-on-container-exit
```

Livebook notes are in [notebooks/demo.livemd](https://github.com/diavasis/diavasi-elixir/blob/v0.1.0/notebooks/demo.livemd).

## Test

`mix coveralls` runs the suite and enforces the floor in `coveralls.json`. Tests tagged `:integration` talk to a Diavasi server and stay out of that run. Include them with:

```bash
mix test --include integration
```

With `DIAVASI_DATA_ADDR`, `DIAVASI_CA`, and `DIAVASI_API_TOKEN` set, the integration test consumes `DIAVASI_TOTAL` records (default 8) from `DIAVASI_GROUP`. Without those variables it passes without connecting.

## Mocks

`Diavasi.Client.put_client/1` swaps in a module that implements `Diavasi.Client.Behaviour`. Add `{:mox, "~> 1.2", only: :test}` to the service and turn the mock on from `test/test_helper.exs`:

```elixir
Mox.defmock(MyApp.DiavasiMock, for: Diavasi.Client.Behaviour)
Diavasi.Client.put_client(MyApp.DiavasiMock)
```

`config :diavasi_client, client: MyApp.DiavasiMock` in `config/test.exs` is the same switch. `Diavasi.Client.reset_client/0` restores the real client. Call `Mox.set_mox_global()` when the code under test runs in another process.

`Diavasi.HTTP.put_client/1` is the lower-level switch. It keeps `Diavasi.Client` and replaces `Mint.HTTP`, which is how this library scripts a Diavasi server. That mock needs `Mox.set_mox_global()` as well, because the GenServer owns the connection.

## Stage 0 bench

The TCP bench task in this tree speaks a different protocol from `data.proto`.

```bash
mise install
mise exec -- mix deps.get
```

```bash
# terminal 1
cargo run -p diavasi --bin diavasi-transport-bench --features transport-bench -- \
  --transport tcp --role server --listen 127.0.0.1:9800 --total-records 200 --smoke

# terminal 2
mise exec -- mix diavasi.bench --connect 127.0.0.1:9800 --total-records 200 \
  --output /tmp/diavasi-bench-results.jsonl
```