Packages
dqe
0.4.5
0.4.15
0.4.14
0.4.13
0.4.12
0.4.11
0.4.10
0.4.9
0.4.8
0.4.7
0.4.6
0.4.5
0.4.4
0.4.3
0.4.2
0.4.1
0.4.0
0.3.25
0.3.24
0.3.23
0.3.22
0.3.21
0.3.20
0.3.19
0.3.18
0.3.17
0.3.16
0.3.15
0.3.14
0.3.13
0.3.12
0.3.11
0.3.10
0.3.9
0.3.8
0.3.7
0.3.6
0.3.5
0.3.4
0.3.3
0.3.2
0.3.1
0.2.2
0.2.1
0.2.0
0.1.36
0.1.35
0.1.34
0.1.33
0.1.22
0.1.21
0.1.20
0.1.10
0.1.9
DalmatinerDB query engine
Current section
Files
Jump to
Current section
Files
src/dqe_events.erl
-module(dqe_events).
-behaviour(dflow).
-export([init/1, describe/1, start/2, emit/3, done/2]).
-record(state, {
bucket :: binary(),
start :: non_neg_integer(),
'end' :: non_neg_integer(),
filter = [],
chunk :: pos_integer()
}).
init([Bucket, Start, End, Filter]) ->
{ok, ChunkMs} = application:get_env(dqe, get_chunk),
Chunk = erlang:convert_time_unit(ChunkMs, milli_seconds, nano_seconds),
{ok, #state{start = Start, 'end' = End, bucket = Bucket, filter = Filter,
chunk = Chunk}, []}.
describe(#state{bucket = Bucket}) ->
[Bucket].
start(run,
State = #state{start = Start, 'end' = End, chunk = Chunk,
bucket = Bucket, filter = Filter}) when
Start + Chunk < End ->
%% We do a bit of cheating here this allows us to loop.
State1 = State#state{start = Start + Chunk},
case ddb_connection:read_events(Bucket, Start, End, Filter) of
{error, _Error} ->
{done, State1};
{ok, Events} ->
dflow:start(self(), run),
{emit, Events, State1}
end;
start(run, State = #state{start = Start, 'end' = End,
bucket = Bucket, filter = Filter}) ->
case ddb_connection:read_events(Bucket, Start, End, Filter) of
{error, _Error} ->
{done, State};
{ok, Data} ->
{done, Data, State}
end.
emit(_Child, _Data, State) ->
{ok, State}.
done(_, State) ->
{done, State}.