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