Packages
extreme
1.1.0-rc8
1.1.4
1.1.3
1.1.2
1.1.1
1.1.1-rc01
1.1.0
1.1.0-rc9
1.1.0-rc8
1.1.0-rc7
1.1.0-rc6
1.1.0-rc5
1.1.0-rc4
1.1.0-rc3
1.1.0-rc2
1.1.0-rc1
1.0.7
1.0.6
1.0.5
1.0.4
1.0.3
1.0.2
1.0.1
1.0.0
0.13.4
0.13.3
0.13.2
0.13.1
0.13.0
0.12.1
0.12.0
0.11.0
0.10.4
0.10.3
0.10.2
0.10.1
0.10.0
0.9.2
0.9.1
0.9.0
0.8.1
0.8.0
0.7.1
0.7.0
0.6.2
0.6.1
0.6.0
0.5.5
0.5.4
0.5.3
0.5.2
0.5.1
0.5.0
0.4.3
0.4.2
0.4.1
Elixir TCP client for EventStore.
Current section
Files
Jump to
Current section
Files
lib/extreme/subscriptions_supervisor.ex
defmodule Extreme.SubscriptionsSupervisor do
use DynamicSupervisor
alias Extreme.{Subscription, ReadingSubscription, PersistentSubscription}
def _name(base_name),
do: Module.concat(base_name, SubscriptionsSupervisor)
def start_link(base_name),
do: DynamicSupervisor.start_link(__MODULE__, :ok, name: _name(base_name))
def init(:ok),
do: DynamicSupervisor.init(strategy: :one_for_one)
def start_subscription(
base_name,
correlation_id,
subscriber,
stream,
resolve_link_tos,
ack_timeout \\ 5_000
) do
base_name
|> _name()
|> DynamicSupervisor.start_child(%{
id: Subscription,
start:
{Subscription, :start_link,
[base_name, correlation_id, subscriber, stream, resolve_link_tos, ack_timeout]},
restart: :temporary
})
end
def start_reading_subscription(base_name, correlation_id, subscriber, read_params) do
base_name
|> _name()
|> DynamicSupervisor.start_child(%{
id: ReadingSubscription,
start:
{ReadingSubscription, :start_link, [base_name, correlation_id, subscriber, read_params]},
restart: :temporary
})
end
def start_persistent_subscription(
base_name,
correlation_id,
subscriber,
stream,
group,
allowed_in_flight_messages
)
when is_binary(stream) and is_binary(group) and is_integer(allowed_in_flight_messages) do
base_name
|> _name()
|> DynamicSupervisor.start_child(%{
id: PersistentSubscription,
start:
{PersistentSubscription, :start_link,
[base_name, correlation_id, subscriber, stream, group, allowed_in_flight_messages]},
restart: :temporary
})
end
def stop_subscription(_base_name, nil),
do: :ok
def stop_subscription(base_name, pid) do
base_name
|> _name()
|> DynamicSupervisor.terminate_child(pid)
end
def kill_all_subscriptions(base_name) do
name = _name(base_name)
name
|> Supervisor.which_children()
|> Enum.each(fn
{_, pid, :worker, _} -> DynamicSupervisor.terminate_child(name, pid)
end)
end
end