191 lines
4.7 KiB
Go
191 lines
4.7 KiB
Go
//go:build mqttintegration
|
|
|
|
package main
|
|
|
|
import (
|
|
"context"
|
|
"io"
|
|
"net"
|
|
"net/url"
|
|
"os"
|
|
"strconv"
|
|
"sync"
|
|
"testing"
|
|
"time"
|
|
|
|
mqtt "github.com/eclipse/paho.mqtt.golang"
|
|
|
|
"git.i3omb.com/gronod/ha-n95-local-control/internal/config"
|
|
"git.i3omb.com/gronod/ha-n95-local-control/internal/mqttbridge"
|
|
"git.i3omb.com/gronod/ha-n95-local-control/internal/robot"
|
|
"git.i3omb.com/gronod/ha-n95-local-control/internal/session"
|
|
)
|
|
|
|
func TestMQTTIntegration(t *testing.T) {
|
|
testURL := os.Getenv("MQTT_TEST_URL")
|
|
if testURL == "" {
|
|
t.Skip("MQTT_TEST_URL is not set; skipping integration test")
|
|
}
|
|
|
|
u, err := url.Parse(testURL)
|
|
if err != nil {
|
|
t.Fatalf("parse MQTT_TEST_URL: %v", err)
|
|
}
|
|
|
|
brokerHost := u.Hostname()
|
|
brokerPortStr := u.Port()
|
|
if brokerPortStr == "" {
|
|
brokerPortStr = "1883"
|
|
}
|
|
if _, err := strconv.Atoi(brokerPortStr); err != nil {
|
|
t.Fatalf("invalid broker port: %v", err)
|
|
}
|
|
brokerAddr := net.JoinHostPort(brokerHost, brokerPortStr)
|
|
|
|
// Start a local TCP proxy to simulate an abrupt network drop (unclean disconnect).
|
|
proxyLn, err := net.Listen("tcp", "127.0.0.1:0")
|
|
if err != nil {
|
|
t.Fatalf("start proxy listener: %v", err)
|
|
}
|
|
defer proxyLn.Close()
|
|
proxyPort := proxyLn.Addr().(*net.TCPAddr).Port
|
|
|
|
var proxyMu sync.Mutex
|
|
var proxyConns []net.Conn
|
|
proxyDone := make(chan struct{})
|
|
|
|
go func() {
|
|
defer close(proxyDone)
|
|
for {
|
|
clientConn, err := proxyLn.Accept()
|
|
if err != nil {
|
|
return
|
|
}
|
|
brokerConn, err := net.Dial("tcp", brokerAddr)
|
|
if err != nil {
|
|
clientConn.Close()
|
|
return
|
|
}
|
|
proxyMu.Lock()
|
|
proxyConns = append(proxyConns, clientConn, brokerConn)
|
|
proxyMu.Unlock()
|
|
|
|
go func() {
|
|
_, _ = io.Copy(brokerConn, clientConn)
|
|
_ = brokerConn.Close()
|
|
}()
|
|
go func() {
|
|
_, _ = io.Copy(clientConn, brokerConn)
|
|
_ = clientConn.Close()
|
|
}()
|
|
}
|
|
}()
|
|
|
|
serial := "testserial123"
|
|
discoveryTopic := "homeassistant/vacuum/ecovacs_" + serial + "/config"
|
|
availabilityTopic := "ecovacs/" + serial + "/availability"
|
|
|
|
var (
|
|
mu sync.Mutex
|
|
discoveryMsg string
|
|
availMsg string
|
|
availHistory []string
|
|
)
|
|
discoveryCh := make(chan string, 10)
|
|
availCh := make(chan string, 10)
|
|
|
|
// Connect test subscriber directly to broker.
|
|
subOpts := mqtt.NewClientOptions()
|
|
subOpts.AddBroker(testURL)
|
|
subOpts.SetClientID("integration-test-subscriber")
|
|
subOpts.SetCleanSession(true)
|
|
subClient := mqtt.NewClient(subOpts)
|
|
if token := subClient.Connect(); !token.WaitTimeout(5*time.Second) || token.Error() != nil {
|
|
t.Fatalf("subscriber connect failed: %v", token.Error())
|
|
}
|
|
defer subClient.Disconnect(250)
|
|
|
|
subClient.Subscribe(discoveryTopic, 0, func(_ mqtt.Client, m mqtt.Message) {
|
|
mu.Lock()
|
|
discoveryMsg = string(m.Payload())
|
|
mu.Unlock()
|
|
discoveryCh <- string(m.Payload())
|
|
})
|
|
|
|
subClient.Subscribe(availabilityTopic, 0, func(_ mqtt.Client, m mqtt.Message) {
|
|
mu.Lock()
|
|
availMsg = string(m.Payload())
|
|
availHistory = append(availHistory, string(m.Payload()))
|
|
mu.Unlock()
|
|
availCh <- string(m.Payload())
|
|
})
|
|
|
|
// Create bridge pointing at the proxy port.
|
|
cfg := config.Config{
|
|
MQTTHost: "127.0.0.1",
|
|
MQTTPort: proxyPort,
|
|
MQTTBase: "ecovacs",
|
|
HADiscoveryPrefix: "homeassistant",
|
|
MQTTClientID: "testbridge",
|
|
ControllerJID: "n95bridge@ecouser.net/homeassistant",
|
|
}
|
|
|
|
ctx, cancel := context.WithCancel(context.Background())
|
|
defer cancel()
|
|
|
|
submit := func(_ context.Context, _ string, _ robot.Command) error {
|
|
return nil
|
|
}
|
|
bridge, err := mqttbridge.New(ctx, cfg, submit)
|
|
if err != nil {
|
|
t.Fatalf("mqttbridge.New: %v", err)
|
|
}
|
|
|
|
// Trigger session announcement so bridge creates slot, connects, and publishes online.
|
|
bridge.AnnounceOK(session.ReadyEvent{
|
|
Serial: serial,
|
|
Generation: 1,
|
|
})
|
|
|
|
// 1. Assert retained discovery payload is published.
|
|
select {
|
|
case disc := <-discoveryCh:
|
|
if disc == "" {
|
|
t.Fatal("empty discovery payload received")
|
|
}
|
|
case <-time.After(5 * time.Second):
|
|
t.Fatal("timed out waiting for retained discovery payload")
|
|
}
|
|
|
|
// 2. Assert retained availability payload "online" is published.
|
|
select {
|
|
case avail := <-availCh:
|
|
if avail != "online" {
|
|
t.Fatalf("expected availability 'online', got: %q", avail)
|
|
}
|
|
case <-time.After(5 * time.Second):
|
|
t.Fatal("timed out waiting for retained availability online")
|
|
}
|
|
|
|
// 3. Disconnect bridge without clean disconnect (kill proxy connections so broker detects drop).
|
|
proxyMu.Lock()
|
|
for _, c := range proxyConns {
|
|
_ = c.Close()
|
|
}
|
|
proxyMu.Unlock()
|
|
_ = proxyLn.Close()
|
|
|
|
// 4. Assert broker fires Last Will and publishes "offline".
|
|
select {
|
|
case avail := <-availCh:
|
|
if avail != "offline" {
|
|
t.Fatalf("expected LWT 'offline', got: %q", avail)
|
|
}
|
|
case <-time.After(5 * time.Second):
|
|
t.Fatal("timed out waiting for broker retained will offline")
|
|
}
|
|
|
|
_ = discoveryMsg
|
|
_ = availMsg
|
|
}
|