Packages
kvs
8.4.0
13.5.22-aleph
13.4.16
13.4.15
13.4.14
13.4.13
13.3.1
13.2.28
11.9.1
10.8.3
10.8.2
10.3.0
9.9.2
9.9.1
9.9.0
9.8.0
9.7.0
9.4.8
9.4.7
9.4.6
9.4.5
9.4.4
9.4.3
9.4.2
9.4.1
9.4.0
8.12.0
8.11.2
8.11.1
8.10.4
8.10.3
8.10.2
8.10.1
8.10.0
8.5.2
8.5.1
8.5.0
8.4.1
8.4.0
8.3.1
8.3.0
7.11.5
7.9.1
7.7.0
7.1.3
7.1.2
7.1.1
6.12.11
6.12.10
6.12.9
6.12.8
6.12.7
6.12.6
6.12.5
6.12.4
6.12.3
6.12.2
6.12.1
6.12.0
6.11.2
6.11.1
6.11.0
6.10.2
6.10.1
6.10.0
6.9.2
6.9.1
6.9.0
6.7.7
6.7.6
6.7.5
6.7.4
6.7.3
6.7.2
6.7.1
6.7.0
6.6.0
2.1.0
0.12.1
retired
KVS Key-Value Store Abstraction Layer
Current section
Files
Jump to
Current section
Files
src/stores/kvs_rocks.erl
-module(kvs_rocks).
-include("backend.hrl").
-include("kvs.hrl").
-include("metainfo.hrl").
-include_lib("stdlib/include/qlc.hrl").
-export(?BACKEND).
-export([ref/0,bt/1,key/2,key/1,fd/1, tb/1]).
-export([seek_it/1, move_it/3, take_it/4]).
e(X,Y) -> element(X,Y).
bt([]) -> [];
bt(X) -> binary_to_term(X).
tb([]) -> [];
tb(T) when is_list(T) -> list_to_binary(T);
tb(T) when is_atom(T) -> atom_to_binary(T, utf8);
tb(T) when is_binary(T) -> T;
tb(T) -> term_to_binary(T).
fmt([]) -> [];
fmt(K) -> tb(K).
% put
key(R) when is_tuple(R) andalso tuple_size(R) > 1 -> key(e(1,R), e(2,R));
key(R) -> key(R,[]).
key(Tab,R) when is_tuple(R) andalso tuple_size(R) > 1 -> key(Tab, e(2,R));
key(Tab,R) -> iolist_to_binary([lists:join(<<"/">>, lists:flatten([<<>>, tb(Tab), fmt(R)]))]).
fd(K) -> Key = tb(K),
End = byte_size(Key),
{S,_} = case binary:matches(Key,[<<"/">>],[]) of
[{0,1}] -> {End,End};
[{0,1},{1,1}] -> {End,End};
[{0,1},{1,1}|T] -> hd(lists:reverse(T));
[{0,1}|T] -> hd(lists:reverse(T));
_ -> {End,End}
end,
binary:part(Key,{0,S}).
o(<<>>,FK,_,_) -> {ok,FK,[],[]};
o(Key,FK,Dir,Fx) ->
S = size(FK),
Infotech = fun (F,K,H,V,Acc) when binary_part(K,{0,S}) == FK -> {F(H,Dir),H,[V|Acc]};
(_,K,H,V,Acc) -> close_it(H),
throw({ok,fd(K),bt(V),[bt(A1)||A1<-Acc]}) end,
Privat = fun(F,K,V,H) -> case F(H,prev) of
{ok,K1,V1} when binary_part(K,{0,S}) == FK -> {{ok,K1,V1},H,[V]};
{ok,K1,V1} -> Infotech(F,K1,H,V1,[]);
E -> E
end end,
It = fun(F,{ok,H}) -> {F(H,{seek,Key}),H};
(F,{{ok,K,V},H})
when Dir =:= prev -> Privat(F,K,V,H);
(F,{{ok,K,V},H}) -> Infotech(F,K,H,V,[]);
(F,{{ok,K,V},H,A}) -> Infotech(F,K,H,V,A);
(_,{{error,_},H,Acc}) -> {{ok,[],[]},H,Acc};
(F,{R,O}) -> F(R,O);
(F,H) -> F(H) end,
catch case lists:foldl(It, {ref(),[]}, Fx) of
{{ok,K,Bin},_,A} -> {ok,fd(K), bt(Bin),[bt(A1)||A1<-A]};
{{ok,K,Bin},_} -> {ok,fd(K), bt(Bin),[]};
{{error,_},_,Acc} -> {ok,fd(FK),bt(shd(Acc)),[bt(A1) ||A1<-Acc]}
end.
start() -> ok.
stop() -> ok.
destroy() -> rocksdb:destroy(application:get_env(kvs,rocks_name,"rocksdb"), []).
version() -> {version,"KVS ROCKSDB"}.
dir() -> [].
leave() -> case ref() of [] -> skip; X -> rocksdb:close(X), application:set_env(kvs,rocks_ref,[]), ok end.
join(_) -> application:start(rocksdb),
leave(), {ok, Ref} = rocksdb:open(application:get_env(kvs,rocks_name,"rocksdb"), [{create_if_missing, true}]),
initialize(),
application:set_env(kvs,rocks_ref,Ref).
initialize() -> [ kvs:initialize(kvs_rocks,Module) || Module <- kvs:modules() ].
ref() -> application:get_env(kvs,rocks_ref,[]).
index(_,_,_) -> [].
close_it(H) -> try rocksdb:iterator_close(H) catch error:badarg -> ok end.
seek_it(K) -> o(K,K,ok,[fun rocksdb:iterator/2,fun rocksdb:iterator_move/2]).
move_it(K,FK,Dir) -> o(K,FK,Dir,[fun rocksdb:iterator/2,fun rocksdb:iterator_move/2,fun rocksdb:iterator_move/2]).
take_it(Key,FK,Dir,N) when is_integer(N) andalso N >= 0 ->
o(Key,FK,Dir,[fun rocksdb:iterator/2,fun rocksdb:iterator_move/2] ++
lists:map(fun(_) -> fun rocksdb:iterator_move/2 end,lists:seq(1,N)));
take_it(Key,FK,Dir,_) -> take_it(Key,FK,Dir,0).
get(Tab, Key) ->
case rocksdb:get(ref(), key(Tab,Key), []) of
not_found -> {error,not_found};
{ok,Bin} -> {ok,bt(Bin)} end.
put(Records) when is_list(Records) -> lists:map(fun(Record) -> put(Record) end, Records);
put(Record) -> rocksdb:put(ref(), key(Record), term_to_binary(Record), [{sync,true}]).
delete(Feed, Id) -> rocksdb:delete(ref(), key(Feed,Id), []).
count(_) -> 0.
all(R) -> kvs_st:feed(R).
shd([]) -> [];
shd(X) -> hd(X).
seq(_,_) ->
case os:type() of
{win32,nt} -> {Mega,Sec,Micro} = erlang:timestamp(), integer_to_list((Mega*1000000+Sec)*1000000+Micro);
_ -> erlang:integer_to_list(element(2,hd(lists:reverse(erlang:system_info(os_monotonic_time_source)))))
end.
create_table(_,_) -> [].
add_table_index(_, _) -> ok.
dump() -> ok.