Packages

A generic event buffer behaviour

Current section

Files

Jump to
gen_buffer lib gen_buffer.ex
Raw

lib/gen_buffer.ex

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