fix stream

This commit is contained in:
anlicheng 2026-07-06 15:13:14 +08:00
parent 2c1af8a208
commit 6e67df6f82
4 changed files with 41 additions and 70 deletions

View File

@ -452,11 +452,11 @@ handle_container_command(Ref, Request, Socket) ->
ok. ok.
-spec handle_stream_frame(term(), term(), #state{}) -> gen_statem:event_handler_result(atom(), #state{}). -spec handle_stream_frame(term(), term(), #state{}) -> gen_statem:event_handler_result(atom(), #state{}).
handle_stream_frame(StreamId, {<<"open">>, Params}, State = #state{streams = Streams}) handle_stream_frame(StreamId, <<"open">>, State = #state{streams = Streams})
when is_integer(StreamId), StreamId > 0, is_map(Params) -> when is_integer(StreamId), StreamId > 0 ->
case valid_iot_stream_id(StreamId) andalso not maps:is_key(StreamId, Streams) of case valid_iot_stream_id(StreamId) andalso not maps:is_key(StreamId, Streams) of
true -> 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}, StreamState = #stream_state{worker_pid = WorkerPid, monitor_ref = MonitorRef},
{keep_state, State#state{streams = maps:put(StreamId, StreamState, Streams)}}; {keep_state, State#state{streams = maps:put(StreamId, StreamState, Streams)}};
false -> false ->

View File

@ -6,28 +6,21 @@
%%%------------------------------------------------------------------- %%%-------------------------------------------------------------------
-module(efka_iot_stream). -module(efka_iot_stream).
-export([start_stream/2]). -export([start_stream/1]).
-export([run/2]). -export([run/1]).
-define(DEFAULT_CONNECT_TIMEOUT, 3000). -define(DEFAULT_CONNECT_TIMEOUT, 3000).
-define(DEFAULT_IDLE_TIMEOUT, 120000). -define(DEFAULT_IDLE_TIMEOUT, 120000).
-type stream_id() :: pos_integer(). -type stream_id() :: pos_integer().
-record(options, { -spec start_stream(StreamId :: integer()) -> {ok, {pid(), reference()}}.
host :: inet:hostname() | inet:ip_address(), start_stream(StreamId) when is_integer(StreamId), StreamId > 0 ->
port :: inet:port_number(), {ok, spawn_monitor(?MODULE, run, [StreamId])}.
connect_timeout = 0 :: timeout(),
idle_timeout = 0 :: timeout()
}).
-spec start_stream(StreamId :: integer(), Params :: map()) -> {ok, {pid(), reference()}}. -spec run(StreamId :: stream_id()) -> ok.
start_stream(StreamId, Params) when is_integer(StreamId), StreamId > 0, is_map(Params) -> run(StreamId) when is_integer(StreamId), StreamId > 0 ->
{ok, spawn_monitor(?MODULE, run, [StreamId, Params])}. try run0(StreamId) of
-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 ->
ok ok
catch catch
@ -40,9 +33,9 @@ run(StreamId, Params) when is_integer(StreamId), StreamId > 0, is_map(Params) ->
efka_iot_client:stream_done(StreamId) efka_iot_client:stream_done(StreamId)
end. end.
-spec run0(stream_id(), map()) -> ok. -spec run0(stream_id()) -> ok.
run0(StreamId, Params) -> run0(StreamId) ->
case open_target_socket(Params) of case open_target_socket() of
{ok, Socket, IdleTimeout} -> {ok, Socket, IdleTimeout} ->
efka_iot_client:send_stream(StreamId, <<"opened">>), efka_iot_client:send_stream(StreamId, <<"opened">>),
ok = inet:setopts(Socket, [{active, once}]), ok = inet:setopts(Socket, [{active, once}]),
@ -52,46 +45,26 @@ run0(StreamId, Params) ->
ok ok
end. end.
-spec open_target_socket(map()) -> {ok, gen_tcp:socket(), timeout()} | {error, term()}. -spec open_target_socket() -> {ok, gen_tcp:socket(), timeout()} | {error, term()}.
open_target_socket(Params) when is_map(Params) -> open_target_socket() ->
case stream_target_options(Params) of {ok, Props} = application:get_env(efka, stream_target),
error -> Host = proplists:get_value(host, Props),
{error, <<"invalid_target">>}; Port = proplists:get_value(port, Props),
{ok, #options{host = Host, port = Port, connect_timeout = ConnectTimeout, idle_timeout = IdleTimeout}} -> ConnectTimeout = proplists:get_value(connect_timeout, Props, ?DEFAULT_CONNECT_TIMEOUT),
SocketOpts = [ IdleTimeout = proplists:get_value(idle_timeout, Props, ?DEFAULT_IDLE_TIMEOUT),
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{}}. SocketOpts = [
stream_target_options(Params = #{<<"target">> := Target}) when is_binary(Target) -> binary,
{ok, StreamTargets} = application:get_env(efka, stream_targets), {packet, raw},
case proplists:get_value(Target, StreamTargets, undefined) of {active, false},
undefined -> {nodelay, true}
error; ],
Props -> case gen_tcp:connect(Host, Port, SocketOpts, ConnectTimeout) of
Host = proplists:get_value(host, Props), {ok, Socket} ->
Port = proplists:get_value(port, Props), {ok, Socket, IdleTimeout};
ConnectTimeout0 = proplists:get_value(connect_timeout, Props, ?DEFAULT_CONNECT_TIMEOUT), {error, Reason} ->
IdleTimeout0 = proplists:get_value(idle_timeout, Props, ?DEFAULT_IDLE_TIMEOUT), {error, Reason}
{ok, #options{ end.
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.
-spec loop(stream_id(), gen_tcp:socket(), timeout()) -> ok. -spec loop(stream_id(), gen_tcp:socket(), timeout()) -> ok.
loop(StreamId, Socket, IdleTimeout) -> loop(StreamId, Socket, IdleTimeout) ->

View File

@ -15,13 +15,11 @@
{udp_port, 24000} {udp_port, 24000}
]}, ]},
{stream_targets, [ {stream_target, [
{<<"manager">>, [ {host, "127.0.0.1"},
{host, "127.0.0.1"}, {port, 18091},
{port, 18091}, {connect_timeout, 3000},
{connect_timeout, 3000}, {idle_timeout, 120000}
{idle_timeout, 120000}
]}
]}, ]},
{heartbeat, [ {heartbeat, [

View File

@ -53,7 +53,7 @@
`Body` 取值: `Body` 取值:
```erlang ```erlang
{<<"open">>, #{<<"target">> => <<"manager">>}} <<"open">>
<<"opened">> <<"opened">>
{<<"open_error">>, Reason} {<<"open_error">>, Reason}
{<<"data">>, Chunk} {<<"data">>, Chunk}
@ -63,7 +63,7 @@
语义: 语义:
- `open`:发起方请求打开一个透明 TCP stream。当前 `efka` 端支持 `target = <<"manager">>`,由本地配置决定实际本机管理程序地址 - `open`:发起方请求打开一个透明 TCP stream。`open` 不携带参数;`efka` 端固定连接本机 `stream_target`,通常指向本机 manager/nginx 入口
- `opened`:本地 TCP 连接建立成功。 - `opened`:本地 TCP 连接建立成功。
- `open_error`:本地 TCP 连接建立失败stream 结束。 - `open_error`:本地 TCP 连接建立失败stream 结束。
- `data`:透明字节数据,`Chunk` 是 binary。HTTP 请求头、请求体、响应头、响应体都只是该 binary 的内容iot/efka 中间协议不解析 HTTP。 - `data`:透明字节数据,`Chunk` 是 binary。HTTP 请求头、请求体、响应头、响应体都只是该 binary 的内容iot/efka 中间协议不解析 HTTP。
@ -73,7 +73,7 @@
典型 iot 发起访问 efka 本机监控服务的流程: 典型 iot 发起访问 efka 本机监控服务的流程:
```erlang ```erlang
iot -> efka: {<<"stream">>, 1, {<<"open">>, #{<<"target">> => <<"manager">>}}} iot -> efka: {<<"stream">>, 1, <<"open">>}
efka -> iot: {<<"stream">>, 1, <<"opened">>} efka -> iot: {<<"stream">>, 1, <<"opened">>}
iot -> efka: {<<"stream">>, 1, {<<"data">>, RawHttpRequestBytes}} iot -> efka: {<<"stream">>, 1, {<<"data">>, RawHttpRequestBytes}}