123 lines
3.3 KiB
Go
123 lines
3.3 KiB
Go
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
|
|
}
|