// Package miot provides MIoT core client for Xiaomi Home devices. // mips_client.go — full MIPS (MIoT Pub/Sub) client ported from py-miot/miot_mips.py. package miot import ( "encoding/binary" "encoding/json" "fmt" "math" "math/rand" "strings" "sync" "time" mqtt "github.com/eclipse/paho.mqtt.golang" "xiaomihome/logger" ) // ============================================================================ // MIPS protocol constants (aligned with Python _MipsClient) // ============================================================================ const ( mipsQoS = 2 mipsReconnectIntervalMin = 10.0 // seconds mipsReconnectIntervalMax = 600.0 // seconds mipsSubPatch = 300 mipsSubInterval = 1 // second mipsMQTTInterval = 1 // second uint32max = 0xFFFFFFFF MipsRequestTimeoutDefault = 10000 // milliseconds ) // ============================================================================ // MIPS message type options (aligned with Python _MipsMsgTypeOptions) // ============================================================================ const ( mipsMsgTypeID = 0 mipsMsgTypeRetTopic = 1 mipsMsgTypePayload = 2 mipsMsgTypeFrom = 3 ) // ============================================================================ // MipsMessage — binary message packing/unpacking // ============================================================================ // MipsMessage represents a packed MIPS protocol message. // Protocol: each segment is [4-byte length LE][1-byte type][len bytes + 1 null] type MipsMessage struct { Mid int MsgFrom string RetTopic string Payload string } // Pack packs a MIPS message into binary format. // Format: + data + null terminator func PackMipsMessage(mid int, payload string, msgFrom, retTopic string) ([]byte, error) { if payload == "" { return nil, &MIoTMipsError{MIoTError{Code: CodeInvalidParams, Message: "invalid mid or payload"}} } var buf []byte // mid segment midBuf := make([]byte, 4) binary.LittleEndian.PutUint32(midBuf, uint32(mid)) buf = append(buf, 0x04, 0x00, 0x00, 0x00) // length = 4 buf = append(buf, byte(mipsMsgTypeID)) buf = append(buf, midBuf...) // msg_from segment (optional) if msgFrom != "" { data := []byte(msgFrom) seg := buildMipsSegment(mipsMsgTypeFrom, data) buf = append(buf, seg...) } // ret_topic segment (optional) if retTopic != "" { data := []byte(retTopic) seg := buildMipsSegment(mipsMsgTypeRetTopic, data) buf = append(buf, seg...) } // payload segment data := []byte(payload) seg := buildMipsSegment(mipsMsgTypePayload, data) buf = append(buf, seg...) return buf, nil } // buildMipsSegment builds a single protocol segment. func buildMipsSegment(msgType int, data []byte) []byte { totalLen := len(data) + 1 // +1 for null terminator lenBuf := make([]byte, 4) binary.LittleEndian.PutUint32(lenBuf, uint32(totalLen)) var seg []byte seg = append(seg, lenBuf...) seg = append(seg, byte(msgType)) seg = append(seg, data...) seg = append(seg, 0) // null terminator return seg } // UnpackMipsMessage unpacks binary data into a MipsMessage. func UnpackMipsMessage(data []byte) *MipsMessage { msg := &MipsMessage{} dataLen := len(data) pos := 0 for pos < dataLen { if pos+5 > dataLen { break } segLen := int(binary.LittleEndian.Uint32(data[pos : pos+4])) segType := int(data[pos+4]) pos += 5 if segLen == 0 || pos+segLen > dataLen { break } segData := data[pos : pos+segLen] // Strip trailing null segData = bytesTrimNull(segData) switch segType { case mipsMsgTypeID: if len(segData) >= 4 { msg.Mid = int(binary.LittleEndian.Uint32(segData)) } case mipsMsgTypeRetTopic: msg.RetTopic = string(segData) case mipsMsgTypePayload: msg.Payload = string(segData) case mipsMsgTypeFrom: msg.MsgFrom = string(segData) } pos += segLen } return msg } func bytesTrimNull(b []byte) []byte { for len(b) > 0 && b[len(b)-1] == 0 { b = b[:len(b)-1] } return b } // String returns a debug representation of the message. func (m *MipsMessage) String() string { return fmt.Sprintf("%d, %s, %s, %s", m.Mid, m.MsgFrom, m.RetTopic, m.Payload) } // ============================================================================ // MIoTDeviceState — device state enum (aligned with Python MIoTDeviceState) // ============================================================================ // MIoTDeviceState represents the online/offline state of a device. type MIoTDeviceState int const ( // DeviceStateDisable means the device is disabled. DeviceStateDisable MIoTDeviceState = 0 // DeviceStateOffline means the device is offline. DeviceStateOffline MIoTDeviceState = 1 // DeviceStateOnline means the device is online. DeviceStateOnline MIoTDeviceState = 2 ) // String returns the string representation of MIoTDeviceState. func (s MIoTDeviceState) String() string { switch s { case DeviceStateDisable: return "disable" case DeviceStateOffline: return "offline" case DeviceStateOnline: return "online" default: return "unknown" } } // ============================================================================ // MipsRequest — pending request tracking // ============================================================================ // MipsRequest tracks a pending MIPS request with timeout. type MipsRequest struct { Mid int OnReply func(payload string, ctx interface{}) OnReplyCtx interface{} Timer *time.Timer } // ============================================================================ // MipsBroadcast — subscription tracking // ============================================================================ // MipsBroadcast represents a subscribed broadcast topic. type MipsBroadcast struct { Topic string Handler func(topic string, payload string, ctx interface{}) HandlerCtx interface{} } // ============================================================================ // mipsClient — core MIPS client (internal) // ============================================================================ // mipsClient is the concrete MIPS client using paho MQTT. type mipsClient struct { mu sync.Mutex lgr logger.Logger clientID string host string port int username string password string caFile string certFile string keyFile string // MQTT connection mqttConn MQTTClient // interface to paho MQTT mqttConnected bool // Reconnect reconnectActive bool reconnectInterval float64 reconnectTimer *time.Timer // Subscriptions pendingSubs map[string]int // topic -> retry count pendingSubMu sync.Mutex subMu sync.Mutex subscriptions map[string]*MipsBroadcast // Requests requestMu sync.Mutex requests map[string]*MipsRequest // mid as string -> request seedID int // State subscribers stateSubMu sync.Mutex stateSubs map[string]func(key string, connected bool) // Events onConnect func() onDisconnect func(err error) // Shutdown stopCh chan struct{} doneCh chan struct{} } // MQTTClient is a small interface for MQTT operations. // This allows decoupling from the concrete paho MQTT library. type MQTTClient interface { Connect() error Disconnect() IsConnected() bool Subscribe(topic string, qos byte, handler func(topic string, payload []byte)) error Unsubscribe(topic string) error Publish(topic string, qos byte, payload []byte) error SetCredentials(username, password string) SetTLS(caFile, certFile, keyFile string) error } // newMipsClient creates a new mipsClient (internal). func newMipsClient(clientID, host string, port int, username, password, caFile, certFile, keyFile string) *mipsClient { c := &mipsClient{ clientID: clientID, host: host, port: port, username: username, password: password, caFile: caFile, certFile: certFile, keyFile: keyFile, lgr: logger.Default(), seedID: rand.Intn(uint32max), pendingSubs: make(map[string]int), subscriptions: make(map[string]*MipsBroadcast), requests: make(map[string]*MipsRequest), stateSubs: make(map[string]func(key string, connected bool)), reconnectActive: true, stopCh: make(chan struct{}), doneCh: make(chan struct{}), } return c } // SetLogger sets a custom logger for mipsClient. func (c *mipsClient) SetLogger(l logger.Logger) { c.lgr = l } // connect starts the MQTT connection. func (c *mipsClient) connect(mqttClient MQTTClient) error { c.mu.Lock() c.mqttConn = mqttClient c.mu.Unlock() if err := c.mqttConn.Connect(); err != nil { return fmt.Errorf("mqtt connect: %w", err) } c.mu.Lock() c.mqttConnected = true c.reconnectActive = true c.reconnectInterval = 0 c.mu.Unlock() // Notify connect callbacks if c.onConnect != nil { c.onConnect() } // Notify state subscribers c.notifyStateSubs(true) return nil } // disconnect stops the MQTT connection. func (c *mipsClient) disconnect() { c.mu.Lock() c.reconnectActive = false if c.reconnectTimer != nil { c.reconnectTimer.Stop() c.reconnectTimer = nil } c.mu.Unlock() if c.mqttConn != nil { c.mqttConn.Disconnect() c.mu.Lock() c.mqttConnected = false c.mu.Unlock() // Notify disconnect if c.onDisconnect != nil { c.onDisconnect(nil) } c.notifyStateSubs(false) } } // isConnected returns the current connection state. func (c *mipsClient) isConnected() bool { c.mu.Lock() defer c.mu.Unlock() if c.mqttConn == nil { return false } return c.mqttConn.IsConnected() && c.mqttConnected } // subscribe adds a subscription to a topic. func (c *mipsClient) subscribe(topic string, handler func(topic string, payload string, ctx interface{}), ctx interface{}) error { c.subMu.Lock() defer c.subMu.Unlock() if c.subscriptions[topic] != nil { c.lgr.Debugf("[mips] re-register subscription: %s", topic) } c.subscriptions[topic] = &MipsBroadcast{ Topic: topic, Handler: handler, HandlerCtx: ctx, } // Add to pending subs for actual MQTT subscribe c.pendingSubMu.Lock() c.pendingSubs[topic] = 0 c.pendingSubMu.Unlock() go c.flushPendingSubs() return nil } // unsubscribe removes a subscription. func (c *mipsClient) unsubscribe(topic string) error { c.subMu.Lock() delete(c.subscriptions, topic) c.subMu.Unlock() if c.mqttConn != nil && c.mqttConn.IsConnected() { if err := c.mqttConn.Unsubscribe(topic); err != nil { c.lgr.Debugf("[mips] unsubscribe error %s: %v", topic, err) } } return nil } // publish sends a message on a topic. func (c *mipsClient) publish(topic string, payload []byte) error { if c.mqttConn == nil || !c.mqttConn.IsConnected() { return fmt.Errorf("mips publish when not connected: %s", topic) } return c.mqttConn.Publish(topic, byte(mipsQoS), payload) } // subState subscribes to connection state changes. func (c *mipsClient) subState(key string, handler func(key string, connected bool)) { c.stateSubMu.Lock() c.stateSubs[key] = handler c.stateSubMu.Unlock() } // unsubState unsubscribes from connection state changes. func (c *mipsClient) unsubState(key string) { c.stateSubMu.Lock() delete(c.stateSubs, key) c.stateSubMu.Unlock() } // notifyStateSubs notifies all state subscribers. func (c *mipsClient) notifyStateSubs(connected bool) { c.stateSubMu.Lock() subs := make(map[string]func(key string, connected bool), len(c.stateSubs)) for k, v := range c.stateSubs { subs[k] = v } c.stateSubMu.Unlock() for key, handler := range subs { if handler != nil { handler(key, connected) } } } // onMQTTConnect is called when the MQTT client connects. func (c *mipsClient) onMQTTConnect() { c.mu.Lock() c.mqttConnected = true c.reconnectInterval = 0 // reset backoff c.mu.Unlock() // Resubscribe all topics c.subMu.Lock() topics := make([]string, 0, len(c.subscriptions)) for topic := range c.subscriptions { topics = append(topics, topic) } c.subMu.Unlock() for _, topic := range topics { c.pendingSubMu.Lock() c.pendingSubs[topic] = 0 c.pendingSubMu.Unlock() } go c.flushPendingSubs() } // onMQTTDisconnect is called when the MQTT client disconnects. func (c *mipsClient) onMQTTDisconnect() { c.mu.Lock() c.mqttConnected = false // Clear pending subs c.pendingSubMu.Lock() c.pendingSubs = make(map[string]int) c.pendingSubMu.Unlock() c.mu.Unlock() // Try reconnect c.tryReconnect(false) } // onMQTTMessage handles an incoming MQTT message. func (c *mipsClient) onMQTTMessage(topic string, payload []byte) { payloadStr := string(payload) // Try exact match first c.subMu.Lock() bc, ok := c.subscriptions[topic] c.subMu.Unlock() if ok && bc != nil && bc.Handler != nil { bc.Handler(topic, payloadStr, bc.HandlerCtx) return } // Wildcard match: check subscriptions ending with # or + for subTopic, bc := range c.copySubs() { if bc == nil || bc.Handler == nil { continue } if matchWildcard(subTopic, topic) { bc.Handler(topic, payloadStr, bc.HandlerCtx) } } } // copySubs returns a copy of subscriptions for safe wildcard iteration. func (c *mipsClient) copySubs() map[string]*MipsBroadcast { c.subMu.Lock() defer c.subMu.Unlock() result := make(map[string]*MipsBroadcast, len(c.subscriptions)) for k, v := range c.subscriptions { result[k] = v } return result } // matchWildcard checks if a topic matches a subscription pattern containing # or +. func matchWildcard(pattern, topic string) bool { // Remove trailing /# from pattern and check prefix if strings.HasSuffix(pattern, "/#") { prefix := strings.TrimSuffix(pattern, "/#") return strings.HasPrefix(topic, prefix+"/") || topic == prefix } return pattern == topic } // flushPendingSubs processes pending subscriptions in batches. // Aligns with Python: __mips_sub_internal_pending_handler. func (c *mipsClient) flushPendingSubs() { if c.mqttConn == nil || !c.mqttConn.IsConnected() { return } c.pendingSubMu.Lock() subbed := 0 for topic, count := range c.pendingSubs { if subbed >= mipsSubPatch { break } if count > 3 { delete(c.pendingSubs, topic) c.lgr.Debugf("[mips] retry sub exceeded: %s", topic) continue } subbed++ if err := c.mqttConn.Subscribe(topic, byte(mipsQoS), func(t string, p []byte) { c.onMQTTMessage(t, p) }); err == nil { delete(c.pendingSubs, topic) c.lgr.Debugf("[mips] sub success: %s", topic) } else { c.pendingSubs[topic] = count + 1 c.lgr.Debugf("[mips] retry sub %d: %s %v", count, topic, err) } } hasMore := len(c.pendingSubs) > 0 c.pendingSubMu.Unlock() if hasMore { time.Sleep(time.Duration(mipsSubInterval) * time.Second) c.flushPendingSubs() } } // sendRequest sends a MIPS request and waits for a reply. func (c *mipsClient) sendRequest(topic string, payload string, timeout time.Duration) (string, error) { mid := c.nextMid() replyCh := make(chan string, 1) req := &MipsRequest{ Mid: mid, OnReply: func(payload string, ctx interface{}) { select { case ctx.(chan string) <- payload: default: } }, OnReplyCtx: replyCh, } c.requestMu.Lock() c.requests[fmt.Sprintf("%d", mid)] = req c.requestMu.Unlock() // Build packed message packed, err := PackMipsMessage(mid, payload, "", topic) if err != nil { c.requestMu.Lock() delete(c.requests, fmt.Sprintf("%d", mid)) c.requestMu.Unlock() return "", err } if err := c.publish(topic, packed); err != nil { c.requestMu.Lock() delete(c.requests, fmt.Sprintf("%d", mid)) c.requestMu.Unlock() return "", err } // Wait for reply or timeout select { case reply := <-replyCh: return reply, nil case <-time.After(timeout): c.requestMu.Lock() delete(c.requests, fmt.Sprintf("%d", mid)) c.requestMu.Unlock() return `{"error":{"code":-10006,"message":"timeout"}}`, nil } } // handleReply processes a reply from mips. func (c *mipsClient) handleReply(mid int, payload string) { key := fmt.Sprintf("%d", mid) c.requestMu.Lock() req, ok := c.requests[key] if ok { delete(c.requests, key) } c.requestMu.Unlock() if req != nil && req.OnReply != nil { req.OnReply(payload, req.OnReplyCtx) } } // nextMid generates the next MIPS message ID. func (c *mipsClient) nextMid() int { c.mu.Lock() mid := c.seedID c.seedID = int((c.seedID + 1) % uint32max) c.mu.Unlock() return mid } // tryReconnect schedules a reconnection attempt with exponential backoff. func (c *mipsClient) tryReconnect(immediate bool) { c.mu.Lock() if !c.reconnectActive { c.mu.Unlock() return } if c.reconnectTimer != nil { c.reconnectTimer.Stop() } var interval time.Duration if !immediate { if c.reconnectInterval < mipsReconnectIntervalMin { c.reconnectInterval = mipsReconnectIntervalMin } else { c.reconnectInterval = math.Min(c.reconnectInterval*2, mipsReconnectIntervalMax) } interval = time.Duration(c.reconnectInterval) * time.Second c.lgr.Debugf("[mips] reconnect after %v", interval) } c.reconnectTimer = time.AfterFunc(interval, func() { if c.mqttConn != nil { if err := c.mqttConn.Connect(); err != nil { c.lgr.Debugf("[mips] reconnect failed: %v", err) c.tryReconnect(false) } } }) c.mu.Unlock() } // ============================================================================ // MipsCloudClient — cloud MIPS client (public API) // ============================================================================ // MipsCloudClient wraps the MIPS client for cloud connections. type MipsCloudClient struct { client *mipsClient onConn func() onDisc func(error) } // SubState subscribes to MIPS cloud connection state changes. // key is a unique subscriber identifier (e.g. "{uid}-{cloudServer}"). func (c *MipsCloudClient) SubState(key string, handler func(key string, connected bool)) { c.client.subState(key, handler) } // UnsubState unsubscribes from MIPS cloud connection state changes. func (c *MipsCloudClient) UnsubState(key string) { c.client.unsubState(key) } // UpdateAccessToken updates the MQTT password (access token) and reconnects. // Aligns with Python: __update_mips_access_token. func (c *MipsCloudClient) UpdateAccessToken(accessToken string) { c.client.mu.Lock() c.client.password = accessToken wasConnected := c.client.mqttConnected c.client.mu.Unlock() if wasConnected { c.Disconnect() c.Connect() } } // NewMipsCloudClient creates a new MipsCloudClient. // broker format: "ssl://cn-ha.mqtt.io.mi.com:8883" func NewMipsCloudClient(broker, clientID, username, password string) *MipsCloudClient { host, port := parseMQTTBroker(broker) c := &MipsCloudClient{ client: newMipsClient(clientID, host, port, username, password, "", "", ""), } return c } // parseMQTTBroker parses a broker URL like "ssl://host:port". func parseMQTTBroker(broker string) (string, int) { // Strip protocol prefix url := broker if strings.HasPrefix(url, "ssl://") { url = url[6:] } else if strings.HasPrefix(url, "tcp://") { url = url[6:] } host := url port := 8883 if idx := strings.LastIndex(url, ":"); idx >= 0 { host = url[:idx] fmt.Sscanf(url[idx+1:], "%d", &port) } return host, port } // Connect connects to the cloud MQTT broker. func (c *MipsCloudClient) Connect() error { mqClient := newPahoMQTTClient(c.client.clientID, c.client.host, c.client.port, c.client.username, c.client.password, "", "", "") mqClient.onConnect = c.client.onMQTTConnect mqClient.onDisconnect = c.client.onMQTTDisconnect mqClient.onMessage = c.client.onMQTTMessage c.client.onConnect = func() { if c.onConn != nil { c.onConn() } } c.client.onDisconnect = func(err error) { if c.onDisc != nil { c.onDisc(err) } } return c.client.connect(mqClient) } // Disconnect disconnects from the MQTT broker. func (c *MipsCloudClient) Disconnect() { c.client.disconnect() } // IsConnected returns true if connected. func (c *MipsCloudClient) IsConnected() bool { return c.client.isConnected() } // Subscribe subscribes to a topic with a handler. func (c *MipsCloudClient) Subscribe(topic string, handler func(string, []byte)) error { return c.client.subscribe(topic, func(_ string, payload string, _ interface{}) { handler(topic, []byte(payload)) }, nil) } // Unsubscribe unsubscribes from a topic. func (c *MipsCloudClient) Unsubscribe(topic string) error { return c.client.unsubscribe(topic) } // Publish publishes a message to a topic. func (c *MipsCloudClient) Publish(topic, payload string) error { return c.client.publish(topic, []byte(payload)) } // SubMipsState registers a handler for MIPS state changes. func (c *MipsCloudClient) SubMipsState(handler func(string, interface{})) { c.onConn = func() { handler("connected", nil) } c.onDisc = func(err error) { handler("disconnected", err) } } // GetDevList sends a get_dev_list request through MIPS. func (c *MipsCloudClient) GetDevList(payload string, timeout time.Duration) (string, error) { return c.client.sendRequest("get_dev_list", payload, timeout) } // ============================================================================ // MipsLocalClient — local MIPS client (public API) // ============================================================================ // MipsLocalClient wraps the MIPS client for local connections. type MipsLocalClient struct { client *mipsClient did string // GroupID is the home group ID this client belongs to. GroupID string // OnDevListChanged is called when the gateway's device list changes. OnDevListChanged func() } // NewMipsLocalClient creates a new MipsLocalClient. func NewMipsLocalClient(did, host string, port int, caFile, certFile, keyFile string) *MipsLocalClient { return &MipsLocalClient{ client: newMipsClient(did, host, port, "", "", caFile, certFile, keyFile), did: did, } } // Connect connects to the local MQTT broker. func (c *MipsLocalClient) Connect() error { mqClient := newPahoMQTTClient(c.client.clientID, c.client.host, c.client.port, "", "", c.client.caFile, c.client.certFile, c.client.keyFile) mqClient.onConnect = c.client.onMQTTConnect mqClient.onDisconnect = c.client.onMQTTDisconnect mqClient.onMessage = func(topic string, payload []byte) { // Try unpacking as MIPS message first msg := UnpackMipsMessage(payload) if msg != nil && msg.Mid != 0 { c.client.handleReply(msg.Mid, msg.Payload) } // Also pass to generic handler c.client.onMQTTMessage(topic, payload) } return c.client.connect(mqClient) } // Disconnect disconnects from the local MQTT broker. func (c *MipsLocalClient) Disconnect() { c.client.disconnect() } // IsConnected returns true if connected. func (c *MipsLocalClient) IsConnected() bool { return c.client.isConnected() } // Subscribe subscribes to a topic. func (c *MipsLocalClient) Subscribe(topic string, handler func(string, []byte)) error { return c.client.subscribe(topic, func(_ string, payload string, _ interface{}) { handler(topic, []byte(payload)) }, nil) } // Unsubscribe unsubscribes from a topic. func (c *MipsLocalClient) Unsubscribe(topic string) error { return c.client.unsubscribe(topic) } // GetDevList returns the device list from the gateway. func (c *MipsLocalClient) GetDevList(payload string, timeout time.Duration) (string, error) { return c.client.sendRequest("getDevList", payload, timeout) } // GetProps fetches property values from the gateway via MIPS request. // params is a JSON-encoded string of property parameters. func (c *MipsLocalClient) GetProps(params []PropParam, timeout time.Duration) ([]PropResult, error) { batch := make([]map[string]interface{}, len(params)) for i, p := range params { batch[i] = map[string]interface{}{ "did": p.DID, "siid": float64(p.Siid), "piid": float64(p.Piid), } } payloadBytes, _ := json.Marshal(map[string]interface{}{ "method": "properties/get", "params": batch, }) payload, err := c.client.sendRequest("properties/get", string(payloadBytes), timeout) if err != nil { return nil, err } // Parse response var result struct { Result []PropResult `json:"result"` } if err := json.Unmarshal([]byte(payload), &result); err != nil { return nil, fmt.Errorf("parse gateway props response: %w", err) } return result.Result, nil } // SubState subscribes to MIPS local connection state changes. func (c *MipsLocalClient) SubState(key string, handler func(key string, connected bool)) { c.client.subState(key, handler) } // UnsubState unsubscribes from MIPS local connection state changes. func (c *MipsLocalClient) UnsubState(key string) { c.client.unsubState(key) } // DID returns the device ID. func (c *MipsLocalClient) DID() string { return c.did } // ============================================================================ // MipsDeviceState — device state tracking // ============================================================================ // MipsDeviceState represents the state of a device's MIPS connection. type MipsDeviceState struct { DID string Online bool Source string Handler func(string, interface{}) } // NewMipsDeviceState creates a new MipsDeviceState. func NewMipsDeviceState(did string, online bool, source string) *MipsDeviceState { return &MipsDeviceState{ DID: did, Online: online, Source: source, } } // ============================================================================ // Paho MQTT Client adapter // ============================================================================ // pahoMQTTClient adapts paho.mqtt.golang to our MQTTClient interface. type pahoMQTTClient struct { clientID string host string port int username string password string caFile string certFile string keyFile string lgr logger.Logger onConnect func() onDisconnect func() onMessage func(topic string, payload []byte) client mqtt.Client connected bool } func newPahoMQTTClient(clientID, host string, port int, username, password, caFile, certFile, keyFile string) *pahoMQTTClient { c := &pahoMQTTClient{ clientID: clientID, host: host, port: port, username: username, password: password, caFile: caFile, certFile: certFile, keyFile: keyFile, lgr: logger.Default(), } return c } // SetLogger sets a custom logger for pahoMQTTClient. func (p *pahoMQTTClient) SetLogger(l logger.Logger) { p.lgr = l } func (p *pahoMQTTClient) Connect() error { opts := mqtt.NewClientOptions() opts.AddBroker(fmt.Sprintf("ssl://%s:%d", p.host, p.port)) opts.SetClientID(p.clientID) opts.SetUsername(p.username) opts.SetPassword(p.password) opts.SetKeepAlive(time.Duration(MIHOME_MQTT_KEEPALIVE) * time.Second) opts.SetAutoReconnect(false) opts.SetCleanSession(true) opts.OnConnect = func(c mqtt.Client) { p.lgr.Debugf("[mips] MQTT connected to %s:%d", p.host, p.port) p.connected = true if p.onConnect != nil { p.onConnect() } } opts.OnConnectionLost = func(c mqtt.Client, err error) { p.lgr.Debugf("[mips] MQTT connection lost: %v", err) p.connected = false if p.onDisconnect != nil { p.onDisconnect() } } opts.DefaultPublishHandler = func(c mqtt.Client, msg mqtt.Message) { p.lgr.Debugf("[mips] MQTT message: %s (%d bytes) data:%s", msg.Topic(), len(msg.Payload()), string(msg.Payload())) if p.onMessage != nil { p.onMessage(msg.Topic(), msg.Payload()) } } p.client = mqtt.NewClient(opts) token := p.client.Connect() if token.Wait() && token.Error() != nil { return token.Error() } p.lgr.Debugf("[mips] MQTT connected to %s:%d as %s", p.host, p.port, p.clientID) return nil } func (p *pahoMQTTClient) Disconnect() { if p.client != nil && p.client.IsConnected() { p.client.Disconnect(250) } p.connected = false if p.onDisconnect != nil { p.onDisconnect() } } func (p *pahoMQTTClient) IsConnected() bool { return p.connected } func (p *pahoMQTTClient) Subscribe(topic string, qos byte, handler func(topic string, payload []byte)) error { if p.client == nil || !p.client.IsConnected() { return fmt.Errorf("not connected") } // Use DefaultPublishHandler for message delivery (set in Connect). // Don't set a per-subscription callback to avoid double-firing with DefaultPublishHandler. token := p.client.Subscribe(topic, qos, nil) if token.Wait() && token.Error() != nil { return token.Error() } return nil } func (p *pahoMQTTClient) Unsubscribe(topic string) error { if p.client == nil || !p.client.IsConnected() { return fmt.Errorf("not connected") } token := p.client.Unsubscribe(topic) if token.Wait() && token.Error() != nil { return token.Error() } p.lgr.Debugf("[mips] unsub success: %s", topic) return nil } func (p *pahoMQTTClient) Publish(topic string, qos byte, payload []byte) error { if p.client == nil || !p.client.IsConnected() { return fmt.Errorf("not connected") } token := p.client.Publish(topic, qos, false, payload) if token.Wait() && token.Error() != nil { return token.Error() } return nil } func (p *pahoMQTTClient) SetCredentials(username, password string) { p.username = username p.password = password } func (p *pahoMQTTClient) SetTLS(caFile, certFile, keyFile string) error { p.caFile = caFile p.certFile = certFile p.keyFile = keyFile return nil }