Packages
electric
1.1.10
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/monitor/cleanup_task_supervisor.ex
defmodule Electric.Shapes.Monitor.CleanupTaskSupervisor do
require Logger
alias Electric.ShapeCache.Storage
@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: {Task.Supervisor, :start_link, [[name: name(stack_id)]]},
type: :supervisor
}
end
def name(stack_id) do
Electric.ProcessRegistry.name(stack_id, __MODULE__)
end
def cleanup_async(
stack_id,
storage_impl,
publication_manager_impl,
shape_status_impl,
shape_handle,
shape,
on_cleanup \\ fn _ -> :ok end
) do
if consumer_alive?(stack_id, shape_handle) do
{:error, "Expected shape #{shape_handle} consumer to not be alive before cleaning shape"}
else
{:ok, _pid} =
Task.Supervisor.start_child(name(stack_id), fn ->
Logger.debug("Cleaning shape data for shape #{inspect(shape_handle)}")
task1 =
Task.Supervisor.async(name(stack_id), fn ->
Logger.metadata(stack_id: stack_id, shape_handle: shape_handle)
cleanup_shape_status(shape_status_impl, shape_handle)
end)
task2 =
Task.Supervisor.async(name(stack_id), fn ->
Logger.metadata(stack_id: stack_id, shape_handle: shape_handle)
cleanup_storage(storage_impl, shape_handle)
end)
task3 =
Task.Supervisor.async(name(stack_id), fn ->
Logger.metadata(stack_id: stack_id, shape_handle: shape_handle)
cleanup_publication_manager(publication_manager_impl, shape_handle, shape)
end)
try do
[task1, task2, task3]
|> Task.await_many(@cleanup_timeout)
catch
:exit, {:timeout, _} ->
Logger.warning(
"Shape cleanup tasks for shape #{shape_handle} timed out after #{@cleanup_timeout}ms"
)
:ok
after
on_cleanup.(shape_handle)
end
end)
:ok
end
end
defp cleanup_shape_status(shape_status_impl, shape_handle) do
{shape_status, shape_status_state} = shape_status_impl
case shape_status.remove_shape(shape_status_state, shape_handle) do
{:ok, _shape} ->
Logger.debug("Deregistered shape #{shape_handle}")
{:error, _reason} ->
# this is actually quite likely as during normal shutdown the shape is removed asap
# this path is just to make sure we do that in case of a crash
Logger.debug(["Shape already de-registered #{shape_handle}"])
end
end
defp cleanup_storage(
storage_impl,
shape_handle
) do
perform_reporting_errors(
fn ->
shape_handle
|> Storage.for_shape(storage_impl)
|> Storage.cleanup!()
end,
"Failed to delete data for shape #{shape_handle}"
)
end
defp cleanup_publication_manager(
publication_manager_impl,
shape_handle,
shape
) do
{publication_manager, publication_manager_opts} = publication_manager_impl
perform_reporting_errors(
fn ->
publication_manager.remove_shape(shape_handle, shape, publication_manager_opts)
end,
"Failed to remove shape #{shape_handle} from publication"
)
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 consumer_alive?(stack_id, shape_handle) do
!is_nil(Electric.Shapes.Consumer.whereis(stack_id, shape_handle))
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