From 0dfd2828723c8120a58c63b9bad747cae00ad9a7 Mon Sep 17 00:00:00 2001 From: anlicheng <244108715@qq.com> Date: Sun, 26 Apr 2026 15:41:06 +0800 Subject: [PATCH] fix --- src/transport/cloud_wire.erl | 81 ----------------------------------- src/transport/efka_client.erl | 70 ++++++++++++++++++++++-------- 2 files changed, 52 insertions(+), 99 deletions(-) delete mode 100644 src/transport/cloud_wire.erl diff --git a/src/transport/cloud_wire.erl b/src/transport/cloud_wire.erl deleted file mode 100644 index a788ee3..0000000 --- a/src/transport/cloud_wire.erl +++ /dev/null @@ -1,81 +0,0 @@ --module(cloud_wire). - --export([encode/1, decode/1]). - --spec encode(term()) -> binary(). -encode(Term) -> - term_to_binary(validate(Term)). - --spec decode(binary()) -> term(). -decode(Bin) when is_binary(Bin) -> - validate(binary_to_term(Bin, [safe])). - --spec validate(term()) -> term(). -validate({request, PacketId, {auth_request, #{uuid := UUID, token := Token, timestamp := Timestamp} = Auth}}) - when is_integer(PacketId), PacketId > 0, is_binary(UUID), is_binary(Token), is_integer(Timestamp) -> - {request, PacketId, {auth_request, Auth}}; -validate({request, PacketId, {container_request, Request}}) - when is_integer(PacketId), PacketId > 0 -> - {request, PacketId, {container_request, validate_container_request(Request)}}; -validate({response, PacketId, {ok, Result}}) - when is_integer(PacketId), PacketId > 0, is_binary(Result) -> - {response, PacketId, {ok, Result}}; -validate({response, PacketId, {error, Code, Message}}) - when is_integer(PacketId), PacketId > 0, is_integer(Code), is_binary(Message) -> - {response, PacketId, {error, Code, Message}}; -validate({message, {data, #{route_key := RouteKey, metric := Metric} = Body}}) - when is_binary(RouteKey), is_binary(Metric) -> - {message, {data, Body}}; -validate({message, {task_event, #{task_id := TaskId, type := Type, stream := Stream} = Body}}) - when is_integer(TaskId), is_binary(Type), is_binary(Stream) -> - {message, {task_event, Body}}; -validate({message, {pub, #{topic := Topic, qos := Qos, content := Content} = Body}}) - when is_binary(Topic), is_integer(Qos), is_binary(Content) -> - {message, {pub, Body}}; -validate({message, {auth_control, Command}}) - when Command =:= activate; Command =:= deactivate -> - {message, {auth_control, Command}}; -validate(Term) -> - error({invalid_wire_term, Term}). - --spec validate_container_request(term()) -> map(). -validate_container_request(#{action := list} = Request) -> - true = is_boolean(maps:get(all, Request, true)), - Request; -validate_container_request(#{action := deploy, task_id := TaskId, params := Params} = Request) - when is_integer(TaskId), is_map(Params) -> - true = is_binary(maps:get(container_name, Params, <<>>)), - true = is_binary(maps:get(container_dir, Params, <<>>)), - true = is_map(maps:get(create, Params, #{})), - Request; -validate_container_request(#{action := start, target := Target} = Request) -> - validate_target(Target), - Request; -validate_container_request(#{action := stop, target := Target, timeout_seconds := TimeoutSeconds} = Request) - when is_integer(TimeoutSeconds), TimeoutSeconds >= 0 -> - validate_target(Target), - Request; -validate_container_request(#{action := kill, target := Target, signal := Signal} = Request) - when is_binary(Signal) -> - validate_target(Target), - Request; -validate_container_request(#{action := remove, target := Target, force := Force, remove_volumes := RemoveVolumes} = Request) - when is_boolean(Force), is_boolean(RemoveVolumes) -> - validate_target(Target), - Request; -validate_container_request(#{action := config, target := Target, config := Config} = Request) - when is_binary(Config) -> - validate_target(Target), - Request; -validate_container_request(Request) -> - error({invalid_container_request, Request}). - --spec validate_target(term()) -> ok. -validate_target(#{name := Name, id := Id}) when is_binary(Name), is_binary(Id) -> - ok; -validate_target(#{name := Name}) when is_binary(Name) -> - ok; -validate_target(#{id := Id}) when is_binary(Id) -> - ok; -validate_target(Target) -> - error({invalid_container_target, Target}). diff --git a/src/transport/efka_client.erl b/src/transport/efka_client.erl index d2aef11..e7b4e1a 100644 --- a/src/transport/efka_client.erl +++ b/src/transport/efka_client.erl @@ -95,7 +95,7 @@ callback_mode() -> %% 异步发送数据, 连接存在时候直接发送;否则缓存到mnesia -spec handle_event(term(), term(), atom(), #state{}) -> term(). handle_event(cast, {metric_data, RouteKey, Metric}, StateName, State = #state{socket = Socket}) -> - Packet = cloud_wire:encode({message, {data, #{route_key => RouteKey, metric => Metric}}}), + Packet = term_to_binary({message, {data, #{route_key => RouteKey, metric => Metric}}}), case StateName of ?STATE_ACTIVATED -> send_packet(Socket, Packet), @@ -108,12 +108,12 @@ handle_event(cast, {metric_data, RouteKey, Metric}, StateName, State = #state{so %% 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]), - Packet = cloud_wire:encode({message, {task_event, #{task_id => TaskId, type => Type, stream => Stream}}}), + Packet = term_to_binary({message, {task_event, #{task_id => TaskId, type => Type, stream => Stream}}}), send_packet(Socket, Packet), {keep_state, State}; handle_event(cast, {close_task_event_stream, TaskId, Reason}, ?STATE_ACTIVATED, State = #state{socket = Socket}) -> - Packet = cloud_wire:encode({message, {task_event, #{task_id => TaskId, type => <<"close">>, stream => Reason}}}), + Packet = term_to_binary({message, {task_event, #{task_id => TaskId, type => <<"close">>, stream => Reason}}}), send_packet(Socket, Packet), {keep_state, State}; @@ -162,9 +162,26 @@ handle_event(info, flush_cache, _, State) -> {keep_state, State}; %% 处理收到的ssl消息 -handle_event(info, {ssl, Socket, PacketBin}, _, State = #state{socket = Socket}) -> - Packet = cloud_wire:decode(PacketBin), - {keep_state, State, [{next_event, internal, {decoded_packet, Packet}}]}; +handle_event(info, {ssl, Socket, PacketBin}, _, State = #state{socket = Socket}) + when is_binary(PacketBin) -> + case binary_to_term(PacketBin, [safe]) of + {request, PacketId, {container_request, Request}} -> + {keep_state, State, [{next_event, internal, {decoded_request, PacketId, Request}}]}; + {request, PacketId, Body} -> + {keep_state, State, [{next_event, internal, {decoded_request_invalid, PacketId, Body}}]}; + {response, AuthPacketId, {ok, Message}} -> + {keep_state, State, [{next_event, internal, {decoded_auth_ok, AuthPacketId, Message}}]}; + {response, AuthPacketId, {error, 1, Message}} -> + {keep_state, State, [{next_event, internal, {decoded_auth_denied, AuthPacketId, Message}}]}; + {response, AuthPacketId, {error, Code, Message}} -> + {keep_state, State, [{next_event, internal, {decoded_auth_error, AuthPacketId, Code, Message}}]}; + {message, {auth_control, Cmd}} -> + {keep_state, State, [{next_event, internal, {decoded_auth_control, Cmd}}]}; + {message, {pub, #{topic := Topic, qos := Qos, content := Content}}} -> + {keep_state, State, [{next_event, internal, {decoded_pub, Topic, Qos, Content}}]}; + Packet -> + {keep_state, State, [{next_event, internal, {decoded_unknown, Packet}}]} + end; handle_event(info, {ssl_error, Socket, Reason}, _, State = #state{socket = Socket}) -> logger:debug("[efka_client] ssl error: ~p", [Reason]), disconnect(Socket), @@ -177,7 +194,7 @@ handle_event(info, {ssl_closed, Socket}, _, State = #state{socket = Socket}) -> %%% 处理内部消息,ssl收到的消息会解析成protobuf的消息格式,并按照internal类型处理 %% 微服务部署 -handle_event(internal, {decoded_packet, {request, PacketId, {container_request, Request}}}, +handle_event(internal, {decoded_request, PacketId, Request}, ?STATE_ACTIVATED, State = #state{socket = Socket}) -> case docker_container_service:handle_request(Request) of ok -> @@ -188,36 +205,50 @@ handle_event(internal, {decoded_packet, {request, PacketId, {container_request, send_error_reply(Socket, PacketId, Reason) end, {keep_state, State}; -handle_event(internal, {decoded_packet, {request, PacketId, _Body}}, +handle_event(internal, {decoded_request_invalid, PacketId, _Body}, ?STATE_RESTRICTED, State = #state{socket = Socket}) -> send_error_reply(Socket, PacketId, <<"agent restricted">>), {keep_state, State}; -handle_event(internal, {decoded_packet, {request, PacketId, _Body}}, +handle_event(internal, {decoded_request_invalid, PacketId, _Body}, + _StateName, State = #state{socket = Socket}) -> + send_error_reply(Socket, PacketId, <<"agent state invalid">>), + {keep_state, State}; +handle_event(internal, {decoded_request, PacketId, _Request}, + ?STATE_RESTRICTED, State = #state{socket = Socket}) -> + send_error_reply(Socket, PacketId, <<"agent restricted">>), + {keep_state, State}; +handle_event(internal, {decoded_request, PacketId, _Request}, _StateName, State = #state{socket = Socket}) -> send_error_reply(Socket, PacketId, <<"agent state invalid">>), {keep_state, State}; -handle_event(internal, {decoded_packet, {response, AuthPacketId, {ok, Message}}}, +handle_event(internal, {decoded_auth_ok, AuthPacketId, 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_packet, {response, AuthPacketId, {error, 1, Message}}}, +handle_event(internal, {decoded_auth_denied, AuthPacketId, 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_packet, {response, AuthPacketId, {error, _Code, Message}}}, +handle_event(internal, {decoded_auth_error, AuthPacketId, _Code, 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_packet, {response, _PacketId, Reply}}, StateName, State) -> +handle_event(internal, {decoded_auth_ok, _PacketId, Reply}, StateName, State) -> + logger:warning("[efka_client] ignore unexpected response in state ~p: ~p", [StateName, Reply]), + {keep_state, State}; +handle_event(internal, {decoded_auth_denied, _PacketId, Reply}, StateName, State) -> + logger:warning("[efka_client] ignore unexpected response in state ~p: ~p", [StateName, Reply]), + {keep_state, State}; +handle_event(internal, {decoded_auth_error, _PacketId, _Code, Reply}, StateName, State) -> logger:warning("[efka_client] ignore unexpected response in state ~p: ~p", [StateName, Reply]), {keep_state, State}; %% 处理命令 -handle_event(internal, {decoded_packet, {message, {auth_control, Cmd}}}, +handle_event(internal, {decoded_auth_control, Cmd}, StateName, State = #state{socket = Socket, next_packet_id = PacketId}) -> logger:debug("[efka_client] auth cmd: ~p", [Cmd]), @@ -234,10 +265,13 @@ handle_event(internal, {decoded_packet, {message, {auth_control, Cmd}}}, end; %% 处理Pub/Sub机制 -handle_event(internal, {decoded_packet, {message, {pub, #{topic := Topic, qos := Qos, content := Content}}}}, ?STATE_ACTIVATED, State) -> +handle_event(internal, {decoded_pub, Topic, Qos, 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(internal, {decoded_unknown, Packet}, _StateName, State) -> + logger:warning("[efka_client] ignore unknown packet: ~p", [Packet]), + {keep_state, State}; handle_event(info, Info, _, State = #state{}) -> logger:notice("[efka_client] get unknown info: ~p", [Info]), @@ -264,7 +298,7 @@ auth_packet(PktId) when is_integer(PktId) -> Token = proplists:get_value(token, AuthInfo), Timestamp = efka_util:timestamp(), - cloud_wire:encode({request, PktId, {auth_request, #{ + term_to_binary({request, PktId, {auth_request, #{ uuid => list_to_binary(UUID), token => list_to_binary(Token), timestamp => Timestamp @@ -388,10 +422,10 @@ delete_oldest_cache_entry() -> -spec send_result_reply(ssl:sslsocket(), integer(), binary()) -> ok. send_result_reply(Socket, PacketId, Payload) when is_binary(Payload) -> - Packet = cloud_wire:encode({response, PacketId, {ok, Payload}}), + Packet = term_to_binary({response, PacketId, {ok, 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 = cloud_wire:encode({response, PacketId, {error, -1, Reason}}), + Packet = term_to_binary({response, PacketId, {error, -1, Reason}}), send_packet(Socket, Packet).