diff --git a/app/src/api.ts b/app/src/api.ts index 1d5e025c..11d8cb6f 100644 --- a/app/src/api.ts +++ b/app/src/api.ts @@ -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("/api/nftqos"), setNftqos: (limit: Partial) => diff --git a/app/src/components/system/MQTTCard.tsx b/app/src/components/system/MQTTCard.tsx index a2eaff4e..ec29a4df 100644 --- a/app/src/components/system/MQTTCard.tsx +++ b/app/src/components/system/MQTTCard.tsx @@ -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) => { @@ -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") @@ -146,6 +163,9 @@ export function MQTTCard({ index = 0 }: { index?: number }) { > {t("mqtt.save")} + ); diff --git a/app/src/demo/index.ts b/app/src/demo/index.ts index ed212e21..a57aef11 100644 --- a/app/src/demo/index.ts +++ b/app/src/demo/index.ts @@ -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); diff --git a/app/src/locales/en.ts b/app/src/locales/en.ts index 0d6c1de3..b6dcf4ba 100644 --- a/app/src/locales/en.ts +++ b/app/src/locales/en.ts @@ -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", }, diff --git a/app/src/locales/es.ts b/app/src/locales/es.ts index 46f0ae4d..48c242c3 100644 --- a/app/src/locales/es.ts +++ b/app/src/locales/es.ts @@ -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", }, diff --git a/internal/executor/executor.go b/internal/executor/executor.go index 7a780b4c..4eb9adce 100644 --- a/internal/executor/executor.go +++ b/internal/executor/executor.go @@ -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"} diff --git a/internal/executor/executor_test.go b/internal/executor/executor_test.go index 0bec9fd0..a03f8eb0 100644 --- a/internal/executor/executor_test.go +++ b/internal/executor/executor_test.go @@ -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]) + } + } +} diff --git a/internal/modules/executor_delegate.go b/internal/modules/executor_delegate.go index c9ed6401..340b4ab2 100644 --- a/internal/modules/executor_delegate.go +++ b/internal/modules/executor_delegate.go @@ -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) } diff --git a/internal/modules/mqtt_ops.go b/internal/modules/mqtt_ops.go new file mode 100644 index 00000000..4e79c27e --- /dev/null +++ b/internal/modules/mqtt_ops.go @@ -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 +} diff --git a/internal/modules/mqtt_ops_test.go b/internal/modules/mqtt_ops_test.go new file mode 100644 index 00000000..db2ddb9a --- /dev/null +++ b/internal/modules/mqtt_ops_test.go @@ -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") + } +} diff --git a/internal/server/server.go b/internal/server/server.go index 9c644002..0d16e5a6 100644 --- a/internal/server/server.go +++ b/internal/server/server.go @@ -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)) @@ -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())