Packages
bpe
2.4.0
13.5.22-aleph
11.4.16
11.4.15
11.4.14
11.4.13
9.9.7
9.9.6
8.12.4
8.12.3
8.12.1
8.12.0
retired
8.2.1
8.2.0
8.1.0
7.11.0
7.10.4
7.10.3
7.10.2
7.10.1
7.9.1
7.9.0
7.8.2
7.8.1
7.8.0
7.6.4
7.6.3
7.6.2
7.6.1
7.6.0
7.5.15
7.5.14
7.5.13
7.5.12
7.5.11
7.5.10
7.5.9
7.5.8
7.5.7
7.5.6
7.5.5
7.5.3
7.5.2
7.5.1
7.5.0
7.4.11
7.4.10
7.4.9
7.4.8
7.4.7
7.4.6
7.4.5
7.4.4
7.4.3
7.4.2
7.4.1
7.4.0
7.3.0
7.2.8
7.2.7
7.2.6
7.1.6
7.1.5
7.1.4
7.1.3
7.1.2
6.12.7
6.12.6
6.12.5
6.12.3
6.12.2
6.12.1
6.12.0
6.11.0
6.10.0
6.5.3
6.5.2
6.5.1
6.5.0
6.4.0
6.3.0
5.12.0
5.11.4
5.11.3
5.11.2
5.11.1
5.11.0
5.8.7
5.8.6
5.8.5
5.8.4
5.8.3
5.8.2
5.8.1
5.8.0
5.7.0
5.6.0
5.5.2
5.5.1
5.4.0
5.2.0
5.1.3
5.1.2
5.1.1
4.12.4
4.12.3
4.12.2
4.12.1
4.12.0
4.11.8
4.11.7
4.11.6
4.11.5
4.11.4
4.11.3
4.11.2
4.11.1
4.11.0
4.10.24
4.10.23
4.10.22
4.10.21
4.10.20
4.10.19
4.10.18
4.10.17
4.10.16
4.10.15
4.10.14
4.10.13
4.10.12
4.10.11
4.10.10
4.10.9
4.10.8
4.10.7
4.10.6
4.10.5
4.10.4
4.10.3
4.10.2
4.10.1
4.10.0
4.9.18
4.9.17
4.9.16
4.9.15
4.9.14
4.9.13
4.9.12
4.9.11
4.9.10
4.9.9
4.9.8
4.9.7
4.9.6
4.9.5
4.9.4
4.9.3
4.9.2
4.9.1
4.9.0
4.8.1
4.8.0
4.7.5
4.7.3
4.7.2
4.7.1
4.7.0
4.6.0
2.4.0
0.7.16
ERP/1: RTP GST WebRTC ICE SDP H.264 H.265 MP4 MPEG-2 HLS HEVC
Current section
Files
Jump to
Current section
Files
src/bpe_proc.erl
-module(bpe_proc).
-author('Maxim Sokhatsky').
-include("bpe.hrl").
-behaviour(gen_server).
-export([start_link/1]).
-export([init/1,handle_call/3,handle_cast/2,handle_info/2,terminate/2,code_change/3]).
-compile(export_all).
start_link(Parameters) -> gen_server:start_link(?MODULE, Parameters, []).
process_event(Event,Proc) ->
Targets = bpe_task:targets(element(#messageEvent.name,Event),Proc),
io:format("Event Targets: ~p",[Targets]),
{Status,{Reason,Target},ProcState} = bpe_event:handle_event(Event,bpe_task:find_flow(Targets),Proc),
kvs:add(#history { id = kvs:next_id("history",1),
feed_id = {history,ProcState#process.id},
name = ProcState#process.name,
time = calendar:local_time(),
task = { event, element(#messageEvent.name,Event) }}),
NewProcState = ProcState#process{task = Target},
FlowReply = fix_reply({Status,{Reason,Target},NewProcState}),
kvs:info(?MODULE,"Process ~p Flow Reply ~tp ",[Proc#process.id,{Status,{Reason,Target}}]),
kvs:put(transient(NewProcState)),
FlowReply.
run(Task,Process) ->
CurrentTask = Process#process.task,
case bpe_proc:process_flow([],Process,false) of
{reply,{complete,Reached},NewProc}
when Reached /= CurrentTask andalso Reached /= Task -> run(Task,NewProc);
Else -> Else end.
process_flow(Stage,Proc) -> process_flow(Stage,Proc,false).
process_flow(Stage,Proc,NoFlow) ->
Curr = Proc#process.task,
Term = [],
Task = bpe:task(Curr,Proc),
Targets = case NoFlow of
true -> noflow;
_ -> bpe_task:targets(Curr,Proc) end,
kvs:info(?MODULE,"Process ~p Task: ~p Targets: ~p",[Proc#process.id, Curr,Targets]),
{Status,{Reason,Target},ProcState} = case {Targets,Proc#process.task,Stage} of
{noflow,_,_} -> {reply,{complete,Curr},Proc};
{[],Term,_} -> bpe_task:already_finished(Proc);
{[],Curr,_} -> bpe_task:handle_task(Task,Curr,Curr,Proc);
{[],_,_} -> bpe_task:denied_flow(Curr,Proc);
{List,_,[]} -> bpe_task:handle_task(Task,Curr,bpe_task:find_flow(Stage,List),Proc);
{List,_,_} -> {reply,{complete,bpe_task:find_flow(Stage,List)},Proc} end,
kvs:add(#history { id = kvs:next_id("history",1),
feed_id = {history,ProcState#process.id},
name = ProcState#process.name,
time = calendar:local_time(),
task = {task, Curr} }),
NewProcState = ProcState#process{task = Target},
FlowReply = fix_reply({Status,{Reason,Target},NewProcState}),
kvs:info(?MODULE,"Process ~p Flow Reply ~tp ",[Proc#process.id,{Status,{Reason,Target}}]),
kvs:put(transient(NewProcState)),
FlowReply.
fix_reply({stop,{Reason,Reply},State}) -> {stop,Reason,Reply,State};
fix_reply(P) -> P.
handle_call({get},_,Proc) -> { reply,Proc,Proc };
handle_call({run},_,Proc) -> run('Finish',Proc);
handle_call({until,Stage},_,Proc) -> run(Stage,Proc);
handle_call({start},_,Proc) -> process_flow([],Proc);
handle_call({complete},_,Proc) -> process_flow([],Proc);
handle_call({complete,Stage},_,Proc) -> process_flow(Stage,Proc);
handle_call({event,Event},_,Proc) -> process_event(Event,Proc);
handle_call({amend,Form,true},_,Proc)
when is_list(Form) -> process_flow([],set_rec_in_proc(Proc,Form),true);
handle_call({amend,Form,true},_,Proc) -> process_flow([],Proc#process{docs=plist_setkey(element(1,Form),1,Proc#process.docs,Form)},true);
handle_call({amend,Form},_,Proc)
when is_list(Form) -> process_flow([],set_rec_in_proc(Proc,Form));
handle_call({amend,Form},_,Proc) -> process_flow([],Proc#process{docs=plist_setkey(element(1,Form),1,Proc#process.docs,Form)});
handle_call(Command,_,Proc) -> { reply,{unknown,Command},Proc }.
init(Process) ->
kvs:info(?MODULE,"Process ~p spawned ~p",[Process#process.id,self()]),
Proc = case kvs:get(process,Process#process.id) of
{ok,Exists} -> Exists;
{error,_} -> Process end,
Till = bpe:till(calendar:local_time(), kvs:config(bpe,ttl,24*60*60)),
bpe:cache({process,Proc#process.id},self(),Till),
[ bpe:reg({messageEvent,Name,Proc#process.id}) || {Name,_} <- bpe:events(Proc) ],
{ok, Proc#process{timer=erlang:send_after(crypto:rand_uniform(1,10000),self(),{timer,ping})}}.
handle_cast(Msg, State) ->
kvs:info(?MODULE,"Unknown API async: ~p", [Msg]),
{stop, {error, {unknown_cast, Msg}}, State}.
timer_restart(Diff) -> {X,Y,Z} = Diff, erlang:send_after(500*(Z+60*Y+60*60*X),self(),{timer,ping}).
ping() -> application:get_env(bpe,ping,{0,0,5}).
handle_info({timer,ping}, State=#process{task=Task,timer=Timer,id=Id,events=Events,notifications=Pid}) ->
case Timer of undefined -> skip; _ -> erlang:cancel_timer(Timer) end,
Wildcard = '*',
Terminal= case lists:keytake(Wildcard,#messageEvent.name,Events) of
{value,Event,_} -> {Wildcard,element(1,Event),element(#messageEvent.timeout,Event)};
false -> {Wildcard,boundaryEvent,{5,ping()}} end,
{Name,Record,{Days,Pattern}} = case lists:keytake(Task,#messageEvent.name,Events) of
{value,Event2,_} -> {Task,element(1,Event2),element(#messageEvent.timeout,Event2)};
false -> Terminal end,
Time2 = calendar:local_time(),
%kvs:info(?MODULE,"Ping: ~p, Task ~p, Event ~p, Record ~p ~n", [Id,Task,Name,Record]),
{DD,Diff} = case bpe:history(Id,1) of
[#history{time=Time1}] -> calendar:time_difference(Time1,Time2);
_ -> {immediate,timeout} end,
case {{DD,Diff} < {Days,Pattern}, Record} of
{true,_} -> {noreply,State#process{timer=timer_restart(ping())}};
{false,timeoutEvent} ->
kvs:info(?MODULE,"BPE process ~p: next step by timeout. ~nTime Diff is ~p~n",[Id,{DD,Diff}]),
case process_flow([],State) of
{reply,_,NewState} -> {noreply,NewState#process{timer=timer_restart(ping())}};
{stop,normal,_,NewState} -> {stop,normal,NewState} end;
{false,_} -> kvs:info(?MODULE,"BPE process ~p: Closing Timeout. ~nTime Diff is ~p~n",[Id,{DD,Diff}]),
case is_pid(Pid) of
true -> Pid ! {direct,{bpe,terminate,{Name,{Days,Pattern}}}};
false -> skip end,
bpe:cache({process,Id},undefined),
{stop,normal,State} end;
handle_info({'DOWN', _MonitorRef, _Type, _Object, _Info} = Msg, State = #process{id=Id}) ->
kvs:info(?MODULE, "connection closed, shutting down session:~p", [Msg]),
bpe:cache({process,Id},undefined),
{stop, normal, State};
handle_info(Info, State=#process{}) ->
kvs:info(?MODULE,"Unrecognized info: ~p", [Info]),
{noreply, State}.
terminate(Reason, #process{id=Id}) ->
kvs:info(?MODULE,"Terminating session Id cache: ~p~n Reason: ~p", [Id,Reason]),
spawn(fun() -> supervisor:delete_child(bpe_sup,Id) end),
bpe:cache({process,Id},undefined),
ok.
code_change(_OldVsn, State, _Extra) ->
{ok, State}.
plist_setkey(Name,Pos,List,New) ->
case lists:keyfind(Name,Pos,List) of
false -> [New|List];
_Element -> lists:keyreplace(Name,Pos,List,New) end.
set_rec_in_proc(Proc, []) -> Proc;
set_rec_in_proc(Proc, [H|T]) ->
ProcNew = Proc#process{ docs=plist_setkey(kvs:rname(element(1,H)),1,Proc#process.docs,H)},
set_rec_in_proc(ProcNew, T).
transient(#process{docs=Docs}=Process) ->
Process#process{docs=lists:filter(
fun (X) -> not lists:member(element(1,X),
kvs:config(bpe,transient,[])) end,Docs)}.