311 lines
13 KiB
Erlang
311 lines
13 KiB
Erlang
%%%-------------------------------------------------------------------
|
||
%%% @author anlicheng
|
||
%%% @copyright (C) 2025, <COMPANY>
|
||
%%% @doc
|
||
%%%
|
||
%%% @end
|
||
%%% Created : 20. 4月 2026 00:00
|
||
%%%-------------------------------------------------------------------
|
||
-module(efka_client).
|
||
-author("anlicheng").
|
||
-include("protocol.hrl").
|
||
-include("message_pb.hrl").
|
||
-include("efka_tables.hrl").
|
||
|
||
-behaviour(gen_statem).
|
||
|
||
%% API
|
||
-export([start_link/0]).
|
||
-export([metric_data/2, ping/13, task_event_stream/3, close_task_event_stream/2]).
|
||
-export([is_activated/0]).
|
||
|
||
%% gen_statem callbacks
|
||
-export([init/1, handle_event/4, terminate/3, code_change/4, callback_mode/0]).
|
||
|
||
-define(SERVER, ?MODULE).
|
||
|
||
%% 标记当前agent的状态,只有在 activated 状态下才可以正常的发送数据
|
||
-define(STATE_DISCONNECTED, disconnected).
|
||
-define(STATE_AUTH, auth).
|
||
%% 不能推送消息到服务,但是可以接受服务器的部分指令
|
||
-define(STATE_RESTRICTED, restricted).
|
||
%% 激活状态下
|
||
-define(STATE_ACTIVATED, activated).
|
||
|
||
-record(state, {
|
||
socket :: undefined | ssl:sslsocket(),
|
||
next_packet_id = 1,
|
||
%% 保存当前auth请求的packet_id,用来建立auth请求和响应的对应关系
|
||
auth_packet_id = 1
|
||
}).
|
||
|
||
%%%===================================================================
|
||
%%% API
|
||
%%%===================================================================
|
||
|
||
%% 发送数据
|
||
-spec metric_data(RouteKey :: binary(), Metric :: binary()) -> no_return().
|
||
metric_data(RouteKey, Metric) when is_binary(RouteKey), is_binary(Metric) ->
|
||
gen_statem:cast(?SERVER, {metric_data, RouteKey, Metric}).
|
||
|
||
-spec task_event_stream(TaskId :: integer(), Type :: binary(), Stream :: binary()) -> no_return().
|
||
task_event_stream(TaskId, Type, Stream) when is_integer(TaskId), is_binary(Type), is_binary(Stream) ->
|
||
gen_statem:cast(?SERVER, {task_event_stream, TaskId, Type, Stream}).
|
||
|
||
-spec close_task_event_stream(TaskId :: integer(), Reason :: binary()) -> no_return().
|
||
close_task_event_stream(TaskId, Reason) when is_integer(TaskId), is_binary(Reason) ->
|
||
gen_statem:cast(?SERVER, {close_task_event_stream, TaskId, Reason}).
|
||
|
||
-spec is_activated() -> boolean().
|
||
is_activated() ->
|
||
gen_statem:call(?SERVER, is_activated).
|
||
|
||
ping(AdCode, BootTime, Province, City, EfkaVersion, KernelArch, Ips, CpuCore, CpuLoad, CpuTemperature, Disk, Memory, Interfaces) ->
|
||
gen_statem:cast(?SERVER, {ping, AdCode, BootTime, Province, City, EfkaVersion, KernelArch, Ips, CpuCore, CpuLoad, CpuTemperature, Disk, Memory, Interfaces}).
|
||
|
||
start_link() ->
|
||
gen_statem:start_link({local, ?SERVER}, ?MODULE, [], []).
|
||
|
||
%%%===================================================================
|
||
%%% gen_statem callbacks
|
||
%%%===================================================================
|
||
|
||
init([]) ->
|
||
erlang:start_timer(0, self(), create_transport),
|
||
{ok, ?STATE_DISCONNECTED, #state{socket = undefined}}.
|
||
|
||
callback_mode() ->
|
||
handle_event_function.
|
||
|
||
%% 异步发送数据, 连接存在时候直接发送;否则缓存到mnesia
|
||
handle_event(cast, {metric_data, RouteKey, Metric}, StateName, State = #state{socket = Socket}) ->
|
||
CastFrame = message_pb:encode_msg(#'CastFrame'{
|
||
body = {data, #'Data'{route_key = RouteKey, metric = Metric}}
|
||
}),
|
||
Packet = <<?FRAME_CAST, CastFrame/binary>>,
|
||
case StateName of
|
||
?STATE_ACTIVATED ->
|
||
send_packet(Socket, Packet);
|
||
_ ->
|
||
ok = cache_model:insert(Packet)
|
||
end,
|
||
{keep_state, State};
|
||
|
||
%% Task的stream流,只做实时的
|
||
handle_event(cast, {task_event_stream, TaskId, Type, Stream}, ?STATE_ACTIVATED, State = #state{socket = Socket}) ->
|
||
logger:debug("[efka_client] event_stream task_id: ~p, stream: ~ts", [TaskId, Stream]),
|
||
EventPacket = message_pb:encode_msg(#'CastFrame'{
|
||
body = {event_stream, #'TaskEventStream'{task_id = TaskId, type = Type, stream = Stream}}
|
||
}),
|
||
send_packet(Socket, <<?FRAME_CAST, EventPacket/binary>>),
|
||
{keep_state, State};
|
||
|
||
handle_event(cast, {close_task_event_stream, TaskId, Reason}, ?STATE_ACTIVATED, State = #state{socket = Socket}) ->
|
||
EventPacket = message_pb:encode_msg(#'CastFrame'{
|
||
body = {event_stream, #'TaskEventStream'{task_id = TaskId, type = <<"close">>, stream = Reason}}
|
||
}),
|
||
send_packet(Socket, <<?FRAME_CAST, EventPacket/binary>>),
|
||
{keep_state, State};
|
||
|
||
%% 其他情况下直接忽略
|
||
handle_event(cast, _, _, State = #state{}) ->
|
||
{keep_state, State};
|
||
|
||
handle_event({call, From}, is_activated, ?STATE_ACTIVATED, State = #state{}) ->
|
||
{keep_state, State, [{reply, From, true}]};
|
||
handle_event({call, From}, is_activated, _StateName, State = #state{}) ->
|
||
{keep_state, State, [{reply, From, false}]};
|
||
|
||
%% 异步建立到服务器的连接
|
||
handle_event(info, {timeout, _, create_transport}, ?STATE_DISCONNECTED, State = #state{next_packet_id = PacketId}) ->
|
||
case connect_socket() of
|
||
{ok, Socket} ->
|
||
AuthPacket = auth_packet(PacketId),
|
||
send_packet(Socket, AuthPacket),
|
||
{next_state, ?STATE_AUTH, State#state{socket = Socket, auth_packet_id = PacketId, next_packet_id = PacketId + 1},
|
||
[{state_timeout, 5000, auth_timeout}]};
|
||
{error, Reason} ->
|
||
logger:debug("[efka_client] connect failed, error: ~p", [Reason]),
|
||
schedule_reconnect(),
|
||
{keep_state, State#state{socket = undefined}}
|
||
end;
|
||
|
||
handle_event(state_timeout, auth_timeout, ?STATE_AUTH, State = #state{socket = Socket}) ->
|
||
logger:debug("[efka_client] auth request timeout"),
|
||
disconnect(Socket),
|
||
schedule_reconnect(),
|
||
{next_state, ?STATE_DISCONNECTED, State#state{socket = undefined}};
|
||
|
||
%% 将缓存中的数据推送到服务器端
|
||
handle_event(info, flush_cache, ?STATE_ACTIVATED, State = #state{socket = Socket}) ->
|
||
case cache_model:fetch_next() of
|
||
{ok, {Id, Packet}} ->
|
||
send_packet(Socket, Packet),
|
||
cache_model:delete(Id),
|
||
{keep_state, State, [{next_event, info, flush_cache}]};
|
||
error ->
|
||
{keep_state, State}
|
||
end;
|
||
handle_event(info, flush_cache, _, State) ->
|
||
{keep_state, State};
|
||
|
||
%% 处理收到的ssl消息
|
||
handle_event(info, {ssl, Socket, <<?FRAME_RESPONSE, PacketBin/binary>>}, _, State = #state{socket = Socket}) ->
|
||
ReplyFrame = message_pb:decode_msg(PacketBin, 'ReplyFrame'),
|
||
{keep_state, State, [{next_event, internal, {decoded_reply, ReplyFrame}}]};
|
||
handle_event(info, {ssl, Socket, <<?FRAME_REQUEST, PacketBin/binary>>}, _, State = #state{socket = Socket}) ->
|
||
RequestFrame = message_pb:decode_msg(PacketBin, 'RequestFrame'),
|
||
{keep_state, State, [{next_event, internal, {decoded_request, RequestFrame}}]};
|
||
handle_event(info, {ssl, Socket, <<?FRAME_CAST, PacketBin/binary>>}, _, State = #state{socket = Socket}) ->
|
||
CastFrame = message_pb:decode_msg(PacketBin, 'CastFrame'),
|
||
{keep_state, State, [{next_event, internal, {decoded_cast, CastFrame}}]};
|
||
handle_event(info, {ssl_error, Socket, Reason}, _, State = #state{socket = Socket}) ->
|
||
logger:debug("[efka_client] ssl error: ~p", [Reason]),
|
||
disconnect(Socket),
|
||
schedule_reconnect(),
|
||
{next_state, ?STATE_DISCONNECTED, State#state{socket = undefined}};
|
||
handle_event(info, {ssl_closed, Socket}, _, State = #state{socket = Socket}) ->
|
||
schedule_reconnect(),
|
||
{next_state, ?STATE_DISCONNECTED, State#state{socket = undefined}};
|
||
|
||
%%% 处理内部消息,ssl收到的消息会解析成protobuf的消息格式,并按照internal类型处理
|
||
|
||
%% 微服务部署
|
||
handle_event(internal, {decoded_request, #'RequestFrame'{packet_id = PacketId, body = {container_request, Request}}},
|
||
?STATE_ACTIVATED, State = #state{socket = Socket}) ->
|
||
case docker_container_service:handle_request(Request) of
|
||
ok ->
|
||
send_result_reply(Socket, PacketId, <<"ok">>);
|
||
{ok, Reply} ->
|
||
send_result_reply(Socket, PacketId, Reply);
|
||
{error, Reason} ->
|
||
send_error_reply(Socket, PacketId, Reason)
|
||
end,
|
||
{keep_state, State};
|
||
handle_event(internal, {decoded_request, #'RequestFrame'{packet_id = PacketId, body = _Body}},
|
||
?STATE_RESTRICTED, State = #state{socket = Socket}) ->
|
||
send_error_reply(Socket, PacketId, <<"agent restricted">>),
|
||
{keep_state, State};
|
||
handle_event(internal, {decoded_request, #'RequestFrame'{packet_id = PacketId, body = _Body}},
|
||
_StateName, State = #state{socket = Socket}) ->
|
||
send_error_reply(Socket, PacketId, <<"agent state invalid">>),
|
||
{keep_state, State};
|
||
|
||
handle_event(internal, {decoded_reply, #'ReplyFrame'{packet_id = AuthPacketId, reply = {result, #'ReplyResult'{data = Message}}}},
|
||
?STATE_AUTH, State = #state{auth_packet_id = AuthPacketId}) ->
|
||
|
||
logger:debug("[efka_client] auth success, message: ~p", [Message]),
|
||
{next_state, ?STATE_ACTIVATED, State, [{next_event, info, flush_cache}]};
|
||
handle_event(internal, {decoded_reply, #'ReplyFrame'{packet_id = AuthPacketId, reply = {error, #'ReplyError'{code = 1, message = Message}}}},
|
||
?STATE_AUTH, State = #state{auth_packet_id = AuthPacketId}) ->
|
||
logger:debug("[efka_client] auth denied, message: ~p", [Message]),
|
||
{next_state, ?STATE_RESTRICTED, State};
|
||
handle_event(internal, {decoded_reply, #'ReplyFrame'{packet_id = AuthPacketId, reply = {error, #'ReplyError'{message = Message}}}},
|
||
?STATE_AUTH, State = #state{socket = Socket, auth_packet_id = AuthPacketId}) ->
|
||
logger:debug("[efka_client] auth failed, message: ~p", [Message]),
|
||
disconnect(Socket),
|
||
schedule_reconnect(),
|
||
{next_state, ?STATE_DISCONNECTED, State#state{socket = undefined}};
|
||
handle_event(internal, {decoded_reply, ReplyFrame}, StateName, State) ->
|
||
logger:warning("[efka_client] ignore unexpected reply in state ~p: ~p", [StateName, ReplyFrame]),
|
||
{keep_state, State};
|
||
|
||
%% 处理命令
|
||
handle_event(internal, {decoded_cast, #'CastFrame'{body = {command, #'Command'{command_type = ?COMMAND_AUTH, command = Auth0}}}},
|
||
StateName, State = #state{socket = Socket, next_packet_id = PacketId}) ->
|
||
|
||
Auth = binary_to_integer(Auth0),
|
||
case {Auth, StateName} of
|
||
{1, ?STATE_ACTIVATED} ->
|
||
{keep_state, State};
|
||
{1, _} ->
|
||
AuthPacket = auth_packet(PacketId),
|
||
send_packet(Socket, AuthPacket),
|
||
{next_state, ?STATE_AUTH, State#state{auth_packet_id = PacketId, next_packet_id = PacketId + 1},
|
||
[{state_timeout, 5000, auth_timeout}]};
|
||
{0, _} ->
|
||
{next_state, ?STATE_RESTRICTED, State}
|
||
end;
|
||
|
||
%% 处理Pub/Sub机制
|
||
handle_event(internal, {decoded_cast, #'CastFrame'{body = {pub, #'Pub'{topic = Topic, qos = Qos, content = Content}}}}, ?STATE_ACTIVATED, State) ->
|
||
logger:debug("[efka_client] get pub topic: ~p, qos: ~p, content: ~p", [Topic, Qos, Content]),
|
||
efka_subscription:publish(Topic, Qos, Content),
|
||
{keep_state, State};
|
||
|
||
handle_event(info, Info, _, State = #state{}) ->
|
||
logger:notice("[efka_client] get unknown info: ~p", [Info]),
|
||
{keep_state, State}.
|
||
|
||
terminate(_Reason, _StateName, _State = #state{socket = Socket}) ->
|
||
disconnect(Socket),
|
||
ok.
|
||
|
||
code_change(_OldVsn, StateName, State = #state{}, _Extra) ->
|
||
{ok, StateName, State}.
|
||
|
||
%%%===================================================================
|
||
%%% Internal functions
|
||
%%%===================================================================
|
||
|
||
-spec auth_packet(PktId :: integer()) -> binary().
|
||
auth_packet(PktId) when is_integer(PktId) ->
|
||
{ok, AuthInfo} = application:get_env(efka, auth),
|
||
UUID = proplists:get_value(uuid, AuthInfo),
|
||
Username = proplists:get_value(username, AuthInfo),
|
||
Salt = proplists:get_value(salt, AuthInfo),
|
||
Token = proplists:get_value(token, AuthInfo),
|
||
message_pb:encode_msg(#'RequestFrame'{
|
||
packet_id = PktId,
|
||
body = {auth_request, #'AuthRequest'{
|
||
uuid = unicode:characters_to_binary(UUID),
|
||
username = unicode:characters_to_binary(Username),
|
||
salt = unicode:characters_to_binary(Salt),
|
||
token = unicode:characters_to_binary(Token),
|
||
timestamp = efka_util:timestamp()
|
||
}}
|
||
}).
|
||
|
||
-spec connect_socket() -> {ok, ssl:sslsocket()} | {error, term()}.
|
||
connect_socket() ->
|
||
{ok, Props} = application:get_env(efka, tls_server),
|
||
Host = proplists:get_value(host, Props),
|
||
Port = proplists:get_value(port, Props),
|
||
SslOptions = [
|
||
binary,
|
||
{active, true},
|
||
{packet, 4},
|
||
{verify, verify_none}
|
||
],
|
||
ssl:connect(Host, Port, SslOptions, 5000).
|
||
|
||
-spec send_packet(ssl:sslsocket(), binary()) -> ok.
|
||
send_packet(Socket, Packet) when is_binary(Packet) ->
|
||
ok = ssl:send(Socket, Packet).
|
||
|
||
-spec disconnect(undefined | ssl:sslsocket()) -> ok.
|
||
disconnect(undefined) ->
|
||
ok;
|
||
disconnect(Socket) ->
|
||
catch ssl:close(Socket),
|
||
ok.
|
||
|
||
-spec schedule_reconnect() -> reference().
|
||
schedule_reconnect() ->
|
||
erlang:start_timer(5000, self(), create_transport).
|
||
|
||
-spec send_result_reply(ssl:sslsocket(), integer(), binary()) -> ok.
|
||
send_result_reply(Socket, PacketId, Payload) when is_binary(Payload) ->
|
||
Packet = message_pb:encode_msg(#'ReplyFrame'{
|
||
packet_id = PacketId,
|
||
reply = {result, #'ReplyResult'{data = Payload}}
|
||
}),
|
||
send_packet(Socket, Packet).
|
||
|
||
-spec send_error_reply(ssl:sslsocket(), integer(), binary()) -> ok.
|
||
send_error_reply(Socket, PacketId, Reason) when is_binary(Reason) ->
|
||
Packet = message_pb:encode_msg(#'ReplyFrame'{
|
||
packet_id = PacketId,
|
||
reply = {error, #'ReplyError'{code = -1, message = Reason}}
|
||
}),
|
||
send_packet(Socket, Packet). |