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