Packages
electric
1.4.13
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/replication/eval/runner.ex
defmodule Electric.Replication.Eval.Runner do
require Logger
alias Electric.Replication.Eval.Walker
alias Electric.Utils
alias Electric.Replication.Eval.Expr
alias Electric.Replication.Eval.Env
alias Electric.Replication.Eval.Parser.{Const, Func, Ref, RowExpr, Array}
@doc """
Generate a ref values object based on the record and a given table name
## Examples
iex> used_refs = %{["id"] => :int8, ["created_at"] => :timestamp}
iex> record_to_ref_values(used_refs, %{"id" => "80", "created_at" => "2020-01-01T11:00:00Z"})
%{
["id"] => 80,
["created_at"] => ~N[2020-01-01 11:00:00]
}
"""
@spec record_to_ref_values(Expr.used_refs(), map(), Env.t()) :: {:ok, map()} | :error
def record_to_ref_values(used_refs, record, env \\ Env.new()) do
used_refs
# Keep only used refs that are pointing to current table
|> Enum.filter(&match?({[_], _}, &1))
|> Enum.reduce_while({:ok, %{}}, fn {[key] = path, type}, {:ok, acc} ->
value = record[key]
case Env.parse_const(env, value, type) do
{:ok, value} ->
{:cont, {:ok, Map.put(acc, path, value)}}
:error ->
Logger.warning(
"Could not parse #{inspect(value)} as #{inspect(type)} while casting #{inspect(key)}"
)
{:halt, :error}
end
end)
end
def execute_for_record(expr, record, extra_refs \\ %{}) do
with {:ok, ref_values} <- record_to_ref_values(expr.used_refs, record),
{:ok, evaluated} <- execute(expr, Map.merge(ref_values, extra_refs)) do
{:ok, evaluated}
end
end
@doc """
Run a PG function parsed by `Electric.Replication.Eval.Parser` based on the inputs
"""
@spec execute(Expr.t(), map()) :: {:ok, term()} | {:error, {%Func{}, [term()]}}
def execute(%Expr{} = tree, ref_values) do
Walker.fold(tree.eval, &do_execute/3, ref_values)
catch
{:could_not_compute, func} -> {:error, func}
end
defp do_execute(%Const{value: value}, _, _), do: {:ok, value}
defp do_execute(%Ref{path: path}, _, refs), do: {:ok, Map.fetch!(refs, path)}
defp do_execute(%Array{}, %{elements: elements}, _), do: {:ok, elements}
defp do_execute(%RowExpr{}, %{elements: elements}, _), do: {:ok, List.to_tuple(elements)}
defp do_execute(%Func{strict?: false} = func, %{args: args}, _) do
# For a non-strict function, we don't care about nil values in the arguments
{:ok, try_apply(func, args)}
end
defp do_execute(%Func{strict?: true, variadic_arg: vararg_position} = func, %{args: args}, _) do
has_nils? =
case vararg_position do
nil -> Enum.any?(args, &is_nil/1)
pos -> Enum.any?(Enum.at(args, pos), &is_nil/1) or Enum.any?(args, &is_nil/1)
end
if has_nils? do
{:ok, nil}
else
{:ok, try_apply(func, args)}
end
end
defp try_apply(
%Func{implementation: impl, map_over_array_in_pos: map_over_array_in_pos} = func,
args
) do
case {impl, map_over_array_in_pos} do
{{module, fun}, nil} ->
apply(module, fun, args)
{fun, nil} ->
apply(fun, args)
{{module, function}, 0} ->
Utils.deep_map(hd(args), &apply(module, function, [&1 | tl(args)]))
{function, 0} ->
Utils.deep_map(hd(args), &apply(function, [&1 | tl(args)]))
{{module, function}, pos} ->
Utils.deep_map(
Enum.at(args, pos),
&apply(module, function, List.replace_at(args, pos, &1))
)
{function, pos} ->
Utils.deep_map(Enum.at(args, pos), &apply(function, List.replace_at(args, pos, &1)))
end
rescue
_ ->
# Anything could have gone wrong here
throw({:could_not_compute, %{func | args: args}})
end
end