-module(jason_stream). -behaviour(gen_statem). -export([start_link/1]). -export([init/1,callback_mode/0,terminate/3]). -export([token/3, sofar/3]). -include("jason.hrl"). -record(stream, {iodevice = undefined % device to read ,bytes = 1024 % bytes to read ,holdt = [] % tokens hold for next parse ,holdc = [] % characters hold for next tokenisation ,opts = [] % jason opts (mode etc...) ,workers = [] % workers seen by server ,depth = 0 % depth of streamin (0, or 1 for streaming of a collection) ,array_cnt = 0 % counter of array (collection) found (for better context of workers) ,context = [] % track worker context (store worker pid and array_cnt it was working on) ,line = 0 % line number of document ,offset = 0 % offset in document }). start_link({Name, IoDevice, Bytes, JasonOpts}) -> gen_statem:start_link({global, Name}, ?MODULE, {Name, IoDevice, Bytes, JasonOpts}, []). init({_Name, IoDevice, Bytes, JasonOpts}) -> O = jason_lib:options(JasonOpts), put(jason_aliases, O#opt.aliases), put(jason_records, O#opt.records), put(jason_binary, O#opt.binary), put(jason_mode, O#opt.mode), Data = #stream{iodevice = IoDevice, bytes = Bytes, opts = JasonOpts}, {ok, token, Data}. callback_mode() -> state_functions. %% %% Update bytes to read at first call token({call, From}, {bytes, NewBytes}, Data) when is_integer(NewBytes) -> erlang:display(?FUNCTION_NAME), gen_statem:reply(From, ok), {next_state, token, Data#stream{bytes = NewBytes, workers = update_workers(From, Data#stream.workers)}}; %% %% Read a single token to see if a collection is starting token({call, From}, read, Data) when Data#stream.depth == 0 -> erlang:display(?FUNCTION_NAME), case file:read(Data#stream.iodevice, Data#stream.bytes) of {ok, Json} -> JsonString = case is_binary(Json) of true -> erlang:binary_to_list(Json); false -> Json end, {Ret, State, DataNew, Depth} = case jason_lex:tokens(Data#stream.holdt, Data#stream.holdc ++ JsonString) of {more, Cont} -> erlang:display({more, Cont}), {'sofar', 'token', Data#stream{holdt = Cont}, 0} ; {done, TokenRet, RestChars} -> erlang:display(TokenRet), case TokenRet of {ok, {'b-a', _, _}, Line} -> % A collection is starting Cnt = Data#stream.array_cnt + 1, Ctxt = lists:keystore(From, 1, Data#stream.context, {From, Cnt}), {'array_start' , 'sofar' , Data#stream{line=Line ,offset = Data#stream.bytes + Data#stream.offset ,depth = 1 ,array_cnt = Cnt ,context = Ctxt ,holdc = Data#stream.holdc ++ RestChars } , 1}; {eof, Line} -> {'sofar' ,'sofar' ,Data#stream{line=Line ,offset = Data#stream.bytes + Data#stream.offset ,holdc = Data#stream.holdc ++ JsonString } ,0}; X -> {ok, _ , Line} = X, {'sofar' ,'sofar' ,Data#stream{line=Line ,offset = Data#stream.bytes + Data#stream.offset ,holdt = [] ,holdc = Data#stream.holdc ++ JsonString } ,0} end end, gen_statem:reply(From, Ret), {next_state, State, DataNew#stream{workers = update_workers(From, Data#stream.workers), depth = Depth}} ; eof when Data#stream.holdt =/= [] -> Sofar = case jason_yec:parse(Data#stream.holdt) of {ok, Ret } -> {'sofar', Ret} ; X -> {error, X} end, gen_statem:reply(From, Sofar), NewData = Data#stream{holdt = [], holdc = [], workers = update_workers(From, Data#stream.workers)}, {next_state, sofar, NewData} ; eof when Data#stream.holdc =/= [] -> % file is less than 1024 gen_statem:reply(From, jason:decode(Data#stream.holdc, Data#stream.opts)), NewData = Data#stream{holdt = [], holdc = [], workers = update_workers(From, Data#stream.workers)}, {next_state, sofar, NewData} ; eof -> gen_statem:reply(From, 'end'), {stop, normal, Data} ; {error, Reason} -> gen_statem:reply(From, {error, Reason}), {stop, {error, Reason}, Data} end. %% %% Update bytes to read sofar({call, From}, {bytes, NewBytes}, Data) when is_integer(NewBytes) -> erlang:display(?FUNCTION_NAME), {next_state, sofar, Data#stream{bytes = NewBytes, workers = update_workers(From, Data#stream.workers)}}; sofar({call, From}, read, Data) ->erlang:display(?FUNCTION_NAME), % Hold contains incomplete preceding incomplete data that should % prefix next grab % io:format("*****************************************~n~p~n", [Data]), case file:read(Data#stream.iodevice, Data#stream.bytes) of {ok, Json} -> JsonString = case is_binary(Json) of true -> erlang:binary_to_list(Json); false -> Json end, %erlang:display(Data#stream.holdc ++ JsonString), case jason_lex:tokens(Data#stream.holdt, Data#stream.holdc ++ JsonString) of {more, Cont} -> %io:format("*** more : ~nHoldt : ~p~nString : ~p~n-> Cont : ~p~n~n", [Data#stream.holdt, JsonString, Cont]), erlang:display({more, Cont}), gen_statem:reply(From, 'sofar'), {next_state, sofar, Data#stream{holdt = Cont}}; {done, TokenRet, RestChars} -> %io:format("*** done :~n-> Tokens : ~p~nRestChars : ~p~n~n", [TokenRet, RestChars]), erlang:display({done, TokenRet, RestChars}), % Try to parse tokens, if ok return Term, otherwise hold data %erlang:display({TokenRet, RestChars}), {Sofar, NewHoldc, NewHoldt} = case (catch jason_yec:parse(TokenRet)) of {ok, Res} -> {{sofar, Res}, RestChars, []} ; X -> io:format("Catch : ~p~n~n", [X]), {'sofar', RestChars, TokenRet} end, %io:format("~p~n~n", [{Sofar, NewHoldc, NewHoldt}]), NewData = Data#stream{holdt = NewHoldt, holdc = NewHoldc, workers = update_workers(From, Data#stream.workers)}, %io:format("Newdata : ~p~n~n", [NewData]), gen_statem:reply(From, Sofar), {next_state, sofar, NewData}; X -> erlang:display({ici, X}) end; eof when Data#stream.holdt =/= [] -> erlang:display(eof1), {tokens,_,_,_,_,_, Holdt,_,_} = Data#stream.holdt, Sofar = case jason_yec:parse(Holdt) of {ok, Ret } -> {'sofar', Ret} ; X -> {error, X} end, gen_statem:reply(From, Sofar), {next_state, sofar, Data#stream{holdt = [], holdc = [], workers = update_workers(From, Data#stream.workers)}} ; eof -> erlang:display(eof2), gen_statem:reply(From, 'end'), {stop, normal, Data#stream{workers = update_workers(From, Data#stream.workers)}} ; {error, Reason} -> gen_statem:reply(From, {error, Reason}), {stop, {error, Reason}, Data} end; sofar(_EventType, _EventContent, Data) -> erlang:display(?FUNCTION_NAME), erlang:display({_EventType, _EventContent, Data}), {keep_state,Data}. terminate(_Reason, _State, _Data) -> ok. %% %% update_workers({P, _}, L) -> case lists:member(P, L) of true -> L ; false -> L ++ [P] end.