%%% -*- erlang -*- %%% %%% This file is part of erlang-redi released under the BSD license. %%% %%% Copyright (c) 2018 Renaud Mariana %%% -module(redi). -behaviour(gen_server). %% API -export([ start_link/0, start_link/1, start_link/2, stop/1, child_spec/1, child_spec/2, set/3, set/4, set_bulk/3, set_bulk/4, get/2, delete/2, delete/3, size/1, get_bucket_name/1, add_bucket/2, add_bucket/3, gc_client/2, gc_client/3, all/1, dump/1, dump/2, get_maybe_update/3 ]). %% gen_server callbacks -export([ init/1, handle_call/3, handle_cast/2, handle_info/2, terminate/2, code_change/3 ]). -define(SERVER, ?MODULE). -define(ts_ms(), erlang:system_time(milli_seconds)). -define(ts_max, 16#7ffffffffffffff). %% 256 non gc ts values -define(ts_non_gc, (?ts_max - 256)). %% default configuration values -define(ENTRY_TTL_MS, timer:hours(1)). -define(GC_INTERVAL_MS, timer:seconds(30)). -define(BUCKET_NAME, redi_keys). -define(BUCKET_TYPE, set). -record(state, {next_gc_ms, entry_ttl_ms, bucket_name, gc_q, gc_client}). -spec start_link() -> {ok, Pid :: pid()} | {error, Error :: {already_started, pid()}} | {error, Error :: term()} | ignore. start_link() -> start_link(#{}). -spec start_link(atom() | map()) -> {ok, Pid :: pid()} | {error, Error :: {already_started, pid()}} | {error, Error :: term()} | ignore. start_link(Name) when is_atom(Name) -> start_link(Name, #{}); start_link(Opts) when is_map(Opts) -> start_link(?SERVER, Opts). stop(Name) -> gen_server:stop(Name). %% @doc creates an REDI cache %% Options are: %% - `bucket_name' name of the ets table (used by get) %% - `bucket_type' type of bucket, default to 'set' %% - `entry_ttl_ms' the time to live of REDI elements %% - `next_gc_ms' next interval time REDI will scan to remove oldest elements %% -spec start_link(atom(), map()) -> {ok, Pid :: pid()} | {error, Error :: {already_started, pid()}} | {error, Error :: term()} | ignore. start_link(Name, Opts) -> gen_server:start_link({local, Name}, ?MODULE, [Opts], []). child_spec(Opts) -> #{ id => ?MODULE, start => {?MODULE, start_link, Opts}, shutdown => 500 }. %% helper when using multiple instances of redi child_spec(BucketName, TTL_ms) when is_atom(BucketName) -> #{ id => BucketName, start => {?MODULE, start_link, [ BucketName, #{bucket_name => BucketName, entry_ttl_ms => TTL_ms} ]} }. -spec delete(Gen_server :: pid(), Key :: term()) -> ok | {error, undefined_bucket}. delete(Pid, Key) -> gen_server:call(Pid, {delete, Key}). -spec delete(Gen_server :: pid(), Key :: term(), atom()) -> ok. delete(Pid, Key, Bucket_name) when is_atom(Bucket_name) -> case ets:info(Bucket_name, size) of undefined -> {error, undefined_bucket}; _ -> gen_server:call(Pid, {delete, Key, Bucket_name}) end. -spec gc_client(Gen_server :: pid(), Client :: pid()) -> ok. gc_client(Redi_pid, Client_pid) when is_pid(Client_pid) -> gc_client(Redi_pid, Client_pid, #{returns => key}). -spec gc_client(Gen_server :: pid(), Client :: pid(), Opts :: maps:maps()) -> ok. gc_client(Redi_pid, Client_pid, Opts) when is_pid(Client_pid), is_map(Opts) -> gen_server:call(Redi_pid, {gc_client, Client_pid, maps:get(returns, Opts)}). %% @doc %% returns ets:lookup(Bucket_name, Key) %% returns [] if no entry is found %% returns [{key,data}, ..] if there is at least an entry %% -spec get(Bucket_name :: atom(), Key :: term()) -> list(). get(Bucket_name, Key) -> ets:lookup(Bucket_name, Key). -spec size(Bucket_name :: atom()) -> pos_integer(). size(Bucket_name) -> ets:info(Bucket_name, size). -spec set(Gen_server :: pid(), Key :: term(), Value :: term()) -> ok. set(Pid, Key, Value) -> gen_server:call(Pid, {set, Key, Value}). -spec set(Gen_server :: pid(), Key :: term(), Value :: term(), Bucket_name :: atom()) -> ok. set(Pid, Key, Value, Bucket_name) when is_atom(Bucket_name) -> case ets:info(Bucket_name, size) of undefined -> {error, undefined_bucket}; _ -> gen_server:call(Pid, {set, Key, Value, Bucket_name}) end. -spec set_bulk(Gen_server :: pid(), Key :: term(), Value :: term()) -> ok. set_bulk(Pid, Key, Value) -> gen_server:call(Pid, {set_bulk, Key, Value}). -spec set_bulk(Gen_server :: pid(), Key :: term(), Value :: term(), Bucket_name :: atom()) -> ok. set_bulk(Pid, Key, Value, Bucket_name) when is_atom(Bucket_name) -> case ets:info(Bucket_name, size) of undefined -> {error, undefined_bucket}; _ -> gen_server:call(Pid, {set_bulk, Key, Value, Bucket_name}) end. %% @doc get bucket name. Default is redi_keys %% the bucket name is required by get/2 -spec get_bucket_name(Gen_server :: pid()) -> atom(). get_bucket_name(Pid) -> gen_server:call(Pid, get_bucket_name). %% debug purpose, returns {list_of_kvalues, gc_entry_length} dump(Pid) -> gen_server:call(Pid, dump). dump(Pid, Bucket_name) -> gen_server:call(Pid, {dump, Bucket_name}). -spec add_bucket(Gen_server :: pid(), atom()) -> atom(). add_bucket(Pid, Bucket_name) -> add_bucket(Pid, Bucket_name, ?BUCKET_TYPE). -spec add_bucket(Gen_server :: pid(), atom(), atom()) -> atom(). add_bucket(Pid, Bucket_name, Bucket_type) -> gen_server:call(Pid, {add_bucket, Bucket_name, Bucket_type}). -spec all(Bucket_name :: atom()) -> list(). all(Bucket_name) -> ets:tab2list(Bucket_name). %% @doc get value from Key / Bucket_name %% if key is missing, call Fallback(Key) to update the cache %% Returns the value -spec get_maybe_update(Bucket_name :: atom(), Key :: term(), Fallback :: fun()) -> term(). get_maybe_update(Key, Bucket, Fallback) when is_atom(Bucket), is_function(Fallback, 1) -> case redi:get(Bucket, Key) of [] -> Val = Fallback(Key), redi:set(Bucket, Key, Val), Val; [{_, Val}] -> Val end. %%==================================================================== %% Internal functions %%==================================================================== %% %% @private -spec init(Args :: term()) -> {ok, State :: term()} | {ok, State :: term(), Timeout :: timeout()} | {ok, State :: term(), hibernate} | {stop, Reason :: term()} | ignore. init([Opts]) -> process_flag(trap_exit, true), Bucket_name = maps:get(bucket_name, Opts, ?BUCKET_NAME), Bucket_type = maps:get(bucket_type, Opts, ?BUCKET_TYPE), Next_gc_ms = maps:get(next_gc_ms, Opts, ?GC_INTERVAL_MS), create_bucket(Bucket_name, Bucket_type), erlang:send_after(Next_gc_ms, self(), refresh_gc), rand:seed(exrop), {ok, #state{ next_gc_ms = Next_gc_ms, bucket_name = Bucket_name, entry_ttl_ms = maps:get(entry_ttl_ms, Opts, ?ENTRY_TTL_MS), gc_q = queue:new() }}. %% @private -spec handle_call(Request :: term(), From :: {pid(), term()}, State :: term()) -> {reply, Reply :: term(), NewState :: term()} | {reply, Reply :: term(), NewState :: term(), Timeout :: timeout()} | {reply, Reply :: term(), NewState :: term(), hibernate} | {noreply, NewState :: term()} | {noreply, NewState :: term(), Timeout :: timeout()} | {noreply, NewState :: term(), hibernate} | {stop, Reason :: term(), Reply :: term(), NewState :: term()} | {stop, Reason :: term(), NewState :: term()}. handle_call({set, Key, Value}, _From, #state{gc_q = GC_q, bucket_name = Bucket_name} = State) -> NewGC_q = do_insert_gc(?ts_ms(), Key, GC_q), ets:insert(Bucket_name, {Key, Value}), {reply, ok, State#state{gc_q = NewGC_q}}; handle_call({set, Key, Value, Bucket_name}, _From, #state{gc_q = GC_q} = State) -> NewGC_q = do_insert_gc(?ts_ms(), {Key, Bucket_name}, GC_q), ets:insert(Bucket_name, {Key, Value}), {reply, ok, State#state{gc_q = NewGC_q}}; handle_call({set_bulk, Key, Value}, From, State) -> handle_call({set_bulk, Key, Value, State#state.bucket_name}, From, State); handle_call( {set_bulk, Key, Value, Bucket_name}, _From, #state{gc_q = GC_q, entry_ttl_ms = TTL} = State ) -> NewGC_q = do_insert_gc(?ts_ms() + rand:uniform(TTL), Key, GC_q), ets:insert(Bucket_name, {Key, Value}), {reply, ok, State#state{gc_q = NewGC_q}}; handle_call(dump, _From, #state{gc_q = GC_q, bucket_name = Bucket_name} = State) -> {reply, {ets:tab2list(Bucket_name), queue:to_list(GC_q)}, State}; handle_call({dump, Bucket_name}, _From, #state{gc_q = GC_q} = State) -> {reply, {ets:tab2list(Bucket_name), queue:to_list(GC_q)}, State}; handle_call(get_bucket_name, _From, #state{bucket_name = Bucket_name} = State) -> {reply, Bucket_name, State}; handle_call({add_bucket, Bucket_name, Bucket_type}, _From, State) -> create_bucket(Bucket_name, Bucket_type), {reply, Bucket_name, State}; handle_call({gc_client, Pid, TypeReturns}, _From, State) -> {reply, ok, State#state{gc_client = {Pid, TypeReturns}}}; handle_call({delete, Key}, _From, #state{gc_q = GC_q, bucket_name = Bucket_name} = State) -> ets:delete(Bucket_name, Key), NewGC_q = queue:from_list(lists:keydelete(Key, 2, queue:to_list(GC_q))), {reply, ok, State#state{gc_q = NewGC_q}}; handle_call({delete, Key, Bucket_name}, _From, #state{gc_q = GC_q} = State) -> ets:delete(Bucket_name, Key), NewGC_q = queue:from_list(lists:keydelete(Key, 2, queue:to_list(GC_q))), {reply, ok, State#state{gc_q = NewGC_q}}. %% @private -spec handle_cast(Request :: term(), State :: term()) -> {noreply, NewState :: term()} | {noreply, NewState :: term(), Timeout :: timeout()} | {noreply, NewState :: term(), hibernate} | {stop, Reason :: term(), NewState :: term()}. handle_cast(_Request, State) -> {noreply, State}. %% @private -spec handle_info(Info :: timeout() | term(), State :: term()) -> {noreply, NewState :: term()} | {noreply, NewState :: term(), Timeout :: timeout()} | {noreply, NewState :: term(), hibernate} | {stop, Reason :: normal | term(), NewState :: term()}. handle_info(refresh_gc, #state{entry_ttl_ms = TTL} = State) -> erlang:send_after(State#state.next_gc_ms, self(), refresh_gc), T_gc_ms = ?ts_ms() - TTL, State2 = clean_older0(T_gc_ms, [], State), {noreply, State2}. %% @private -spec terminate( Reason :: normal | shutdown | {shutdown, term()} | term(), State :: term() ) -> any(). terminate(_Reason, _State) -> ok. %% @private -spec code_change( OldVsn :: term() | {down, term()}, State :: term(), Extra :: term() ) -> {ok, NewState :: term()} | {error, Reason :: term()}. code_change(_OldVsn, State, _Extra) -> {ok, State}. do_insert_gc(Ts, Key, GC_q) when Ts < ?ts_non_gc -> queue:in({Ts, Key}, GC_q); do_insert_gc(_Ts, _Key, GC_q) -> GC_q. extract_key_bucket({Key, Bucket}, _State) -> {Key, Bucket}; extract_key_bucket(Key, State) -> {Key, State#state.bucket_name}. clean_older0(T_gc_ms, Acc, #state{gc_client = undefined} = State) -> case queue:peek(State#state.gc_q) of empty -> State; {value, _V = {Ts, {Key, Bucket_name}}} when Ts < T_gc_ms -> ets:delete(Bucket_name, Key), GC1_q = queue:drop(State#state.gc_q), clean_older0(T_gc_ms, Acc, State#state{gc_q = GC1_q}); {value, _V = {Ts, Key}} when Ts < T_gc_ms -> ets:delete(State#state.bucket_name, Key), GC1_q = queue:drop(State#state.gc_q), clean_older0(T_gc_ms, Acc, State#state{gc_q = GC1_q}); {value, _V} -> State end; clean_older0(T_gc_ms, Acc, #state{gc_client = {_, TypeReturns}} = State) -> case queue:peek(State#state.gc_q) of empty -> terminate_clean_older(Acc, State); {value, _V = {Ts, KeyM}} when Ts < T_gc_ms -> {Key, Bucket_name} = extract_key_bucket(KeyM, State), Gc_ed_keys = case ets:lookup(Bucket_name, Key) of [] -> Acc; [Key_value] -> ets:delete(Bucket_name, Key), if TypeReturns == key_value -> [Key_value | Acc]; true -> [Key | Acc] end; Key_values -> ets:delete(Bucket_name, Key), if TypeReturns == key_value -> lists:append(Key_values, Acc); true -> [Key | Acc] end end, GC1_q = queue:drop(State#state.gc_q), clean_older0(T_gc_ms, Gc_ed_keys, State#state{gc_q = GC1_q}); {value, _V} -> terminate_clean_older(Acc, State) end. terminate_clean_older(_Gc_ed_keys, #state{gc_client = {undefined, _}} = State) -> State; terminate_clean_older([], State) -> State; terminate_clean_older(Gc_ed_keys, #state{gc_client = {Pid, _}, bucket_name = Bucket} = State) -> erlang:send(Pid, {redi_gc, Bucket, Gc_ed_keys}), State. create_bucket(Bucket_name, Bucket_type) -> ets:new(Bucket_name, [Bucket_type, named_table]).