diff --git a/src/docker/docker_commands.erl b/src/docker/docker_commands.erl index 5afd872..2c747e3 100644 --- a/src/docker/docker_commands.erl +++ b/src/docker/docker_commands.erl @@ -99,7 +99,7 @@ start_container(ContainerName) when is_binary(ContainerName) -> {ok, 304, _Headers, _} -> {error, <<"container already started">>}; {ok, _StatusCode, _Header, ErrorResp} -> - case catch jiffy:decode(ErrorResp) of + case catch jiffy:decode(ErrorResp, [return_maps]) of #{<<"message">> := Msg} -> {error, Msg}; _ -> @@ -125,7 +125,7 @@ stop_container(ContainerName, TimeoutSeconds) when is_binary(ContainerName), is_ {ok, 304, _Headers, _} -> {error, <<"container already stopped">>}; {ok, _StatusCode, _Header, ErrorResp} -> - case catch jiffy:decode(ErrorResp) of + case catch jiffy:decode(ErrorResp, [return_maps]) of #{<<"message">> := Msg} -> {error, Msg}; _ -> @@ -149,7 +149,7 @@ kill_container(ContainerName, Signal) when is_binary(ContainerName), is_binary(S {ok, 204, _Headers, _} -> ok; {ok, _StatusCode, _Header, ErrorResp} -> - case catch jiffy:decode(ErrorResp) of + case catch jiffy:decode(ErrorResp, [return_maps]) of #{<<"message">> := Msg} -> {error, Msg}; _ -> @@ -176,7 +176,7 @@ remove_container(ContainerName, Force, RemoveVolumes) {ok, 304, _Headers, _} -> {error, <<"container already stopped">>}; {ok, _StatusCode, _Header, ErrorResp} -> - case catch jiffy:decode(ErrorResp) of + case catch jiffy:decode(ErrorResp, [return_maps]) of #{<<"message">> := Msg} -> {error, Msg}; _ -> @@ -197,7 +197,7 @@ get_containers() -> Containers = jiffy:decode(ContainersBin, [return_maps]), {ok, Containers}; {ok, _StatusCode, _Header, ErrorResp} -> - case catch jiffy:decode(ErrorResp) of + case catch jiffy:decode(ErrorResp, [return_maps]) of #{<<"message">> := Msg} -> {error, Msg}; _ -> @@ -217,7 +217,7 @@ inspect_container(ContainerId) when is_binary(ContainerId) -> {ok, 200, _Headers, Resp} -> decode_container_inspect_summary(Resp); {ok, _StatusCode, _Header, ErrorResp} -> - case catch jiffy:decode(ErrorResp) of + case catch jiffy:decode(ErrorResp, [return_maps]) of #{<<"message">> := Msg} -> {error, Msg}; _ -> diff --git a/src/transport/efka_client.erl b/src/transport/efka_client.erl index 2327a7e..303729a 100644 --- a/src/transport/efka_client.erl +++ b/src/transport/efka_client.erl @@ -87,12 +87,11 @@ handle_event(cast, {metric_data, RouteKey, Metric}, StateName, State = #state{so CastFrame = message_pb:encode_msg(#'CastFrame'{ body = {data, #'CastFrame.Data'{route_key = RouteKey, metric = Metric}} }), - Packet = <>, case StateName of ?STATE_ACTIVATED -> - send_packet(Socket, Packet); + send_packet(Socket, [?FRAME_CAST, CastFrame]); _ -> - ok = cache_model:insert(Packet) + ok = cache_model:insert(<>) end, {keep_state, State}; @@ -102,14 +101,14 @@ handle_event(cast, {task_event_stream, TaskId, Type, Stream}, ?STATE_ACTIVATED, EventPacket = message_pb:encode_msg(#'CastFrame'{ body = {task_event, #'CastFrame.TaskEvent'{task_id = TaskId, type = Type, stream = Stream}} }), - send_packet(Socket, <>), + send_packet(Socket, [?FRAME_CAST, EventPacket]), {keep_state, State}; handle_event(cast, {close_task_event_stream, TaskId, Reason}, ?STATE_ACTIVATED, State = #state{socket = Socket}) -> EventPacket = message_pb:encode_msg(#'CastFrame'{ body = {task_event, #'CastFrame.TaskEvent'{task_id = TaskId, type = <<"close">>, stream = Reason}} }), - send_packet(Socket, <>), + send_packet(Socket, [?FRAME_CAST, EventPacket]), {keep_state, State}; %% 其他情况下直接忽略 @@ -126,8 +125,7 @@ handle_event(info, {timeout, _, create_transport}, ?STATE_DISCONNECTED, State = case connect_socket() of {ok, Socket} -> AuthPacket = auth_packet(PacketId), - AuthRequestFrame = <>, - send_packet(Socket, AuthRequestFrame), + send_packet(Socket, [?FRAME_REQUEST, AuthPacket]), {next_state, ?STATE_AUTH, State#state{socket = Socket, auth_packet_id = PacketId, next_packet_id = PacketId + 1}, [{state_timeout, 5000, auth_timeout}]}; {error, Reason} -> @@ -226,8 +224,7 @@ handle_event(internal, {decoded_cast, #'CastFrame'{body = {command, #'CastFrame. {keep_state, State}; {1, _} -> AuthPacket = auth_packet(PacketId), - AuthRequestFrame = <>, - send_packet(Socket, AuthRequestFrame), + send_packet(Socket, [?FRAME_REQUEST, AuthPacket]), {next_state, ?STATE_AUTH, State#state{auth_packet_id = PacketId, next_packet_id = PacketId + 1}, [{state_timeout, 5000, auth_timeout}]}; {0, _} -> @@ -288,8 +285,8 @@ connect_socket() -> ], ssl:connect(Host, Port, SslOptions, 5000). --spec send_packet(ssl:sslsocket(), binary()) -> ok. -send_packet(Socket, Packet) when is_binary(Packet) -> +-spec send_packet(ssl:sslsocket(), erlang:iodata()) -> ok. +send_packet(Socket, Packet) -> ok = ssl:send(Socket, Packet). -spec disconnect(undefined | ssl:sslsocket()) -> ok. @@ -309,7 +306,7 @@ send_result_reply(Socket, PacketId, Payload) when is_binary(Payload) -> packet_id = PacketId, reply = {result, Payload} }), - send_packet(Socket, Packet). + send_packet(Socket, [?FRAME_REPLY, Packet]). -spec send_error_reply(ssl:sslsocket(), integer(), binary()) -> ok. send_error_reply(Socket, PacketId, Reason) when is_binary(Reason) -> @@ -317,4 +314,4 @@ send_error_reply(Socket, PacketId, Reason) when is_binary(Reason) -> packet_id = PacketId, reply = {error, #'ReplyFrame.Error'{code = -1, message = Reason}} }), - send_packet(Socket, Packet). + send_packet(Socket, [?FRAME_REPLY, Packet]).