Packages
spawn
2.0.0-RC2
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/entity/lifecycle/stream_consumer.ex
defmodule Actors.Actor.Entity.Lifecycle.StreamConsumer do
@moduledoc false
use Broadway
alias Broadway.Message
alias Spawn.Utils.Nats
alias Spawn.Fact
alias Google.Protobuf.Timestamp
@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
},
concurrency: build_concurrency(opts)
],
processors: [
default: [concurrency: build_concurrency(opts)]
],
batchers: [
default: [
concurrency: build_concurrency(opts),
# Avoi 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
message
|> build_fact()
|> Message.configure_ack(on_success: :term)
end
@spec handle_batch(any(), Message.t(), any(), opts()) :: list(Message.t())
def handle_batch(_, messages, _, context) do
GenServer.cast(context.projection_pid, {:process_projection_events, messages})
messages
end
@spec build_fact(Message.t()) :: Message.t()
defp build_fact(message) do
# %Broadway.Message{data: "{\"ACTION\":\"KEY_ADDED\",\"KEY\":\"MYKEY\",\"VALUE\":\"MYVALUE\"}", metadata: %{headers: [], topic: "actors.mike"}, acknowledger: {OffBroadway.Jetstream.Acknowledger, #Reference<0.743380651.807927811.227242>, %{on_success: :term, reply_to: "$JS.ACK.newtest.projectionviewertest.1.11.11.1725657673932595345.21"}}, batcher: :default, batch_key: :default, batch_mode: :bulk, status: :ok}
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()
end