dets.erl

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

ERL
2,178
字号
    erlang:send_after(Millis, self(), ?DETS_CALL(self(), auto_save)).%% Version 9: Peek the message queue and try to evaluate several%% lookup requests in parallel. Evalute delete_object, delete and%% insert as well.stream_op(Op, Pid, Pids, Head, N) ->    stream_op(Head, Pids, [], N, Pid, Op, Head#head.fixed).stream_loop(Head, Pids, C, N, false = Fxd) ->    receive	?DETS_CALL(From, Message) ->	    stream_op(Head, Pids, C, N, From, Message, Fxd)    after 0 ->	    stream_end(Head, Pids, C, N, no_more)    end;stream_loop(Head, Pids, C, N, _Fxd) ->    stream_end(Head, Pids, C, N, no_more).stream_op(Head, Pids, C, N, Pid, {lookup_keys,Keys}, Fxd) ->    NC = [{{lookup,Pid},Keys} | C],    stream_loop(Head, Pids, NC, N, Fxd);stream_op(Head, Pids, C, N, Pid, {insert, _Objects} = Op, Fxd) ->    NC = [Op | C],    stream_loop(Head, [Pid | Pids], NC, N, Fxd);stream_op(Head, Pids, C, N, Pid, {insert_new, _Objects} = Op, Fxd) ->    NC = [Op | C],    stream_loop(Head, [Pid | Pids], NC, N, Fxd);stream_op(Head, Pids, C, N, Pid, {delete_key, _Keys} = Op, Fxd) ->    NC = [Op | C],    stream_loop(Head, [Pid | Pids], NC, N, Fxd);stream_op(Head, Pids, C, N, Pid, {delete_object, _Objects} = Op, Fxd) ->    NC = [Op | C],    stream_loop(Head, [Pid | Pids], NC, N, Fxd);stream_op(Head, Pids, C, N, Pid, {member, Key}, Fxd) ->    NC = [{{lookup,[Pid]},[Key]} | C],    stream_loop(Head, Pids, NC, N, Fxd);stream_op(Head, Pids, C, N, Pid, Op, _Fxd) ->    stream_end(Head, Pids, C, N, {Pid,Op}).stream_end(Head, Pids0, C, N, Next) ->    case catch update_cache(Head, lists:reverse(C)) of	{Head1, [], PwriteList} ->	    stream_end1(Pids0, Next, N, C, Head1, PwriteList);	{Head1, Found, PwriteList} ->	    %% Possibly an optimization: reply to lookup requests	    %% first, then write stuff. This makes it possible for	    %% clients to continue while the disk is accessed.	    %% (Replies to lookup requests are sent earlier than	    %%  replies to delete and insert requests even if the	    %%  latter requests were made before the lookup requests,	    %%  which can be confusing.)	    lookup_replies(Found),	    stream_end1(Pids0, Next, N, C, Head1, PwriteList);	Head1 when is_record(Head1, head) ->	    stream_end2(Pids0, Pids0, Next, N, C, Head1, ok);	    	{Head1, Error} when is_record(Head1, head) ->	    %% Dig out the processes that did lookup or member.	    Fun = fun({{lookup,[Pid]},_Keys}, L) -> [Pid | L];		     ({{lookup,Pid},_Keys}, L) -> [Pid | L];		     (_, L) -> L		  end,	    LPs0 = lists:foldl(Fun, [], C),	    LPs = lists:usort(lists:flatten(LPs0)),	    stream_end2(Pids0 ++ LPs, Pids0, Next, N, C, Head1, Error);        DetsError ->            DetsError    end.stream_end1(Pids, Next, N, C, Head, []) ->    stream_end2(Pids, Pids, Next, N, C, Head, ok);stream_end1(Pids, Next, N, C, Head, PwriteList) ->    {Head1, PR} = (catch dets_utils:pwrite(Head, PwriteList)),    stream_end2(Pids, Pids, Next, N, C, Head1, PR).stream_end2([Pid | Pids], Ps, Next, N, C, Head, Reply) ->    Pid ! {self(), Reply},    stream_end2(Pids, Ps, Next, N+1, C, Head, Reply);stream_end2([], Ps, no_more, N, C, Head, _Reply) ->    penalty(Head, Ps, C),    {N, Head};stream_end2([], _Ps, {From, Op}, N, _C, Head, _Reply) ->    apply_op(Op, From, Head, N).penalty(H, _Ps, _C) when H#head.fixed =:= false ->    ok;penalty(_H, _Ps, [{{lookup,_Pids},_Keys}]) ->    ok;penalty(#head{fixed = {_,[{Pid,_}]}}, [Pid], _C) ->    ok;penalty(_H, _Ps, _C) ->    timer:sleep(1).lookup_replies([{P,O}]) ->    lookup_reply(P, O);lookup_replies(Q) ->    [{P,O} | L] = dets_utils:family(Q),    lookup_replies(P, lists:append(O), L).lookup_replies(P, O, []) ->    lookup_reply(P, O);lookup_replies(P, O, [{P2,O2} | L]) ->    lookup_reply(P, O),    lookup_replies(P2, lists:append(O2), L).%% If a list of Pid then op was {member, Key}. Inlined.lookup_reply([P], O) ->    P ! {self(), O =/= []};lookup_reply(P, O) ->    P ! {self(), O}.%%-----------------------------------------------------------------%% Callback functions for system messages handling.%%-----------------------------------------------------------------system_continue(_Parent, _, Head) ->    open_file_loop(Head).system_terminate(Reason, _Parent, _, Head) ->    _NewHead = do_stop(Head),    exit(Reason).%%-----------------------------------------------------------------%% Code for upgrade.%%-----------------------------------------------------------------system_code_change(State, _Module, _OldVsn, _Extra) ->    {ok, State}.%%%----------------------------------------------------------------------%%% Internal functions%%%----------------------------------------------------------------------constants(FH, FileName) ->    Version = FH#fileheader.version,    if         Version =< 8 ->            dets_v8:constants();        Version =:= 9 ->            dets_v9:constants();        true ->            throw({error, {not_a_dets_file, FileName}})    end.%% -> {ok, Fd, fileheader()} | throw(Error)read_file_header(FileName, Access, RamFile) ->    BF = if	     RamFile ->		 case file:read_file(FileName) of		     {ok, B} -> B;		     Err -> dets_utils:file_error(FileName, Err)		 end;	     true ->		 FileName	 end,    {ok, Fd} = dets_utils:open(BF, open_args(Access, RamFile)),    {ok, <<Version:32>>} =         dets_utils:pread_close(Fd, FileName, ?FILE_FORMAT_VERSION_POS, 4),    if         Version =< 8 ->            dets_v8:read_file_header(Fd, FileName);        Version =:= 9 ->            dets_v9:read_file_header(Fd, FileName);        true ->            throw({error, {not_a_dets_file, FileName}})    end.fclose(Head) ->    {Head1, Res} = perform_save(Head, false),    case Head1#head.ram_file of	true -> 	    ignore;	false -> 	    file:close(Head1#head.fptr)    end,    Res.%% -> {NewHead, Res}perform_save(Head, DoSync) when Head#head.update_mode =:= dirty;				Head#head.update_mode =:= new_dirty ->    case catch begin                   {Head1, []} = write_cache(Head),                   {Head2, ok} = (Head1#head.mod):do_perform_save(Head1),                   ok = ensure_written(Head2, DoSync),                   {Head2#head{update_mode = saved}, ok}               end of        {NewHead, _} = Reply when is_record(NewHead, head) ->            Reply    end;perform_save(Head, _DoSync) ->    {Head, status(Head)}.ensure_written(Head, DoSync) when Head#head.ram_file ->    {ok, EOF} = dets_utils:position(Head, eof),    {ok, Bin} = dets_utils:pread(Head, 0, EOF, 0),    if	DoSync ->	    dets_utils:write_file(Head, Bin);	not DoSync ->            case file:write_file(Head#head.filename, Bin) of                ok ->                     ok;                Error ->                     dets_utils:corrupt_file(Head, Error)            end    end;ensure_written(Head, true) when not Head#head.ram_file ->    dets_utils:sync(Head);ensure_written(Head, false) when not Head#head.ram_file ->    ok.%% -> {NewHead, {cont(), [binary()]}} | {NewHead, Error}do_bchunk_init(Head, Tab) ->    case catch write_cache(Head) of	{H2, []} ->	    case (H2#head.mod):table_parameters(H2) of		undefined ->		    {H2, {error, old_version}};		Parms ->                    L = dets_utils:all_allocated(H2),                    C0 = #dets_cont{no_objs = default, bin = <<>>, alloc = L},		    BinParms = term_to_binary(Parms),		    {H2, {C0#dets_cont{tab = Tab, what = bchunk}, [BinParms]}}	    end;	{NewHead, _} = HeadError when is_record(NewHead, head) ->	    HeadError    end.%% -> {NewHead, {cont(), [binary()]}} | {NewHead, Error}do_bchunk(Head, State) ->    case dets_v9:read_bchunks(Head, State#dets_cont.alloc) of	{error, Reason} ->	    dets_utils:corrupt_reason(Head, Reason);	{finished, Bins} ->	    {Head, {State#dets_cont{bin = eof}, Bins}};	{Bins, NewL} ->	    {Head, {State#dets_cont{alloc = NewL}, Bins}}    end.%% -> {NewHead, Result}fdelete_all_objects(Head) when Head#head.fixed =:= false ->    case catch do_delete_all_objects(Head) of	{ok, NewHead} ->	    start_auto_save_timer(NewHead),	    {NewHead, ok};	{error, Reason} ->	    dets_utils:corrupt_reason(Head, Reason)    end;fdelete_all_objects(Head) ->    {Head, fixed}.do_delete_all_objects(Head) ->    #head{fptr = Fd, name = Tab, filename = Fname, type = Type, keypos = Kp, 	  ram_file = Ram, auto_save = Auto, min_no_slots = MinSlots,	  max_no_slots = MaxSlots, cache = Cache} = Head,     CacheSz = dets_utils:cache_size(Cache),    ok = dets_utils:truncate(Fd, Fname, bof),    (Head#head.mod):initiate_file(Fd, Tab, Fname, Type, Kp, MinSlots, MaxSlots,				  Ram, CacheSz, Auto, true).%% -> {NewHead, Reply}, Reply = ok | Error.fdelete_key(Head, Keys) ->    do_delete(Head, Keys, delete_key).%% -> {NewHead, Reply}, Reply = ok | badarg | Error.fdelete_object(Head, Objects) ->    do_delete(Head, Objects, delete_object).ffirst(H) ->    Ref = make_ref(),    case catch {Ref, ffirst1(H)} of	{Ref, {NH, R}} -> 	    {NH, {ok, R}};	{NH, R} when is_record(NH, head) -> 	    {NH, {error, R}}    end.ffirst1(H) ->    {NH, []} = write_cache(H),    ffirst(NH, 0).ffirst(H, Slot) ->    case (H#head.mod):slot_objs(H, Slot) of	'$end_of_table' -> {H, '$end_of_table'};	[] -> ffirst(H, Slot+1);	[X|_] -> {H, element(H#head.keypos, X)}    end.%% -> {NewHead, Reply}, Reply = ok | badarg | Error.finsert(Head, Objects) ->    case catch update_cache(Head, Objects, insert) of	{NewHead, []} ->	    {NewHead, ok};	{NewHead, _} = HeadError when is_record(NewHead, head) ->	    HeadError    end.%% -> {NewHead, Reply}, Reply = ok | badarg | Error.finsert_new(Head, Objects) ->    KeyPos = Head#head.keypos,    case catch lists:map(fun(Obj) -> element(KeyPos, Obj) end, Objects) of        Keys when is_list(Keys) ->            case catch update_cache(Head, Keys, {lookup, nopid}) of                {Head1, PidObjs} when is_list(PidObjs) ->                    case lists:all(fun({_P,OL}) -> OL =:= [] end, PidObjs) of                        true ->                            case catch update_cache(Head1, Objects, insert) of                                {NewHead, []} ->                                    {NewHead, true};                                {NewHead, Error} when is_record(NewHead, head) ->                                    {NewHead, Error}                            end;                        false=Reply ->                            {Head1, Reply}                    end;                {NewHead, _} = HeadError when is_record(NewHead, head) ->                    HeadError            end;        _ ->            {Head, badarg}    end.do_safe_fixtable(Head, Pid, true) ->    case Head#head.fixed of 	false -> 	    link(Pid),	    Fixed = {erlang:now(), [{Pid, 1}]},	    Ftab = dets_utils:get_freelists(Head),	    Head#head{fixed = Fixed, freelists = {Ftab, Ftab}};	{TimeStamp, Counters} ->	    case lists:keysearch(Pid, 1, Counters) of		{value, {Pid, Counter}} -> % when Counter > 1		    NewCounters = lists:keyreplace(Pid, 1, Counters, 						   {Pid, Counter+1}),		    Head#head{fixed = {TimeStamp, NewCounters}};		false ->		    link(Pid),		    Fixed = {TimeStamp, [{Pid, 1} | Counters]},		    Head#head{fixed = Fixed}	    end    end;do_safe_fixtable(Head, Pid, false) ->    remove_fix(Head, Pid, false).remove_fix(Head, Pid, How) ->    case Head#head.fixed of 	false -> 	    Head;	{TimeStamp, Counters} ->	    case lists:keysearch(Pid, 1, Counters) of		%% How =:= close when Pid closes the table.		{value, {Pid, Counter}} when Counter =:= 1; How =:= close ->		    unlink(Pid),		    case lists:keydelete(Pid, 1, Counters) of			[] -> 			    check_growth(Head),			    erlang:garbage_collect(),			    Head#head{fixed = false, 				  freelists = dets_utils:get_freelists(Head)};			NewCounters ->			    Head#head{fixed = {TimeStamp, NewCounters}}		    end;		{value, {Pid, Counter}} ->		    NewCounters = lists:keyreplace(Pid, 1, Counters, 						   {Pid, Counter-1}),		    Head#head{fixed = {TimeStamp, NewCounters}};		false ->		    Head	    end    end.do_stop(Head) ->    unlink_fixing_procs(Head),    fclose(Head).unlink_fixing_procs(Head) ->    case Head#head.fixed of	false ->	    Head;	{_, Counters} ->	    lists:map(fun({Pid, _Counter}) -> unlink(Pid) end, Counters),	    Head#head{fixed = false, 		      freelists = dets_utils:get_freelists(Head)}    end.%% -> NewHeaddo_fixtable(Head, _Pid, false) when Head#head.fixed =/= false ->    check_growth(Head),    unlink_fixing_procs(Head);do_fixtable(Head, Pid, true) ->    do_safe_fixtable(Head, Pid, true);do_fixtable(Head, _Pid, _Value) ->    Head.check_growth(#head{access = read}) ->    ok;check_growth(Head) ->    NoThings = no_things(Head),    if	NoThings > Head#head.next ->	    erlang:send_after(200, self(), 			      ?DETS_CALL(self(), may_grow)); % Catch up.	true ->	    ok    end.finfo(H) ->     case catch write_cache(H) of	{H2, []} ->	    Info = (catch [{type, H2#head.type}, 			   {keypos, H2#head.keypos}, 			   {size, H2#head.no_objects},			   {file_size, 			    file_size(H2#head.fptr, H2#head.filename)},			   {filename, H2#head.filename}]),	    {H2, Info};	{H2, _} = HeadError when is_record(H2, head) ->	    HeadError    end.finfo(H, access) -> {H, H#head.access};finfo(H, auto_save) -> {H, H#head.auto_save};finfo(H, bchunk_format) ->     case catch write_cache(H) of        {H2, []} ->            case (H2#head.mod):table_parameters(H2) of                undefined = Undef ->                    {H2, Undef};                Parms ->                    {H2, term_to_binary(Parms)}            end;        {H2, _} = HeadError when is_record(H2, head) ->            HeadError    end;finfo(H, delayed_write) -> % undocumented    {H, dets_utils:cache_size(H#head.cache)};finfo(H, filename) -> {H, H#head.filename};finfo(H, file_size) ->

⌨️ 快捷键说明

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