Packages
spawn
1.1.1
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/pubsub.ex
defmodule Actors.Actor.Pubsub do
use GenServer
require Logger
alias Eigr.Functions.Protocol.InvocationRequest
alias Eigr.Functions.Protocol.Actors.ActorSystem
alias Eigr.Functions.Protocol.Actors.Actor
alias Eigr.Functions.Protocol.Actors.ActorId
@default_pubsub_group :actor_channel
@pubsub Application.compile_env(:spawn, :pubsub_group, @default_pubsub_group)
def start_link(_opts) do
GenServer.start_link(__MODULE__, nil, name: __MODULE__)
end
def publish(topic, payload, request) do
Phoenix.PubSub.broadcast(
@pubsub,
topic,
{:receive, payload, request.actor},
Actors.Actor.PubsubDispatcher
)
end
@doc """
Subscribes a specific actor to a topic
"""
def subscribe(topic, actor_name, system, action_handler \\ nil) do
GenServer.cast(__MODULE__, {:subscribe, topic, actor_name, system, action_handler})
end
@impl true
def init(_opts) do
{:ok, %{}}
end
@impl true
def handle_cast({:subscribe, topic, actor_name, system, action_handler}, state) do
metadata = %{actor_name: actor_name, system: system, action: action_handler}
key = :erlang.phash2(metadata)
if Map.get(state, key) do
{:noreply, state}
else
Phoenix.PubSub.subscribe(@pubsub, topic, metadata: metadata)
{:noreply, Map.put(state, key, true)}
end
end
@impl true
def handle_info({{:receive, payload, caller}, metadata}, state) do
action =
Map.get(metadata, :action)
|> case do
nil -> "receive"
"" -> "receive"
action -> action
end
actor_name = Map.get(metadata, :actor_name)
system = Map.get(metadata, :system)
Logger.debug(
"Actor [#{actor_name}] Received Broadcast Event to perform Action [#{action}] from caller #{inspect(caller)}"
)
invocation = %InvocationRequest{
system: %ActorSystem{name: system},
actor: %Actor{
id: %ActorId{name: actor_name, system: system}
},
action_name: action,
payload: payload,
caller: caller,
async: true
}
{:ok, :async} = Actors.invoke(invocation)
{:noreply, state}
end
end