Packages
quantum
3.2.0
3.5.3
3.5.2
retired
3.5.1
retired
3.5.0
3.4.0
3.3.0
3.2.0
3.1.0
3.0.2
3.0.1
3.0.0
3.0.0-rc.3
3.0.0-rc.2
3.0.0-rc.1
2.4.0
2.3.4
2.3.3
2.3.2
2.3.1
retired
2.3.0
retired
2.2.7
2.2.6
2.2.5
retired
2.2.4
retired
2.2.3
retired
2.2.2
retired
2.2.1
retired
2.2.0
retired
2.1.3
retired
2.1.2
retired
2.1.1
retired
2.1.0
retired
2.1.0-beta.1
retired
2.0.4
retired
2.0.3
retired
2.0.2
retired
2.0.1
retired
2.0.0
retired
2.0.0-beta.2
retired
2.0.0-beta.1
retired
1.9.3
retired
1.9.2
retired
1.9.1
retired
1.9.0
retired
1.8.1
retired
1.8.0
retired
1.7.1
retired
1.7.0
retired
1.6.1
retired
1.6.0
retired
1.5.0
retired
1.4.0
retired
1.3.2
retired
1.3.1
retired
1.3.0
retired
1.2.4
retired
1.2.3
retired
1.2.2
retired
1.2.1
retired
1.2.0
retired
1.1.0
retired
1.0.4
retired
1.0.3
retired
1.0.2
retired
1.0.1
retired
1.0.0
retired
Cron-like job scheduler for Elixir.
Current section
Files
Jump to
Current section
Files
lib/quantum/task_registry.ex
defmodule Quantum.TaskRegistry do
@moduledoc false
# Registry to check if a task is already running on a node.
require Logger
alias __MODULE__.StartOpts
alias Quantum.Job
# Start the registry
@spec start_link(StartOpts.t()) :: GenServer.on_start()
def start_link(%StartOpts{name: name, listeners: listeners}) do
[keys: :unique, name: name, listeners: listeners]
|> Registry.start_link()
|> case do
{:ok, pid} ->
{:ok, pid}
{:error, {:already_started, pid}} ->
Process.monitor(pid)
{:ok, pid}
{:error, _reason} = error ->
error
end
end
@spec child_spec(options :: StartOpts.t()) :: Supervisor.child_spec()
def child_spec(options),
do:
[]
|> Registry.child_spec()
|> Map.put(:start, {__MODULE__, :start_link, [options]})
# Mark a task as Running
#
# ### Examples
#
# iex> Quantum.TaskRegistry.mark_running(server, running_job.name, Node.self())
# :already_running
#
# iex> Quantum.TaskRegistry.mark_running(server, not_running_job.name, Node.self())
# :marked_running
@spec mark_running(server :: atom, task :: Job.name(), node :: Node.t()) ::
:already_running | :marked_running
def mark_running(server, task, node) do
server
|> Registry.register({task, node}, true)
|> case do
{:ok, _pid} -> :marked_running
{:error, {:already_registered, _other_pid}} -> :already_running
end
end
# Mark a task as Finished
#
# ### Examples
#
# iex> Quantum.TaskRegistry.mark_running(server, running_job.name, Node.self())
# :ok
#
# iex> Quantum.TaskRegistry.mark_running(server, not_running_job.name, Node.self())
# :ok
@spec mark_finished(server :: atom, task :: Job.name(), node :: Node.t()) :: :ok
def mark_finished(server, task, node) do
Registry.unregister(server, {task, node})
end
end