Packages
electric
1.5.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/log_items.ex
defmodule Electric.LogItems do
alias Electric.Replication.Changes
alias Electric.Replication.LogOffset
alias Electric.Shapes.Shape
@moduledoc """
Defines the structure and how to create the items in the log that the electric client reads.
The log_item() data structure is a map for ease of consumption in the Elixir code,
however when JSON encoded (not done in this module) it's the format that the electric
client accepts.
"""
defp put_if_true(map, key, value) do
if value, do: Map.put(map, key, value), else: map
end
defp put_if_true(map, key, condition, value) do
if condition, do: Map.put(map, key, value), else: map
end
@type log_item ::
{LogOffset.t(),
%{
key: String.t(),
value: map(),
headers: map()
}}
@spec from_change(
Changes.data_change(),
txids :: nil | non_neg_integer() | [non_neg_integer(), ...],
pk_cols :: [String.t()],
replica :: Shape.replica()
) :: [log_item(), ...]
def from_change(%Changes.NewRecord{} = change, txids, _, _replica) do
[
{change.log_offset,
%{
key: change.key,
value: change.record,
headers:
%{
operation: :insert,
txids: List.wrap(txids),
relation: Tuple.to_list(change.relation),
lsn: to_string(change.log_offset.tx_offset),
op_position: change.log_offset.op_offset
}
|> put_if_true(:last, change.last?)
|> put_if_true(:tags, change.move_tags != [], change.move_tags)
}}
]
end
def from_change(%Changes.DeletedRecord{} = change, txids, pk_cols, replica) do
[
{change.log_offset,
%{
key: change.key,
value: take_pks_or_all(change.old_record, pk_cols, replica),
headers:
%{
operation: :delete,
txids: List.wrap(txids),
relation: Tuple.to_list(change.relation),
lsn: to_string(change.log_offset.tx_offset),
op_position: change.log_offset.op_offset
}
|> put_if_true(:last, change.last?)
|> put_if_true(:tags, change.move_tags != [], change.move_tags)
}}
]
end
# `old_key` is nil when it's unchanged. This is not possible when there is no PK defined.
def from_change(%Changes.UpdatedRecord{old_key: nil} = change, txids, pk_cols, replica) do
[
{change.log_offset,
%{
key: change.key,
headers:
%{
operation: :update,
txids: List.wrap(txids),
relation: Tuple.to_list(change.relation),
lsn: to_string(change.log_offset.tx_offset),
op_position: change.log_offset.op_offset
}
|> put_if_true(:last, change.last?)
|> put_if_true(:tags, change.move_tags != [], change.move_tags)
|> put_if_true(:removed_tags, change.move_tags != [], change.removed_move_tags)
}
|> Map.merge(put_update_values(change, pk_cols, replica))}
]
end
def from_change(%Changes.UpdatedRecord{} = change, txids, pk_cols, replica) do
new_offset = LogOffset.increment(change.log_offset)
[
{change.log_offset,
%{
key: change.old_key,
value: take_pks_or_all(change.old_record, pk_cols, replica),
headers:
%{
operation: :delete,
txids: List.wrap(txids),
relation: Tuple.to_list(change.relation),
key_change_to: change.key,
lsn: to_string(change.log_offset.tx_offset),
op_position: change.log_offset.op_offset
}
|> put_if_true(
:tags,
change.move_tags != [],
change.move_tags ++ change.removed_move_tags
)
}},
{new_offset,
%{
key: change.key,
value: change.record,
headers:
%{
operation: :insert,
txids: List.wrap(txids),
relation: Tuple.to_list(change.relation),
key_change_from: change.old_key,
lsn: to_string(new_offset.tx_offset),
op_position: new_offset.op_offset
}
|> put_if_true(:last, change.last?)
|> put_if_true(:tags, change.move_tags != [], change.move_tags)
}}
]
end
def expected_offset_after_split(%Changes.UpdatedRecord{old_key: x, log_offset: offset})
when not is_nil(x),
do: LogOffset.increment(offset)
def expected_offset_after_split(%{log_offset: offset}), do: offset
defp take_pks_or_all(record, _pks, :full), do: record
defp take_pks_or_all(record, [], :default), do: record
defp take_pks_or_all(record, pks, :default), do: Map.take(record, pks)
defp put_update_values(%{record: record, changed_columns: changed_columns}, pk_cols, :default) do
%{value: Map.take(record, Enum.concat(pk_cols, changed_columns))}
end
defp put_update_values(
%{record: record, old_record: old_record, changed_columns: changed_columns},
_pk_cols,
:full
) do
%{value: record, old_value: Map.take(old_record, MapSet.to_list(changed_columns))}
end
def merge_updates(u1, u2) when is_map_key(u1, "old_value") or is_map_key(u2, "old_value") do
%{
"key" => u1["key"],
"headers" => Map.take(u1["headers"], ["operation", "relation"]),
"value" => Map.merge(u1["value"], u2["value"]),
# When merging old values, we give preference to the older u1
"old_value" => Map.merge(u2["old_value"] || %{}, u1["old_value"] || %{})
}
end
def merge_updates(u1, u2) do
%{
"key" => u1["key"],
"headers" => Map.take(u1["headers"], ["operation", "relation"]),
"value" => Map.merge(u1["value"], u2["value"])
}
end
def keep_generic_headers(item) do
Map.update!(item, "headers", &Map.take(&1, ["operation", "relation"]))
end
end