Packages
electric_client
0.1.0-dev-13
0.10.3
0.10.2
0.10.1
0.10.1-beta-1
0.10.0
0.9.5-beta-1
0.9.4
0.9.4-beta-1
0.9.3
0.9.2
0.9.1
0.9.0
0.8.3
0.8.3-beta-1
0.8.2
0.8.1
0.8.0
0.8.0-beta-1
0.7.3
0.7.2
0.7.1
0.7.0
0.6.5
0.6.5-beta-5
0.6.5-beta-4
0.6.5-beta-3
0.6.5-beta-2
0.6.5-beta-1
0.6.4
0.6.3
0.6.2
0.6.1
0.6.0
0.5.0
0.5.0-beta-1
0.4.1
0.4.0
0.3.2
0.3.1
0.3.0
0.3.0-beta.4
0.3.0-beta.3
0.3.0-beta.2
0.2.6-pre-1
retired
0.2.6-beta.1
0.2.6-beta.0
0.2.5
0.2.4
0.2.4-pre-8
0.2.4-pre-7
0.2.4-pre-6
0.2.4-pre-5
0.2.4-pre-4
0.2.4-pre-3
0.2.4-pre-2
0.2.4-pre-1
0.2.3
0.2.3-rc-1
0.2.2
0.2.2-rc-1
0.2.1
0.2.1-rc-3
0.2.1-rc-2
0.2.1-rc-1
0.2.0
0.1.2
0.1.1
0.1.0
0.1.0-dev-9
0.1.0-dev-8
0.1.0-dev-7
0.1.0-dev-6
0.1.0-dev-5
0.1.0-dev-4
0.1.0-dev-3
0.1.0-dev-2
0.1.0-dev-17
0.1.0-dev-16
0.1.0-dev-15
0.1.0-dev-14
0.1.0-dev-13
0.1.0-dev-12
0.1.0-dev-11
0.1.0-dev-10
0.1.0-dev
Elixir client for ElectricSQL
Current section
Files
Jump to
Current section
Files
lib/electric/client/mock.ex
defmodule Electric.Client.Mock do
@moduledoc """
Allows for mocking stream messages.
## Usage
``` elixir
{:ok, client} = Electric.Client.Mock.new()
users = [
%{id: 1, name: "User 1"},
%{id: 2, name: "User 2"},
%{id: 3, name: "User 3"}
]
ref = Electric.Client.Mock.async_response(client,
status: 200,
schema: %{id: %{type: "int8"}, name: %{type: "text"}},
last_offset: Client.Offset.first(),
shape_id: "users-1",
body: Electric.Client.Mock.transaction(users, operation: :insert)
)
messages = Electric.Client.stream(client, "users", live: false) |> Enum.into([])
request = Electric.Client.Mock.async_await(ref, timeout = 5_000)
```
"""
alias Electric.Client
alias Electric.Client.Fetch
alias Electric.Client.Offset
@behaviour Electric.Client.Fetch
defmodule Endpoint do
@moduledoc false
use GenServer
def start_link(parent) do
GenServer.start_link(__MODULE__, parent)
end
def init(parent) do
{:ok, %{parent: parent, from: nil, request: nil, response: nil}}
end
def request(pid, request) do
try do
GenServer.call(pid, {:request, request}, :infinity)
catch
:exit, _reason -> {:error, :exit}
end
end
def response(pid, response) do
GenServer.call(pid, {:response, response})
end
def handle_call({:request, request}, from, %{response: nil} = state) do
{:noreply, %{state | from: from, request: request}}
end
def handle_call({:request, request}, _from, %{from: from, response: %{} = response} = state) do
GenServer.reply(from, {:ok, request})
{:reply, {:ok, response}, %{state | from: nil, response: nil}}
end
def handle_call({:response, response}, from, %{from: nil} = state) do
{:noreply, %{state | from: from, response: response}}
end
def handle_call({:response, response}, _from, %{from: from} = state) when not is_nil(from) do
GenServer.reply(from, {:ok, response})
{:reply, {:ok, state.request}, %{state | from: nil, request: nil}}
end
end
@type response_opt ::
{:status, pos_integer()}
| {:headers, %{String.t() => String.t() | [String.t(), ...]}}
| {:body, [map()]}
| {:schema, Client.schema()}
| {:shape_id, Client.shape_id()}
| {:last_offset, Client.Offset.t()}
@type response_opts :: [response_opt()]
@type change_opt ::
{:value, map()}
| {:operation, :insert | :update | :delete}
| {:offset, Client.Offset.t()}
@type change_opts :: [change_opt()]
@type transaction_opt :: {:lsn, non_neg_integer()} | {:up_to_date, boolean()}
@type transaction_opts :: [transaction_opt() | change_opt()]
@impl Electric.Client.Fetch
def fetch(%Fetch.Request{} = request, opts) do
{:ok, endpoint} = Keyword.fetch(opts, :endpoint)
Endpoint.request(endpoint, request)
end
@doc """
Create a new mock client, linked to the `parent` process, `self()` by default.
"""
@spec new(pid()) :: {:ok, Client.t()}
def new(parent \\ self()) do
{:ok, endpoint} = Endpoint.start_link(parent)
Client.new(
base_url: "http://mock.electric",
fetch: {Electric.Client.Mock, endpoint: endpoint}
)
end
@spec response(Client.t(), response_opts()) :: {:ok, Fetch.Request.t()}
def response(%Client{fetch: {__MODULE__, opts}}, response) when is_list(response) do
{:ok, endpoint} = Keyword.fetch(opts, :endpoint)
Endpoint.response(endpoint, build_response(response))
end
@spec response(Client.t(), response_opts()) :: reference()
def async_response(client, response) do
parent = self()
ref = make_ref()
Task.start_link(fn ->
{:ok, request} = response(client, response)
send(parent, {__MODULE__, ref, request})
end)
ref
end
@spec async_await(reference(), pos_integer() | :infinity) :: Fetch.Request.t()
def async_await(ref, timeout \\ 5000) do
receive do
{__MODULE__, ^ref, request} -> request
after
timeout ->
raise "No request received"
end
end
@spec up_to_date() :: map()
def up_to_date(_opts \\ []) do
%{"headers" => %{"control" => "up-to-date"}}
end
@doc """
Wrap the given `values` in `Client.Messages.ChangeMessage` structs at the
given `:lsn`.
By default this will append an `up-to-date` control message to the end of the
liist of changes. Pass `up_to_date: false` to disable this.
"""
@spec transaction(values :: [map()], transaction_opts()) :: [map()]
def transaction(values, opts \\ []) do
tx_offset = Keyword.get(opts, :lsn, 0)
up_to_date =
if Keyword.get(opts, :up_to_date, true) do
[up_to_date()]
else
[]
end
values
|> Enum.with_index()
|> Enum.map(fn {value, op_offset} ->
opts
|> Keyword.merge(value: value, offset: Offset.new(tx_offset, op_offset))
|> change()
end)
|> Enum.concat(up_to_date)
end
@spec change(change_opts()) :: map()
def change(opts) do
%{
value: opts[:value] || %{},
headers: change_headers(opts[:operation] || :insert),
offset: Offset.to_string(opts[:offset] || Offset.first())
}
|> jsonify()
end
defp change_headers(operation) do
jsonify(%{operation: to_string(operation)})
end
defp build_response(opts) do
%Fetch.Response{
status: Keyword.get(opts, :status, 200),
headers: headers(opts[:headers] || []),
body: jsonify(opts[:body] || []),
schema: Keyword.get(opts, :schema, nil),
shape_id: Keyword.get(opts, :shape_id, nil),
last_offset: Keyword.get(opts, :last_offset, nil)
}
end
@spec headers([
{:shape_id, Client.shape_id()}
| {:last_offset, Client.Offset.t()}
| {:schema, Client.schema()}
]) :: %{String.t() => [String.t()]}
def headers(args) do
%{}
|> put_optional_header("electric-shape-id", args[:shape_id])
|> put_optional_header(
"electric-chunk-last-offset",
args[:last_offset],
&Client.Offset.to_string/1
)
|> put_optional_header("electric-schema", args[:schema], &Jason.encode!/1)
end
defp put_optional_header(headers, header, value, encoder \\ & &1)
defp put_optional_header(headers, _header, nil, _encoder) do
headers
end
defp put_optional_header(headers, header, value, encoder) do
Map.put(headers, header, [encoder.(value)])
end
defp jsonify(value) do
value |> Jason.encode!() |> Jason.decode!()
end
end