Packages
bpe
4.6.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.erl
-module(bpe).
-author('Maxim Sokhatsky').
-include("bpe.hrl").
-include_lib("kvs/include/cursors.hrl").
-include("api.hrl").
-compile(export_all).
-define(TIMEOUT, application:get_env(bpe,timeout,60000)).
load(#process{id = ProcName}) -> {ok,Proc} = kvs:get(process,ProcName), Proc;
load(ProcName) -> {ok,Proc} = kvs:get(process,ProcName), Proc.
cleanup(P) ->
[ kvs:delete({hist,P},Id) || #hist{id=Id} <- bpe:hist(P) ],
kvs:delete(writer,{hist,P}),
kvs:delete(process,P).
start(Proc0, Options) ->
Pid = proplists:get_value(notification,Options,undefined),
Proc = case Proc0#process.id == [] of
true -> Id = kvs:seq([],[]),
Proc0#process{id=Id,task=Proc0#process.beginEvent,
options = Options,notifications = Pid,
started=calendar:local_time()};
_ -> Proc0#process{started=calendar:local_time()} end,
kvs:append(Proc, process),
Key = {hist,Proc#process.id},
kvs:ensure(#writer{id=Key}),
kvs:append(#hist{ id = 0,
name = Proc#process.name,
time = Proc#process.started,
docs = Proc#process.docs,
task = { event, Proc#process.beginEvent }}, Key),
Restart = transient,
Shutdown = ?TIMEOUT,
ChildSpec = { Proc#process.id,
{bpe_proc, start_link, [Proc]},
Restart, Shutdown, worker, [bpe_proc] },
case supervisor:start_child(bpe_otp,ChildSpec) of
{ok,_} -> {ok,Proc#process.id};
{ok,_,_} -> {ok,Proc#process.id};
{error,_} -> {error,Proc#process.id} end.
find_pid(Id) -> bpe:cache({process,Id}).
proc(ProcId) -> gen_server:call(find_pid(ProcId),{get}, ?TIMEOUT).
complete(ProcId) -> gen_server:call(find_pid(ProcId),{complete}, ?TIMEOUT).
run(ProcId) -> gen_server:call(find_pid(ProcId),{run}, ?TIMEOUT).
until(ProcId,Task) -> gen_server:call(find_pid(ProcId),{until,Task}, ?TIMEOUT).
complete(Stage,ProcId) -> gen_server:call(find_pid(ProcId),{complete,Stage}, ?TIMEOUT).
amend(ProcId,Form) -> gen_server:call(find_pid(ProcId),{amend,Form}, ?TIMEOUT).
amend(ProcId,Form,noflow) -> gen_server:call(find_pid(ProcId),{amend,Form,true},?TIMEOUT).
event(ProcId,Event) -> gen_server:call(find_pid(ProcId),{event,Event}, ?TIMEOUT).
delete_tasks(Proc, Tasks) ->
Proc#process { tasks = [ Task || Task <- Proc#process.tasks,
lists:member(Task#task.name,Tasks) ] }.
hist(ProcId) -> kvs:all({hist,ProcId}).
hist(ProcId,N) -> case kvs:get({hist,ProcId},N) of
{ok,Res} -> Res;
{error,_Reason} -> [] end.
source(Name, Proc) ->
case [ Task || Task <- events(Proc), element(#task.name,Task) == Name] of
[T] -> T;
[] -> #beginEvent{};
E -> E end.
step(Name, Proc) ->
case [ Task || Task <- tasks(Proc), element(#task.name,Task) == Name] of
[T] -> T;
[] -> #task{};
E -> E end.
doc(Rec, Proc) ->
case [ Doc || Doc <- docs(Proc), element(1,Doc) == element(1,Rec)] of
[D] -> D;
[] -> [];
E -> E end.
docs (Proc) -> Proc#process.docs.
tasks (Proc) -> Proc#process.tasks.
events(Proc) -> Proc#process.events.
% Process Schema
new_task(Proc,GivenTask) ->
Existed = [ Task || Task<- Proc#process.tasks, Task#task.name == GivenTask#task.name],
case Existed of
[] -> Proc#process{tasks=[GivenTask|Proc#process.tasks]};
_ -> {error,exist,Existed} end.
delete(_Proc) -> ok.
val(Document,Proc,Cond) -> val(Document,Proc,Cond,fun(_,_)-> ok end).
val(Document,Proc,Cond,Action) ->
case Cond(Document,Proc) of
true -> Action(Document,Proc), {reply,Proc};
{false,Message} -> {{reply,Message},Proc#process.task,Proc};
ErrorList -> io:format("BPE:val/4 failed: ~tp~n",[ErrorList]),
{{reply,ErrorList},Proc#process.task,Proc} end.
cache(Key, undefined) -> ets:delete(processes,Key);
cache(Key, Value) -> ets:insert(processes,{Key,till(calendar:local_time(), ttl()),Value}), Value.
cache(Key, Value, Till) -> ets:insert(processes,{Key,Till,Value}), Value.
cache(Key) ->
Res = ets:lookup(processes,Key),
Val = case Res of [] -> undefined; [Value] -> Value; Values -> Values end,
case Val of undefined -> undefined;
{_,infinity,X} -> X;
{_,Expire,X} -> case Expire < calendar:local_time() of
true -> ets:delete(processes,Key), undefined;
false -> X end end.
ttl() -> application:get_env(bpe,ttl,60*15).
till(Now,TTL) ->
calendar:gregorian_seconds_to_datetime(
calendar:datetime_to_gregorian_seconds(Now) + TTL).
send(Pool, Message) -> syn:publish(term_to_binary(Pool),Message).
reg(Pool) -> reg(Pool,undefined).
reg(Pool, Value) ->
case get({pool,Pool}) of
undefined -> syn:register(term_to_binary(Pool),self(),Value),
syn:join(term_to_binary(Pool),self()),
erlang:put({pool,Pool},Pool);
_Defined -> skip end.
unreg(Pool) ->
case get({pool,Pool}) of
undefined -> skip;
_Defined -> syn:leave(Pool, self()),
erlang:erase({pool,Pool}) end.
reload(Module) ->
{Module, Binary, Filename} = code:get_object_code(Module),
case code:load_binary(Module, Filename, Binary) of
{module, Module} ->
{reloaded, Module};
{error, Reason} ->
{load_error, Module, Reason}
end.