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 }