Packages
electric
1.6.2
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/postgres/replication_client.ex
defmodule Electric.Postgres.ReplicationClient do
@moduledoc """
A client module for Postgres logical replication.
"""
use Electric.Postgres.ReplicationConnection
alias Electric.Postgres.LogicalReplication.Decoder
alias Electric.Postgres.Lsn
alias Electric.LsnTracker
alias Electric.Postgres.ReplicationClient.MessageConverter
alias Electric.Postgres.ReplicationClient.ConnectionSetup
alias Electric.Replication.Changes.TransactionFragment
alias Electric.Replication.Changes.Relation
alias Electric.Telemetry.OpenTelemetry
alias Electric.Telemetry.Sampler
require Logger
require MessageConverter
@type step ::
:disconnected
| :connected
| :identify_system
| :query_pg_info
| :acquire_lock
| :create_publication
| :check_if_publication_exists
| :drop_slot
| :create_slot
| :query_slot_flushed_lsn
| :set_display_setting
| :ready_to_stream
| :start_streaming
| :streaming
defmodule State do
@enforce_keys [:handle_event, :publication_name]
defstruct [
:stack_id,
:connection_manager,
:handle_event,
:publication_name,
:lock_acquired?,
:try_creating_publication?,
:recreate_slot?,
:start_streaming?,
:pg_version,
:slot_name,
:slot_temporary?,
:display_settings,
:message_converter,
:publication_owner?,
:replication_idle_timeout,
wal_sender_timeout: 60_000,
step: :disconnected,
wait_for_active_ref: nil,
pending_event: nil,
received_wal: 0,
flushed_wal: 0,
last_seen_txn_lsn: Lsn.from_integer(0),
last_seen_txn_timestamp: nil,
flush_up_to_date?: true
]
@type t() :: %__MODULE__{
stack_id: String.t(),
connection_manager: pid(),
handle_event: {module(), atom(), [term()]},
publication_name: String.t(),
try_creating_publication?: boolean(),
recreate_slot?: boolean(),
start_streaming?: boolean(),
pg_version: non_neg_integer(),
slot_name: String.t(),
slot_temporary?: boolean(),
display_settings: [String.t()],
message_converter: MessageConverter.t(),
publication_owner?: boolean(),
replication_idle_timeout: non_neg_integer(),
wal_sender_timeout: non_neg_integer(),
step: Electric.Postgres.ReplicationClient.step(),
wait_for_active_ref: {reference(), term()} | nil,
pending_event: {reference(), term(), non_neg_integer(), integer()} | nil,
received_wal: non_neg_integer(),
flushed_wal: non_neg_integer(),
last_seen_txn_lsn: Lsn.t(),
last_seen_txn_timestamp: integer(),
flush_up_to_date?: boolean()
}
@opts_schema NimbleOptions.new!(
stack_id: [required: true, type: :string],
connection_manager: [required: true, type: :pid],
handle_event: [required: true, type: :mfa],
publication_name: [required: true, type: :string],
try_creating_publication?: [required: true, type: :boolean],
start_streaming?: [type: :boolean, default: true],
slot_name: [required: true, type: :string],
slot_temporary?: [type: :boolean, default: false],
replication_idle_timeout: [type: :non_neg_integer, default: 0],
# Set a reasonable limit for the maximum size of a transaction that
# we can handle, above which we would exit as we run the risk of running
# out of memmory.
# TODO: stream out transactions and collect on disk to avoid this
max_txn_size: [type: {:or, [:non_neg_integer, nil]}, default: nil],
# Maximum number of changes to buffer before flushing a transaction fragment.
# Smaller values result in more message passing overhead but lower memory usage.
# The minimum allowed value is 2.
max_batch_size: [type: :non_neg_integer, default: 100]
)
@spec new(Access.t()) :: t()
def new(opts) do
opts = NimbleOptions.validate!(opts, @opts_schema)
settings = [display_settings: Electric.Postgres.display_settings()]
opts = settings ++ opts
{max_txn_size, opts} = Keyword.pop!(opts, :max_txn_size)
{max_batch_size, opts} = Keyword.pop!(opts, :max_batch_size)
# Assert the implicit requirement
true = max_batch_size >= 2
struct!(
__MODULE__,
opts ++
[
message_converter:
MessageConverter.new(
max_tx_size: max_txn_size,
max_batch_size: max_batch_size
)
]
)
end
end
# @type state :: State.t()
@repl_msg_x_log_data ?w
@repl_msg_primary_keepalive ?k
@repl_msg_standby_status_update ?r
@default_connect_timeout 30_000
@idle_check_interval Electric.Config.min_replication_idle_timeout()
# Maximum keepalive interval. Caps the derived interval (wal_sender_timeout/3)
# so we stay responsive even if wal_sender_timeout is set very high or changes
# on the source PG after we've connected.
@max_keepalive_interval 15_000
# Delay before retrying a failed event dispatch.
@event_retry_delay 50
# Maximum time to spend retrying a crashed event handler before giving up.
@max_event_retry_time 10 * 60_000
@spec start_link(Keyword.t()) :: :gen_statem.start_ret()
def start_link(opts) do
config = Map.new(opts)
# Disable the reconnection logic in Postgex.ReplicationConnection to force it to exit with
# the connection error. Without this, we may observe undesirable restarts in tests between
# one test process exiting and the next one starting.
start_opts =
[
name: name(config.stack_id),
timeout: Access.get(opts, :timeout, @default_connect_timeout),
auto_reconnect: false,
sync_connect: false
] ++ Electric.Utils.deobfuscate_password(config.replication_opts[:connection_opts])
Electric.Postgres.ReplicationConnection.start_link(
__MODULE__,
Keyword.delete(config.replication_opts, :connection_opts),
start_opts
)
end
def name(stack_id) do
Electric.ProcessRegistry.name(stack_id, __MODULE__)
end
# This is a send() and not a call() to prevent the caller (the Connection.Manager process) from
# getting blocked when the replication connection is blocked some replication slot condition
# that doesn't let it start streaming immediately.
def start_streaming(client) do
send(client, :start_streaming)
end
def stop(client, reason) do
Electric.Postgres.ReplicationConnection.call(client, {:stop, reason})
end
# The `Postgrex.ReplicationConnection` behaviour does not follow the gen server conventions and
# establishes its own instead. Unless the `sync_connect: false` option is passed to `start_link()`, the
# connection process will try opening a replication connection to Postgres before returning
# from its `init()` callback.
#
# The callbacks `init()`, `handle_connect()` and `handle_result()` defined in this module
# would all be invoked inside the connection process' `init()` callback in that case. Once
# any of the callbacks return `{:stream, ...}`, the connection process finishes its
# initialization and switches into the logical streaming mode to start receiving logical
# messages from Postgres, invoking the `handle_data()` callback for each one.
#
# TODO(alco): this needs additional info about :noreply and :query return tuples.
@impl true
def init(replication_opts) do
state = State.new(replication_opts)
Process.set_label({:replication_client, state.stack_id})
Logger.metadata(stack_id: state.stack_id, is_connection_process?: true)
Electric.Telemetry.Sentry.set_tags_context(stack_id: state.stack_id)
{:ok, state}
end
# `Postgrex.ReplicationConnection` opens a new replication connection to Postgres and then
# gives us a chance to execute one or more queries before switching into the logical
# streaming mode. It doesn't give us the connection socket but instead takes the query returned
# by one of our `handle_connect/1`, `handle_result/2` or `handle_info/2` callbacks, executes
# it, invokes the `handle_result/2` callback on the result which may return another query to
# execute, executes that, and so it goes on and on, recursively, until a callback returns
# `{:noreply, ...}` or `{:streaming, ...}`.
#
# To execute a series of queries one after the other, we define an ad-hoc state
# machine that starts from the :connected state in `handle_connect/1`, then transitions to
# the next step and returns the appropriate query to `Postgrex.ReplicationConnection` for execution,
# This is all implemented in a separate module named `Electric.Postgres.ReplicationClient.ConnectionSetup`
# to separate the connection setup logic from logical streaming.
@impl true
def handle_connect(state) do
%{state | step: :connected}
|> notify_connection_opened()
|> ConnectionSetup.start()
end
@impl true
def handle_result(result_list_or_error, state) do
{current_step, next_step, extra_info, updated_state, return_val} =
ConnectionSetup.process_query_result(result_list_or_error, state)
if current_step == :identify_system,
do: notify_system_identified(state, extra_info)
if current_step == :query_pg_info,
do: notify_pg_info_obtained(state, extra_info)
if current_step == :acquire_lock do
case extra_info do
:lock_acquired -> notify_lock_acquired(state)
{:lock_acquisition_failed, error} -> notify_lock_acquisition_error(state, error)
end
end
# for new slots, always reset the last processed LSN
if current_step == :create_slot and extra_info == :created_new_slot do
Electric.LsnTracker.set_last_processed_lsn(state.stack_id, updated_state.flushed_wal)
notify_created_new_slot(state)
end
# for existing slots, populate the last processed LSN if not present
if current_step == :query_slot_flushed_lsn,
do:
Electric.LsnTracker.initialize_last_processed_lsn(
state.stack_id,
updated_state.flushed_wal
)
if next_step == :ready_to_stream,
do: notify_ready_to_stream(state)
return_val
end
@impl true
def handle_call({:stop, reason}, from, _state) do
Logger.notice(
"Replication client #{inspect(self())} is stopping after receiving stop request from #{inspect(elem(from, 0))} with reason #{inspect(reason)}"
)
{:disconnect, reason}
end
@impl true
def handle_info({:flush_boundary_updated, lsn}, state) do
state =
if Lsn.from_integer(lsn) == state.last_seen_txn_lsn do
%{
state
| flush_up_to_date?: true,
flushed_wal: state.received_wal,
received_wal: max(lsn, state.received_wal)
}
else
%{state | flushed_wal: max(lsn, state.flushed_wal), received_wal: state.received_wal}
end
{:noreply, [encode_standby_status_update(state)], state}
end
@impl true
def handle_info(:start_streaming, %State{step: :ready_to_stream} = state) do
ConnectionSetup.start_streaming(state)
end
def handle_info(:start_streaming, %State{step: step} = state) do
Logger.debug("Replication client requested to start streaming while step=#{step}")
{:noreply, state}
end
def handle_info(:check_if_idle, %State{last_seen_txn_timestamp: txn_ts} = state) do
time_diff = System.convert_time_unit(System.monotonic_time() - txn_ts, :native, :millisecond)
if time_diff >= state.replication_idle_timeout do
{:disconnect, {:shutdown, {:connection_idle, time_diff}}}
else
{:noreply, state}
end
end
# Periodic keepalive: send a StandbyStatusUpdate to prevent wal_sender_timeout.
# This fires every @keepalive_interval ms regardless of whether the socket is
# paused. Sending the same LSN without advancement is safe — it only resets
# PostgreSQL's last_reply_timestamp.
def handle_info(:send_keepalive, %State{step: :streaming} = state) do
{:noreply, [encode_standby_status_update(state)], state}
end
def handle_info(:send_keepalive, state) do
{:noreply, state}
end
# Event processing messages — see dispatch_event/2 and apply_event/3 below.
def handle_info({:process_event, event, time_remaining}, state),
do: apply_event(event, time_remaining, state)
# StatusMonitor notification: stack became active — retry the pending event.
# The delay prevents spinning if the handler returns :not_ready again despite
# StatusMonitor reporting :active (a brief race during startup).
def handle_info(
{{Electric.StatusMonitor, ref}, {:ok, :active}},
%State{wait_for_active_ref: {ref, event}} = state
) do
Process.send_after(self(), {:process_event, event, @max_event_retry_time}, @event_retry_delay)
{:noreply, %{state | wait_for_active_ref: nil}}
end
# Stale or unexpected StatusMonitor notification (e.g. after a retry already
# succeeded and cleared wait_ref). Discard silently.
def handle_info({{Electric.StatusMonitor, _ref}, _result}, state) do
{:noreply, state}
end
# Async event handler replied :ok — demonitor, ack transaction, resume socket.
def handle_info(
{ref, :ok},
%State{pending_event: {ref, event, _time_remaining, _start_time}} = state
)
when is_reference(ref) do
Process.demonitor(ref, [:flush])
state = %{state | pending_event: nil}
state = maybe_update_flush_up_to_date(state)
{acks, state} = acknowledge_transaction(event, state)
{:noreply_and_resume, acks, state}
end
# Async event handler replied with a recoverable error — wait and retry.
def handle_info(
{ref, {:error, error}},
%State{pending_event: {ref, event, time_remaining, start_time}} = state
)
when is_reference(ref) and error in [:not_ready, :connection_not_available] do
Process.demonitor(ref, [:flush])
remaining = time_remaining - (System.monotonic_time(:millisecond) - start_time)
state = %{state | pending_event: nil}
wait_for_active_and_retry(event, remaining, state)
end
# Async event handler crashed — retry with budget.
def handle_info(
{:DOWN, ref, :process, _pid, reason},
%State{pending_event: {ref, event, time_remaining, start_time}} = state
) do
remaining = time_remaining - (System.monotonic_time(:millisecond) - start_time)
state = %{state | pending_event: nil}
if remaining > 0 do
Logger.error(
"Error processing replication event (#{remaining}ms retry budget left): " <>
inspect(reason)
)
Process.send_after(self(), {:process_event, event, remaining}, @event_retry_delay)
{:noreply, state}
else
Logger.error("Exhausted retry budget processing replication event: " <> inspect(reason))
exit(reason)
end
end
# This callback is invoked when the connection process receives a shutdown signal.
def handle_info({:EXIT, _pid, :shutdown}, _state) do
Logger.debug("Replication client #{inspect(self())} received shutdown signal, stopping")
{:disconnect, :shutdown}
end
# Some other exit reason we're not expecting: disconnect and shut down.
def handle_info({:EXIT, _pid, reason}, _state) do
{:disconnect, reason}
end
# The implementation of Postgrex.ReplicationConnection doesn't give us a convenient way to
# check whether the START_REPLICATION_SLOT statement succeeded before switching the
# connection into streaming mode. Returning {:query, "START_REPLICATION_SLOT ...", state}
# works fine when the query result is an error: it is then passed to the handle_result()
# callback. But if streaming starts without issues, a function clause error is encountered
# inside Postgrex.ReplicationConnection because it expects the connection to already have
# been switched into streaming mode by returning {:stream, "START_REPLICATION_SLOT ...", [], state}.
#
# Hence this function clause of `handle_data()` that notifies the connection manager about
# successful streaming start as soon as it receives the first replication message from
# Postgres.
@impl true
@spec handle_data(binary(), State.t()) ::
{:noreply, State.t()}
| {:noreply, list(binary()), State.t()}
| {:noreply_and_pause, list(binary()), State.t()}
| {:disconnect, term()}
def handle_data(data, %State{step: :start_streaming} = state) do
# Modify the state as if we've just seen a transaction so that in the future we have a
# starting point to check how long the stream has been idle for.
state = %{state | step: :streaming, last_seen_txn_timestamp: System.monotonic_time()}
if state.replication_idle_timeout > 0 do
:timer.send_interval(@idle_check_interval, :check_if_idle)
end
# Start a periodic keepalive timer. This sends StandbyStatusUpdate messages
# to PostgreSQL at regular intervals, preventing wal_sender_timeout from
# firing even when the socket is paused for backpressure.
#
# The interval is derived from PostgreSQL's wal_sender_timeout (queried during
# connection setup): timeout/3 provides a safe margin, matching the heuristic
# used by pg_recvlogical and other replication clients.
keepalive_interval = keepalive_interval(state.wal_sender_timeout)
:timer.send_interval(keepalive_interval, :send_keepalive)
Logger.debug(
"Keepalive interval set to #{keepalive_interval}ms (wal_sender_timeout=#{state.wal_sender_timeout}ms)"
)
notify_seen_first_message(state)
handle_data(data, state)
end
def handle_data(<<@repl_msg_primary_keepalive, wal_end::64, _clock::64, reply>>, state) do
Logger.debug(fn ->
"Primary Keepalive: wal_end=#{wal_end} (#{Lsn.from_integer(wal_end)}) reply=#{reply}"
end)
in_transaction? = MessageConverter.in_transaction?(state.message_converter)
# Broadcast the server's latest WAL position to consumers with materializer
# dependencies. This covers silent-postgres periods where no transactions
# arrive but buffered move-in results may be ready to splice.
unless in_transaction? do
LsnTracker.broadcast_last_seen_lsn(state.stack_id, wal_end)
end
case reply do
1 when in_transaction? ->
{:noreply, [encode_standby_status_update(state)], state}
# if we are not in a transaction, advance the replication slot
# with keepalives to avoid it getting filled with irrelevant changes, like
# heartbeats from the database provider
1 ->
state = update_stored_wals(state, wal_end)
{:noreply, [encode_standby_status_update(state)], state}
0 when in_transaction? ->
{:noreply, [], state}
0 ->
state = update_stored_wals(state, wal_end)
{:noreply, [], state}
end
end
def handle_data(
<<@repl_msg_x_log_data, _wal_start::64, _server_wal_end::64, _clock::64, data::binary>>,
%State{} = state
) do
msg = Decoder.decode(data)
case MessageConverter.convert(msg, state.message_converter) do
{:error, reason} ->
{:disconnect, {:irrecoverable_slot, reason}}
{:buffering, converter} ->
{:noreply, %{state | message_converter: converter}}
{:ok, event, converter} ->
state = %{state | message_converter: converter}
dispatch_event(event, state)
end
end
# Dispatch event processing asynchronously. Pauses the socket so we don't
# receive more data until processing completes. The gen_statem remains
# responsive to handle_info messages (keepalive timer, flush_boundary_updated,
# EXIT signals) while providing backpressure to the replication stream.
#
# maybe_update_flush_up_to_date and acknowledge_transaction are intentionally
# deferred to apply_event's success path, preserving the original semantics
# where they only ran after handle_event succeeded.
defp dispatch_event(event, state) do
send(self(), {:process_event, event, @max_event_retry_time})
{:noreply_and_pause, [], state}
end
# Dispatch the event handler as a non-blocking $gen_call. The MFA returns a
# monitor ref; the gen_statem returns immediately and handles the reply (or
# :DOWN on crash) in handle_info. This keeps the gen_statem responsive to
# keepalive timers while the handler processes the event.
defp apply_event(event, time_remaining, state) do
{m, f, args} = state.handle_event
start_time = System.monotonic_time(:millisecond)
try do
ref = apply(m, f, [event | args])
{:noreply, %{state | pending_event: {ref, event, time_remaining, start_time}}}
catch
kind, reason ->
remaining = time_remaining - (System.monotonic_time(:millisecond) - start_time)
if remaining > 0 do
Logger.error(
"Error dispatching replication event (#{remaining}ms retry budget left): " <>
Exception.format(kind, reason, __STACKTRACE__)
)
Process.send_after(self(), {:process_event, event, remaining}, @event_retry_delay)
{:noreply, state}
else
Logger.error(
"Exhausted retry budget dispatching replication event: " <>
Exception.format(kind, reason, __STACKTRACE__)
)
:erlang.raise(kind, reason, __STACKTRACE__)
end
end
end
# Downstream returned :not_ready — subscribe to StatusMonitor for notification
# when the stack becomes active, then retry. This replaces the old blocking
# wait_until_active(timeout: :infinity) call with an async notification.
# The keepalive timer prevents wal_sender_timeout during the wait.
# The remaining budget is intentionally discarded — a fresh @max_event_retry_time
# is used after the stack becomes active, matching the old apply_with_retries
# behavior which reset the retry timer after wait_until_active returned.
defp wait_for_active_and_retry(event, _remaining, state) do
ref = Electric.StatusMonitor.wait_until_async(state.stack_id, :active)
{:noreply, %{state | wait_for_active_ref: {ref, event}}}
end
defp acknowledge_transaction(%TransactionFragment{commit: nil}, state), do: {[], state}
defp acknowledge_transaction(%TransactionFragment{lsn: lsn, commit: commit}, state) do
if Sampler.sample_metrics?() do
alias Electric.Replication.Changes.Commit
OpenTelemetry.execute(
[:electric, :postgres, :replication, :transaction_received],
%{
monotonic_time: System.monotonic_time(),
receive_lag: Commit.calculate_final_receive_lag(commit, System.monotonic_time()),
bytes: commit.transaction_size,
count: 1,
operations: commit.txn_change_count
},
%{stack_id: state.stack_id}
)
end
state =
%{
state
| last_seen_txn_lsn: lsn,
last_seen_txn_timestamp: System.monotonic_time()
}
|> update_received_wal(Lsn.to_integer(lsn))
{[encode_standby_status_update(state)], state}
end
defp acknowledge_transaction(%Relation{}, state), do: {[], state}
defp maybe_update_flush_up_to_date(state) do
if MessageConverter.in_transaction?(state.message_converter) do
%{state | flush_up_to_date?: false}
else
state
end
end
defp encode_standby_status_update(state) do
Logger.debug(fn ->
"Standby status update: received_wal=#{Lsn.from_integer(state.received_wal)}, flushed_wal=#{Lsn.from_integer(state.flushed_wal)}"
end)
<<
@repl_msg_standby_status_update,
state.received_wal + 1::64,
state.flushed_wal + 1::64,
state.flushed_wal + 1::64,
current_time()::64,
0
>>
end
# Derive keepalive interval from PostgreSQL's wal_sender_timeout.
# Uses min(timeout/3, 15s): timeout/3 provides a safe margin for low timeouts,
# while the 15s cap ensures responsiveness even if wal_sender_timeout is very
# high or changes on the source PG after we've connected.
defp keepalive_interval(0), do: @max_keepalive_interval
defp keepalive_interval(wal_sender_timeout_ms),
do: min(div(wal_sender_timeout_ms, 3), @max_keepalive_interval)
@epoch DateTime.to_unix(~U[2000-01-01 00:00:00Z], :microsecond)
defp current_time(), do: System.os_time(:microsecond) - @epoch
defp update_stored_wals(
%{
received_wal: received_wal,
flushed_wal: flushed_wal,
flush_up_to_date?: flush_up_to_date?
} = state,
wal
) do
received_wal = max(received_wal, wal)
flushed_wal = if flush_up_to_date?, do: max(flushed_wal, wal), else: flushed_wal
%{state | received_wal: received_wal, flushed_wal: flushed_wal}
end
defp update_received_wal(state, wal) when is_number(wal) and wal >= state.received_wal,
do: %{state | received_wal: wal}
defp update_received_wal(state, wal) when is_number(wal), do: state
defp notify_connection_opened(%State{connection_manager: manager} = state) do
:ok = Electric.Connection.Manager.replication_client_started(manager)
state
end
defp notify_system_identified(%State{connection_manager: manager} = state, info) do
:ok = Electric.Connection.Manager.pg_system_identified(manager, info)
state
end
defp notify_pg_info_obtained(%State{connection_manager: manager} = state, pg_info) do
:ok = Electric.Connection.Manager.pg_info_obtained(manager, pg_info)
state
end
defp notify_lock_acquisition_error(%State{connection_manager: manager} = state, error) do
:ok = Electric.Connection.Manager.replication_client_lock_acquisition_failed(manager, error)
state
end
defp notify_lock_acquired(%State{connection_manager: manager} = state) do
:ok = Electric.Connection.Manager.replication_client_lock_acquired(manager)
state
end
defp notify_created_new_slot(%State{connection_manager: manager} = state) do
:ok = Electric.Connection.Manager.replication_client_created_new_slot(manager)
state
end
defp notify_ready_to_stream(%State{connection_manager: manager} = state) do
:ok = Electric.Connection.Manager.replication_client_ready_to_stream(manager)
state
end
defp notify_seen_first_message(%State{connection_manager: manager} = state) do
:ok = Electric.Connection.Manager.replication_client_streamed_first_message(manager)
state
end
end