Adds the commercial Oban Pro package (vendored from the towerops-web2 vendor tree) so this project can use its workers, plugins, and Smart engine features.
102 lines
2.7 KiB
Elixir
102 lines
2.7 KiB
Elixir
defmodule Oban.Pro.Limiters.Global do
|
|
@moduledoc false
|
|
|
|
@behaviour Oban.Pro.Limiter
|
|
|
|
alias Oban.Pro.Partition
|
|
|
|
@impl Oban.Pro.Limiter
|
|
def check(_repo, %{producer: producer} = changes) when is_map(producer.meta.global_limit) do
|
|
tracked =
|
|
[producer | changes.all_producers]
|
|
|> Enum.uniq_by(& &1.uuid)
|
|
|> Enum.map(& &1.meta.global_limit)
|
|
|> Enum.filter(&is_map/1)
|
|
|> Enum.reduce(%{}, &merge/2)
|
|
|
|
allowed = producer.meta.global_limit.allowed
|
|
tracked = update_legacy_tracked(tracked)
|
|
|
|
if is_nil(producer.meta.global_limit.partition) do
|
|
{:ok, combined_demands(allowed, tracked)}
|
|
else
|
|
{:ok, separate_demands(changes.conf, allowed, producer, tracked)}
|
|
end
|
|
end
|
|
|
|
def check(_repo, _changes), do: {:ok, nil}
|
|
|
|
defp merge(global_limit, acc) do
|
|
Map.merge(global_limit.tracked, acc, fn _, v1, v2 -> v1 + v2 end)
|
|
end
|
|
|
|
defp combined_demands(allowed, tracked) do
|
|
total = tracked |> Map.values() |> Enum.sum()
|
|
|
|
max(allowed - total, 0)
|
|
end
|
|
|
|
defp separate_demands(conf, allowed, producer, tracked) do
|
|
untracked =
|
|
for key <- Partition.available_keys(conf, producer.queue, producer.meta.local_limit),
|
|
not is_map_key(tracked, key),
|
|
into: %{},
|
|
do: {key, allowed}
|
|
|
|
allowed =
|
|
if producer.meta.global_limit.burst do
|
|
size =
|
|
tracked
|
|
|> Map.merge(untracked)
|
|
|> Map.delete("none")
|
|
|> map_size()
|
|
|> max(1)
|
|
|
|
producer.meta.local_limit
|
|
|> div(size)
|
|
|> max(allowed)
|
|
else
|
|
allowed
|
|
end
|
|
|
|
Enum.reduce(tracked, untracked, fn {key, total}, acc ->
|
|
case max(allowed - total, 0) do
|
|
demand when demand > 0 -> Map.put(acc, key, demand)
|
|
_ -> acc
|
|
end
|
|
end)
|
|
end
|
|
|
|
# Tracking must be updated after acking and before fetching in order to have accurate limits for
|
|
# the current fetch transaction.
|
|
def prepare_tracked(%{global_limit: global} = meta, keys) do
|
|
tracked =
|
|
Enum.reduce(keys, global.tracked, fn key, acc ->
|
|
if is_map_key(acc, key) do
|
|
Map.update!(acc, key, &(&1 - 1))
|
|
else
|
|
acc
|
|
end
|
|
end)
|
|
|
|
put_in(meta.global_limit.tracked, tracked)
|
|
end
|
|
|
|
@impl Oban.Pro.Limiter
|
|
def track(%{global_limit: global} = meta, jobs) when is_map(global) do
|
|
tracked = Enum.frequencies_by(jobs, &Partition.get_key(&1, "*"))
|
|
|
|
put_in(meta.global_limit.tracked, tracked)
|
|
end
|
|
|
|
def track(meta, _jobs), do: meta
|
|
|
|
defp update_legacy_tracked(%{"count" => count}), do: %{"none" => count}
|
|
|
|
defp update_legacy_tracked(tracked) do
|
|
Map.new(tracked, fn
|
|
{key, %{"count" => count}} -> {key, count}
|
|
tuple -> tuple
|
|
end)
|
|
end
|
|
end
|