Packages
electric
1.6.5
1.7.8
1.7.7
1.7.6
1.7.5
1.7.4
1.7.3
1.7.2
1.7.1
1.7.0
1.6.10
1.6.9
1.6.8
1.6.7
1.6.6
1.6.5
1.6.4
1.6.3
1.6.2
1.6.1
1.6.0
1.5.1
1.5.0
1.4.16
1.4.16-beta-1
1.4.15
1.4.14
1.4.13
1.4.12
1.4.11
1.4.10
1.4.8
1.4.7
1.4.6
1.4.5
1.4.4
1.4.3
1.4.2
1.4.1
1.4.0
1.3.4
1.3.3
1.3.2
1.2.4
1.2.3
1.2.2
1.2.1
1.2.0
1.1.14
1.1.13
1.1.12
1.1.11
1.1.10
1.1.9
1.1.8
1.1.7
1.1.6
retired
1.1.5
retired
1.1.4
retired
1.1.3
retired
1.1.2
1.1.1
1.1.0
1.0.24
1.0.23
1.0.22
1.0.21
1.0.20
1.0.19
1.0.18
1.0.17
1.0.15
1.0.13
1.0.12
1.0.11
1.0.10
1.0.9
1.0.5
1.0.4
1.0.3
1.0.2
1.0.1
1.0.0
1.0.0-beta.23
1.0.0-beta.22
1.0.0-beta.20
1.0.0-beta.19
1.0.0-beta.18
1.0.0-beta.17
1.0.0-beta.16
1.0.0-beta.15
1.0.0-beta.14
1.0.0-beta.13
1.0.0-beta.12
1.0.0-beta.11
1.0.0-beta.10
1.0.0-beta.9
1.0.0-beta.8
1.0.0-beta.7
1.0.0-beta.6
1.0.0-beta.5
1.0.0-beta.4
1.0.0-beta.3
1.0.0-beta.2
1.0.0-beta.1
0.9.5
0.9.4
0.9.3
0.9.2
0.9.1
0.9.0
0.8.1
0.8.0
0.7.7
0.7.6
0.7.5
0.7.4
0.7.3
0.7.2
0.7.1
0.7.0
0.6.3
0.6.2
0.6.1
0.5.2
0.4.4
Postgres sync engine. Sync little subsets of your Postgres data into local apps and services.
Current section
Files
Jump to
Current section
Files
lib/electric/shapes/api/encoder.ex
defmodule Electric.Shapes.Api.Encoder do
@callback message(term()) :: Enum.t()
@callback log(term()) :: Enum.t()
@callback subset(term()) :: Enum.t()
def validate!(impl) do
case impl do
module when is_atom(module) ->
# just assume that the module implements the behaviour
module
invalid ->
raise ArgumentError,
message:
"Expected a module that implements the #{inspect(__MODULE__)} protocol. Got #{inspect(invalid)}"
end
end
end
defmodule Electric.Shapes.Api.Encoder.JSON do
@behaviour Electric.Shapes.Api.Encoder
@impl Electric.Shapes.Api.Encoder
def message(message) when is_binary(message) do
[message]
end
def message(term) do
Stream.map([term], &Jason.encode_to_iodata!/1)
end
@impl Electric.Shapes.Api.Encoder
# the log is streamed from storage as a stream of json-encoded messages
def log(item_stream) do
item_stream |> Stream.map(&ensure_json/1) |> to_json_stream()
end
@impl Electric.Shapes.Api.Encoder
def subset({metadata, item_stream}) do
metadata =
metadata
|> Map.update!(:xmin, &to_string/1)
|> Map.update!(:xmax, &to_string/1)
|> Map.update!(:xip_list, &Enum.map(&1, fn xid -> to_string(xid) end))
Stream.concat([
[
~s|{"metadata":|,
Jason.encode_to_iodata!(metadata),
~s|, "data": |
],
to_json_stream(item_stream),
[~s|}|]
])
end
defp ensure_json(json) when is_binary(json) do
json
end
defp ensure_json(term) do
Jason.encode_to_iodata!(term)
end
@json_list_start "["
@json_list_end "]"
@json_item_separator ","
defp to_json_stream(items) do
Stream.concat([
[@json_list_start],
Stream.intersperse(items, @json_item_separator),
[@json_list_end]
])
|> Stream.chunk_every(500)
end
end
defmodule Electric.Shapes.Api.Encoder.SSE do
@behaviour Electric.Shapes.Api.Encoder
@impl Electric.Shapes.Api.Encoder
def log(item_stream) do
# Note that, unlike the JSON log encoder, this doesn't currently use
# `Stream.chunk_every/1`.
#
# This is because it's only handling live events and is usually used
# for small updates (the point of enabling SSE mode is to avoid request
# overhead when consuming small changes).
item_stream
|> Stream.flat_map(&message/1)
end
@impl Electric.Shapes.Api.Encoder
def message(message) do
["data: ", ensure_json(message), "\n\n"]
end
@impl Electric.Shapes.Api.Encoder
def subset(_), do: raise("Subset encoding not supported for SSE")
defp ensure_json(json) when is_binary(json) do
json
end
defp ensure_json(term) do
Jason.encode_to_iodata!(term)
end
end
defmodule Electric.Shapes.Api.Encoder.Term do
@behaviour Electric.Shapes.Api.Encoder
@impl Electric.Shapes.Api.Encoder
def message(message) when is_binary(message) do
[Jason.decode!(message)]
end
def message(term) do
[term]
end
@impl Electric.Shapes.Api.Encoder
# the log is streamed from storage as a stream of json-encoded messages
def log(item_stream) do
Stream.map(item_stream, &maybe_decode_json!/1)
end
@impl Electric.Shapes.Api.Encoder
def subset({metadata, item_stream}) do
{metadata, Stream.map(item_stream, &maybe_decode_json!/1)}
end
defp maybe_decode_json!(json) when is_binary(json) do
Jason.decode!(json)
end
defp maybe_decode_json!(term) do
term
end
end