Packages
electric
1.1.11
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.ex
defmodule Electric.Shapes.Monitor do
use Supervisor
alias __MODULE__.RefCounter
@type stack_id :: Electric.stack_id()
@type shape_handle :: Electric.ShapeCache.shape_handle()
@schema NimbleOptions.new!(
stack_id: [type: :string, required: true],
storage: [type: :mod_arg, required: true],
shape_status: [type: :mod_arg, required: true],
publication_manager: [type: :mod_arg, required: true],
on_remove: [type: {:or, [nil, {:fun, 2}]}],
on_cleanup: [type: {:or, [nil, {:fun, 1}]}]
)
def name(stack_id) do
Electric.ProcessRegistry.name(stack_id, __MODULE__)
end
def start_link(args) do
with {:ok, config} <- NimbleOptions.validate(Map.new(args), @schema) do
Supervisor.start_link(__MODULE__, config, name: name(config.stack_id))
end
end
@doc """
Register the current process as a reader of the given shape.
"""
@spec register_reader(stack_id(), shape_handle(), pid()) :: :ok
defdelegate register_reader(stack_id, shape_handle, pid \\ self()), to: RefCounter
@doc """
Unregister the current process as a reader of the given shape.
"""
@spec unregister_reader(stack_id(), shape_handle(), pid()) :: :ok
defdelegate unregister_reader(stack_id, shape_handle, pid \\ self()), to: RefCounter
@doc """
Register the current process as a writer (consumer) of the given shape.
"""
@spec register_writer(stack_id(), shape_handle(), pid()) :: :ok | {:error, term()}
defdelegate register_writer(stack_id, shape_handle, shape, pid \\ self()), to: RefCounter
@doc """
The number of active readers of the given shape.
"""
@spec reader_count(stack_id(), shape_handle()) :: {:ok, non_neg_integer()}
defdelegate reader_count(stack_id, shape_handle), to: RefCounter
@doc """
The number of active readers of all shapes.
"""
@spec reader_count(stack_id()) :: {:ok, non_neg_integer()}
defdelegate reader_count(stack_id), to: RefCounter
@doc """
The number of active readers of all shapes.
"""
@spec reader_count!(stack_id()) :: non_neg_integer()
defdelegate reader_count!(stack_id), to: RefCounter
@doc """
Request a message when all readers of the given handle have finished or terminated.
Sends `{Electric.Shapes.Monitor, :reader_termination, shape_handle, reason}`
to the registered `pid` when the reader count on a shape is `0`.
"""
@spec notify_reader_termination(stack_id(), shape_handle(), term(), pid()) :: :ok
defdelegate notify_reader_termination(stack_id, shape_handle, reason, pid \\ self()),
to: RefCounter
@doc """
clean up the state of a non-running consumer.
"""
@spec purge_shape(stack_id(), shape_handle(), Electric.Shapes.Shape.t()) :: :ok
defdelegate purge_shape(stack_id, shape_handle, shape), to: RefCounter
# used in tests to validate internal state
@doc false
defdelegate termination_watchers(stack_id, shape_handle), to: RefCounter
def init(opts) do
%{
stack_id: stack_id,
storage: storage,
publication_manager: publication_manager,
shape_status: shape_status
} = opts
Process.set_label({:shapes_monitor, stack_id})
children = [
{__MODULE__.CleanupTaskSupervisor, stack_id: stack_id},
{__MODULE__.RefCounter,
stack_id: stack_id,
storage: storage,
publication_manager: publication_manager,
shape_status: shape_status,
on_remove: Map.get(opts, :on_remove),
on_cleanup: Map.get(opts, :on_cleanup)}
]
Supervisor.init(children, strategy: :one_for_one)
end
end