%%%------------------------------------------------------------------- %%% @author anlicheng %%% @copyright (C) 2025, %%% @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 = <>, 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, <>), {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, <>), {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, <>}, _, 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, <>}, _, 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, <>}, _, 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).