Packages

Simple memory database for Erlang application using MVCC to store the data.

Current section

Files

Jump to
memstore src memstore.erl
Raw

src/memstore.erl

%%%-------------------------------------------------------------------
%%% @author benoitc
%%% @copyright (C) 2017, <COMPANY>
%%% @doc
%%%
%%% @end
%%% Created : 03. Nov 2017 12:11
%%%-------------------------------------------------------------------
-module(memstore).
-author("benoitc").
%% API
-export([
start_link/2,
open/2,
close/1,
put/3,
delete/2,
write_batch/2,
get/2, get/3,
new_snapshot/1,
release_snapshot/1,
get_snapshot_sequence/1,
iterator/2,
iterator_close/1,
iterator_move/2,
get_last_sequence/1
]).
-export([
init/1
]).
%% default write buffer size: 64 MB
-define(WRITER_BUFFER_SIZE, 64 bsl 20).
%% default compaction interval: 1s
-define(COMPACTION_INTERVAL, 1000).
-type db() :: atom().
-type key() :: term().
-type value() :: term().
-type snapshot() :: term().
-type iterator() :: pid().
-type batch() :: [
{put, Key :: key(), Value :: value()} | {delete, Key :: key()}
].
-type db_opts() :: [
{write_buffer_size, non_neg_integer()} |
{compaction_interval, non_neg_integer()}
].
-type read_options() :: [{snapshot, snapshot()}].
-type iterator_action() ::
{seek, Key :: term()} |
{seek_for_prev, term()} |
first | last | next | prev.
-type iterator_result() :: {ok, key(), value()} | false.
-export_types([
db/0,
key/0,
value/0,
batch/0,
snapshot/0,
db_opts/0,
read_options/0,
iterator/0,
iterator_action/0,
iterator_result/0
]).
%% @doc open a database
-spec open(Db :: db(), Options :: db_opts()) -> ok | {error, any()}.
open(Name, Options) ->
case supervisor:start_child(memstore_sup, [Name, Options]) of
{ok, _Pid} -> ok;
{error, {already_started, _Pid}} -> ok;
Error -> Error
end.
%% %% @doc close a database
-spec close(Db :: db()) -> ok.
close(Name) ->
supervisor:terminate_child(memstore_sup, whereis(Name)).
%% @doc update a key Key with value Value
-spec put(db(), key(), value()) -> ok | {error, any()}.
put(Name, Key, Value) ->
write_batch(Name, [{put, Key, Value}]).
%% @doc delete a key
-spec delete(db(), key()) -> ok |{error, any()}.
delete(Name, Key) ->
write_batch(Name, [{delete, Key}]).
%% @doc atomically do some update/delete operations
-spec write_batch(db(), batch()) -> ok | {error, term()}.
write_batch(Name, Batch) when is_list(Batch) ->
Writer = ets:lookup_element(Name, writer, 2),
memstore_writer:write_batch(Writer, Batch).
%% @doc fetch a key
-spec get(db(), key()) -> {ok, value()} | not_found | {error, term()}.
get(Name, Key) -> get(Name, Key, []).
-spec get(db(), key(), read_options()) ->{ok, value()} | not_found | {error, term()}.
get(Name, Key, Options) ->
Tab = ets:lookup_element(Name, tab, 2),
case proplists:get_value(snapshot, Options) of
#{ seq := Seq } ->
memstore_iterator:lookup(Tab, Key, Seq);
_ ->
memstore_iterator:lookup_last(Tab, Key)
end.
%% @doc get a new snapshot of the database. Snapshots provide consistent
%% read-only views over the entire state of the key-value store.
-spec new_snapshot(db()) -> snapshot().
new_snapshot(Name) ->
Mgr = ets:lookup_element(Name, version_mgr, 2),
gen_server:call(Mgr, new_snapshot).
%% @doc release a snapshot
-spec release_snapshot(snapshot()) -> ok.
release_snapshot(undefined) -> ok;
release_snapshot(#{ name := Name, snapshot_id := _Id} = Snapshot) ->
Mgr = ets:lookup_element(Name, version_mgr, 2),
gen_server:call(Mgr, {release_snapshot, Snapshot});
release_snapshot(_) ->
erlang:error(badarg).
%% @doc get the database version linked to this snapshots
-spec get_snapshot_sequence(snapshot()) -> non_neg_integer().
get_snapshot_sequence(#{ snapshot_id := _Id, seq := Seq}) -> Seq;
get_snapshot_sequence(_) -> erlang:error(badarg).
%% @doc create an iterator
-spec iterator(db(), read_options()) -> {ok, iterator()} |{error, any()}.
iterator(Name, Options) ->
Sup = ets:lookup_element(Name, iterator_sup, 2),
supervisor:start_child(Sup, [Name, Options]).
%% @doc close an iteratror
-spec iterator_close(iterator()) -> ok.
iterator_close(Itr) ->
memstore_iterator:iterator_close(Itr).
%% @doc move in the iterator
-spec iterator_move(iterator(), iterator_action()) -> iterator_result().
iterator_move(Itr, Action) ->
memstore_iterator:iterator_move(Itr, Action).
%% @doc return the last sequence
-spec get_last_sequence(db()) -> non_neg_integer().
get_last_sequence(Name) ->
Mgr = ets:lookup_element(Name, version_mgr, 2),
gen_server:call(Mgr, get_last_seq).
start_link(Name, Options) ->
supervisor:start_link({local, Name}, ?MODULE, [Name, Options]).
init([Name, Options]) ->
VersionMgrName = version_mgr_name(Name),
WriterName = writer_name(Name),
DataTab = table_name(Name),
IteratorSup = iterator_sup(Name),
WriteBufferSize = proplists:get_value(write_buffer_size, Options,?WRITER_BUFFER_SIZE),
Interval = proplists:get_value(compaction_interval, Options, ?COMPACTION_INTERVAL),
%% init metadata table
_ = ets:new(Name, [ordered_set, named_table, public, {read_concurrency, true}]),
ets:insert(
Name,
[
{sup, self()},
{tab, DataTab},
{writer, WriterName},
{version_mgr, VersionMgrName},
{iterator_sup, IteratorSup},
{last_seq, 0}
]
),
Specs = [
%% versions manager
#{
id => versions_mgr,
start => {memstore_versions, start_link, [VersionMgrName, Name]},
restart => transient,
shutdown => 2000,
type => worker,
modules => [memstore_versions]
},
%% writer
#{
id => writer,
start => {memstore_writer, start_link, [WriterName, Name, DataTab, WriteBufferSize]},
restart => transient,
shutdown => 2000,
type => worker,
modules => [memstore_writer]
},
%% iterator supervisor
#{
id => iterator_sup,
start => {memstore_iterator_sup, start_link, [IteratorSup]},
restart => transient,
shutdown => 2000,
type => worker,
modules => [memstore_writer]
},
%% background compaction
#{
id => gc,
start => {memstore_gc, start_link, [Name, DataTab, Interval]},
restart => transient,
shutdown => 2000,
type => worker,
modules => [memstore_gc]
}
],
%% finally start the supervisor
SupFlags = #{ strategy => one_for_all, intensity => 4, period => 3600},
{ok, {SupFlags, Specs}}.
version_mgr_name(Name) ->
list_to_atom(atom_to_list(Name) ++ "_versions_mgr").
writer_name(Name) ->
list_to_atom(atom_to_list(Name) ++ "_writer").
table_name(Name) ->
list_to_atom(atom_to_list(Name) ++ "_data").
iterator_sup(Name) ->
list_to_atom(atom_to_list(Name) ++ "_iterator_sup").