diff --git a/apps/efka/src/iot/efka_iot_client.erl b/apps/efka/src/iot/efka_iot_client.erl index 565632b..958e9ed 100644 --- a/apps/efka/src/iot/efka_iot_client.erl +++ b/apps/efka/src/iot/efka_iot_client.erl @@ -452,11 +452,11 @@ handle_container_command(Ref, Request, Socket) -> 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) -> +handle_stream_frame(StreamId, <<"open">>, State = #state{streams = Streams}) + when is_integer(StreamId), StreamId > 0 -> 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), + {ok, {WorkerPid, MonitorRef}} = efka_iot_stream:start_stream(StreamId), StreamState = #stream_state{worker_pid = WorkerPid, monitor_ref = MonitorRef}, {keep_state, State#state{streams = maps:put(StreamId, StreamState, Streams)}}; false -> diff --git a/apps/efka/src/iot/efka_iot_stream.erl b/apps/efka/src/iot/efka_iot_stream.erl index af8dbed..1d7f768 100644 --- a/apps/efka/src/iot/efka_iot_stream.erl +++ b/apps/efka/src/iot/efka_iot_stream.erl @@ -6,28 +6,21 @@ %%%------------------------------------------------------------------- -module(efka_iot_stream). --export([start_stream/2]). --export([run/2]). +-export([start_stream/1]). +-export([run/1]). -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()) -> {ok, {pid(), reference()}}. +start_stream(StreamId) when is_integer(StreamId), StreamId > 0 -> + {ok, spawn_monitor(?MODULE, run, [StreamId])}. --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 +-spec run(StreamId :: stream_id()) -> ok. +run(StreamId) when is_integer(StreamId), StreamId > 0 -> + try run0(StreamId) of ok -> ok catch @@ -40,9 +33,9 @@ run(StreamId, Params) when is_integer(StreamId), StreamId > 0, is_map(Params) -> efka_iot_client:stream_done(StreamId) end. --spec run0(stream_id(), map()) -> ok. -run0(StreamId, Params) -> - case open_target_socket(Params) of +-spec run0(stream_id()) -> ok. +run0(StreamId) -> + case open_target_socket() of {ok, Socket, IdleTimeout} -> efka_iot_client:send_stream(StreamId, <<"opened">>), ok = inet:setopts(Socket, [{active, once}]), @@ -52,46 +45,26 @@ run0(StreamId, Params) -> 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 open_target_socket() -> {ok, gen_tcp:socket(), timeout()} | {error, term()}. +open_target_socket() -> + {ok, Props} = application:get_env(efka, stream_target), + Host = proplists:get_value(host, Props), + Port = proplists:get_value(port, Props), + ConnectTimeout = proplists:get_value(connect_timeout, Props, ?DEFAULT_CONNECT_TIMEOUT), + IdleTimeout = proplists:get_value(idle_timeout, Props, ?DEFAULT_IDLE_TIMEOUT), --spec stream_target_options(map()) -> error | {ok, #options{}}. -stream_target_options(Params = #{<<"target">> := Target}) when is_binary(Target) -> - {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; -stream_target_options(_Params) -> - error. + 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. -spec loop(stream_id(), gen_tcp:socket(), timeout()) -> ok. loop(StreamId, Socket, IdleTimeout) -> diff --git a/config/sys.config.src b/config/sys.config.src index 2ddbf85..372313c 100644 --- a/config/sys.config.src +++ b/config/sys.config.src @@ -15,13 +15,11 @@ {udp_port, 24000} ]}, - {stream_targets, [ - {<<"manager">>, [ - {host, "127.0.0.1"}, - {port, 18091}, - {connect_timeout, 3000}, - {idle_timeout, 120000} - ]} + {stream_target, [ + {host, "127.0.0.1"}, + {port, 18091}, + {connect_timeout, 3000}, + {idle_timeout, 120000} ]}, {heartbeat, [ diff --git a/docs/efka_iot_protocol.md b/docs/efka_iot_protocol.md index d7d5845..07c0990 100644 --- a/docs/efka_iot_protocol.md +++ b/docs/efka_iot_protocol.md @@ -53,7 +53,7 @@ `Body` 取值: ```erlang -{<<"open">>, #{<<"target">> => <<"manager">>}} +<<"open">> <<"opened">> {<<"open_error">>, Reason} {<<"data">>, Chunk} @@ -63,7 +63,7 @@ 语义: -- `open`:发起方请求打开一个透明 TCP stream。当前 `efka` 端支持 `target = <<"manager">>`,由本地配置决定实际本机管理程序地址。 +- `open`:发起方请求打开一个透明 TCP stream。`open` 不携带参数;`efka` 端固定连接本机 `stream_target`,通常指向本机 manager/nginx 入口。 - `opened`:本地 TCP 连接建立成功。 - `open_error`:本地 TCP 连接建立失败,stream 结束。 - `data`:透明字节数据,`Chunk` 是 binary。HTTP 请求头、请求体、响应头、响应体都只是该 binary 的内容,iot/efka 中间协议不解析 HTTP。 @@ -73,7 +73,7 @@ 典型 iot 发起访问 efka 本机监控服务的流程: ```erlang -iot -> efka: {<<"stream">>, 1, {<<"open">>, #{<<"target">> => <<"manager">>}}} +iot -> efka: {<<"stream">>, 1, <<"open">>} efka -> iot: {<<"stream">>, 1, <<"opened">>} iot -> efka: {<<"stream">>, 1, {<<"data">>, RawHttpRequestBytes}}