Packages
electric
1.7.7
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/transaction_builder.ex
defmodule Electric.Replication.TransactionBuilder do
@moduledoc """
Builds complete transactions from a stream of TransactionFragments.
Takes TransactionFragments containing begin, commit, and changes,
and builds up Transaction structs. Returns complete transactions
when a fragment with a commit is seen.
"""
alias Electric.Replication.Changes.{
Transaction,
TransactionFragment
}
defstruct transaction: nil
@type t() :: %__MODULE__{
transaction: nil | Transaction.t()
}
def new, do: %__MODULE__{}
@doc """
Build transactions from a TransactionFragment.
Returns a tuple of {results, state} where results is a list of
complete transactions, and state is the updated builder state
containing any partial transaction.
"""
@spec build(TransactionFragment.t(), t()) :: {[Transaction.t()], t()}
def build(%TransactionFragment{} = fragment, state) do
state
|> maybe_start_transaction(fragment)
|> add_changes(fragment)
|> maybe_complete_transaction(fragment)
end
defp maybe_start_transaction(state, %TransactionFragment{has_begin?: false}), do: state
defp maybe_start_transaction(
%__MODULE__{} = state,
%TransactionFragment{xid: xid, has_begin?: true}
) do
txn = %Transaction{
xid: xid,
changes: [],
commit_timestamp: nil,
lsn: nil,
last_log_offset: nil
}
%{state | transaction: txn}
end
defp add_changes(%{transaction: txn} = state, %TransactionFragment{} = fragment) do
txn = %{
txn
| changes: Enum.reverse(fragment.changes) ++ txn.changes,
num_changes: txn.num_changes + fragment.change_count
}
%{state | transaction: txn}
end
defp maybe_complete_transaction(state, %TransactionFragment{commit: nil}) do
{[], state}
end
defp maybe_complete_transaction(
%__MODULE__{transaction: txn},
%TransactionFragment{
lsn: lsn,
commit: commit,
last_log_offset: last_log_offset
}
) do
completed_txn =
%{
txn
| lsn: lsn,
commit_timestamp: commit.commit_timestamp,
changes: Enum.reverse(txn.changes),
# The transaction may have had some changes filtered
# out, so we need to set the last_log_offset from the fragment
last_log_offset: last_log_offset
}
{[completed_txn], %__MODULE__{transaction: nil}}
end
end