Packages
bpe
9.9.7
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").
-include_lib("kvs/include/kvs.hrl").
-include_lib("kvs/include/cursors.hrl").
-behaviour(gen_server).
-export([start_link/1]).
-export([init/1,
handle_call/3,
handle_cast/2,
handle_continue/2,
handle_info/2,
terminate/2,
code_change/3]).
-export([debug/6,
process_task/2,
process_task/3]).
-define(REQUEST_LIMIT, application:get_env(bpe, request_limit, 500)).
start_link(Parameters) ->
gen_server:start_link(?MODULE, Parameters, []).
debug(Proc, Name, Targets, Target, Status, Reason) ->
case application:get_env(bpe, debug, true) of
true ->
logger:notice("BPE: ~ts [~ts:~ts] ~p/~p ~p",
[Proc#process.id,
Name,
Target,
Status,
Reason,
Targets]);
false -> skip
end.
process_event(EventType, Event, #process{id=Pid, module=Module} = Proc) ->
Name = kvs:field(Event, name),
#result{type=T, state = NewProc} = Res =
case Module:action(Event, Proc) of
#result{type = reply} = R when EventType == async -> R#result{type = noreply, reply = []};
#result{} = R when EventType == async -> R#result{reply = []};
R -> R
end,
debug(Proc, atom_to_list(element(1, Event)), "", Name, T, ""),
kvs:append(NewProc, "/bpe/proc"),
kvs:remove(Event, bpe:key("/bpe/messages/queue/", Pid)),
kvs:append(erlang:setelement(2, Event, kvs:seq([], [])), bpe:key("/bpe/messages/hist/", Pid)),
case Res of
#result{continue=C} = X when C =/= [] ->
bpe:constructResult(X#result{opt={continue, C}});
X -> bpe:constructResult(X)
end.
process_task(Stage, Proc) ->
process_task(Stage, Proc, false).
process_task(Stage, Proc, NoFlow) ->
{_, Curr} = bpe:current_task(Proc),
Task = bpe:step(Proc, Curr),
Targets = case NoFlow of
true -> noflow;
_ -> bpe_task:targets(Curr, Proc)
end,
{Status, {Reason, Target}, ProcState} = case {Targets,
Curr,
Stage}
of
{noflow, _, _} -> {reply, {complete, Curr}, Proc};
{[], [], _} -> 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,
case Status == stop orelse NoFlow == true of
true -> {stop, Reason, Target, ProcState};
_ ->
bpe:add_trace(ProcState, [], Target),
debug(ProcState, Curr, Targets, Target, Status, Reason),
{Status, {Reason, Target}, ProcState}
end.
behavior_fun_args(handle_call, Req, From, St) -> [Req, From, St];
behavior_fun_args(_, Req, _, St) -> [Req, St].
behavior_fun(asyncEvent, _) -> handle_cast;
behavior_fun(broadcastEvent, {_, broadcastEvent, #broadcastEvent{type = immediate}}) -> handle_cast;
behavior_fun(_, _) -> handle_call.
convert_api_args(proc, [Id, _ProcId]) -> {Id, get};
convert_api_args(update, [Id, _ProcId, State]) -> {Id, set, State};
convert_api_args(assign, [Id, _ProcId]) -> {Id, ensure_mon};
convert_api_args(Fn, [Id, _ | Args]) ->
list_to_tuple([Id, Fn | Args]).
continueId([#continue{} | _] = X, Id) ->
lists:map(fun (C) -> C#continue{id = Id} end, X);
continueId(_, _) -> [].
delete_lock(Id) -> kvs:delete(terminateLock, Id, #kvs{mod=kvs_mnesia}).
handle_result(call, {noreply, S}) -> {reply, ok, S};
handle_result(call, {noreply, S, X}) -> {reply, ok, S, X};
handle_result(_, X) -> X.
handle_result(T, Id, X, #process{id = Pid} = DefState) ->
handle_result(T, terminate_check(Id, X, kvs:index(terminateLock, pid, Pid, #kvs{mod=kvs_mnesia}), ?REQUEST_LIMIT, DefState)).
terminate_check(_, _, L, Limit, #process{id = Pid} = DefState) when length(L) > Limit ->
logger:error("TERMINATE LOCK LIMIT: ~tp", [Pid]),
bpe:cache(terminateLocks, {terminate, Pid}, self()),
{stop, normal, DefState};
terminate_check(Id, X, [#terminateLock{id=Id}], _, #process{id = Pid}) ->
bpe:cache(terminateLocks, {terminate, Pid}, case element(1, X) of stop -> self(); _ -> undefined end),
delete_lock(Id), X;
terminate_check(Id, {stop, normal, S}, [#terminateLock{} | _], _, #process{}) ->
delete_lock(Id), {noreply, S};
terminate_check(Id, {stop, normal, Reply, S}, [#terminateLock{} | _], _, #process{}) ->
delete_lock(Id), {reply, Reply, S};
terminate_check(Id, X, _, _, #process{id = Pid}) ->
bpe:cache(terminateLocks, {terminate, Pid}, case element(1, X) of stop -> self(); _ -> undefined end),
delete_lock(Id), X.
handleContinue({noreply, State}, Continue, Id) when Continue =/= [] ->
{noreply, State, {continue, continueId(lists:flatten(Continue), Id)}};
handleContinue({reply, Reply, State, {continue, C}}, Continue, Id) ->
{reply, Reply, State, {continue, continueId(lists:flatten(Continue) ++ C, Id)}};
handleContinue({reply, Reply, State, _}, Continue, Id) when Continue =/= [] ->
{reply, Reply, State, {continue, continueId(lists:flatten(Continue), Id)}};
handleContinue({noreply, State, {continue, C}}, Continue, Id) ->
{noreply, State, {continue, continueId(lists:flatten(Continue) ++ C, Id)}};
handleContinue({noreply, State, _}, Continue, Id) when Continue =/= [] ->
{noreply, State, {continue, continueId(lists:flatten(Continue), Id)}};
handleContinue({reply, Reply, State}, Continue, Id) when Continue =/= [] ->
{reply, Reply, State, {continue, continueId(lists:flatten(Continue), Id)}};
handleContinue({stop, _, State}, Continue, Id) when Continue =/= [] ->
{noreply, State, {continue, continueId(lists:flatten(Continue) ++ [#continue{type=stop}], Id)}};
handleContinue({stop, _, Reply, State}, Continue, Id) when Continue =/= [] ->
{reply, Reply, State, {continue, continueId(lists:flatten(Continue) ++ [#continue{type=stop}], Id)}};
handleContinue(X, _, _) -> X.
% BPMN 2.0 Инфотех
handle_call({stop, Reason}, _, Proc) -> {stop, Reason, Proc};
handle_call({Id, mon_link, MID}, _, Proc) ->
logger:notice("BPE MON_LINK: ~p.", [Proc#process.id]),
ProcNew = Proc#process{monitor = MID},
handle_result(call, Id, {reply, ProcNew, ProcNew}, Proc);
handle_call({Id, ensure_mon}, _, Proc) ->
logger:notice("BPE ENSURE_MON: ~p.", [Proc#process.id]),
{Mon, ProcNew} = bpe:ensure_mon(Proc),
handle_result(call, Id, {stop, normal, Mon, ProcNew}, Proc);
handle_call({Id, get}, _, Proc) -> logger:notice("BPE GET: ~p.", [Proc#process.id]), handle_result(call, Id, {stop, normal, Proc, Proc}, Proc);
handle_call({Id, set, State}, _, Proc) ->
logger:notice("BPE SET: ~p.", [Proc#process.id]),
handle_result(call, Id, {stop, normal, Proc, State}, State);
handle_call({Id, persist, State}, _, #process{} = Proc) ->
logger:notice("BPE PERSIST: ~p.", [Proc#process.id]),
kvs:append(State, "/bpe/proc"),
handle_result(call, Id, {stop, normal, State, State}, State);
handle_call({Id, next}, _, #process{} = Proc) ->
try handle_result(call, Id, handleContinue(bpe:processFlow(Proc), [], Id), Proc) catch
_X:_Y:Z -> {stop, {error, 'next/1', Z}, {error, 'next/1', Z}, Proc}
end;
handle_call({Id, next, [#continue{} | _] = Continue}, _, #process{} = Proc) ->
try handle_result(call, Id, handleContinue(bpe:processFlow(Proc), Continue, Id), Proc) catch
_X:_Y:Z -> {stop, {error, 'next/1', Z}, {error, 'next/1'}, Proc}
end;
handle_call({Id, next, Stage}, _, Proc) ->
try handle_result(call, Id, handleContinue(bpe:processFlow(Stage, Proc), [], Id), Proc) catch
_X:_Y:Z -> {stop, {error, 'next/2', Z}, {error, 'next/2', Z}, Proc}
end;
handle_call({Id, amend, Form}, _, Proc) ->
try handle_result(call, Id, handleContinue(bpe:processFlow(bpe_env:append(env, Proc, Form)), [], Id), Proc)
catch
_X:_Y:Z -> {stop, {error, 'amend/2', Z}, {error, 'amend/2', Z}, Proc}
end;
handle_call({Id, discard, Form}, _, Proc) ->
try handle_result(call, Id, handleContinue(bpe:processFlow(bpe_env:remove(env, Proc, Form)), [], Id), Proc)
catch
_X:_Y:Z -> {stop, {error, 'amend/2', Z}, {error, 'discard/2', Z}, Proc}
end;
handle_call({Id, messageEvent, Event}, _, Proc) ->
try handle_result(call, Id, handleContinue(process_event(sync, Event, Proc), [], Id), Proc) catch
_X:_Y:Z -> {stop, {error, 'messageEvent/2', Z}, {error, 'messageEvent/2', Z}, Proc}
end;
handle_call({Id, messageEvent, Event, [#continue{} | _] = Continue}, _, Proc) ->
try handle_result(call, Id, handleContinue(process_event(sync, Event, Proc), Continue, Id), Proc) catch
_X:_Y:Z -> {stop, {error, 'messageEvent/3', Z}, {error, 'messageEvent/3', Z}, Proc}
end;
% BPMN 1.0 ПриватБанк
handle_call({Id, complete}, _, Proc) ->
try handle_result(call, Id, handleContinue(process_task([], Proc), [], Id), Proc) catch
_X:_Y:Z -> {stop, {error, 'complete/1', Z}, {error, 'complete/1', Z}, Proc}
end;
handle_call({Id, complete, [#continue{} | _] = Continue}, _, Proc) ->
try handle_result(call, Id, handleContinue(process_task([], Proc), Continue, Id), Proc) catch
_X:_Y:Z -> {stop, {error, 'complete/2', Z}, {error, 'complete/2', Z}, Proc}
end;
handle_call({Id, complete, Stage}, _, Proc) ->
try handle_result(call, Id, handleContinue(process_task(Stage, Proc), [], Id), Proc) catch
_X:_Y:Z -> {stop, {error, 'complete/2', Z}, {error, 'complete/2', Z}, Proc}
end;
handle_call({Id, modify, Form, append}, _, Proc) ->
try handle_result(call, Id, handleContinue(process_task([],
bpe_env:append(env, Proc, Form),
true), [], Id), Proc)
catch
_X:_Y:Z -> {stop, {error, 'append/2', Z}, {error, 'append/2', Z}, Proc}
end;
handle_call({Id, modify, Form, remove}, _, Proc) ->
try handle_result(call, Id, handleContinue(process_task([],
bpe_env:remove(env, Proc, Form),
true), [], Id), Proc)
catch
_X:_Y:Z -> {stop, {error, 'remove/2', Z}, {error, 'remove/2', Z}, Proc}
end;
handle_call({Id, mon_link, MID, Continue}, _, Proc) ->
logger:notice("BPE MON_LINK CONTINUE: ~p.", [Proc#process.id]),
ProcNew = Proc#process{monitor = MID},
handle_result(call, Id, handleContinue({stop, normal, ProcNew, ProcNew}, Continue, Id), Proc);
handle_call({Id, ensure_mon, Continue}, _, Proc) ->
logger:notice("BPE ENSURE_MON CONTINUE: ~p.", [Proc#process.id]),
{Mon, ProcNew} = bpe:ensure_mon(Proc),
handle_result(call, Id, handleContinue({stop, normal, Mon, ProcNew}, Continue, Id), Proc);
handle_call({Id, set, State, Continue}, _, Proc) ->
logger:notice("BPE SET CONTINUE: ~p.", [Proc#process.id]),
handle_result(call, Id, handleContinue({stop, normal, Proc, State}, Continue, Id), State);
handle_call({Id, persist, State, Continue}, _, #process{} = Proc) ->
logger:notice("BPE PERSIST CONTINUE: ~p.", [Proc#process.id]),
kvs:append(State, "/bpe/proc"),
handle_result(call, Id, handleContinue({stop, normal, State, State}, Continue, Id), State);
handle_call({Id, next, Stage, Continue}, _, Proc) ->
try handle_result(call, Id, handleContinue(bpe:processFlow(Stage, Proc), Continue, Id), Proc) catch
_X:_Y:Z -> {stop, {error, 'next/2', Z}, {error, 'next/2', Z}, Proc}
end;
handle_call({Id, amend, Form, Continue}, _, Proc) ->
try handle_result(call, Id, handleContinue(bpe:processFlow(bpe_env:append(env, Proc, Form)), Continue, Id), Proc)
catch
_X:_Y:Z -> {stop, {error, 'amend/2', Z}, {error, 'amend/2', Z}, Proc}
end;
handle_call({Id, discard, Form, Continue}, _, Proc) ->
try handle_result(call, Id, handleContinue(bpe:processFlow(bpe_env:remove(env, Proc, Form)), Continue, Id), Proc)
catch
_X:_Y:Z -> {stop, {error, 'discard/2', Z}, {error, 'discard/2', Z}, Proc}
end;
handle_call({Id, complete, Stage, Continue}, _, Proc) ->
try handle_result(call, Id, handleContinue(process_task(Stage, Proc), Continue, Id), Proc) catch
_X:_Y:Z -> {stop, {error, 'complete/2', Z}, {error, 'complete/2', Z}, Proc}
end;
handle_call({Id, modify, Form, append, Continue}, _, Proc) ->
try handle_result(call, Id, handleContinue(process_task([],
bpe_env:append(env, Proc, Form),
true), Continue, Id), Proc)
catch
_X:_Y:Z -> {stop, {error, 'append/2', Z}, {error, 'append/2', Z}, Proc}
end;
handle_call({Id, modify, Form, remove, Continue}, _, Proc) ->
try handle_result(call, Id, handleContinue(process_task([],
bpe_env:remove(env, Proc, Form),
true), Continue, Id), Proc)
catch
_X:_Y:Z -> {stop, {error, 'remove/2', Z}, {error, 'remove/2', Z}, Proc}
end;
handle_call(Command, _, Proc) ->
handle_result(call, kvs:seq([], []), {stop, unknown, {unknown, Command}, Proc}, Proc).
handle_continue([#continue{type=spawn, module=Module, fn=Fn, args=Args} | T], Proc) ->
spawn(fun() -> apply(Module, Fn, Args) end),
{noreply, Proc, {continue, T}};
handle_continue([#continue{id = Id, type=bpe, fn=Fn, args=Args} | T], Proc) ->
BpeArgs = convert_api_args(Fn, [Id | Args]),
Result = try apply(bpe_proc, behavior_fun(Fn, BpeArgs), behavior_fun_args(behavior_fun(Fn, BpeArgs), BpeArgs, [], Proc))
catch
_X:_Y:Z -> {stop, {error, Z}, Proc}
end,
case Result of
{noreply, State} ->
{noreply, State, {continue, T}};
{reply, _, State, {continue, C}} ->
{noreply, State, {continue, T ++ continueId(C, Id)}};
{reply, _, State, _} ->
{noreply, State, {continue, T}};
{noreply, State, {continue, C}} ->
{noreply, State, {continue, T ++ continueId(C, Id)}};
{noreply, State, _} ->
{noreply, State, {continue, T}};
{reply, _, State} ->
{noreply, State, {continue, T}};
{stop, _, State} ->
{noreply, State, {continue, T ++ [#continue{id = Id, type=stop}]}};
{stop, _, _, State} ->
{noreply, State, {continue, T ++ [#continue{id = Id, type=stop}]}};
X -> X
end;
handle_continue([#continue{id=Id, type=stop}], Proc) ->
handle_result(continue, Id, {stop, normal, Proc}, Proc);
handle_continue([#continue{type=stop} = X | T], Proc) ->
case lists:keymember(stop, #continue.type, T) of
true -> {noreply, Proc, {continue, T}};
_ -> {noreply, Proc, {continue, T ++ [X]}}
end;
handle_continue([], Proc) ->
handle_result(continue, [], {stop, normal, Proc}, Proc).
init(Process) ->
process_flag(trap_exit, true),
bpe:cache(terminateLocks, {terminate, Process#process.id}, undefined),
Proc = bpe:load(Process#process.id, Process),
logger:notice("BPE: ~ts spawned as ~p",
[Proc#process.id, self()]),
Till = bpe:till(calendar:local_time(),
application:get_env(bpe, ttl, 24 * 60 * 60)),
bpe:cache(processes, {process, Proc#process.id}, self(), Till),
[bpe:reg({messageEvent,
element(1, EventRec),
Proc#process.id})
|| EventRec <- bpe:events(Proc)],
{ok,
Proc#process{timer =
erlang:send_after(rand:uniform(10000),
self(),
{timer, ping})}}.
handle_cast({Id, asyncEvent, Event}, Proc) ->
try handle_result(cast, Id, handleContinue(process_event(async, Event, Proc), [], Id), Proc) catch
_X:_Y:Z -> {stop, {error, 'asyncEvent/2', Z}, Proc}
end;
handle_cast({Id, asyncEvent, Event, [#continue{} | _] = Continue}, Proc) ->
try handle_result(cast, Id, handleContinue(process_event(async, Event, Proc), Continue, Id), Proc) catch
_X:_Y:Z -> {stop, {error, 'asyncEvent/3', Z}, Proc}
end;
handle_cast({Id, broadcastEvent, Event}, Proc) ->
try handle_result(cast, Id, handleContinue(process_event(async, Event, Proc), [], Id), Proc) catch
_X:_Y:Z -> {stop, {error, 'broadcastEvent/2', Z}, Proc}
end;
handle_cast(Msg, State) ->
logger:notice("BPE: Unknown API async: ~p.", [Msg]),
handle_result(cast, kvs:seq([], []), {stop, {error, {unknown_cast, Msg}}, State}, State).
handle_info({timer, ping},
State = #process{timer = _Timer, id = Id,
events = _Events, notifications = _Pid}) ->
lists:foreach(fun (Event) ->
try process_event(async, Event, State) catch
_X:_Y:Z -> Z
end
end, kvs:all(bpe:key("/bpe/messages/queue/", Id))),
case application:get_env(bpe, ping_discipline, bpe_ping) of
undefined -> {noreply, State};
M -> M:ping(State)
end;
handle_info({'DOWN', _, _, _, _} = Msg, State) ->
logger:notice("BPE: Connection closed, shutting down "
"session: ~p.",
[Msg]),
handle_result(info, kvs:seq([], []), {stop, normal, State}, State);
handle_info({'EXIT', _, Reason} = Msg, State) ->
logger:notice("BPE EXIT: ~p.", [Msg]),
handle_result(info, kvs:seq([], []), {stop, Reason, State}, State);
handle_info(Info, State = #process{}) ->
logger:notice("BPE: Unrecognized info: ~p", [Info]),
handle_result(info, kvs:seq([], []), {stop, unknown_info, State}, State).
terminate(Reason, #process{id = Id} = Proc) ->
bpe:cache(terminateLocks, {terminate, Id}, self()),
logger:notice("BPE: ~ts terminate Reason: ~p ~tp", [Id, Reason, self()]),
lists:foreach(fun (Event) ->
try process_event(async, Event, Proc) catch
_X:_Y:Z -> Z
end
end, kvs:all(bpe:key("/bpe/messages/queue/", Id))),
bpe:cache(processes, {process, Id}, undefined),
ok.
code_change(_OldVsn, State, _Extra) -> {ok, State}.