diff --git a/apps/docker/src/docker_deploy_manager.erl b/apps/docker/src/docker_deploy_manager.erl deleted file mode 100644 index 9cecaa1..0000000 --- a/apps/docker/src/docker_deploy_manager.erl +++ /dev/null @@ -1,90 +0,0 @@ -%%%------------------------------------------------------------------- -%%% @author anlicheng -%%% @copyright (C) 2026, -%%% @doc -%%% -%%% @end -%%% Created : 20. 4月 2026 -%%%------------------------------------------------------------------- --module(docker_deploy_manager). --author("anlicheng"). - --behaviour(gen_server). - -%% API --export([start_link/0, deploy/2]). - -%% gen_server callbacks --export([init/1, handle_call/3, handle_cast/2, handle_info/2, terminate/2, code_change/3]). - --define(SERVER, ?MODULE). - --record(state, { - root_dir :: string(), - task_map = #{} -}). - -%%%=================================================================== -%%% API -%%%=================================================================== - --spec start_link() -> {ok, pid()} | ignore | {error, term()}. -start_link() -> - gen_server:start_link({local, ?SERVER}, ?MODULE, [], []). - --spec deploy(integer(), map()) -> ok | {error, binary()}. -deploy(TaskId, Params) when is_integer(TaskId), is_map(Params) -> - gen_server:call(?SERVER, {deploy, TaskId, Params}). - -%%%=================================================================== -%%% gen_server callbacks -%%%=================================================================== - --spec init(list()) -> {ok, #state{}}. -init([]) -> - erlang:process_flag(trap_exit, true), - {ok, RootDir} = docker_helper:root_dir(), - {ok, #state{root_dir = RootDir}}. - --spec handle_call(term(), {pid(), term()}, #state{}) -> {reply, term(), #state{}}. -handle_call({deploy, TaskId, Params}, _From, State = #state{root_dir = RootDir, task_map = TaskMap}) - when is_map(Params) -> - ContainerName = maps:get(<<"container_name">>, Params), - {ok, ContainerDir} = docker_helper:ensure_container_dir(RootDir, ContainerName), - {ok, {TaskPid, _Ref}} = docker_deployer:start_monitor(TaskId, ContainerDir, Params), - logger:debug("[docker_deploy_manager] start deploy task_id: ~p, params: ~p", [TaskId, Params]), - {reply, ok, State#state{task_map = maps:put(TaskPid, TaskId, TaskMap)}}; -handle_call(_Request, _From, State) -> - {reply, ok, State}. - --spec handle_cast(term(), #state{}) -> {noreply, #state{}}. -handle_cast(_Request, State) -> - {noreply, State}. - --spec handle_info(term(), #state{}) -> {noreply, #state{}}. -handle_info({'DOWN', _Ref, process, TaskPid, Reason}, State = #state{task_map = TaskMap}) -> - case maps:take(TaskPid, TaskMap) of - error -> - {noreply, State}; - {TaskId, NTaskMap} -> - case Reason of - normal -> - logger:debug("[docker_deploy_manager] task_id: ~p, exit normal", [TaskId]); - Error0 -> - Error = iolist_to_binary(io_lib:format("~p", [Error0])), - docker_task_reporter:stream(TaskId, <<"error">>, <<"任务失败: "/utf8, Error/binary>>), - docker_task_reporter:close(TaskId, <<"task exited">>), - logger:notice("[docker_deploy_manager] task_id: ~p, exit with error: ~p", [TaskId, Error]) - end, - {noreply, State#state{task_map = NTaskMap}} - end; -handle_info(_Info, State) -> - {noreply, State}. - --spec terminate(term(), #state{}) -> ok. -terminate(_Reason, _State) -> - ok. - --spec code_change(term(), #state{}, term()) -> {ok, #state{}}. -code_change(_OldVsn, State, _Extra) -> - {ok, State}. diff --git a/apps/docker/src/docker_deployer.erl b/apps/docker/src/docker_deployer.erl index ae8fa78..fac97be 100644 --- a/apps/docker/src/docker_deployer.erl +++ b/apps/docker/src/docker_deployer.erl @@ -11,26 +11,17 @@ -dialyzer([{nowarn_function, normalize_image/1}]). %% API --export([start_monitor/3]). --export([deploy/3]). +-export([deploy/4]). -define(TASK_SUCCESS, <<"success">>). -define(TASK_FAIL, <<"fail">>). +-type reporter() :: {stream, pos_integer()}. + %%%=================================================================== %%% API %%%=================================================================== --spec(start_monitor(TaskId :: integer(), ContainerDir :: string(), Params :: map()) -> - {ok, {pid(), reference()}}). -start_monitor(TaskId, ContainerDir, Params) - when is_integer(TaskId), is_list(ContainerDir), is_map(Params) -> - {ok, spawn_monitor(?MODULE, deploy, [TaskId, ContainerDir, Params])}. - -%%%=================================================================== -%%% Internal functions -%%%=================================================================== - %{ % "image": "nginx:latest", % "container_name": "my_nginx", @@ -41,34 +32,35 @@ start_monitor(TaskId, ContainerDir, Params) % "command": ["nginx", "-g", "daemon off;"], % "restart": "always" %} --spec deploy(TaskId :: integer(), ContainerDir :: string(), Params :: map()) -> ok. -deploy(TaskId, ContainerDir, Params) when is_integer(TaskId), is_list(ContainerDir), is_map(Params) -> +-spec deploy(TaskId :: integer(), ContainerDir :: string(), Params :: map(), Reporter :: reporter()) -> ok. +deploy(TaskId, ContainerDir, Params, Reporter) when is_integer(TaskId), is_list(ContainerDir), is_map(Params) -> ContainerName = deploy_container_name(Params), Image0 = deploy_image(Params), - report_task_event(TaskId, <<"info">>, <<"开始部署容器:"/utf8, ContainerName/binary>>), + report_stream_event(Reporter, TaskId, <<"info">>, <<"开始部署容器:"/utf8, ContainerName/binary>>), try - ok = ensure_container_absent(TaskId, ContainerName), - {ok, Image} = ensure_image_ready(TaskId, Image0), - {ok, ContainerId} = create_container_and_config(TaskId, ContainerDir, Params), + ok = ensure_container_absent(Reporter, TaskId, ContainerName), + {ok, Image} = ensure_image_ready(Reporter, TaskId, Image0), + {ok, ContainerId} = create_container_and_config(Reporter, TaskId, ContainerDir, Params), ShortContainerId = short_container_id(ContainerId), - report_task_event(TaskId, <<"info">>, <<"容器创建成功: "/utf8, ShortContainerId/binary>>), - report_task_event(TaskId, <<"info">>, <<"任务完成"/utf8>>), + report_stream_event(Reporter, TaskId, <<"container_id">>, ContainerId), + report_stream_event(Reporter, TaskId, <<"info">>, <<"容器创建成功: "/utf8, ShortContainerId/binary>>), + report_stream_event(Reporter, TaskId, <<"info">>, <<"任务完成"/utf8>>), write_task_summary(TaskId, <<"success">>, ContainerName, Image, ContainerId), - docker_task_reporter:close(TaskId, ?TASK_SUCCESS) + close_task(Reporter, TaskId, ?TASK_SUCCESS) catch throw:{deploy_error, Reason} -> - report_task_event(TaskId, <<"error">>, Reason), - report_task_event(TaskId, <<"error">>, <<"任务失败"/utf8>>), + report_stream_event(Reporter, TaskId, <<"error">>, Reason), + report_stream_event(Reporter, TaskId, <<"error">>, <<"任务失败"/utf8>>), write_task_summary(TaskId, <<"fail">>, ContainerName, Image0, undefined), write_task_failure_reason(TaskId, Reason), - docker_task_reporter:close(TaskId, ?TASK_FAIL); + close_task(Reporter, TaskId, ?TASK_FAIL); Class:Reason:Stacktrace -> Error = iolist_to_binary(io_lib:format("deploy crashed: ~p:~p ~p", [Class, Reason, Stacktrace])), - report_task_event(TaskId, <<"error">>, Error), - report_task_event(TaskId, <<"error">>, <<"任务失败"/utf8>>), + report_stream_event(Reporter, TaskId, <<"error">>, Error), + report_stream_event(Reporter, TaskId, <<"error">>, <<"任务失败"/utf8>>), write_task_summary(TaskId, <<"fail">>, ContainerName, Image0, undefined), write_task_failure_reason(TaskId, Error), - docker_task_reporter:close(TaskId, ?TASK_FAIL) + close_task(Reporter, TaskId, ?TASK_FAIL) end. -spec normalize_image(binary()) -> binary(). @@ -81,9 +73,33 @@ normalize_image(Image) when is_binary(Image) -> end, iolist_to_binary(lists:join(<<"/">>, PrefixParts ++ [NormalizedLast])). --spec report_task_event(TaskId :: integer(), Level :: binary(), Msg :: binary()) -> ok. -report_task_event(TaskId, Level, Msg) when is_integer(TaskId), is_binary(Level), is_binary(Msg) -> - docker_task_reporter:stream(TaskId, Level, Msg). +-spec report_stream_event(reporter(), TaskId :: integer(), Level :: binary(), Msg :: binary()) -> ok. +report_stream_event({stream, StreamId}, _TaskId, Level, Msg) when is_integer(StreamId), is_binary(Level), is_binary(Msg) -> + efka_iot_client:send_stream(StreamId, {data, encode_stream_event(Level, Msg)}). + +-spec close_task(reporter(), integer(), binary()) -> ok. +close_task({stream, StreamId}, _TaskId, Reason) when is_integer(StreamId), is_binary(Reason) -> + ok = efka_iot_client:send_stream(StreamId, {data, encode_close_event(Reason)}), + efka_iot_client:send_stream(StreamId, fin). + +-spec encode_stream_event(binary(), binary()) -> binary(). +encode_stream_event(<<"container_id">>, ContainerId) -> + iolist_to_binary(json:encode(#{ + <<"type">> => <<"container_id">>, + <<"container_id">> => ContainerId + })); +encode_stream_event(Type, Msg) -> + iolist_to_binary(json:encode(#{ + <<"type">> => Type, + <<"message">> => Msg + })). + +-spec encode_close_event(binary()) -> binary(). +encode_close_event(Reason) -> + iolist_to_binary(json:encode(#{ + <<"type">> => <<"close">>, + <<"reason">> => Reason + })). -spec write_task_summary(integer(), binary(), binary(), binary(), undefined | binary()) -> ok. write_task_summary(TaskId, Status, ContainerName, Image, ContainerId) @@ -113,55 +129,55 @@ write_task_failure_reason(TaskId, Reason) when is_integer(TaskId), is_binary(Rea ]), efka_logger:write(Info). --spec ensure_container_absent(TaskId :: integer(), ContainerName :: binary()) -> ok. -ensure_container_absent(TaskId, ContainerName) when is_integer(TaskId), is_binary(ContainerName) -> - report_task_event(TaskId, <<"info">>, <<"开始创建容器: "/utf8, ContainerName/binary>>), +-spec ensure_container_absent(reporter(), TaskId :: integer(), ContainerName :: binary()) -> ok. +ensure_container_absent(Reporter, TaskId, ContainerName) when is_integer(TaskId), is_binary(ContainerName) -> + report_stream_event(Reporter, TaskId, <<"info">>, <<"开始创建容器: "/utf8, ContainerName/binary>>), ok. --spec ensure_image_ready(TaskId :: integer(), Image0 :: binary()) -> {ok, binary()}. -ensure_image_ready(TaskId, Image0) when is_integer(TaskId), is_binary(Image0) -> +-spec ensure_image_ready(reporter(), TaskId :: integer(), Image0 :: binary()) -> {ok, binary()}. +ensure_image_ready(Reporter, TaskId, Image0) when is_integer(TaskId), is_binary(Image0) -> Image = normalize_image(Image0), - report_task_event(TaskId, <<"info">>, <<"使用镜像:"/utf8, Image/binary>>), + report_stream_event(Reporter, TaskId, <<"info">>, <<"使用镜像:"/utf8, Image/binary>>), case docker_commands:check_image_exist(Image) of true -> - report_task_event(TaskId, <<"info">>, <<"本地镜像已存在,跳过拉取:"/utf8, Image/binary>>), + report_stream_event(Reporter, TaskId, <<"info">>, <<"本地镜像已存在,跳过拉取:"/utf8, Image/binary>>), {ok, Image}; false -> - report_task_event(TaskId, <<"info">>, <<"开始拉取镜像:"/utf8, Image/binary>>), + report_stream_event(Reporter, TaskId, <<"info">>, <<"开始拉取镜像:"/utf8, Image/binary>>), case docker_commands:pull_image(Image) of {ok, Ref, Pid, MRef} -> - await_pull_image(TaskId, Ref, Pid, MRef), + await_pull_image(Reporter, TaskId, Ref, Pid, MRef), {ok, Image} end end. --spec await_pull_image(TaskId :: integer(), Ref :: reference(), Pid :: pid(), MRef :: reference()) -> ok. -await_pull_image(TaskId, Ref, Pid, MRef) -> +-spec await_pull_image(reporter(), TaskId :: integer(), Ref :: reference(), Pid :: pid(), MRef :: reference()) -> ok. +await_pull_image(Reporter, TaskId, Ref, Pid, MRef) -> receive {docker_client, Ref, {response, _Status, _Headers}} -> - await_pull_image(TaskId, Ref, Pid, MRef); + await_pull_image(Reporter, TaskId, Ref, Pid, MRef); {docker_client, Ref, {data, Data}} -> - report_task_event(TaskId, <<"info">>, Data), - await_pull_image(TaskId, Ref, Pid, MRef); + report_stream_event(Reporter, TaskId, <<"info">>, Data), + await_pull_image(Reporter, TaskId, Ref, Pid, MRef); {docker_client, Ref, done} -> erlang:demonitor(MRef, [flush]), ok; {docker_client, Ref, {error, Reason}} -> erlang:demonitor(MRef, [flush]), - report_task_event(TaskId, <<"error">>, Reason), + report_stream_event(Reporter, TaskId, <<"error">>, Reason), throw({deploy_error, <<"镜像拉取失败: "/utf8, Reason/binary>>}); {'DOWN', MRef, process, Pid, Reason0} -> Reason = iolist_to_binary(io_lib:format("~p", [Reason0])), throw({deploy_error, <<"镜像拉取失败: "/utf8, Reason/binary>>}) end. --spec create_container_and_config(TaskId :: integer(), ContainerDir :: string(), Params :: map()) -> +-spec create_container_and_config(reporter(), TaskId :: integer(), ContainerDir :: string(), Params :: map()) -> {ok, binary()}. -create_container_and_config(TaskId, ContainerDir, Params) +create_container_and_config(Reporter, TaskId, ContainerDir, Params) when is_integer(TaskId), is_list(ContainerDir), is_map(Params) -> case docker_commands:create_container(ContainerDir, Params) of {ok, ContainerId} -> - ok = create_config_file(TaskId, ContainerDir), + ok = create_config_file(Reporter, TaskId, ContainerDir), {ok, ContainerId}; {error, Reason} when is_binary(Reason) -> throw({deploy_error, format_create_container_error(Reason)}); @@ -170,8 +186,8 @@ create_container_and_config(TaskId, ContainerDir, Params) throw({deploy_error, Error}) end. --spec create_config_file(TaskId :: integer(), ContainerDir :: string()) -> ok. -create_config_file(TaskId, ContainerDir) when is_integer(TaskId), is_list(ContainerDir) -> +-spec create_config_file(reporter(), TaskId :: integer(), ContainerDir :: string()) -> ok. +create_config_file(Reporter, TaskId, ContainerDir) when is_integer(TaskId), is_list(ContainerDir) -> ConfigFile = docker_helper:get_config_file(ContainerDir), case file:open(ConfigFile, [write, exclusive]) of {ok, FD} -> @@ -180,7 +196,7 @@ create_config_file(TaskId, ContainerDir) when is_integer(TaskId), is_list(Contai ok; {error, Reason} -> ReasonBin = list_to_binary(io_lib:format("~p", [Reason])), - report_task_event(TaskId, <<"notice">>, <<"创建配置文件失败: "/utf8, ReasonBin/binary>>), + report_stream_event(Reporter, TaskId, <<"notice">>, <<"创建配置文件失败: "/utf8, ReasonBin/binary>>), ok end. diff --git a/apps/docker/src/docker_sup.erl b/apps/docker/src/docker_sup.erl index 2fd5ee3..8c094b8 100644 --- a/apps/docker/src/docker_sup.erl +++ b/apps/docker/src/docker_sup.erl @@ -19,24 +19,6 @@ start_link() -> -spec init(list()) -> {ok, {supervisor:sup_flags(), [supervisor:child_spec()]}}. init([]) -> SupFlags = #{strategy => one_for_one, intensity => 1000, period => 3600}, - ChildSpecs = [ - #{ - id => docker_task_reporter, - start => {docker_task_reporter, start_link, []}, - restart => permanent, - shutdown => 2000, - type => worker, - modules => [docker_task_reporter] - }, - - #{ - id => docker_deploy_manager, - start => {docker_deploy_manager, start_link, []}, - restart => permanent, - shutdown => 2000, - type => worker, - modules => [docker_deploy_manager] - } - ], + ChildSpecs = [], {ok, {SupFlags, ChildSpecs}}. diff --git a/apps/docker/src/docker_task_reporter.erl b/apps/docker/src/docker_task_reporter.erl deleted file mode 100644 index 897f1dd..0000000 --- a/apps/docker/src/docker_task_reporter.erl +++ /dev/null @@ -1,121 +0,0 @@ -%%%------------------------------------------------------------------- -%%% @author anlicheng -%%% @copyright (C) 2026, -%%% @doc -%%% -%%% @end -%%% Created : 20. 4月 2026 -%%%------------------------------------------------------------------- --module(docker_task_reporter). --author("anlicheng"). - --behaviour(gen_server). - -%% API --export([start_link/0]). --export([stream/3, close/2]). - -%% gen_server callbacks --export([init/1, handle_call/3, handle_cast/2, handle_info/2, terminate/2, code_change/3]). - --define(SERVER, ?MODULE). --define(FLUSH_INTERVAL, 1000). - --record(state, { - pending = queue:new(), - flush_ref = undefined -}). - -%%%=================================================================== -%%% API -%%%=================================================================== - --spec stream(TaskId :: integer(), Type :: binary(), Stream :: binary()) -> ok. -stream(TaskId, Type, Stream) when is_integer(TaskId), is_binary(Type), is_binary(Stream) -> - gen_server:cast(?SERVER, {stream, TaskId, Type, Stream}). - --spec close(TaskId :: integer(), Reason :: binary()) -> ok. -close(TaskId, Reason) when is_integer(TaskId), is_binary(Reason) -> - gen_server:cast(?SERVER, {close, TaskId, Reason}). - --spec start_link() -> {ok, pid()} | ignore | {error, term()}. -start_link() -> - gen_server:start_link({local, ?SERVER}, ?MODULE, [], []). - -%%%=================================================================== -%%% gen_server callbacks -%%%=================================================================== - --spec init(list()) -> {ok, #state{}}. -init([]) -> - {ok, #state{}}. - --spec handle_call(term(), {pid(), term()}, #state{}) -> {reply, ok, #state{}}. -handle_call(_Request, _From, State) -> - {reply, ok, State}. - --spec handle_cast(term(), #state{}) -> {noreply, #state{}}. -handle_cast({stream, TaskId, Type, Stream}, State0 = #state{pending = Pending0}) -> - Pending = queue:in({stream, TaskId, Type, Stream}, Pending0), - {noreply, ensure_flush_timer(State0#state{pending = Pending}, 0)}; -handle_cast({close, TaskId, Reason}, State0 = #state{pending = Pending0}) -> - Pending = queue:in({close, TaskId, Reason}, Pending0), - {noreply, ensure_flush_timer(State0#state{pending = Pending}, 0)}; -handle_cast(_Request, State) -> - {noreply, State}. - --spec handle_info(term(), #state{}) -> {noreply, #state{}}. -handle_info(flush, State0 = #state{}) -> - State1 = State0#state{flush_ref = undefined}, - {noreply, flush_pending(State1)}; -handle_info(_Info, State) -> - {noreply, State}. - --spec terminate(term(), #state{}) -> ok. -terminate(_Reason, _State) -> - ok. - --spec code_change(term(), #state{}, term()) -> {ok, #state{}}. -code_change(_OldVsn, State, _Extra) -> - {ok, State}. - -%%%=================================================================== -%%% Internal functions -%%%=================================================================== - --spec ensure_flush_timer(#state{}, non_neg_integer()) -> #state{}. -ensure_flush_timer(State = #state{flush_ref = undefined}, Delay) when is_integer(Delay), Delay >= 0 -> - Ref = erlang:send_after(Delay, self(), flush), - State#state{flush_ref = Ref}; -ensure_flush_timer(State = #state{}, _Delay) -> - State. - --spec flush_pending(#state{}) -> #state{}. -flush_pending(State = #state{pending = Pending0}) -> - case queue:out(Pending0) of - {empty, _} -> - State; - {{value, Event}, Pending1} -> - case send_event(Event) of - ok -> - flush_pending(State#state{pending = Pending1}); - not_ready -> - ensure_flush_timer(State#state{pending = Pending0}, ?FLUSH_INTERVAL) - end - end. - --spec send_event({stream, integer(), binary(), binary()} | {close, integer(), binary()}) -> ok | not_ready. -send_event({stream, TaskId, Type, Stream}) -> - maybe_send(fun() -> efka_iot_client:task_event_stream(TaskId, Type, Stream) end); -send_event({close, TaskId, Reason}) -> - maybe_send(fun() -> efka_iot_client:close_task_event_stream(TaskId, Reason) end). - --spec maybe_send(fun(() -> any())) -> ok | not_ready. -maybe_send(SendFun) -> - case catch efka_iot_client:is_activated() of - true -> - _ = catch SendFun(), - ok; - _ -> - not_ready - end. diff --git a/apps/efka/include/message.hrl b/apps/efka/include/message.hrl index bf3a5cb..4b3f3b6 100644 --- a/apps/efka/include/message.hrl +++ b/apps/efka/include/message.hrl @@ -17,3 +17,9 @@ -define(CLASS_COMMAND_RESPONSE, 4). -define(CLASS_MESSAGE, 5). -define(CLASS_STREAM, 6). + +%%-------------------------------------------------------------------- +%% Stream targets +%%-------------------------------------------------------------------- +-define(STREAM_TARGET_MANAGER, 1). +-define(STREAM_TARGET_CONTAINER_DEPLOY, 2). diff --git a/apps/efka/include/message_pb.hrl b/apps/efka/include/message_pb.hrl index bd78133..06f1be1 100644 --- a/apps/efka/include/message_pb.hrl +++ b/apps/efka/include/message_pb.hrl @@ -51,7 +51,7 @@ -ifndef('COMMAND.CONTAINER_PB_H'). -define('COMMAND.CONTAINER_PB_H', true). -record('Command.Container', - {action :: {list, message_pb:'Command.Container.ContainerList'()} | {deploy, message_pb:'Command.Container.ContainerDeploy'()} | {start, message_pb:'Command.Container.ContainerStart'()} | {stop, message_pb:'Command.Container.ContainerStop'()} | {kill, message_pb:'Command.Container.ContainerKill'()} | {remove, message_pb:'Command.Container.ContainerRemove'()} | {config, message_pb:'Command.Container.ContainerConfig'()} | undefined % oneof + {action :: {list, message_pb:'Command.Container.ContainerList'()} | {start, message_pb:'Command.Container.ContainerStart'()} | {stop, message_pb:'Command.Container.ContainerStop'()} | {kill, message_pb:'Command.Container.ContainerKill'()} | {remove, message_pb:'Command.Container.ContainerRemove'()} | {config, message_pb:'Command.Container.ContainerConfig'()} | undefined % oneof }). -endif. @@ -95,14 +95,6 @@ }). -endif. --ifndef('COMMAND.CONTAINER.CONTAINERDEPLOY_PB_H'). --define('COMMAND.CONTAINER.CONTAINERDEPLOY_PB_H', true). --record('Command.Container.ContainerDeploy', - {task_id = 0 :: non_neg_integer() | undefined, % = 1, optional, 32 bits - params = <<>> :: iodata() | undefined % = 2, optional - }). --endif. - -ifndef('COMMAND.CONTAINER.CONTAINERLIST_PB_H'). -define('COMMAND.CONTAINER.CONTAINERLIST_PB_H', true). -record('Command.Container.ContainerList', @@ -173,26 +165,17 @@ }). -endif. --ifndef('MESSAGE.TASKEVENT_PB_H'). --define('MESSAGE.TASKEVENT_PB_H', true). --record('Message.TaskEvent', - {task_id = 0 :: non_neg_integer() | undefined, % = 1, optional, 32 bits - type = <<>> :: iodata() | undefined, % = 2, optional - stream = <<>> :: iodata() | undefined % = 3, optional - }). --endif. - -ifndef('MESSAGE_PB_H'). -define('MESSAGE_PB_H', true). -record('Message', - {body :: {ping, message_pb:'Message.Ping'()} | {pong, message_pb:'Message.Pong'()} | {pub, message_pb:'Message.Pub'()} | {metric_data, message_pb:'Message.MetricData'()} | {task_event, message_pb:'Message.TaskEvent'()} | undefined % oneof + {body :: {ping, message_pb:'Message.Ping'()} | {pong, message_pb:'Message.Pong'()} | {pub, message_pb:'Message.Pub'()} | {metric_data, message_pb:'Message.MetricData'()} | undefined % oneof }). -endif. -ifndef('STREAM.OPEN_PB_H'). -define('STREAM.OPEN_PB_H', true). -record('Stream.Open', - { + {target = 0 :: non_neg_integer() | undefined % = 1, optional, 32 bits }). -endif. diff --git a/apps/efka/proto/message.proto b/apps/efka/proto/message.proto index b5603f9..edc5401 100644 --- a/apps/efka/proto/message.proto +++ b/apps/efka/proto/message.proto @@ -60,13 +60,6 @@ message Command { } - message ContainerDeploy { - uint32 task_id = 1; - // The current deploy config is still an application-level structured - // payload. Keep it opaque until the Docker create schema is finalized. - bytes params = 2; - } - message ContainerStart { ContainerTarget target = 1; } @@ -94,7 +87,6 @@ message Command { oneof action { ContainerList list = 1; - ContainerDeploy deploy = 2; ContainerStart start = 3; ContainerStop stop = 4; ContainerKill kill = 5; @@ -143,18 +135,11 @@ message Message { bytes metric = 2; } - message TaskEvent { - uint32 task_id = 1; - bytes type = 2; - bytes stream = 3; - } - oneof body { Ping ping = 10; Pong pong = 11; Pub pub = 12; MetricData metric_data = 13; - TaskEvent task_event = 14; } } @@ -162,6 +147,7 @@ message Stream { uint32 stream_id = 1; message Open { + uint32 target = 1; } message Opened { diff --git a/apps/efka/src/iot/efka_iot_client.erl b/apps/efka/src/iot/efka_iot_client.erl index 0355dd3..006e978 100644 --- a/apps/efka/src/iot/efka_iot_client.erl +++ b/apps/efka/src/iot/efka_iot_client.erl @@ -16,7 +16,7 @@ %% API -export([start_link/0]). --export([metric_data/2, ping/13, task_event_stream/3, close_task_event_stream/2]). +-export([metric_data/2, ping/13]). -export([send_stream/2, stream_done/1]). -export([is_activated/0, dropped_message_count/0]). @@ -63,14 +63,6 @@ metric_data(RouteKey, Metric) when is_binary(RouteKey), is_binary(Metric) -> gen_statem:cast(?SERVER, {metric_data, RouteKey, Metric}). --spec task_event_stream(TaskId :: integer(), Type :: binary(), Stream :: binary()) -> ok. -task_event_stream(TaskId, Type, Stream) when is_integer(TaskId), is_binary(Type), is_binary(Stream) -> - gen_statem:cast(?SERVER, {task_event_stream, TaskId, Type, Stream}). - --spec close_task_event_stream(TaskId :: integer(), Reason :: binary()) -> ok. -close_task_event_stream(TaskId, Reason) when is_integer(TaskId), is_binary(Reason) -> - gen_statem:cast(?SERVER, {close_task_event_stream, TaskId, Reason}). - -spec send_stream(StreamId :: stream_id(), Body :: term()) -> ok. send_stream(StreamId, Body) when is_integer(StreamId), StreamId > 0 -> gen_statem:cast(?SERVER, {send_stream, StreamId, Body}). @@ -147,18 +139,6 @@ handle_event(cast, {metric_data, RouteKey, Metric}, StateName, State = #state{so end end; -%% Task的stream流,只做实时的 -handle_event(cast, {task_event_stream, TaskId, Type, Stream}, ?STATE_ACTIVATED, State = #state{socket = Socket}) -> - logger:debug("[efka_iot_client] event_stream task_id: ~p, stream: ~ts", [TaskId, Stream]), - Packet = encode_message_frame({task_event, TaskId, Type, Stream}), - ok = ssl:send(Socket, Packet), - {keep_state, State}; - -handle_event(cast, {close_task_event_stream, TaskId, Reason}, ?STATE_ACTIVATED, State = #state{socket = Socket}) -> - Packet = encode_message_frame({task_event, TaskId, <<"close">>, Reason}), - ok = ssl:send(Socket, Packet), - {keep_state, State}; - handle_event(cast, {send_stream, StreamId, Body}, ?STATE_ACTIVATED, State = #state{socket = Socket}) when is_integer(StreamId), StreamId > 0 -> Packet = encode_stream_frame(StreamId, Body), @@ -277,14 +257,19 @@ handle_event(info, {'DOWN', MonitorRef, process, WorkerPid, Reason}, _StateName, {keep_state, State}; {StreamId, StreamState, NStreams} -> demonitor_stream(StreamState), - logger:warning("[efka_iot_client] stream worker down, stream_id: ~p, reason: ~p", [StreamId, Reason]), - case State#state.socket of - undefined -> + case Reason of + normal -> {keep_state, State#state{streams = NStreams}}; - Socket -> - Packet = encode_stream_frame(StreamId, {reset, {worker_down, Reason}}), - ok = ssl:send(Socket, Packet), - {keep_state, State#state{streams = NStreams}} + _ -> + logger:warning("[efka_iot_client] stream worker down, stream_id: ~p, reason: ~p", [StreamId, Reason]), + case State#state.socket of + undefined -> + {keep_state, State#state{streams = NStreams}}; + Socket -> + Packet = encode_stream_frame(StreamId, {reset, {worker_down, Reason}}), + ok = ssl:send(Socket, Packet), + {keep_state, State#state{streams = NStreams}} + end end end; @@ -319,8 +304,8 @@ handle_event(internal, #'CommandResponse'{} = Reply, StateName, State) -> %% 透明 TCP stream 多路复用。 handle_event(internal, #'Stream'{stream_id = StreamId, payload = Payload}, ?STATE_ACTIVATED, State) -> Body = case Payload of - {open, #'Stream.Open'{}} -> - open; + {open, #'Stream.Open'{target = Target}} -> + {open, Target}; {opened, #'Stream.Opened'{}} -> opened; {open_error, #'Stream.OpenError'{reason = Reason}} -> @@ -393,10 +378,6 @@ encode_message_frame({metric_data, RouteKey, Metric}) -> encode_transport_frame(?CLASS_MESSAGE, #'Message'{ body = {metric_data, #'Message.MetricData'{route_key = RouteKey, metric = Metric}} }); -encode_message_frame({task_event, TaskId, Type, Stream}) -> - encode_transport_frame(?CLASS_MESSAGE, #'Message'{ - body = {task_event, #'Message.TaskEvent'{task_id = TaskId, type = Type, stream = Stream}} - }); encode_message_frame(Body) -> error({unsupported_message_body, Body}). @@ -404,7 +385,9 @@ encode_message_frame(Body) -> encode_stream_frame(StreamId, Body) -> Payload = case Body of open -> - {open, #'Stream.Open'{}}; + {open, #'Stream.Open'{target = ?STREAM_TARGET_MANAGER}}; + {open, Target} when is_integer(Target), Target > 0 -> + {open, #'Stream.Open'{target = Target}}; opened -> {opened, #'Stream.Opened'{}}; {open_error, Reason} -> @@ -513,10 +496,6 @@ handle_container_command(PktId, #'Command.Container'{action = {list, #'Command.C Reply = docker_commands:get_containers(), send_container_response(Socket, PktId, Reply), ok; -handle_container_command(PktId, #'Command.Container'{action = {deploy, #'Command.Container.ContainerDeploy'{task_id = TaskId, params = Params}}}, Socket) -> - Reply = docker_deploy_manager:deploy(TaskId, decode_json_payload(Params)), - send_container_response(Socket, PktId, Reply), - ok; handle_container_command(PktId, #'Command.Container'{action = {start, #'Command.Container.ContainerStart'{target = Target}}}, Socket) -> Reply = docker_commands:start_container(container_target(Target)), send_container_response(Socket, PktId, Reply), @@ -542,30 +521,29 @@ handle_container_command(PktId, ContainerCommand, Socket) -> send_container_response(Socket, PktId, {error, <<"agent invalid">>}), ok. --spec decode_json_payload(iodata()) -> term(). -decode_json_payload(Payload) -> - Bin = iolist_to_binary(Payload), - try json:decode(Bin) of - Term -> - Term - catch - _:_ -> - Bin - end. - -%% iot 发起 open 时不携带参数;efka 固定按本地 stream_target 建连。 +%% iot 发起 open 时携带 target;efka 根据本地 stream_targets 配置建连。 -spec handle_stream_frame(term(), term(), #state{}) -> gen_statem:event_handler_result(atom(), #state{}). -handle_stream_frame(StreamId, open, State = #state{streams = Streams}) - when is_integer(StreamId), StreamId > 0 -> +handle_stream_frame(StreamId, {open, Target}, State = #state{streams = Streams}) + when is_integer(StreamId), StreamId > 0, is_integer(Target), Target > 0 -> case valid_iot_stream_id(StreamId) andalso not maps:is_key(StreamId, Streams) of true -> - {ok, {WorkerPid, MonitorRef}} = efka_iot_stream:start_stream(StreamId), - StreamState = #stream_state{worker_pid = WorkerPid, monitor_ref = MonitorRef}, - {keep_state, State#state{streams = maps:put(StreamId, StreamState, Streams)}}; + case start_stream_worker(StreamId, Target) of + {ok, {WorkerPid, MonitorRef}} -> + StreamState = #stream_state{worker_pid = WorkerPid, monitor_ref = MonitorRef}, + {keep_state, State#state{streams = maps:put(StreamId, StreamState, Streams)}}; + {error, Reason} -> + send_stream(StreamId, {open_error, Reason}), + {keep_state, State} + end; false -> send_stream(StreamId, {reset, <<"invalid stream open">>}), {keep_state, State} end; +handle_stream_frame(StreamId, {open, Target}, State) + when is_integer(StreamId), StreamId > 0 -> + logger:warning("[efka_iot_client] invalid stream target, stream_id: ~p, target: ~p", [StreamId, Target]), + send_stream(StreamId, {open_error, <<"invalid stream target">>}), + {keep_state, State}; handle_stream_frame(StreamId, Body, State = #state{streams = Streams}) when is_integer(StreamId), StreamId > 0 -> case maps:get(StreamId, Streams, undefined) of @@ -586,6 +564,14 @@ handle_stream_frame(StreamId, Body, State) -> logger:warning("[efka_iot_client] invalid stream frame, stream_id: ~p, body: ~p", [StreamId, Body]), {keep_state, State}. +-spec start_stream_worker(stream_id(), pos_integer()) -> {ok, {pid(), reference()}} | {error, term()}. +start_stream_worker(StreamId, ?STREAM_TARGET_MANAGER) -> + efka_iot_stream:start_stream(StreamId, ?STREAM_TARGET_MANAGER); +start_stream_worker(StreamId, ?STREAM_TARGET_CONTAINER_DEPLOY) -> + efka_iot_deploy_stream:start_stream(StreamId); +start_stream_worker(_StreamId, Target) -> + {error, {unknown_stream_target, Target}}. + -spec maybe_reset_unknown_stream(stream_id(), term()) -> ok. maybe_reset_unknown_stream(_StreamId, {reset, _Reason}) -> ok; diff --git a/apps/efka/src/iot/efka_iot_deploy_stream.erl b/apps/efka/src/iot/efka_iot_deploy_stream.erl new file mode 100644 index 0000000..15e3338 --- /dev/null +++ b/apps/efka/src/iot/efka_iot_deploy_stream.erl @@ -0,0 +1,119 @@ +%%%------------------------------------------------------------------- +%%% @doc One container deploy feedback stream from iot. +%%% @end +%%%------------------------------------------------------------------- +-module(efka_iot_deploy_stream). + +-export([start_stream/1]). +-export([run/1]). + +-define(REQUEST_TIMEOUT, 10000). +-define(MAX_REQUEST_BYTES, 16 * 1024 * 1024). + +-type stream_id() :: pos_integer(). + +-spec start_stream(StreamId :: stream_id()) -> {ok, {pid(), reference()}}. +start_stream(StreamId) when is_integer(StreamId), StreamId > 0 -> + {ok, spawn_monitor(?MODULE, run, [StreamId])}. + +-spec run(StreamId :: stream_id()) -> ok. +run(StreamId) when is_integer(StreamId), StreamId > 0 -> + try run0(StreamId) of + ok -> + ok + catch + Class:Reason:Stack -> + logger:warning("[efka_iot_deploy_stream] stream_id: ~p crashed, class: ~p, reason: ~p, stack: ~p", + [StreamId, Class, Reason, Stack]), + send_error_and_close(StreamId, iolist_to_binary(io_lib:format("~p:~p", [Class, Reason]))), + ok + after + efka_iot_client:stream_done(StreamId) + end. + +-spec run0(stream_id()) -> ok. +run0(StreamId) -> + efka_iot_client:send_stream(StreamId, opened), + case receive_request(StreamId, <<>>) of + {ok, Request} -> + handle_request(StreamId, Request); + {error, reset} -> + ok; + {error, Reason} -> + send_error_and_close(StreamId, reason_to_binary(Reason)) + end. + +-spec receive_request(stream_id(), binary()) -> {ok, binary()} | {error, term()}. +receive_request(StreamId, Acc) -> + receive + {stream, StreamId, {data, Data}} when is_binary(Data) -> + NAcc = <>, + case byte_size(NAcc) =< ?MAX_REQUEST_BYTES of + true -> + receive_request(StreamId, NAcc); + false -> + {error, request_too_large} + end; + {stream, StreamId, fin} -> + {ok, Acc}; + {stream, StreamId, {reset, _Reason}} -> + {error, reset}; + Info -> + logger:debug("[efka_iot_deploy_stream] stream_id: ~p ignore unknown info: ~p", [StreamId, Info]), + receive_request(StreamId, Acc) + after ?REQUEST_TIMEOUT -> + {error, request_timeout} + end. + +-spec handle_request(stream_id(), binary()) -> ok. +handle_request(StreamId, RequestBin) -> + case decode_request(RequestBin) of + {ok, TaskId, Params} -> + deploy(StreamId, TaskId, Params); + {error, Reason} -> + send_error_and_close(StreamId, reason_to_binary(Reason)) + end. + +-spec decode_request(binary()) -> {ok, non_neg_integer(), map()} | {error, term()}. +decode_request(RequestBin) -> + try json:decode(RequestBin) of + #{<<"task_id">> := TaskId, <<"params">> := Params} + when is_integer(TaskId), TaskId >= 0, is_map(Params) -> + {ok, TaskId, Params}; + _ -> + {error, invalid_deploy_request} + catch + Class:Reason -> + {error, {bad_json, Class, Reason}} + end. + +-spec deploy(stream_id(), non_neg_integer(), map()) -> ok. +deploy(StreamId, TaskId, Params) -> + try + ContainerName = maps:get(<<"container_name">>, Params), + {ok, RootDir} = docker_helper:root_dir(), + {ok, ContainerDir} = docker_helper:ensure_container_dir(RootDir, ContainerName), + docker_deployer:deploy(TaskId, ContainerDir, Params, {stream, StreamId}) + catch + Class:Reason -> + Error = iolist_to_binary(io_lib:format("deploy stream failed: ~p:~p", [Class, Reason])), + send_error_and_close(StreamId, Error) + end. + +-spec send_error_and_close(stream_id(), binary()) -> ok. +send_error_and_close(StreamId, Reason) -> + ok = efka_iot_client:send_stream(StreamId, {data, iolist_to_binary(json:encode(#{ + <<"type">> => <<"error">>, + <<"message">> => Reason + }))}), + ok = efka_iot_client:send_stream(StreamId, {data, iolist_to_binary(json:encode(#{ + <<"type">> => <<"close">>, + <<"reason">> => <<"fail">> + }))}), + efka_iot_client:send_stream(StreamId, fin). + +-spec reason_to_binary(term()) -> binary(). +reason_to_binary(Reason) when is_binary(Reason) -> + Reason; +reason_to_binary(Reason) -> + unicode:characters_to_binary(io_lib:format("~p", [Reason])). diff --git a/apps/efka/src/iot/efka_iot_stream.erl b/apps/efka/src/iot/efka_iot_stream.erl index 111563f..7d461d4 100644 --- a/apps/efka/src/iot/efka_iot_stream.erl +++ b/apps/efka/src/iot/efka_iot_stream.erl @@ -5,22 +5,24 @@ %%% @end %%%------------------------------------------------------------------- -module(efka_iot_stream). +-include("message.hrl"). --export([start_stream/1]). --export([run/1]). +-export([start_stream/2]). +-export([run/2]). -define(DEFAULT_CONNECT_TIMEOUT, 3000). -define(DEFAULT_IDLE_TIMEOUT, 120000). -type stream_id() :: pos_integer(). +-type stream_target() :: pos_integer(). --spec start_stream(StreamId :: integer()) -> {ok, {pid(), reference()}}. -start_stream(StreamId) when is_integer(StreamId), StreamId > 0 -> - {ok, spawn_monitor(?MODULE, run, [StreamId])}. +-spec start_stream(StreamId :: integer(), Target :: stream_target()) -> {ok, {pid(), reference()}}. +start_stream(StreamId, Target) when is_integer(StreamId), StreamId > 0, is_integer(Target), Target > 0 -> + {ok, spawn_monitor(?MODULE, run, [StreamId, Target])}. --spec run(StreamId :: stream_id()) -> ok. -run(StreamId) when is_integer(StreamId), StreamId > 0 -> - try run0(StreamId) of +-spec run(StreamId :: stream_id(), Target :: stream_target()) -> ok. +run(StreamId, Target) when is_integer(StreamId), StreamId > 0, is_integer(Target), Target > 0 -> + try run0(StreamId, Target) of ok -> ok catch @@ -33,9 +35,9 @@ run(StreamId) when is_integer(StreamId), StreamId > 0 -> efka_iot_client:stream_done(StreamId) end. --spec run0(stream_id()) -> ok. -run0(StreamId) -> - case open_target_socket() of +-spec run0(stream_id(), stream_target()) -> ok. +run0(StreamId, Target) -> + case open_target_socket(Target) of {ok, Socket, IdleTimeout} -> efka_iot_client:send_stream(StreamId, opened), ok = inet:setopts(Socket, [{active, once}]), @@ -45,9 +47,18 @@ run0(StreamId) -> ok end. --spec open_target_socket() -> {ok, gen_tcp:socket(), timeout()} | {error, term()}. -open_target_socket() -> - {ok, Props} = application:get_env(efka, stream_target), +-spec open_target_socket(stream_target()) -> {ok, gen_tcp:socket(), timeout()} | {error, term()}. +open_target_socket(Target) -> + case stream_target_props(Target) of + {ok, Props} -> + open_target_socket(Target, Props); + {error, Reason} -> + {error, Reason} + end. + +-spec open_target_socket(stream_target(), proplists:proplist()) -> + {ok, gen_tcp:socket(), timeout()} | {error, term()}. +open_target_socket(Target, Props) -> Host = proplists:get_value(host, Props), Port = proplists:get_value(port, Props), ConnectTimeout = proplists:get_value(connect_timeout, Props, ?DEFAULT_CONNECT_TIMEOUT), @@ -63,7 +74,28 @@ open_target_socket() -> {ok, Socket} -> {ok, Socket, IdleTimeout}; {error, Reason} -> - {error, Reason} + {error, {connect_failed, Target, Reason}} + end. + +-spec stream_target_props(stream_target()) -> {ok, proplists:proplist()} | {error, term()}. +stream_target_props(Target) -> + case application:get_env(efka, stream_targets) of + {ok, Targets} -> + case proplists:get_value(Target, Targets) of + undefined -> + {error, {unknown_stream_target, Target}}; + Props -> + {ok, Props} + end; + undefined when Target =:= ?STREAM_TARGET_MANAGER -> + case application:get_env(efka, stream_target) of + {ok, Props} -> + {ok, Props}; + undefined -> + {error, {unknown_stream_target, Target}} + end; + undefined -> + {error, {unknown_stream_target, Target}} end. -spec loop(stream_id(), gen_tcp:socket(), timeout()) -> ok. diff --git a/apps/efka/src/protobuf/message_pb.erl b/apps/efka/src/protobuf/message_pb.erl index 8976653..fc2197a 100644 --- a/apps/efka/src/protobuf/message_pb.erl +++ b/apps/efka/src/protobuf/message_pb.erl @@ -76,8 +76,6 @@ -type 'Command.Container.ContainerStart'() :: #'Command.Container.ContainerStart'{}. --type 'Command.Container.ContainerDeploy'() :: #'Command.Container.ContainerDeploy'{}. - -type 'Command.Container.ContainerList'() :: #'Command.Container.ContainerList'{}. -type 'Command.Container.ContainerTarget'() :: #'Command.Container.ContainerTarget'{}. @@ -96,8 +94,6 @@ -type 'Message.MetricData'() :: #'Message.MetricData'{}. --type 'Message.TaskEvent'() :: #'Message.TaskEvent'{}. - -type 'Message'() :: #'Message'{}. -type 'Stream.Open'() :: #'Stream.Open'{}. @@ -114,9 +110,9 @@ -type 'Stream'() :: #'Stream'{}. --export_type(['Request.AuthRequest'/0, 'Request'/0, 'Response.Error'/0, 'Response.AuthResponse'/0, 'Response'/0, 'Command.Container'/0, 'Command.Container.ContainerConfig'/0, 'Command.Container.ContainerRemove'/0, 'Command.Container.ContainerKill'/0, 'Command.Container.ContainerStop'/0, 'Command.Container.ContainerStart'/0, 'Command.Container.ContainerDeploy'/0, 'Command.Container.ContainerList'/0, 'Command.Container.ContainerTarget'/0, 'Command'/0, 'CommandResponse.Error'/0, 'CommandResponse'/0, 'Message.Ping'/0, 'Message.Pong'/0, 'Message.Pub'/0, 'Message.MetricData'/0, 'Message.TaskEvent'/0, 'Message'/0, 'Stream.Open'/0, 'Stream.Opened'/0, 'Stream.OpenError'/0, 'Stream.Data'/0, 'Stream.Fin'/0, 'Stream.Reset'/0, 'Stream'/0]). --type '$msg_name'() :: 'Request.AuthRequest' | 'Request' | 'Response.Error' | 'Response.AuthResponse' | 'Response' | 'Command.Container' | 'Command.Container.ContainerConfig' | 'Command.Container.ContainerRemove' | 'Command.Container.ContainerKill' | 'Command.Container.ContainerStop' | 'Command.Container.ContainerStart' | 'Command.Container.ContainerDeploy' | 'Command.Container.ContainerList' | 'Command.Container.ContainerTarget' | 'Command' | 'CommandResponse.Error' | 'CommandResponse' | 'Message.Ping' | 'Message.Pong' | 'Message.Pub' | 'Message.MetricData' | 'Message.TaskEvent' | 'Message' | 'Stream.Open' | 'Stream.Opened' | 'Stream.OpenError' | 'Stream.Data' | 'Stream.Fin' | 'Stream.Reset' | 'Stream'. --type '$msg'() :: 'Request.AuthRequest'() | 'Request'() | 'Response.Error'() | 'Response.AuthResponse'() | 'Response'() | 'Command.Container'() | 'Command.Container.ContainerConfig'() | 'Command.Container.ContainerRemove'() | 'Command.Container.ContainerKill'() | 'Command.Container.ContainerStop'() | 'Command.Container.ContainerStart'() | 'Command.Container.ContainerDeploy'() | 'Command.Container.ContainerList'() | 'Command.Container.ContainerTarget'() | 'Command'() | 'CommandResponse.Error'() | 'CommandResponse'() | 'Message.Ping'() | 'Message.Pong'() | 'Message.Pub'() | 'Message.MetricData'() | 'Message.TaskEvent'() | 'Message'() | 'Stream.Open'() | 'Stream.Opened'() | 'Stream.OpenError'() | 'Stream.Data'() | 'Stream.Fin'() | 'Stream.Reset'() | 'Stream'(). +-export_type(['Request.AuthRequest'/0, 'Request'/0, 'Response.Error'/0, 'Response.AuthResponse'/0, 'Response'/0, 'Command.Container'/0, 'Command.Container.ContainerConfig'/0, 'Command.Container.ContainerRemove'/0, 'Command.Container.ContainerKill'/0, 'Command.Container.ContainerStop'/0, 'Command.Container.ContainerStart'/0, 'Command.Container.ContainerList'/0, 'Command.Container.ContainerTarget'/0, 'Command'/0, 'CommandResponse.Error'/0, 'CommandResponse'/0, 'Message.Ping'/0, 'Message.Pong'/0, 'Message.Pub'/0, 'Message.MetricData'/0, 'Message'/0, 'Stream.Open'/0, 'Stream.Opened'/0, 'Stream.OpenError'/0, 'Stream.Data'/0, 'Stream.Fin'/0, 'Stream.Reset'/0, 'Stream'/0]). +-type '$msg_name'() :: 'Request.AuthRequest' | 'Request' | 'Response.Error' | 'Response.AuthResponse' | 'Response' | 'Command.Container' | 'Command.Container.ContainerConfig' | 'Command.Container.ContainerRemove' | 'Command.Container.ContainerKill' | 'Command.Container.ContainerStop' | 'Command.Container.ContainerStart' | 'Command.Container.ContainerList' | 'Command.Container.ContainerTarget' | 'Command' | 'CommandResponse.Error' | 'CommandResponse' | 'Message.Ping' | 'Message.Pong' | 'Message.Pub' | 'Message.MetricData' | 'Message' | 'Stream.Open' | 'Stream.Opened' | 'Stream.OpenError' | 'Stream.Data' | 'Stream.Fin' | 'Stream.Reset' | 'Stream'. +-type '$msg'() :: 'Request.AuthRequest'() | 'Request'() | 'Response.Error'() | 'Response.AuthResponse'() | 'Response'() | 'Command.Container'() | 'Command.Container.ContainerConfig'() | 'Command.Container.ContainerRemove'() | 'Command.Container.ContainerKill'() | 'Command.Container.ContainerStop'() | 'Command.Container.ContainerStart'() | 'Command.Container.ContainerList'() | 'Command.Container.ContainerTarget'() | 'Command'() | 'CommandResponse.Error'() | 'CommandResponse'() | 'Message.Ping'() | 'Message.Pong'() | 'Message.Pub'() | 'Message.MetricData'() | 'Message'() | 'Stream.Open'() | 'Stream.Opened'() | 'Stream.OpenError'() | 'Stream.Data'() | 'Stream.Fin'() | 'Stream.Reset'() | 'Stream'(). -export_type(['$msg_name'/0, '$msg'/0]). -if(?OTP_RELEASE >= 24). @@ -154,7 +150,6 @@ encode_msg(Msg, MsgName, Opts) -> 'Command.Container.ContainerKill' -> 'encode_msg_Command.Container.ContainerKill'(id(Msg, TrUserData), TrUserData); 'Command.Container.ContainerStop' -> 'encode_msg_Command.Container.ContainerStop'(id(Msg, TrUserData), TrUserData); 'Command.Container.ContainerStart' -> 'encode_msg_Command.Container.ContainerStart'(id(Msg, TrUserData), TrUserData); - 'Command.Container.ContainerDeploy' -> 'encode_msg_Command.Container.ContainerDeploy'(id(Msg, TrUserData), TrUserData); 'Command.Container.ContainerList' -> 'encode_msg_Command.Container.ContainerList'(id(Msg, TrUserData), TrUserData); 'Command.Container.ContainerTarget' -> 'encode_msg_Command.Container.ContainerTarget'(id(Msg, TrUserData), TrUserData); 'Command' -> encode_msg_Command(id(Msg, TrUserData), TrUserData); @@ -164,7 +159,6 @@ encode_msg(Msg, MsgName, Opts) -> 'Message.Pong' -> 'encode_msg_Message.Pong'(id(Msg, TrUserData), TrUserData); 'Message.Pub' -> 'encode_msg_Message.Pub'(id(Msg, TrUserData), TrUserData); 'Message.MetricData' -> 'encode_msg_Message.MetricData'(id(Msg, TrUserData), TrUserData); - 'Message.TaskEvent' -> 'encode_msg_Message.TaskEvent'(id(Msg, TrUserData), TrUserData); 'Message' -> encode_msg_Message(id(Msg, TrUserData), TrUserData); 'Stream.Open' -> 'encode_msg_Stream.Open'(id(Msg, TrUserData), TrUserData); 'Stream.Opened' -> 'encode_msg_Stream.Opened'(id(Msg, TrUserData), TrUserData); @@ -282,7 +276,6 @@ encode_msg_Response(#'Response'{packet_id = F1, body = F2}, Bin, TrUserData) -> true -> case id(F1, TrUserData) of {list, TF1} -> begin TrTF1 = id(TF1, TrUserData), 'e_mfield_Command.Container_list'(TrTF1, <>, TrUserData) end; - {deploy, TF1} -> begin TrTF1 = id(TF1, TrUserData), 'e_mfield_Command.Container_deploy'(TrTF1, <>, TrUserData) end; {start, TF1} -> begin TrTF1 = id(TF1, TrUserData), 'e_mfield_Command.Container_start'(TrTF1, <>, TrUserData) end; {stop, TF1} -> begin TrTF1 = id(TF1, TrUserData), 'e_mfield_Command.Container_stop'(TrTF1, <>, TrUserData) end; {kill, TF1} -> begin TrTF1 = id(TF1, TrUserData), 'e_mfield_Command.Container_kill'(TrTF1, <>, TrUserData) end; @@ -408,30 +401,6 @@ encode_msg_Response(#'Response'{packet_id = F1, body = F2}, Bin, TrUserData) -> end end. -'encode_msg_Command.Container.ContainerDeploy'(Msg, TrUserData) -> 'encode_msg_Command.Container.ContainerDeploy'(Msg, <<>>, TrUserData). - - -'encode_msg_Command.Container.ContainerDeploy'(#'Command.Container.ContainerDeploy'{task_id = F1, params = F2}, Bin, TrUserData) -> - B1 = if F1 == undefined -> Bin; - true -> - begin - TrF1 = id(F1, TrUserData), - if TrF1 =:= 0 -> Bin; - true -> e_varint(TrF1, <>, TrUserData) - end - end - end, - if F2 == undefined -> B1; - true -> - begin - TrF2 = id(F2, TrUserData), - case iolist_size(TrF2) of - 0 -> B1; - _ -> e_type_bytes(TrF2, <>, TrUserData) - end - end - end. - 'encode_msg_Command.Container.ContainerList'(_Msg, _TrUserData) -> <<>>. 'encode_msg_Command.Container.ContainerTarget'(Msg, TrUserData) -> 'encode_msg_Command.Container.ContainerTarget'(Msg, <<>>, TrUserData). @@ -584,40 +553,6 @@ encode_msg_CommandResponse(#'CommandResponse'{packet_id = F1, body = F2}, Bin, T end end. -'encode_msg_Message.TaskEvent'(Msg, TrUserData) -> 'encode_msg_Message.TaskEvent'(Msg, <<>>, TrUserData). - - -'encode_msg_Message.TaskEvent'(#'Message.TaskEvent'{task_id = F1, type = F2, stream = F3}, Bin, TrUserData) -> - B1 = if F1 == undefined -> Bin; - true -> - begin - TrF1 = id(F1, TrUserData), - if TrF1 =:= 0 -> Bin; - true -> e_varint(TrF1, <>, TrUserData) - end - end - end, - B2 = if F2 == undefined -> B1; - true -> - begin - TrF2 = id(F2, TrUserData), - case iolist_size(TrF2) of - 0 -> B1; - _ -> e_type_bytes(TrF2, <>, TrUserData) - end - end - end, - if F3 == undefined -> B2; - true -> - begin - TrF3 = id(F3, TrUserData), - case iolist_size(TrF3) of - 0 -> B2; - _ -> e_type_bytes(TrF3, <>, TrUserData) - end - end - end. - encode_msg_Message(Msg, TrUserData) -> encode_msg_Message(Msg, <<>>, TrUserData). @@ -628,12 +563,23 @@ encode_msg_Message(#'Message'{body = F1}, Bin, TrUserData) -> {ping, TF1} -> begin TrTF1 = id(TF1, TrUserData), e_mfield_Message_ping(TrTF1, <>, TrUserData) end; {pong, TF1} -> begin TrTF1 = id(TF1, TrUserData), e_mfield_Message_pong(TrTF1, <>, TrUserData) end; {pub, TF1} -> begin TrTF1 = id(TF1, TrUserData), e_mfield_Message_pub(TrTF1, <>, TrUserData) end; - {metric_data, TF1} -> begin TrTF1 = id(TF1, TrUserData), e_mfield_Message_metric_data(TrTF1, <>, TrUserData) end; - {task_event, TF1} -> begin TrTF1 = id(TF1, TrUserData), e_mfield_Message_task_event(TrTF1, <>, TrUserData) end + {metric_data, TF1} -> begin TrTF1 = id(TF1, TrUserData), e_mfield_Message_metric_data(TrTF1, <>, TrUserData) end end end. -'encode_msg_Stream.Open'(_Msg, _TrUserData) -> <<>>. +'encode_msg_Stream.Open'(Msg, TrUserData) -> 'encode_msg_Stream.Open'(Msg, <<>>, TrUserData). + + +'encode_msg_Stream.Open'(#'Stream.Open'{target = F1}, Bin, TrUserData) -> + if F1 == undefined -> Bin; + true -> + begin + TrF1 = id(F1, TrUserData), + if TrF1 =:= 0 -> Bin; + true -> e_varint(TrF1, <>, TrUserData) + end + end + end. 'encode_msg_Stream.Opened'(_Msg, _TrUserData) -> <<>>. @@ -723,11 +669,6 @@ e_mfield_Response_error(Msg, Bin, TrUserData) -> 'e_mfield_Command.Container_list'(_Msg, Bin, _TrUserData) -> <>. -'e_mfield_Command.Container_deploy'(Msg, Bin, TrUserData) -> - SubBin = 'encode_msg_Command.Container.ContainerDeploy'(Msg, <<>>, TrUserData), - Bin2 = e_varint(byte_size(SubBin), Bin), - <>. - 'e_mfield_Command.Container_start'(Msg, Bin, TrUserData) -> SubBin = 'encode_msg_Command.Container.ContainerStart'(Msg, <<>>, TrUserData), Bin2 = e_varint(byte_size(SubBin), Bin), @@ -802,13 +743,11 @@ e_mfield_Message_metric_data(Msg, Bin, TrUserData) -> Bin2 = e_varint(byte_size(SubBin), Bin), <>. -e_mfield_Message_task_event(Msg, Bin, TrUserData) -> - SubBin = 'encode_msg_Message.TaskEvent'(Msg, <<>>, TrUserData), +e_mfield_Stream_open(Msg, Bin, TrUserData) -> + SubBin = 'encode_msg_Stream.Open'(Msg, <<>>, TrUserData), Bin2 = e_varint(byte_size(SubBin), Bin), <>. -e_mfield_Stream_open(_Msg, Bin, _TrUserData) -> <>. - e_mfield_Stream_opened(_Msg, Bin, _TrUserData) -> <>. e_mfield_Stream_open_error(Msg, Bin, TrUserData) -> @@ -961,7 +900,6 @@ decode_msg_2_doit('Command.Container.ContainerRemove', Bin, TrUserData) -> id('d decode_msg_2_doit('Command.Container.ContainerKill', Bin, TrUserData) -> id('decode_msg_Command.Container.ContainerKill'(Bin, TrUserData), TrUserData); decode_msg_2_doit('Command.Container.ContainerStop', Bin, TrUserData) -> id('decode_msg_Command.Container.ContainerStop'(Bin, TrUserData), TrUserData); decode_msg_2_doit('Command.Container.ContainerStart', Bin, TrUserData) -> id('decode_msg_Command.Container.ContainerStart'(Bin, TrUserData), TrUserData); -decode_msg_2_doit('Command.Container.ContainerDeploy', Bin, TrUserData) -> id('decode_msg_Command.Container.ContainerDeploy'(Bin, TrUserData), TrUserData); decode_msg_2_doit('Command.Container.ContainerList', Bin, TrUserData) -> id('decode_msg_Command.Container.ContainerList'(Bin, TrUserData), TrUserData); decode_msg_2_doit('Command.Container.ContainerTarget', Bin, TrUserData) -> id('decode_msg_Command.Container.ContainerTarget'(Bin, TrUserData), TrUserData); decode_msg_2_doit('Command', Bin, TrUserData) -> id(decode_msg_Command(Bin, TrUserData), TrUserData); @@ -971,7 +909,6 @@ decode_msg_2_doit('Message.Ping', Bin, TrUserData) -> id('decode_msg_Message.Pin decode_msg_2_doit('Message.Pong', Bin, TrUserData) -> id('decode_msg_Message.Pong'(Bin, TrUserData), TrUserData); decode_msg_2_doit('Message.Pub', Bin, TrUserData) -> id('decode_msg_Message.Pub'(Bin, TrUserData), TrUserData); decode_msg_2_doit('Message.MetricData', Bin, TrUserData) -> id('decode_msg_Message.MetricData'(Bin, TrUserData), TrUserData); -decode_msg_2_doit('Message.TaskEvent', Bin, TrUserData) -> id('decode_msg_Message.TaskEvent'(Bin, TrUserData), TrUserData); decode_msg_2_doit('Message', Bin, TrUserData) -> id(decode_msg_Message(Bin, TrUserData), TrUserData); decode_msg_2_doit('Stream.Open', Bin, TrUserData) -> id('decode_msg_Stream.Open'(Bin, TrUserData), TrUserData); decode_msg_2_doit('Stream.Opened', Bin, TrUserData) -> id('decode_msg_Stream.Opened'(Bin, TrUserData), TrUserData); @@ -1268,7 +1205,6 @@ skip_64_Response(<<_:64, Rest/binary>>, Z1, Z2, F, F@_1, F@_2, TrUserData) -> df 'decode_msg_Command.Container'(Bin, TrUserData) -> 'dfp_read_field_def_Command.Container'(Bin, 0, 0, 0, id(undefined, TrUserData), TrUserData). 'dfp_read_field_def_Command.Container'(<<10, Rest/binary>>, Z1, Z2, F, F@_1, TrUserData) -> 'd_field_Command.Container_list'(Rest, Z1, Z2, F, F@_1, TrUserData); -'dfp_read_field_def_Command.Container'(<<18, Rest/binary>>, Z1, Z2, F, F@_1, TrUserData) -> 'd_field_Command.Container_deploy'(Rest, Z1, Z2, F, F@_1, TrUserData); 'dfp_read_field_def_Command.Container'(<<26, Rest/binary>>, Z1, Z2, F, F@_1, TrUserData) -> 'd_field_Command.Container_start'(Rest, Z1, Z2, F, F@_1, TrUserData); 'dfp_read_field_def_Command.Container'(<<34, Rest/binary>>, Z1, Z2, F, F@_1, TrUserData) -> 'd_field_Command.Container_stop'(Rest, Z1, Z2, F, F@_1, TrUserData); 'dfp_read_field_def_Command.Container'(<<42, Rest/binary>>, Z1, Z2, F, F@_1, TrUserData) -> 'd_field_Command.Container_kill'(Rest, Z1, Z2, F, F@_1, TrUserData); @@ -1282,7 +1218,6 @@ skip_64_Response(<<_:64, Rest/binary>>, Z1, Z2, F, F@_1, F@_2, TrUserData) -> df Key = X bsl N + Acc, case Key of 10 -> 'd_field_Command.Container_list'(Rest, 0, 0, 0, F@_1, TrUserData); - 18 -> 'd_field_Command.Container_deploy'(Rest, 0, 0, 0, F@_1, TrUserData); 26 -> 'd_field_Command.Container_start'(Rest, 0, 0, 0, F@_1, TrUserData); 34 -> 'd_field_Command.Container_stop'(Rest, 0, 0, 0, F@_1, TrUserData); 42 -> 'd_field_Command.Container_kill'(Rest, 0, 0, 0, F@_1, TrUserData); @@ -1313,20 +1248,6 @@ skip_64_Response(<<_:64, Rest/binary>>, Z1, Z2, F, F@_1, F@_2, TrUserData) -> df end, TrUserData). -'d_field_Command.Container_deploy'(<<1:1, X:7, Rest/binary>>, N, Acc, F, F@_1, TrUserData) when N < 57 -> 'd_field_Command.Container_deploy'(Rest, N + 7, X bsl N + Acc, F, F@_1, TrUserData); -'d_field_Command.Container_deploy'(<<0:1, X:7, Rest/binary>>, N, Acc, F, Prev, TrUserData) -> - {NewFValue, RestF} = begin Len = X bsl N + Acc, <> = Rest, {id('decode_msg_Command.Container.ContainerDeploy'(Bs, TrUserData), TrUserData), Rest2} end, - 'dfp_read_field_def_Command.Container'(RestF, - 0, - 0, - F, - case Prev of - undefined -> id({deploy, NewFValue}, TrUserData); - {deploy, MVPrev} -> id({deploy, 'merge_msg_Command.Container.ContainerDeploy'(MVPrev, NewFValue, TrUserData)}, TrUserData); - _ -> id({deploy, NewFValue}, TrUserData) - end, - TrUserData). - 'd_field_Command.Container_start'(<<1:1, X:7, Rest/binary>>, N, Acc, F, F@_1, TrUserData) when N < 57 -> 'd_field_Command.Container_start'(Rest, N + 7, X bsl N + Acc, F, F@_1, TrUserData); 'd_field_Command.Container_start'(<<0:1, X:7, Rest/binary>>, N, Acc, F, Prev, TrUserData) -> {NewFValue, RestF} = begin Len = X bsl N + Acc, <> = Rest, {id('decode_msg_Command.Container.ContainerStart'(Bs, TrUserData), TrUserData), Rest2} end, @@ -1712,57 +1633,6 @@ skip_64_Response(<<_:64, Rest/binary>>, Z1, Z2, F, F@_1, F@_2, TrUserData) -> df 'skip_64_Command.Container.ContainerStart'(<<_:64, Rest/binary>>, Z1, Z2, F, F@_1, TrUserData) -> 'dfp_read_field_def_Command.Container.ContainerStart'(Rest, Z1, Z2, F, F@_1, TrUserData). -'decode_msg_Command.Container.ContainerDeploy'(Bin, TrUserData) -> 'dfp_read_field_def_Command.Container.ContainerDeploy'(Bin, 0, 0, 0, id(0, TrUserData), id(<<>>, TrUserData), TrUserData). - -'dfp_read_field_def_Command.Container.ContainerDeploy'(<<8, Rest/binary>>, Z1, Z2, F, F@_1, F@_2, TrUserData) -> 'd_field_Command.Container.ContainerDeploy_task_id'(Rest, Z1, Z2, F, F@_1, F@_2, TrUserData); -'dfp_read_field_def_Command.Container.ContainerDeploy'(<<18, Rest/binary>>, Z1, Z2, F, F@_1, F@_2, TrUserData) -> 'd_field_Command.Container.ContainerDeploy_params'(Rest, Z1, Z2, F, F@_1, F@_2, TrUserData); -'dfp_read_field_def_Command.Container.ContainerDeploy'(<<>>, 0, 0, _, F@_1, F@_2, _) -> #'Command.Container.ContainerDeploy'{task_id = F@_1, params = F@_2}; -'dfp_read_field_def_Command.Container.ContainerDeploy'(Other, Z1, Z2, F, F@_1, F@_2, TrUserData) -> 'dg_read_field_def_Command.Container.ContainerDeploy'(Other, Z1, Z2, F, F@_1, F@_2, TrUserData). - -'dg_read_field_def_Command.Container.ContainerDeploy'(<<1:1, X:7, Rest/binary>>, N, Acc, F, F@_1, F@_2, TrUserData) when N < 32 - 7 -> 'dg_read_field_def_Command.Container.ContainerDeploy'(Rest, N + 7, X bsl N + Acc, F, F@_1, F@_2, TrUserData); -'dg_read_field_def_Command.Container.ContainerDeploy'(<<0:1, X:7, Rest/binary>>, N, Acc, _, F@_1, F@_2, TrUserData) -> - Key = X bsl N + Acc, - case Key of - 8 -> 'd_field_Command.Container.ContainerDeploy_task_id'(Rest, 0, 0, 0, F@_1, F@_2, TrUserData); - 18 -> 'd_field_Command.Container.ContainerDeploy_params'(Rest, 0, 0, 0, F@_1, F@_2, TrUserData); - _ -> - case Key band 7 of - 0 -> 'skip_varint_Command.Container.ContainerDeploy'(Rest, 0, 0, Key bsr 3, F@_1, F@_2, TrUserData); - 1 -> 'skip_64_Command.Container.ContainerDeploy'(Rest, 0, 0, Key bsr 3, F@_1, F@_2, TrUserData); - 2 -> 'skip_length_delimited_Command.Container.ContainerDeploy'(Rest, 0, 0, Key bsr 3, F@_1, F@_2, TrUserData); - 3 -> 'skip_group_Command.Container.ContainerDeploy'(Rest, 0, 0, Key bsr 3, F@_1, F@_2, TrUserData); - 5 -> 'skip_32_Command.Container.ContainerDeploy'(Rest, 0, 0, Key bsr 3, F@_1, F@_2, TrUserData) - end - end; -'dg_read_field_def_Command.Container.ContainerDeploy'(<<>>, 0, 0, _, F@_1, F@_2, _) -> #'Command.Container.ContainerDeploy'{task_id = F@_1, params = F@_2}. - -'d_field_Command.Container.ContainerDeploy_task_id'(<<1:1, X:7, Rest/binary>>, N, Acc, F, F@_1, F@_2, TrUserData) when N < 57 -> 'd_field_Command.Container.ContainerDeploy_task_id'(Rest, N + 7, X bsl N + Acc, F, F@_1, F@_2, TrUserData); -'d_field_Command.Container.ContainerDeploy_task_id'(<<0:1, X:7, Rest/binary>>, N, Acc, F, _, F@_2, TrUserData) -> - {NewFValue, RestF} = {id((X bsl N + Acc) band 4294967295, TrUserData), Rest}, - 'dfp_read_field_def_Command.Container.ContainerDeploy'(RestF, 0, 0, F, NewFValue, F@_2, TrUserData). - -'d_field_Command.Container.ContainerDeploy_params'(<<1:1, X:7, Rest/binary>>, N, Acc, F, F@_1, F@_2, TrUserData) when N < 57 -> 'd_field_Command.Container.ContainerDeploy_params'(Rest, N + 7, X bsl N + Acc, F, F@_1, F@_2, TrUserData); -'d_field_Command.Container.ContainerDeploy_params'(<<0:1, X:7, Rest/binary>>, N, Acc, F, F@_1, _, TrUserData) -> - {NewFValue, RestF} = begin Len = X bsl N + Acc, <> = Rest, Bytes2 = binary:copy(Bytes), {id(Bytes2, TrUserData), Rest2} end, - 'dfp_read_field_def_Command.Container.ContainerDeploy'(RestF, 0, 0, F, F@_1, NewFValue, TrUserData). - -'skip_varint_Command.Container.ContainerDeploy'(<<1:1, _:7, Rest/binary>>, Z1, Z2, F, F@_1, F@_2, TrUserData) -> 'skip_varint_Command.Container.ContainerDeploy'(Rest, Z1, Z2, F, F@_1, F@_2, TrUserData); -'skip_varint_Command.Container.ContainerDeploy'(<<0:1, _:7, Rest/binary>>, Z1, Z2, F, F@_1, F@_2, TrUserData) -> 'dfp_read_field_def_Command.Container.ContainerDeploy'(Rest, Z1, Z2, F, F@_1, F@_2, TrUserData). - -'skip_length_delimited_Command.Container.ContainerDeploy'(<<1:1, X:7, Rest/binary>>, N, Acc, F, F@_1, F@_2, TrUserData) when N < 57 -> 'skip_length_delimited_Command.Container.ContainerDeploy'(Rest, N + 7, X bsl N + Acc, F, F@_1, F@_2, TrUserData); -'skip_length_delimited_Command.Container.ContainerDeploy'(<<0:1, X:7, Rest/binary>>, N, Acc, F, F@_1, F@_2, TrUserData) -> - Length = X bsl N + Acc, - <<_:Length/binary, Rest2/binary>> = Rest, - 'dfp_read_field_def_Command.Container.ContainerDeploy'(Rest2, 0, 0, F, F@_1, F@_2, TrUserData). - -'skip_group_Command.Container.ContainerDeploy'(Bin, _, Z2, FNum, F@_1, F@_2, TrUserData) -> - {_, Rest} = read_group(Bin, FNum), - 'dfp_read_field_def_Command.Container.ContainerDeploy'(Rest, 0, Z2, FNum, F@_1, F@_2, TrUserData). - -'skip_32_Command.Container.ContainerDeploy'(<<_:32, Rest/binary>>, Z1, Z2, F, F@_1, F@_2, TrUserData) -> 'dfp_read_field_def_Command.Container.ContainerDeploy'(Rest, Z1, Z2, F, F@_1, F@_2, TrUserData). - -'skip_64_Command.Container.ContainerDeploy'(<<_:64, Rest/binary>>, Z1, Z2, F, F@_1, F@_2, TrUserData) -> 'dfp_read_field_def_Command.Container.ContainerDeploy'(Rest, Z1, Z2, F, F@_1, F@_2, TrUserData). - 'decode_msg_Command.Container.ContainerList'(Bin, TrUserData) -> 'dfp_read_field_def_Command.Container.ContainerList'(Bin, 0, 0, 0, TrUserData). 'dfp_read_field_def_Command.Container.ContainerList'(<<>>, 0, 0, _, _) -> #'Command.Container.ContainerList'{}; @@ -2205,71 +2075,12 @@ skip_64_CommandResponse(<<_:64, Rest/binary>>, Z1, Z2, F, F@_1, F@_2, TrUserData 'skip_64_Message.MetricData'(<<_:64, Rest/binary>>, Z1, Z2, F, F@_1, F@_2, TrUserData) -> 'dfp_read_field_def_Message.MetricData'(Rest, Z1, Z2, F, F@_1, F@_2, TrUserData). -'decode_msg_Message.TaskEvent'(Bin, TrUserData) -> 'dfp_read_field_def_Message.TaskEvent'(Bin, 0, 0, 0, id(0, TrUserData), id(<<>>, TrUserData), id(<<>>, TrUserData), TrUserData). - -'dfp_read_field_def_Message.TaskEvent'(<<8, Rest/binary>>, Z1, Z2, F, F@_1, F@_2, F@_3, TrUserData) -> 'd_field_Message.TaskEvent_task_id'(Rest, Z1, Z2, F, F@_1, F@_2, F@_3, TrUserData); -'dfp_read_field_def_Message.TaskEvent'(<<18, Rest/binary>>, Z1, Z2, F, F@_1, F@_2, F@_3, TrUserData) -> 'd_field_Message.TaskEvent_type'(Rest, Z1, Z2, F, F@_1, F@_2, F@_3, TrUserData); -'dfp_read_field_def_Message.TaskEvent'(<<26, Rest/binary>>, Z1, Z2, F, F@_1, F@_2, F@_3, TrUserData) -> 'd_field_Message.TaskEvent_stream'(Rest, Z1, Z2, F, F@_1, F@_2, F@_3, TrUserData); -'dfp_read_field_def_Message.TaskEvent'(<<>>, 0, 0, _, F@_1, F@_2, F@_3, _) -> #'Message.TaskEvent'{task_id = F@_1, type = F@_2, stream = F@_3}; -'dfp_read_field_def_Message.TaskEvent'(Other, Z1, Z2, F, F@_1, F@_2, F@_3, TrUserData) -> 'dg_read_field_def_Message.TaskEvent'(Other, Z1, Z2, F, F@_1, F@_2, F@_3, TrUserData). - -'dg_read_field_def_Message.TaskEvent'(<<1:1, X:7, Rest/binary>>, N, Acc, F, F@_1, F@_2, F@_3, TrUserData) when N < 32 - 7 -> 'dg_read_field_def_Message.TaskEvent'(Rest, N + 7, X bsl N + Acc, F, F@_1, F@_2, F@_3, TrUserData); -'dg_read_field_def_Message.TaskEvent'(<<0:1, X:7, Rest/binary>>, N, Acc, _, F@_1, F@_2, F@_3, TrUserData) -> - Key = X bsl N + Acc, - case Key of - 8 -> 'd_field_Message.TaskEvent_task_id'(Rest, 0, 0, 0, F@_1, F@_2, F@_3, TrUserData); - 18 -> 'd_field_Message.TaskEvent_type'(Rest, 0, 0, 0, F@_1, F@_2, F@_3, TrUserData); - 26 -> 'd_field_Message.TaskEvent_stream'(Rest, 0, 0, 0, F@_1, F@_2, F@_3, TrUserData); - _ -> - case Key band 7 of - 0 -> 'skip_varint_Message.TaskEvent'(Rest, 0, 0, Key bsr 3, F@_1, F@_2, F@_3, TrUserData); - 1 -> 'skip_64_Message.TaskEvent'(Rest, 0, 0, Key bsr 3, F@_1, F@_2, F@_3, TrUserData); - 2 -> 'skip_length_delimited_Message.TaskEvent'(Rest, 0, 0, Key bsr 3, F@_1, F@_2, F@_3, TrUserData); - 3 -> 'skip_group_Message.TaskEvent'(Rest, 0, 0, Key bsr 3, F@_1, F@_2, F@_3, TrUserData); - 5 -> 'skip_32_Message.TaskEvent'(Rest, 0, 0, Key bsr 3, F@_1, F@_2, F@_3, TrUserData) - end - end; -'dg_read_field_def_Message.TaskEvent'(<<>>, 0, 0, _, F@_1, F@_2, F@_3, _) -> #'Message.TaskEvent'{task_id = F@_1, type = F@_2, stream = F@_3}. - -'d_field_Message.TaskEvent_task_id'(<<1:1, X:7, Rest/binary>>, N, Acc, F, F@_1, F@_2, F@_3, TrUserData) when N < 57 -> 'd_field_Message.TaskEvent_task_id'(Rest, N + 7, X bsl N + Acc, F, F@_1, F@_2, F@_3, TrUserData); -'d_field_Message.TaskEvent_task_id'(<<0:1, X:7, Rest/binary>>, N, Acc, F, _, F@_2, F@_3, TrUserData) -> - {NewFValue, RestF} = {id((X bsl N + Acc) band 4294967295, TrUserData), Rest}, - 'dfp_read_field_def_Message.TaskEvent'(RestF, 0, 0, F, NewFValue, F@_2, F@_3, TrUserData). - -'d_field_Message.TaskEvent_type'(<<1:1, X:7, Rest/binary>>, N, Acc, F, F@_1, F@_2, F@_3, TrUserData) when N < 57 -> 'd_field_Message.TaskEvent_type'(Rest, N + 7, X bsl N + Acc, F, F@_1, F@_2, F@_3, TrUserData); -'d_field_Message.TaskEvent_type'(<<0:1, X:7, Rest/binary>>, N, Acc, F, F@_1, _, F@_3, TrUserData) -> - {NewFValue, RestF} = begin Len = X bsl N + Acc, <> = Rest, Bytes2 = binary:copy(Bytes), {id(Bytes2, TrUserData), Rest2} end, - 'dfp_read_field_def_Message.TaskEvent'(RestF, 0, 0, F, F@_1, NewFValue, F@_3, TrUserData). - -'d_field_Message.TaskEvent_stream'(<<1:1, X:7, Rest/binary>>, N, Acc, F, F@_1, F@_2, F@_3, TrUserData) when N < 57 -> 'd_field_Message.TaskEvent_stream'(Rest, N + 7, X bsl N + Acc, F, F@_1, F@_2, F@_3, TrUserData); -'d_field_Message.TaskEvent_stream'(<<0:1, X:7, Rest/binary>>, N, Acc, F, F@_1, F@_2, _, TrUserData) -> - {NewFValue, RestF} = begin Len = X bsl N + Acc, <> = Rest, Bytes2 = binary:copy(Bytes), {id(Bytes2, TrUserData), Rest2} end, - 'dfp_read_field_def_Message.TaskEvent'(RestF, 0, 0, F, F@_1, F@_2, NewFValue, TrUserData). - -'skip_varint_Message.TaskEvent'(<<1:1, _:7, Rest/binary>>, Z1, Z2, F, F@_1, F@_2, F@_3, TrUserData) -> 'skip_varint_Message.TaskEvent'(Rest, Z1, Z2, F, F@_1, F@_2, F@_3, TrUserData); -'skip_varint_Message.TaskEvent'(<<0:1, _:7, Rest/binary>>, Z1, Z2, F, F@_1, F@_2, F@_3, TrUserData) -> 'dfp_read_field_def_Message.TaskEvent'(Rest, Z1, Z2, F, F@_1, F@_2, F@_3, TrUserData). - -'skip_length_delimited_Message.TaskEvent'(<<1:1, X:7, Rest/binary>>, N, Acc, F, F@_1, F@_2, F@_3, TrUserData) when N < 57 -> 'skip_length_delimited_Message.TaskEvent'(Rest, N + 7, X bsl N + Acc, F, F@_1, F@_2, F@_3, TrUserData); -'skip_length_delimited_Message.TaskEvent'(<<0:1, X:7, Rest/binary>>, N, Acc, F, F@_1, F@_2, F@_3, TrUserData) -> - Length = X bsl N + Acc, - <<_:Length/binary, Rest2/binary>> = Rest, - 'dfp_read_field_def_Message.TaskEvent'(Rest2, 0, 0, F, F@_1, F@_2, F@_3, TrUserData). - -'skip_group_Message.TaskEvent'(Bin, _, Z2, FNum, F@_1, F@_2, F@_3, TrUserData) -> - {_, Rest} = read_group(Bin, FNum), - 'dfp_read_field_def_Message.TaskEvent'(Rest, 0, Z2, FNum, F@_1, F@_2, F@_3, TrUserData). - -'skip_32_Message.TaskEvent'(<<_:32, Rest/binary>>, Z1, Z2, F, F@_1, F@_2, F@_3, TrUserData) -> 'dfp_read_field_def_Message.TaskEvent'(Rest, Z1, Z2, F, F@_1, F@_2, F@_3, TrUserData). - -'skip_64_Message.TaskEvent'(<<_:64, Rest/binary>>, Z1, Z2, F, F@_1, F@_2, F@_3, TrUserData) -> 'dfp_read_field_def_Message.TaskEvent'(Rest, Z1, Z2, F, F@_1, F@_2, F@_3, TrUserData). - decode_msg_Message(Bin, TrUserData) -> dfp_read_field_def_Message(Bin, 0, 0, 0, id(undefined, TrUserData), TrUserData). dfp_read_field_def_Message(<<82, Rest/binary>>, Z1, Z2, F, F@_1, TrUserData) -> d_field_Message_ping(Rest, Z1, Z2, F, F@_1, TrUserData); dfp_read_field_def_Message(<<90, Rest/binary>>, Z1, Z2, F, F@_1, TrUserData) -> d_field_Message_pong(Rest, Z1, Z2, F, F@_1, TrUserData); dfp_read_field_def_Message(<<98, Rest/binary>>, Z1, Z2, F, F@_1, TrUserData) -> d_field_Message_pub(Rest, Z1, Z2, F, F@_1, TrUserData); dfp_read_field_def_Message(<<106, Rest/binary>>, Z1, Z2, F, F@_1, TrUserData) -> d_field_Message_metric_data(Rest, Z1, Z2, F, F@_1, TrUserData); -dfp_read_field_def_Message(<<114, Rest/binary>>, Z1, Z2, F, F@_1, TrUserData) -> d_field_Message_task_event(Rest, Z1, Z2, F, F@_1, TrUserData); dfp_read_field_def_Message(<<>>, 0, 0, _, F@_1, _) -> #'Message'{body = F@_1}; dfp_read_field_def_Message(Other, Z1, Z2, F, F@_1, TrUserData) -> dg_read_field_def_Message(Other, Z1, Z2, F, F@_1, TrUserData). @@ -2281,7 +2092,6 @@ dg_read_field_def_Message(<<0:1, X:7, Rest/binary>>, N, Acc, _, F@_1, TrUserData 90 -> d_field_Message_pong(Rest, 0, 0, 0, F@_1, TrUserData); 98 -> d_field_Message_pub(Rest, 0, 0, 0, F@_1, TrUserData); 106 -> d_field_Message_metric_data(Rest, 0, 0, 0, F@_1, TrUserData); - 114 -> d_field_Message_task_event(Rest, 0, 0, 0, F@_1, TrUserData); _ -> case Key band 7 of 0 -> skip_varint_Message(Rest, 0, 0, Key bsr 3, F@_1, TrUserData); @@ -2349,20 +2159,6 @@ d_field_Message_metric_data(<<0:1, X:7, Rest/binary>>, N, Acc, F, Prev, TrUserDa end, TrUserData). -d_field_Message_task_event(<<1:1, X:7, Rest/binary>>, N, Acc, F, F@_1, TrUserData) when N < 57 -> d_field_Message_task_event(Rest, N + 7, X bsl N + Acc, F, F@_1, TrUserData); -d_field_Message_task_event(<<0:1, X:7, Rest/binary>>, N, Acc, F, Prev, TrUserData) -> - {NewFValue, RestF} = begin Len = X bsl N + Acc, <> = Rest, {id('decode_msg_Message.TaskEvent'(Bs, TrUserData), TrUserData), Rest2} end, - dfp_read_field_def_Message(RestF, - 0, - 0, - F, - case Prev of - undefined -> id({task_event, NewFValue}, TrUserData); - {task_event, MVPrev} -> id({task_event, 'merge_msg_Message.TaskEvent'(MVPrev, NewFValue, TrUserData)}, TrUserData); - _ -> id({task_event, NewFValue}, TrUserData) - end, - TrUserData). - skip_varint_Message(<<1:1, _:7, Rest/binary>>, Z1, Z2, F, F@_1, TrUserData) -> skip_varint_Message(Rest, Z1, Z2, F, F@_1, TrUserData); skip_varint_Message(<<0:1, _:7, Rest/binary>>, Z1, Z2, F, F@_1, TrUserData) -> dfp_read_field_def_Message(Rest, Z1, Z2, F, F@_1, TrUserData). @@ -2380,39 +2176,49 @@ skip_32_Message(<<_:32, Rest/binary>>, Z1, Z2, F, F@_1, TrUserData) -> dfp_read_ skip_64_Message(<<_:64, Rest/binary>>, Z1, Z2, F, F@_1, TrUserData) -> dfp_read_field_def_Message(Rest, Z1, Z2, F, F@_1, TrUserData). -'decode_msg_Stream.Open'(Bin, TrUserData) -> 'dfp_read_field_def_Stream.Open'(Bin, 0, 0, 0, TrUserData). +'decode_msg_Stream.Open'(Bin, TrUserData) -> 'dfp_read_field_def_Stream.Open'(Bin, 0, 0, 0, id(0, TrUserData), TrUserData). -'dfp_read_field_def_Stream.Open'(<<>>, 0, 0, _, _) -> #'Stream.Open'{}; -'dfp_read_field_def_Stream.Open'(Other, Z1, Z2, F, TrUserData) -> 'dg_read_field_def_Stream.Open'(Other, Z1, Z2, F, TrUserData). +'dfp_read_field_def_Stream.Open'(<<8, Rest/binary>>, Z1, Z2, F, F@_1, TrUserData) -> 'd_field_Stream.Open_target'(Rest, Z1, Z2, F, F@_1, TrUserData); +'dfp_read_field_def_Stream.Open'(<<>>, 0, 0, _, F@_1, _) -> #'Stream.Open'{target = F@_1}; +'dfp_read_field_def_Stream.Open'(Other, Z1, Z2, F, F@_1, TrUserData) -> 'dg_read_field_def_Stream.Open'(Other, Z1, Z2, F, F@_1, TrUserData). -'dg_read_field_def_Stream.Open'(<<1:1, X:7, Rest/binary>>, N, Acc, F, TrUserData) when N < 32 - 7 -> 'dg_read_field_def_Stream.Open'(Rest, N + 7, X bsl N + Acc, F, TrUserData); -'dg_read_field_def_Stream.Open'(<<0:1, X:7, Rest/binary>>, N, Acc, _, TrUserData) -> +'dg_read_field_def_Stream.Open'(<<1:1, X:7, Rest/binary>>, N, Acc, F, F@_1, TrUserData) when N < 32 - 7 -> 'dg_read_field_def_Stream.Open'(Rest, N + 7, X bsl N + Acc, F, F@_1, TrUserData); +'dg_read_field_def_Stream.Open'(<<0:1, X:7, Rest/binary>>, N, Acc, _, F@_1, TrUserData) -> Key = X bsl N + Acc, - case Key band 7 of - 0 -> 'skip_varint_Stream.Open'(Rest, 0, 0, Key bsr 3, TrUserData); - 1 -> 'skip_64_Stream.Open'(Rest, 0, 0, Key bsr 3, TrUserData); - 2 -> 'skip_length_delimited_Stream.Open'(Rest, 0, 0, Key bsr 3, TrUserData); - 3 -> 'skip_group_Stream.Open'(Rest, 0, 0, Key bsr 3, TrUserData); - 5 -> 'skip_32_Stream.Open'(Rest, 0, 0, Key bsr 3, TrUserData) + case Key of + 8 -> 'd_field_Stream.Open_target'(Rest, 0, 0, 0, F@_1, TrUserData); + _ -> + case Key band 7 of + 0 -> 'skip_varint_Stream.Open'(Rest, 0, 0, Key bsr 3, F@_1, TrUserData); + 1 -> 'skip_64_Stream.Open'(Rest, 0, 0, Key bsr 3, F@_1, TrUserData); + 2 -> 'skip_length_delimited_Stream.Open'(Rest, 0, 0, Key bsr 3, F@_1, TrUserData); + 3 -> 'skip_group_Stream.Open'(Rest, 0, 0, Key bsr 3, F@_1, TrUserData); + 5 -> 'skip_32_Stream.Open'(Rest, 0, 0, Key bsr 3, F@_1, TrUserData) + end end; -'dg_read_field_def_Stream.Open'(<<>>, 0, 0, _, _) -> #'Stream.Open'{}. +'dg_read_field_def_Stream.Open'(<<>>, 0, 0, _, F@_1, _) -> #'Stream.Open'{target = F@_1}. -'skip_varint_Stream.Open'(<<1:1, _:7, Rest/binary>>, Z1, Z2, F, TrUserData) -> 'skip_varint_Stream.Open'(Rest, Z1, Z2, F, TrUserData); -'skip_varint_Stream.Open'(<<0:1, _:7, Rest/binary>>, Z1, Z2, F, TrUserData) -> 'dfp_read_field_def_Stream.Open'(Rest, Z1, Z2, F, TrUserData). +'d_field_Stream.Open_target'(<<1:1, X:7, Rest/binary>>, N, Acc, F, F@_1, TrUserData) when N < 57 -> 'd_field_Stream.Open_target'(Rest, N + 7, X bsl N + Acc, F, F@_1, TrUserData); +'d_field_Stream.Open_target'(<<0:1, X:7, Rest/binary>>, N, Acc, F, _, TrUserData) -> + {NewFValue, RestF} = {id((X bsl N + Acc) band 4294967295, TrUserData), Rest}, + 'dfp_read_field_def_Stream.Open'(RestF, 0, 0, F, NewFValue, TrUserData). -'skip_length_delimited_Stream.Open'(<<1:1, X:7, Rest/binary>>, N, Acc, F, TrUserData) when N < 57 -> 'skip_length_delimited_Stream.Open'(Rest, N + 7, X bsl N + Acc, F, TrUserData); -'skip_length_delimited_Stream.Open'(<<0:1, X:7, Rest/binary>>, N, Acc, F, TrUserData) -> +'skip_varint_Stream.Open'(<<1:1, _:7, Rest/binary>>, Z1, Z2, F, F@_1, TrUserData) -> 'skip_varint_Stream.Open'(Rest, Z1, Z2, F, F@_1, TrUserData); +'skip_varint_Stream.Open'(<<0:1, _:7, Rest/binary>>, Z1, Z2, F, F@_1, TrUserData) -> 'dfp_read_field_def_Stream.Open'(Rest, Z1, Z2, F, F@_1, TrUserData). + +'skip_length_delimited_Stream.Open'(<<1:1, X:7, Rest/binary>>, N, Acc, F, F@_1, TrUserData) when N < 57 -> 'skip_length_delimited_Stream.Open'(Rest, N + 7, X bsl N + Acc, F, F@_1, TrUserData); +'skip_length_delimited_Stream.Open'(<<0:1, X:7, Rest/binary>>, N, Acc, F, F@_1, TrUserData) -> Length = X bsl N + Acc, <<_:Length/binary, Rest2/binary>> = Rest, - 'dfp_read_field_def_Stream.Open'(Rest2, 0, 0, F, TrUserData). + 'dfp_read_field_def_Stream.Open'(Rest2, 0, 0, F, F@_1, TrUserData). -'skip_group_Stream.Open'(Bin, _, Z2, FNum, TrUserData) -> +'skip_group_Stream.Open'(Bin, _, Z2, FNum, F@_1, TrUserData) -> {_, Rest} = read_group(Bin, FNum), - 'dfp_read_field_def_Stream.Open'(Rest, 0, Z2, FNum, TrUserData). + 'dfp_read_field_def_Stream.Open'(Rest, 0, Z2, FNum, F@_1, TrUserData). -'skip_32_Stream.Open'(<<_:32, Rest/binary>>, Z1, Z2, F, TrUserData) -> 'dfp_read_field_def_Stream.Open'(Rest, Z1, Z2, F, TrUserData). +'skip_32_Stream.Open'(<<_:32, Rest/binary>>, Z1, Z2, F, F@_1, TrUserData) -> 'dfp_read_field_def_Stream.Open'(Rest, Z1, Z2, F, F@_1, TrUserData). -'skip_64_Stream.Open'(<<_:64, Rest/binary>>, Z1, Z2, F, TrUserData) -> 'dfp_read_field_def_Stream.Open'(Rest, Z1, Z2, F, TrUserData). +'skip_64_Stream.Open'(<<_:64, Rest/binary>>, Z1, Z2, F, F@_1, TrUserData) -> 'dfp_read_field_def_Stream.Open'(Rest, Z1, Z2, F, F@_1, TrUserData). 'decode_msg_Stream.Opened'(Bin, TrUserData) -> 'dfp_read_field_def_Stream.Opened'(Bin, 0, 0, 0, TrUserData). @@ -2837,7 +2643,6 @@ merge_msgs(Prev, New, MsgName, Opts) -> 'Command.Container.ContainerKill' -> 'merge_msg_Command.Container.ContainerKill'(Prev, New, TrUserData); 'Command.Container.ContainerStop' -> 'merge_msg_Command.Container.ContainerStop'(Prev, New, TrUserData); 'Command.Container.ContainerStart' -> 'merge_msg_Command.Container.ContainerStart'(Prev, New, TrUserData); - 'Command.Container.ContainerDeploy' -> 'merge_msg_Command.Container.ContainerDeploy'(Prev, New, TrUserData); 'Command.Container.ContainerList' -> 'merge_msg_Command.Container.ContainerList'(Prev, New, TrUserData); 'Command.Container.ContainerTarget' -> 'merge_msg_Command.Container.ContainerTarget'(Prev, New, TrUserData); 'Command' -> merge_msg_Command(Prev, New, TrUserData); @@ -2847,7 +2652,6 @@ merge_msgs(Prev, New, MsgName, Opts) -> 'Message.Pong' -> 'merge_msg_Message.Pong'(Prev, New, TrUserData); 'Message.Pub' -> 'merge_msg_Message.Pub'(Prev, New, TrUserData); 'Message.MetricData' -> 'merge_msg_Message.MetricData'(Prev, New, TrUserData); - 'Message.TaskEvent' -> 'merge_msg_Message.TaskEvent'(Prev, New, TrUserData); 'Message' -> merge_msg_Message(Prev, New, TrUserData); 'Stream.Open' -> 'merge_msg_Stream.Open'(Prev, New, TrUserData); 'Stream.Opened' -> 'merge_msg_Stream.Opened'(Prev, New, TrUserData); @@ -2919,7 +2723,6 @@ merge_msg_Response(#'Response'{packet_id = PFpacket_id, body = PFbody}, #'Respon #'Command.Container'{action = case {PFaction, NFaction} of {{list, OPFaction}, {list, ONFaction}} -> {list, 'merge_msg_Command.Container.ContainerList'(OPFaction, ONFaction, TrUserData)}; - {{deploy, OPFaction}, {deploy, ONFaction}} -> {deploy, 'merge_msg_Command.Container.ContainerDeploy'(OPFaction, ONFaction, TrUserData)}; {{start, OPFaction}, {start, ONFaction}} -> {start, 'merge_msg_Command.Container.ContainerStart'(OPFaction, ONFaction, TrUserData)}; {{stop, OPFaction}, {stop, ONFaction}} -> {stop, 'merge_msg_Command.Container.ContainerStop'(OPFaction, ONFaction, TrUserData)}; {{kill, OPFaction}, {kill, ONFaction}} -> {kill, 'merge_msg_Command.Container.ContainerKill'(OPFaction, ONFaction, TrUserData)}; @@ -2990,17 +2793,6 @@ merge_msg_Response(#'Response'{packet_id = PFpacket_id, body = PFbody}, #'Respon NFtarget == undefined -> PFtarget end}. --compile({nowarn_unused_function,'merge_msg_Command.Container.ContainerDeploy'/3}). -'merge_msg_Command.Container.ContainerDeploy'(#'Command.Container.ContainerDeploy'{task_id = PFtask_id, params = PFparams}, #'Command.Container.ContainerDeploy'{task_id = NFtask_id, params = NFparams}, _) -> - #'Command.Container.ContainerDeploy'{task_id = - if NFtask_id =:= undefined -> PFtask_id; - true -> NFtask_id - end, - params = - if NFparams =:= undefined -> PFparams; - true -> NFparams - end}. - -compile({nowarn_unused_function,'merge_msg_Command.Container.ContainerList'/3}). 'merge_msg_Command.Container.ContainerList'(_Prev, New, _TrUserData) -> New. @@ -3084,21 +2876,6 @@ merge_msg_CommandResponse(#'CommandResponse'{packet_id = PFpacket_id, body = PFb true -> NFmetric end}. --compile({nowarn_unused_function,'merge_msg_Message.TaskEvent'/3}). -'merge_msg_Message.TaskEvent'(#'Message.TaskEvent'{task_id = PFtask_id, type = PFtype, stream = PFstream}, #'Message.TaskEvent'{task_id = NFtask_id, type = NFtype, stream = NFstream}, _) -> - #'Message.TaskEvent'{task_id = - if NFtask_id =:= undefined -> PFtask_id; - true -> NFtask_id - end, - type = - if NFtype =:= undefined -> PFtype; - true -> NFtype - end, - stream = - if NFstream =:= undefined -> PFstream; - true -> NFstream - end}. - -compile({nowarn_unused_function,merge_msg_Message/3}). merge_msg_Message(#'Message'{body = PFbody}, #'Message'{body = NFbody}, TrUserData) -> #'Message'{body = @@ -3107,13 +2884,16 @@ merge_msg_Message(#'Message'{body = PFbody}, #'Message'{body = NFbody}, TrUserDa {{pong, OPFbody}, {pong, ONFbody}} -> {pong, 'merge_msg_Message.Pong'(OPFbody, ONFbody, TrUserData)}; {{pub, OPFbody}, {pub, ONFbody}} -> {pub, 'merge_msg_Message.Pub'(OPFbody, ONFbody, TrUserData)}; {{metric_data, OPFbody}, {metric_data, ONFbody}} -> {metric_data, 'merge_msg_Message.MetricData'(OPFbody, ONFbody, TrUserData)}; - {{task_event, OPFbody}, {task_event, ONFbody}} -> {task_event, 'merge_msg_Message.TaskEvent'(OPFbody, ONFbody, TrUserData)}; {_, undefined} -> PFbody; _ -> NFbody end}. -compile({nowarn_unused_function,'merge_msg_Stream.Open'/3}). -'merge_msg_Stream.Open'(_Prev, New, _TrUserData) -> New. +'merge_msg_Stream.Open'(#'Stream.Open'{target = PFtarget}, #'Stream.Open'{target = NFtarget}, _) -> + #'Stream.Open'{target = + if NFtarget =:= undefined -> PFtarget; + true -> NFtarget + end}. -compile({nowarn_unused_function,'merge_msg_Stream.Opened'/3}). 'merge_msg_Stream.Opened'(_Prev, New, _TrUserData) -> New. @@ -3182,7 +2962,6 @@ verify_msg(Msg, MsgName, Opts) -> 'Command.Container.ContainerKill' -> 'v_msg_Command.Container.ContainerKill'(Msg, [MsgName], TrUserData); 'Command.Container.ContainerStop' -> 'v_msg_Command.Container.ContainerStop'(Msg, [MsgName], TrUserData); 'Command.Container.ContainerStart' -> 'v_msg_Command.Container.ContainerStart'(Msg, [MsgName], TrUserData); - 'Command.Container.ContainerDeploy' -> 'v_msg_Command.Container.ContainerDeploy'(Msg, [MsgName], TrUserData); 'Command.Container.ContainerList' -> 'v_msg_Command.Container.ContainerList'(Msg, [MsgName], TrUserData); 'Command.Container.ContainerTarget' -> 'v_msg_Command.Container.ContainerTarget'(Msg, [MsgName], TrUserData); 'Command' -> v_msg_Command(Msg, [MsgName], TrUserData); @@ -3192,7 +2971,6 @@ verify_msg(Msg, MsgName, Opts) -> 'Message.Pong' -> 'v_msg_Message.Pong'(Msg, [MsgName], TrUserData); 'Message.Pub' -> 'v_msg_Message.Pub'(Msg, [MsgName], TrUserData); 'Message.MetricData' -> 'v_msg_Message.MetricData'(Msg, [MsgName], TrUserData); - 'Message.TaskEvent' -> 'v_msg_Message.TaskEvent'(Msg, [MsgName], TrUserData); 'Message' -> v_msg_Message(Msg, [MsgName], TrUserData); 'Stream.Open' -> 'v_msg_Stream.Open'(Msg, [MsgName], TrUserData); 'Stream.Opened' -> 'v_msg_Stream.Opened'(Msg, [MsgName], TrUserData); @@ -3288,7 +3066,6 @@ v_msg_Response(X, Path, _TrUserData) -> mk_type_error({expected_msg, 'Response'} case F1 of undefined -> ok; {list, OF1} -> 'v_submsg_Command.Container.ContainerList'(OF1, [list, action | Path], TrUserData); - {deploy, OF1} -> 'v_submsg_Command.Container.ContainerDeploy'(OF1, [deploy, action | Path], TrUserData); {start, OF1} -> 'v_submsg_Command.Container.ContainerStart'(OF1, [start, action | Path], TrUserData); {stop, OF1} -> 'v_submsg_Command.Container.ContainerStop'(OF1, [stop, action | Path], TrUserData); {kill, OF1} -> 'v_submsg_Command.Container.ContainerKill'(OF1, [kill, action | Path], TrUserData); @@ -3379,22 +3156,6 @@ v_msg_Response(X, Path, _TrUserData) -> mk_type_error({expected_msg, 'Response'} ok; 'v_msg_Command.Container.ContainerStart'(X, Path, _TrUserData) -> mk_type_error({expected_msg, 'Command.Container.ContainerStart'}, X, Path). --compile({nowarn_unused_function,'v_submsg_Command.Container.ContainerDeploy'/3}). --dialyzer({nowarn_function,'v_submsg_Command.Container.ContainerDeploy'/3}). -'v_submsg_Command.Container.ContainerDeploy'(Msg, Path, TrUserData) -> 'v_msg_Command.Container.ContainerDeploy'(Msg, Path, TrUserData). - --compile({nowarn_unused_function,'v_msg_Command.Container.ContainerDeploy'/3}). --dialyzer({nowarn_function,'v_msg_Command.Container.ContainerDeploy'/3}). -'v_msg_Command.Container.ContainerDeploy'(#'Command.Container.ContainerDeploy'{task_id = F1, params = F2}, Path, TrUserData) -> - if F1 == undefined -> ok; - true -> v_type_uint32(F1, [task_id | Path], TrUserData) - end, - if F2 == undefined -> ok; - true -> v_type_bytes(F2, [params | Path], TrUserData) - end, - ok; -'v_msg_Command.Container.ContainerDeploy'(X, Path, _TrUserData) -> mk_type_error({expected_msg, 'Command.Container.ContainerDeploy'}, X, Path). - -compile({nowarn_unused_function,'v_submsg_Command.Container.ContainerList'/3}). -dialyzer({nowarn_function,'v_submsg_Command.Container.ContainerList'/3}). 'v_submsg_Command.Container.ContainerList'(Msg, Path, TrUserData) -> 'v_msg_Command.Container.ContainerList'(Msg, Path, TrUserData). @@ -3518,25 +3279,6 @@ v_msg_CommandResponse(X, Path, _TrUserData) -> mk_type_error({expected_msg, 'Com ok; 'v_msg_Message.MetricData'(X, Path, _TrUserData) -> mk_type_error({expected_msg, 'Message.MetricData'}, X, Path). --compile({nowarn_unused_function,'v_submsg_Message.TaskEvent'/3}). --dialyzer({nowarn_function,'v_submsg_Message.TaskEvent'/3}). -'v_submsg_Message.TaskEvent'(Msg, Path, TrUserData) -> 'v_msg_Message.TaskEvent'(Msg, Path, TrUserData). - --compile({nowarn_unused_function,'v_msg_Message.TaskEvent'/3}). --dialyzer({nowarn_function,'v_msg_Message.TaskEvent'/3}). -'v_msg_Message.TaskEvent'(#'Message.TaskEvent'{task_id = F1, type = F2, stream = F3}, Path, TrUserData) -> - if F1 == undefined -> ok; - true -> v_type_uint32(F1, [task_id | Path], TrUserData) - end, - if F2 == undefined -> ok; - true -> v_type_bytes(F2, [type | Path], TrUserData) - end, - if F3 == undefined -> ok; - true -> v_type_bytes(F3, [stream | Path], TrUserData) - end, - ok; -'v_msg_Message.TaskEvent'(X, Path, _TrUserData) -> mk_type_error({expected_msg, 'Message.TaskEvent'}, X, Path). - -compile({nowarn_unused_function,v_msg_Message/3}). -dialyzer({nowarn_function,v_msg_Message/3}). v_msg_Message(#'Message'{body = F1}, Path, TrUserData) -> @@ -3546,7 +3288,6 @@ v_msg_Message(#'Message'{body = F1}, Path, TrUserData) -> {pong, OF1} -> 'v_submsg_Message.Pong'(OF1, [pong, body | Path], TrUserData); {pub, OF1} -> 'v_submsg_Message.Pub'(OF1, [pub, body | Path], TrUserData); {metric_data, OF1} -> 'v_submsg_Message.MetricData'(OF1, [metric_data, body | Path], TrUserData); - {task_event, OF1} -> 'v_submsg_Message.TaskEvent'(OF1, [task_event, body | Path], TrUserData); _ -> mk_type_error(invalid_oneof, F1, [body | Path]) end, ok; @@ -3558,7 +3299,11 @@ v_msg_Message(X, Path, _TrUserData) -> mk_type_error({expected_msg, 'Message'}, -compile({nowarn_unused_function,'v_msg_Stream.Open'/3}). -dialyzer({nowarn_function,'v_msg_Stream.Open'/3}). -'v_msg_Stream.Open'(#'Stream.Open'{}, _Path, _) -> ok; +'v_msg_Stream.Open'(#'Stream.Open'{target = F1}, Path, TrUserData) -> + if F1 == undefined -> ok; + true -> v_type_uint32(F1, [target | Path], TrUserData) + end, + ok; 'v_msg_Stream.Open'(X, Path, _TrUserData) -> mk_type_error({expected_msg, 'Stream.Open'}, X, Path). -compile({nowarn_unused_function,'v_submsg_Stream.Opened'/3}). @@ -3719,7 +3464,6 @@ get_msg_defs() -> [#gpb_oneof{name = action, rnum = 2, fields = [#field{name = list, fnum = 1, rnum = 2, type = {msg, 'Command.Container.ContainerList'}, occurrence = optional, opts = []}, - #field{name = deploy, fnum = 2, rnum = 2, type = {msg, 'Command.Container.ContainerDeploy'}, occurrence = optional, opts = []}, #field{name = start, fnum = 3, rnum = 2, type = {msg, 'Command.Container.ContainerStart'}, occurrence = optional, opts = []}, #field{name = stop, fnum = 4, rnum = 2, type = {msg, 'Command.Container.ContainerStop'}, occurrence = optional, opts = []}, #field{name = kill, fnum = 5, rnum = 2, type = {msg, 'Command.Container.ContainerKill'}, occurrence = optional, opts = []}, @@ -3737,7 +3481,6 @@ get_msg_defs() -> {{msg, 'Command.Container.ContainerStop'}, [#field{name = target, fnum = 1, rnum = 2, type = {msg, 'Command.Container.ContainerTarget'}, occurrence = optional, opts = []}, #field{name = timeout_seconds, fnum = 2, rnum = 3, type = uint32, occurrence = optional, opts = []}]}, {{msg, 'Command.Container.ContainerStart'}, [#field{name = target, fnum = 1, rnum = 2, type = {msg, 'Command.Container.ContainerTarget'}, occurrence = optional, opts = []}]}, - {{msg, 'Command.Container.ContainerDeploy'}, [#field{name = task_id, fnum = 1, rnum = 2, type = uint32, occurrence = optional, opts = []}, #field{name = params, fnum = 2, rnum = 3, type = bytes, occurrence = optional, opts = []}]}, {{msg, 'Command.Container.ContainerList'}, []}, {{msg, 'Command.Container.ContainerTarget'}, [#field{name = name, fnum = 1, rnum = 2, type = bytes, occurrence = optional, opts = []}, #field{name = id, fnum = 2, rnum = 3, type = bytes, occurrence = optional, opts = []}]}, {{msg, 'Command'}, @@ -3755,20 +3498,15 @@ get_msg_defs() -> #field{name = qos, fnum = 2, rnum = 3, type = uint32, occurrence = optional, opts = []}, #field{name = content, fnum = 3, rnum = 4, type = bytes, occurrence = optional, opts = []}]}, {{msg, 'Message.MetricData'}, [#field{name = route_key, fnum = 1, rnum = 2, type = bytes, occurrence = optional, opts = []}, #field{name = metric, fnum = 2, rnum = 3, type = bytes, occurrence = optional, opts = []}]}, - {{msg, 'Message.TaskEvent'}, - [#field{name = task_id, fnum = 1, rnum = 2, type = uint32, occurrence = optional, opts = []}, - #field{name = type, fnum = 2, rnum = 3, type = bytes, occurrence = optional, opts = []}, - #field{name = stream, fnum = 3, rnum = 4, type = bytes, occurrence = optional, opts = []}]}, {{msg, 'Message'}, [#gpb_oneof{name = body, rnum = 2, fields = [#field{name = ping, fnum = 10, rnum = 2, type = {msg, 'Message.Ping'}, occurrence = optional, opts = []}, #field{name = pong, fnum = 11, rnum = 2, type = {msg, 'Message.Pong'}, occurrence = optional, opts = []}, #field{name = pub, fnum = 12, rnum = 2, type = {msg, 'Message.Pub'}, occurrence = optional, opts = []}, - #field{name = metric_data, fnum = 13, rnum = 2, type = {msg, 'Message.MetricData'}, occurrence = optional, opts = []}, - #field{name = task_event, fnum = 14, rnum = 2, type = {msg, 'Message.TaskEvent'}, occurrence = optional, opts = []}], + #field{name = metric_data, fnum = 13, rnum = 2, type = {msg, 'Message.MetricData'}, occurrence = optional, opts = []}], opts = []}]}, - {{msg, 'Stream.Open'}, []}, + {{msg, 'Stream.Open'}, [#field{name = target, fnum = 1, rnum = 2, type = uint32, occurrence = optional, opts = []}]}, {{msg, 'Stream.Opened'}, []}, {{msg, 'Stream.OpenError'}, [#field{name = reason, fnum = 1, rnum = 2, type = bytes, occurrence = optional, opts = []}]}, {{msg, 'Stream.Data'}, [#field{name = bytes, fnum = 1, rnum = 2, type = bytes, occurrence = optional, opts = []}]}, @@ -3799,7 +3537,6 @@ get_msg_names() -> 'Command.Container.ContainerKill', 'Command.Container.ContainerStop', 'Command.Container.ContainerStart', - 'Command.Container.ContainerDeploy', 'Command.Container.ContainerList', 'Command.Container.ContainerTarget', 'Command', @@ -3809,7 +3546,6 @@ get_msg_names() -> 'Message.Pong', 'Message.Pub', 'Message.MetricData', - 'Message.TaskEvent', 'Message', 'Stream.Open', 'Stream.Opened', @@ -3835,7 +3571,6 @@ get_msg_or_group_names() -> 'Command.Container.ContainerKill', 'Command.Container.ContainerStop', 'Command.Container.ContainerStart', - 'Command.Container.ContainerDeploy', 'Command.Container.ContainerList', 'Command.Container.ContainerTarget', 'Command', @@ -3845,7 +3580,6 @@ get_msg_or_group_names() -> 'Message.Pong', 'Message.Pub', 'Message.MetricData', - 'Message.TaskEvent', 'Message', 'Stream.Open', 'Stream.Opened', @@ -3888,7 +3622,6 @@ find_msg_def('Command.Container') -> [#gpb_oneof{name = action, rnum = 2, fields = [#field{name = list, fnum = 1, rnum = 2, type = {msg, 'Command.Container.ContainerList'}, occurrence = optional, opts = []}, - #field{name = deploy, fnum = 2, rnum = 2, type = {msg, 'Command.Container.ContainerDeploy'}, occurrence = optional, opts = []}, #field{name = start, fnum = 3, rnum = 2, type = {msg, 'Command.Container.ContainerStart'}, occurrence = optional, opts = []}, #field{name = stop, fnum = 4, rnum = 2, type = {msg, 'Command.Container.ContainerStop'}, occurrence = optional, opts = []}, #field{name = kill, fnum = 5, rnum = 2, type = {msg, 'Command.Container.ContainerKill'}, occurrence = optional, opts = []}, @@ -3906,7 +3639,6 @@ find_msg_def('Command.Container.ContainerKill') -> find_msg_def('Command.Container.ContainerStop') -> [#field{name = target, fnum = 1, rnum = 2, type = {msg, 'Command.Container.ContainerTarget'}, occurrence = optional, opts = []}, #field{name = timeout_seconds, fnum = 2, rnum = 3, type = uint32, occurrence = optional, opts = []}]; find_msg_def('Command.Container.ContainerStart') -> [#field{name = target, fnum = 1, rnum = 2, type = {msg, 'Command.Container.ContainerTarget'}, occurrence = optional, opts = []}]; -find_msg_def('Command.Container.ContainerDeploy') -> [#field{name = task_id, fnum = 1, rnum = 2, type = uint32, occurrence = optional, opts = []}, #field{name = params, fnum = 2, rnum = 3, type = bytes, occurrence = optional, opts = []}]; find_msg_def('Command.Container.ContainerList') -> []; find_msg_def('Command.Container.ContainerTarget') -> [#field{name = name, fnum = 1, rnum = 2, type = bytes, occurrence = optional, opts = []}, #field{name = id, fnum = 2, rnum = 3, type = bytes, occurrence = optional, opts = []}]; find_msg_def('Command') -> @@ -3924,20 +3656,15 @@ find_msg_def('Message.Pub') -> #field{name = qos, fnum = 2, rnum = 3, type = uint32, occurrence = optional, opts = []}, #field{name = content, fnum = 3, rnum = 4, type = bytes, occurrence = optional, opts = []}]; find_msg_def('Message.MetricData') -> [#field{name = route_key, fnum = 1, rnum = 2, type = bytes, occurrence = optional, opts = []}, #field{name = metric, fnum = 2, rnum = 3, type = bytes, occurrence = optional, opts = []}]; -find_msg_def('Message.TaskEvent') -> - [#field{name = task_id, fnum = 1, rnum = 2, type = uint32, occurrence = optional, opts = []}, - #field{name = type, fnum = 2, rnum = 3, type = bytes, occurrence = optional, opts = []}, - #field{name = stream, fnum = 3, rnum = 4, type = bytes, occurrence = optional, opts = []}]; find_msg_def('Message') -> [#gpb_oneof{name = body, rnum = 2, fields = [#field{name = ping, fnum = 10, rnum = 2, type = {msg, 'Message.Ping'}, occurrence = optional, opts = []}, #field{name = pong, fnum = 11, rnum = 2, type = {msg, 'Message.Pong'}, occurrence = optional, opts = []}, #field{name = pub, fnum = 12, rnum = 2, type = {msg, 'Message.Pub'}, occurrence = optional, opts = []}, - #field{name = metric_data, fnum = 13, rnum = 2, type = {msg, 'Message.MetricData'}, occurrence = optional, opts = []}, - #field{name = task_event, fnum = 14, rnum = 2, type = {msg, 'Message.TaskEvent'}, occurrence = optional, opts = []}], + #field{name = metric_data, fnum = 13, rnum = 2, type = {msg, 'Message.MetricData'}, occurrence = optional, opts = []}], opts = []}]; -find_msg_def('Stream.Open') -> []; +find_msg_def('Stream.Open') -> [#field{name = target, fnum = 1, rnum = 2, type = uint32, occurrence = optional, opts = []}]; find_msg_def('Stream.Opened') -> []; find_msg_def('Stream.OpenError') -> [#field{name = reason, fnum = 1, rnum = 2, type = bytes, occurrence = optional, opts = []}]; find_msg_def('Stream.Data') -> [#field{name = bytes, fnum = 1, rnum = 2, type = bytes, occurrence = optional, opts = []}]; @@ -4023,7 +3750,6 @@ fqbin_to_msg_name(<<"Command.Container.ContainerRemove">>) -> 'Command.Container fqbin_to_msg_name(<<"Command.Container.ContainerKill">>) -> 'Command.Container.ContainerKill'; fqbin_to_msg_name(<<"Command.Container.ContainerStop">>) -> 'Command.Container.ContainerStop'; fqbin_to_msg_name(<<"Command.Container.ContainerStart">>) -> 'Command.Container.ContainerStart'; -fqbin_to_msg_name(<<"Command.Container.ContainerDeploy">>) -> 'Command.Container.ContainerDeploy'; fqbin_to_msg_name(<<"Command.Container.ContainerList">>) -> 'Command.Container.ContainerList'; fqbin_to_msg_name(<<"Command.Container.ContainerTarget">>) -> 'Command.Container.ContainerTarget'; fqbin_to_msg_name(<<"Command">>) -> 'Command'; @@ -4033,7 +3759,6 @@ fqbin_to_msg_name(<<"Message.Ping">>) -> 'Message.Ping'; fqbin_to_msg_name(<<"Message.Pong">>) -> 'Message.Pong'; fqbin_to_msg_name(<<"Message.Pub">>) -> 'Message.Pub'; fqbin_to_msg_name(<<"Message.MetricData">>) -> 'Message.MetricData'; -fqbin_to_msg_name(<<"Message.TaskEvent">>) -> 'Message.TaskEvent'; fqbin_to_msg_name(<<"Message">>) -> 'Message'; fqbin_to_msg_name(<<"Stream.Open">>) -> 'Stream.Open'; fqbin_to_msg_name(<<"Stream.Opened">>) -> 'Stream.Opened'; @@ -4056,7 +3781,6 @@ msg_name_to_fqbin('Command.Container.ContainerRemove') -> <<"Command.Container.C msg_name_to_fqbin('Command.Container.ContainerKill') -> <<"Command.Container.ContainerKill">>; msg_name_to_fqbin('Command.Container.ContainerStop') -> <<"Command.Container.ContainerStop">>; msg_name_to_fqbin('Command.Container.ContainerStart') -> <<"Command.Container.ContainerStart">>; -msg_name_to_fqbin('Command.Container.ContainerDeploy') -> <<"Command.Container.ContainerDeploy">>; msg_name_to_fqbin('Command.Container.ContainerList') -> <<"Command.Container.ContainerList">>; msg_name_to_fqbin('Command.Container.ContainerTarget') -> <<"Command.Container.ContainerTarget">>; msg_name_to_fqbin('Command') -> <<"Command">>; @@ -4066,7 +3790,6 @@ msg_name_to_fqbin('Message.Ping') -> <<"Message.Ping">>; msg_name_to_fqbin('Message.Pong') -> <<"Message.Pong">>; msg_name_to_fqbin('Message.Pub') -> <<"Message.Pub">>; msg_name_to_fqbin('Message.MetricData') -> <<"Message.MetricData">>; -msg_name_to_fqbin('Message.TaskEvent') -> <<"Message.TaskEvent">>; msg_name_to_fqbin('Message') -> <<"Message">>; msg_name_to_fqbin('Stream.Open') -> <<"Stream.Open">>; msg_name_to_fqbin('Stream.Opened') -> <<"Stream.Opened">>; @@ -4117,7 +3840,6 @@ get_msg_containment("message") -> ['Command', 'Command.Container', 'Command.Container.ContainerConfig', - 'Command.Container.ContainerDeploy', 'Command.Container.ContainerKill', 'Command.Container.ContainerList', 'Command.Container.ContainerRemove', @@ -4131,7 +3853,6 @@ get_msg_containment("message") -> 'Message.Ping', 'Message.Pong', 'Message.Pub', - 'Message.TaskEvent', 'Request', 'Request.AuthRequest', 'Response', @@ -4175,7 +3896,6 @@ get_proto_by_msg_name_as_fqbin(<<"Stream.Reset">>) -> "message"; get_proto_by_msg_name_as_fqbin(<<"Stream.Opened">>) -> "message"; get_proto_by_msg_name_as_fqbin(<<"Request.AuthRequest">>) -> "message"; get_proto_by_msg_name_as_fqbin(<<"Request">>) -> "message"; -get_proto_by_msg_name_as_fqbin(<<"Message.TaskEvent">>) -> "message"; get_proto_by_msg_name_as_fqbin(<<"Command.Container.ContainerTarget">>) -> "message"; get_proto_by_msg_name_as_fqbin(<<"Command.Container.ContainerStart">>) -> "message"; get_proto_by_msg_name_as_fqbin(<<"Command.Container.ContainerList">>) -> "message"; @@ -4188,7 +3908,6 @@ get_proto_by_msg_name_as_fqbin(<<"Command.Container.ContainerRemove">>) -> "mess get_proto_by_msg_name_as_fqbin(<<"Message.Pong">>) -> "message"; get_proto_by_msg_name_as_fqbin(<<"Message.Ping">>) -> "message"; get_proto_by_msg_name_as_fqbin(<<"Command.Container.ContainerConfig">>) -> "message"; -get_proto_by_msg_name_as_fqbin(<<"Command.Container.ContainerDeploy">>) -> "message"; get_proto_by_msg_name_as_fqbin(<<"Command.Container.ContainerKill">>) -> "message"; get_proto_by_msg_name_as_fqbin(<<"Stream">>) -> "message"; get_proto_by_msg_name_as_fqbin(<<"Stream.Open">>) -> "message"; diff --git a/config/sys.config.src b/config/sys.config.src index abd8c5e..fa77a3b 100644 --- a/config/sys.config.src +++ b/config/sys.config.src @@ -15,11 +15,14 @@ {udp_port, 24000} ]}, - {stream_target, [ + {stream_targets, [ + {1, [ + {name, manager}, {host, "127.0.0.1"}, {port, 81}, {connect_timeout, 3000}, {idle_timeout, 120000} + ]} ]}, {heartbeat, [