Packages
spawn
1.0.0-rc.36
2.0.0-RC9
2.0.0-RC8
2.0.0-RC7
2.0.0-RC6
2.0.0-RC5
2.0.0-RC4
2.0.0-RC3
2.0.0-RC2
2.0.0-RC14
2.0.0-RC13
2.0.0-RC12
2.0.0-RC11
2.0.0-RC10
2.0.0-RC1
1.4.3
1.4.2
1.4.1
1.4.0
1.3.3
1.3.2
1.3.1
1.3.0
1.2.1
1.2.0
1.1.1
1.1.0
1.0.1
1.0.0
1.0.0-rc3
1.0.0-rc16
1.0.0-rc1
1.0.0-rc.38
1.0.0-rc.37
1.0.0-rc.36
1.0.0-rc.35
1.0.0-rc.34
1.0.0-rc.33
1.0.0-rc.32
1.0.0-rc.31
1.0.0-rc.30
1.0.0-rc.29
1.0.0-rc.28
1.0.0-rc.27
1.0.0-rc.26
1.0.0-rc.25
1.0.0-rc.24
1.0.0-rc.23
1.0.0-rc.22
1.0.0-rc.21
1.0.0-rc.20
1.0.0-rc.19
1.0.0-rc.18
1.0.0-rc.17
1.0.0-rc.2
0.6.3
0.6.2
0.6.1
0.6.0
0.5.5
0.5.4
0.5.3
0.5.1
0.5.0
0.5.0-rc.13
0.5.0-rc.12
0.5.0-rc.11
0.5.0-rc.10
0.5.0-rc.9
0.5.0-rc.8
0.5.0-rc.7
0.5.0-rc.6
0.5.0-rc.5
0.5.0-rc.3
0.5.0-alpha.13
0.5.0-alpha.12
0.5.0-alpha.11
0.5.0-alpha.10
0.5.0-alpha.9
0.5.0-alpha.8
0.5.0-alpha.7
0.5.0-alpha.6
0.5.0-alpha.5
0.5.0-alpha.4
0.5.0-alpha.3
0.5.0-alpha.2
0.5.0-alpha.1
0.1.0
Spawn is the core lib for Spawn Actors System
Current section
Files
Jump to
Current section
Files
lib/actors/actor/invocation_scheduler.ex
defmodule Actors.Actor.InvocationScheduler do
@moduledoc """
`InvocationScheduler` is a global process for the cluster that controls
all Actions of type Schedule.
This process is global to allow that even after restarts of a process or restart
of an application we will still be able to perform invocations to actors,
without the need for persistent storage such as a database.
"""
use GenServer, restart: :transient
use Retry
require Logger
alias Actors.Registry.ActorRegistry
alias Eigr.Functions.Protocol.InvocationRequest
@hibernate_delay 20_000
@hibernate_jitter 30_000
@impl true
def init(_arg) do
Process.flag(:trap_exit, true)
Process.flag(:message_queue_data, :off_heap)
{:ok, %{}, {:continue, :init_invocations}}
end
@impl true
def handle_continue(:init_invocations, state) do
# TODO: Fix this module
# schedule_hibernate()
# stored_invocations = ActorRegistry.get_all_invocations()
# Enum.each(stored_invocations, &call_invoke/1)
{:noreply, state}
end
@impl true
def terminate(reason, _state) do
Logger.debug("InvocationScheduler down with reason (#{inspect(reason)})")
end
@impl true
def handle_info({:invoke, decoded_request}, state) do
ActorRegistry.remove_invocation_request(
decoded_request.actor.id,
InvocationRequest.encode(decoded_request)
)
spawn(fn ->
request_to_invoke = %InvocationRequest{decoded_request | scheduled_to: nil, async: true}
Actors.invoke(request_to_invoke)
end)
{:noreply, state}
end
def handle_info(:hibernate, state) do
schedule_hibernate()
{:noreply, state, :hibernate}
end
@impl true
def handle_cast({:schedule, request}, state) do
encoded_request = InvocationRequest.encode(request)
spawn(fn ->
ActorRegistry.register_invocation_request(request.actor.id, encoded_request)
end)
call_invoke(request)
{:noreply, state}
end
defp call_invoke(encoded_request) when is_binary(encoded_request) do
InvocationRequest.decode(encoded_request) |> call_invoke()
end
defp call_invoke(%InvocationRequest{} = decoded_request) do
delay_in_ms =
decoded_request.scheduled_to
|> DateTime.from_unix!(:millisecond)
|> DateTime.diff(DateTime.utc_now(), :millisecond)
if delay_in_ms <= 0 do
Logger.warn("Received negative delayed invocation request (#{delay_in_ms}), invoking now")
Process.send(self(), {:invoke, decoded_request}, [])
else
Process.send_after(self(), {:invoke, decoded_request}, delay_in_ms)
end
end
defp schedule_hibernate() do
Process.send_after(self(), :hibernate, next_hibernate_delay())
end
def next_hibernate_delay(), do: @hibernate_delay + :rand.uniform(@hibernate_jitter)
# Client
def schedule_invoke(%InvocationRequest{} = invocation_request) do
GenServer.cast({:global, __MODULE__}, {:schedule, invocation_request})
end
def child_spec do
%{
id: __MODULE__,
start: {__MODULE__, :start_link, []},
restart: :transient
}
end
def start_link() do
GenServer.start_link(__MODULE__, [], name: {:global, __MODULE__})
end
end