From 7eb799c0100e7e309fcc6c9b5aacdee186c5e776 Mon Sep 17 00:00:00 2001 From: fcraviolatti Date: Thu, 26 Mar 2026 20:13:20 +0100 Subject: [PATCH 1/3] feat: add VerneMQ 2.x support and update default version to 2.1.2-alpine - Allow major version 2 in statefulset version check - Update defaultVerneMQVersion to 2.1.2-alpine --- controllers/statefulset.go | 2 ++ controllers/vernemq_controller.go | 2 +- 2 files changed, 3 insertions(+), 1 deletion(-) diff --git a/controllers/statefulset.go b/controllers/statefulset.go index 11b02c0d..b349035e 100644 --- a/controllers/statefulset.go +++ b/controllers/statefulset.go @@ -151,6 +151,8 @@ func makeStatefulSetSpec(instance *vernemqv1alpha1.VerneMQ) (*appsv1.StatefulSet if version.Minor < 7 { return nil, pkgerr.Errorf("unsupported VerneMQ minor version %s", version) } + case 2: + // VerneMQ 2.x supported default: return nil, pkgerr.Errorf("unsupported VerneMQ major version %s", version) } diff --git a/controllers/vernemq_controller.go b/controllers/vernemq_controller.go index 7fbf2083..fb2bd76d 100644 --- a/controllers/vernemq_controller.go +++ b/controllers/vernemq_controller.go @@ -37,7 +37,7 @@ const ( secretsDir = "/vernemq/etc/secrets/" sSetInputHashName = "vernemq-operator-input-hash" - defaultVerneMQVersion = "1.13.0-alpine" + defaultVerneMQVersion = "2.1.2-alpine" defaultVerneMQBaseImage = "vernemq/vernemq" defaultBundlerBaseImage = "vernemq/vmq-plugin-bundler" defaultBundlerVersion = "latest" From 4fa023dfc3f16dd45da5996f8be8ced1ddf089a4 Mon Sep 17 00:00:00 2001 From: fcraviolatti Date: Fri, 27 Mar 2026 09:38:46 +0100 Subject: [PATCH 2/3] add vmq_k8s Erlang plugin sources with OTP logger migration - Recover erlang/src/ files from old commit 44adcef7 (plugin was removed) - Migrate lager -> OTP logger throughout vmq_k8s_reloader.erl - Fix missing comma bug in cluster leave command args - Add rebar.config with yamerl dep (no lager) - Update bundler rebar dep to reference this fork instead of upstream - Update Dockerfile: golang:1.18 -> golang:1.23 - Add Dockerfile.prebuilt for pre-built binary workflow --- Dockerfile | 2 +- Dockerfile.prebuilt | 5 + controllers/deployment.go | 2 +- erlang/rebar.config | 5 + erlang/src/vmq_k8s.app.src | 17 ++ erlang/src/vmq_k8s_app.erl | 26 +++ erlang/src/vmq_k8s_reloader.erl | 321 ++++++++++++++++++++++++++++++++ erlang/src/vmq_k8s_sup.erl | 38 ++++ 8 files changed, 414 insertions(+), 2 deletions(-) create mode 100644 Dockerfile.prebuilt create mode 100644 erlang/rebar.config create mode 100644 erlang/src/vmq_k8s.app.src create mode 100644 erlang/src/vmq_k8s_app.erl create mode 100644 erlang/src/vmq_k8s_reloader.erl create mode 100644 erlang/src/vmq_k8s_sup.erl diff --git a/Dockerfile b/Dockerfile index 5a355c21..ff0ca306 100644 --- a/Dockerfile +++ b/Dockerfile @@ -1,5 +1,5 @@ # Build the manager binary -FROM golang:1.18 as builder +FROM golang:1.23 AS builder WORKDIR /workspace # Copy the Go Modules manifests diff --git a/Dockerfile.prebuilt b/Dockerfile.prebuilt new file mode 100644 index 00000000..4b322271 --- /dev/null +++ b/Dockerfile.prebuilt @@ -0,0 +1,5 @@ +FROM gcr.io/distroless/static:nonroot +WORKDIR / +COPY manager . +USER 65532:65532 +ENTRYPOINT ["/manager"] diff --git a/controllers/deployment.go b/controllers/deployment.go index adf0f7e9..7cd45725 100644 --- a/controllers/deployment.go +++ b/controllers/deployment.go @@ -103,7 +103,7 @@ func makeBundlerConfig(instance *vernemqv1alpha1.VerneMQ) string { config = config + fmt.Sprintf("{%s, {git, \"%s\", {%s, \"%s\"}}},\n", p.ApplicationName, p.RepoURL, p.VersionType, p.Version) } config = config + ` - {vmq_k8s, {git, "https://github.com/vernemq/vmq-operator", {branch, "master"}}} + {vmq_k8s, {git, "https://github.com/fcraviolatti/vmq-operator", {branch, "master"}}} ]}. ` return config diff --git a/erlang/rebar.config b/erlang/rebar.config new file mode 100644 index 00000000..22a1c6f3 --- /dev/null +++ b/erlang/rebar.config @@ -0,0 +1,5 @@ +{erl_opts, [debug_info]}. + +{deps, [ + {yamerl, "0.10.0"} +]}. diff --git a/erlang/src/vmq_k8s.app.src b/erlang/src/vmq_k8s.app.src new file mode 100644 index 00000000..505d386b --- /dev/null +++ b/erlang/src/vmq_k8s.app.src @@ -0,0 +1,17 @@ +{application, vmq_k8s, + [{description, "VMQ Kubernetes Helper Plugin"}, + {vsn, git}, + {registered, []}, + {mod, { vmq_k8s_app, []}}, + {applications, + [kernel, + stdlib, + yamerl + ]}, + {env,[]}, + {modules, []}, + + {maintainers, []}, + {licenses, []}, + {links, []} + ]}. diff --git a/erlang/src/vmq_k8s_app.erl b/erlang/src/vmq_k8s_app.erl new file mode 100644 index 00000000..73daba07 --- /dev/null +++ b/erlang/src/vmq_k8s_app.erl @@ -0,0 +1,26 @@ +%%%------------------------------------------------------------------- +%% @doc vmq_k8s public API +%% @end +%%%------------------------------------------------------------------- + +-module(vmq_k8s_app). + +-behaviour(application). + +%% Application callbacks +-export([start/2, stop/1]). + +%%==================================================================== +%% API +%%==================================================================== + +start(_StartType, _StartArgs) -> + vmq_k8s_sup:start_link(). + +%%-------------------------------------------------------------------- +stop(_State) -> + ok. + +%%==================================================================== +%% Internal functions +%%==================================================================== diff --git a/erlang/src/vmq_k8s_reloader.erl b/erlang/src/vmq_k8s_reloader.erl new file mode 100644 index 00000000..9cc9a0db --- /dev/null +++ b/erlang/src/vmq_k8s_reloader.erl @@ -0,0 +1,321 @@ +%%%------------------------------------------------------------------- +%%% @doc +%%% VerneMQ Kubernetes Helper Plugin — reloader +%%% Reads clusterview and config from ConfigMap-mounted files, +%%% applies VerneMQ config changes and manages cluster membership. +%%% @end +%%%------------------------------------------------------------------- +-module(vmq_k8s_reloader). + +-behaviour(gen_server). + +%% API +-export([start_link/0]). + +%% 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, {map=#{}, clustering_state=[]}). + +%%%=================================================================== +%%% API +%%%=================================================================== + +start_link() -> + gen_server:start_link({local, ?SERVER}, ?MODULE, [], []). + +%%%=================================================================== +%%% gen_server callbacks +%%%=================================================================== + +init([]) -> + erlang:send_after(1000, self(), check_config), + {ok, #state{}}. + +handle_call(_Request, _From, State) -> + {reply, ok, State}. + +handle_cast(_Msg, State) -> + {noreply, State}. + +handle_info(check_config, #state{map=ConfigState0, clustering_state=ClusteringState0} = State0) -> + ConfigFile = os:getenv("VMQ_CONFIGMAP", "/vernemq/etc/config.yaml"), + ConfigState1 = + try yamerl_constr:file(ConfigFile) of + Config -> + apply_config(Config, ConfigState0) + catch + throw:{yamerl_exception, Exception} -> + logger:error("Can't parse YAML config ~p", [Exception]), + ConfigState0; + E:R -> + logger:error("Error while parsing config ~p ~p", [E, R]), + ConfigState0 + end, + ClusterviewFile = os:getenv("VMQ_CLUSTERVIEW", "/vernemq/etc/configmaps/clusterview/clusterview.yaml"), + Ret = + case file:read_file(ClusterviewFile) of + {ok, Content} -> + Nodes = [N || N <- re:split(Content, ";"), N =/= <<>>], + check_clustering(Nodes, ClusteringState0); + {error, Reason} -> + logger:error("Can't read Clusterview File ~p due to ~p", [ClusterviewFile, Reason]), + ClusteringState0 + end, + case Ret of + its_over -> + {stop, normal, State0}; + _ -> + erlang:send_after(1000, self(), check_config), + {noreply, State0#state{map=ConfigState1, clustering_state=Ret}} + end. + +apply_config([Config|_], CurrentState) -> + State0 = #{}, + State1 = apply_plugins_config(proplists:get_value("plugins", Config, []), State0, CurrentState), + State2 = apply_listener_config(proplists:get_value("listeners", Config, []), State1, CurrentState), + State3 = apply_value_config(proplists:get_value("configs", Config, []), State2, CurrentState), + State3; +apply_config([], CurrentState) -> + CurrentState. + +apply_plugins_config([PluginConfig|Rest], Acc, CurrentState) -> + case proplists:get_value("name", PluginConfig) of + undefined -> + logger:error("Can't apply plugin config, missing 'name' in ~p", [PluginConfig]), + apply_plugins_config(Rest, Acc, CurrentState); + Name -> + case maps:is_key({plugin, Name}, CurrentState) of + true -> + apply_plugins_config(Rest, maps:put({plugin, Name}, PluginConfig, Acc), CurrentState); + false -> + PreCmds = proplists:get_value("preStart", PluginConfig, []), + ok = exec_commands(PreCmds), + Acc1 = + case proplists:get_value("path", PluginConfig) of + undefined -> + command(["plugin", "enable", "-n", Name], + succf({plugin, Name}, PluginConfig), Acc); + Path -> + command(["plugin", "enable", "-n", Name, "-p", Path], + succf({plugin, Name}, PluginConfig), Acc) + end, + PostCmds = proplists:get_value("postStart", PluginConfig, []), + ok = exec_commands(PostCmds), + apply_plugins_config(Rest, Acc1, CurrentState) + end + end; +apply_plugins_config([], NewState, OldState) -> + New = [Name || {plugin, Name} <- maps:keys(NewState)], + Old = [Name || {plugin, Name} <- maps:keys(OldState)], + ToBeDisabled = Old -- New, + lists:foreach(fun(Name) -> + Cfg = maps:get({plugin, Name}, OldState, []), + PreCmds = proplists:get_value("preStop", Cfg, []), + exec_commands(PreCmds), + command(["plugin", "disable", "-n", Name]), + PostCmds = proplists:get_value("postStop", Cfg, []), + exec_commands(PostCmds) + end, ToBeDisabled), + NewState; +apply_plugins_config(null, NewState, _OldState) -> + NewState. + +exec_commands([]) -> ok; +exec_commands([CmdConfig|Rest]) -> + TimeoutSeconds = proplists:get_value("timeoutSeconds", CmdConfig, 5), + Cmd = proplists:get_value("cmd", CmdConfig), + Ref = make_ref(), + Self = self(), + {Pid, MRef} = spawn_monitor( + fun() -> + Res = os:cmd(Cmd), + Self ! {exec_cmd_res, Ref, Cmd, Res} + end), + TimeoutMs = TimeoutSeconds*1000, + receive + {exec_cmd_res, Ref, Cmd, Res} -> + demonitor(MRef, [flush]), + logger:info("Execute \"~s\" \"~s\"", [Cmd, string:trim(Res)]); + {'DOWN', MRef, process, _, Info} -> + logger:info("Execute \"~s\" abnormally terminated (~p)", [Cmd, Info]) + after + TimeoutMs -> + exit(Pid, kill), + receive + {exec_cmd_res, Ref, Cmd, Res} -> + demonitor(MRef, [flush]), + logger:info("Execute \"~s\" \"~s\"", [Cmd, string:trim(Res)]); + {'DOWN', MRef, process, _, killed} -> + logger:info("Execute \"~s\" aborted due to timeout (~ps)", [Cmd, TimeoutSeconds]); + {'DOWN', MRef, process, _, Info} -> + logger:info("Execute \"~s\" abnormally terminated (~p)", [Cmd, Info]) + end + end, + exec_commands(Rest). + +to_snake_case(S) -> + string:lowercase(re:replace(S, "[A-Z]", "_&", [{return, list}, global])). + +maybe_addr_from_interface(AddrOrIf) -> + case inet:parse_address(AddrOrIf) of + {ok, _} -> + AddrOrIf; + {error, _} -> + {ok, Interfaces} = inet:getifaddrs(), + Interface = proplists:get_value(AddrOrIf, Interfaces, []), + inet:ntoa(proplists:get_value(addr, Interface, {127, 0, 0, 1})) + end. + +apply_listener_config([ListenerConfig|Rest], Acc, State) -> + case {proplists:get_value("address", ListenerConfig), + proplists:get_value("port", ListenerConfig)} of + {AddrOrIf, IPort} when AddrOrIf =/= undefined, IPort =/= undefined -> + Addr = maybe_addr_from_interface(AddrOrIf), + Port = integer_to_list(IPort), + ListenerConfig1 = [C || C = {K, _} <- ListenerConfig, + not lists:member(K, ["address", "port", "tlsConfig"])], + ListenerConfig2 = ListenerConfig1 ++ proplists:get_value("tlsConfig", ListenerConfig, []), + Flags0 = lists:foldl(fun({K, V}, CAcc) -> + case V of + true -> + ["--" ++ to_snake_case(K) | CAcc]; + false -> + CAcc; + Int when is_integer(Int) -> + ["--" ++ to_snake_case(K) ++ "=" ++ integer_to_list(V) | CAcc]; + _ -> + ["--" ++ to_snake_case(K) ++ "=" ++ V | CAcc] + end; + (_, CAcc) -> + CAcc + end, [], ListenerConfig2), + Flags1 = lists:usort(Flags0), + case maps:get({listener, {Addr, Port}}, State, undefined) of + Flags1 -> + apply_listener_config(Rest, maps:put({listener, {Addr, Port}}, Flags1, Acc), State); + _UndefOldFlags -> + DeleteCommand = ["listener", "delete", "address=" ++ Addr, "port=" ++ Port], + command(DeleteCommand), + StartCommand = ["listener", "start", "address=" ++ Addr, "port=" ++ Port] ++ Flags1, + Acc1 = command(StartCommand, succf({listener, {Addr, Port}}, Flags1), Acc), + apply_listener_config(Rest, Acc1, State) + end; + _ -> + logger:error("address or port not set in ~p", [ListenerConfig]), + apply_listener_config(Rest, Acc, State) + end; +apply_listener_config([], NewState, OldState) -> + New = [AddrPort || {listener, AddrPort} <- maps:keys(NewState)], + Old = [AddrPort || {listener, AddrPort} <- maps:keys(OldState)], + ToBeDeleted = Old -- New, + lists:foreach(fun({Addr, Port}) -> + command(["listener", "delete", "address=" ++ Addr, "port=" ++ Port]) + end, ToBeDeleted), + NewState; +apply_listener_config(null, NewState, _OldState) -> + NewState. + +apply_value_config([[{"name", ConfigKey0}, {"value", ConfigValue}]|Rest], Acc, CurrentState) -> + ConfigKey1 = to_snake_case(ConfigKey0), + case maps:get({config, ConfigKey1}, CurrentState, undefined) of + undefined -> + case default_val(ConfigKey1) of + {ok, Default} -> + Acc1 = command(["set", ConfigKey1 ++ "=" ++ ConfigValue], + succf({config, ConfigKey1}, {ConfigValue, Default}), Acc), + apply_value_config(Rest, Acc1, CurrentState); + {error, invalid_key} -> + logger:error("Invalid config key ~p", [ConfigKey1]), + apply_value_config(Rest, Acc, CurrentState) + end; + {ConfigValue, Default} -> + apply_value_config(Rest, maps:put({config, ConfigKey1}, {ConfigValue, Default}, Acc), CurrentState); + {_Other, Default} -> + Acc1 = command(["set", ConfigKey1 ++ "=" ++ ConfigValue], + succf({config, ConfigKey1}, {ConfigValue, Default}), Acc), + apply_value_config(Rest, Acc1, CurrentState) + end; +apply_value_config([], NewState, OldState) -> + New = [ConfigKey || {config, ConfigKey} <- maps:keys(NewState)], + Old = [ConfigKey || {config, ConfigKey} <- maps:keys(OldState)], + ToBeReset = Old -- New, + lists:foreach(fun(ConfigKey) -> + {_, DefaultValue} = maps:get(ConfigKey, OldState), + command(["set", ConfigKey ++ "=" ++ DefaultValue]) + end, ToBeReset), + NewState. + +default_val(ConfigKey) -> + case clique_config:show([ConfigKey], []) of + [{table, [Res]}] -> + {ok, proplists:get_value(ConfigKey, Res)}; + {error, {invalid_config_keys, _}} -> + {error, invalid_key} + end. + +check_clustering(CurNodes, OldNodes) -> + MySelf = atom_to_binary(node(), utf8), + case {lists:member(MySelf, CurNodes), + lists:member(MySelf, OldNodes)} of + {true, false} -> + case CurNodes -- [MySelf] of + [] -> + CurNodes; + [FirstNode|_] -> + command(["cluster", "join", "discovery-node=" ++ binary_to_list(FirstNode)], + fun(_) -> CurNodes end, OldNodes) + end; + {false, true} -> + command(["cluster", "leave", "node=" ++ binary_to_list(MySelf), "--timeout=3600", "--kill_sessions"]), + its_over; + _ -> + CurNodes + end. + +succf(K, V) -> + fun(A) -> maps:put(K, V, A) end. + +command(Args) -> + command(Args, fun(_) -> ignore end, ignore). + +command(Args, SuccessFun, Acc) -> + Cmd = ["vmq-admin" | Args], + try vmq_server_cli:command(Cmd, false) of + {ok, Ret} -> + logger:info("Execute: ~p ~p", [Cmd, proplists:get_value(text, Ret, "Done")]), + SuccessFun(Acc); + {error, [{alert, [{text, Txt}]}]} -> + Text = lists:flatten(Txt), + case string:find(Text, "already_enabled") of + nomatch -> + logger:error("Execute error: ~p ~p", [Cmd, Text]), + Acc; + _ -> + logger:info("Execute: ~p ~p", [Cmd, Text]), + SuccessFun(Acc) + end; + {error, Error} -> + logger:error("Execute error: ~p ~p", [Cmd, Error]), + Acc; + Other -> + logger:error("Execute error: ~p ~p", [Cmd, Other]), + Acc + catch + E:R -> + logger:error("Execute error: ~p ~p ~p", [Cmd, E, R]), + Acc + end. + +terminate(_Reason, _State) -> + ok. + +code_change(_OldVsn, State, _Extra) -> + {ok, State}. diff --git a/erlang/src/vmq_k8s_sup.erl b/erlang/src/vmq_k8s_sup.erl new file mode 100644 index 00000000..76879427 --- /dev/null +++ b/erlang/src/vmq_k8s_sup.erl @@ -0,0 +1,38 @@ +%%%------------------------------------------------------------------- +%% @doc vmq_k8s top level supervisor. +%% @end +%%%------------------------------------------------------------------- + +-module(vmq_k8s_sup). + +-behaviour(supervisor). + +%% API +-export([start_link/0]). + +%% Supervisor callbacks +-export([init/1]). + +-define(SERVER, ?MODULE). + +%%==================================================================== +%% API functions +%%==================================================================== + +start_link() -> + supervisor:start_link({local, ?SERVER}, ?MODULE, []). + +%%==================================================================== +%% Supervisor callbacks +%%==================================================================== + +%% Child :: {Id,StartFunc,Restart,Shutdown,Type,Modules} +init([]) -> + Reloader = {vmq_k8s_reloader, {vmq_k8s_reloader, start_link, []}, + permanent, 5000, worker, [vmq_k8s_reloader]}, + + {ok, { {one_for_all, 0, 1}, [Reloader]} }. + +%%==================================================================== +%% Internal functions +%%==================================================================== From 97beb999b0c7a13dedf69bd01fd09cab0f4316e0 Mon Sep 17 00:00:00 2001 From: "fabio.craviolatti" Date: Thu, 30 Apr 2026 17:57:16 +0200 Subject: [PATCH 3/3] fix: consolidate vmq_k8s Erlang sources to repo root and update upstream reference - Move Erlang sources from erlang/src/ to src/ at repo root for rebar3 compatibility (rebar3 fetches deps from root by default) - Move rebar.config to repo root accordingly - Update vmq_k8s dep reference to vernemq/vmq-operator --- controllers/deployment.go | 2 +- erlang/rebar.config => rebar.config | 0 {erlang/src => src}/vmq_k8s.app.src | 0 {erlang/src => src}/vmq_k8s_app.erl | 0 {erlang/src => src}/vmq_k8s_reloader.erl | 0 {erlang/src => src}/vmq_k8s_sup.erl | 0 6 files changed, 1 insertion(+), 1 deletion(-) rename erlang/rebar.config => rebar.config (100%) rename {erlang/src => src}/vmq_k8s.app.src (100%) rename {erlang/src => src}/vmq_k8s_app.erl (100%) rename {erlang/src => src}/vmq_k8s_reloader.erl (100%) rename {erlang/src => src}/vmq_k8s_sup.erl (100%) diff --git a/controllers/deployment.go b/controllers/deployment.go index 7cd45725..adf0f7e9 100644 --- a/controllers/deployment.go +++ b/controllers/deployment.go @@ -103,7 +103,7 @@ func makeBundlerConfig(instance *vernemqv1alpha1.VerneMQ) string { config = config + fmt.Sprintf("{%s, {git, \"%s\", {%s, \"%s\"}}},\n", p.ApplicationName, p.RepoURL, p.VersionType, p.Version) } config = config + ` - {vmq_k8s, {git, "https://github.com/fcraviolatti/vmq-operator", {branch, "master"}}} + {vmq_k8s, {git, "https://github.com/vernemq/vmq-operator", {branch, "master"}}} ]}. ` return config diff --git a/erlang/rebar.config b/rebar.config similarity index 100% rename from erlang/rebar.config rename to rebar.config diff --git a/erlang/src/vmq_k8s.app.src b/src/vmq_k8s.app.src similarity index 100% rename from erlang/src/vmq_k8s.app.src rename to src/vmq_k8s.app.src diff --git a/erlang/src/vmq_k8s_app.erl b/src/vmq_k8s_app.erl similarity index 100% rename from erlang/src/vmq_k8s_app.erl rename to src/vmq_k8s_app.erl diff --git a/erlang/src/vmq_k8s_reloader.erl b/src/vmq_k8s_reloader.erl similarity index 100% rename from erlang/src/vmq_k8s_reloader.erl rename to src/vmq_k8s_reloader.erl diff --git a/erlang/src/vmq_k8s_sup.erl b/src/vmq_k8s_sup.erl similarity index 100% rename from erlang/src/vmq_k8s_sup.erl rename to src/vmq_k8s_sup.erl