Current section

Files

Jump to
spawn lib actors actor entity lifecycle stream_consumer.ex
Raw

lib/actors/actor/entity/lifecycle/stream_consumer.ex

defmodule Actors.Actor.Entity.Lifecycle.StreamConsumer do
@moduledoc false
use Broadway
require Logger
require OpenTelemetry.Tracer, as: Tracer
alias Broadway.Message
alias Spawn.Utils.Nats
alias Spawn.Fact
alias Google.Protobuf.Timestamp
alias Sidecar.GracefulShutdown
alias Spawn.Actors.ActorSystem
alias Spawn.Actors.Actor
alias Spawn.Actors.ActorId
alias Spawn.InvocationRequest
@type fact :: %Fact{}
@type opts :: %{
projection_pid: pid(),
actor_name: String.t(),
strict_ordering: boolean()
}
@spec start_link(opts :: opts()) :: :ignore | {:error, any()} | {:ok, pid()}
def start_link(opts) do
Broadway.start_link(
__MODULE__,
# there will be not a lot so probably fine to convert to atom
name: String.to_atom(opts.actor_name),
context: opts,
producer: [
module: {
OffBroadway.Jetstream.Producer,
connection_name: Nats.connection_name(),
stream_name: opts.actor_name,
consumer_name: opts.actor_name,
receive_interval: 2_000
},
concurrency: build_concurrency(opts)
],
processors: [
default: [
concurrency: build_concurrency(opts)
]
],
batchers: [
default: [
concurrency: build_concurrency(opts),
# Avoid big batches, micro batches is better
batch_size: 10,
batch_timeout: 2_000
]
]
)
end
@spec handle_message(any(), Message.t(), any()) :: Message.t()
def handle_message(_processor_name, message, _context) do
if GracefulShutdown.running?() do
message
|> build_fact()
|> Message.configure_ack(on_success: :term)
else
message
|> Message.failed("Failed to deliver because app is draining")
end
end
@spec handle_batch(any(), Message.t(), any(), opts()) :: list(Message.t())
def handle_batch(_, messages, _, context) do
Enum.map(messages, fn message ->
try do
process_message(message, context)
message
rescue
error ->
# Let Broadway handle the retry by marking the message as failed
Logger.error("Error processing message: #{inspect(error)}")
Message.failed(message, "#{inspect(error)}")
catch
error ->
Logger.error("Error processing message: #{inspect(error)}")
Message.failed(message, "#{inspect(error)}")
end
end)
end
@spec build_fact(Message.t()) :: Message.t()
defp build_fact(message) do
message
|> Message.put_data(process_data(message))
end
@spec process_data(Message.t()) :: fact()
defp process_data(message) do
payload = message.data
metadata =
Enum.reduce(message.metadata.headers, %{}, fn {key, value}, acc ->
Map.put(acc, key, value)
end)
|> Map.put("topic", message.metadata.topic)
time = DateTime.utc_now() |> DateTime.to_unix(:seconds)
%Fact{
uuid: UUID.uuid4(:hex),
metadata: metadata,
state: payload,
timestamp: %Timestamp{seconds: time}
}
end
# Projections are like long-lasting threads and therefore concurrency should be avoided
# if the intention is to have some notion of ordering.
defp build_concurrency(%{strict_ordering: true}), do: 1
defp build_concurrency(%{strict_ordering: false}), do: System.schedulers_online()
# Process a single message and invoke the actor
defp process_message(%Broadway.Message{data: %Fact{} = message}, state) do
actor_name = state.actor_name |> String.split("-") |> List.last()
actor_settings = :persistent_term.get("actor-#{actor_name}")
system_name = Map.get(message.metadata, "spawn-system")
parent = Map.get(message.metadata, "actor-parent")
name = Map.get(message.metadata, "actor-name")
source_action = Map.get(message.metadata, "actor-action")
action_metadata =
case Map.get(message.metadata, "action-metadata") do
nil -> %{}
metadata -> Jason.decode!(metadata)
end
action =
actor_settings.subjects
|> Enum.find(fn subject -> subject.source_action == source_action end)
|> Map.get(:action)
invocation = %InvocationRequest{
system: %ActorSystem{name: system_name},
actor: %Actor{id: %ActorId{name: actor_name, system: system_name}},
metadata: action_metadata,
action_name: action,
payload: {:value, Google.Protobuf.Any.decode(message.state)},
caller: %ActorId{name: name, system: system_name, parent: parent}
}
# If this raises or throws, the error will be caught in handle_batch
# and the message will be marked as failed for Broadway to retry
{:ok, _response} = Actors.invoke(invocation, span_ctx: Tracer.current_span_ctx())
end
end