fix
This commit is contained in:
parent
a46614b77f
commit
0dfd282872
@ -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}).
|
|
||||||
@ -95,7 +95,7 @@ callback_mode() ->
|
|||||||
%% 异步发送数据, 连接存在时候直接发送;否则缓存到mnesia
|
%% 异步发送数据, 连接存在时候直接发送;否则缓存到mnesia
|
||||||
-spec handle_event(term(), term(), atom(), #state{}) -> term().
|
-spec handle_event(term(), term(), atom(), #state{}) -> term().
|
||||||
handle_event(cast, {metric_data, RouteKey, Metric}, StateName, State = #state{socket = Socket}) ->
|
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
|
case StateName of
|
||||||
?STATE_ACTIVATED ->
|
?STATE_ACTIVATED ->
|
||||||
send_packet(Socket, Packet),
|
send_packet(Socket, Packet),
|
||||||
@ -108,12 +108,12 @@ handle_event(cast, {metric_data, RouteKey, Metric}, StateName, State = #state{so
|
|||||||
%% Task的stream流,只做实时的
|
%% Task的stream流,只做实时的
|
||||||
handle_event(cast, {task_event_stream, TaskId, Type, Stream}, ?STATE_ACTIVATED, State = #state{socket = Socket}) ->
|
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]),
|
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),
|
send_packet(Socket, Packet),
|
||||||
{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}) ->
|
||||||
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),
|
send_packet(Socket, Packet),
|
||||||
{keep_state, State};
|
{keep_state, State};
|
||||||
|
|
||||||
@ -162,9 +162,26 @@ handle_event(info, flush_cache, _, State) ->
|
|||||||
{keep_state, State};
|
{keep_state, State};
|
||||||
|
|
||||||
%% 处理收到的ssl消息
|
%% 处理收到的ssl消息
|
||||||
handle_event(info, {ssl, Socket, PacketBin}, _, State = #state{socket = Socket}) ->
|
handle_event(info, {ssl, Socket, PacketBin}, _, State = #state{socket = Socket})
|
||||||
Packet = cloud_wire:decode(PacketBin),
|
when is_binary(PacketBin) ->
|
||||||
{keep_state, State, [{next_event, internal, {decoded_packet, Packet}}]};
|
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}) ->
|
handle_event(info, {ssl_error, Socket, Reason}, _, State = #state{socket = Socket}) ->
|
||||||
logger:debug("[efka_client] ssl error: ~p", [Reason]),
|
logger:debug("[efka_client] ssl error: ~p", [Reason]),
|
||||||
disconnect(Socket),
|
disconnect(Socket),
|
||||||
@ -177,7 +194,7 @@ handle_event(info, {ssl_closed, Socket}, _, State = #state{socket = Socket}) ->
|
|||||||
%%% 处理内部消息,ssl收到的消息会解析成protobuf的消息格式,并按照internal类型处理
|
%%% 处理内部消息,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}) ->
|
?STATE_ACTIVATED, State = #state{socket = Socket}) ->
|
||||||
case docker_container_service:handle_request(Request) of
|
case docker_container_service:handle_request(Request) of
|
||||||
ok ->
|
ok ->
|
||||||
@ -188,36 +205,50 @@ handle_event(internal, {decoded_packet, {request, PacketId, {container_request,
|
|||||||
send_error_reply(Socket, PacketId, Reason)
|
send_error_reply(Socket, PacketId, Reason)
|
||||||
end,
|
end,
|
||||||
{keep_state, State};
|
{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}) ->
|
?STATE_RESTRICTED, State = #state{socket = Socket}) ->
|
||||||
send_error_reply(Socket, PacketId, <<"agent restricted">>),
|
send_error_reply(Socket, PacketId, <<"agent restricted">>),
|
||||||
{keep_state, State};
|
{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}) ->
|
_StateName, State = #state{socket = Socket}) ->
|
||||||
send_error_reply(Socket, PacketId, <<"agent state invalid">>),
|
send_error_reply(Socket, PacketId, <<"agent state invalid">>),
|
||||||
{keep_state, State};
|
{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}) ->
|
?STATE_AUTH, State = #state{auth_packet_id = AuthPacketId}) ->
|
||||||
|
|
||||||
logger:debug("[efka_client] auth success, message: ~p", [Message]),
|
logger:debug("[efka_client] auth success, message: ~p", [Message]),
|
||||||
{next_state, ?STATE_ACTIVATED, State, [{next_event, info, flush_cache}]};
|
{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}) ->
|
?STATE_AUTH, State = #state{auth_packet_id = AuthPacketId}) ->
|
||||||
logger:debug("[efka_client] auth denied, message: ~p", [Message]),
|
logger:debug("[efka_client] auth denied, message: ~p", [Message]),
|
||||||
{next_state, ?STATE_RESTRICTED, State};
|
{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}) ->
|
?STATE_AUTH, State = #state{socket = Socket, auth_packet_id = AuthPacketId}) ->
|
||||||
logger:debug("[efka_client] auth failed, message: ~p", [Message]),
|
logger:debug("[efka_client] auth failed, message: ~p", [Message]),
|
||||||
disconnect(Socket),
|
disconnect(Socket),
|
||||||
schedule_reconnect(),
|
schedule_reconnect(),
|
||||||
{next_state, ?STATE_DISCONNECTED, State#state{socket = undefined}};
|
{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]),
|
logger:warning("[efka_client] ignore unexpected response in state ~p: ~p", [StateName, Reply]),
|
||||||
{keep_state, State};
|
{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}) ->
|
StateName, State = #state{socket = Socket, next_packet_id = PacketId}) ->
|
||||||
|
|
||||||
logger:debug("[efka_client] auth cmd: ~p", [Cmd]),
|
logger:debug("[efka_client] auth cmd: ~p", [Cmd]),
|
||||||
@ -234,10 +265,13 @@ handle_event(internal, {decoded_packet, {message, {auth_control, Cmd}}},
|
|||||||
end;
|
end;
|
||||||
|
|
||||||
%% 处理Pub/Sub机制
|
%% 处理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]),
|
logger:debug("[efka_client] get pub topic: ~p, qos: ~p, content: ~p", [Topic, Qos, Content]),
|
||||||
efka_subscription:publish(Topic, Qos, Content),
|
efka_subscription:publish(Topic, Qos, Content),
|
||||||
{keep_state, State};
|
{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{}) ->
|
handle_event(info, Info, _, State = #state{}) ->
|
||||||
logger:notice("[efka_client] get unknown info: ~p", [Info]),
|
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),
|
Token = proplists:get_value(token, AuthInfo),
|
||||||
|
|
||||||
Timestamp = efka_util:timestamp(),
|
Timestamp = efka_util:timestamp(),
|
||||||
cloud_wire:encode({request, PktId, {auth_request, #{
|
term_to_binary({request, PktId, {auth_request, #{
|
||||||
uuid => list_to_binary(UUID),
|
uuid => list_to_binary(UUID),
|
||||||
token => list_to_binary(Token),
|
token => list_to_binary(Token),
|
||||||
timestamp => Timestamp
|
timestamp => Timestamp
|
||||||
@ -388,10 +422,10 @@ delete_oldest_cache_entry() ->
|
|||||||
|
|
||||||
-spec send_result_reply(ssl:sslsocket(), integer(), binary()) -> ok.
|
-spec send_result_reply(ssl:sslsocket(), integer(), binary()) -> ok.
|
||||||
send_result_reply(Socket, PacketId, Payload) when is_binary(Payload) ->
|
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).
|
send_packet(Socket, 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) ->
|
||||||
Packet = cloud_wire:encode({response, PacketId, {error, -1, Reason}}),
|
Packet = term_to_binary({response, PacketId, {error, -1, Reason}}),
|
||||||
send_packet(Socket, Packet).
|
send_packet(Socket, Packet).
|
||||||
|
|||||||
Loading…
x
Reference in New Issue
Block a user