add internal/infra/mqtt.go

This commit is contained in:
2026-09-30 21:59:32 +02:00
parent c89600e6b2
commit ddcb8fc600
+122
View File
@@ -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
}