Packages
spawn
0.5.0-rc.8
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/state_manager.ex
if Code.ensure_loaded?(Statestores.Supervisor) do
defmodule Actors.Actor.StateManager do
@behaviour Actors.Actor.StateManager.Behaviour
require Logger
alias Eigr.Functions.Protocol.Actors.ActorState
alias Google.Protobuf.Any
alias Statestores.Schemas.Event
alias Statestores.Manager.StateManager, as: StateStoreManager
@impl true
def is_new?(_old_hash, new_state) when is_nil(new_state), do: false
def is_new?(old_hash, new_state) do
with bytes_from_state <- Any.encode(new_state),
hash <- :crypto.hash(:sha256, bytes_from_state) do
old_hash != hash
else
_ ->
false
end
catch
_kind, error ->
{:error, error}
end
@impl true
@spec load(String.t()) :: {:ok, any}
def load(name) do
case StateStoreManager.load(name) do
%Event{revision: _rev, tags: tags, data_type: type, data: data} = _event ->
{:ok,
ActorState.new(tags: tags, state: Google.Protobuf.Any.new(type_url: type, value: data))}
_ ->
{:not_found, %{}}
end
catch
_kind, error ->
{:error, error}
end
@impl true
@spec save(String.t(), Eigr.Functions.Protocol.Actors.ActorState.t()) ::
{:ok, Eigr.Functions.Protocol.Actors.ActorState.t()}
| {:error, any(), Eigr.Functions.Protocol.Actors.ActorState.t()}
def save(_name, nil), do: {:ok, nil}
def save(_name, %ActorState{state: actor_state} = _state)
when is_nil(actor_state) or actor_state == %{},
do: {:ok, actor_state}
def save(name, %ActorState{tags: tags, state: actor_state} = _state) do
Logger.debug("Saving state for actor #{name}")
with bytes_from_state <- Any.encode(actor_state),
hash <- :crypto.hash(:sha256, bytes_from_state) do
%Event{
actor: name,
revision: 0,
tags: tags,
data_type: actor_state.type_url,
data: actor_state.value
}
|> StateStoreManager.save()
|> case do
{:ok, _event} ->
{:ok, actor_state, hash}
{:error, changeset} ->
{:error, changeset, actor_state, hash}
other ->
{:error, other, actor_state}
end
end
catch
_kind, error ->
{:error, error, actor_state}
end
@impl true
@spec save_async(String.t(), Eigr.Functions.Protocol.Actors.ActorState.t()) ::
{:ok, Eigr.Functions.Protocol.Actors.ActorState.t()}
| {:error, any(), Eigr.Functions.Protocol.Actors.ActorState.t()}
def save_async(name, state, timeout \\ 5000)
def save_async(_name, nil, _timeout), do: {:ok, %{}}
def save_async(_name, %ActorState{state: actor_state} = _state, _timeout)
when is_nil(actor_state) or actor_state == %{},
do: {:ok, actor_state}
def save_async(name, %ActorState{tags: tags, state: actor_state} = _state, timeout) do
parent = self()
persist_data_task =
Task.async(fn ->
Logger.debug("Saving state for actor #{name}")
%Event{
actor: name,
revision: 0,
tags: tags,
data_type: actor_state.type_url,
data: actor_state.value
}
|> StateStoreManager.save()
end)
try do
res = Task.await(persist_data_task, timeout)
with bytes_from_state <- Any.encode(actor_state),
hash <- :crypto.hash(:sha256, bytes_from_state) do
if inserted_successfully?(parent, persist_data_task.pid) do
case res do
{:ok, _event} ->
{:ok, actor_state, hash}
{:error, changeset} ->
{:error, changeset, actor_state, hash}
other ->
{:error, other, actor_state}
end
else
{:error, :unsuccessfully, hash}
end
end
catch
_kind, error ->
Task.shutdown(persist_data_task, :brutal_kill)
{:error, error, actor_state}
end
end
defp inserted_successfully?(ref, pid) do
receive do
{^ref, :ok} -> true
{^ref, _} -> false
{:EXIT, ^pid, _} -> false
end
end
end
else
defmodule Actors.Actor.StateManager do
@behaviour Actors.Actor.StateManager.Behaviour
@not_loaded_message """
Statestores not loaded properly
If you are creating actors with flag `persistent: true` consider adding :spawn_statestores to your deps list
"""
def is_new?(_old_hash, _new_state), do: raise(@not_loaded_message)
def load(_key), do: raise(@not_loaded_message)
def save(_name, _state), do: raise(@not_loaded_message)
def save_async(_name, _state, _timeout), do: raise(@not_loaded_message)
end
end