Packages
electric
1.0.21
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_supervisor.ex
defmodule Electric.Shapes.ConsumerSupervisor do
use Supervisor, restart: :transient
require Logger
@name_schema_tuple {:tuple, [:atom, :atom, :any]}
@genserver_name_schema {:or, [:atom, @name_schema_tuple]}
# TODO: unify these with ShapeCache
@schema NimbleOptions.new!(
stack_id: [type: :any, required: true],
shape_handle: [type: :string, required: true],
shape: [type: {:struct, Electric.Shapes.Shape}, required: true],
inspector: [type: :mod_arg, required: true],
log_producer: [type: @genserver_name_schema, required: true],
registry: [type: :atom, required: true],
shape_status: [type: :mod_arg, required: true],
storage: [type: :mod_arg, required: true],
publication_manager: [type: :mod_arg, required: true],
chunk_bytes_threshold: [type: :non_neg_integer, required: true],
run_with_conn_fn: [type: {:fun, 2}, default: &DBConnection.run/2],
db_pool: [type: {:or, [:atom, :pid, @name_schema_tuple]}, required: true],
create_snapshot_fn: [
type: {:fun, 7},
default: &Electric.Shapes.Consumer.Snapshotter.query_in_readonly_txn/7
],
otel_ctx: [type: :any, required: false]
)
def name(stack_id, shape_handle) when is_binary(shape_handle) do
Electric.ProcessRegistry.name(stack_id, __MODULE__, shape_handle)
end
def name(%{
stack_id: stack_id,
shape_handle: shape_handle
}) do
name(stack_id, shape_handle)
end
def whereis(stack_id, shape_handle) do
GenServer.whereis(name(stack_id, shape_handle))
end
def start_link(opts) do
with {:ok, opts} <- NimbleOptions.validate(opts, @schema) do
config = Map.new(opts)
Supervisor.start_link(__MODULE__, config, name: name(config))
end
end
def stop_and_clean(%{
stack_id: stack_id,
shape_handle: shape_handle
}) do
stop_and_clean(stack_id, shape_handle)
end
def stop_and_clean(stack_id, shape_handle) do
# if consumer is present, terminate it gracefully, otherwise terminate supervisor
consumer = Electric.Shapes.Consumer.name(stack_id, shape_handle)
case GenServer.whereis(consumer) do
nil ->
try do
Supervisor.stop(name(stack_id, shape_handle))
:noproc
catch
:exit, {:noproc, _} -> :noproc
end
consumer_pid when is_pid(consumer_pid) ->
GenServer.call(consumer_pid, :stop_and_clean, 30_000)
end
end
def init(config) when is_map(config) do
%{shape_handle: shape_handle, storage: {_, _} = storage, shape: shape} = config
Process.set_label({:consumer_supervisor, shape_handle})
metadata = [stack_id: config.stack_id, shape_handle: shape_handle]
Logger.metadata(metadata)
Electric.Telemetry.Sentry.set_tags_context(metadata)
shape_storage = Electric.ShapeCache.Storage.for_shape(shape_handle, storage)
shape_config = %{config | storage: shape_storage}
children = [
{Electric.ShapeCache.Storage, shape_storage},
{Electric.Shapes.Consumer, shape_config},
{Electric.Shapes.Consumer.Snapshotter, shape_config}
]
children =
if should_enable_compaction?(shape, config) do
children ++
[{Electric.ShapeCache.CompactionRunner, [{:storage, shape_storage} | metadata]}]
else
children
end
Supervisor.init(children, strategy: :one_for_one, auto_shutdown: :any_significant)
end
defp should_enable_compaction?(%{storage: %{compaction: :enabled}}, _config), do: true
defp should_enable_compaction?(%{storage: %{compaction: :disabled}}, _config), do: false
# Old shapes don't get compaction by default.
defp should_enable_compaction?(_, _), do: false
end