Packages
electric
1.3.4
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/initial_snapshot.ex
defmodule Electric.Shapes.Consumer.InitialSnapshot do
@moduledoc false
# Internal module, used as a part of the consumer state, dealing
# with the initial snapshot state and the waiting for the snapshot to start.
alias Electric.Postgres.Xid
alias Electric.Replication.Changes.Transaction
alias Electric.ShapeCache.Storage
alias Electric.Shapes.Consumer.State
defstruct filtering?: true,
snapshot_started?: false,
pg_snapshot: nil,
awaiting_snapshot_start: []
@type t() :: %__MODULE__{
filtering?: boolean(),
snapshot_started?: boolean(),
pg_snapshot: nil | State.pg_snapshot(),
awaiting_snapshot_start: list(GenServer.from())
}
@spec new(Storage.pg_snapshot() | nil) :: t()
def new(nil), do: %__MODULE__{filtering?: true}
def new(%{xmin: xmin, xmax: xmax, xip_list: xip_list} = snapshot) do
%__MODULE__{
filtering?: Map.get(snapshot, :filter_txns?, true),
pg_snapshot: {xmin, xmax, xip_list}
}
end
def add_waiter(%__MODULE__{} = state, from) do
%{state | awaiting_snapshot_start: [from | state.awaiting_snapshot_start]}
end
def reply_to_waiters(%__MODULE__{} = state, reply) do
for client <- List.wrap(state.awaiting_snapshot_start),
not is_nil(client),
do: GenServer.reply(client, reply)
%{state | awaiting_snapshot_start: []}
end
def needs_buffering?(%__MODULE__{pg_snapshot: snapshot}), do: is_nil(snapshot)
def maybe_stop_initial_filtering(
%__MODULE__{pg_snapshot: {xmin, xmax, xip_list} = snapshot} = state,
storage,
%Transaction{xid: xid}
) do
if Xid.after_snapshot?(xid, snapshot) do
Storage.set_pg_snapshot(
%{xmin: xmin, xmax: xmax, xip_list: xip_list, filter_txns?: false},
storage
)
%{state | filtering?: false}
else
state
end
end
@spec set_initial_snapshot(t(), Storage.shape_storage(), State.pg_snapshot()) :: t()
def set_initial_snapshot(
%__MODULE__{pg_snapshot: nil} = state,
storage,
{xmin, xmax, xip_list} = snapshot
) do
# We're not changing snapshot storage format for backwards compatibility.
Storage.set_pg_snapshot(
%{xmin: xmin, xmax: xmax, xip_list: xip_list, filter_txns?: true},
storage
)
%{state | pg_snapshot: snapshot, filtering?: true}
end
def mark_snapshot_started(
%__MODULE__{snapshot_started?: true} = state,
_stack_id,
_shape_handle,
_
),
do: state
def mark_snapshot_started(%__MODULE__{} = state, stack_id, shape_handle, storage) do
Electric.Shapes.mark_snapshot_started(storage, stack_id, shape_handle)
state = reply_to_waiters(state, :started)
%{state | snapshot_started?: true}
end
def filter(state, storage, %Transaction{} = txn) do
if Transaction.visible_in_snapshot?(txn, state.pg_snapshot) do
{:consider_flushed, state}
else
state = maybe_stop_initial_filtering(state, storage, txn)
{:continue, state}
end
end
end