defmodule Oban.Pro.Limiters.Rate 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.rate_limit) do %{allowed: allowed, period: period, window_time: prev_time} = producer.meta.rate_limit curr_time = unix_now() weight = div(max(period - curr_time - prev_time, 0), period) windows = [producer | changes.all_producers] |> Enum.uniq_by(& &1.uuid) |> Enum.map(& &1.meta.rate_limit) |> Enum.filter(&(is_map(&1) and &1.window_time >= curr_time - &1.period)) |> Enum.reduce(%{}, &merge/2) windows = update_legacy_windows(windows) if is_nil(producer.meta.rate_limit.partition) do {:ok, combined_demands(weight, allowed, windows)} else {:ok, separate_demands(changes.conf, weight, allowed, producer, windows)} end end def check(_repo, _changes), do: {:ok, nil} defp merge(rate_limit, acc) do Map.merge(rate_limit.windows, acc, fn _key, v1, v2 -> %{"curr_count" => c1, "prev_count" => p1} = v1 %{"curr_count" => c2, "prev_count" => p2} = v2 %{"curr_count" => c1 + c2, "prev_count" => p1 + p2} end) end defp update_legacy_windows(%{"args" => _} = window), do: %{"none" => window} defp update_legacy_windows(windows), do: windows defp combined_demands(weight, allowed, windows) do total = Enum.reduce(windows, 0, fn {_, val}, acc -> acc + (val["prev_count"] * weight + val["curr_count"]) end) max(allowed - total, 0) end defp separate_demands(conf, weight, allowed, producer, windows) do untracked = for key <- Partition.available_keys(conf, producer.queue, producer.meta.local_limit), not is_map_key(windows, key), into: %{}, do: {key, allowed} Enum.reduce(windows, untracked, fn {key, val}, acc -> %{"curr_count" => curr, "prev_count" => prev} = val case max(allowed - (prev * weight + curr), 0) do demand when demand > 0 -> Map.put(acc, key, demand) _ -> acc end end) end @impl Oban.Pro.Limiter def track(%{rate_limit: rate_limit} = meta, jobs) when is_map(rate_limit) do new_mapping = Enum.reduce(jobs, %{}, fn job, acc -> key = Partition.get_key(job, "*") map = %{"curr_count" => 1, "prev_count" => 0} Map.update(acc, key, map, &%{&1 | "curr_count" => &1["curr_count"] + 1}) end) new_windows = for {key, val} <- new_mapping, not is_map_key(rate_limit.windows, key), into: %{}, do: {key, val} prev_unix = rate_limit.window_time unix_diff = unix_now() - prev_unix {mode, next_time} = cond do unix_diff > rate_limit.period * 2 -> {:dump_prev, unix_now()} unix_diff > rate_limit.period -> {:swap_prev, unix_now()} true -> {:bump_curr, prev_unix} end all_windows = for {key, window} <- rate_limit.windows, not match?(%{"curr_count" => 0, "prev_count" => 0}, window), into: new_windows, do: {key, put_next_count(key, window, mode, new_mapping)} %{meta | rate_limit: %{rate_limit | windows: all_windows, window_time: next_time}} end def track(meta, _jobs), do: meta defp put_next_count(key, window, mode, mapping) do %{"curr_count" => curr_count, "prev_count" => prev_count} = window next_count = get_in(mapping, [key, "curr_count"]) || 0 {curr_count, prev_count} = case mode do :dump_prev -> {next_count, 0} :swap_prev -> {next_count, curr_count} :bump_curr -> {curr_count + next_count, prev_count} end %{window | "curr_count" => curr_count, "prev_count" => prev_count} end defp unix_now do DateTime.to_unix(DateTime.utc_now(), :second) end end