Packages
electric
1.7.1
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.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(), keyword()) :: {:ok, term()} | {:error, {%Func{}, [term()]}}
def execute(%Expr{} = tree, ref_values, opts \\ []) do
ctx = %{refs: ref_values, subquery_member?: Keyword.get(opts, :subquery_member?)}
execute_node(tree.eval, ctx)
catch
{:could_not_compute, func} -> {:error, func}
end
defp execute_node(%Const{value: value}, _ctx), do: {:ok, value}
defp execute_node(
%Ref{path: ["$sublink", _] = path},
%{refs: refs, subquery_member?: subquery_member?}
)
when is_function(subquery_member?, 2) do
{:ok, Map.get(refs, path, {:subquery_ref, path})}
end
defp execute_node(%Ref{path: path}, %{refs: refs}), do: {:ok, Map.fetch!(refs, path)}
defp execute_node(%Array{elements: elements}, ctx) do
Utils.map_while_ok(elements, &execute_node(&1, ctx))
end
defp execute_node(%RowExpr{elements: elements}, ctx) do
with {:ok, elements} <- Utils.map_while_ok(elements, &execute_node(&1, ctx)) do
{:ok, List.to_tuple(elements)}
end
end
defp execute_node(
%Func{name: "coalesce", variadic_arg: 0, args: [%Array{elements: elements}]},
ctx
) do
execute_coalesce(elements, ctx)
end
defp execute_node(
%Func{name: "sublink_membership_check"} = func,
%{subquery_member?: subquery_member?} = ctx
)
when is_function(subquery_member?, 2) do
with {:ok, [value, {:subquery_ref, path}]} <-
Utils.map_while_ok(func.args, &execute_node(&1, ctx)) do
{:ok,
try do
subquery_member?.(path, value)
rescue
_ -> throw({:could_not_compute, %{func | args: [value, {:subquery_ref, path}]}})
end}
end
end
defp execute_node(%Func{args: args} = func, ctx) do
with {:ok, args} <- Utils.map_while_ok(args, &execute_node(&1, ctx)) do
apply_func(func, args)
end
end
defp execute_coalesce([], _ctx), do: {:ok, nil}
defp execute_coalesce([arg | rest], ctx) do
case execute_node(arg, ctx) do
{:ok, nil} -> execute_coalesce(rest, ctx)
{:ok, value} -> {:ok, value}
end
end
defp apply_func(%Func{strict?: false} = func, args) do
# For a non-strict function, we don't care about nil values in the arguments
{:ok, try_apply(func, args)}
end
defp apply_func(%Func{strict?: true, variadic_arg: vararg_position} = func, 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