defmodule GenBuffer do @moduledoc """ Documentation for GenBuffer. ## Usage defmodule ExdStreams.Processing.Writer do use GenBuffer, block: true, interval: 1000, limit: 10 @impl true def flush(events) do :ok end end """ defmacro __using__(opts) do quote bind_quoted: [opts: opts] do use GenServer @static_opts opts {otp_app, adapter} = GenBuffer.Config.compile_config(__MODULE__, opts) @otp_app otp_app @adapter adapter def __adapter__, do: @adapter def child_spec(opts) do %{ id: __MODULE__, start: {__MODULE__, :start_link, [opts]}, type: :worker } end ## Client @doc """ """ def start_link(opts \\ []) do name = Keyword.get(opts, :name, __MODULE__) GenServer.start_link(__MODULE__, opts, name: name) end @doc """ """ def add(event) when not is_list(event) do GenServer.call(__MODULE__, {:add, [event]}) end def add(events) do GenServer.call(__MODULE__, {:add, events}) end ## Server @impl true def init(opts) do limit = Keyword.fetch!(@static_opts, :limit) interval = Keyword.fetch!(@static_opts, :interval) Process.send_after(self(), :flush, interval) queue = :queue.new() state = %{ queue: queue, count: 0, limit: limit, interval: interval, blocks: [] } {:ok, state} end @impl true def handle_call({:add, events}, from, state) do blocks = state.blocks ++ [from] queue = Enum.reduce(events, state.queue, fn event, queue -> :queue.in(event, queue) end) count = state.count + length(events) new_state = %{state | queue: queue, count: count, blocks: blocks} if count >= state.limit do handle_info(:flush, new_state) else {:noreply, new_state} end end @impl true def handle_info(:flush, state) do flushed = :queue.to_list(state.queue) if length(flushed) > 0, do: __MODULE__.flush(flushed) for pid <- state.blocks, do: GenServer.reply(pid, :ok) Process.send_after(self(), :flush, state.interval) {:noreply, %{state | queue: :queue.new(), count: 0, blocks: []}} end end end end