Packages
electric_client
0.6.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/ecto_adapter.ex
if Code.ensure_loaded?(Ecto) do
defmodule Electric.Client.EctoAdapter do
@moduledoc false
alias Electric.Client.ShapeDefinition
@behaviour Electric.Client.ValueMapper
def shape!(schema, opts \\ [])
def shape!(schema, opts) when is_atom(schema), do: shape_from_query!(schema, opts)
def shape!(%Ecto.Query{} = query, opts), do: shape_from_query!(query, opts)
def shape!(%Ecto.Changeset{} = changeset, opts), do: shape_from_changeset!(changeset, opts)
def shape!(changeset_fun, opts) when is_function(changeset_fun, 1),
do: shape_from_changeset!(changeset_fun, opts)
@doc false
@spec shape_from_query!(Ecto.Queryable.t()) :: ShapeDefinition.t()
def shape_from_query!(queryable, opts \\ []) do
query = Ecto.Queryable.to_query(queryable)
validate_query!(query)
{table_name, namespace, struct} = table_name(query)
# it's possible that the ecto schema does not contain all the columns in
# the table so, since we know the columns we want, let's specify them
# explicitly
columns = query_columns(query)
where = where(query)
ShapeDefinition.new!(
table_name,
merge_shape_opts(
[namespace: namespace, where: where, columns: columns, parser: {__MODULE__, struct}],
opts
)
)
end
@shape_from_changeset_opts_schema ShapeDefinition.schema_definition()
|> Keyword.drop([:columns])
|> NimbleOptions.new!()
@type shape_from_changeset_opts() :: [
unquote(NimbleOptions.option_typespec(@shape_from_changeset_opts_schema))
]
@spec shape_from_changeset!(
(map() -> Ecto.Changeset.t()) | Ecto.Changeset.t(),
shape_from_changeset_opts()
) :: ShapeDefinition.t()
def shape_from_changeset!(changeset_or_fun, opts \\ [])
def shape_from_changeset!(%Ecto.Changeset{} = changeset, opts) do
generate_shape_from_changeset(changeset, opts)
end
def shape_from_changeset!(changeset_fun, opts) when is_function(changeset_fun, 1) do
case changeset_fun.(%{}) do
%Ecto.Changeset{} = changeset ->
generate_shape_from_changeset(changeset, opts)
invalid ->
raise ArgumentError,
message:
"Changeset function returned #{inspect(invalid)}, was expecting an Ecto.Changeset struct"
end
end
defp generate_shape_from_changeset(%Ecto.Changeset{data: %schema{}} = changeset, opts) do
if !function_exported?(schema, :__schema__, 1),
do:
raise(ArgumentError,
message:
"cannot generate a shape from a schema-less changeset. Use #{inspect(ShapeDefinition)}.new/2"
)
table_name = schema.__schema__(:source)
namespace = schema.__schema__(:prefix)
pks = schema.__schema__(:primary_key)
columns =
[pks, changeset.required, Keyword.keys(changeset.validations)]
|> Enum.concat()
|> Enum.uniq()
|> Enum.map(&schema.__schema__(:field_source, &1))
|> Enum.map(&to_string/1)
ShapeDefinition.new!(
table_name,
merge_shape_opts(
[columns: columns, parser: {__MODULE__, schema}, namespace: namespace],
opts
)
)
end
defp merge_shape_opts(base, overrides) do
Keyword.merge(base, overrides, fn _key, base, over ->
if is_nil(over), do: base, else: over
end)
end
defp table_name(%{
prefix: query_prefix,
from: %{prefix: source_prefix, source: {table_name, struct}}
}) do
{table_name, query_prefix || source_prefix, struct}
end
defp query_columns(%{from: %{source: {_table_name, struct}}}) do
Enum.map(
struct.__schema__(:fields),
&to_string(struct.__schema__(:field_source, &1))
)
end
@doc false
def where(%{wheres: []} = _query) do
nil
end
def where(query) do
%{from: %{source: {table_name, struct}}} = query
{query, bindings, _key} =
Ecto.Query.Planner.plan(query, :all, Ecto.Adapters.Postgres)
{query, _} = Ecto.Query.Planner.normalize(query, :all, Ecto.Adapters.Postgres, 1)
query
|> Electric.Client.EctoAdapter.Postgres.where(
{{quote_table(table_name), quote_table(table_name), struct}, []},
bindings
)
|> case do
[] -> nil
iodata -> IO.iodata_to_binary(iodata)
end
end
defp quote_table(table_name) do
[34, table_name, 34]
end
@impl Electric.Client.ValueMapper
def for_schema(_schema, module) do
fields =
Enum.map(module.__schema__(:fields), fn field ->
{module.__schema__(:field_source, field) |> to_string(), field,
cast_to(module.__schema__(:type, field))}
end)
&map_values(&1, module, fields)
end
defp map_values(values, module, fields) do
struct(
module,
Enum.map(fields, fn {column, field, fun} ->
{field, fun.(Map.get(values, column))}
end)
)
end
defp cast_to(:boolean) do
fn
nil -> nil
"t" -> true
"f" -> false
value -> Ecto.Type.cast!(:boolean, value)
end
end
defp cast_to({:parameterized, {_, _}} = type) do
fn
nil ->
nil
value ->
decoded_value =
case Jason.decode(value) do
{:ok, decoded} -> decoded
{:error, _} -> value
end
case Ecto.Type.load(type, decoded_value) do
{:ok, loaded} -> loaded
:error -> raise Ecto.CastError, type: type, value: value
end
end
end
defp cast_to(type) do
fn
nil -> nil
value -> Ecto.Type.cast!(type, value)
end
end
@problematic_clauses [
joins: "JOIN",
updates: "UPDATE",
order_bys: "ORDER BY",
havings: "HAVING",
group_bys: "GROUP BY",
distinct: "DISTINCT"
]
for {key, name} <- @problematic_clauses do
defp validate_query!(%{unquote(key) => [_ | _]}),
do:
raise(ArgumentError,
message: "Electric does not support streaming queries with #{unquote(name)} clauses"
)
end
defp validate_query!(_query), do: :ok
end
end