Packages
electric
1.0.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/lock_connection.ex
defmodule Electric.Postgres.LockConnection do
@moduledoc """
A Postgres connection that ensures an advisory lock is held for its entire duration,
useful for ensuring only a single sync service instance can be using a single
replication slot at any given time.
The connection attempts to grab the lock and waits on it until it acquires it.
When it does, it fires off an :exclusive_connection_lock_acquired message to the specified
`Electric.Connection.Manager` such that the required setup can acquired now that
the service is sure to be the only one operating on this replication stream.
"""
require Logger
@behaviour Postgrex.SimpleConnection
@type option ::
{:connection_opts, Keyword.t()}
| {:connection_manager, GenServer.server()}
| {:lock_name, String.t()}
@type options :: [option]
defmodule State do
defstruct [
:connection_manager,
:lock_acquired,
:lock_name,
:backoff
]
end
@spec start_link(options()) :: {:ok, pid()} | {:error, Postgrex.Error.t() | term()}
def start_link(opts) do
{connection_opts, init_opts} = Keyword.pop(opts, :connection_opts)
Postgrex.SimpleConnection.start_link(
__MODULE__,
init_opts,
[timeout: :infinity, auto_reconnect: false] ++
Electric.Utils.deobfuscate_password(connection_opts)
)
end
@impl true
def init(opts) do
send(self(), :acquire_lock)
metadata = [
lock_name: Keyword.fetch!(opts, :lock_name),
stack_id: Keyword.fetch!(opts, :stack_id)
]
Logger.metadata(metadata)
Electric.Telemetry.Sentry.set_tags_context(metadata)
{:ok,
%State{
connection_manager: Keyword.fetch!(opts, :connection_manager),
lock_name: Keyword.fetch!(opts, :lock_name),
lock_acquired: false,
backoff: {:backoff.init(1000, 10_000), nil}
}}
end
@impl true
def handle_info(:acquire_lock, state) do
if state.lock_acquired do
notify_lock_acquired(state)
{:noreply, state}
else
Logger.info("Acquiring lock from postgres with name #{state.lock_name}")
{:query, lock_query(state), state}
end
end
def handle_info({:timeout, tref, msg}, %{backoff: {backoff, tref}} = state) do
handle_info(msg, %{state | backoff: {backoff, nil}})
end
@impl true
def handle_result([%Postgrex.Result{columns: ["pg_advisory_lock"]}], state) do
Logger.info("Lock acquired from postgres with name #{state.lock_name}")
notify_lock_acquired(state)
{:noreply, %{state | lock_acquired: true}}
end
def handle_result(%Postgrex.Error{} = error, %State{backoff: {backoff, _}} = state) do
{time, backoff} = :backoff.fail(backoff)
tref = :erlang.start_timer(time, self(), :acquire_lock)
Logger.error(
"Failed to acquire lock #{state.lock_name} with reason #{inspect(error)} - retrying in #{inspect(time)}ms."
)
{:noreply, %{state | lock_acquired: false, backoff: {backoff, tref}}}
end
defp notify_lock_acquired(%State{connection_manager: connection_manager} = _state) do
Electric.Connection.Manager.exclusive_connection_lock_acquired(connection_manager)
end
defp lock_query(%State{lock_name: name} = _state) do
"SELECT pg_advisory_lock(hashtext('#{name}'))"
end
@impl true
def notify(_channel, _payload, _state) do
:ok
end
end