Packages
altworx_runbox
11.0.1
25.0.0
24.0.0
23.1.0
23.0.0
22.2.0
22.1.0
22.0.0
21.2.0
21.1.2
21.1.1
21.1.0
21.0.0
20.0.0
19.0.0
18.0.0
17.2.0
17.1.0
17.0.1
17.0.0
16.2.0
16.1.0
16.0.0
15.0.0
14.1.0
14.0.1
14.0.0
13.0.3
13.0.2
13.0.1
13.0.0
12.1.0
12.0.0
11.0.1
11.0.0
10.0.0
9.0.0
8.0.0
7.0.1
7.0.0
6.0.0
5.0.0
4.0.0
3.0.0
2.1.0
2.0.0
1.4.1
1.4.0
1.3.0
1.2.0
1.1.0
1.0.0
0.1.3
0.1.2
0.1.1
0.1.0
Runbox is a library for running Altworx scenarios.
Current section
Files
Jump to
Current section
Files
lib/runbox/runtime/simple/timezip.ex
defmodule Runbox.Runtime.Simple.Timezip do
@moduledoc group: :internal
@moduledoc """
Timezip component responsible for merging multiple message streams for Simple runtime.
In a Simple scenario, there is only need for a Timezip if reading from multiple input topics. Then
we need to merge these streams into a single stream for the template.
"""
# Disable code duplication checks, since this is mostly copied from Stage. This is OK though,
# since we aim to remove Stage later, so generalizing this would be a waste of time.
# credo:disable-for-this-file Credo.Check.Design.DuplicatedCode
alias Runbox.Deduplicator
alias Runbox.RunStartContext
alias Runbox.Runtime.Timezip.MultiQueue
alias Toolbox.Message
use GenStage
require Logger
@default_max_demand 1000
@doc """
Returns component name.
There are no parameters, since in Simple scenario runtime there is at most one Timezip.
"""
def component_name do
:timezip
end
@doc "Starts the GenStage."
def start_link(args, _runbox_ctx, start_ctx) do
GenStage.start_link(__MODULE__, {args, start_ctx})
end
@impl true
def init(
{%{run_id: run_id, config: %{subscribe_to: subscribe_to, scenario_id: scenario_id}},
start_ctx}
) do
subscribe_to =
Enum.map(subscribe_to, fn component ->
{RunStartContext.component_pid(start_ctx, component), []}
end)
Logger.metadata(run_id: run_id, scenario_id: scenario_id)
{:producer_consumer,
%{
multi_queue: MultiQueue.new(),
deduplicator: Deduplicator.new([], &message_timestamp/1),
run_id: run_id
}, subscribe_to: subscribe_to}
end
@impl true
def handle_events(events, from, %{multi_queue: multi_queue, deduplicator: dedup} = state) do
case zip_through_multiqueue(multi_queue, from, events) do
{:ok, multi_queue, msgs} ->
{msgs, dedup} = deduplicate(msgs, dedup)
{:noreply, msgs, %{state | multi_queue: multi_queue, deduplicator: dedup}}
error ->
{:stop, error, state}
end
end
@impl true
def handle_subscribe(:producer, options, from, state) do
new_multi_queue = MultiQueue.add_queue(state.multi_queue, from)
init_demand = options[:max_demand] || @default_max_demand
GenStage.ask(from, init_demand)
{:manual, %{state | multi_queue: new_multi_queue}}
end
@impl true
def handle_subscribe(:consumer, _options, _from, state) do
{:automatic, state}
end
@impl true
def handle_cancel(_, from, %{multi_queue: multi_queue} = state) do
# well this (cancel) really cant happen during run, and when stopping it may as well be :stop
new_multi_queue = MultiQueue.remove_queue(multi_queue, from)
{:noreply, [], %{state | multi_queue: new_multi_queue}}
end
defp zip_through_multiqueue(multi_queue, queue_id, messages) do
with {:ok, multi_queue} <- MultiQueue.enqueue(multi_queue, queue_id, messages),
{:ok, multi_queue, demands, msgs} <-
MultiQueue.dequeue_all(multi_queue, &message_comparator/2) do
ask_for_more(demands)
{:ok, multi_queue, msgs}
end
end
defp ask_for_more(demand_distribution) do
for {id, number} <- demand_distribution do
GenStage.ask(id, number)
end
end
defp deduplicate(msgs, dedup) do
{msgs, dedup} =
Enum.reduce(msgs, {[], dedup}, fn msg, {msgs, dedup} ->
case Deduplicator.deduplicate(msg, dedup) do
{:new, dedup} ->
{[msg | msgs], dedup}
{:duplicate, dedup} ->
{msgs, dedup}
{:old, dedup} ->
# Should not happen if there is not error in scenario or timezip.
# This may indicate, that a message with timestamp lower than
# timestamps in previous messages was passed to timezip.
Logger.error("Deduplicator reported as old following message #{inspect(msg)}")
{[msg | msgs], dedup}
end
end)
{Enum.reverse(msgs), dedup}
end
defp message_comparator(%Message{timestamp: ts1}, %Message{timestamp: ts2}) do
ts1 < ts2
end
defp message_timestamp(%Message{timestamp: ts}) do
ts
end
end