From 9bdf5fcac2ffd5064a5cdb7bca3bfea053346e95 Mon Sep 17 00:00:00 2001 From: anlicheng <244108715@qq.com> Date: Mon, 20 Apr 2026 15:11:44 +0800 Subject: [PATCH] fix efka_client --- include/message.hrl | 32 -------------------------------- src/efka_client.erl | 42 ++++++++++++++++++++++-------------------- 2 files changed, 22 insertions(+), 52 deletions(-) delete mode 100644 include/message.hrl diff --git a/include/message.hrl b/include/message.hrl deleted file mode 100644 index fdb9a52..0000000 --- a/include/message.hrl +++ /dev/null @@ -1,32 +0,0 @@ -%%%------------------------------------------------------------------- -%%% @author anlicheng -%%% @copyright (C) 2025, -%%% @doc -%%% 扩展部分, 1: 支持基于topic的pub/sub机制; 2. 基于target的单点通讯和广播 -%%% @end -%%% Created : 21. 4月 2025 17:28 -%%%------------------------------------------------------------------- --author("anlicheng"). - -%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%% -%%%% 二级分类定义 -%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%% - -%% 主机端上报数据类型标识 --define(MESSAGE_AUTH_REQUEST, 16#01). --define(MESSAGE_AUTH_REPLY, 16#02). - --define(MESSAGE_COMMAND, 16#03). --define(MESSAGE_PUB, 16#05). - --define(MESSAGE_DATA, 16#06). - -%% efka主动上报的event-stream流, 单向消息,主要是: docker-create的实时处理逻辑上报 --define(MESSAGE_EVENT_STREAM, 16#08). - --define(MESSAGE_JSONRPC_REQUEST, 16#F0). --define(MESSAGE_JSONRPC_REPLY, 16#F1). - -%%%% 命令类型子分类, 不需要返回值 -%% 授权 --define(COMMAND_AUTH, 16#08). diff --git a/src/efka_client.erl b/src/efka_client.erl index 696c57e..c9ec4b6 100644 --- a/src/efka_client.erl +++ b/src/efka_client.erl @@ -8,7 +8,7 @@ %%%------------------------------------------------------------------- -module(efka_client). -author("anlicheng"). --include("message.hrl"). +-include("protocol.hrl"). -include("message_pb.hrl"). -include("efka_tables.hrl"). @@ -118,6 +118,27 @@ handle_event(info, {timeout, _, create_transport}, ?STATE_DENIED, State) -> handle_event(info, {timeout, _, create_transport}, _, State) -> {keep_state, State}; +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_DENIED, 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, PacketBin}, ?STATE_AUTH, State = #state{socket = Socket}) -> #'ResponseFrame'{packet_id = 1, body = {auth_reply, #'AuthReply'{code = Code, payload = Message}}} = message_pb:decode_msg(PacketBin, 'ResponseFrame'), @@ -140,25 +161,6 @@ handle_event(info, {ssl, Socket, PacketBin}, ?STATE_AUTH, State = #state{socket {next_state, ?STATE_DENIED, 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_DENIED, 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}; - handle_event(info, {ssl, Socket, <<8, _/binary>> = PacketBin}, StateName, State = #state{socket = Socket}) when StateName =:= ?STATE_ACTIVATED; StateName =:= ?STATE_RESTRICTED -> #'RequestFrame'{packet_id = PacketId, body = Body} = message_pb:decode_msg(PacketBin, 'RequestFrame'),