support stream
This commit is contained in:
parent
59563870f7
commit
8f801fd29c
@ -15,6 +15,7 @@
|
||||
%% API
|
||||
-export([start_link/0]).
|
||||
-export([metric_data/2, ping/13, task_event_stream/3, close_task_event_stream/2]).
|
||||
-export([send_stream/2, stream_done/1]).
|
||||
-export([is_activated/0, dropped_message_count/0]).
|
||||
|
||||
%% gen_statem callbacks
|
||||
@ -36,12 +37,20 @@
|
||||
-record(state, {
|
||||
socket :: undefined | ssl:sslsocket(),
|
||||
outbox :: efka_iot_outbox:outbox(),
|
||||
streams = #{},
|
||||
%% 保存当前auth请求的ref,用来建立auth请求和响应的对应关系
|
||||
auth_ref = undefined :: undefined | binary(),
|
||||
ping_timer_ref = undefined :: undefined | reference(),
|
||||
dropped_message_count = 0 :: non_neg_integer()
|
||||
}).
|
||||
|
||||
-record(stream_state, {
|
||||
worker_pid :: pid(),
|
||||
monitor_ref :: reference()
|
||||
}).
|
||||
|
||||
-type stream_id() :: pos_integer().
|
||||
|
||||
%%%===================================================================
|
||||
%%% API
|
||||
%%%===================================================================
|
||||
@ -59,6 +68,14 @@ task_event_stream(TaskId, Type, Stream) when is_integer(TaskId), is_binary(Type)
|
||||
close_task_event_stream(TaskId, Reason) when is_integer(TaskId), is_binary(Reason) ->
|
||||
gen_statem:cast(?SERVER, {close_task_event_stream, TaskId, Reason}).
|
||||
|
||||
-spec send_stream(StreamId :: stream_id(), Body :: term()) -> ok.
|
||||
send_stream(StreamId, Body) when is_integer(StreamId), StreamId > 0 ->
|
||||
gen_statem:cast(?SERVER, {send_stream, StreamId, Body}).
|
||||
|
||||
-spec stream_done(StreamId :: stream_id()) -> ok.
|
||||
stream_done(StreamId) when is_integer(StreamId), StreamId > 0 ->
|
||||
gen_statem:cast(?SERVER, {stream_done, StreamId}).
|
||||
|
||||
-spec is_activated() -> boolean().
|
||||
is_activated() ->
|
||||
gen_statem:call(?SERVER, is_activated).
|
||||
@ -139,6 +156,24 @@ handle_event(cast, {close_task_event_stream, TaskId, Reason}, ?STATE_ACTIVATED,
|
||||
ok = ssl:send(Socket, Packet),
|
||||
{keep_state, State};
|
||||
|
||||
handle_event(cast, {send_stream, StreamId, Body}, ?STATE_ACTIVATED, State = #state{socket = Socket})
|
||||
when is_integer(StreamId), StreamId > 0 ->
|
||||
Packet = term_to_binary({<<"stream">>, StreamId, Body}),
|
||||
ok = ssl:send(Socket, Packet),
|
||||
{keep_state, State};
|
||||
handle_event(cast, {send_stream, _StreamId, _Body}, _StateName, State) ->
|
||||
{keep_state, State};
|
||||
|
||||
handle_event(cast, {stream_done, StreamId}, _StateName, State = #state{streams = Streams})
|
||||
when is_integer(StreamId), StreamId > 0 ->
|
||||
case maps:take(StreamId, Streams) of
|
||||
error ->
|
||||
{keep_state, State};
|
||||
{StreamState, NStreams} ->
|
||||
demonitor_stream(StreamState),
|
||||
{keep_state, State#state{streams = NStreams}}
|
||||
end;
|
||||
|
||||
%% 其他情况下直接忽略
|
||||
handle_event(cast, _, _, State = #state{}) ->
|
||||
{keep_state, State};
|
||||
@ -179,7 +214,8 @@ handle_event(info, {timeout, TimerRef, ssl_ping}, ?STATE_ACTIVATED, State = #sta
|
||||
logger:warning("[efka_iot_client] send ssl ping failed, reason: ~p", [Reason]),
|
||||
disconnect(Socket),
|
||||
schedule_reconnect(),
|
||||
{next_state, ?STATE_DISCONNECTED, State#state{socket = undefined, auth_ref = undefined, ping_timer_ref = undefined}}
|
||||
NState = close_all_streams({send_ping_failed, Reason}, State),
|
||||
{next_state, ?STATE_DISCONNECTED, NState#state{socket = undefined, auth_ref = undefined, ping_timer_ref = undefined}}
|
||||
end;
|
||||
handle_event(info, {timeout, _TimerRef, ssl_ping}, _StateName, State) ->
|
||||
{keep_state, State};
|
||||
@ -216,20 +252,40 @@ handle_event(info, {ssl, Socket, PacketBin}, _, State = #state{socket = Socket})
|
||||
disconnect(Socket),
|
||||
cancel_ssl_ping(State),
|
||||
schedule_reconnect(),
|
||||
{next_state, ?STATE_DISCONNECTED, State#state{socket = undefined, auth_ref = undefined, ping_timer_ref = undefined}}
|
||||
NState = close_all_streams(bad_packet, State),
|
||||
{next_state, ?STATE_DISCONNECTED, NState#state{socket = undefined, auth_ref = undefined, ping_timer_ref = undefined}}
|
||||
end;
|
||||
handle_event(info, {ssl_error, Socket, Reason}, _, State = #state{}) ->
|
||||
logger:debug("[efka_iot_client] ssl error: ~p", [Reason]),
|
||||
disconnect(Socket),
|
||||
cancel_ssl_ping(State),
|
||||
schedule_reconnect(),
|
||||
{next_state, ?STATE_DISCONNECTED, State#state{socket = undefined, auth_ref = undefined, ping_timer_ref = undefined}};
|
||||
NState = close_all_streams(ssl_error, State),
|
||||
{next_state, ?STATE_DISCONNECTED, NState#state{socket = undefined, auth_ref = undefined, ping_timer_ref = undefined}};
|
||||
handle_event(info, {ssl_closed, Socket}, _, State = #state{}) ->
|
||||
logger:debug("[efka_iot_client] ssl closed"),
|
||||
disconnect(Socket),
|
||||
cancel_ssl_ping(State),
|
||||
schedule_reconnect(),
|
||||
{next_state, ?STATE_DISCONNECTED, State#state{socket = undefined, auth_ref = undefined, ping_timer_ref = undefined}};
|
||||
NState = close_all_streams(ssl_closed, State),
|
||||
{next_state, ?STATE_DISCONNECTED, NState#state{socket = undefined, auth_ref = undefined, ping_timer_ref = undefined}};
|
||||
|
||||
handle_event(info, {'DOWN', MonitorRef, process, WorkerPid, Reason}, _StateName, State = #state{streams = Streams}) ->
|
||||
case take_stream_by_monitor(MonitorRef, WorkerPid, Streams) of
|
||||
error ->
|
||||
{keep_state, State};
|
||||
{StreamId, StreamState, NStreams} ->
|
||||
demonitor_stream(StreamState),
|
||||
logger:warning("[efka_iot_client] stream worker down, stream_id: ~p, reason: ~p", [StreamId, Reason]),
|
||||
case State#state.socket of
|
||||
undefined ->
|
||||
{keep_state, State#state{streams = NStreams}};
|
||||
Socket ->
|
||||
Packet = term_to_binary({<<"stream">>, StreamId, {<<"reset">>, safe_term({worker_down, Reason})}}),
|
||||
ok = ssl:send(Socket, Packet),
|
||||
{keep_state, State#state{streams = NStreams}}
|
||||
end
|
||||
end;
|
||||
|
||||
%%% 处理内部消息,ssl收到的消息会先 binary_to_term,再由这里按协议结构模式匹配
|
||||
|
||||
@ -259,6 +315,14 @@ handle_event(internal, {<<"command_response">>, _Ref, Reply}, StateName, State)
|
||||
logger:warning("[efka_iot_client] ignore unexpected command_response in state ~p: ~p", [StateName, Reply]),
|
||||
{keep_state, State};
|
||||
|
||||
%% 透明 TCP stream 多路复用,StreamId 只出现在新顶层 <<"stream">> 帧里。
|
||||
handle_event(internal, {<<"stream">>, StreamId, Body}, ?STATE_ACTIVATED, State) ->
|
||||
handle_stream_frame(StreamId, Body, State);
|
||||
handle_event(internal, {<<"stream">>, StreamId, Body}, StateName, State) ->
|
||||
logger:warning("[efka_iot_client] ignore stream frame in state ~p, stream_id: ~p, body: ~p",
|
||||
[StateName, StreamId, Body]),
|
||||
{keep_state, State};
|
||||
|
||||
%% 处理Pub/Sub机制
|
||||
handle_event(internal, {<<"message">>, <<"pong">>}, ?STATE_ACTIVATED, State) ->
|
||||
{keep_state, State};
|
||||
@ -274,43 +338,10 @@ handle_event(info, Info, _, State = #state{}) ->
|
||||
logger:notice("[efka_iot_client] get unknown info: ~p", [Info]),
|
||||
{keep_state, State}.
|
||||
|
||||
-spec handle_container_command(binary(), term(), ssl:sslsocket()) -> ok.
|
||||
handle_container_command(Ref, #{<<"action">> := <<"list">>}, Socket) ->
|
||||
Reply = docker_commands:get_containers(),
|
||||
send_container_response(Socket, Ref, Reply),
|
||||
ok;
|
||||
handle_container_command(Ref, #{<<"action">> := <<"deploy">>, <<"task_id">> := TaskId, <<"params">> := Params}, Socket) ->
|
||||
Reply = docker_deploy_manager:deploy(TaskId, Params),
|
||||
send_container_response(Socket, Ref, Reply),
|
||||
ok;
|
||||
handle_container_command(Ref, #{<<"action">> := <<"start">>, <<"target">> := Target}, Socket) ->
|
||||
Reply = docker_commands:start_container(container_target(Target)),
|
||||
send_container_response(Socket, Ref, Reply),
|
||||
ok;
|
||||
handle_container_command(Ref, #{<<"action">> := <<"stop">>, <<"target">> := Target, <<"timeout_seconds">> := TimeoutSeconds}, Socket) ->
|
||||
Reply = docker_commands:stop_container(container_target(Target), TimeoutSeconds),
|
||||
send_container_response(Socket, Ref, Reply),
|
||||
ok;
|
||||
handle_container_command(Ref, #{<<"action">> := <<"kill">>, <<"target">> := Target, <<"signal">> := Signal}, Socket) ->
|
||||
Reply = docker_commands:kill_container(container_target(Target), to_binary(Signal)),
|
||||
send_container_response(Socket, Ref, Reply),
|
||||
ok;
|
||||
handle_container_command(Ref, #{<<"action">> := <<"remove">>, <<"target">> := Target, <<"force">> := Force, <<"remove_volumes">> := RemoveVolumes}, Socket) ->
|
||||
Reply = docker_commands:remove_container(container_target(Target), to_bool(Force), to_bool(RemoveVolumes)),
|
||||
send_container_response(Socket, Ref, Reply),
|
||||
ok;
|
||||
handle_container_command(Ref, #{<<"action">> := <<"config">>, <<"target">> := Target, <<"config">> := Config}, Socket) ->
|
||||
Reply = docker_helper:update_container_config(container_target(Target), iolist_to_binary(Config)),
|
||||
send_container_response(Socket, Ref, Reply),
|
||||
ok;
|
||||
handle_container_command(Ref, Request, Socket) ->
|
||||
logger:notice("[efka_iot_client] get an invalid command: ~p, agent invalid", [Request]),
|
||||
send_container_response(Socket, Ref, {error, <<"agent invalid">>}),
|
||||
ok.
|
||||
|
||||
-spec terminate(term(), atom(), #state{}) -> ok.
|
||||
terminate(Reason, _StateName, State = #state{socket = Socket, outbox = Outbox}) ->
|
||||
cancel_ssl_ping(State),
|
||||
_ = close_all_streams(Reason, State),
|
||||
disconnect(Socket),
|
||||
efka_iot_outbox:close(Outbox),
|
||||
logger:notice("[efka_iot_client] terminate with reason: ~p", [Reason]),
|
||||
@ -386,6 +417,114 @@ send_container_response(Socket, Ref, Reply) ->
|
||||
Packet = term_to_binary({<<"command_response">>, Ref, {<<"container">>, safe_reply(Reply)}}),
|
||||
ok = ssl:send(Socket, Packet).
|
||||
|
||||
-spec handle_container_command(binary(), term(), ssl:sslsocket()) -> ok.
|
||||
handle_container_command(Ref, #{<<"action">> := <<"list">>}, Socket) ->
|
||||
Reply = docker_commands:get_containers(),
|
||||
send_container_response(Socket, Ref, Reply),
|
||||
ok;
|
||||
handle_container_command(Ref, #{<<"action">> := <<"deploy">>, <<"task_id">> := TaskId, <<"params">> := Params}, Socket) ->
|
||||
Reply = docker_deploy_manager:deploy(TaskId, Params),
|
||||
send_container_response(Socket, Ref, Reply),
|
||||
ok;
|
||||
handle_container_command(Ref, #{<<"action">> := <<"start">>, <<"target">> := Target}, Socket) ->
|
||||
Reply = docker_commands:start_container(container_target(Target)),
|
||||
send_container_response(Socket, Ref, Reply),
|
||||
ok;
|
||||
handle_container_command(Ref, #{<<"action">> := <<"stop">>, <<"target">> := Target, <<"timeout_seconds">> := TimeoutSeconds}, Socket) ->
|
||||
Reply = docker_commands:stop_container(container_target(Target), TimeoutSeconds),
|
||||
send_container_response(Socket, Ref, Reply),
|
||||
ok;
|
||||
handle_container_command(Ref, #{<<"action">> := <<"kill">>, <<"target">> := Target, <<"signal">> := Signal}, Socket) ->
|
||||
Reply = docker_commands:kill_container(container_target(Target), to_binary(Signal)),
|
||||
send_container_response(Socket, Ref, Reply),
|
||||
ok;
|
||||
handle_container_command(Ref, #{<<"action">> := <<"remove">>, <<"target">> := Target, <<"force">> := Force, <<"remove_volumes">> := RemoveVolumes}, Socket) ->
|
||||
Reply = docker_commands:remove_container(container_target(Target), to_bool(Force), to_bool(RemoveVolumes)),
|
||||
send_container_response(Socket, Ref, Reply),
|
||||
ok;
|
||||
handle_container_command(Ref, #{<<"action">> := <<"config">>, <<"target">> := Target, <<"config">> := Config}, Socket) ->
|
||||
Reply = docker_helper:update_container_config(container_target(Target), iolist_to_binary(Config)),
|
||||
send_container_response(Socket, Ref, Reply),
|
||||
ok;
|
||||
handle_container_command(Ref, Request, Socket) ->
|
||||
logger:notice("[efka_iot_client] get an invalid command: ~p, agent invalid", [Request]),
|
||||
send_container_response(Socket, Ref, {error, <<"agent invalid">>}),
|
||||
ok.
|
||||
|
||||
-spec handle_stream_frame(term(), term(), #state{}) -> gen_statem:event_handler_result(atom(), #state{}).
|
||||
handle_stream_frame(StreamId, {<<"open">>, Params}, State = #state{streams = Streams})
|
||||
when is_integer(StreamId), StreamId > 0, is_map(Params) ->
|
||||
case valid_iot_stream_id(StreamId) andalso not maps:is_key(StreamId, Streams) of
|
||||
true ->
|
||||
{ok, {WorkerPid, MonitorRef}} = efka_iot_stream:start_stream(StreamId, Params),
|
||||
StreamState = #stream_state{worker_pid = WorkerPid, monitor_ref = MonitorRef},
|
||||
{keep_state, State#state{streams = maps:put(StreamId, StreamState, Streams)}};
|
||||
false ->
|
||||
send_stream(StreamId, {<<"reset">>, <<"invalid stream open">>}),
|
||||
{keep_state, State}
|
||||
end;
|
||||
handle_stream_frame(StreamId, Body, State = #state{streams = Streams})
|
||||
when is_integer(StreamId), StreamId > 0 ->
|
||||
case maps:get(StreamId, Streams, undefined) of
|
||||
undefined ->
|
||||
maybe_reset_unknown_stream(StreamId, Body),
|
||||
{keep_state, State};
|
||||
StreamState = #stream_state{worker_pid = WorkerPid} ->
|
||||
WorkerPid ! {stream, StreamId, Body},
|
||||
case Body of
|
||||
{<<"reset">>, _Reason} ->
|
||||
demonitor_stream(StreamState),
|
||||
{keep_state, State#state{streams = maps:remove(StreamId, Streams)}};
|
||||
_ ->
|
||||
{keep_state, State}
|
||||
end
|
||||
end;
|
||||
handle_stream_frame(StreamId, Body, State) ->
|
||||
logger:warning("[efka_iot_client] invalid stream frame, stream_id: ~p, body: ~p", [StreamId, Body]),
|
||||
{keep_state, State}.
|
||||
|
||||
-spec maybe_reset_unknown_stream(stream_id(), term()) -> ok.
|
||||
maybe_reset_unknown_stream(_StreamId, {<<"reset">>, _Reason}) ->
|
||||
ok;
|
||||
maybe_reset_unknown_stream(_StreamId, <<"fin">>) ->
|
||||
ok;
|
||||
maybe_reset_unknown_stream(_StreamId, {<<"open_error">>, _Reason}) ->
|
||||
ok;
|
||||
maybe_reset_unknown_stream(StreamId, _Body) ->
|
||||
send_stream(StreamId, {<<"reset">>, <<"unknown stream">>}).
|
||||
|
||||
-spec valid_iot_stream_id(stream_id()) -> boolean().
|
||||
valid_iot_stream_id(StreamId) ->
|
||||
StreamId rem 2 =:= 1.
|
||||
|
||||
-spec take_stream_by_monitor(reference(), pid(), map()) ->
|
||||
{stream_id(), #stream_state{}, map()} | error.
|
||||
take_stream_by_monitor(MonitorRef, WorkerPid, Streams) ->
|
||||
take_stream_by_monitor(MonitorRef, WorkerPid, maps:iterator(Streams), Streams).
|
||||
|
||||
take_stream_by_monitor(MonitorRef, WorkerPid, Iter, Streams) ->
|
||||
case maps:next(Iter) of
|
||||
none ->
|
||||
error;
|
||||
{StreamId, StreamState = #stream_state{worker_pid = WorkerPid, monitor_ref = MonitorRef}, _NextIter} ->
|
||||
{StreamId, StreamState, maps:remove(StreamId, Streams)};
|
||||
{_StreamId, _StreamState, NextIter} ->
|
||||
take_stream_by_monitor(MonitorRef, WorkerPid, NextIter, Streams)
|
||||
end.
|
||||
|
||||
-spec close_all_streams(term(), #state{}) -> #state{}.
|
||||
close_all_streams(Reason, State = #state{streams = Streams}) ->
|
||||
maps:foreach(fun(StreamId, StreamState = #stream_state{worker_pid = WorkerPid}) ->
|
||||
demonitor_stream(StreamState),
|
||||
WorkerPid ! {stream, StreamId, {<<"reset">>, {<<"channel_closed">>, safe_term(Reason)}}}
|
||||
end, Streams),
|
||||
State#state{streams = #{}}.
|
||||
|
||||
-spec demonitor_stream(#stream_state{}) -> ok.
|
||||
demonitor_stream(#stream_state{monitor_ref = MonitorRef}) ->
|
||||
erlang:demonitor(MonitorRef, [flush]),
|
||||
ok.
|
||||
|
||||
-spec safe_reply(term()) -> term().
|
||||
safe_reply(ok) ->
|
||||
<<"ok">>;
|
||||
|
||||
150
apps/efka/src/iot/efka_iot_stream.erl
Normal file
150
apps/efka/src/iot/efka_iot_stream.erl
Normal file
@ -0,0 +1,150 @@
|
||||
%%%-------------------------------------------------------------------
|
||||
%%% @doc One transparent TCP stream from iot to the local manager service.
|
||||
%%% The iot/efka protocol does not inspect HTTP; all HTTP bytes are carried
|
||||
%%% in <<"data">> frames.
|
||||
%%% @end
|
||||
%%%-------------------------------------------------------------------
|
||||
-module(efka_iot_stream).
|
||||
|
||||
-export([start_stream/2]).
|
||||
-export([run/2]).
|
||||
|
||||
-define(DEFAULT_CONNECT_TIMEOUT, 3000).
|
||||
-define(DEFAULT_IDLE_TIMEOUT, 120000).
|
||||
|
||||
-type stream_id() :: pos_integer().
|
||||
|
||||
-record(options, {
|
||||
host :: inet:hostname() | inet:ip_address(),
|
||||
port :: inet:port_number(),
|
||||
connect_timeout = 0 :: timeout(),
|
||||
idle_timeout = 0 :: timeout()
|
||||
}).
|
||||
|
||||
-spec start_stream(StreamId :: integer(), Params :: map()) -> {ok, {pid(), reference()}}.
|
||||
start_stream(StreamId, Params) when is_integer(StreamId), StreamId > 0, is_map(Params) ->
|
||||
{ok, spawn_monitor(?MODULE, run, [StreamId, Params])}.
|
||||
|
||||
-spec run(StreamId :: stream_id(), Params :: map()) -> ok.
|
||||
run(StreamId, Params) when is_integer(StreamId), StreamId > 0, is_map(Params) ->
|
||||
try run0(StreamId, Params) of
|
||||
ok ->
|
||||
ok
|
||||
catch
|
||||
Class:Reason:Stack ->
|
||||
logger:warning("[efka_iot_stream] stream_id: ~p crashed, class: ~p, reason: ~p, stack: ~p",
|
||||
[StreamId, Class, Reason, Stack]),
|
||||
efka_iot_client:send_stream(StreamId, {<<"reset">>, safe_term({Class, Reason})}),
|
||||
ok
|
||||
after
|
||||
efka_iot_client:stream_done(StreamId)
|
||||
end.
|
||||
|
||||
-spec run0(stream_id(), map()) -> ok.
|
||||
run0(StreamId, Params) ->
|
||||
case open_target_socket(Params) of
|
||||
{ok, Socket, IdleTimeout} ->
|
||||
efka_iot_client:send_stream(StreamId, <<"opened">>),
|
||||
ok = inet:setopts(Socket, [{active, once}]),
|
||||
loop(StreamId, Socket, IdleTimeout);
|
||||
{error, Reason} ->
|
||||
efka_iot_client:send_stream(StreamId, {<<"open_error">>, safe_term(Reason)}),
|
||||
ok
|
||||
end.
|
||||
|
||||
-spec open_target_socket(map()) -> {ok, gen_tcp:socket(), timeout()} | {error, term()}.
|
||||
open_target_socket(Params) when is_map(Params) ->
|
||||
case stream_target_options(Params) of
|
||||
error ->
|
||||
{error, <<"invalid_target">>};
|
||||
{ok, #options{host = Host, port = Port, connect_timeout = ConnectTimeout, idle_timeout = IdleTimeout}} ->
|
||||
SocketOpts = [
|
||||
binary,
|
||||
{packet, raw},
|
||||
{active, false},
|
||||
{nodelay, true}
|
||||
],
|
||||
case gen_tcp:connect(Host, Port, SocketOpts, ConnectTimeout) of
|
||||
{ok, Socket} ->
|
||||
{ok, Socket, IdleTimeout};
|
||||
{error, Reason} ->
|
||||
{error, Reason}
|
||||
end
|
||||
end.
|
||||
|
||||
-spec stream_target_options(map()) -> error | {ok, #options{}}.
|
||||
stream_target_options(Params = #{<<"target">> := Target0}) ->
|
||||
Target = binary_to_atom(Target0),
|
||||
{ok, StreamTargets} = application:get_env(efka, stream_targets),
|
||||
case proplists:get_value(Target, StreamTargets, undefined) of
|
||||
undefined ->
|
||||
error;
|
||||
Props ->
|
||||
Host = proplists:get_value(host, Props),
|
||||
Port = proplists:get_value(port, Props),
|
||||
ConnectTimeout0 = proplists:get_value(connect_timeout, Props, ?DEFAULT_CONNECT_TIMEOUT),
|
||||
IdleTimeout0 = proplists:get_value(idle_timeout, Props, ?DEFAULT_IDLE_TIMEOUT),
|
||||
{ok, #options{
|
||||
host = Host,
|
||||
port = Port,
|
||||
connect_timeout = maps:get(<<"connect_timeout">>, Params, ConnectTimeout0),
|
||||
idle_timeout = maps:get(<<"idle_timeout">>, Params, IdleTimeout0)
|
||||
}}
|
||||
end.
|
||||
|
||||
-spec loop(stream_id(), gen_tcp:socket(), timeout()) -> ok.
|
||||
loop(StreamId, Socket, IdleTimeout) ->
|
||||
receive
|
||||
{stream, StreamId, {<<"data">>, Data}} when is_binary(Data) ->
|
||||
case gen_tcp:send(Socket, Data) of
|
||||
ok ->
|
||||
loop(StreamId, Socket, IdleTimeout);
|
||||
{error, Reason} ->
|
||||
efka_iot_client:send_stream(StreamId, {<<"reset">>, safe_term(Reason)}),
|
||||
close_socket(Socket)
|
||||
end;
|
||||
{stream, StreamId, <<"fin">>} ->
|
||||
_ = gen_tcp:shutdown(Socket, write),
|
||||
loop(StreamId, Socket, IdleTimeout);
|
||||
{stream, StreamId, {<<"reset">>, _Reason}} ->
|
||||
close_socket(Socket);
|
||||
{tcp, Socket, Data} ->
|
||||
efka_iot_client:send_stream(StreamId, {<<"data">>, Data}),
|
||||
ok = inet:setopts(Socket, [{active, once}]),
|
||||
loop(StreamId, Socket, IdleTimeout);
|
||||
{tcp_closed, Socket} ->
|
||||
efka_iot_client:send_stream(StreamId, <<"fin">>),
|
||||
close_socket(Socket);
|
||||
{tcp_error, Socket, Reason} ->
|
||||
efka_iot_client:send_stream(StreamId, {<<"reset">>, safe_term(Reason)}),
|
||||
close_socket(Socket);
|
||||
Info ->
|
||||
logger:debug("[efka_iot_stream] stream_id: ~p ignore unknown info: ~p", [StreamId, Info]),
|
||||
loop(StreamId, Socket, IdleTimeout)
|
||||
after IdleTimeout ->
|
||||
efka_iot_client:send_stream(StreamId, {<<"reset">>, <<"idle_timeout">>}),
|
||||
close_socket(Socket)
|
||||
end.
|
||||
|
||||
-spec close_socket(gen_tcp:socket()) -> ok.
|
||||
close_socket(Socket) ->
|
||||
catch gen_tcp:close(Socket),
|
||||
ok.
|
||||
|
||||
-spec safe_term(term()) -> term().
|
||||
safe_term(true) ->
|
||||
true;
|
||||
safe_term(false) ->
|
||||
false;
|
||||
safe_term(undefined) ->
|
||||
undefined;
|
||||
safe_term(Value) when is_atom(Value) ->
|
||||
atom_to_binary(Value, utf8);
|
||||
safe_term(Value) when is_map(Value) ->
|
||||
maps:from_list([{safe_term(K), safe_term(V)} || {K, V} <- maps:to_list(Value)]);
|
||||
safe_term(Value) when is_list(Value) ->
|
||||
[safe_term(Item) || Item <- Value];
|
||||
safe_term(Value) when is_tuple(Value) ->
|
||||
list_to_tuple([safe_term(Item) || Item <- tuple_to_list(Value)]);
|
||||
safe_term(Value) ->
|
||||
Value.
|
||||
@ -15,6 +15,15 @@
|
||||
{udp_port, 24000}
|
||||
]},
|
||||
|
||||
{stream_targets, [
|
||||
{manager, [
|
||||
{host, "127.0.0.1"},
|
||||
{port, 18091},
|
||||
{connect_timeout, 3000},
|
||||
{idle_timeout, 120000}
|
||||
]}
|
||||
]},
|
||||
|
||||
{heartbeat, [
|
||||
{interval, 5000}
|
||||
]},
|
||||
|
||||
@ -21,6 +21,7 @@
|
||||
{<<"command">>, Ref, {Domain, Payload}}
|
||||
{<<"command_response">>, Ref, {Domain, Reply}}
|
||||
{<<"message">>, Body}
|
||||
{<<"stream">>, StreamId, Body}
|
||||
```
|
||||
|
||||
语义说明:
|
||||
@ -32,11 +33,56 @@
|
||||
| `{<<"command">>, Ref, {Domain, Payload}}` | iot -> efka | iot 下发命令,需要 efka 回复 |
|
||||
| `{<<"command_response">>, Ref, {Domain, Reply}}` | efka -> iot | efka 对 iot command 的回复 |
|
||||
| `{<<"message">>, Body}` | 双向 | 异步消息,不要求回复 |
|
||||
| `{<<"stream">>, StreamId, Body}` | 双向 | 透明 TCP 字节流多路复用 |
|
||||
|
||||
`command` 和 `command_response` 的 `Domain` 表示业务域,目前支持:
|
||||
|
||||
- `<<"container">>`
|
||||
|
||||
## 透明 TCP Stream
|
||||
|
||||
`stream` 用于在同一条 TLS 长连接上复用多条透明 TCP 字节流。旧的 `request`、`response`、`command`、`command_response`、`message` 帧格式保持不变;`StreamId = 0` 只作为保留概念,不出现在新的 `stream` 帧里。
|
||||
|
||||
`StreamId` 规则:
|
||||
|
||||
- `0`:保留给现有控制流。
|
||||
- 奇数:`iot` 发起。
|
||||
- 偶数:`efka` 发起,当前预留。
|
||||
- 同一条 TLS 连接内 `StreamId` 不复用。
|
||||
|
||||
`Body` 取值:
|
||||
|
||||
```erlang
|
||||
{<<"open">>, #{<<"target">> => <<"manager">>}}
|
||||
<<"opened">>
|
||||
{<<"open_error">>, Reason}
|
||||
{<<"data">>, Chunk}
|
||||
<<"fin">>
|
||||
{<<"reset">>, Reason}
|
||||
```
|
||||
|
||||
语义:
|
||||
|
||||
- `open`:发起方请求打开一个透明 TCP stream。当前 `efka` 端支持 `target = <<"manager">>`,由本地配置决定实际本机管理程序地址。
|
||||
- `opened`:本地 TCP 连接建立成功。
|
||||
- `open_error`:本地 TCP 连接建立失败,stream 结束。
|
||||
- `data`:透明字节数据,`Chunk` 是 binary。HTTP 请求头、请求体、响应头、响应体都只是该 binary 的内容,iot/efka 中间协议不解析 HTTP。
|
||||
- `fin`:半关闭,发送方后续不再发送 `data`,但仍可继续接收对端 `data`。
|
||||
- `reset`:异常关闭,双方应立即释放该 `StreamId` 的资源。
|
||||
|
||||
典型 iot 发起访问 efka 本机监控服务的流程:
|
||||
|
||||
```erlang
|
||||
iot -> efka: {<<"stream">>, 1, {<<"open">>, #{<<"target">> => <<"manager">>}}}
|
||||
efka -> iot: {<<"stream">>, 1, <<"opened">>}
|
||||
|
||||
iot -> efka: {<<"stream">>, 1, {<<"data">>, RawHttpRequestBytes}}
|
||||
iot -> efka: {<<"stream">>, 1, <<"fin">>}
|
||||
|
||||
efka -> iot: {<<"stream">>, 1, {<<"data">>, RawHttpResponseBytes}}
|
||||
efka -> iot: {<<"stream">>, 1, <<"fin">>}
|
||||
```
|
||||
|
||||
## 鉴权请求
|
||||
|
||||
初始连接由 `efka` 发起鉴权 request。每条 TLS 连接只允许一次鉴权;`iot` 侧鉴权成功后会在 `ssl_channel` 标记该连接已鉴权,如果同一连接再次发送 `auth_request`,`iot` 会直接关闭连接。
|
||||
|
||||
Loading…
x
Reference in New Issue
Block a user