Packages
electric
1.2.3
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/shape_log_collector/affected_columns.ex
defmodule Electric.Replication.ShapeLogCollector.AffectedColumns do
@moduledoc false
require Logger
alias Electric.Replication.Changes.Relation
def init(%{id_to_table_info: id_to_table_info, table_to_id: table_to_id})
when is_map(id_to_table_info) and is_map(table_to_id) do
{:ok, %{id_to_table_info: id_to_table_info, table_to_id: table_to_id}}
end
def transform_relation(
%Relation{schema: schema, table: table, id: id} = rel,
%{
id_to_table_info: id_to_table_info,
table_to_id: table_to_id
} = state
) do
schema_table = {schema, table}
existing_id = Map.get(table_to_id, schema_table)
existing_rel = Map.get(id_to_table_info, id)
case {existing_id, existing_rel} do
# New relation, register it
{nil, nil} ->
{rel, add_relation(state, id, rel)}
# Relation identity matches known, let's compare columns
{^id, %Relation{schema: ^schema, table: ^table}} ->
case find_differing_columns(existing_rel, rel) do
# No (noticable) changes to the relation, continue as-is
[] ->
{rel, state}
affected_cols ->
updated_rel = %{rel | affected_columns: affected_cols}
{updated_rel, add_relation(state, id, rel)}
end
# Some part of identity changed, update the state and pass it through
{_, _} ->
Logger.debug(fn ->
"Relation identity changed: #{existing_id}/#{inspect(existing_rel)} -> #{inspect(rel)}"
end)
{rel,
state
|> delete_tracked_relation(schema_table_key(existing_rel), existing_id)
|> add_relation(id, rel)}
end
end
defp schema_table_key(%Relation{schema: schema, table: table}), do: {schema, table}
defp schema_table_key(nil), do: nil
defp add_relation(state, id, rel) do
state
|> put_in([:table_to_id, schema_table_key(rel)], id)
|> put_in([:id_to_table_info, id], rel)
end
defp delete_tracked_relation(state, schema_table, id) do
state
|> update_in([:table_to_id], &Map.delete(&1, schema_table))
|> update_in([:id_to_table_info], &Map.delete(&1, id))
end
defp find_differing_columns(%Relation{columns: old_cols}, %Relation{columns: new_cols})
when old_cols == new_cols,
do: []
defp find_differing_columns(%Relation{columns: old_cols}, %Relation{columns: new_cols}) do
(old_cols ++ new_cols)
|> Enum.reduce(%{}, fn
%{name: name, type_oid: type_oid}, acc when is_map_key(acc, name) ->
# We're seeing column with this name for a second time, so we can remove it from the diff if type oid is the same
if acc[name] == type_oid, do: Map.delete(acc, name), else: acc
%{name: name, type_oid: type_oid}, acc ->
# If we're seeing column with this name for a first time, it'll either stay if it's present only in one set,
# or be deleted if seen again
Map.put(acc, name, type_oid)
end)
|> Map.keys()
end
end