Packages
riptide
0.2.79
0.5.2
0.5.1
0.5.0-beta9
0.5.0-beta8
0.5.0-beta7
0.5.0-beta6
0.5.0-beta5
0.5.0-beta4
0.5.0-beta3
0.5.0-beta2
0.5.0-beta11
0.5.0-beta10
0.5.0-beta
0.4.6
0.4.5
0.4.4
0.4.3
0.4.2
0.4.1
0.4.0
0.3.13
0.3.12
0.3.11
0.3.10
0.3.9
0.3.8
0.3.7
0.3.6
0.3.5
0.3.4
0.3.3
0.3.2
0.3.1
0.3.0
0.3.0-bd63a38
0.2.79
0.2.78
0.2.74
0.2.4
0.2.3
0.2.2
0.2.1
0.2.0
0.1.15
0.1.14
0.1.13
0.1.12
0.1.11
0.1.10
0.1.9
0.1.8
0.1.7
0.1.6
0.1.5
0.1.4
0.1.3
0.1.2
0.1.1
0.1.0
A data first framework for building realtime applications
Current section
Files
Jump to
Current section
Files
lib/riptide/scheduler/dispatch.ex
defmodule Riptide.Scheduler.Dispatch do
require Logger
use GenServer
def start_link(_opts) do
GenServer.start_link(__MODULE__, [], name: __MODULE__)
end
def init(_) do
:timer.send_interval(:timer.seconds(1), :poll)
{:ok, %{active: %{}}}
end
def handle_info(:poll, state) do
case active?() and Riptide.Config.riptide_scheduler() do
true ->
now = :os.system_time(:millisecond)
execute =
Riptide.Scheduler.stream()
|> Stream.filter(fn {task, _info} -> Dynamic.get(state, [:active, task]) == nil end)
|> Stream.filter(fn {_task, info} -> info["timestamp"] <= now end)
|> Stream.map(fn {task, _info} -> task end)
|> Enum.to_list()
Enum.each(execute, fn task ->
nodes()
|> Enum.random()
|> case do
node ->
Task.Supervisor.async_nolink(
{Riptide.Scheduler, node},
__MODULE__,
:execute,
[task]
)
end
end)
{:noreply,
%{
state
| active:
execute
|> Stream.map(fn task -> {task, true} end)
|> Enum.into(state.active)
}}
false ->
{:noreply, state}
end
end
def handle_info({_ref, {:finish, task, _result}}, state) do
{:noreply, %{state | active: Map.delete(state.active, task)}}
end
def handle_info(_msg, state) do
{:noreply, state}
end
def execute(task) do
{:finish, task, Riptide.Scheduler.execute(task)}
end
def nodes() do
[Node.self() | Node.list()]
|> Enum.sort()
end
def active?() do
Enum.at(nodes(), 0) == Node.self()
end
end