fix frame

This commit is contained in:
anlicheng 2026-04-23 20:11:55 +08:00
parent c51045d667
commit 00bb22711b
2 changed files with 16 additions and 19 deletions

View File

@ -99,7 +99,7 @@ start_container(ContainerName) when is_binary(ContainerName) ->
{ok, 304, _Headers, _} -> {ok, 304, _Headers, _} ->
{error, <<"container already started">>}; {error, <<"container already started">>};
{ok, _StatusCode, _Header, ErrorResp} -> {ok, _StatusCode, _Header, ErrorResp} ->
case catch jiffy:decode(ErrorResp) of case catch jiffy:decode(ErrorResp, [return_maps]) of
#{<<"message">> := Msg} -> #{<<"message">> := Msg} ->
{error, Msg}; {error, Msg};
_ -> _ ->
@ -125,7 +125,7 @@ stop_container(ContainerName, TimeoutSeconds) when is_binary(ContainerName), is_
{ok, 304, _Headers, _} -> {ok, 304, _Headers, _} ->
{error, <<"container already stopped">>}; {error, <<"container already stopped">>};
{ok, _StatusCode, _Header, ErrorResp} -> {ok, _StatusCode, _Header, ErrorResp} ->
case catch jiffy:decode(ErrorResp) of case catch jiffy:decode(ErrorResp, [return_maps]) of
#{<<"message">> := Msg} -> #{<<"message">> := Msg} ->
{error, Msg}; {error, Msg};
_ -> _ ->
@ -149,7 +149,7 @@ kill_container(ContainerName, Signal) when is_binary(ContainerName), is_binary(S
{ok, 204, _Headers, _} -> {ok, 204, _Headers, _} ->
ok; ok;
{ok, _StatusCode, _Header, ErrorResp} -> {ok, _StatusCode, _Header, ErrorResp} ->
case catch jiffy:decode(ErrorResp) of case catch jiffy:decode(ErrorResp, [return_maps]) of
#{<<"message">> := Msg} -> #{<<"message">> := Msg} ->
{error, Msg}; {error, Msg};
_ -> _ ->
@ -176,7 +176,7 @@ remove_container(ContainerName, Force, RemoveVolumes)
{ok, 304, _Headers, _} -> {ok, 304, _Headers, _} ->
{error, <<"container already stopped">>}; {error, <<"container already stopped">>};
{ok, _StatusCode, _Header, ErrorResp} -> {ok, _StatusCode, _Header, ErrorResp} ->
case catch jiffy:decode(ErrorResp) of case catch jiffy:decode(ErrorResp, [return_maps]) of
#{<<"message">> := Msg} -> #{<<"message">> := Msg} ->
{error, Msg}; {error, Msg};
_ -> _ ->
@ -197,7 +197,7 @@ get_containers() ->
Containers = jiffy:decode(ContainersBin, [return_maps]), Containers = jiffy:decode(ContainersBin, [return_maps]),
{ok, Containers}; {ok, Containers};
{ok, _StatusCode, _Header, ErrorResp} -> {ok, _StatusCode, _Header, ErrorResp} ->
case catch jiffy:decode(ErrorResp) of case catch jiffy:decode(ErrorResp, [return_maps]) of
#{<<"message">> := Msg} -> #{<<"message">> := Msg} ->
{error, Msg}; {error, Msg};
_ -> _ ->
@ -217,7 +217,7 @@ inspect_container(ContainerId) when is_binary(ContainerId) ->
{ok, 200, _Headers, Resp} -> {ok, 200, _Headers, Resp} ->
decode_container_inspect_summary(Resp); decode_container_inspect_summary(Resp);
{ok, _StatusCode, _Header, ErrorResp} -> {ok, _StatusCode, _Header, ErrorResp} ->
case catch jiffy:decode(ErrorResp) of case catch jiffy:decode(ErrorResp, [return_maps]) of
#{<<"message">> := Msg} -> #{<<"message">> := Msg} ->
{error, Msg}; {error, Msg};
_ -> _ ->

View File

@ -87,12 +87,11 @@ handle_event(cast, {metric_data, RouteKey, Metric}, StateName, State = #state{so
CastFrame = message_pb:encode_msg(#'CastFrame'{ CastFrame = message_pb:encode_msg(#'CastFrame'{
body = {data, #'CastFrame.Data'{route_key = RouteKey, metric = Metric}} body = {data, #'CastFrame.Data'{route_key = RouteKey, metric = Metric}}
}), }),
Packet = <<?FRAME_CAST, CastFrame/binary>>,
case StateName of case StateName of
?STATE_ACTIVATED -> ?STATE_ACTIVATED ->
send_packet(Socket, Packet); send_packet(Socket, [?FRAME_CAST, CastFrame]);
_ -> _ ->
ok = cache_model:insert(Packet) ok = cache_model:insert(<<?FRAME_CAST, CastFrame/binary>>)
end, end,
{keep_state, State}; {keep_state, State};
@ -102,14 +101,14 @@ handle_event(cast, {task_event_stream, TaskId, Type, Stream}, ?STATE_ACTIVATED,
EventPacket = message_pb:encode_msg(#'CastFrame'{ EventPacket = message_pb:encode_msg(#'CastFrame'{
body = {task_event, #'CastFrame.TaskEvent'{task_id = TaskId, type = Type, stream = Stream}} body = {task_event, #'CastFrame.TaskEvent'{task_id = TaskId, type = Type, stream = Stream}}
}), }),
send_packet(Socket, <<?FRAME_CAST, EventPacket/binary>>), send_packet(Socket, [?FRAME_CAST, EventPacket]),
{keep_state, State}; {keep_state, State};
handle_event(cast, {close_task_event_stream, TaskId, Reason}, ?STATE_ACTIVATED, State = #state{socket = Socket}) -> handle_event(cast, {close_task_event_stream, TaskId, Reason}, ?STATE_ACTIVATED, State = #state{socket = Socket}) ->
EventPacket = message_pb:encode_msg(#'CastFrame'{ EventPacket = message_pb:encode_msg(#'CastFrame'{
body = {task_event, #'CastFrame.TaskEvent'{task_id = TaskId, type = <<"close">>, stream = Reason}} body = {task_event, #'CastFrame.TaskEvent'{task_id = TaskId, type = <<"close">>, stream = Reason}}
}), }),
send_packet(Socket, <<?FRAME_CAST, EventPacket/binary>>), send_packet(Socket, [?FRAME_CAST, EventPacket]),
{keep_state, State}; {keep_state, State};
%% %%
@ -126,8 +125,7 @@ handle_event(info, {timeout, _, create_transport}, ?STATE_DISCONNECTED, State =
case connect_socket() of case connect_socket() of
{ok, Socket} -> {ok, Socket} ->
AuthPacket = auth_packet(PacketId), AuthPacket = auth_packet(PacketId),
AuthRequestFrame = <<?FRAME_REQUEST, AuthPacket/binary>>, send_packet(Socket, [?FRAME_REQUEST, AuthPacket]),
send_packet(Socket, AuthRequestFrame),
{next_state, ?STATE_AUTH, State#state{socket = Socket, auth_packet_id = PacketId, next_packet_id = PacketId + 1}, {next_state, ?STATE_AUTH, State#state{socket = Socket, auth_packet_id = PacketId, next_packet_id = PacketId + 1},
[{state_timeout, 5000, auth_timeout}]}; [{state_timeout, 5000, auth_timeout}]};
{error, Reason} -> {error, Reason} ->
@ -226,8 +224,7 @@ handle_event(internal, {decoded_cast, #'CastFrame'{body = {command, #'CastFrame.
{keep_state, State}; {keep_state, State};
{1, _} -> {1, _} ->
AuthPacket = auth_packet(PacketId), AuthPacket = auth_packet(PacketId),
AuthRequestFrame = <<?FRAME_REQUEST, AuthPacket/binary>>, send_packet(Socket, [?FRAME_REQUEST, AuthPacket]),
send_packet(Socket, AuthRequestFrame),
{next_state, ?STATE_AUTH, State#state{auth_packet_id = PacketId, next_packet_id = PacketId + 1}, {next_state, ?STATE_AUTH, State#state{auth_packet_id = PacketId, next_packet_id = PacketId + 1},
[{state_timeout, 5000, auth_timeout}]}; [{state_timeout, 5000, auth_timeout}]};
{0, _} -> {0, _} ->
@ -288,8 +285,8 @@ connect_socket() ->
], ],
ssl:connect(Host, Port, SslOptions, 5000). ssl:connect(Host, Port, SslOptions, 5000).
-spec send_packet(ssl:sslsocket(), binary()) -> ok. -spec send_packet(ssl:sslsocket(), erlang:iodata()) -> ok.
send_packet(Socket, Packet) when is_binary(Packet) -> send_packet(Socket, Packet) ->
ok = ssl:send(Socket, Packet). ok = ssl:send(Socket, Packet).
-spec disconnect(undefined | ssl:sslsocket()) -> ok. -spec disconnect(undefined | ssl:sslsocket()) -> ok.
@ -309,7 +306,7 @@ send_result_reply(Socket, PacketId, Payload) when is_binary(Payload) ->
packet_id = PacketId, packet_id = PacketId,
reply = {result, Payload} reply = {result, Payload}
}), }),
send_packet(Socket, Packet). send_packet(Socket, [?FRAME_REPLY, Packet]).
-spec send_error_reply(ssl:sslsocket(), integer(), binary()) -> ok. -spec send_error_reply(ssl:sslsocket(), integer(), binary()) -> ok.
send_error_reply(Socket, PacketId, Reason) when is_binary(Reason) -> 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, packet_id = PacketId,
reply = {error, #'ReplyFrame.Error'{code = -1, message = Reason}} reply = {error, #'ReplyFrame.Error'{code = -1, message = Reason}}
}), }),
send_packet(Socket, Packet). send_packet(Socket, [?FRAME_REPLY, Packet]).