From 38511c7b7a472b22e1c7ea324f68c9971169b4eb Mon Sep 17 00:00:00 2001 From: anlicheng <244108715@qq.com> Date: Mon, 20 Apr 2026 16:47:28 +0800 Subject: [PATCH] fix efka_client --- src/efka_client.erl | 38 +++++++++++++++++++++++++------------- 1 file changed, 25 insertions(+), 13 deletions(-) diff --git a/src/efka_client.erl b/src/efka_client.erl index 4ae03a3..5edd6fc 100644 --- a/src/efka_client.erl +++ b/src/efka_client.erl @@ -32,7 +32,10 @@ -define(STATE_ACTIVATED, activated). -record(state, { - socket :: undefined | ssl:sslsocket() + socket :: undefined | ssl:sslsocket(), + next_packet_id = 1, + %% 保存当前auth请求的packet_id,用来建立auth请求和响应的对应关系 + auth_packet_id = 1 }). %%%=================================================================== @@ -105,12 +108,13 @@ handle_event(cast, {task_event_stream, _TaskId, _Stream}, _, State = #state{}) - {keep_state, State}; %% 异步建立到服务器的连接 -handle_event(info, {timeout, _, create_transport}, ?STATE_DISCONNECTED, State) -> +handle_event(info, {timeout, _, create_transport}, ?STATE_DISCONNECTED, State = #state{next_packet_id = PacketId}) -> case connect_socket() of {ok, Socket} -> - AuthPacket = auth_packet(), + AuthPacket = auth_packet(PacketId), send_packet(Socket, AuthPacket), - {next_state, ?STATE_AUTH, State#state{socket = Socket}, [{state_timeout, 5000, auth_timeout}]}; + {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(), @@ -180,30 +184,38 @@ handle_event(internal, {decoded_request, #'RequestFrame'{packet_id = PacketId, b send_error_reply(Socket, PacketId, <<"agent state invalid">>), {keep_state, State}; -handle_event(internal, {decoded_reply, #'ReplyFrame'{packet_id = 1, reply = {result, #'ReplyResult'{data = Message}}}}, ?STATE_AUTH, 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 = 1, reply = {error, #'ReplyError'{code = 1, message = Message}}}}, ?STATE_AUTH, State) -> +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 = 1, reply = {error, #'ReplyError'{message = Message}}}}, ?STATE_AUTH, State = #state{socket = Socket}) -> +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}) -> +}}, 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(), + AuthPacket = auth_packet(PacketId), send_packet(Socket, AuthPacket), - {next_state, ?STATE_AUTH, State, [{state_timeout, 5000, auth_timeout}]}; + {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; @@ -231,15 +243,15 @@ code_change(_OldVsn, StateName, State = #state{}, _Extra) -> %%% Internal functions %%%=================================================================== --spec auth_packet() -> binary(). -auth_packet() -> +-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 = 1, + packet_id = PktId, body = {auth_request, #'AuthRequest'{ uuid = unicode:characters_to_binary(UUID), username = unicode:characters_to_binary(Username),