Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
6 changes: 6 additions & 0 deletions app/src/api.ts
Original file line number Diff line number Diff line change
Expand Up @@ -793,6 +793,12 @@ const realApi = {
headers: { "Content-Type": "application/json" },
body: JSON.stringify(cfg),
}),
testMqtt: (cfg: import("./types").MQTTSet) =>
request<{ ok: boolean; error?: string }>("/api/mqtt/test", {
method: "POST",
headers: { "Content-Type": "application/json" },
body: JSON.stringify(cfg),
}),
nftqos: () =>
request<import("./types").NftQoSProbe>("/api/nftqos"),
setNftqos: (limit: Partial<import("./types").NftQoSLimit>) =>
Expand Down
20 changes: 20 additions & 0 deletions app/src/components/system/MQTTCard.tsx
Original file line number Diff line number Diff line change
Expand Up @@ -23,6 +23,7 @@ export function MQTTCard({ index = 0 }: { index?: number }) {
const [intervalSec, setIntervalSec] = useState(60);
const [enabled, setEnabled] = useState(false);
const [saving, setSaving] = useState(false);
const [testing, setTesting] = useState(false);

const refresh = () => {
api.mqtt().then((s) => {
Expand Down Expand Up @@ -57,6 +58,22 @@ export function MQTTCard({ index = 0 }: { index?: number }) {
}
};

// test comprueba los valores del formulario contra el broker sin
// guardarlos (#409).
const test = async () => {
setTesting(true);
try {
const r = await api.testMqtt({ enabled, host, port, user, pass, nodeId, interval: intervalSec });
push(r.ok
? { tone: "ok", text: t("mqtt.testOk") }
: { tone: "danger", text: r.error || t("mqtt.testFail") });
} catch (e) {
push({ tone: "danger", text: e instanceof Error ? e.message : String(e) });
} finally {
setTesting(false);
}
};

const dotTone = !state?.enabled ? "muted" : state.connected ? "ok" : "danger";
const statusText = !state?.enabled
? t("mqtt.idle")
Expand Down Expand Up @@ -146,6 +163,9 @@ export function MQTTCard({ index = 0 }: { index?: number }) {
>
{t("mqtt.save")}
</Button>
<Button variant="secondary" onClick={test} loading={testing} disabled={!host}>
{t("mqtt.test")}
</Button>
</div>
</Card>
);
Expand Down
6 changes: 6 additions & 0 deletions app/src/demo/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -784,6 +784,12 @@ export const demoApi: typeof api = {
version: "v0.72.18",
};
},
testMqtt: async (cfg) => {
await wait(600, 1200);
if (!cfg.host) return { ok: false, error: "host is required" };
if (cfg.host.includes("fail")) return { ok: false, error: "connection refused" };
return { ok: true };
},
nftqos: () => get({ ...state.nftqos }),
setNftqos: async (limit) => {
await wait(800, 1500);
Expand Down
3 changes: 3 additions & 0 deletions app/src/locales/en.ts
Original file line number Diff line number Diff line change
Expand Up @@ -2321,6 +2321,9 @@ export default {
save: "Save",
saved: "Configuration saved",
idle: "MQTT disabled",
test: "Test connection",
testOk: "Broker connection succeeded",
testFail: "Could not connect to the broker",
connected: "Connected to the broker",
disconnected: "No connection to the broker",
},
Expand Down
3 changes: 3 additions & 0 deletions app/src/locales/es.ts
Original file line number Diff line number Diff line change
Expand Up @@ -2321,6 +2321,9 @@ export default {
save: "Guardar",
saved: "Configuración guardada",
idle: "MQTT desactivado",
test: "Probar conexión",
testOk: "Conexión con el broker correcta",
testFail: "No se pudo conectar con el broker",
connected: "Conectado al broker",
disconnected: "Sin conexión con el broker",
},
Expand Down
3 changes: 3 additions & 0 deletions internal/executor/executor.go
Original file line number Diff line number Diff line change
Expand Up @@ -81,6 +81,9 @@ func decodeOpArgs(kind string, raw json.RawMessage) ([]string, error) {
return []string{p}, nil
}
return []string{get("name")}, nil
case "mqtt.configure":
// Orden fijo: enabled host port user pass nodeId interval (#409).
return []string{get("enabled"), get("host"), get("port"), get("user"), get("pass"), get("nodeId"), get("interval")}, nil
default:
// Generic: collect the object values in a stable key order.
keys := []string{"config", "section", "option", "value", "service", "action"}
Expand Down
17 changes: 17 additions & 0 deletions internal/executor/executor_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -104,3 +104,20 @@ func TestUCIDeleteOtherErrorStillFails(t *testing.T) {
t.Fatal("a delete failure other than 'Entry not found' must be reported")
}
}

func TestDecodeOpArgsMQTTConfigure(t *testing.T) {
raw := json.RawMessage(`{"enabled":"1","host":"10.0.0.10","port":"1883","user":"u","pass":"p","nodeId":"rt3","interval":"60"}`)
args, err := decodeOpArgs("mqtt.configure", raw)
if err != nil {
t.Fatalf("decodeOpArgs: %v", err)
}
want := []string{"1", "10.0.0.10", "1883", "u", "p", "rt3", "60"}
if len(args) != len(want) {
t.Fatalf("args = %v, want %v", args, want)
}
for i := range want {
if args[i] != want[i] {
t.Fatalf("args[%d] = %q, want %q", i, args[i], want[i])
}
}
}
2 changes: 2 additions & 0 deletions internal/modules/executor_delegate.go
Original file line number Diff line number Diff line change
Expand Up @@ -72,6 +72,8 @@ func runExecutorOp(op executor.Op) error {
return fmt.Errorf("unsupported service action: %s", action)
case "install", "apk_install":
return executor.Run(executor.Op{Kind: "pkg_add", Args: op.Args})
case "mqtt.configure":
return applyMQTTConfigureOp(op)
default:
return fmt.Errorf("op kind %q not supported by NetGrip executor", op.Kind)
}
Expand Down
103 changes: 103 additions & 0 deletions internal/modules/mqtt_ops.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,103 @@
// mqtt_ops.go: prueba de conexión del broker y op de orquestación
// mqtt.configure (#409), que NetPulse envía a los routers NetGrip para
// configurarles el MQTT sin tocarlos a mano.
package modules

import (
"context"
"fmt"
"strconv"
"strings"
"time"

"github.com/gnacho/netgrip/internal/executor"
"github.com/gonzalop/mq"
)

// mqttProbe comprueba que el broker responde y acepta las credenciales: hace
// un CONNECT y desconecta. Sirve para el botón "probar conexión" de la tarjeta
// y para verificar un mqtt.configure antes de darlo por bueno.
func mqttProbe(ctx context.Context, cfg MQTTConfig, node string) error {
addr := fmt.Sprintf("tcp://%s:%d", cfg.Host, cfg.Port)
opts := []mq.Option{
mq.WithProtocolVersion(mq.ProtocolV311),
mq.WithClientID("netgrip-probe-" + node),
mq.WithConnectTimeout(5 * time.Second),
}
if cfg.User != "" {
opts = append(opts, mq.WithCredentials(cfg.User, cfg.Pass))
}
pctx, cancel := context.WithTimeout(ctx, 6*time.Second)
defer cancel()
client, err := mq.DialContext(pctx, addr, opts...)
if err != nil {
return err
}
client.Disconnect(pctx)
return nil
}

// MQTTTestConnection expone la sonda para la API: prueba los valores del
// formulario sin guardarlos.
func MQTTTestConnection(cfg MQTTConfig) error {
if strings.TrimSpace(cfg.Host) == "" {
return fmt.Errorf("host is required")
}
ctx, cancel := context.WithTimeout(context.Background(), 8*time.Second)
defer cancel()
return mqttProbe(ctx, cfg, mqttNodeID(cfg))
}

// mqttOpEnvPath es la ruta del env que escribe el op; inyectable en tests.
var mqttOpEnvPath = mqttEnvFile

// applyMQTTConfigureOp aplica un mqtt.configure venido de NetPulse: persiste
// /etc/netgrip/mqtt.env y lo aplica en caliente. Si el broker nuevo no
// responde, restaura la configuración anterior y devuelve error.
func applyMQTTConfigureOp(op executor.Op) error {
// Args: enabled host port user pass nodeId interval
if len(op.Args) != 7 {
return fmt.Errorf("mqtt.configure needs 7 args, got %d", len(op.Args))
}
enabled := op.Args[0] == "1" || strings.EqualFold(op.Args[0], "true")
port := mqttDefaultPort
if op.Args[2] != "" {
p, err := strconv.Atoi(op.Args[2])
if err != nil || p <= 0 || p > 65535 {
return fmt.Errorf("invalid port %q", op.Args[2])
}
port = p
}
interval, err := strconv.Atoi(op.Args[6])
if err != nil || interval <= 0 {
interval = mqttDefaultInterval
}
next := MQTTConfig{
Enabled: enabled,
Host: op.Args[1],
Port: port,
User: op.Args[3],
Pass: op.Args[4],
NodeID: op.Args[5],
Interval: interval,
}

prev, prevErr := ReadMQTTConfig(mqttOpEnvPath)
if err := setMQTTConfigAt(mqttOpEnvPath, next); err != nil {
return fmt.Errorf("write mqtt config: %w", err)
}
if !next.Enabled {
return nil
}
ctx, cancel := context.WithTimeout(context.Background(), 8*time.Second)
defer cancel()
if err := mqttProbe(ctx, next, mqttNodeID(next)); err != nil {
if prevErr == nil {
if rerr := setMQTTConfigAt(mqttOpEnvPath, prev); rerr != nil {
return fmt.Errorf("broker unreachable (%v) and rollback failed: %w", err, rerr)
}
}
return fmt.Errorf("broker unreachable with the new settings: %w", err)
}
return nil
}
54 changes: 54 additions & 0 deletions internal/modules/mqtt_ops_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,54 @@
package modules

import (
"os"
"path/filepath"
"testing"

"github.com/gnacho/netgrip/internal/executor"
)

func TestApplyMQTTConfigureOp(t *testing.T) {
dir := t.TempDir()
old := mqttOpEnvPath
defer func() { mqttOpEnvPath = old }()
mqttOpEnvPath = filepath.Join(dir, "mqtt.env")

// Configuración previa: activado.
if err := os.WriteFile(mqttOpEnvPath, []byte("MQTT_ENABLED=1\nMQTT_HOST=10.0.0.10\nMQTT_PORT=1883\n"), 0o600); err != nil {
t.Fatal(err)
}

// Args incompletos.
if err := applyMQTTConfigureOp(executor.Op{Kind: "mqtt.configure", Args: []string{"1", "h"}}); err == nil {
t.Fatal("must fail with incomplete args")
}
// Puerto inválido.
if err := applyMQTTConfigureOp(executor.Op{Kind: "mqtt.configure", Args: []string{"1", "h", "x", "", "", "", "60"}}); err == nil {
t.Fatal("must fail with an invalid port")
}

// Desactivado: aplica sin sondear el broker.
if err := applyMQTTConfigureOp(executor.Op{Kind: "mqtt.configure", Args: []string{"0", "", "", "", "", "", "60"}}); err != nil {
t.Fatalf("disable: %v", err)
}
cfg, err := ReadMQTTConfig(mqttOpEnvPath)
if err != nil {
t.Fatal(err)
}
if cfg.Enabled {
t.Fatal("must end up disabled")
}

// Activado con un broker inalcanzable: falla y revierte a lo anterior.
if err := applyMQTTConfigureOp(executor.Op{Kind: "mqtt.configure", Args: []string{"1", "127.0.0.1", "1", "", "", "", "60"}}); err == nil {
t.Fatal("must fail with an unreachable broker")
}
cfg, err = ReadMQTTConfig(mqttOpEnvPath)
if err != nil {
t.Fatal(err)
}
if cfg.Enabled {
t.Fatal("must have rolled back to the previous (disabled) config")
}
}
17 changes: 17 additions & 0 deletions internal/server/server.go
Original file line number Diff line number Diff line change
Expand Up @@ -236,6 +236,7 @@ func New(rpcdURL, version string) *Server {
s.mux.HandleFunc("POST /api/netpulse", s.requireAuth(s.handleNetPulseSet))
s.mux.HandleFunc("GET /api/mqtt", s.requireAuth(s.handleMQTTGet))
s.mux.HandleFunc("PUT /api/mqtt", s.requireAuth(s.handleMQTTSet))
s.mux.HandleFunc("POST /api/mqtt/test", s.requireAuth(s.handleMQTTTest))
s.mux.HandleFunc("POST /api/reboot", s.requireAuth(s.handleReboot))
s.mux.HandleFunc("GET /api/nftqos", s.requireAuth(s.handleNftQoSGet))
s.mux.HandleFunc("POST /api/nftqos", s.requireAuth(s.handleNftQoSSet))
Expand Down Expand Up @@ -2557,6 +2558,22 @@ func (s *Server) handleMQTTSet(w http.ResponseWriter, r *http.Request) {
writeJSON(w, modules.MQTTInfoNow())
}

// handleMQTTTest prueba los valores del formulario contra el broker sin
// guardarlos (#409).
func (s *Server) handleMQTTTest(w http.ResponseWriter, r *http.Request) {
var req mqttSetRequest
if err := json.NewDecoder(r.Body).Decode(&req); err != nil {
writeError(w, http.StatusBadRequest, "invalid json")
return
}
cfg := modules.MQTTConfig{Host: req.Host, Port: req.Port, User: req.User, Pass: req.Pass}
if err := modules.MQTTTestConnection(cfg); err != nil {
writeJSON(w, map[string]any{"ok": false, "error": err.Error()})
return
}
writeJSON(w, map[string]any{"ok": true})
}

func (s *Server) handleReboot(w http.ResponseWriter, _ *http.Request) {
if err := modules.RebootRouter(); err != nil {
writeError(w, http.StatusInternalServerError, err.Error())
Expand Down
Loading