Current section
Files
Jump to
Current section
Files
src/grisp_connect_client.erl
%% @doc Client to interact with grisp.io
%%
%% This module contains a state machine to ensure connectivity with grisp.io.
%% JsonRPC traffic is managed here.
%% @end
-module(grisp_connect_client).
% External API
-export([start_link/0]).
-export([connect/0]).
-export([is_connected/0]).
-export([request/3]).
-export([notify/3]).
% Internal API
-export([connected/0]).
-export([disconnected/0]).
-export([handle_message/1]).
-export([reboot/0]).
-behaviour(gen_statem).
-export([init/1, terminate/3, code_change/4, callback_mode/0]).
% State Functions
-export([idle/3]).
-export([waiting_ip/3]).
-export([connecting/3]).
-export([connected/3]).
-include_lib("kernel/include/logger.hrl").
-define(STD_TIMEOUT, 1000).
-define(HANDLE_COMMON,
?FUNCTION_NAME(EventType, EventContent, Data) ->
handle_common(EventType, EventContent, ?FUNCTION_NAME, Data)).
-record(data, {
requests = #{}
}).
% API
start_link() ->
gen_statem:start_link({local, ?MODULE}, ?MODULE, [], []).
connect() ->
gen_statem:cast(?MODULE, ?FUNCTION_NAME).
is_connected() ->
gen_statem:call(?MODULE, ?FUNCTION_NAME).
request(Method, Type, Params) ->
gen_statem:call(?MODULE, {?FUNCTION_NAME, Method, Type, Params}).
notify(Method, Type, Params) ->
gen_statem:cast(?MODULE, {?FUNCTION_NAME, Method, Type, Params}).
connected() ->
gen_statem:cast(?MODULE, ?FUNCTION_NAME).
disconnected() ->
gen_statem:cast(?MODULE, ?FUNCTION_NAME).
handle_message(Payload) ->
gen_statem:cast(?MODULE, {?FUNCTION_NAME, Payload}).
reboot() ->
erlang:send_after(1000, ?MODULE, reboot).
% gen_statem CALLBACKS ---------------------------------------------------------
init([]) ->
{ok, Connect} = application:get_env(grisp_connect, connect),
NextState = case Connect of
true -> waiting_ip;
false -> idle
end,
{ok, NextState, #data{}}.
terminate(_Reason, _State, _Data) -> ok.
code_change(_Vsn, State, Data, _Extra) -> {ok, State, Data}.
callback_mode() -> [state_functions, state_enter].
%%% STATE CALLBACKS ------------------------------------------------------------
idle(enter, _OldState, _Data) ->
keep_state_and_data;
idle(cast, connect, Data) ->
{next_state, waiting_ip, Data};
?HANDLE_COMMON.
waiting_ip(enter, _OldState, _Data) ->
{keep_state_and_data, [{state_timeout, 0, retry}]};
waiting_ip(state_timeout, retry, Data) ->
case check_inet_ipv4() of
{ok, IP} ->
?LOG_INFO(#{event => checked_ip, ip => IP}),
{next_state, connecting, Data};
invalid ->
?LOG_DEBUG(#{event => waiting_ip}),
{next_state, waiting_ip, Data, [{state_timeout, ?STD_TIMEOUT, retry}]}
end;
?HANDLE_COMMON.
connecting(enter, _OldState, _Data) ->
{ok, Domain} = application:get_env(grisp_connect, domain),
{ok, Port} = application:get_env(grisp_connect, port),
?LOG_NOTICE(#{event => connecting, domain => Domain, port => Port}),
grisp_connect_ws:connect(Domain, Port),
keep_state_and_data;
connecting(cast, connected, Data) ->
?LOG_NOTICE(#{event => connected}),
{next_state, connected, Data};
connecting(cast, disconnected, _Data) ->
repeat_state_and_data;
?HANDLE_COMMON.
connected(enter, _OldState, _Data) ->
grisp_connect_log_server:start(),
keep_state_and_data;
connected({call, From}, is_connected, _) ->
{keep_state_and_data, [{reply, From, true}]};
connected(cast, disconnected, Data) ->
?LOG_WARNING(#{event => disconnected}),
grisp_connect_log_server:stop(),
{next_state, waiting_ip, Data};
connected(cast, {handle_message, Payload}, #data{requests = Requests} = Data) ->
Responses = grisp_connect_api:handle_msg(Payload),
% A reduce operation is needed to support jsonrpc batch comunications
case Responses of
[] ->
keep_state_and_data;
[{send_response, Response}] -> % Response for a GRiSP.io request
grisp_connect_ws:send(Response),
keep_state_and_data;
[{handle_response, ID, Response}] -> % handle a GRiSP.io response
{OtherRequests, Actions} = dispatch_response(ID, Response, Requests),
{keep_state, Data#data{requests = OtherRequests}, Actions}
end;
connected({call, From}, {request, Method, Type, Params},
#data{requests = Requests} = Data) ->
{ID, Payload} = grisp_connect_api:request(Method, Type, Params),
grisp_connect_ws:send(Payload),
NewRequests = Requests#{ID => From},
{keep_state,
Data#data{requests = NewRequests},
[{{timeout, ID}, request_timeout(), request}]};
connected(cast, {notify, Method, Type, Params}, _Data) ->
Payload = grisp_connect_api:notify(Method, Type, Params),
grisp_connect_ws:send(Payload),
keep_state_and_data;
?HANDLE_COMMON.
% Common event handling appended as last match case to each state_function
handle_common(cast, connect, State, _Data) when State =/= idle ->
keep_state_and_data;
handle_common({call, From}, is_connected, State, _) when State =/= connected ->
{keep_state_and_data, [{reply, From, false}]};
handle_common({call, From}, {request, _, _, _}, State, _Data)
when State =/= connected ->
{keep_state_and_data, [{reply, From, {error, disconnected}}]};
handle_common(cast, {notify, _Method, _Type, _Params}, _State, _Data) ->
% We ignore notifications sent while disconnected
keep_state_and_data;
handle_common({timeout, ID}, request, _, #data{requests = Requests} = Data) ->
Caller = maps:get(ID, Requests),
{keep_state,
Data#data{requests = maps:remove(ID, Requests)},
[{reply, Caller, {error, timeout}}]};
handle_common(info, reboot, _, _) ->
init:stop(),
keep_state_and_data;
handle_common(cast, Cast, _, _) ->
error({unexpected_cast, Cast});
handle_common({call, _}, Call, _, _) ->
error({unexpected_call, Call});
handle_common(info, Info, State, Data) ->
?LOG_ERROR(#{event => unexpected_info,
info => Info,
state => State,
data => Data}),
keep_state_and_data.
% INTERNALS --------------------------------------------------------------------
dispatch_response(ID, Response, Requests) ->
case maps:take(ID, Requests) of
{Caller, OtherRequests} ->
Actions = [{{timeout, ID}, cancel}, {reply, Caller, Response}],
{OtherRequests, Actions};
error ->
?LOG_DEBUG(#{event => ?FUNCTION_NAME, reason => {missing_id, ID},
data => Response}),
{Requests, []}
end.
request_timeout() ->
{ok, V} = application:get_env(grisp_connect, ws_requests_timeout),
V.
% IP check functions
check_inet_ipv4() ->
case get_ip_of_valid_interfaces() of
{IP1,_,_,_} = IP when IP1 =/= 127 -> {ok, IP};
_ -> invalid
end.
get_ipv4_from_opts([]) ->
undefined;
get_ipv4_from_opts([{addr, {_1, _2, _3, _4}} | _]) ->
{_1, _2, _3, _4};
get_ipv4_from_opts([_ | TL]) ->
get_ipv4_from_opts(TL).
has_ipv4(Opts) ->
get_ipv4_from_opts(Opts) =/= undefined.
flags_are_ok(Flags) ->
lists:member(up, Flags) and
lists:member(running, Flags) and
not lists:member(loopback, Flags).
get_valid_interfaces() ->
{ok, Interfaces} = inet:getifaddrs(),
[
Opts
|| {_Name, [{flags, Flags} | Opts]} <- Interfaces,
flags_are_ok(Flags),
has_ipv4(Opts)
].
get_ip_of_valid_interfaces() ->
case get_valid_interfaces() of
[Opts | _] -> get_ipv4_from_opts(Opts);
_ -> undefined
end.