Packages
finitomata
0.23.5
0.41.0
0.40.0
0.35.0
0.34.0
0.33.0
0.32.0
0.31.1
0.30.3
0.30.2
0.30.1
0.30.0
0.29.10
0.29.9
0.29.8
0.29.7
0.29.6
0.29.5
0.29.4
0.29.3
0.29.2
0.29.1
0.29.0
0.28.1
0.28.0
0.27.1
0.27.0
0.26.4
0.26.3
0.26.2
0.26.1
0.26.0
0.25.0
0.24.4
0.24.3
0.24.2
0.24.1
0.24.0
0.23.7
0.23.6
0.23.5
0.23.4
0.23.3
0.23.2
0.23.1
0.23.0
0.22.1
0.22.0
0.21.4
0.21.3
0.21.2
0.21.1
0.21.0
0.20.2
0.20.1
0.20.0
0.19.6
0.19.5
0.19.4
0.19.3
0.19.2
0.19.1
0.19.0
0.18.4
0.18.3
0.18.2
0.18.1
0.18.0
0.17.1
0.17.0
0.16.0
0.15.1
0.15.0
0.14.6
0.14.5
0.14.4
0.14.3
0.14.2
0.14.1
0.14.0
0.13.0
0.12.1
0.12.0
0.11.3
0.11.2
0.11.1
0.11.0
0.10.0
0.9.1
0.9.0
0.8.2
0.8.1
0.8.0
0.7.2
0.7.1
0.7.0
0.6.3
0.6.2
0.6.1
0.6.0
0.5.2
0.5.1
0.5.0
0.4.0
0.3.0
0.2.0
0.1.1
0.1.0
The FSM implementation generated from PlantUML textual representation.
Current section
Files
Jump to
Current section
Files
lib/finitomata/throttler/consumer.ex
defmodule Finitomata.Throttler.Consumer do
@moduledoc false
use GenStage
alias Finitomata.Throttler
@throttler_options Application.compile_env(:finitomata, :throttler, [])
@max_demand Keyword.get(@throttler_options, :max_demand, 5)
@interval Keyword.get(@throttler_options, :interval, 400)
def start_link(initial \\ :ok),
do: GenStage.start_link(__MODULE__, initial)
@impl GenStage
def init(opts) do
max_demand = Keyword.get(opts, :max_demand, @max_demand)
interval = Keyword.get(opts, :interval, @interval)
{:consumer, %{__throttler_options__: %{max_demand: max_demand, interval: interval}}}
end
@impl GenStage
def handle_subscribe(:producer, opts, from, producers) do
max_demand =
Keyword.get_lazy(opts, :max_demand, fn ->
get_in(producers, ~w|__throttler_options__ max_demand|a)
end)
interval =
Keyword.get_lazy(opts, :interval, fn ->
get_in(producers, ~w|__throttler_options__ interval|a)
end)
producers =
producers
|> Map.put(from, {max_demand, interval})
|> ask_and_schedule(from)
# `manual` to control over the demand
{:manual, producers}
end
@impl GenStage
def handle_cancel(_, from, producers),
do: {:noreply, [], Map.delete(producers, from)}
@impl GenStage
def handle_events(events, from, producers) do
producers =
Map.update!(producers, from, fn {pending, interval} ->
{pending + length(events), interval}
end)
perform(events)
{:noreply, [], producers}
end
@impl GenStage
def handle_info({:ask, from}, producers),
do: {:noreply, [], ask_and_schedule(producers, from)}
defp ask_and_schedule(
%{__throttler_options__: %{max_demand: max_demand, interval: interval}} = producers,
from
) do
case producers do
%{^from => {0, _interval}} ->
GenStage.ask(from, max_demand)
Process.send_after(self(), {:ask, from}, interval)
producers
%{^from => {pending, interval}} ->
GenStage.ask(from, pending)
Process.send_after(self(), {:ask, from}, interval)
Map.put(producers, from, {0, interval})
%{} ->
producers
end
end
@spec perform([Throttler.t()]) :: :ok
defp perform(events) do
Enum.each(events, fn %Throttler{from: from, fun: fun, args: args} = throttler ->
result =
case fun do
f when is_function(f, 0) -> f.()
f when is_function(f, 1) -> f.(args)
{mod, fun} when is_atom(mod) and is_atom(fun) -> apply(mod, fun, args)
end
case from do
{pid, _alias_ref} when is_pid(pid) ->
GenStage.reply(from, %Throttler{throttler | result: result})
nil ->
Throttler.debug(%Throttler{throttler | result: result}, label: "Malformed owner")
end
end)
end
end