cosnotification_eventdb.erl
来自「OTP是开放电信平台的简称」· ERL 代码 · 共 1,350 行 · 第 1/4 页
ERL
1,350 行
%% Returns : A modified return from now().%%------------------------------------------------------------extract_deadline(_, _, _, _, false) -> false;extract_deadline(Event, DefaultT, StopTSupported, TRef, MappingVal) -> extract_deadline(Event, DefaultT, StopTSupported, TRef, MappingVal, now()).extract_deadline(_, _, _, _, false, _) -> false;extract_deadline(#'CosNotification_StructuredEvent' {header = #'CosNotification_EventHeader' {variable_header = VH}}, DefaultT, StopTSupported, TRef, undefined, Now) -> DL = case extract_value(VH, ?not_Timeout, undefined) of undefined when StopTSupported == true, TRef =/= undefined -> case extract_value(VH, ?not_StopTime, undefined) of undefined -> DefaultT; DefinedTime -> DefinedTime end; undefined -> DefaultT; DefinedTime -> DefinedTime end, convert_time(DL, TRef, Now);%% Maybe a unstructured event.extract_deadline(_, Time, _, TRef, undefined, Now) -> convert_time(Time, TRef, Now);extract_deadline(_, _, _, TRef, DOverride, Now) -> %% Must have an associated MappingFilter defining a Deadline. convert_time(DOverride, TRef, Now).convert_time(0, _, _) -> false;convert_time(UTC, TRef, {M,S,U}) when record(UTC, 'TimeBase_UtcT') -> case catch get_time_diff(UTC, TRef) of {'EXCEPTION', _} -> false; {'EXIT', _} -> false; DL -> MicroSecs = round(DL/10), Secs = round(MicroSecs/1000000), MegaSecs = round(Secs/1000000), {-M-MegaSecs, -S-Secs+MegaSecs, -U-MicroSecs+Secs} end;convert_time(DL, _, {M,S,U}) when integer(DL) -> MicroSecs = round(DL/10), Secs = round(MicroSecs/1000000), MegaSecs = round(Secs/1000000), {-M-MegaSecs, -S-Secs+MegaSecs, -U-MicroSecs+Secs};convert_time(_, _, _) -> false.get_time_diff(UTC, TRef) -> UTO = 'CosTime_TimeService':universal_time(TRef), UTO2 = 'CosTime_TimeService':uto_from_utc(TRef, UTC), TIO = 'CosTime_UTO':time_to_interval(UTO, UTO2), #'TimeBase_IntervalT'{lower_bound=LB, upper_bound = UB} = 'CosTime_TIO':'_get_time_interval'(TIO), UB-LB.check_deadline(DL) when tuple(DL) -> {M,S,U} = now(), DL >= {-M,-S,-U};check_deadline(_DL) -> %% This case will cover if no timeout is set. false.check_start_time(ST) when tuple(ST) -> {M,S,U} = now(), ST >= {-M,-S,-U};check_start_time(_ST) -> %% This case will cover if no earliest delivery time is set. true.%%------------------------------------------------------------%% function : extract_value%% Arguments: A Property Sequence%% ID - wanted property string()%% Other - default-value.%% Returns : Value associated with given ID or default value.%%------------------------------------------------------------extract_value([], _, Other) -> Other;extract_value([#'CosNotification_Property'{name=ID, value=V}|_], ID, _) -> any:get_value(V);extract_value([_H|T], ID, Other) -> extract_value(T, ID, Other).%%------------------------------------------------------------%% function : get_event%% Arguments: %% Returns : %%------------------------------------------------------------get_event(DBRef) -> get_event(DBRef, true).get_event(DBRef, Delete) -> case get_events(DBRef, 1, Delete) of {[], false} -> {[], false}; {[], false, Keys} -> {[], false, Keys}; {[Event], Bool} -> {Event, Bool}; {[Event], Bool, Keys} -> {Event, Bool, Keys} end.%%------------------------------------------------------------%% function : get_events%% Arguments: %% Returns : A list of events (possibly empty) and a boolean%% indicating if event found.%% Comments : Try to extract Max events from the database.%%------------------------------------------------------------get_events(#dbRef{orderRef = ORef, discardRef = DRef}, Max) -> event_loop(ets:last(ORef), ORef, DRef, Max, [], [], true).get_events(#dbRef{orderRef = ORef, discardRef = DRef}, Max, Delete) -> event_loop(ets:last(ORef), ORef, DRef, Max, [], [], Delete).event_loop('$end_of_table', _, _, _, [], _, true) -> {[], false};event_loop('$end_of_table', _, _, _, [], [], _) -> {[], false, []};event_loop('$end_of_table', _ORef, _, _, Accum, _Keys, true) -> {lists:reverse(Accum), true};event_loop('$end_of_table', _ORef, _, _, Accum, Keys, _) -> {lists:reverse(Accum), true, Keys};event_loop(_, _ORef, _, 0, [], _Keys, true) -> %% Only possible if some tries to pull a sequence of 0 events. %% Should we really test for this case? {[], false};event_loop(_, _ORef, _, 0, [], Keys, _) -> {[], false, Keys};event_loop(_, _ORef, _, 0, Accum, _Keys, true) -> {lists:reverse(Accum), true};event_loop(_, _ORef, _, 0, Accum, Keys, _) -> {lists:reverse(Accum), true, Keys};event_loop(Key, ORef, undefined, Left, Accum, Keys, Delete) -> [{_,DL,ST,_PO,Event}]=ets:lookup(ORef, Key), case check_deadline(DL) of true -> ets:delete(ORef, Key), event_loop(ets:prev(ORef, Key), ORef, undefined, Left, Accum, Keys, Delete); false -> case check_start_time(ST) of true when Delete == true -> ets:delete(ORef, Key), event_loop(ets:prev(ORef, Key), ORef, undefined, Left-1, [Event|Accum], Keys, Delete); true -> event_loop(ets:prev(ORef, Key), ORef, undefined, Left-1, [Event|Accum], [{ORef, Key}|Keys], Delete); false -> event_loop(ets:prev(ORef, Key), ORef, undefined, Left, Accum, Keys, Delete) end end;event_loop({Key1, Key2}, ORef, DRef, Left, Accum, Keys, Delete) -> [{_,DL,ST,_PO,Event}]=ets:lookup(ORef, {Key1, Key2}), case check_deadline(DL) of true -> ets:delete(ORef, {Key1, Key2}), ets:delete(DRef, {Key2, Key1}), event_loop(ets:prev(ORef, {Key1, Key2}), ORef, DRef, Left, Accum, Keys, Delete); false -> case check_start_time(ST) of true when Delete == true -> ets:delete(ORef, {Key1, Key2}), ets:delete(DRef, {Key2, Key1}), event_loop(ets:prev(ORef, {Key1, Key2}), ORef, DRef, Left-1, [Event|Accum], Keys, Delete); true -> event_loop(ets:prev(ORef, {Key1, Key2}), ORef, DRef, Left-1, [Event|Accum], [{ORef, {Key1, Key2}}, {DRef, {Key2, Key1}}|Keys], Delete); false -> event_loop(ets:prev(ORef, {Key1, Key2}), ORef, DRef, Left, Accum, Keys, Delete) end end; event_loop({Key1, Key2, Key3}, ORef, DRef, Left, Accum, Keys, Delete) -> [{_,DL,ST,_PO,Event}]=ets:lookup(ORef, {Key1, Key2, Key3}), case check_deadline(DL) of true -> ets:delete(ORef, {Key1, Key2, Key3}), ets:delete(DRef, {Key3, Key2, Key1}), event_loop(ets:prev(ORef, {Key1, Key2, Key3}), ORef, DRef, Left, Accum, Keys, Delete); false -> case check_start_time(ST) of true when Delete == true -> ets:delete(ORef, {Key1, Key2, Key3}), ets:delete(DRef, {Key3, Key2, Key1}), event_loop(ets:prev(ORef, {Key1, Key2, Key3}), ORef, DRef, Left-1, [Event|Accum], Keys, Delete); true -> event_loop(ets:prev(ORef, {Key1, Key2, Key3}), ORef, DRef, Left-1, [Event|Accum], [{ORef, {Key1, Key2, Key3}}, {DRef, {Key3, Key2, Key1}}|Keys], Delete); false -> event_loop(ets:prev(ORef, {Key1, Key2, Key3}), ORef, DRef, Left, Accum, Keys, Delete) end end.%%------------------------------------------------------------%% function : delete_events%% Arguments: EventList - what's returned by get_event, get_events%% and add_and_get_event.%% Returns : %% Comment : Shall be invoked when it's safe to premanently remove%% the events found in the EventList.%% %%------------------------------------------------------------delete_events([]) -> ok;delete_events([{DB, Key}|T]) -> ets:delete(DB, Key), delete_events(T).%%------------------------------------------------------------%% function : update%% Arguments: %% Returns : %% Comment : As default we shall deliver Events in Priority order.%% Hence, if AnyOrder set we will still deliver in%% Priority order.%%------------------------------------------------------------update(undefined, _QoS) -> ok;update(DBRef, QoS) -> update(DBRef, QoS, undefined, undefined).update(DBRef, QoS, LifeFilter, PrioFilter) -> case updated_order(DBRef, ?not_GetOrderPolicy(QoS)) of false -> case updated_discard(DBRef, ?not_GetDiscardPolicy(QoS)) of false -> DBR2 = ?set_DefPriority(DBRef, ?not_GetPriority(QoS)), DBR3 = ?set_MaxEvents(DBR2, ?not_GetMaxEventsPerConsumer(QoS)), DBR4 = ?set_DefStopT(DBR3, ?not_GetTimeout(QoS)), DBR5 = ?set_StartTsupport(DBR4, ?not_GetStartTimeSupported(QoS)), DBR6 = ?set_StopTsupport(DBR5, ?not_GetStopTimeSupported(QoS)), case ets:info(?get_OrderRef(DBR6), size) of N when N =< ?get_MaxEvents(DBR6) -> %% Even if the QoS MaxEvents have been changed %% we don't reach the limit. DBR6; N -> %% The QoS MaxEvents must have been decreased. discard_events(DBR6, N-?get_MaxEvents(DBR6)), DBR6 end; true -> destroy_discard_db(DBRef), NewDBRef = create_db(QoS, ?get_GCTime(DBRef), ?get_GCLimit(DBRef), ?get_TimeRef(DBRef)), move_events(DBRef, NewDBRef, ets:first(?get_OrderRef(DBRef)), LifeFilter, PrioFilter) end; true -> destroy_discard_db(DBRef), NewDBRef = create_db(QoS, ?get_GCTime(DBRef), ?get_GCLimit(DBRef), ?get_TimeRef(DBRef)), move_events(DBRef, NewDBRef, ets:first(?get_OrderRef(DBRef)), LifeFilter, PrioFilter) end.updated_order(#dbRef{orderPolicy = Equal}, Equal) -> false;updated_order(#dbRef{orderPolicy = ?not_PriorityOrder}, ?not_AnyOrder) -> false;updated_order(#dbRef{orderPolicy = ?not_AnyOrder}, ?not_PriorityOrder) -> false;updated_order(_, _) -> true.updated_discard(#dbRef{discardPolicy = Equal}, Equal) -> false;updated_discard(#dbRef{discardPolicy = ?not_RejectNewEvents}, ?not_AnyOrder) -> false;updated_discard(#dbRef{discardPolicy = ?not_AnyOrder}, ?not_RejectNewEvents) -> false;updated_discard(_, _) -> true.move_events(DBRef, NewDBRef, '$end_of_table', _, _) -> destroy_order_db(DBRef), case ets:info(?get_OrderRef(NewDBRef), size) of N when N =< ?get_MaxEvents(NewDBRef) -> %% Even if the QoS MaxEvents have been changed %% we don't reach the limit. NewDBRef; N -> %% The QoS MaxEvents must have been decreased. discard_events(DBRef, N-?get_MaxEvents(NewDBRef)), NewDBRef end;move_events(DBRef, NewDBRef, Key, LifeFilter, PrioFilter) -> [{Keys, DeadLine, StartTime, PriorityOverride, Event}] = ets:lookup(?get_OrderRef(DBRef), Key), case check_deadline(DeadLine) of true -> ok; _-> write_event(?get_OrderP(DBRef), {Keys, DeadLine, StartTime, PriorityOverride, Event}, DBRef, NewDBRef, Key, LifeFilter, PrioFilter) end, ets:delete(?get_OrderRef(DBRef), Key), move_events(DBRef, NewDBRef, ets:next(?get_OrderRef(DBRef), Key), LifeFilter, PrioFilter).%% We cannot use do_add_event directly since we MUST lookup the timestamp (TS).write_event(?not_DeadlineOrder, {{_, TS, _Prio}, DL, ST, PO, Event}, _DBRef, NewDBRef, _Key, _LifeFilter, _PrioFilter) -> StartT = update_starttime(NewDBRef, Event, ST), %% Deadline and Priority previously extracted. do_add_event(NewDBRef, Event, TS, DL, StartT, PO);write_event(?not_DeadlineOrder, {{_, TS}, DL, _ST, PO, Event}, _DBRef, NewDBRef, _Key, _LifeFilter, PrioFilter) -> %% Priority not previously extracted. POverride = update_priority(NewDBRef, PrioFilter, Event, PO), StartT = extract_start_time(Event, ?get_StartTsupport(NewDBRef), ?get_TimeRef(NewDBRef)), do_add_event(NewDBRef, Event, TS, DL, StartT, POverride);write_event(?not_FifoOrder, {{TS, _PorD}, DL, ST, PO, Event}, _DBRef, NewDBRef, _Key, LifeFilter, PrioFilter) -> %% Priority or Deadline have been extracted before but we cannot tell which. POverride = update_priority(NewDBRef, PrioFilter, Event, PO), DeadL = update_deadline(NewDBRef, LifeFilter, Event, TS, DL), StartT = update_starttime(NewDBRef, Event, ST), do_add_event(NewDBRef, Event, TS, DeadL, StartT, POverride);write_event(?not_FifoOrder, {TS, DL, ST, PO, Event}, _DBRef, NewDBRef, _Key, LifeFilter, PrioFilter) -> %% Priority and Deadline not extracetd before. Do it now.
⌨️ 快捷键说明
复制代码Ctrl + C
搜索代码Ctrl + F
全屏模式F11
增大字号Ctrl + =
减小字号Ctrl + -
显示快捷键?