Packages
riptide
0.5.0-beta9
0.5.2
0.5.1
0.5.0-beta9
0.5.0-beta8
0.5.0-beta7
0.5.0-beta6
0.5.0-beta5
0.5.0-beta4
0.5.0-beta3
0.5.0-beta2
0.5.0-beta11
0.5.0-beta10
0.5.0-beta
0.4.6
0.4.5
0.4.4
0.4.3
0.4.2
0.4.1
0.4.0
0.3.13
0.3.12
0.3.11
0.3.10
0.3.9
0.3.8
0.3.7
0.3.6
0.3.5
0.3.4
0.3.3
0.3.2
0.3.1
0.3.0
0.3.0-bd63a38
0.2.79
0.2.78
0.2.74
0.2.4
0.2.3
0.2.2
0.2.1
0.2.0
0.1.15
0.1.14
0.1.13
0.1.12
0.1.11
0.1.10
0.1.9
0.1.8
0.1.7
0.1.6
0.1.5
0.1.4
0.1.3
0.1.2
0.1.1
0.1.0
A data first framework for building realtime applications
Current section
Files
Jump to
Current section
Files
lib/riptide/store/tree_lmdb.ex
defmodule Riptide.Store.TreeLMDB do
@behaviour Riptide.Store
@delimiter "×"
@impl true
def init(opts) do
directory = opts_directory(opts)
{:ok, env} = Bridge.LMDB.open_env(directory)
:persistent_term.put({:riptide_lmdb, directory}, env)
:ok
end
@impl true
def mutation(merges, deletes, opts) do
tree = opts_tree(opts)
env = env(opts)
merges
|> Stream.map(fn {path, val} ->
branch = tree.for_path(path)
{columns, extra} = Enum.split(path, Enum.count(branch.columns))
{columns, extra, val}
end)
|> Enum.group_by(
fn {columns, _path, _val} ->
columns
end,
fn {_columns, path, val} ->
{path, val}
end
)
|> Stream.map(fn {columns, values} ->
existing =
env
|> iterate(columns, %{})
|> Enum.at(0)
|> case do
nil ->
nil
result ->
Jason.decode!(result)
end
next =
values
|> Enum.reduce(existing, fn
{[], val}, _collect ->
val
{path, val}, collect ->
case collect do
nil -> %{}
result when is_map(result) -> result
_ -> %{}
end
|> Dynamic.put(path, val)
end)
end)
end
def iterate(env, path, opts) do
combined = encode_path(path)
{min, max} = Riptide.Store.Prefix.range(combined, opts)
min = Enum.join(min, @delimiter)
max = Enum.join(max, @delimiter)
{:ok, tx} = Bridge.LMDB.txn_read_new(env)
exact =
tx
|> Bridge.LMDB.get(min)
|> case do
{:ok, value} -> [{min, value}]
_ -> []
end
:ok = Bridge.LMDB.txn_read_abort(tx)
Stream.concat(
exact,
Bridge.LMDB.stream(env, min <> @delimiter, max)
)
end
defp encode_path(path) do
Enum.join(path, @delimiter)
end
defp decode_path(input) do
String.split(input, @delimiter)
end
defp env(opts) do
directory = opts_directory(opts)
:persistent_term.get({:riptide_lmdb, directory})
end
defp opts_tree(opts), do: Keyword.get(opts, :tree)
defp opts_directory(opts), do: Keyword.get(opts, :directory)
end