dets.erl

来自「OTP是开放电信平台的简称」· ERL 代码 · 共 2,178 行 · 第 1/5 页

ERL
2,178
字号
	    {value, {min_no_slots, I}} ->		case catch is_min_no_slots(I) of		    {'EXIT', _} -> badarg;		    MinNoSlots -> {ok, MinNoSlots}		end;            {value, {n_objects, default}} ->                {ok, default_option(Key)};            {value, {n_objects, NObjs}} when is_integer(NObjs),                                             NObjs >= 1 ->                {ok, NObjs};            {value, {traverse, select}} ->                {ok, select};            {value, {traverse, {select, MS}}} ->                {ok, {select, MS}};            {value, {traverse, first_next}} ->                {ok, first_next};	    {value, {Key, _}} ->		badarg;	    false ->		Default = default_option(Key),		{ok, Default}	end,    case V of	badarg ->	    {badarg, Key};	{ok, Value} ->	    NewOptions = lists:keydelete(Key, 1, Options),	    options(NewOptions, Keys, [Value | L])    end;options([], [], L) ->    lists:reverse(L);options(Options, _, _L) ->    {badarg,Options}.default_option(format) -> term;default_option(min_no_slots) -> default;default_option(traverse) -> select;default_option(n_objects) -> default.listify(L) when is_list(L) ->    L;listify(T) ->    [T].treq(Tab, R) ->    case catch dets_server:get_pid(Tab) of	Pid when is_pid(Pid) ->	    req(Pid, R);	_ ->	    badarg    end.req(Proc, R) ->    Ref = erlang:monitor(process, Proc),    Proc ! ?DETS_CALL(self(), R),    receive 	{'DOWN', Ref, process, Proc, _Info} ->            badarg;	{Proc, Reply} ->	    erlang:demonitor(Ref),	    receive 		{'DOWN', Ref, process, Proc, _Reason} ->                    Reply	    after 0 ->                    Reply	    end    end.%% Inlined.einval({error, {file_error, _, einval}}, A) ->    erlang:error(badarg, A);einval(Reply, _A) ->    Reply.%% Inlined.badarg(badarg, A) ->    erlang:error(badarg, A);badarg(Reply, _A) ->    Reply.%% Inlined.undefined(badarg) ->    undefined;undefined(Reply) ->    Reply.%% Inlined.badarg_exit(badarg, A) ->    erlang:error(badarg, A);badarg_exit({ok, Reply}, _A) ->    Reply;badarg_exit(Reply, _A) ->    exit(Reply).%%%-----------------------------------------------------------------%%% Server functions%%%-----------------------------------------------------------------init(Parent, Server) ->    process_flag(trap_exit, true),    open_file_loop(#head{parent = Parent, server = Server}).open_file_loop(Head) ->    open_file_loop(Head, 0).open_file_loop(Head, N) when element(1, Head#head.update_mode) =:= error ->    open_file_loop2(Head, N);open_file_loop(Head, N) ->    receive         %% When the table is fixed it can be assumed that at least one        %% traversal is in progress. To speed the traversal up three        %% things have been done:        %% - prioritize match_init, bchunk, next, and match_delete_init;        %% - do not peek the message queue for updates;        %% - wait 1 ms after each update.        %% next is normally followed by lookup, but since lookup is also        %% used when not traversing the table, it is not prioritized.        ?DETS_CALL(From, {match_init, _State} = Op) ->            do_apply_op(Op, From, Head, N);        ?DETS_CALL(From, {bchunk, _State} = Op) ->            do_apply_op(Op, From, Head, N);        ?DETS_CALL(From, {next, _Key} = Op) ->            do_apply_op(Op, From, Head, N);                    ?DETS_CALL(From, {match_delete_init, _MP, _Spec} = Op) ->            do_apply_op(Op, From, Head, N);        {'EXIT', Pid, Reason} when Pid =:= Head#head.parent ->            %% Parent orders shutdown.            _NewHead = do_stop(Head),            exit(Reason);        {'EXIT', Pid, Reason} when Pid =:= Head#head.server ->            %% The server is gone.            _NewHead = do_stop(Head),            exit(Reason);        {'EXIT', Pid, _Reason} ->            %% A process fixing the table exits.            H2 = remove_fix(Head, Pid, close),            open_file_loop(H2, N);        {system, From, Req} ->            sys:handle_system_msg(Req, From, Head#head.parent,                                   ?MODULE, [], Head)    after 0 ->            open_file_loop2(Head, N)    end.open_file_loop2(Head, N) ->    receive        ?DETS_CALL(From, Op) ->            do_apply_op(Op, From, Head, N);        {'EXIT', Pid, Reason} when Pid =:= Head#head.parent ->            %% Parent orders shutdown.            _NewHead = do_stop(Head),            exit(Reason);        {'EXIT', Pid, Reason} when Pid =:= Head#head.server ->            %% The server is gone.            _NewHead = do_stop(Head),            exit(Reason);        {'EXIT', Pid, _Reason} ->            %% A process fixing the table exits.            H2 = remove_fix(Head, Pid, close),            open_file_loop(H2, N);        {system, From, Req} ->            sys:handle_system_msg(Req, From, Head#head.parent,                                   ?MODULE, [], Head);        Message ->            error_logger:format("** dets: unexpected message"                                "(ignored): ~w~n", [Message]),            open_file_loop(Head, N)    end.do_apply_op(Op, From, Head, N) ->    try apply_op(Op, From, Head, N) of        ok ->             open_file_loop(Head, N);        {N2, H2} when is_record(H2, head), is_integer(N2) ->            open_file_loop(H2, N2);        H2 when is_record(H2, head) ->            open_file_loop(H2, N)    catch         exit:normal ->             exit(normal);        _:Bad ->             %% If stream_op/5 found more requests, this is not            %% the last operation.            Name = Head#head.name,            error_logger:format              ("** dets: Bug was found when accessing table ~w,~n"               "** dets: operation was ~p and reply was ~w.~n"               "** dets: Stacktrace: ~w~n",               [Name, Op, Bad, erlang:get_stacktrace()]),            if                From =/= self() ->                    From ! {self(), {error, {dets_bug, Name, Op, Bad}}};                true -> % auto_save | may_grow | {delayed_write, _}                    ok            end,            open_file_loop(Head, N)    end.apply_op(Op, From, Head, N) ->    case Op of	{add_user, Tab, OpenArgs}->            #open_args{file = Fname, type = Type, keypos = Keypos,                        ram_file = Ram, access = Access, 		       version = Version} = OpenArgs,            VersionOK = (Version =:= default) or                         (Head#head.version =:= Version),	    %% min_no_slots and max_no_slots are not tested	    Res = if		      Tab =:= Head#head.name,		      Head#head.keypos =:= Keypos,		      Head#head.type =:= Type,		      Head#head.ram_file =:= Ram,		      Head#head.access =:= Access,		      VersionOK,		      Fname =:= Head#head.filename ->			  ok;		      true ->			  err({error, incompatible_arguments})		  end,	    From ! {self(), Res},	    ok;	auto_save ->	    case Head#head.update_mode of		saved ->		    Head;		{error, _Reason} ->		    Head;		_Dirty when N =:= 0 -> % dirty or new_dirty		    %% The updates seems to have declined		    dets_utils:vformat("** dets: Auto save of ~p\n",                                        [Head#head.name]), 		    {NewHead, _Res} = perform_save(Head, true),		    erlang:garbage_collect(),		    {0, NewHead};		dirty -> 		    %% Reset counter and try later		    start_auto_save_timer(Head),		    {0, Head}	    end;	close  ->	    From ! {self(), fclose(Head)},	    _NewHead = unlink_fixing_procs(Head),	    ?PROFILE(ep:done()),	    exit(normal);	{close, Pid} -> 	    %% Used from dets_server when Pid has closed the table,	    %% but the table is still opened by some process.	    NewHead = remove_fix(Head, Pid, close),	    From ! {self(), status(NewHead)},	    NewHead;	{corrupt, Reason} ->	    {H2, Error} = dets_utils:corrupt_reason(Head, Reason),	    From ! {self(), Error},	    H2;	{delayed_write, WrTime} ->	    delayed_write(Head, WrTime);	info ->	    {H2, Res} = finfo(Head),	    From ! {self(), Res},	    H2;	{info, Tag} ->	    {H2, Res} = finfo(Head, Tag),	    From ! {self(), Res},	    H2;        {is_compatible_bchunk_format, Term} ->            Res = test_bchunk_format(Head, Term),            From ! {self(), Res},            ok;	{internal_open, Ref, Args} ->	    ?PROFILE(ep:do()),	    case do_open_file(Args, Head#head.parent, Head#head.server,Ref) of		{ok, H2} -> 		    From ! {self(), ok},		    H2;		Error -> 		    From ! {self(), Error},		    exit(normal)	    end;	may_grow when Head#head.update_mode =/= saved ->	    if		Head#head.update_mode =:= dirty ->		    %% Won't grow more if the table is full.		    {H2, _Res} = 			(Head#head.mod):may_grow(Head, 0, many_times),		    {N + 1, H2};		true -> 		    ok	    end;	{set_verbose, What} ->	    set_verbose(What), 	    From ! {self(), ok},	    ok;	{where, Object} ->	    {H2, Res} = where_is_object(Head, Object),	    From ! {self(), Res},	    H2;	_Message when element(1, Head#head.update_mode) =:= error ->	    From ! {self(), status(Head)},	    ok;	%% The following messages assume that the status of the table is OK.	{bchunk_init, Tab} ->	    {H2, Res} = do_bchunk_init(Head, Tab),	    From ! {self(), Res},	    H2;	{bchunk, State} ->	    {H2, Res} = do_bchunk(Head, State),	    From ! {self(), Res},	    H2;	delete_all_objects ->	    {H2, Res} = fdelete_all_objects(Head),	    From ! {self(), Res},	    erlang:garbage_collect(),	    {0, H2};	{delete_key, Keys} when Head#head.update_mode =:= dirty ->	    if		Head#head.version =:= 8 ->		    {H2, Res} = fdelete_key(Head, Keys),		    From ! {self(), Res},		    {N + 1, H2};		true ->		    stream_op(Op, From, [], Head, N)	    end;	{delete_object, Objs} when Head#head.update_mode =:= dirty ->	    case check_objects(Objs, Head#head.keypos) of		true when Head#head.version =:= 8 ->		    {H2, Res} = fdelete_object(Head, Objs),		    From ! {self(), Res},		    {N + 1, H2};		true ->		    stream_op(Op, From, [], Head, N);		false ->		    From ! {self(), badarg},		    ok	    end;	first ->	    {H2, Res} = ffirst(Head),	    From ! {self(), Res},	    H2;	{fixtable, Fixed} ->	    NewHead = do_fixtable(Head, From, Fixed),	    From ! {self(), ok},	    NewHead;        {initialize, InitFun, Format, MinNoSlots} ->            {H2, Res} = finit(Head, InitFun, Format, MinNoSlots),            From ! {self(), Res},	    erlang:garbage_collect(),            H2;	{insert, Objs} when Head#head.update_mode =:= dirty ->	    case check_objects(Objs, Head#head.keypos) of		true when Head#head.version =:= 8 ->		    {H2, Res} = finsert(Head, Objs),		    From ! {self(), Res},		    {N + 1, H2};		true ->		    stream_op(Op, From, [], Head, N);		false ->		    From ! {self(), badarg},		    ok	    end;	{insert_new, Objs} when Head#head.update_mode =:= dirty ->            {H2, Res} = finsert_new(Head, Objs),            From ! {self(), Res},            {N + 1, H2};	{lookup_keys, Keys}  when Head#head.version =:= 8 ->	    {H2, Res} = flookup_keys(Head, Keys),	    From ! {self(), Res},	    H2;	{lookup_keys, _Keys} ->	    stream_op(Op, From, [], Head, N);	{match_init, State} ->	    {H2, Res} = fmatch_init(Head, State),	    From ! {self(), Res},	    H2;	{match, MP, Spec, NObjs} ->	    {H2, Res} = fmatch(Head, MP, Spec, NObjs),	    From ! {self(), Res},	    H2;	{member, Key} when Head#head.version =:= 8 ->	    {H2, Res} = fmember(Head, Key),	    From ! {self(), Res},	    H2;	{member, _Key} = Op ->	    stream_op(Op, From, [], Head, N);	{next, Key} ->	    {H2, Res} = fnext(Head, Key),	    From ! {self(), Res},	    H2;	{match_delete, State} when Head#head.update_mode =:= dirty ->	    {H2, Res} = fmatch_delete(Head, State),	    From ! {self(), Res},	    {N + 1, H2};	{match_delete_init, MP, Spec} when Head#head.update_mode =:= dirty ->	    {H2, Res} = fmatch_delete_init(Head, MP, Spec),	    From ! {self(), Res},	    {N + 1, H2};	{safe_fixtable, Bool} ->	    NewHead = do_safe_fixtable(Head, From, Bool),	    From ! {self(), ok},	    NewHead;	{slot, Slot} ->	    {H2, Res} = fslot(Head, Slot),	    From ! {self(), Res},	    H2;	sync ->	    {NewHead, Res} = perform_save(Head, true),	    From ! {self(), Res},	    erlang:garbage_collect(),	    {0, NewHead};	{update_counter, Key, Incr} when Head#head.update_mode =:= dirty ->	    {NewHead, Res} = do_update_counter(Head, Key, Incr),	    From ! {self(), Res},	    {N + 1, NewHead};	WriteOp when Head#head.update_mode =:= new_dirty ->	    H2 = Head#head{update_mode = dirty},	    apply_op(WriteOp, From, H2, 0);	WriteOp when Head#head.access =:= read_write,		     Head#head.update_mode =:= saved ->	    case catch (Head#head.mod):mark_dirty(Head) of		ok ->		    start_auto_save_timer(Head),		    H2 = Head#head{update_mode = dirty},		    apply_op(WriteOp, From, H2, 0);		{NewHead, Error} when is_record(NewHead, head) ->		    From ! {self(), Error},		    NewHead	    end;	WriteOp when is_tuple(WriteOp), Head#head.access =:= read ->	    Reason = {access_mode, Head#head.filename},	    From ! {self(), err({error, Reason})},	    ok    end.start_auto_save_timer(Head) when Head#head.auto_save =:= infinity ->    ok;start_auto_save_timer(Head) ->    Millis = Head#head.auto_save,

⌨️ 快捷键说明

复制代码Ctrl + C
搜索代码Ctrl + F
全屏模式F11
增大字号Ctrl + =
减小字号Ctrl + -
显示快捷键?