Packages
electric
1.1.14
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/connection/manager/supervisor.ex
defmodule Electric.Connection.Manager.Supervisor do
@moduledoc """
Intermediate supervisor that helps tie Connection.Manager's lifetime to that of Replication.Supervisor.
"""
use Supervisor
def name(opts) do
Electric.ProcessRegistry.name(opts[:stack_id], __MODULE__)
end
def start_link(opts) do
Supervisor.start_link(__MODULE__, opts, name: name(opts))
end
def init(opts) do
Process.set_label({:connection_manager_supervisor, opts[:stack_id]})
Logger.metadata(stack_id: opts[:stack_id])
Electric.Telemetry.Sentry.set_tags_context(stack_id: opts[:stack_id])
children = [
{Electric.Connection.Manager, opts},
{Electric.Connection.Manager.ConnectionResolver, stack_id: opts[:stack_id]}
]
# Electric.Connection.Manager is a permanent child of the supervisor, so when it dies, the
# :one_for_all strategy will kick in and restart the other children.
# This is not the case for Electric.Replication.Supervisor which needs to be a temporary
# child such that Electric.Connection.Manager decides when it starts. Because of this, when
# Electric.Replication.Supervisor dies, even due to an error, it doesn't activate the
# :one_for_all strategy.
# We work around this by marking Electric.Replication.Supervisor as significant and
# configuring this supervisor with [auto_shutdown: :any_significant].
Supervisor.init(children, strategy: :one_for_all, auto_shutdown: :any_significant)
end
@doc """
This function is supposed to be called from Connection.Manager at the right point in its
initialization sequence.
Replication.Supervisor is started as a temporary child so that, when it dies, it is up to the
Connection.Manager process to restart it again at the right point in time.
"""
def start_replication_supervisor(opts) do
stack_id = Keyword.fetch!(opts, :stack_id)
shape_cache_opts = Keyword.fetch!(opts, :shape_cache_opts)
replication_opts = Keyword.fetch!(opts, :replication_opts)
inspector = Keyword.fetch!(shape_cache_opts, :inspector)
persistent_kv = Keyword.fetch!(opts, :persistent_kv)
tweaks = Keyword.fetch!(opts, :tweaks)
shape_status_owner_spec =
{Electric.ShapeCache.ShapeStatusOwner,
[stack_id: stack_id, shape_status: Keyword.fetch!(shape_cache_opts, :shape_status)]}
consumer_supervisor_spec = {Electric.Shapes.DynamicConsumerSupervisor, [stack_id: stack_id]}
shape_cleaner_spec =
{Electric.ShapeCache.ShapeCleaner,
stack_id: stack_id, shape_status: Keyword.fetch!(shape_cache_opts, :shape_status)}
shape_cache_spec = {Electric.ShapeCache, shape_cache_opts}
publication_manager_spec =
{Electric.Replication.PublicationManager,
stack_id: stack_id,
publication_name: Keyword.fetch!(replication_opts, :publication_name),
can_alter_publication?: Keyword.fetch!(opts, :can_alter_publication?),
manual_table_publishing?: Keyword.fetch!(opts, :manual_table_publishing?),
db_pool: Electric.Connection.Manager.admin_pool(stack_id),
update_debounce_timeout: Keyword.get(tweaks, :publication_alter_debounce_ms, 0),
refresh_period: Keyword.get(tweaks, :publication_refresh_period, 60_000)}
shape_log_collector_spec =
{Electric.Replication.ShapeLogCollector,
stack_id: stack_id, inspector: inspector, persistent_kv: persistent_kv}
schema_reconciler_spec =
{Electric.Replication.SchemaReconciler,
stack_id: stack_id,
inspector: inspector,
period: Keyword.get(tweaks, :schema_reconciler_period, 60_000)}
expiry_manager_spec =
{Electric.ShapeCache.ExpiryManager,
max_shapes: Keyword.fetch!(opts, :max_shapes),
expiry_batch_size: Keyword.fetch!(opts, :expiry_batch_size),
stack_id: stack_id,
shape_status: Keyword.fetch!(shape_cache_opts, :shape_status)}
child_spec =
Supervisor.child_spec(
{
Electric.Replication.Supervisor,
stack_id: stack_id,
shape_status_owner: shape_status_owner_spec,
consumer_supervisor: consumer_supervisor_spec,
shape_cleaner: shape_cleaner_spec,
shape_cache: shape_cache_spec,
publication_manager: publication_manager_spec,
log_collector: shape_log_collector_spec,
schema_reconciler: schema_reconciler_spec,
expiry_manager: expiry_manager_spec
},
restart: :temporary,
significant: true
)
Supervisor.start_child(name(opts), child_spec)
end
end