ekfa/apps/efka/src/iot/efka_iot_stream.erl
2026-07-06 13:51:01 +08:00

152 lines
5.8 KiB
Erlang

%%%-------------------------------------------------------------------
%%% @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">> := 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.
-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.