fix efka_client

This commit is contained in:
anlicheng 2026-04-20 16:47:28 +08:00
parent 437f5cf1f7
commit 38511c7b7a

View File

@ -32,7 +32,10 @@
-define(STATE_ACTIVATED, activated). -define(STATE_ACTIVATED, activated).
-record(state, { -record(state, {
socket :: undefined | ssl:sslsocket() socket :: undefined | ssl:sslsocket(),
next_packet_id = 1,
%% auth请求的packet_idauth请求和响应的对应关系
auth_packet_id = 1
}). }).
%%%=================================================================== %%%===================================================================
@ -105,12 +108,13 @@ handle_event(cast, {task_event_stream, _TaskId, _Stream}, _, State = #state{}) -
{keep_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 case connect_socket() of
{ok, Socket} -> {ok, Socket} ->
AuthPacket = auth_packet(), AuthPacket = auth_packet(PacketId),
send_packet(Socket, AuthPacket), 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} -> {error, Reason} ->
logger:debug("[efka_client] connect failed, error: ~p", [Reason]), logger:debug("[efka_client] connect failed, error: ~p", [Reason]),
schedule_reconnect(), schedule_reconnect(),
@ -180,30 +184,38 @@ handle_event(internal, {decoded_request, #'RequestFrame'{packet_id = PacketId, b
send_error_reply(Socket, PacketId, <<"agent state invalid">>), send_error_reply(Socket, PacketId, <<"agent state invalid">>),
{keep_state, State}; {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]), logger:debug("[efka_client] auth success, message: ~p", [Message]),
{next_state, ?STATE_ACTIVATED, State, [{next_event, info, flush_cache}]}; {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]), logger:debug("[efka_client] auth denied, message: ~p", [Message]),
{next_state, ?STATE_RESTRICTED, State}; {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]), logger:debug("[efka_client] auth failed, message: ~p", [Message]),
disconnect(Socket), disconnect(Socket),
schedule_reconnect(), schedule_reconnect(),
{next_state, ?STATE_DISCONNECTED, State#state{socket = undefined}}; {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'{ handle_event(internal, {decoded_cast, #'CastFrame'{
body = {command, #'Command'{command_type = ?COMMAND_AUTH, command = Auth0}} 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), Auth = binary_to_integer(Auth0),
case {Auth, StateName} of case {Auth, StateName} of
{1, ?STATE_ACTIVATED} -> {1, ?STATE_ACTIVATED} ->
{keep_state, State}; {keep_state, State};
{1, _} -> {1, _} ->
AuthPacket = auth_packet(), AuthPacket = auth_packet(PacketId),
send_packet(Socket, AuthPacket), 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, _} -> {0, _} ->
{next_state, ?STATE_RESTRICTED, State} {next_state, ?STATE_RESTRICTED, State}
end; end;
@ -231,15 +243,15 @@ code_change(_OldVsn, StateName, State = #state{}, _Extra) ->
%%% Internal functions %%% Internal functions
%%%=================================================================== %%%===================================================================
-spec auth_packet() -> binary(). -spec auth_packet(PktId :: integer()) -> binary().
auth_packet() -> auth_packet(PktId) when is_integer(PktId) ->
{ok, AuthInfo} = application:get_env(efka, auth), {ok, AuthInfo} = application:get_env(efka, auth),
UUID = proplists:get_value(uuid, AuthInfo), UUID = proplists:get_value(uuid, AuthInfo),
Username = proplists:get_value(username, AuthInfo), Username = proplists:get_value(username, AuthInfo),
Salt = proplists:get_value(salt, AuthInfo), Salt = proplists:get_value(salt, AuthInfo),
Token = proplists:get_value(token, AuthInfo), Token = proplists:get_value(token, AuthInfo),
message_pb:encode_msg(#'RequestFrame'{ message_pb:encode_msg(#'RequestFrame'{
packet_id = 1, packet_id = PktId,
body = {auth_request, #'AuthRequest'{ body = {auth_request, #'AuthRequest'{
uuid = unicode:characters_to_binary(UUID), uuid = unicode:characters_to_binary(UUID),
username = unicode:characters_to_binary(Username), username = unicode:characters_to_binary(Username),