Packages

ExWal is a project that aims to provide a solution for managing write-ahead log (WAL) in Elixir.

Current section

Files

Jump to
ex_wal lib ex_wal manager standalone.ex
Raw

lib/ex_wal/manager/standalone.ex

defmodule ExWal.Manager.Standalone do
@moduledoc false
use GenServer, restart: :transient
alias ExWal.FS
alias ExWal.Manager.Options
alias ExWal.Models
alias ExWal.Models.Deletable
alias ExWal.Models.VirtualLog
alias ExWal.Recycler
require Logger
@type t :: %__MODULE__{
name: GenServer.name(),
recycler: ExWal.Recycler.t(),
dynamic_sup: atom(),
registry: atom(),
fs: ExWal.FS.t(),
dirname: binary(),
queue: [Models.VirtualLog.t()],
initial_obsolete: [Models.Deletable.t()]
}
defstruct name: nil,
recycler: nil,
dynamic_sup: nil,
registry: nil,
fs: nil,
dirname: "",
queue: [],
initial_obsolete: []
@spec start_link({
name :: GenServer.name(),
dynamic_sup :: GenServer.name(),
registry :: GenServer.name(),
opts :: Options.t()
}) :: GenServer.on_start()
def start_link({name, dynamic_sup, registry, opts}) do
GenServer.start_link(
__MODULE__,
{name, dynamic_sup, registry, opts},
name: name
)
end
@spec stop(GenServer.name()) :: :ok | {:error, reason :: any()}
def stop(name) do
GenServer.stop(name)
end
@spec get(name :: GenServer.name()) :: t()
def get(name), do: GenServer.call(name, :get_self)
@spec create(name :: GenServer.name(), log_num :: ExWal.Models.VirtualLog.log_num()) ::
{:ok, ExWal.LogWriter.t()} | {:error, reason :: any()}
def create(name, log_num) do
GenServer.call(name, {:create, log_num})
end
@spec obsolete(
name :: GenServer.name(),
min_log_num :: ExWal.Models.VirtualLog.log_num(),
recycle? :: boolean()
) ::
{:ok, [Deletable.t()]} | {:error, reason :: any()}
def obsolete(name, min_log_num, recycle?) do
GenServer.call(name, {:obsolete, min_log_num, recycle?})
end
@spec list(name :: GenServer.name()) :: {:ok, [Models.VirtualLog.t()]} | {:error, reason :: any()}
def list(name) do
GenServer.call(name, :list)
end
# --------------------- server -----------------
@impl GenServer
def init({name, dynamic_sup, registry, opts}) do
%Options{primary: primary} = opts
{:ok,
%__MODULE__{
name: name,
dynamic_sup: dynamic_sup,
registry: registry,
fs: primary[:fs],
dirname: primary[:dir]
}, {:continue, {:recycler, opts}}}
end
@impl GenServer
def terminate(reason, %__MODULE__{recycler: r}) do
may_log_reason(reason)
Recycler.stop(r)
end
@impl GenServer
def handle_continue({:recycler, opts}, state) do
%__MODULE__{
dirname: dirname,
registry: registry
} = state
recycler =
with do
rn = {:via, Registry, {registry, {:recycler, dirname}}}
{:ok, _} = Recycler.ETS.start_link(rn)
Recycler.ETS.get(rn)
end
{:noreply, %__MODULE__{state | recycler: recycler}, {:continue, {:initialize, opts}}}
end
def handle_continue({:initialize, opts}, state) do
%__MODULE__{dirname: dirname, fs: fs, recycler: recycler} = state
FS.mkdir_all(fs, dirname)
with do
%Options{max_num_recyclable_logs: ml} = opts
:ok = Recycler.initialize(recycler, ml)
end
:ok = FS.mkdir_all(fs, dirname)
{:ok, files} = FS.list(fs, dirname)
# build logs
logs =
with do
files
|> Enum.map(fn f -> VirtualLog.parse_filename(f) end)
|> Enum.group_by(fn {log_num, _} -> log_num end)
|> Enum.map(fn {log_num, ms} ->
ss =
ms
|> Enum.map(fn {_, index} -> %Models.Segment{index: index, dir: dirname, fs: fs} end)
|> Enum.sort_by(fn %Models.Segment{index: i} -> i end)
%VirtualLog{log_num: log_num, segments: ss}
end)
end
# set recycler min
Enum.each(logs, fn %VirtualLog{log_num: n} ->
recycler
|> Recycler.get_min()
|> Kernel.<=(n)
|> if do
Recycler.set_min(recycler, n + 1)
end
end)
# add to initial obsolete
initial_obsolete =
Enum.map(logs, fn %VirtualLog{log_num: log_num} ->
%Models.Deletable{fs: fs, path: Path.join(dirname, VirtualLog.filename(log_num, 0)), log_num: log_num}
end)
{:noreply, %__MODULE__{state | initial_obsolete: initial_obsolete}}
end
@impl GenServer
def handle_call(:get_self, _from, state) do
{:reply, state, state}
end
def handle_call({:create, log_num}, _from, state) do
%__MODULE__{
dirname: dirname,
registry: registry,
dynamic_sup: dynamic_sup,
queue: q,
fs: fs
} = state
# new file
writer =
with do
log_name = Path.join(dirname, Models.VirtualLog.filename(log_num, 0))
{:ok, file} = create_or_reuse(log_name, state)
writer_name = {:via, Registry, {registry, {:writer, log_num}}}
{:ok, _} = DynamicSupervisor.start_child(dynamic_sup, {ExWal.LogWriter.Single, {writer_name, file, log_num}})
ExWal.LogWriter.Single.get(writer_name)
end
# add to queue
q = [
%VirtualLog{
log_num: log_num,
segments: [%Models.Segment{index: 0, dir: dirname, fs: fs}]
}
| q
]
{
:reply,
{:ok, writer},
%__MODULE__{state | queue: q}
}
end
def handle_call({:obsolete, min_log_num, true}, _from, state) do
%__MODULE__{
recycler: recycler,
queue: q,
initial_obsolete: init_ob
} = state
to_del =
with do
q
|> Enum.filter(fn %VirtualLog{log_num: n} -> n < min_log_num end)
|> Enum.reduce(
init_ob,
fn l, acc ->
recycler
|> Recycler.add(l)
|> if do
[from_log(l) | acc]
else
acc
end
end
)
end
queue = Enum.reject(q, fn %VirtualLog{log_num: n} -> n < min_log_num end)
{
:reply,
{:ok, to_del},
%__MODULE__{
state
| queue: queue,
initial_obsolete: []
}
}
end
def handle_call({:obsolete, min_log_num, false}, _from, state) do
%__MODULE__{queue: q, initial_obsolete: init_ob} = state
q = Enum.reject(q, fn %VirtualLog{log_num: n} -> n < min_log_num end)
{
:reply,
{:ok, init_ob},
%__MODULE__{
state
| queue: q,
initial_obsolete: []
}
}
end
def handle_call(:list, _from, state) do
%__MODULE__{queue: q} = state
{:reply, {:ok, q}, state}
end
defp create_or_reuse(log_name, %__MODULE__{fs: fs, recycler: recycler, dirname: dirname}) do
recycler
|> Recycler.pop()
|> case do
nil ->
FS.create(fs, log_name)
recycled_log ->
%VirtualLog{log_num: n} = recycled_log
old_name = Path.join(dirname, Models.VirtualLog.filename(n, 0))
FS.reuse_for_write(fs, old_name, log_name)
end
end
@spec from_log(Models.VirtualLog.t()) :: Models.Deletable.t()
defp from_log(%VirtualLog{log_num: l, segments: [%Models.Segment{dir: dir, fs: fs}]}) do
%Models.Deletable{
fs: fs,
path: Path.join(dir, VirtualLog.filename(l, 0)),
log_num: l
}
end
defp may_log_reason(:normal), do: :pass
defp may_log_reason(reason), do: Logger.error("standalone manager terminated, reason: #{inspect(reason)}")
end
defimpl ExWal.Manager, for: ExWal.Manager.Standalone do
alias ExWal.Manager.Standalone
@spec list(ExWal.Manager.t()) :: {:ok, [ExWal.Models.VirtualLog.t()]} | {:error, reason :: any()}
def list(%Standalone{name: name}), do: Standalone.list(name)
@spec obsolete(
impl :: ExWal.Manager.t(),
min_log_num :: ExWal.Models.VirtualLog.log_num(),
recycle? :: boolean()
) ::
{:ok, [ExWal.Models.Deletable.t()]} | {:error, reason :: any()}
def obsolete(%Standalone{name: name}, min_log_num, recycle?), do: Standalone.obsolete(name, min_log_num, recycle?)
@spec create(impl :: ExWal.Manager.t(), log_num :: ExWal.Models.VirtualLog.log_num()) ::
{:ok, ExWal.LogWriter.t()} | {:error, reason :: any()}
def create(%Standalone{name: name}, log_num), do: Standalone.create(name, log_num)
@spec close(impl :: ExWal.Manager.t()) :: :ok | {:error, reason :: any()}
def close(%Standalone{name: name}), do: Standalone.stop(name)
end