ekfa/apps/efka/src/iot/efka_iot_stream.erl
2026-07-09 10:38:16 +08:00

157 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).
-include("message.hrl").
-export([start_stream/2]).
-export([run/2]).
-define(DEFAULT_CONNECT_TIMEOUT, 3000).
-define(DEFAULT_IDLE_TIMEOUT, 120000).
-type stream_id() :: pos_integer().
-type stream_target() :: pos_integer().
-spec start_stream(StreamId :: integer(), Target :: stream_target()) -> {ok, {pid(), reference()}}.
start_stream(StreamId, Target) when is_integer(StreamId), StreamId > 0, is_integer(Target), Target > 0 ->
{ok, spawn_monitor(?MODULE, run, [StreamId, Target])}.
-spec run(StreamId :: stream_id(), Target :: stream_target()) -> ok.
run(StreamId, Target) when is_integer(StreamId), StreamId > 0, is_integer(Target), Target > 0 ->
try run0(StreamId, Target) 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(), stream_target()) -> ok.
run0(StreamId, Target) ->
case open_target_socket(Target) 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(stream_target()) -> {ok, gen_tcp:socket(), timeout()} | {error, term()}.
open_target_socket(Target) ->
case stream_target_props(Target) of
{ok, Props} ->
open_target_socket(Target, Props);
{error, Reason} ->
{error, Reason}
end.
-spec open_target_socket(stream_target(), proplists:proplist()) ->
{ok, gen_tcp:socket(), timeout()} | {error, term()}.
open_target_socket(Target, Props) ->
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),
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, {connect_failed, Target, Reason}}
end.
-spec stream_target_props(stream_target()) -> {ok, proplists:proplist()} | {error, term()}.
stream_target_props(Target) ->
case application:get_env(efka, stream_targets) of
{ok, Targets} ->
case proplists:get_value(Target, Targets) of
undefined ->
{error, {unknown_stream_target, Target}};
Props ->
{ok, Props}
end;
undefined when Target =:= ?STREAM_TARGET_MANAGER ->
case application:get_env(efka, stream_target) of
{ok, Props} ->
{ok, Props};
undefined ->
{error, {unknown_stream_target, Target}}
end;
undefined ->
{error, {unknown_stream_target, Target}}
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.