Skip to content
This repository was archived by the owner on Nov 2, 2021. It is now read-only.
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
19 changes: 16 additions & 3 deletions components/proto_json/src/proto_json_rpc.erl
Original file line number Diff line number Diff line change
Expand Up @@ -22,7 +22,7 @@

-define(SERVER, ?MODULE).
-export([start_json_server/0]).
-export([send_message/8,
-export([send_message/8, send_message/9,
receive_message/3]).

-record(st, {
Expand All @@ -42,6 +42,10 @@ start_json_server() ->
rvi_common:start_json_rpc_server(protocol, ?MODULE, proto_json_sup).


send_message(CompSpec, TID, ServiceName, Timeout, ProtoOpts,
DataLinkMod, DataLinkOpts, Parameters) ->
send_message(CompSpec, TID, ServiceName, Timeout, ProtoOpts,
DataLinkMod, DataLinkOpts, [], Parameters).

send_message(CompSpec,
TID,
Expand All @@ -50,6 +54,7 @@ send_message(CompSpec,
ProtoOpts,
DataLinkMod,
DataLinkOpts,
Extra,
Parameters) ->
rvi_common:request(protocol, ?MODULE, send_message,
[{ transaction_id, TID },
Expand All @@ -58,6 +63,7 @@ send_message(CompSpec,
{ protocol_opts, ProtoOpts },
{ data_link_mod, DataLinkMod },
{ data_link_opts, DataLinkOpts },
{ extra, Extra },
{ parameters, Parameters }],
[ status ], CompSpec).

Expand All @@ -79,6 +85,7 @@ handle_rpc(<<"send_message">>, Args) ->
{ok, ProtoOpts} = rvi_common:get_json_element(["protocol_opts"], Args),
{ok, DataLinkMod} = rvi_common:get_json_element(["data_link_mod"], Args),
{ok, DataLinkOpts} = rvi_common:get_json_element(["data_link_opts"], Args),
Extra = rvi_common:get_opt_json_element(["extra"], [], Args),
{ok, Parameters} = rvi_common:get_json_element(["parameters"], Args),
[ ok ] = gen_server:call(?SERVER, { rvi, send_message,
[TID,
Expand All @@ -87,6 +94,7 @@ handle_rpc(<<"send_message">>, Args) ->
ProtoOpts,
DataLinkMod,
DataLinkOpts,
Extra,
Parameters,
LogId]}),
{ok, [ {status, rvi_common:json_rpc_status(ok)} ]};
Expand Down Expand Up @@ -121,17 +129,20 @@ handle_call({rvi, send_message,
ProtoOpts,
DataLinkMod,
DataLinkOpts,
Extra,
Parameters | _LogId]}, _From, St) ->
?debug(" protocol:send(): transaction id: ~p~n", [TID]),
?debug(" protocol:send(): service name: ~p~n", [ServiceName]),
?debug(" protocol:send(): timeout: ~p~n", [Timeout]),
?debug(" protocol:send(): opts: ~p~n", [ProtoOpts]),
?debug(" protocol:send(): data_link_mod: ~p~n", [DataLinkMod]),
?debug(" protocol:send(): data_link_opts: ~p~n", [DataLinkOpts]),
?debug(" protocol:send(): extra: ~p~n", [Extra]),
?debug(" protocol:send(): parameters: ~p~n", [Parameters]),
Data = [{ <<"service">>, ServiceName },
{ <<"timeout">>, Timeout },
{ <<"parameters">>, Parameters }
| Extra
],
RviOpts = rvi_common:rvi_options(Parameters),
Res = DataLinkMod:send_data(
Expand All @@ -146,17 +157,19 @@ handle_call(Other, _From, St) ->
handle_cast({rvi, receive_message, [Elems, IP, Port | _LogId]} = Msg, St) ->
?debug("~p:handle_cast(~p)", [?MODULE, Msg]),

[ ServiceName, Timeout, Parameters ] =
opts([<<"service">>, <<"timeout">>, <<"parameters">>],
[ ServiceName, Timeout, Src, Parameters ] =
opts([<<"service">>, <<"timeout">>, <<"src">>, <<"parameters">>],
Elems, undefined),

?debug(" protocol:rcv(): service name: ~p~n", [ServiceName]),
?debug(" protocol:rcv(): timeout: ~p~n", [Timeout]),
?debug(" protocol:rcv(): src: ~p~n", [Src]),
?debug(" protocol:rcv(): remote IP/Port: ~p~n", [{IP, Port}]),
service_edge_rpc:handle_remote_message(St#st.cs,
{IP, Port},
ServiceName,
Timeout,
[{<<"src">>, Src}],
Parameters),
{noreply, St};

Expand Down
25 changes: 18 additions & 7 deletions components/proto_msgpack/src/proto_msgpack_rpc.erl
Original file line number Diff line number Diff line change
Expand Up @@ -22,7 +22,7 @@

-define(SERVER, ?MODULE).
-export([start_json_server/0]).
-export([send_message/8,
-export([send_message/8, send_message/9,
receive_message/3]).

-record(st, {
Expand All @@ -44,7 +44,10 @@ init([]) ->
start_json_server() ->
rvi_common:start_json_rpc_server(protocol, ?MODULE, proto_msgpack_sup).


send_message(CompSpec, TID, ServiceName, Timeout, ProtoOpts,
DataLinkMod, DataLinkOpts, Parameters) ->
send_message(CompSpec, TID, ServiceName, Timeout, ProtoOpts,
DataLinkMod, DataLinkOpts, [], Parameters).

send_message(CompSpec,
TID,
Expand All @@ -53,6 +56,7 @@ send_message(CompSpec,
ProtoOpts,
DataLinkMod,
DataLinkOpts,
Extra,
Parameters) ->
rvi_common:request(protocol, ?MODULE, send_message,
[{ transaction_id, TID },
Expand All @@ -61,6 +65,7 @@ send_message(CompSpec,
{ protocol_opts, ProtoOpts },
{ data_link_mod, DataLinkMod },
{ data_link_opts, DataLinkOpts },
{ extra, Extra },
{ parameters, Parameters }],
[ status ], CompSpec).

Expand All @@ -82,6 +87,7 @@ handle_rpc(<<"send_message">>, Args) ->
{ok, ProtoOpts} = rvi_common:get_json_element(["protocol_opts"], Args),
{ok, DataLinkMod} = rvi_common:get_json_element(["data_link_mod"], Args),
{ok, DataLinkOpts} = rvi_common:get_json_element(["data_link_opts"], Args),
Extra = rvi_common:get_opt_json_element(["extra"], [], Args),
{ok, Parameters} = rvi_common:get_json_element(["parameters"], Args),
[ ok ] = gen_server:call(?SERVER, { rvi, send_message,
[TID,
Expand All @@ -90,6 +96,7 @@ handle_rpc(<<"send_message">>, Args) ->
ProtoOpts,
DataLinkMod,
DataLinkOpts,
Extra,
Parameters,
LogId]}),
{ok, [ {status, rvi_common:json_rpc_status(ok)} ]};
Expand Down Expand Up @@ -123,18 +130,21 @@ handle_call({rvi, send_message,
ProtoOpts,
DataLinkMod,
DataLinkOpts,
Extra,
Parameters
| LogId]}, _From, St) ->
| _LogId]}, _From, St) ->
?debug(" protocol:send(): transaction id: ~p~n", [TID]),
?debug(" protocol:send(): service name: ~p~n", [ServiceName]),
?debug(" protocol:send(): timeout: ~p~n", [Timeout]),
?debug(" protocol:send(): opts: ~p~n", [ProtoOpts]),
?debug(" protocol:send(): data_link_mod: ~p~n", [DataLinkMod]),
?debug(" protocol:send(): data_link_opts: ~p~n", [DataLinkOpts]),
?debug(" protocol:send(): extra: ~p~n", [Extra]),
?debug(" protocol:send(): parameters: ~p~n", [Parameters]),
Data = [ { <<"service">>, ServiceName },
{ <<"timeout">>, Timeout },
{ <<"parameters">>, Parameters } ],
{ <<"parameters">>, Parameters }
| Extra ],
RviOpts = rvi_common:rvi_options(Parameters),
Res = DataLinkMod:send_data(
St#st.cs, ?MODULE, ServiceName, RviOpts ++ DataLinkOpts, Data),
Expand All @@ -146,11 +156,11 @@ handle_call(Other, _From, St) ->


%% Convert list-based data to binary.
handle_cast({rvi, receive_message, [Elems, IP, Port | LogId]} = Msg, St) ->
handle_cast({rvi, receive_message, [Elems, IP, Port | _LogId]} = Msg, St) ->
?debug("~p:handle_cast(~p)", [?MODULE, Msg]),

[ ServiceName, Timeout, Parameters ] =
opts([<<"service">>, <<"timeout">>, <<"parameters">>],
[ ServiceName, Timeout, Src, Parameters ] =
opts([<<"service">>, <<"timeout">>, <<"src">>, <<"parameters">>],
Elems, undefined),

?debug(" protocol:rcv(): service name: ~p~n", [ServiceName]),
Expand All @@ -161,6 +171,7 @@ handle_cast({rvi, receive_message, [Elems, IP, Port | LogId]} = Msg, St) ->
{IP, Port},
ServiceName,
Timeout,
[{<<"src">>, Src}],
Parameters),
{noreply, St};

Expand Down
17 changes: 14 additions & 3 deletions components/schedule/src/schedule_rpc.erl
Original file line number Diff line number Diff line change
Expand Up @@ -13,7 +13,7 @@
-include_lib("rvi_common/include/rvi_common.hrl").
%% API
-export([start_link/0]).
-export([schedule_message/4]).
-export([schedule_message/4, schedule_message/5]).

%% Invoked by service discovery
%% FIXME: Should be rvi_service_discovery behavior
Expand Down Expand Up @@ -64,6 +64,7 @@
routes, %% Routes retrieved for this
timeout_tref, %% Reference to erlang timer associated with this message.
log_id,
extra = [],
parameters
}).

Expand Down Expand Up @@ -120,20 +121,24 @@ init([]) ->
start_json_server() ->
rvi_common:start_json_rpc_server(schedule, ?MODULE, schedule_sup).

schedule_message(CompSpec, SvcName, Timeout, Parameters) ->
schedule_message(CompSpec, SvcName, Timeout, [], Parameters).

schedule_message(CompSpec,
SvcName,
Timeout,
Extra,
Parameters) ->

rvi_common:request(schedule, ?MODULE,
schedule_message,
[{ service, SvcName },
{ timeout, Timeout },
{ extra, Extra },
{ parameters, Parameters }],
[status, transaction_id], CompSpec).



service_available(CompSpec, SvcName, DataLinkModule) ->

rvi_common:notification(schedule, ?MODULE,
Expand All @@ -156,6 +161,7 @@ handle_rpc(<<"schedule_message">>, Args) ->
{ok, SvcName} = rvi_common:get_json_element(["service"], Args),
{ok, Timeout} = rvi_common:get_json_element(["timeout"], Args),
{ok, Parameters} = rvi_common:get_json_element(["parameters"], Args),
Extra = rvi_common:get_opt_json_element(["extra"], [], Args),
LogId = rvi_common:get_json_log_id(Args),

?debug("schedule_rpc:schedule_request(): service: ~p", [ SvcName]),
Expand All @@ -165,6 +171,7 @@ handle_rpc(<<"schedule_message">>, Args) ->
[ok, TransID] = gen_server:call(?SERVER, { rvi, schedule_message,
[ SvcName,
Timeout,
Extra,
Parameters,
LogId ]}),

Expand Down Expand Up @@ -208,6 +215,7 @@ handle_notification(Other, _Args) ->
handle_call( { rvi, schedule_message,
[SvcName,
Timeout,
Extra,
Parameters | LogId] }, _From, St) ->

?debug("sched:sched_msg(): service: ~p", [SvcName]),
Expand All @@ -222,6 +230,7 @@ handle_call( { rvi, schedule_message,
Msg = #message{transaction_id = TransID,
service = SvcName,
timeout = Timeout,
extra = Extra,
parameters = Parameters,
log_id = LogId},
{_, NSt2 }= queue_message(Msg,
Expand Down Expand Up @@ -495,6 +504,7 @@ send_message(local, _, _, _, Msg, St) ->
service_edge_rpc:handle_remote_message(St#st.cs,
Msg#message.service,
Msg#message.timeout,
Msg#message.extra,
Msg#message.parameters),
{ok, St};

Expand All @@ -514,6 +524,7 @@ send_message(DataLinkMod, DataLinkOpts,
ProtoOpts,
DataLinkMod,
DataLinkOpts,
Msg#message.extra,
Msg#message.parameters) of

%% Success
Expand Down Expand Up @@ -755,7 +766,7 @@ create_transaction_id(St) ->
%% Calculate a relative timeout based on the Msec UnixTime TS we are
%% provided with.
calc_relative_tout(UnixTimeMS) ->
{ Mega, Sec, Micro } = now(),
{ Mega, Sec, Micro } = os:timestamp(),
Now = Mega * 1000000000 + Sec * 1000 + trunc(Micro / 1000) ,
?debug("sched:calc_relative_tout(): TimeoutUnixMS(~p) - Now(~p) = ~p",
[ UnixTimeMS, Now, UnixTimeMS - Now ]),
Expand Down
Loading