Packages
electric
1.1.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/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.
"""
alias Electric.Connection
require Logger
@behaviour Postgrex.SimpleConnection
@type option ::
{:connection_opts, Keyword.t()}
| {:connection_manager, GenServer.server()}
| {:lock_name, String.t()}
@type options :: [option]
@default_timeout 30_000
defmodule State do
defstruct [
:connection_manager,
:lock_acquired,
:lock_name,
:backoff,
:stack_id
]
end
def name(stack_id) do
Electric.ProcessRegistry.name(stack_id, __MODULE__)
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)
# Start the lock connection in logical replication mode to side-step any connection pooler
# that may be sitting between us and the Postgres server.
#
# We cannot get desired semantics of session-level advisory locks when connecting to
# Postgres through a pooler that runs in transaction mode (such as PGBouncer running in
# front of Neon). Starting a connection in replication mode ensures that it will be a
# direct connection to the database, so it can take a session-level advisory lock whose
# lifetime will be tied to the connection's lifetime.
connection_opts =
connection_opts
|> Electric.Utils.deobfuscate_password()
|> connection_opts_with_logical_replication()
stack_id = Keyword.fetch!(opts, :stack_id)
Postgrex.SimpleConnection.start_link(
__MODULE__,
init_opts,
[
timeout: Access.get(opts, :timeout, @default_timeout),
auto_reconnect: false,
sync_connect: false,
name: name(stack_id)
] ++
connection_opts
)
end
@impl true
def init(opts) do
opts = Map.new(opts)
Process.set_label({:lock_connection, opts.stack_id})
metadata = [
lock_name: opts.lock_name,
# flag used for error filtering
is_connection_process?: true,
stack_id: opts.stack_id
]
Logger.metadata(metadata)
Electric.Telemetry.Sentry.set_tags_context(metadata)
Logger.debug("Opening lock connection")
state = %State{
connection_manager: opts.connection_manager,
lock_name: opts.lock_name,
lock_acquired: false,
backoff: {:backoff.init(1000, 10_000), nil}
}
{:ok, state}
end
@impl true
def handle_connect(state) do
notify_connection_opened(state)
# Verify that the connection has been opened in replication mode.
#
# If there's a pooler running in front of the Postgres server, it may have simply ignored
# the replication=database connection parameter, defeating the purpose of us requesting the
# replication mode in the first place which is to get the desired session-level locking
# semantics.
#
# Issuing a statement that would cause a syntax error on a regular connection is a surefire
# way to ensure the connection is running in the correct mode.
send(self(), :identify_system)
{:noreply, state}
end
@impl true
def handle_info(:identify_system, state) do
{:query, "IDENTIFY_SYSTEM", state}
end
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{command: :identify} = result], state) do
# [db] postgres:postgres=> IDENTIFY_SYSTEM;
# systemid │ timeline │ xlogpos │ dbname
# ─────────────────────┼──────────┼───────────┼──────────
# 7506979529870965272 │ 1 │ 0/220AE10 │ postgres
# (1 row)
[[systemid, timeline, xlogpos, _dbname]] = result.rows
notify_system_identified(state, %{
system_identifier: systemid,
timeline_id: timeline,
current_wal_flush_lsn: xlogpos
})
# Now proceed to the actual lock acquisition.
send(self(), :acquire_lock)
{:noreply, state}
end
def handle_result(%Postgrex.Error{postgres: %{code: :syntax_error}} = error, _state) do
# Postgrex.SimpleConnection does not support {:stop, ...} or {:shutdown, ...} return values
# from callback functions, so we raise here and let the connection manager handle the error.
raise error
end
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)
if not is_expected_error?(error),
do:
Logger.error(
"Failed to acquire lock #{state.lock_name} with reason #{inspect(error)} - retrying in #{inspect(time)}ms."
)
notify_lock_acquisition_error(error, state)
{:noreply, %{state | lock_acquired: false, backoff: {backoff, tref}}}
end
defp notify_connection_opened(%State{connection_manager: manager}) do
Connection.Manager.lock_connection_started(manager)
end
defp notify_system_identified(%State{connection_manager: manager}, info) do
Connection.Manager.pg_system_info_obtained(manager, info)
end
defp notify_lock_acquisition_error(error, %State{connection_manager: manager}) do
Connection.Manager.exclusive_connection_lock_acquisition_failed(manager, error)
end
defp notify_lock_acquired(%State{connection_manager: manager}) do
Connection.Manager.exclusive_connection_lock_acquired(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
defp is_expected_error?(%Postgrex.Error{
postgres: %{
code: :query_canceled,
pg_code: "57014",
message: "canceling statement due to statement timeout"
}
}),
do: true
defp is_expected_error?(_), do: false
defp connection_opts_with_logical_replication(connection_opts) do
update_in(
connection_opts,
[:parameters],
fn params -> params |> List.wrap() |> Keyword.put(:replication, "database") end
)
end
end