fix transport protocol
This commit is contained in:
parent
a687abdefd
commit
c9dd1cf90d
@ -35,9 +35,8 @@
|
|||||||
|
|
||||||
-record(state, {
|
-record(state, {
|
||||||
socket :: undefined | ssl:sslsocket(),
|
socket :: undefined | ssl:sslsocket(),
|
||||||
next_packet_id = 1,
|
%% 保存当前auth请求的ref,用来建立auth请求和响应的对应关系
|
||||||
%% 保存当前auth请求的packet_id,用来建立auth请求和响应的对应关系
|
auth_ref = undefined :: undefined | reference(),
|
||||||
auth_packet_id = 1,
|
|
||||||
dropped_message_count = 0 :: non_neg_integer()
|
dropped_message_count = 0 :: non_neg_integer()
|
||||||
}).
|
}).
|
||||||
|
|
||||||
@ -129,12 +128,13 @@ handle_event({call, From}, dropped_message_count, _StateName, State = #state{dro
|
|||||||
{keep_state, State, [{reply, From, DroppedCount}]};
|
{keep_state, State, [{reply, From, DroppedCount}]};
|
||||||
|
|
||||||
%% 异步建立到服务器的连接
|
%% 异步建立到服务器的连接
|
||||||
handle_event(info, {timeout, _, create_transport}, ?STATE_DISCONNECTED, State = #state{next_packet_id = PacketId}) ->
|
handle_event(info, {timeout, _, create_transport}, ?STATE_DISCONNECTED, State = #state{}) ->
|
||||||
case connect_socket() of
|
case connect_socket() of
|
||||||
{ok, Socket} ->
|
{ok, Socket} ->
|
||||||
AuthPacket = auth_packet(PacketId),
|
Ref = make_ref(),
|
||||||
|
AuthPacket = auth_packet(Ref),
|
||||||
ok = ssl:send(Socket, AuthPacket),
|
ok = ssl:send(Socket, AuthPacket),
|
||||||
{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_ref = Ref},
|
||||||
[{state_timeout, 5000, auth_timeout}]};
|
[{state_timeout, 5000, auth_timeout}]};
|
||||||
{error, _Reason} ->
|
{error, _Reason} ->
|
||||||
%logger:debug("[efka_client] connect failed"),
|
%logger:debug("[efka_client] connect failed"),
|
||||||
@ -146,7 +146,7 @@ handle_event(state_timeout, auth_timeout, ?STATE_AUTH, State = #state{socket = S
|
|||||||
logger:debug("[efka_client] auth request timeout"),
|
logger:debug("[efka_client] auth request timeout"),
|
||||||
disconnect(Socket),
|
disconnect(Socket),
|
||||||
schedule_reconnect(),
|
schedule_reconnect(),
|
||||||
{next_state, ?STATE_DISCONNECTED, State#state{socket = undefined}};
|
{next_state, ?STATE_DISCONNECTED, State#state{socket = undefined, auth_ref = undefined}};
|
||||||
|
|
||||||
%% 将缓存中的数据推送到服务器端
|
%% 将缓存中的数据推送到服务器端
|
||||||
handle_event(info, flush_cache, ?STATE_ACTIVATED, State = #state{socket = Socket}) ->
|
handle_event(info, flush_cache, ?STATE_ACTIVATED, State = #state{socket = Socket}) ->
|
||||||
@ -169,89 +169,90 @@ handle_event(info, {ssl_error, Socket, Reason}, _, State = #state{socket = Socke
|
|||||||
logger:debug("[efka_client] ssl error: ~p", [Reason]),
|
logger:debug("[efka_client] ssl error: ~p", [Reason]),
|
||||||
disconnect(Socket),
|
disconnect(Socket),
|
||||||
schedule_reconnect(),
|
schedule_reconnect(),
|
||||||
{next_state, ?STATE_DISCONNECTED, State#state{socket = undefined}};
|
{next_state, ?STATE_DISCONNECTED, State#state{socket = undefined, auth_ref = undefined}};
|
||||||
handle_event(info, {ssl_closed, Socket}, _, State = #state{socket = Socket}) ->
|
handle_event(info, {ssl_closed, Socket}, _, State = #state{socket = Socket}) ->
|
||||||
schedule_reconnect(),
|
schedule_reconnect(),
|
||||||
{next_state, ?STATE_DISCONNECTED, State#state{socket = undefined}};
|
{next_state, ?STATE_DISCONNECTED, State#state{socket = undefined, auth_ref = undefined}};
|
||||||
|
|
||||||
%%% 处理内部消息,ssl收到的消息会先 binary_to_term,再由这里按协议结构模式匹配
|
%%% 处理内部消息,ssl收到的消息会先 binary_to_term,再由这里按协议结构模式匹配
|
||||||
|
|
||||||
%% 容器管理请求
|
%% 容器管理请求
|
||||||
handle_event(internal, {request, PacketId, {container_request, #{action := list}}},
|
handle_event(internal, {request, Ref, {container_request, #{action := list}}},
|
||||||
?STATE_ACTIVATED, State = #state{socket = Socket}) ->
|
?STATE_ACTIVATED, State = #state{socket = Socket}) ->
|
||||||
Reply = docker_commands:get_containers(),
|
Reply = docker_commands:get_containers(),
|
||||||
handle_container_reply(Socket, PacketId, Reply),
|
handle_container_reply(Socket, Ref, Reply),
|
||||||
{keep_state, State};
|
{keep_state, State};
|
||||||
handle_event(internal, {request, PacketId, {container_request, #{action := deploy, task_id := TaskId, params := Params}}},
|
handle_event(internal, {request, Ref, {container_request, #{action := deploy, task_id := TaskId, params := Params}}},
|
||||||
?STATE_ACTIVATED, State = #state{socket = Socket}) ->
|
?STATE_ACTIVATED, State = #state{socket = Socket}) ->
|
||||||
Reply = docker_deploy_manager:deploy(TaskId, Params),
|
Reply = docker_deploy_manager:deploy(TaskId, Params),
|
||||||
handle_container_reply(Socket, PacketId, Reply),
|
handle_container_reply(Socket, Ref, Reply),
|
||||||
{keep_state, State};
|
{keep_state, State};
|
||||||
handle_event(internal, {request, PacketId, {container_request, #{action := start, target := Target}}},
|
handle_event(internal, {request, Ref, {container_request, #{action := start, target := Target}}},
|
||||||
?STATE_ACTIVATED, State = #state{socket = Socket}) ->
|
?STATE_ACTIVATED, State = #state{socket = Socket}) ->
|
||||||
Reply = docker_commands:start_container(container_target(Target)),
|
Reply = docker_commands:start_container(container_target(Target)),
|
||||||
handle_container_reply(Socket, PacketId, Reply),
|
handle_container_reply(Socket, Ref, Reply),
|
||||||
{keep_state, State};
|
{keep_state, State};
|
||||||
handle_event(internal, {request, PacketId, {container_request, #{action := stop, target := Target, timeout_seconds := TimeoutSeconds}}},
|
handle_event(internal, {request, Ref, {container_request, #{action := stop, target := Target, timeout_seconds := TimeoutSeconds}}},
|
||||||
?STATE_ACTIVATED, State = #state{socket = Socket}) ->
|
?STATE_ACTIVATED, State = #state{socket = Socket}) ->
|
||||||
Reply = docker_commands:stop_container(container_target(Target), TimeoutSeconds),
|
Reply = docker_commands:stop_container(container_target(Target), TimeoutSeconds),
|
||||||
handle_container_reply(Socket, PacketId, Reply),
|
handle_container_reply(Socket, Ref, Reply),
|
||||||
{keep_state, State};
|
{keep_state, State};
|
||||||
handle_event(internal, {request, PacketId, {container_request, #{action := kill, target := Target, signal := Signal}}},
|
handle_event(internal, {request, Ref, {container_request, #{action := kill, target := Target, signal := Signal}}},
|
||||||
?STATE_ACTIVATED, State = #state{socket = Socket}) ->
|
?STATE_ACTIVATED, State = #state{socket = Socket}) ->
|
||||||
Reply = docker_commands:kill_container(container_target(Target), to_binary(Signal)),
|
Reply = docker_commands:kill_container(container_target(Target), to_binary(Signal)),
|
||||||
handle_container_reply(Socket, PacketId, Reply),
|
handle_container_reply(Socket, Ref, Reply),
|
||||||
{keep_state, State};
|
{keep_state, State};
|
||||||
handle_event(internal, {request, PacketId, {container_request, #{action := remove, target := Target, force := Force, remove_volumes := RemoveVolumes}}},
|
handle_event(internal, {request, Ref, {container_request, #{action := remove, target := Target, force := Force, remove_volumes := RemoveVolumes}}},
|
||||||
?STATE_ACTIVATED, State = #state{socket = Socket}) ->
|
?STATE_ACTIVATED, State = #state{socket = Socket}) ->
|
||||||
Reply = docker_commands:remove_container(container_target(Target), to_bool(Force), to_bool(RemoveVolumes)),
|
Reply = docker_commands:remove_container(container_target(Target), to_bool(Force), to_bool(RemoveVolumes)),
|
||||||
handle_container_reply(Socket, PacketId, Reply),
|
handle_container_reply(Socket, Ref, Reply),
|
||||||
{keep_state, State};
|
{keep_state, State};
|
||||||
handle_event(internal, {request, PacketId, {container_request, #{action := config, target := Target, config := Config}}},
|
handle_event(internal, {request, Ref, {container_request, #{action := config, target := Target, config := Config}}},
|
||||||
?STATE_ACTIVATED, State = #state{socket = Socket}) ->
|
?STATE_ACTIVATED, State = #state{socket = Socket}) ->
|
||||||
Reply = update_container_config(container_target(Target), iolist_to_binary(Config)),
|
Reply = update_container_config(container_target(Target), iolist_to_binary(Config)),
|
||||||
handle_container_reply(Socket, PacketId, Reply),
|
handle_container_reply(Socket, Ref, Reply),
|
||||||
{keep_state, State};
|
{keep_state, State};
|
||||||
handle_event(internal, {request, PacketId, {container_request, Request}}, ?STATE_RESTRICTED, State = #state{socket = Socket}) ->
|
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]),
|
logger:notice("[efka_client] get a invalid request: ~p, agent restricted", [Request]),
|
||||||
handle_container_reply(Socket, PacketId, {error, <<"agent restricted">>}),
|
handle_container_reply(Socket, Ref, {error, <<"agent restricted">>}),
|
||||||
{keep_state, State};
|
{keep_state, State};
|
||||||
handle_event(internal, {request, PacketId, {container_request, Request}}, _StateName, State = #state{socket = Socket}) ->
|
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]),
|
logger:notice("[efka_client] get a invalid request: ~p, agent restricted", [Request]),
|
||||||
handle_container_reply(Socket, PacketId, {error, <<"agent invalid">>}),
|
handle_container_reply(Socket, Ref, {error, <<"agent invalid">>}),
|
||||||
{keep_state, State};
|
{keep_state, State};
|
||||||
|
|
||||||
%% 处理response
|
%% 处理response
|
||||||
handle_event(internal, {response, AuthPacketId, {auth_response, {ok, Message}}},
|
handle_event(internal, {response, AuthRef, {auth_response, {ok, Message}}},
|
||||||
?STATE_AUTH, State = #state{auth_packet_id = AuthPacketId}) ->
|
?STATE_AUTH, State = #state{auth_ref = AuthRef}) ->
|
||||||
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#state{auth_ref = undefined}, [{next_event, info, flush_cache}]};
|
||||||
handle_event(internal, {response, AuthPacketId, {auth_response, {error, 1, Message}}},
|
handle_event(internal, {response, AuthRef, {auth_response, {error, {denied, Message}}}},
|
||||||
?STATE_AUTH, State = #state{auth_packet_id = AuthPacketId}) ->
|
?STATE_AUTH, State = #state{auth_ref = AuthRef}) ->
|
||||||
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#state{auth_ref = undefined}};
|
||||||
handle_event(internal, {response, AuthPacketId, {auth_response, {error, _Code, Message}}},
|
handle_event(internal, {response, AuthRef, {auth_response, {error, Reason}}},
|
||||||
?STATE_AUTH, State = #state{socket = Socket, auth_packet_id = AuthPacketId}) ->
|
?STATE_AUTH, State = #state{socket = Socket, auth_ref = AuthRef}) ->
|
||||||
logger:debug("[efka_client] auth failed, message: ~p", [Message]),
|
logger:debug("[efka_client] auth failed, reason: ~p", [Reason]),
|
||||||
disconnect(Socket),
|
disconnect(Socket),
|
||||||
schedule_reconnect(),
|
schedule_reconnect(),
|
||||||
{next_state, ?STATE_DISCONNECTED, State#state{socket = undefined}};
|
{next_state, ?STATE_DISCONNECTED, State#state{socket = undefined, auth_ref = undefined}};
|
||||||
handle_event(internal, {response, _PacketId, Reply}, StateName, State) ->
|
handle_event(internal, {response, _Ref, 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, {message, {auth_control, Cmd}},
|
handle_event(internal, {message, {auth_control, Cmd}},
|
||||||
StateName, State = #state{socket = Socket, next_packet_id = PacketId}) ->
|
StateName, State = #state{socket = Socket}) ->
|
||||||
|
|
||||||
logger:debug("[efka_client] auth cmd: ~p", [Cmd]),
|
logger:debug("[efka_client] auth cmd: ~p", [Cmd]),
|
||||||
case {Cmd, StateName} of
|
case {Cmd, StateName} of
|
||||||
{activate, ?STATE_ACTIVATED} ->
|
{activate, ?STATE_ACTIVATED} ->
|
||||||
{keep_state, State};
|
{keep_state, State};
|
||||||
{activate, _} ->
|
{activate, _} ->
|
||||||
AuthPacket = auth_packet(PacketId),
|
Ref = make_ref(),
|
||||||
|
AuthPacket = auth_packet(Ref),
|
||||||
ok = ssl:send(Socket, AuthPacket),
|
ok = ssl:send(Socket, AuthPacket),
|
||||||
{next_state, ?STATE_AUTH, State#state{auth_packet_id = PacketId, next_packet_id = PacketId + 1},
|
{next_state, ?STATE_AUTH, State#state{auth_ref = Ref},
|
||||||
[{state_timeout, 5000, auth_timeout}]};
|
[{state_timeout, 5000, auth_timeout}]};
|
||||||
{deactivate, _} ->
|
{deactivate, _} ->
|
||||||
{next_state, ?STATE_RESTRICTED, State}
|
{next_state, ?STATE_RESTRICTED, State}
|
||||||
@ -284,14 +285,14 @@ code_change(_OldVsn, StateName, State = #state{}, _Extra) ->
|
|||||||
%%% Internal functions
|
%%% Internal functions
|
||||||
%%%===================================================================
|
%%%===================================================================
|
||||||
|
|
||||||
-spec auth_packet(PktId :: integer()) -> binary().
|
-spec auth_packet(reference()) -> binary().
|
||||||
auth_packet(PktId) when is_integer(PktId) ->
|
auth_packet(Ref) when is_reference(Ref) ->
|
||||||
{ok, AuthInfo} = application:get_env(efka, auth),
|
{ok, AuthInfo} = application:get_env(efka, auth),
|
||||||
UUID = proplists:get_value(uuid, AuthInfo),
|
UUID = proplists:get_value(uuid, AuthInfo),
|
||||||
Token = proplists:get_value(token, AuthInfo),
|
Token = proplists:get_value(token, AuthInfo),
|
||||||
|
|
||||||
Timestamp = efka_util:timestamp(),
|
Timestamp = efka_util:timestamp(),
|
||||||
term_to_binary({request, PktId, {auth_request, #{
|
term_to_binary({request, Ref, {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
|
||||||
@ -409,15 +410,15 @@ delete_oldest_cache_entry() ->
|
|||||||
dets:delete(?CACHE_TAB, Key)
|
dets:delete(?CACHE_TAB, Key)
|
||||||
end.
|
end.
|
||||||
|
|
||||||
-spec handle_container_reply(ssl:sslsocket(), integer(), term()) -> ok.
|
-spec handle_container_reply(ssl:sslsocket(), reference(), term()) -> ok.
|
||||||
handle_container_reply(Socket, PacketId, ok) ->
|
handle_container_reply(Socket, Ref, ok) ->
|
||||||
Packet = term_to_binary({response, PacketId, {container_response, {ok, <<"ok">>}}}),
|
Packet = term_to_binary({response, Ref, {container_response, {ok, <<"ok">>}}}),
|
||||||
ok = ssl:send(Socket, Packet);
|
ok = ssl:send(Socket, Packet);
|
||||||
handle_container_reply(Socket, PacketId, {ok, Response}) ->
|
handle_container_reply(Socket, Ref, {ok, Response}) ->
|
||||||
Packet = term_to_binary({response, PacketId, {container_response, {ok, Response}}}),
|
Packet = term_to_binary({response, Ref, {container_response, {ok, Response}}}),
|
||||||
ok = ssl:send(Socket, Packet);
|
ok = ssl:send(Socket, Packet);
|
||||||
handle_container_reply(Socket, PacketId, {error, Reason}) ->
|
handle_container_reply(Socket, Ref, {error, Reason}) ->
|
||||||
Packet = term_to_binary({response, PacketId, {container_response, {error, -1, Reason}}}),
|
Packet = term_to_binary({response, Ref, {container_response, {error, Reason}}}),
|
||||||
ok = ssl:send(Socket, Packet).
|
ok = ssl:send(Socket, Packet).
|
||||||
|
|
||||||
-spec container_target(map()) -> binary().
|
-spec container_target(map()) -> binary().
|
||||||
|
|||||||
Loading…
x
Reference in New Issue
Block a user