Packages
electric
1.6.8
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/shape_cache/shape_cleaner/cleanup_task_supervisor.ex
defmodule Electric.ShapeCache.ShapeCleaner.CleanupTaskSupervisor do
require Logger
@env Mix.env()
# set a high timeout (except for tests) for the shape cleanup to terminate
# we don't want to see errors due to e.g. a slow filesystem.
# any actual errors in the processes will be caught and reported
@cleanup_timeout if @env != :test, do: 60_000, else: 3_000
def child_spec(opts) do
{:ok, stack_id} = Keyword.fetch(opts, :stack_id)
%{
id: {__MODULE__, stack_id},
start: {__MODULE__, :start_link, [opts]},
type: :supervisor
}
end
def start_link(opts) do
{:ok, stack_id} = Keyword.fetch(opts, :stack_id)
if on_cleanup_callback = Keyword.get(opts, :on_cleanup, nil) do
Electric.StackConfig.put(
stack_id,
{Electric.ShapeCache.ShapeCleaner, :on_cleanup},
on_cleanup_callback
)
end
Task.Supervisor.start_link(name: name(stack_id))
end
def name(stack_id) do
Electric.ProcessRegistry.name(stack_id, __MODULE__)
end
def perform_async(stack_id, fun) do
with {:ok, _pid} <-
Task.Supervisor.start_child(name(stack_id), fn ->
set_task_metadata(stack_id)
fun.()
end) do
:ok
end
end
def cleanup_async(stack_id, shape_handles) when is_list(shape_handles) do
perform_async(stack_id, fn ->
cleanup_callback = on_cleanup_callback(stack_id)
set_task_metadata(stack_id)
tasks = [
async(stack_id, shape_handles, ¬ify_shape_rotation/2),
async(stack_id, shape_handles, &cleanup_publication_manager/2)
]
try do
Task.await_many(tasks, @cleanup_timeout)
catch
:exit, {:timeout, _} ->
Logger.warning(
"Shape cleanup tasks for #{length(shape_handles)} shapes timed out after #{@cleanup_timeout}ms"
)
:ok
after
Enum.each(shape_handles, cleanup_callback)
end
end)
end
defp notify_shape_rotation(stack_id, shape_handles) do
Enum.each(shape_handles, fn shape_handle ->
Registry.dispatch(
Electric.StackSupervisor.registry_name(stack_id),
shape_handle,
fn registered ->
Logger.debug(fn ->
"Notifying ~#{length(registered)} clients about removal of shape #{shape_handle}"
end)
for {pid, ref} <- registered, do: send(pid, {ref, :shape_rotation})
end
)
end)
end
defp cleanup_publication_manager(stack_id, shape_handles) do
Enum.each(shape_handles, fn shape_handle ->
perform_reporting_errors(
fn ->
Electric.Replication.PublicationManager.remove_shape(stack_id, shape_handle)
end,
"Failed to remove shape #{shape_handle} from publication"
)
end)
end
defp async(stack_id, shape_handles, fun) do
Task.Supervisor.async(name(stack_id), fn ->
set_task_metadata(stack_id)
fun.(stack_id, shape_handles)
end)
end
defp set_task_metadata(stack_id) do
Logger.metadata(stack_id: stack_id)
Electric.Telemetry.Sentry.set_tags_context(stack_id: stack_id)
end
defp perform_reporting_errors(fun, message) do
try do
fun.()
catch
kind, reason when kind in [:exit, :error] ->
log_error(kind, [message, ": ", Exception.format(kind, reason, __STACKTRACE__)])
{:error, reason}
end
end
defp on_cleanup_callback(stack_id) do
Electric.StackConfig.lookup(stack_id, {Electric.ShapeCache.ShapeCleaner, :on_cleanup}, fn _ ->
:ok
end)
end
if @env == :test do
# don't spam test logs with failures due to process shutdown
defp log_error(:exit, _message) do
:ok
end
else
# don't spam sentry with errors caused by shutdown order (i.e. when the
# publication manager has been shutdown)
defp log_error(:exit, message) do
Logger.log(:warning, message)
end
end
defp log_error(:error, message) do
Logger.log(:error, message)
end
end