Packages
electric_client
0.3.2
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/message.ex
defmodule Electric.Client.Message do
@moduledoc false
alias Electric.Client
alias Electric.Client.Offset
defmodule Headers do
defstruct [:operation, :relation, :handle]
@type operation :: :insert | :update | :delete
@type relation :: [String.t(), ...]
@type t :: %__MODULE__{
operation: operation(),
relation: relation(),
handle: Client.shape_handle()
}
@doc false
def from_message(msg, handle) do
%{"operation" => operation} = msg
%__MODULE__{
relation: msg["relation"],
operation: parse_operation(operation),
handle: handle
}
end
defp parse_operation("insert"), do: :insert
defp parse_operation("update"), do: :update
defp parse_operation("delete"), do: :delete
def insert(opts \\ []), do: struct(%__MODULE__{operation: :insert}, opts)
def update(opts \\ []), do: struct(%__MODULE__{operation: :update}, opts)
def delete(opts \\ []), do: struct(%__MODULE__{operation: :delete}, opts)
end
defmodule ControlMessage do
defstruct [:control, :global_last_seen_lsn, :handle, :request_timestamp]
@type control :: :must_refetch | :up_to_date
@type t :: %__MODULE__{
control: control(),
global_last_seen_lsn: pos_integer(),
handle: Client.shape_handle(),
request_timestamp: DateTime.t()
}
def from_message(
%{"headers" => %{"control" => control} = headers},
handle
) do
%__MODULE__{
control: control_atom(control),
global_last_seen_lsn: global_last_seen_lsn(headers),
handle: handle
}
end
def from_message(
%{headers: %{control: control} = headers},
handle
) do
%__MODULE__{
control: control_atom(control),
global_last_seen_lsn: global_last_seen_lsn(headers),
handle: handle
}
end
defp control_atom("must-refetch"), do: :must_refetch
defp control_atom("up-to-date"), do: :up_to_date
defp control_atom(a) when is_atom(a), do: a
defp global_last_seen_lsn(headers) do
parse_lsn(headers["global_last_seen_lsn"] || headers[:global_last_seen_lsn])
end
defp parse_lsn(nil), do: nil
defp parse_lsn(lsn) when is_binary(lsn), do: String.to_integer(lsn)
defp parse_lsn(lsn) when is_integer(lsn), do: lsn
def up_to_date, do: %__MODULE__{control: :up_to_date}
def must_refetch, do: %__MODULE__{control: :must_refetch}
end
defmodule ChangeMessage do
defstruct [:key, :value, :headers, :request_timestamp]
@type key :: String.t()
@type value :: %{String.t() => binary()}
@type t :: %__MODULE__{
key: key(),
value: value(),
headers: Headers.t(),
request_timestamp: DateTime.t()
}
require Logger
def from_message(msg, handle, value_mapping_fun) do
%{
"headers" => headers,
"value" => raw_value
} = msg
value =
try do
value_mapping_fun.(raw_value)
rescue
exception ->
Logger.error(
"Unable to cast field values: #{Exception.format(:error, exception, __STACKTRACE__)}"
)
reraise exception, __STACKTRACE__
end
%__MODULE__{
key: msg["key"],
headers: Headers.from_message(headers, handle),
value: value
}
end
end
defmodule ResumeMessage do
@moduledoc """
Emitted by the synchronisation stream before terminating early. If passed
as an option to [`Client.stream/3`](`Electric.Client.stream/3`) allows for
resuming a shape stream at the given point.
E.g.
```
# passing `live: false` means the stream will terminate once it receives an
# `up-to-date` message from the server
messages = Electric.Client.stream(client, "my_table", live: false) |> Enum.to_list()
%ResumeMessage{} = resume = List.last(messages)
# `stream` will resume from whatever point the initial one finished
stream = Electric.Client.stream(client, "my_table", resume: resume)
```
"""
@enforce_keys [:shape_handle, :offset, :schema]
defstruct [:shape_handle, :offset, :schema]
@type t :: %__MODULE__{
shape_handle: Client.shape_handle(),
offset: Offset.t(),
schema: Client.schema()
}
end
defguard is_insert(msg) when is_struct(msg, ChangeMessage) and msg.headers.operation == :insert
def parse(%{"value" => _} = msg, shape_handle, value_mapper_fun) do
[ChangeMessage.from_message(msg, shape_handle, value_mapper_fun)]
end
def parse(%{"headers" => %{"control" => _}} = msg, shape_handle, _value_mapper_fun) do
[ControlMessage.from_message(msg, shape_handle)]
end
def parse(%{headers: %{control: _}} = msg, shape_handle, _value_mapper_fun) do
[ControlMessage.from_message(msg, shape_handle)]
end
def parse("", _handle, _value_mapper_fun) do
[]
end
end