Packages
electric
1.7.2
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/consumer/subqueries/move_queue.ex
defmodule Electric.Shapes.Consumer.Subqueries.MoveQueue do
@moduledoc """
Multi-dependency move queue. Tracks move_in/move_out operations per dependency index,
with deduplication and redundancy elimination scoped per dependency.
Move-outs from any dependency are drained before move-ins from any dependency.
Each per-dep batch also accumulates the upstream Postgres transaction ids that
contributed to it, so move-in/move-out broadcasts can carry `txids` for client
attribution (mirroring `Electric.LogItems.from_change/4`).
"""
@type move_value() :: {term(), term()}
@type txid() :: pos_integer()
@type entry() :: {[move_value()], MapSet.t(txid())}
# move_out/move_in are maps from dep_index to {[move_value], MapSet<txid>}
defstruct move_out: %{}, move_in: %{}
@type t() :: %__MODULE__{
move_out: %{non_neg_integer() => entry()},
move_in: %{non_neg_integer() => entry()}
}
@type batch_kind() :: :move_out | :move_in
@type batch() :: {batch_kind(), non_neg_integer(), [move_value()], [txid()]}
@spec new() :: t()
def new, do: %__MODULE__{}
@spec length(t()) :: non_neg_integer()
def length(%__MODULE__{move_out: move_out, move_in: move_in}) do
count_values(move_out) + count_values(move_in)
end
defp count_values(map) do
Enum.reduce(map, 0, fn {_, {vs, _}}, acc -> acc + Kernel.length(vs) end)
end
@doc """
Enqueue a materializer payload for a specific dependency.
`dep_view` is the current view for this dependency, used for redundancy elimination.
The payload may include a `:txids` key listing the upstream xids that produced
the moves. Those xids are unioned with any already accumulated for this dep.
"""
@spec enqueue(t(), non_neg_integer(), map() | keyword(), MapSet.t()) :: t()
def enqueue(%__MODULE__{} = queue, dep_index, payload, %MapSet{} = dep_view)
when is_map(payload) or is_list(payload) do
payload = Map.new(payload)
new_txids = payload |> Map.get(:txids, []) |> MapSet.new()
{existing_outs, existing_out_txids} = Map.get(queue.move_out, dep_index, {[], MapSet.new()})
{existing_ins, existing_in_txids} = Map.get(queue.move_in, dep_index, {[], MapSet.new()})
ops =
Enum.map(existing_outs, &{:move_out, &1}) ++
Enum.map(existing_ins, &{:move_in, &1}) ++
payload_to_ops(payload)
{new_outs, new_ins} = reduce(ops, dep_view)
%__MODULE__{
move_out:
put_or_delete(
queue.move_out,
dep_index,
new_outs,
MapSet.union(existing_out_txids, new_txids)
),
move_in:
put_or_delete(
queue.move_in,
dep_index,
new_ins,
MapSet.union(existing_in_txids, new_txids)
)
}
end
@doc """
Pop the next batch of operations. Returns move-out batches (any dep) before move-in batches.
Returns `{batch, updated_queue}` or `nil` if the queue is empty.
"""
@spec pop_next(t()) :: {batch(), t()} | nil
def pop_next(%__MODULE__{move_out: move_out} = queue) when move_out != %{} do
{dep_index, {values, txids}} = Enum.min_by(move_out, &elem(&1, 0))
{{:move_out, dep_index, values, sorted_txids(txids)},
%{queue | move_out: Map.delete(move_out, dep_index)}}
end
def pop_next(%__MODULE__{move_out: move_out, move_in: move_in} = queue)
when move_out == %{} and move_in != %{} do
{dep_index, {values, txids}} = Enum.min_by(move_in, &elem(&1, 0))
{{:move_in, dep_index, values, sorted_txids(txids)},
%{queue | move_in: Map.delete(move_in, dep_index)}}
end
def pop_next(%__MODULE__{}), do: nil
defp sorted_txids(%MapSet{} = txids), do: Enum.sort(txids)
defp payload_to_ops(payload) do
Enum.map(Map.get(payload, :move_out, []), &{:move_out, &1}) ++
Enum.map(Map.get(payload, :move_in, []), &{:move_in, &1})
end
defp reduce(ops, base_view) do
terminal_ops =
ops
|> Enum.with_index()
|> Enum.reduce(%{}, fn {{kind, move_value}, index}, acc ->
Map.put(acc, elem(move_value, 0), %{kind: kind, move_value: move_value, index: index})
end)
|> Map.values()
|> Enum.reject(&redundant?(&1, base_view))
|> Enum.sort_by(& &1.index)
{
for(%{kind: :move_out, move_value: move_value} <- terminal_ops, do: move_value),
for(%{kind: :move_in, move_value: move_value} <- terminal_ops, do: move_value)
}
end
defp redundant?(%{kind: :move_in, move_value: {value, _}}, base_view) do
MapSet.member?(base_view, value)
end
defp redundant?(%{kind: :move_out, move_value: {value, _}}, base_view) do
not MapSet.member?(base_view, value)
end
defp put_or_delete(map, key, [], _txids), do: Map.delete(map, key)
defp put_or_delete(map, key, values, txids), do: Map.put(map, key, {values, txids})
end