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 + -
显示快捷键?