From 6d6a9c45853560698eb5b3e59ca19bdce26cc321 Mon Sep 17 00:00:00 2001 From: anlicheng <244108715@qq.com> Date: Mon, 27 Apr 2026 11:25:24 +0800 Subject: [PATCH] =?UTF-8?q?=E7=BB=9F=E4=B8=80=E6=B6=88=E6=81=AF=E6=A0=BC?= =?UTF-8?q?=E5=BC=8F?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- src/transport/efka_client.erl | 34 ++++++++++++++-------------------- 1 file changed, 14 insertions(+), 20 deletions(-) diff --git a/src/transport/efka_client.erl b/src/transport/efka_client.erl index 602ead4..2496649 100644 --- a/src/transport/efka_client.erl +++ b/src/transport/efka_client.erl @@ -180,51 +180,51 @@ handle_event(info, {ssl_closed, Socket}, _, State = #state{socket = Socket}) -> handle_event(internal, {request, Ref, {container_request, #{action := list}}}, ?STATE_ACTIVATED, State = #state{socket = Socket}) -> Reply = docker_commands:get_containers(), - handle_container_reply(Socket, Ref, Reply), + handle_container_response(Socket, Ref, Reply), {keep_state, State}; handle_event(internal, {request, Ref, {container_request, #{action := deploy, task_id := TaskId, params := Params}}}, ?STATE_ACTIVATED, State = #state{socket = Socket}) -> Reply = docker_deploy_manager:deploy(TaskId, Params), - handle_container_reply(Socket, Ref, Reply), + handle_container_response(Socket, Ref, Reply), {keep_state, State}; handle_event(internal, {request, Ref, {container_request, #{action := start, target := Target}}}, ?STATE_ACTIVATED, State = #state{socket = Socket}) -> Reply = docker_commands:start_container(container_target(Target)), - handle_container_reply(Socket, Ref, Reply), + handle_container_response(Socket, Ref, Reply), {keep_state, State}; handle_event(internal, {request, Ref, {container_request, #{action := stop, target := Target, timeout_seconds := TimeoutSeconds}}}, ?STATE_ACTIVATED, State = #state{socket = Socket}) -> Reply = docker_commands:stop_container(container_target(Target), TimeoutSeconds), - handle_container_reply(Socket, Ref, Reply), + handle_container_response(Socket, Ref, Reply), {keep_state, State}; handle_event(internal, {request, Ref, {container_request, #{action := kill, target := Target, signal := Signal}}}, ?STATE_ACTIVATED, State = #state{socket = Socket}) -> Reply = docker_commands:kill_container(container_target(Target), to_binary(Signal)), - handle_container_reply(Socket, Ref, Reply), + handle_container_response(Socket, Ref, Reply), {keep_state, State}; handle_event(internal, {request, Ref, {container_request, #{action := remove, target := Target, force := Force, remove_volumes := RemoveVolumes}}}, ?STATE_ACTIVATED, State = #state{socket = Socket}) -> Reply = docker_commands:remove_container(container_target(Target), to_bool(Force), to_bool(RemoveVolumes)), - handle_container_reply(Socket, Ref, Reply), + handle_container_response(Socket, Ref, Reply), {keep_state, State}; handle_event(internal, {request, Ref, {container_request, #{action := config, target := Target, config := Config}}}, ?STATE_ACTIVATED, State = #state{socket = Socket}) -> Reply = update_container_config(container_target(Target), iolist_to_binary(Config)), - handle_container_reply(Socket, Ref, Reply), + handle_container_response(Socket, Ref, Reply), {keep_state, State}; handle_event(internal, {request, Ref, {container_request, Request}}, ?STATE_RESTRICTED, State = #state{socket = Socket}) -> logger:notice("[efka_client] get a invalid request: ~p, agent restricted", [Request]), - handle_container_reply(Socket, Ref, {error, <<"agent restricted">>}), + handle_container_response(Socket, Ref, {error, <<"agent restricted">>}), {keep_state, State}; handle_event(internal, {request, Ref, {container_request, Request}}, _StateName, State = #state{socket = Socket}) -> logger:notice("[efka_client] get a invalid request: ~p, agent restricted", [Request]), - handle_container_reply(Socket, Ref, {error, <<"agent invalid">>}), + handle_container_response(Socket, Ref, {error, <<"agent invalid">>}), {keep_state, State}; %% 处理response -handle_event(internal, {response, AuthRef, {auth_response, {ok, Message}}}, +handle_event(internal, {response, AuthRef, {auth_response, ok}}, ?STATE_AUTH, State = #state{auth_ref = AuthRef}) -> - logger:debug("[efka_client] auth success, message: ~p", [Message]), + logger:debug("[efka_client] auth success"), {next_state, ?STATE_ACTIVATED, State#state{auth_ref = undefined}, [{next_event, info, flush_cache}]}; handle_event(internal, {response, AuthRef, {auth_response, {error, {denied, Message}}}}, ?STATE_AUTH, State = #state{auth_ref = AuthRef}) -> @@ -410,15 +410,9 @@ delete_oldest_cache_entry() -> dets:delete(?CACHE_TAB, Key) end. --spec handle_container_reply(ssl:sslsocket(), reference(), term()) -> ok. -handle_container_reply(Socket, Ref, ok) -> - Packet = term_to_binary({response, Ref, {container_response, {ok, <<"ok">>}}}), - ok = ssl:send(Socket, Packet); -handle_container_reply(Socket, Ref, {ok, Response}) -> - Packet = term_to_binary({response, Ref, {container_response, {ok, Response}}}), - ok = ssl:send(Socket, Packet); -handle_container_reply(Socket, Ref, {error, Reason}) -> - Packet = term_to_binary({response, Ref, {container_response, {error, Reason}}}), +-spec handle_container_response(ssl:sslsocket(), reference(), term()) -> ok. +handle_container_response(Socket, Ref, Reply) -> + Packet = term_to_binary({response, Ref, {container_response, Reply}}), ok = ssl:send(Socket, Packet). -spec container_target(map()) -> binary().