From ddcb8fc600c0fad3d57f0fc008925181e2aa5339 Mon Sep 17 00:00:00 2001 From: hermes Date: Wed, 30 Sep 2026 21:59:32 +0200 Subject: [PATCH] add internal/infra/mqtt.go --- internal/infra/mqtt.go | 122 +++++++++++++++++++++++++++++++++++++++++ 1 file changed, 122 insertions(+) create mode 100644 internal/infra/mqtt.go diff --git a/internal/infra/mqtt.go b/internal/infra/mqtt.go new file mode 100644 index 0000000..9c8868b --- /dev/null +++ b/internal/infra/mqtt.go @@ -0,0 +1,122 @@ +package infra + +import ( + "bytes" + "fmt" + "net" + "time" +) + +// Minimal MQTT 3.1.1 client (QoS 0 only) — enough for now-playing +// announcements without pulling in a third-party dependency. +// +// Wire format references: OASIS MQTT v3.1.1, sections 2-3. +// All multi-byte integers are big-endian; strings are 2-byte length prefixed. + +// mqttEncodeRemainingLength encodes the variable-length integer used in every +// MQTT fixed header (7 bits per byte, high bit = "more bytes follow"). +func mqttEncodeRemainingLength(n int) []byte { + var out []byte + for { + b := byte(n % 128) + n /= 128 + if n > 0 { + b |= 0x80 + } + out = append(out, b) + if n == 0 { + return out + } + } +} + +// mqttStr encodes a UTF-8 string field (2-byte length prefix + bytes). +func mqttStr(s string) []byte { + return append([]byte{byte(len(s) >> 8), byte(len(s))}, s...) +} + +// mqttConnect builds a CONNECT packet (protocol 3.1.1, clean session, +// 30s keepalive). Username/password fields are included only when non-empty. +func mqttConnect(clientID, user, pass string) []byte { + var v bytes.Buffer + v.Write(mqttStr("MQTT")) // protocol name + v.WriteByte(4) // protocol level 3.1.1 + flags := byte(0x02) // clean session + if user != "" { + flags |= 0x80 + } + if pass != "" { + flags |= 0x40 + } + v.WriteByte(flags) + v.Write([]byte{0x00, 30}) // keepalive 30s + v.Write(mqttStr(clientID)) + if user != "" { + v.Write(mqttStr(user)) + } + if pass != "" { + v.Write(mqttStr(pass)) + } + body := v.Bytes() + return append([]byte{0x10}, append(mqttEncodeRemainingLength(len(body)), body...)...) +} + +// mqttPublish builds a QoS 0 PUBLISH packet (no packet id on the wire). +func mqttPublish(topic string, payload []byte) []byte { + var v bytes.Buffer + v.Write(mqttStr(topic)) + v.Write(payload) + body := v.Bytes() + return append([]byte{0x30}, append(mqttEncodeRemainingLength(len(body)), body...)...) +} + +// MQTTNowPlaying publishes now-playing JSON to a broker, one short-lived +// connection per publish (fire-and-forget, QoS 0). Safe to call every song. +type MQTTNowPlaying struct { + Addr string // host:port + User string + Pass string + ClientID string + Topic string // e.g. "airstudio/nowplaying" +} + +// Publish sends payload to the configured topic. Returns a descriptive error +// on any network failure; callers should log and continue — the stream must +// never stop because MQTT is down. +func (m *MQTTNowPlaying) Publish(payload []byte) error { + conn, err := net.DialTimeout("tcp", m.Addr, 5*time.Second) + if err != nil { + return fmt.Errorf("mqtt dial %s: %w", m.Addr, err) + } + defer conn.Close() + _ = conn.SetDeadline(time.Now().Add(5 * time.Second)) + + if _, err := conn.Write(mqttConnect(m.ClientID, m.User, m.Pass)); err != nil { + return fmt.Errorf("mqtt connect: %w", err) + } + // read CONNACK (4 bytes); verify return code 0 + ack := make([]byte, 4) + if _, err := readFull(conn, ack); err != nil { + return fmt.Errorf("mqtt connack: %w", err) + } + if ack[0] != 0x20 || ack[3] != 0x00 { + return fmt.Errorf("mqtt broker refused, code %d", ack[3]) + } + if _, err := conn.Write(mqttPublish(m.Topic, payload)); err != nil { + return fmt.Errorf("mqtt publish: %w", err) + } + return nil +} + +// readFull fills buf exactly or errors. +func readFull(r net.Conn, buf []byte) (int, error) { + got := 0 + for got < len(buf) { + n, err := r.Read(buf[got:]) + got += n + if err != nil { + return got, err + } + } + return got, nil +}