Packages
spawn
1.0.0-rc.20
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/spawn/cluster/node/server.ex
defmodule Spawn.Cluster.Node.Server do
@moduledoc """
Node subscriber
"""
use Gnat.Server
require Logger
require OpenTelemetry.Tracer, as: Tracer
alias Eigr.Functions.Protocol.{InvocationRequest, ActorInvocationResponse}
def request(%{topic: topic, body: body, reply_to: reply_to} = req)
when is_binary(body) do
Logger.info("Received Actor Invocation via Nats on #{topic}")
headers = Map.get(req, :headers, [])
handle_request(topic, body, reply_to, headers)
end
def request(%{topic: topic, body: _, reply_to: _reply_to} = _req) do
Logger.warning("Received Invalid Actor Invocation via Nats on #{topic}")
{:reply, {:error, :bad_request}}
end
def error(%{gnat: gnat, reply_to: reply_to}, error) do
Logger.error(
"Error on #{inspect(__MODULE__)} during handle incoming message. Error #{inspect(error)}"
)
Gnat.pub(gnat, reply_to, error)
end
defp handle_request(_topic, body, reply_to, headers) do
opts = headers_to_opts(headers)
Tracer.with_span opts[:span_ctx], "Handle Actor Invoke", kind: :server do
request = InvocationRequest.decode(body)
Actors.Actor.CallerConsumer.invoke_with_span(request, opts)
|> handle_reply(reply_to)
end
end
defp headers_to_opts(headers) do
ctx =
if Enum.any?(headers, fn {k, _v} -> k == "traceparent" end) do
:otel_propagator_text_map.extract(headers)
else
OpenTelemetry.Ctx.new()
end
Keyword.put([], :span_ctx, ctx)
end
defp handle_reply(_response, nil), do: :ok
defp handle_reply(response, _reply_to) do
case response do
{:ok, :async} ->
{:reply, "async"}
{:ok, response} ->
{:reply, ActorInvocationResponse.encode(response)}
{:error, error} ->
{:reply, error}
end
end
end