From 863156667c12518ba63728f553896a4c300fa7df Mon Sep 17 00:00:00 2001
From: anlicheng <244108715@qq.com>
Date: Mon, 20 Apr 2026 14:18:26 +0800
Subject: [PATCH] fix
---
src/docker/docker_commands.erl | 119 +++++++++---
src/docker/docker_helper.erl | 53 ++++--
src/docker/docker_manager.erl | 36 ++--
src/efka_remote_agent.erl | 319 ++++++++++++++++++++++++++++-----
4 files changed, 441 insertions(+), 86 deletions(-)
diff --git a/src/docker/docker_commands.erl b/src/docker/docker_commands.erl
index 21da27a..5448963 100644
--- a/src/docker/docker_commands.erl
+++ b/src/docker/docker_commands.erl
@@ -12,7 +12,7 @@
%% API
-export([pull_image/2, check_image_exist/1]).
-export([create_container/3, check_container_exist/1, is_container_running/1,
- start_container/1, stop_container/1, remove_container/1, kill_container/1,
+ start_container/1, stop_container/1, stop_container/2, remove_container/1, remove_container/3, kill_container/1, kill_container/2,
get_containers/0]).
-spec pull_image(Image :: binary(), Callback :: fun((Msg :: any()) -> no_return())) -> ok | {error, Reason :: any()}.
@@ -112,7 +112,11 @@ start_container(ContainerName) when is_binary(ContainerName) ->
-spec stop_container(ContainerName :: binary()) -> ok | {error, Reason :: binary()}.
stop_container(ContainerName) when is_binary(ContainerName) ->
- Url = lists:flatten(io_lib:format("/containers/~s/stop", [binary_to_list(ContainerName)])),
+ stop_container(ContainerName, 0).
+
+-spec stop_container(ContainerName :: binary(), TimeoutSeconds :: non_neg_integer()) -> ok | {error, Reason :: binary()}.
+stop_container(ContainerName, TimeoutSeconds) when is_binary(ContainerName), is_integer(TimeoutSeconds), TimeoutSeconds >= 0 ->
+ Url = build_stop_container_url(ContainerName, TimeoutSeconds),
Headers = [
{<<"Content-Type">>, <<"application/json">>}
],
@@ -134,7 +138,11 @@ stop_container(ContainerName) when is_binary(ContainerName) ->
-spec kill_container(ContainerName :: binary()) -> ok | {error, Reason :: binary()}.
kill_container(ContainerName) when is_binary(ContainerName) ->
- Url = lists:flatten(io_lib:format("/containers/~s/kill", [binary_to_list(ContainerName)])),
+ kill_container(ContainerName, <<>>).
+
+-spec kill_container(ContainerName :: binary(), Signal :: binary()) -> ok | {error, Reason :: binary()}.
+kill_container(ContainerName, Signal) when is_binary(ContainerName), is_binary(Signal) ->
+ Url = build_kill_container_url(ContainerName, Signal),
Headers = [
{<<"Content-Type">>, <<"application/json">>}
],
@@ -154,7 +162,12 @@ kill_container(ContainerName) when is_binary(ContainerName) ->
-spec remove_container(ContainerName :: binary()) -> ok | {error, Reason :: binary()}.
remove_container(ContainerName) when is_binary(ContainerName) ->
- Url = lists:flatten(io_lib:format("/containers/~s", [binary_to_list(ContainerName)])),
+ remove_container(ContainerName, false, false).
+
+-spec remove_container(ContainerName :: binary(), Force :: boolean(), RemoveVolumes :: boolean()) -> ok | {error, Reason :: binary()}.
+remove_container(ContainerName, Force, RemoveVolumes)
+ when is_binary(ContainerName), is_boolean(Force), is_boolean(RemoveVolumes) ->
+ Url = build_remove_container_url(ContainerName, Force, RemoveVolumes),
Headers = [
{<<"Content-Type">>, <<"application/json">>}
],
@@ -264,7 +277,7 @@ build_expose(Config) ->
case Ports of
[] -> #{};
_ ->
- maps:from_list([{<
>, #{}} || P <- Ports])
+ maps:from_list([{normalize_expose_port(P), #{}} || P <- Ports])
end.
build_volumes(Config) ->
@@ -274,7 +287,7 @@ build_volumes(Config) ->
#{};
_ ->
maps:from_list(lists:map(fun(V) ->
- [_Host, Cont] = binary:split(V, <<":">>, []),
+ [_Host, Cont | _Modes] = binary:split(V, <<":">>, [global]),
{Cont, #{}}
end, Vols))
end.
@@ -328,17 +341,25 @@ build_healthcheck(Config) ->
end.
parse_duration(Bin) ->
- %% "30s" -> 30000000000
- Sz = byte_size(Bin),
- NBin = binary:part(Bin, {0, Sz-1}),
- N = list_to_integer(binary_to_list(NBin)),
- case binary:last(Bin) of
- $s ->
- N * 1000000000;
- $m ->
- N * 60000000000;
- _ ->
- N
+ case re:run(Bin, <<"^(\\d+)(ns|us|ms|s|m|h)?$">>, [{capture, all_but_first, binary}]) of
+ {match, [NumberBin, Unit]} ->
+ Number = binary_to_integer(NumberBin),
+ case Unit of
+ <<"ns">> ->
+ Number;
+ <<"us">> ->
+ Number * 1000;
+ <<"ms">> ->
+ Number * 1000000;
+ <<"s">> ->
+ Number * 1000000000;
+ <<"m">> ->
+ Number * 60000000000;
+ <<"h">> ->
+ Number * 3600000000000
+ end;
+ {match, [NumberBin]} ->
+ binary_to_integer(NumberBin)
end.
%% --- 构建子字段 ---
@@ -348,7 +369,15 @@ build_restart(Config) ->
undefined ->
#{};
Policy ->
- #{<<"RestartPolicy">> => #{<<"Name">> => Policy}}
+ case binary:split(Policy, <<":">>) of
+ [Name, RetryCountBin] ->
+ #{<<"RestartPolicy">> => #{
+ <<"Name">> => Name,
+ <<"MaximumRetryCount">> => binary_to_integer(RetryCountBin)
+ }};
+ [Name] ->
+ #{<<"RestartPolicy">> => #{<<"Name">> => Name}}
+ end
end.
build_privileged(Config) ->
@@ -376,9 +405,9 @@ build_devices(Config) ->
#{};
_ ->
DevObjs = [#{<<"PathOnHost">> => H, <<"PathInContainer">> => C,
- <<"CgroupPermissions">> => <<"rwm">>}
+ <<"CgroupPermissions">> => P}
|| D <- Devs,
- [H, C] <- [binary:split(D, <<":">>, [])]],
+ {H, C, P} <- [parse_device_mapping(D)]],
#{<<"Devices">> => DevObjs}
end.
@@ -448,7 +477,7 @@ build_tmpfs(Config) ->
[] ->
#{};
_ ->
- #{<<"Tmpfs">> => maps:from_list([{T, <<>>} || T <- Tmp])}
+ #{<<"Tmpfs">> => maps:from_list([parse_tmpfs_mount(T) || T <- Tmp])}
end.
build_extra_hosts(Config) ->
@@ -464,3 +493,51 @@ build_extra_hosts(Config) ->
display_options(Options) when is_map(Options) ->
logger:debug("deploy options: ~p", [jiffy:encode(Options, [force_utf8])]),
lists:foreach(fun({K, V}) -> logger:debug("~p => ~p", [K, V]) end, maps:to_list(Options)).
+
+normalize_expose_port(Port) when is_binary(Port) ->
+ case binary:match(Port, <<"/">>) of
+ nomatch ->
+ <>;
+ _ ->
+ Port
+ end.
+
+parse_device_mapping(Device) when is_binary(Device) ->
+ case binary:split(Device, <<":">>, [global]) of
+ [HostPath, ContainerPath] ->
+ {HostPath, ContainerPath, <<"rwm">>};
+ [HostPath, ContainerPath, Permissions] ->
+ {HostPath, ContainerPath, Permissions}
+ end.
+
+parse_tmpfs_mount(Tmpfs) when is_binary(Tmpfs) ->
+ case binary:split(Tmpfs, <<":">>) of
+ [Path] ->
+ {Path, <<>>};
+ [Path, Options] ->
+ {Path, Options}
+ end.
+
+build_stop_container_url(ContainerName, 0) ->
+ lists:flatten(io_lib:format("/containers/~s/stop", [binary_to_list(ContainerName)]));
+build_stop_container_url(ContainerName, TimeoutSeconds) ->
+ lists:flatten(io_lib:format("/containers/~s/stop?t=~B", [binary_to_list(ContainerName), TimeoutSeconds])).
+
+build_kill_container_url(ContainerName, <<>>) ->
+ lists:flatten(io_lib:format("/containers/~s/kill", [binary_to_list(ContainerName)]));
+build_kill_container_url(ContainerName, Signal) ->
+ lists:flatten(io_lib:format("/containers/~s/kill?signal=~s", [binary_to_list(ContainerName), binary_to_list(Signal)])).
+
+build_remove_container_url(ContainerName, Force, RemoveVolumes) ->
+ ForceValue = boolean_to_query_value(Force),
+ RemoveVolumesValue = boolean_to_query_value(RemoveVolumes),
+ lists:flatten(io_lib:format("/containers/~s?force=~s&v=~s", [
+ binary_to_list(ContainerName),
+ ForceValue,
+ RemoveVolumesValue
+ ])).
+
+boolean_to_query_value(true) ->
+ "true";
+boolean_to_query_value(false) ->
+ "false".
diff --git a/src/docker/docker_helper.erl b/src/docker/docker_helper.erl
index 1df50da..bfe10d7 100644
--- a/src/docker/docker_helper.erl
+++ b/src/docker/docker_helper.erl
@@ -10,14 +10,22 @@
-author("anlicheng").
%% API
--export([ensure_container_dir/2, get_container_dir/2, get_config_file/1]).
+-export([ensure_container_dir/2, ensure_container_dir/3, get_container_dir/2, get_config_file/1]).
-spec ensure_container_dir(RootDir :: string(), ContainerName :: binary()) -> {ok, ServerRootDir :: string()}.
ensure_container_dir(RootDir, ContainerName) when is_list(RootDir), is_binary(ContainerName) ->
- %% 根目录
- ContainerRootDir = RootDir ++ "/" ++ binary_to_list(ContainerName) ++ "/",
- ok = filelib:ensure_dir(ContainerRootDir),
- {ok, ContainerRootDir}.
+ ensure_container_dir(RootDir, ContainerName, <<>>).
+
+-spec ensure_container_dir(RootDir :: string(), ContainerName :: binary(), ContainerDir :: binary()) ->
+ {ok, ServerRootDir :: string()}.
+ensure_container_dir(RootDir, ContainerName, ContainerDir) when is_list(RootDir), is_binary(ContainerName), is_binary(ContainerDir) ->
+ PointerDir = default_container_dir(RootDir, ContainerName),
+ PointerFile = container_dir_pointer_file(PointerDir),
+ ActualContainerDir = resolve_container_dir(RootDir, ContainerName, ContainerDir),
+ ok = filelib:ensure_dir(PointerDir),
+ ok = filelib:ensure_dir(ActualContainerDir),
+ ok = file:write_file(PointerFile, unicode:characters_to_binary(ActualContainerDir), [write]),
+ {ok, ActualContainerDir}.
-spec get_config_file(ContainerDir :: string()) -> ConfigFile :: string().
get_config_file(ContainerDir) when is_list(ContainerDir) ->
@@ -26,11 +34,36 @@ get_config_file(ContainerDir) when is_list(ContainerDir) ->
-spec get_container_dir(RootDir :: string(), ContainerName :: binary()) -> {ok, ServerRootDir :: string()} | error.
get_container_dir(RootDir, ContainerName) when is_list(RootDir), is_binary(ContainerName) ->
- %% 根目录
- ContainerRootDir = RootDir ++ "/" ++ binary_to_list(ContainerName) ++ "/",
- case filelib:is_dir(ContainerRootDir) of
+ ContainerRootDir = default_container_dir(RootDir, ContainerName),
+ PointerFile = container_dir_pointer_file(ContainerRootDir),
+ ResolvedContainerDir = case file:read_file(PointerFile) of
+ {ok, ActualContainerDirBin} ->
+ normalize_container_dir(binary_to_list(ActualContainerDirBin));
+ {error, _} ->
+ ContainerRootDir
+ end,
+ case filelib:is_dir(ResolvedContainerDir) of
true ->
- {ok, ContainerRootDir};
+ {ok, ResolvedContainerDir};
false ->
error
- end.
\ No newline at end of file
+ end.
+
+default_container_dir(RootDir, ContainerName) ->
+ normalize_container_dir(RootDir ++ "/" ++ binary_to_list(ContainerName)).
+
+container_dir_pointer_file(ContainerDir) ->
+ ContainerDir ++ ".container_dir".
+
+resolve_container_dir(RootDir, ContainerName, <<>>) ->
+ default_container_dir(RootDir, ContainerName);
+resolve_container_dir(_RootDir, _ContainerName, ContainerDir) ->
+ normalize_container_dir(binary_to_list(ContainerDir)).
+
+normalize_container_dir(ContainerDir) ->
+ case lists:last(ContainerDir) of
+ $/ ->
+ ContainerDir;
+ _ ->
+ ContainerDir ++ "/"
+ end.
diff --git a/src/docker/docker_manager.erl b/src/docker/docker_manager.erl
index c74b96d..29186de 100644
--- a/src/docker/docker_manager.erl
+++ b/src/docker/docker_manager.erl
@@ -15,7 +15,7 @@
%% API
-export([start_link/0]).
--export([deploy/2, start_container/1, stop_container/1, config_container/2, kill_container/1, remove_container/1]).
+-export([deploy/2, start_container/1, stop_container/1, stop_container/2, config_container/2, kill_container/1, kill_container/2, remove_container/1, remove_container/3]).
-export([get_containers/0]).
%% gen_server callbacks
@@ -51,15 +51,28 @@ start_container(ContainerId) when is_binary(ContainerId) ->
-spec stop_container(ServiceId :: binary()) -> ok | {error, Reason :: term()}.
stop_container(ContainerId) when is_binary(ContainerId) ->
- gen_server:call(?SERVER, {stop_container, ContainerId}).
+ stop_container(ContainerId, 0).
+
+-spec stop_container(ServiceId :: binary(), TimeoutSeconds :: non_neg_integer()) -> ok | {error, Reason :: term()}.
+stop_container(ContainerId, TimeoutSeconds) when is_binary(ContainerId), is_integer(TimeoutSeconds), TimeoutSeconds >= 0 ->
+ gen_server:call(?SERVER, {stop_container, ContainerId, TimeoutSeconds}).
-spec kill_container(ServiceId :: binary()) -> ok | {error, Reason :: term()}.
kill_container(ContainerId) when is_binary(ContainerId) ->
- gen_server:call(?SERVER, {kill_container, ContainerId}).
+ kill_container(ContainerId, <<>>).
+
+-spec kill_container(ServiceId :: binary(), Signal :: binary()) -> ok | {error, Reason :: term()}.
+kill_container(ContainerId, Signal) when is_binary(ContainerId), is_binary(Signal) ->
+ gen_server:call(?SERVER, {kill_container, ContainerId, Signal}).
-spec remove_container(ServiceId :: binary()) -> ok | {error, Reason :: term()}.
remove_container(ContainerId) when is_binary(ContainerId) ->
- gen_server:call(?SERVER, {remove_container, ContainerId}).
+ remove_container(ContainerId, false, false).
+
+-spec remove_container(ServiceId :: binary(), Force :: boolean(), RemoveVolumes :: boolean()) -> ok | {error, Reason :: term()}.
+remove_container(ContainerId, Force, RemoveVolumes)
+ when is_binary(ContainerId), is_boolean(Force), is_boolean(RemoveVolumes) ->
+ gen_server:call(?SERVER, {remove_container, ContainerId, Force, RemoveVolumes}).
%% @doc Spawns the server and registers the local name (unique)
-spec(start_link() ->
@@ -93,7 +106,8 @@ init([]) ->
{stop, Reason :: term(), NewState :: #state{}}).
handle_call({deploy, TaskId, Config = #{<<"container_name">> := ContainerName}}, _From, State = #state{root_dir = RootDir, task_map = TaskMap}) ->
%% 创建目录
- {ok, ContainerDir} = docker_helper:ensure_container_dir(RootDir, ContainerName),
+ ContainerDir0 = maps:get(<<"container_dir">>, Config, <<>>),
+ {ok, ContainerDir} = docker_helper:ensure_container_dir(RootDir, ContainerName, ContainerDir0),
{ok, {TaskPid, _Ref}} = docker_deployer:start_monitor(TaskId, ContainerDir, Config),
logger:debug("[docker_manager] start deploy task_id: ~p, config: ~p", [TaskId, Config]),
{reply, ok, State#state{task_map = maps:put(TaskPid, TaskId, TaskMap)}};
@@ -127,8 +141,8 @@ handle_call({start_container, ContainerId}, _From, State) ->
%% 停止服务, 主动停止的时候会改变服务配置的status字段
-handle_call({stop_container, ContainerId}, _From, State = #state{}) ->
- case docker_commands:stop_container(ContainerId) of
+handle_call({stop_container, ContainerId, TimeoutSeconds}, _From, State = #state{}) ->
+ case docker_commands:stop_container(ContainerId, TimeoutSeconds) of
ok ->
{reply, ok, State};
{error, Reason} ->
@@ -136,8 +150,8 @@ handle_call({stop_container, ContainerId}, _From, State = #state{}) ->
end;
%% 停止服务, 主动停止的时候会改变服务配置的status字段
-handle_call({kill_container, ContainerId}, _From, State = #state{}) ->
- case docker_commands:kill_container(ContainerId) of
+handle_call({kill_container, ContainerId, Signal}, _From, State = #state{}) ->
+ case docker_commands:kill_container(ContainerId, Signal) of
ok ->
{reply, ok, State};
{error, Reason} ->
@@ -154,8 +168,8 @@ handle_call(get_containers, _From, State = #state{}) ->
end;
%% 停止服务, 主动停止的时候会改变服务配置的status字段
-handle_call({remove_container, ContainerId}, _From, State = #state{}) ->
- case docker_commands:remove_container(ContainerId) of
+handle_call({remove_container, ContainerId, Force, RemoveVolumes}, _From, State = #state{}) ->
+ case docker_commands:remove_container(ContainerId, Force, RemoveVolumes) of
ok ->
{reply, ok, State};
{error, Reason} ->
diff --git a/src/efka_remote_agent.erl b/src/efka_remote_agent.erl
index 94cd9a8..ad444c7 100644
--- a/src/efka_remote_agent.erl
+++ b/src/efka_remote_agent.erl
@@ -201,23 +201,22 @@ handle_event(info, flush_cache, ?STATE_ACTIVATED, State = #state{transport_pid =
handle_event(info, flush_cache, _, State) ->
{keep_state, State};
-handle_event(info, {server_packet, <<8, _/binary>> = PacketBin}, ?STATE_ACTIVATED, State) ->
- #'RequestFrame'{packet_id = PacketId, body = {rpc_request, Request}} =
- message_pb:decode_msg(PacketBin, 'RequestFrame'),
+handle_event(info, {server_packet, <<8, _/binary>> = PacketBin}, StateName, State)
+ when StateName =:= ?STATE_ACTIVATED; StateName =:= ?STATE_RESTRICTED ->
+ #'RequestFrame'{packet_id = PacketId, body = Body} = message_pb:decode_msg(PacketBin, 'RequestFrame'),
true = is_integer(PacketId) andalso PacketId > 0,
- {keep_state, State, [{next_event, info, {server_rpc, PacketId, Request}}]};
+ case Body of
+ {rpc_request, Request} ->
+ {keep_state, State, [{next_event, info, {server_rpc, PacketId, Request}}]};
+ {container_request, Request} ->
+ {keep_state, State, [{next_event, info, {container_request, PacketId, Request}}]}
+ end;
handle_event(info, {server_packet, <<10, _/binary>> = PacketBin}, ?STATE_ACTIVATED, State) ->
#'CastFrame'{body = {pub, Pub}} = message_pb:decode_msg(PacketBin, 'CastFrame'),
{keep_state, State, [{next_event, info, {server_cast, Pub}}]};
handle_event(info, {server_packet, <<18, _/binary>> = PacketBin}, ?STATE_ACTIVATED, State) ->
#'CastFrame'{body = {command, Command}} = message_pb:decode_msg(PacketBin, 'CastFrame'),
{keep_state, State, [{next_event, info, {server_cast, Command}}]};
-
-handle_event(info, {server_packet, <<8, _/binary>> = PacketBin}, ?STATE_RESTRICTED, State) ->
- #'RequestFrame'{packet_id = PacketId, body = {rpc_request, Request}} =
- message_pb:decode_msg(PacketBin, 'RequestFrame'),
- true = is_integer(PacketId) andalso PacketId > 0,
- {keep_state, State, [{next_event, info, {server_rpc, PacketId, Request}}]};
handle_event(info, {server_packet, <<10, _/binary>> = PacketBin}, ?STATE_RESTRICTED, State) ->
#'CastFrame'{body = {pub, Pub}} = message_pb:decode_msg(PacketBin, 'CastFrame'),
{keep_state, State, [{next_event, info, {server_cast, Pub}}]};
@@ -229,8 +228,8 @@ handle_event(info, {server_packet, <<18, _/binary>> = PacketBin}, ?STATE_RESTRIC
%% 激活消息
%% 微服务部署
-handle_event(info, {server_rpc, PacketId, #'RpcRequest'{method = <<"get_containers">>}}, ?STATE_ACTIVATED, State = #state{transport_pid = TransportPid}) ->
- %% 短暂的等待,efka_inetd收到消息后就立即返回了
+handle_event(info, {container_request, PacketId, #'ContainerRequest'{action = {list, #'ContainerRequest.List'{all = _All}}}},
+ ?STATE_ACTIVATED, State = #state{transport_pid = TransportPid}) ->
case docker_manager:get_containers() of
{ok, Containers} ->
Packet = message_pb:encode_msg(#'ResponseFrame'{
@@ -251,10 +250,11 @@ handle_event(info, {server_rpc, PacketId, #'RpcRequest'{method = <<"get_containe
end,
{keep_state, State};
-%% 微服务部署
-handle_event(info, {server_rpc, PacketId, #'RpcRequest'{method = <<"deploy">>, params = ParamsBin}}, ?STATE_ACTIVATED, State = #state{transport_pid = TransportPid}) ->
- #{<<"task_id">> := TaskId, <<"config">> := Config} = decode_rpc_payload(ParamsBin),
- %% 短暂的等待,efka_inetd收到消息后就立即返回了
+handle_event(info, {container_request, PacketId, #'ContainerRequest'{action = {deploy, #'ContainerRequest.Deploy'{
+ task_id = TaskId,
+ params = Params
+}}}}, ?STATE_ACTIVATED, State = #state{transport_pid = TransportPid}) ->
+ Config = container_deploy_config(Params),
case docker_manager:deploy(TaskId, Config) of
ok ->
Packet = message_pb:encode_msg(#'ResponseFrame'{
@@ -275,11 +275,10 @@ handle_event(info, {server_rpc, PacketId, #'RpcRequest'{method = <<"deploy">>, p
end,
{keep_state, State};
-%% 启动微服务
-handle_event(info, {server_rpc, PacketId, #'RpcRequest'{method = <<"start_container">>, params = ParamsBin}}, ?STATE_ACTIVATED, State = #state{transport_pid = TransportPid}) ->
- #{<<"container_name">> := ContainerName} = decode_rpc_payload(ParamsBin),
- %% 短暂的等待,efka_inetd收到消息后就立即返回了
- case docker_manager:start_container(ContainerName) of
+handle_event(info, {container_request, PacketId, #'ContainerRequest'{action = {start, #'ContainerRequest.Start'{target = Target}}}},
+ ?STATE_ACTIVATED, State = #state{transport_pid = TransportPid}) ->
+ ContainerTarget = container_target(Target),
+ case docker_manager:start_container(ContainerTarget) of
ok ->
Packet = message_pb:encode_msg(#'ResponseFrame'{
packet_id = PacketId,
@@ -299,11 +298,12 @@ handle_event(info, {server_rpc, PacketId, #'RpcRequest'{method = <<"start_contai
end,
{keep_state, State};
-%% 停止微服务
-handle_event(info, {server_rpc, PacketId, #'RpcRequest'{method = <<"stop_container">>, params = ParamsBin}}, ?STATE_ACTIVATED, State = #state{transport_pid = TransportPid}) ->
- #{<<"container_name">> := ContainerName} = decode_rpc_payload(ParamsBin),
- %% 短暂的等待,efka_inetd收到消息后就立即返回了
- case docker_manager:stop_container(ContainerName) of
+handle_event(info, {container_request, PacketId, #'ContainerRequest'{action = {stop, #'ContainerRequest.Stop'{
+ target = Target,
+ timeout_seconds = TimeoutSeconds
+}}}}, ?STATE_ACTIVATED, State = #state{transport_pid = TransportPid}) ->
+ ContainerTarget = container_target(Target),
+ case docker_manager:stop_container(ContainerTarget, TimeoutSeconds) of
ok ->
Packet = message_pb:encode_msg(#'ResponseFrame'{
packet_id = PacketId,
@@ -323,10 +323,12 @@ handle_event(info, {server_rpc, PacketId, #'RpcRequest'{method = <<"stop_contain
end,
{keep_state, State};
-handle_event(info, {server_rpc, PacketId, #'RpcRequest'{method = <<"kill_container">>, params = ParamsBin}}, ?STATE_ACTIVATED, State = #state{transport_pid = TransportPid}) ->
- #{<<"container_name">> := ContainerName} = decode_rpc_payload(ParamsBin),
- %% 短暂的等待,efka_inetd收到消息后就立即返回了
- case docker_manager:kill_container(ContainerName) of
+handle_event(info, {container_request, PacketId, #'ContainerRequest'{action = {kill, #'ContainerRequest.Kill'{
+ target = Target,
+ signal = Signal
+}}}}, ?STATE_ACTIVATED, State = #state{transport_pid = TransportPid}) ->
+ ContainerTarget = container_target(Target),
+ case docker_manager:kill_container(ContainerTarget, to_binary(Signal)) of
ok ->
Packet = message_pb:encode_msg(#'ResponseFrame'{
packet_id = PacketId,
@@ -346,10 +348,13 @@ handle_event(info, {server_rpc, PacketId, #'RpcRequest'{method = <<"kill_contain
end,
{keep_state, State};
-handle_event(info, {server_rpc, PacketId, #'RpcRequest'{method = <<"remove_container">>, params = ParamsBin}}, ?STATE_ACTIVATED, State = #state{transport_pid = TransportPid}) ->
- #{<<"container_name">> := ContainerName} = decode_rpc_payload(ParamsBin),
- %% 短暂的等待,efka_inetd收到消息后就立即返回了
- case docker_manager:remove_container(ContainerName) of
+handle_event(info, {container_request, PacketId, #'ContainerRequest'{action = {remove, #'ContainerRequest.Remove'{
+ target = Target,
+ force = Force,
+ remove_volumes = RemoveVolumes
+}}}}, ?STATE_ACTIVATED, State = #state{transport_pid = TransportPid}) ->
+ ContainerTarget = container_target(Target),
+ case docker_manager:remove_container(ContainerTarget, to_bool(Force), to_bool(RemoveVolumes)) of
ok ->
Packet = message_pb:encode_msg(#'ResponseFrame'{
packet_id = PacketId,
@@ -369,10 +374,12 @@ handle_event(info, {server_rpc, PacketId, #'RpcRequest'{method = <<"remove_conta
end,
{keep_state, State};
-%% config.json配置信息
-handle_event(info, {server_rpc, PacketId, #'RpcRequest'{method = <<"config_container">>, params = ParamsBin}}, ?STATE_ACTIVATED, State = #state{transport_pid = TransportPid}) ->
- #{<<"container_name">> := ContainerName, <<"config">> := Config} = decode_rpc_payload(ParamsBin),
- case docker_manager:config_container(ContainerName, Config) of
+handle_event(info, {container_request, PacketId, #'ContainerRequest'{action = {config, #'ContainerRequest.Config'{
+ target = Target,
+ config = Config
+}}}}, ?STATE_ACTIVATED, State = #state{transport_pid = TransportPid}) ->
+ ContainerTarget = container_target(Target),
+ case docker_manager:config_container(ContainerTarget, iolist_to_binary(Config)) of
ok ->
Packet = message_pb:encode_msg(#'ResponseFrame'{
packet_id = PacketId,
@@ -392,6 +399,46 @@ handle_event(info, {server_rpc, PacketId, #'RpcRequest'{method = <<"config_conta
end,
{keep_state, State};
+handle_event(info, {container_request, PacketId, _Request}, ?STATE_ACTIVATED, State = #state{transport_pid = TransportPid}) ->
+ Packet = message_pb:encode_msg(#'ResponseFrame'{
+ packet_id = PacketId,
+ body = {rpc_reply, #'RpcReply'{
+ reply = {error, #'RpcReply.RpcError'{code = -1, message = <<"unsupported container request">>}}
+ }}
+ }),
+ efka_transport:send(TransportPid, Packet),
+ {keep_state, State};
+
+handle_event(info, {container_request, PacketId, _Request}, ?STATE_RESTRICTED, State = #state{transport_pid = TransportPid}) ->
+ Packet = message_pb:encode_msg(#'ResponseFrame'{
+ packet_id = PacketId,
+ body = {rpc_reply, #'RpcReply'{
+ reply = {error, #'RpcReply.RpcError'{code = -1, message = <<"agent restricted">>}}
+ }}
+ }),
+ efka_transport:send(TransportPid, Packet),
+ {keep_state, State};
+
+handle_event(info, {server_rpc, PacketId, #'RpcRequest'{method = Method}}, ?STATE_ACTIVATED, State = #state{transport_pid = TransportPid}) ->
+ Packet = message_pb:encode_msg(#'ResponseFrame'{
+ packet_id = PacketId,
+ body = {rpc_reply, #'RpcReply'{
+ reply = {error, #'RpcReply.RpcError'{code = -1, message = <<"unsupported rpc request: ", Method/binary>>}}
+ }}
+ }),
+ efka_transport:send(TransportPid, Packet),
+ {keep_state, State};
+
+handle_event(info, {server_rpc, PacketId, _Request}, ?STATE_RESTRICTED, State = #state{transport_pid = TransportPid}) ->
+ Packet = message_pb:encode_msg(#'ResponseFrame'{
+ packet_id = PacketId,
+ body = {rpc_reply, #'RpcReply'{
+ reply = {error, #'RpcReply.RpcError'{code = -1, message = <<"agent restricted">>}}
+ }}
+ }),
+ efka_transport:send(TransportPid, Packet),
+ {keep_state, State};
+
%% 处理task_log
%handle_event(info, {server_async_call, PacketId, <>}, ?STATE_ACTIVATED, State = #state{transport_pid = TransportPid}) ->
% #fetch_task_log{task_id = TaskId} = message_pb:decode_msg(TaskLogBin, fetch_task_log),
@@ -479,12 +526,196 @@ auth_packet() ->
}}
}).
--spec decode_rpc_payload(binary()) -> any().
-decode_rpc_payload(<<>>) ->
- #{};
-decode_rpc_payload(Bin) when is_binary(Bin) ->
- jiffy:decode(Bin, [return_maps]).
-
-spec encode_rpc_payload(any()) -> binary().
encode_rpc_payload(Payload) ->
jiffy:encode(Payload, [force_utf8]).
+
+-spec container_deploy_config(message_pb:'ContainerDeployParams'()) -> map().
+container_deploy_config(#'ContainerDeployParams'{
+ container_name = ContainerName,
+ container_dir = ContainerDir,
+ spec = #'ContainerSpec'{
+ image = Image,
+ command = Command,
+ entrypoint = Entrypoint,
+ env = Env,
+ labels = Labels,
+ volumes = Volumes,
+ user = User,
+ working_dir = WorkingDir,
+ hostname = Hostname,
+ expose = Expose,
+ networks = Networks,
+ network_mode = NetworkMode,
+ healthcheck = Healthcheck,
+ restart = Restart,
+ privileged = Privileged,
+ cap_add = CapAdd,
+ cap_drop = CapDrop,
+ devices = Devices,
+ resources = Resources,
+ ulimits = Ulimits,
+ tmpfs = Tmpfs,
+ sysctls = Sysctls,
+ extra_hosts = ExtraHosts
+ }
+}) ->
+ BaseConfig = #{
+ <<"container_name">> => to_binary(ContainerName),
+ <<"container_dir">> => to_binary(ContainerDir),
+ <<"image">> => to_binary(Image),
+ <<"command">> => [to_binary(Item) || Item <- Command],
+ <<"entrypoint">> => [to_binary(Item) || Item <- Entrypoint],
+ <<"envs">> => [to_binary(Item) || Item <- Env],
+ <<"labels">> => maps:from_list([{to_binary(Key), to_binary(Value)} || {Key, Value} <- Labels]),
+ <<"volumes">> => [volume_bind(Volume) || Volume <- Volumes],
+ <<"user">> => to_binary(User),
+ <<"working_dir">> => to_binary(WorkingDir),
+ <<"hostname">> => to_binary(Hostname),
+ <<"expose">> => [port_expose(Port) || Port <- Expose],
+ <<"networks">> => [to_binary(Network) || Network <- Networks],
+ <<"network_mode">> => to_binary(NetworkMode),
+ <<"restart">> => restart_policy(Restart),
+ <<"privileged">> => to_bool(Privileged),
+ <<"cap_add">> => [to_binary(Item) || Item <- CapAdd],
+ <<"cap_drop">> => [to_binary(Item) || Item <- CapDrop],
+ <<"devices">> => [device_mapping(Device) || Device <- Devices],
+ <<"ulimits">> => maps:from_list([{to_binary(Name), ulimit_spec(Ulimit)} || Ulimit = #'Ulimit'{name = Name} <- Ulimits]),
+ <<"tmpfs">> => [tmpfs_mount(Mount) || Mount <- Tmpfs],
+ <<"sysctls">> => maps:from_list([{to_binary(Key), to_binary(Value)} || {Key, Value} <- Sysctls]),
+ <<"extra_hosts">> => [to_binary(Host) || Host <- ExtraHosts]
+ },
+ HealthcheckConfig = case Healthcheck of
+ undefined ->
+ #{};
+ _ ->
+ #{<<"healthcheck">> => healthcheck_config(Healthcheck)}
+ end,
+ ResourceConfig = case Resources of
+ undefined ->
+ #{};
+ _ ->
+ resource_limits_config(Resources)
+ end,
+ maps:merge(BaseConfig, maps:merge(HealthcheckConfig, ResourceConfig)).
+
+container_target(#'ContainerRef'{name = Name, id = Id}) ->
+ NameBin = to_binary(Name),
+ IdBin = to_binary(Id),
+ case NameBin of
+ <<>> ->
+ true = IdBin =/= <<>>,
+ IdBin;
+ _ ->
+ NameBin
+ end.
+
+healthcheck_config(#'Healthcheck'{test = Test, interval_ns = IntervalNs, timeout_ns = TimeoutNs, retries = Retries}) ->
+ #{
+ <<"test">> => [to_binary(Item) || Item <- Test],
+ <<"interval">> => integer_to_binary(IntervalNs),
+ <<"timeout">> => integer_to_binary(TimeoutNs),
+ <<"retries">> => Retries
+ }.
+
+resource_limits_config(#'ResourceLimits'{
+ memory_bytes = MemoryBytes,
+ memory_reservation_bytes = ReservationBytes,
+ nano_cpus = NanoCpus,
+ cpu_shares = CpuShares
+}) ->
+ Config0 = #{},
+ Config1 = case MemoryBytes of
+ 0 ->
+ Config0;
+ _ ->
+ Config0#{<<"mem_limit">> => integer_to_binary(MemoryBytes)}
+ end,
+ Config2 = case ReservationBytes of
+ 0 ->
+ Config1;
+ _ ->
+ Config1#{<<"mem_reservation">> => integer_to_binary(ReservationBytes)}
+ end,
+ Config3 = case NanoCpus of
+ 0 ->
+ Config2;
+ _ ->
+ Config2#{<<"cpus">> => nano_cpus_to_cpus(NanoCpus)}
+ end,
+ case CpuShares of
+ 0 ->
+ Config3;
+ _ ->
+ Config3#{<<"cpu_shares">> => CpuShares}
+ end.
+
+restart_policy(#'RestartPolicy'{name = Name, maximum_retry_count = 0}) ->
+ to_binary(Name);
+restart_policy(#'RestartPolicy'{name = Name, maximum_retry_count = RetryCount}) ->
+ NameBin = to_binary(Name),
+ RetryCountBin = integer_to_binary(RetryCount),
+ <>.
+
+volume_bind(#'VolumeBind'{host_path = HostPath, container_path = ContainerPath, read_only = true}) ->
+ HostPathBin = to_binary(HostPath),
+ ContainerPathBin = to_binary(ContainerPath),
+ <>;
+volume_bind(#'VolumeBind'{host_path = HostPath, container_path = ContainerPath, read_only = ReadOnly}) ->
+ false = to_bool(ReadOnly),
+ HostPathBin = to_binary(HostPath),
+ ContainerPathBin = to_binary(ContainerPath),
+ <>.
+
+port_expose(#'PortExpose'{container_port = ContainerPort, protocol = Protocol}) ->
+ PortBin = integer_to_binary(ContainerPort),
+ ProtocolBin = to_binary(Protocol),
+ case ProtocolBin of
+ <<>> ->
+ PortBin;
+ <<"tcp">> ->
+ PortBin;
+ _ ->
+ <>
+ end.
+
+device_mapping(#'DeviceMapping'{host_path = HostPath, container_path = ContainerPath, cgroup_permissions = <<"rwm">>}) ->
+ HostPathBin = to_binary(HostPath),
+ ContainerPathBin = to_binary(ContainerPath),
+ <>;
+device_mapping(#'DeviceMapping'{host_path = HostPath, container_path = ContainerPath, cgroup_permissions = Permissions}) ->
+ HostPathBin = to_binary(HostPath),
+ ContainerPathBin = to_binary(ContainerPath),
+ PermissionsBin = to_binary(Permissions),
+ <>.
+
+ulimit_spec(#'Ulimit'{soft = Soft, hard = Hard}) ->
+ SoftBin = integer_to_binary(Soft),
+ HardBin = integer_to_binary(Hard),
+ <>.
+
+tmpfs_mount(#'TmpfsMount'{path = Path, options = <<>>}) ->
+ to_binary(Path);
+tmpfs_mount(#'TmpfsMount'{path = Path, options = Options}) ->
+ PathBin = to_binary(Path),
+ OptionsBin = to_binary(Options),
+ <>.
+
+nano_cpus_to_cpus(NanoCpus) when NanoCpus rem 1000000000 =:= 0 ->
+ NanoCpus div 1000000000;
+nano_cpus_to_cpus(NanoCpus) ->
+ NanoCpus / 1000000000.
+
+to_binary(Value) when is_binary(Value) ->
+ Value;
+to_binary(Value) when is_list(Value) ->
+ unicode:characters_to_binary(Value).
+
+to_bool(true) ->
+ true;
+to_bool(1) ->
+ true;
+to_bool(false) ->
+ false;
+to_bool(0) ->
+ false.