// Package miot provides MIoT core client for Xiaomi Home devices. // miot_lan.go — MIoT LAN device control, ported from py-miot/miot_lan.py. // // Only supports MIoT SPEC-v2 WiFi devices on the local network. package miot import ( "crypto/aes" "crypto/cipher" "crypto/md5" "crypto/rand" "encoding/binary" "encoding/json" "fmt" "math/big" "net" "strings" "sync" "sync/atomic" "time" "xiaomihome/logger" ) // ============================================================================ // Protocol constants (aligned with Python _MIoTLanDevice / MIoTLan) // ============================================================================ const ( otHeader = 0x2131 otHeaderLen = 32 otPort = 54321 otProbeLen = 32 otMsgLen = 1400 otSupportWildcardSub = 0xFE // Keep-alive intervals (seconds) kaIntervalMin = 10.0 kaIntervalMax = 50.0 fastPingInterval = 5.0 constructStatePending = 15.0 // Network instability detection networkUnstableCntTh = 10 networkUnstableTimeTh = 120.0 networkUnstableResumeTh = 300.0 // Scan intervals (seconds) otProbeIntervalMin = 5.0 otProbeIntervalMax = 45.0 ) // ============================================================================ // LanDeviceState — device keep-alive state machine // ============================================================================ // LanDeviceState represents the keep-alive state of a LAN device. type LanDeviceState int const ( LanDeviceFresh LanDeviceState = iota // 0 — just received a ping response LanDevicePing1 // 1 — first fast ping attempt LanDevicePing2 // 2 — second fast ping attempt LanDevicePing3 // 3 — third fast ping attempt LanDeviceDead // 4 — device considered offline ) // String returns a human-readable name. func (s LanDeviceState) String() string { switch s { case LanDeviceFresh: return "fresh" case LanDevicePing1: return "ping1" case LanDevicePing2: return "ping2" case LanDevicePing3: return "ping3" case LanDeviceDead: return "dead" default: return "unknown" } } // ============================================================================ // LanDevice — single LAN device with encryption // ============================================================================ // LanDevice manages the encrypted UDP communication and keep-alive for one // MIoT SPEC-v2 WiFi device on the local network. type LanDevice struct { DID string Token []byte IP string IfName string manager *MIoTLan cipher cipher.Block aesIV []byte offset int subscribed bool subTS int supportedWildcardSub bool subLocked bool state LanDeviceState online bool onlineHistory []onlineEvent unstableTimer *time.Timer kaTimer *time.Timer kaInterval float64 } type onlineEvent struct { ts int online bool } // newLanDevice creates a new LanDevice and initializes its AES cipher. func newLanDevice(mgr *MIoTLan, did, tokenHex, ip string) (*LanDevice, error) { token, err := hexDecode(tokenHex) if err != nil { return nil, fmt.Errorf("invalid token hex: %w", err) } if len(token) != 16 { return nil, fmt.Errorf("token must be 16 bytes (32 hex chars), got %d", len(token)) } aesKey := md5Sum(token) aesIV := md5Sum(append(aesKey[:], token...)) block, err := aes.NewCipher(aesKey[:]) if err != nil { return nil, fmt.Errorf("create aes cipher: %w", err) } d := &LanDevice{ DID: did, Token: token, IP: ip, manager: mgr, cipher: block, aesIV: aesIV[:], state: LanDeviceDead, online: false, kaInterval: kaIntervalMin, onlineHistory: make([]onlineEvent, 0), } // Schedule initial keep-alive transition after a randomized delay. delay := randomizeFloat(constructStatePending, 0.5) d.kaTimer = time.AfterFunc(time.Duration(delay*float64(time.Second)), func() { d.updateKeepAlive(LanDeviceDead) }) mgr.lgr.Infof("[lan] device added: did=%s", d.DID) return d, nil } // KeepAlive is called when a probe response is received from the device. func (d *LanDevice) KeepAlive(ip, ifName string) { d.IP = ip if d.IfName != ifName { d.IfName = ifName d.manager.lgr.Infof("[lan] device if_name changed: %s, did=%s", d.IfName, d.DID) } d.updateKeepAlive(LanDeviceFresh) } // Online returns whether the device is currently considered online. func (d *LanDevice) Online() bool { return d.online } // setOnline transitions the online state and broadcasts the change. func (d *LanDevice) setOnline(v bool) { if d.online == v { return } d.online = v d.manager.broadcastDeviceState(d.DID, map[string]interface{}{ "online": d.online, "push_available": d.subscribed, }) } // GenPacket encrypts a JSON message and writes the wire-format packet into out. // Returns the number of bytes written. // Wire format: [0x2131(2)][len(2)][did(8)][timestamp(4)][md5(16)] + encrypted_body func (d *LanDevice) GenPacket(out []byte, msg map[string]interface{}) (int, error) { clearBytes, err := json.Marshal(msg) if err != nil { return 0, fmt.Errorf("marshal message: %w", err) } // PKCS7 padding blockSize := aes.BlockSize padding := blockSize - len(clearBytes)%blockSize padded := make([]byte, len(clearBytes)+padding) copy(padded, clearBytes) for i := len(clearBytes); i < len(padded); i++ { padded[i] = byte(padding) } if len(padded)+otHeaderLen > len(out) { return 0, fmt.Errorf("packet too long: %d > %d", len(padded)+otHeaderLen, len(out)) } // Encrypt enc := cipher.NewCBCEncrypter(d.cipher, d.aesIV) enc.CryptBlocks(padded, padded) dataLen := len(padded) + otHeaderLen // Header (bytes 0-15) binary.BigEndian.PutUint16(out[0:2], otHeader) binary.BigEndian.PutUint16(out[2:4], uint16(dataLen)) binary.BigEndian.PutUint64(out[4:12], mustParseUint64(d.DID)) binary.BigEndian.PutUint32(out[12:16], uint32(int(time.Now().Unix())-d.offset)) copy(out[16:32], d.Token) // Body copy(out[32:dataLen], padded) // MD5 over [0:dataLen] md5Hash := md5Sum(out[:dataLen]) copy(out[16:32], md5Hash[:]) return dataLen, nil } // DecryptPacket decrypts and authenticates a received packet. func (d *LanDevice) DecryptPacket(data []byte) (map[string]interface{}, error) { if len(data) < otHeaderLen { return nil, fmt.Errorf("packet too short: %d", len(data)) } dataLen := int(binary.BigEndian.Uint16(data[2:4])) if dataLen > len(data) { return nil, fmt.Errorf("declared len %d exceeds buffer %d", dataLen, len(data)) } // Verify MD5 origMD5 := make([]byte, 16) copy(origMD5, data[16:32]) // Replace MD5 slot with token for verification saved := make([]byte, 16) copy(saved, data[16:32]) copy(data[16:32], d.Token) calcMD5 := md5Sum(data[:dataLen]) // Restore copy(data[16:32], saved) if string(origMD5) != string(calcMD5[:]) { return nil, fmt.Errorf("invalid md5") } // Decrypt encrypted := data[32:dataLen] if len(encrypted) == 0 { return nil, fmt.Errorf("empty payload") } dec := cipher.NewCBCDecrypter(d.cipher, d.aesIV) decrypted := make([]byte, len(encrypted)) dec.CryptBlocks(decrypted, encrypted) // PKCS7 unpadding if len(decrypted) == 0 { return nil, fmt.Errorf("empty decrypted data") } paddingLen := int(decrypted[len(decrypted)-1]) if paddingLen > aes.BlockSize || paddingLen == 0 || paddingLen > len(decrypted) { return nil, fmt.Errorf("invalid PKCS7 padding: %d", paddingLen) } decrypted = decrypted[:len(decrypted)-paddingLen] // Some devices add a trailing null byte decrypted = trimNull(decrypted) var msg map[string]interface{} if err := json.Unmarshal(decrypted, &msg); err != nil { return nil, fmt.Errorf("unmarshal decrypted: %w", err) } return msg, nil } // Subscribe sends a miIO.sub request to the device. func (d *LanDevice) Subscribe() { if d.subLocked { return } d.subLocked = true defer func() { d.subLocked = false }() subTS := int(time.Now().Unix()) d.manager.sendToDevice(d.DID, map[string]interface{}{ "method": "miIO.sub", "params": map[string]interface{}{ "version": "2.0", "did": d.manager.VirtualDID(), "update_ts": subTS, "sub_method": ".", }, }, func(msg map[string]interface{}, ctx interface{}) { d.subscribeHandler(msg, subTS) }, nil, 5000) } func (d *LanDevice) subscribeHandler(msg map[string]interface{}, subTS int) { result, _ := msg["result"].(map[string]interface{}) if result == nil { d.manager.lgr.Infof("[lan] subscribe error: no result, did=%s, msg=%v", d.DID, msg) return } code, _ := result["code"].(float64) if code != 0 { d.manager.lgr.Infof("[lan] subscribe error: code=%v, did=%s, msg=%v", code, d.DID, msg) return } d.subscribed = true d.subTS = subTS d.manager.broadcastDeviceState(d.DID, map[string]interface{}{ "online": d.online, "push_available": d.subscribed, }) d.manager.lgr.Infof("[lan] subscribe success: if=%s, did=%s", d.IfName, d.DID) } // Unsubscribe sends a miIO.unsub request to the device. func (d *LanDevice) Unsubscribe() { if !d.subscribed { return } subTS := d.subTS d.manager.sendToDevice(d.DID, map[string]interface{}{ "method": "miIO.unsub", "params": map[string]interface{}{ "version": "2.0", "did": d.manager.VirtualDID(), "update_ts": subTS, "sub_method": ".", }, }, func(msg map[string]interface{}, ctx interface{}) { result, _ := msg["result"].(map[string]interface{}) if result != nil { if code, _ := result["code"].(float64); code != 0 { d.manager.lgr.Infof("[lan] unsubscribe error: code=%v, did=%s, msg=%v", code, d.DID, msg) return } } d.manager.lgr.Infof("[lan] unsubscribe success: if=%s, did=%s", d.IfName, d.DID) }, nil, 5000) d.subscribed = false d.manager.broadcastDeviceState(d.DID, map[string]interface{}{ "online": d.online, "push_available": d.subscribed, }) } // OnDelete cleans up timers when the device is removed. func (d *LanDevice) OnDelete() { if d.kaTimer != nil { d.kaTimer.Stop() d.kaTimer = nil } if d.unstableTimer != nil { d.unstableTimer.Stop() d.unstableTimer = nil } d.manager.lgr.Infof("[lan] device deleted: did=%s", d.DID) } // UpdateInfo updates the device token if it has changed. func (d *LanDevice) UpdateInfo(info map[string]interface{}) { tokenStr, ok := info["token"].(string) if !ok || len(tokenStr) != 32 { return } newToken, err := hexDecode(tokenStr) if err != nil { return } if strings.EqualFold(fmt.Sprintf("%x", d.Token), tokenStr) { return } // Token changed — recreate cipher d.Token = newToken aesKey := md5Sum(d.Token) aesIV := md5Sum(append(aesKey[:], d.Token...)) block, err := aes.NewCipher(aesKey[:]) if err != nil { d.manager.lgr.Infof("[lan] update token cipher error: %v, did=%s", err, d.DID) return } d.cipher = block d.aesIV = aesIV[:] d.manager.lgr.Infof("[lan] token updated: did=%s", d.DID) } // ============================================================================ // Keep-alive state machine // ============================================================================ func (d *LanDevice) updateKeepAlive(state LanDeviceState) { lastState := d.state d.state = state if d.state != LanDeviceFresh { d.manager.lgr.Infof("[lan] device status: did=%s, state=%s", d.DID, d.state) } // Cancel existing timer if d.kaTimer != nil { d.kaTimer.Stop() d.kaTimer = nil } switch state { case LanDeviceFresh: if lastState == LanDeviceDead { d.kaInterval = kaIntervalMin d.changeOnline(true) } nextTimeout := d.nextKATimeout() d.kaTimer = time.AfterFunc(time.Duration(nextTimeout*float64(time.Second)), func() { d.updateKeepAlive(LanDevicePing1) }) case LanDevicePing1, LanDevicePing2, LanDevicePing3: // Set timer first (matching Python's approach) nextState := LanDeviceState(int(state) + 1) d.kaTimer = time.AfterFunc(time.Duration(fastPingInterval*float64(time.Second)), func() { d.updateKeepAlive(nextState) }) // Fast ping if d.IfName == "" { d.manager.lgr.Infof("[lan] if_name not set for device: did=%s", d.DID) return } if d.IP == "" { d.manager.lgr.Infof("[lan] ip not set for device: did=%s", d.DID) return } d.manager.ping(d.IfName, d.IP) case LanDeviceDead: if lastState == LanDevicePing3 { d.kaInterval = kaIntervalMin d.changeOnline(false) } } } func (d *LanDevice) nextKATimeout() float64 { d.kaInterval = minFloat(d.kaInterval*2, kaIntervalMax) return randomizeFloat(d.kaInterval, 0.1) } func (d *LanDevice) changeOnline(online bool) { d.manager.lgr.Infof("[lan] change online: did=%s, online=%v", d.DID, online) tsNow := int(time.Now().Unix()) d.onlineHistory = append(d.onlineHistory, onlineEvent{ts: tsNow, online: online}) if len(d.onlineHistory) > networkUnstableCntTh { d.onlineHistory = d.onlineHistory[1:] } if d.unstableTimer != nil { d.unstableTimer.Stop() d.unstableTimer = nil } if !online { d.setOnline(false) } else { if len(d.onlineHistory) < networkUnstableCntTh || float64(tsNow-d.onlineHistory[0].ts) > networkUnstableTimeTh { d.setOnline(true) } else { d.manager.lgr.Infof("[lan] unstable device detected: did=%s", d.DID) d.unstableTimer = time.AfterFunc(time.Duration(networkUnstableResumeTh*float64(time.Second)), func() { d.manager.lgr.Infof("[lan] unstable resume threshold passed: did=%s", d.DID) d.setOnline(true) }) } } } // ============================================================================ // MIoTLan — LAN device control manager // ============================================================================ // MIoTLan manages MIoT LAN device discovery, communication, and subscriptions. type MIoTLan struct { mu sync.Mutex netIFs map[string]bool network *MIoTNetwork mipsService *MipsService enableSub bool virtualDID string // UDP sockets per interface socks map[string]*net.UDPConn localPort int // Devices devices map[string]*LanDevice // Message ID sequencing msgIDSeq int32 pendingReqs map[int]*pendingRequest // Subscriptions deviceStateSubs map[string]*deviceStateCallback matcher *MIoTMatcher // Duplicate message filter replyBuf map[string]*time.Timer // Control stopCh chan struct{} cmdCh chan func() initDone bool initMu sync.Mutex // Network state availableIFs map[string]bool scanTimer *time.Timer lastScanInterval float64 // LAN state votes lanStateSubs map[string]func(bool) lanCtrlVotes map[string]bool // Profile models (loaded from YAML, models that should NOT use LAN control) profileModels map[string]bool lgr logger.Logger } type pendingRequest struct { msgID int handler func(map[string]interface{}, interface{}) ctx interface{} timer *time.Timer } type deviceStateCallback struct { key string handler func(did string, state map[string]interface{}, ctx interface{}) ctx interface{} } // NewMIoTLan creates a new MIoTLan manager. // netIFs: allowed network interface names (matched against MIoTNetwork). // network: the MIoTNetwork monitor for interface change notifications. // mipsService: the MipsService for central hub gateway detection. func NewMIoTLan(netIFs []string, network *MIoTNetwork, mipsService *MipsService, enableSub bool, virtualDID int64) (*MIoTLan, error) { if network == nil { return nil, fmt.Errorf("network is required") } if mipsService == nil { return nil, fmt.Errorf("mipsService is required") } vDID := fmt.Sprintf("%d", virtualDID) if virtualDID == 0 { n, err := rand.Int(rand.Reader, big.NewInt(1<<62)) if err != nil { return nil, fmt.Errorf("generate virtual did: %w", err) } vDID = fmt.Sprintf("%d", n.Uint64()) } // Build probe message probe := make([]byte, otProbeLen) copy(probe[0:20], []byte("!1\x00\x20\xFF\xFF\xFF\xFF\xFF\xFF\xFF\xFF\xFF\xFF\xFF\xFFMDID")) binary.BigEndian.PutUint64(probe[20:28], mustParseUint64(vDID)) // bytes 28:32 are zeros m := &MIoTLan{ netIFs: make(map[string]bool), network: network, mipsService: mipsService, enableSub: enableSub, virtualDID: vDID, socks: make(map[string]*net.UDPConn), devices: make(map[string]*LanDevice), msgIDSeq: genMsgIDSeed(), pendingReqs: make(map[int]*pendingRequest), deviceStateSubs: make(map[string]*deviceStateCallback), matcher: NewMIoTMatcher(), replyBuf: make(map[string]*time.Timer), stopCh: make(chan struct{}), cmdCh: make(chan func(), 64), availableIFs: make(map[string]bool), lanStateSubs: make(map[string]func(bool)), lanCtrlVotes: make(map[string]bool), profileModels: make(map[string]bool), lgr: logger.Default(), } for _, ifName := range netIFs { m.netIFs[ifName] = true } // Subscribe to network interface changes network.SubNetworkInfo("miot_lan", m.onNetworkInfoChange) // Subscribe to mDNS service changes mipsService.SubServiceChange("miot_lan", "*", m.onMipsServiceChange) // Auto-init if no central hub gateway exists and we have interfaces if len(mipsService.GetServices("")) == 0 && len(netIFs) > 0 { m.lgr.Infof("[lan] no central hub gateway service, scheduling init") go func() { m.Init() }() } return m, nil } // SetLogger sets a custom logger for MIoTLan. func (m *MIoTLan) SetLogger(l logger.Logger) { m.lgr = l } // VirtualDID returns the virtual device ID used for LAN communication. func (m *MIoTLan) VirtualDID() string { return m.virtualDID } // InitDone returns whether the LAN controller is initialized. func (m *MIoTLan) InitDone() bool { m.mu.Lock() defer m.mu.Unlock() return m.initDone } // ============================================================================ // Init / Deinit // ============================================================================ // Init starts the LAN controller. Safe to call multiple times. func (m *MIoTLan) Init() { m.initMu.Lock() defer m.initMu.Unlock() m.mu.Lock() if m.initDone { m.mu.Unlock() m.lgr.Infof("[lan] already initialized") return } if len(m.netIFs) == 0 { m.mu.Unlock() m.lgr.Infof("[lan] no net_ifs configured") return } // Check if anyone voted for LAN control hasVote := false for _, v := range m.lanCtrlVotes { if v { hasVote = true break } } if !hasVote { m.mu.Unlock() m.lgr.Infof("[lan] no vote for lan ctrl") return } if len(m.mipsService.GetServices("")) > 0 { m.mu.Unlock() m.lgr.Infof("[lan] central hub gateway service exists, skip LAN init") return } // Build available interface set for ifName := range m.network.GetNetworkInfo() { m.availableIFs[ifName] = true } if len(m.availableIFs) == 0 { m.mu.Unlock() m.lgr.Infof("[lan] no available network interfaces") return } // Check that configured netIFs intersect with available IFs hasValid := false for ifName := range m.netIFs { if m.availableIFs[ifName] { hasValid = true break } } if !hasValid { m.mu.Unlock() m.lgr.Infof("[lan] no valid net_ifs matching configured list") return } m.mu.Unlock() // Start internal goroutine go m.internalLoop() m.mu.Lock() m.initDone = true m.mu.Unlock() // Notify LAN state subscribers m.mu.Lock() handlers := make([]func(bool), 0, len(m.lanStateSubs)) for _, h := range m.lanStateSubs { handlers = append(handlers, h) } m.mu.Unlock() for _, h := range handlers { go h(true) } m.lgr.Infof("[lan] initialized: netIFs=%v, available=%v", m.netIFs, m.availableIFs) } // Deinit stops the LAN controller. func (m *MIoTLan) Deinit() { m.initMu.Lock() defer m.initMu.Unlock() m.mu.Lock() if !m.initDone { m.mu.Unlock() m.lgr.Infof("[lan] not initialized") return } m.initDone = false m.mu.Unlock() // Signal the internal goroutine to stop m.cmdCh <- nil // sentinel: stop m.mu.Lock() // Clean up m.profileModels = make(map[string]bool) m.devices = make(map[string]*LanDevice) m.socks = make(map[string]*net.UDPConn) m.localPort = 0 m.scanTimer = nil m.lastScanInterval = 0 m.msgIDSeq = genMsgIDSeed() m.pendingReqs = make(map[int]*pendingRequest) m.matcher = NewMIoTMatcher() m.deviceStateSubs = make(map[string]*deviceStateCallback) m.replyBuf = make(map[string]*time.Timer) handlers := make([]func(bool), 0, len(m.lanStateSubs)) for _, h := range m.lanStateSubs { handlers = append(handlers, h) } m.mu.Unlock() for _, h := range handlers { go h(false) } m.lgr.Infof("[lan] deinitialized") } // ============================================================================ // Internal loop (runs in a dedicated goroutine) // ============================================================================ func (m *MIoTLan) internalLoop() { m.lgr.Infof("[lan] internal loop started") m.initSockets() // Start scan timer scanDelay := randomizeFloat(3.0, 1.0) // 0-6s random delay m.scanTimer = time.AfterFunc(time.Duration(scanDelay*float64(time.Second)), m.scanDevices) for { select { case cmd := <-m.cmdCh: if cmd == nil { // Stop signal m.deinitInternal() m.lgr.Infof("[lan] internal loop stopped") return } cmd() } } } func (m *MIoTLan) deinitInternal() { if m.scanTimer != nil { m.scanTimer.Stop() m.scanTimer = nil } for _, dev := range m.devices { dev.OnDelete() } m.devices = make(map[string]*LanDevice) for _, req := range m.pendingReqs { if req.timer != nil { req.timer.Stop() } } m.pendingReqs = make(map[int]*pendingRequest) for _, t := range m.replyBuf { t.Stop() } m.replyBuf = make(map[string]*time.Timer) m.deinitSockets() } // ============================================================================ // Socket management // ============================================================================ func (m *MIoTLan) initSockets() { m.mu.Lock() defer m.mu.Unlock() m.deinitSocketsLocked() for ifName := range m.netIFs { if !m.availableIFs[ifName] { continue } m.createSocketLocked(ifName) } } func (m *MIoTLan) deinitSockets() { m.mu.Lock() defer m.mu.Unlock() m.deinitSocketsLocked() } func (m *MIoTLan) deinitSocketsLocked() { for ifName := range m.socks { m.destroySocketLocked(ifName) } m.socks = make(map[string]*net.UDPConn) } func (m *MIoTLan) createSocketLocked(ifName string) { if _, ok := m.socks[ifName]; ok { return } // Get the interface IP from MIoTNetwork ni := m.network.GetNetworkInfo() info, ok := ni[ifName] if !ok || info.IP == "" { m.lgr.Infof("[lan] no IP for interface: %s", ifName) return } conn, err := createBoundSocket(info.IP, m.localPort) if err != nil { m.lgr.Infof("[lan] create socket error: if=%s, err=%v", ifName, err) return } m.socks[ifName] = conn if m.localPort == 0 { addr := conn.LocalAddr().(*net.UDPAddr) m.localPort = addr.Port } // Start reader goroutine for this socket go m.socketReader(ifName, conn) m.lgr.Infof("[lan] socket created: if=%s, ip=%s, port=%d", ifName, info.IP, m.localPort) } func (m *MIoTLan) destroySocketLocked(ifName string) { conn, ok := m.socks[ifName] if !ok { return } delete(m.socks, ifName) conn.Close() m.lgr.Infof("[lan] socket destroyed: if=%s", ifName) } // createBoundSocket creates a UDP socket bound to a specific IP address. func createBoundSocket(ip string, port int) (*net.UDPConn, error) { addr := &net.UDPAddr{ IP: net.ParseIP(ip), Port: port, } return net.ListenUDP("udp4", addr) } // socketReader reads from a UDP socket and dispatches raw messages. func (m *MIoTLan) socketReader(ifName string, conn *net.UDPConn) { buf := make([]byte, otMsgLen) for { n, remoteAddr, err := conn.ReadFromUDP(buf) if err != nil { // Socket closed or error return } if n <= 0 { continue } if remoteAddr.Port != otPort { continue } m.handleRawMessage(buf[:n], n, remoteAddr.IP.String(), ifName) } } // ============================================================================ // Sending // ============================================================================ // send broadcasts a probe/data to all or a specific interface. func (m *MIoTLan) send(ifName, targetIP string, data []byte) { m.mu.Lock() defer m.mu.Unlock() if ifName == "" { // Broadcast to all sockets for ifn, conn := range m.socks { conn.WriteToUDP(data, &net.UDPAddr{IP: net.ParseIP(targetIP), Port: otPort}) _ = ifn } } else { conn, ok := m.socks[ifName] if !ok { m.lgr.Infof("[lan] invalid socket: if=%s", ifName) return } conn.WriteToUDP(data, &net.UDPAddr{IP: net.ParseIP(targetIP), Port: otPort}) } } func (m *MIoTLan) ping(ifName, targetIP string) { if targetIP == "" { return } m.send(ifName, targetIP, m.buildProbe()) } func (m *MIoTLan) buildProbe() []byte { probe := make([]byte, otProbeLen) binary.BigEndian.PutUint16(probe[0:2], otHeader) binary.BigEndian.PutUint16(probe[2:4], otProbeLen) binary.BigEndian.PutUint64(probe[4:12], mustParseUint64(m.virtualDID)) // timestamp at 12:16 binary.BigEndian.PutUint32(probe[12:16], uint32(time.Now().Unix())) // bytes 16:32 remain zero (will be filled by receiver md5) return probe } // sendToDevice sends an encrypted message to a specific device. func (m *MIoTLan) sendToDevice(did string, msg map[string]interface{}, handler func(map[string]interface{}, interface{}), ctx interface{}, timeoutMs int) { m.mu.Lock() dev, ok := m.devices[did] m.mu.Unlock() if !ok { if handler != nil { go handler(map[string]interface{}{ "code": CodeInternalError, "error": "device not found", }, ctx) } return } if dev.cipher == nil { if handler != nil { go handler(map[string]interface{}{ "code": CodeInternalError, "error": "device cipher not initialized", }, ctx) } return } if dev.IfName == "" || dev.IP == "" { if handler != nil { go handler(map[string]interface{}{ "code": CodeInternalError, "error": "device network info missing", }, ctx) } return } msgID := m.genMsgID() msg["id"] = msgID out := make([]byte, otMsgLen) dataLen, err := dev.GenPacket(out, msg) if err != nil { if handler != nil { go handler(map[string]interface{}{ "code": CodeInternalError, "error": err.Error(), }, ctx) } return } m.makeRequest(msgID, out[:dataLen], dev.IfName, dev.IP, handler, ctx, timeoutMs) } func (m *MIoTLan) makeRequest(msgID int, data []byte, ifName, ip string, handler func(map[string]interface{}, interface{}), ctx interface{}, timeoutMs int) { req := &pendingRequest{ msgID: msgID, handler: handler, ctx: ctx, } if timeoutMs > 0 && handler != nil { req.timer = time.AfterFunc(time.Duration(timeoutMs)*time.Millisecond, func() { m.mu.Lock() if _, ok := m.pendingReqs[msgID]; ok { delete(m.pendingReqs, msgID) m.mu.Unlock() handler(map[string]interface{}{ "code": CodeTimeout, "error": "timeout", }, ctx) } else { m.mu.Unlock() } }) } m.mu.Lock() m.pendingReqs[msgID] = req m.mu.Unlock() m.send(ifName, ip, data) } // ============================================================================ // Message ID generation // ============================================================================ func (m *MIoTLan) genMsgID() int { for { id := int(atomic.AddInt32(&m.msgIDSeq, 1)) if id > 0x80000000 { atomic.StoreInt32(&m.msgIDSeq, 1) continue } return id } } func genMsgIDSeed() int32 { n, _ := rand.Int(rand.Reader, big.NewInt(0x7FFFFFFF)) return int32(n.Int64()) } // ============================================================================ // Raw message handling // ============================================================================ func (m *MIoTLan) handleRawMessage(data []byte, dataLen int, ip, ifName string) { if dataLen < 4 || binary.BigEndian.Uint16(data[0:2]) != otHeader { return } did := fmt.Sprintf("%d", binary.BigEndian.Uint64(data[4:12])) m.mu.Lock() dev, ok := m.devices[did] m.mu.Unlock() if !ok { return } // Update clock offset if dataLen >= 16 { timestamp := binary.BigEndian.Uint32(data[12:16]) dev.offset = int(time.Now().Unix()) - int(timestamp) } // Keep alive on probe or if subscribed if dataLen == otProbeLen || dev.subscribed { dev.KeepAlive(ip, ifName) } // Check for subscription handshake in probe if m.enableSub && dataLen == otProbeLen && dataLen >= 29 { if string(data[16:20]) == "MSUB" && string(data[24:27]) == "PUB" { dev.supportedWildcardSub = data[28] == otSupportWildcardSub subTS := int(binary.BigEndian.Uint32(data[20:24])) subType := data[27] if dev.supportedWildcardSub && (subType == 0 || subType == 1 || subType == 4) && subTS != dev.subTS { dev.subscribed = false dev.Subscribe() } } } // Decrypt and handle device messages if dataLen > otProbeLen { msg, err := dev.DecryptPacket(data[:dataLen]) if err != nil { m.lgr.Infof("[lan] decrypt error: did=%s, err=%v", did, err) return } m.handleMessage(did, msg) } } func (m *MIoTLan) handleMessage(did string, msg map[string]interface{}) { msgID, ok := msg["id"] if !ok { m.lgr.Infof("[lan] message without id: did=%s, msg=%v", did, msg) return } id, ok := toInt(msgID) if !ok { return } // Check if this is a reply to a pending request m.mu.Lock() req, ok := m.pendingReqs[id] if ok { delete(m.pendingReqs, id) } m.mu.Unlock() if ok { if req.timer != nil { req.timer.Stop() } if req.handler != nil { go req.handler(msg, req.ctx) } return } // Uplink message — check for method method, ok := msg["method"].(string) if !ok { return } // Filter duplicate messages if m.filterDup(did, id) { m.sendToDevice(did, map[string]interface{}{ "id": id, "result": map[string]interface{}{"code": 0}, }, nil, nil, 0) return } m.lgr.Infof("[lan] message: did=%s, method=%s", did, method) switch method { case "properties_changed": params, _ := msg["params"].([]interface{}) for _, p := range params { param, ok := p.(map[string]interface{}) if !ok { continue } siid, _ := param["siid"] piid, _ := param["piid"] if siid == nil || piid == nil { continue } key := fmt.Sprintf("%s/p/%v/%v", did, siid, piid) for _, entry := range m.matcher.Match(key) { if entry.Handler != nil { go entry.Handler(param, entry.Ctx) } } } case "event_occured": params, _ := msg["params"].(map[string]interface{}) if params != nil { siid, _ := params["siid"] eiid, _ := params["eiid"] if siid != nil && eiid != nil { key := fmt.Sprintf("%s/e/%v/%v", did, siid, eiid) for _, entry := range m.matcher.Match(key) { if entry.Handler != nil { go entry.Handler(params, entry.Ctx) } } } } } // Acknowledge m.sendToDevice(did, map[string]interface{}{ "id": id, "result": map[string]interface{}{"code": 0}, }, nil, nil, 0) } func (m *MIoTLan) filterDup(did string, msgID int) bool { m.mu.Lock() defer m.mu.Unlock() filterID := fmt.Sprintf("%s.%d", did, msgID) if _, ok := m.replyBuf[filterID]; ok { return true } m.replyBuf[filterID] = time.AfterFunc(5*time.Second, func() { m.mu.Lock() delete(m.replyBuf, filterID) m.mu.Unlock() }) return false } // ============================================================================ // Device management (internal loop commands) // ============================================================================ // broadcastDeviceState sends device state to all subscribers. func (m *MIoTLan) broadcastDeviceState(did string, state map[string]interface{}) { m.mu.Lock() callbacks := make([]*deviceStateCallback, 0, len(m.deviceStateSubs)) for _, cb := range m.deviceStateSubs { callbacks = append(callbacks, cb) } m.mu.Unlock() for _, cb := range callbacks { if cb.handler != nil { go cb.handler(did, state, cb.ctx) } } } // ============================================================================ // Scan devices // ============================================================================ func (m *MIoTLan) scanDevices() { if m.scanTimer != nil { m.scanTimer.Stop() m.scanTimer = nil } m.send("", "255.255.255.255", m.buildProbe()) scanTime := m.nextScanTime() m.scanTimer = time.AfterFunc(time.Duration(scanTime*float64(time.Second)), m.scanDevices) m.lgr.Infof("[lan] next scan in %.0fs", scanTime) } func (m *MIoTLan) nextScanTime() float64 { if m.lastScanInterval == 0 { m.lastScanInterval = otProbeIntervalMin } m.lastScanInterval = minFloat(m.lastScanInterval*2, otProbeIntervalMax) return m.lastScanInterval } // ============================================================================ // Network interface change handler // ============================================================================ func (m *MIoTLan) onNetworkInfoChange(status InterfaceStatus, info *NetworkInfo) { m.lgr.Infof("[lan] network info change: status=%s, if=%s", status, info.Name) m.mu.Lock() // Rebuild available IFs m.availableIFs = make(map[string]bool) for ifName := range m.network.GetNetworkInfo() { m.availableIFs[ifName] = true } if len(m.availableIFs) == 0 { m.mu.Unlock() go m.Deinit() return } valid := false for ifName := range m.netIFs { if m.availableIFs[ifName] { valid = true break } } if !valid { m.mu.Unlock() m.lgr.Infof("[lan] no valid net_ifs after change") go m.Deinit() return } if !m.initDone { m.mu.Unlock() go m.Init() return } m.mu.Unlock() // Send command to internal loop m.cmdCh <- func() { switch status { case InterfaceAdd: m.mu.Lock() m.availableIFs[info.Name] = true if m.netIFs[info.Name] { m.createSocketLocked(info.Name) } m.mu.Unlock() case InterfaceRemove: m.mu.Lock() delete(m.availableIFs, info.Name) m.destroySocketLocked(info.Name) m.mu.Unlock() } } } // ============================================================================ // mDNS service change handler // ============================================================================ func (m *MIoTLan) onMipsServiceChange(groupID string, state MipsServiceState, data map[string]interface{}) { m.lgr.Infof("[lan] mips service change: groupID=%s, state=%s, data=%v", groupID, state, data) if len(m.mipsService.GetServices("")) > 0 { m.lgr.Infof("[lan] central hub gateway found, deinit LAN") go m.Deinit() } else { m.lgr.Infof("[lan] no central hub gateway, init LAN") go m.Init() } } // ============================================================================ // Public API // ============================================================================ // VoteForLanCtrl votes for or against LAN control. func (m *MIoTLan) VoteForLanCtrl(key string, vote bool) { m.lgr.Infof("[lan] vote for lan ctrl: key=%s, vote=%v", key, vote) m.mu.Lock() m.lanCtrlVotes[key] = vote hasVote := false for _, v := range m.lanCtrlVotes { if v { hasVote = true break } } m.mu.Unlock() if !hasVote { go m.Deinit() } else { go m.Init() } } // UpdateSubscribeOption enables or disables device subscriptions. func (m *MIoTLan) UpdateSubscribeOption(enable bool) { m.lgr.Infof("[lan] update subscribe option: %v", enable) m.mu.Lock() oldEnable := m.enableSub m.enableSub = enable m.mu.Unlock() if oldEnable && !enable { // Unsubscribe all devices m.cmdCh <- func() { m.mu.Lock() defer m.mu.Unlock() for _, dev := range m.devices { dev.Unsubscribe() } } } } // UpdateDevices registers or updates devices. func (m *MIoTLan) UpdateDevices(devices map[string]map[string]interface{}) { if !m.InitDone() { return } m.cmdCh <- func() { m.mu.Lock() defer m.mu.Unlock() for did, info := range devices { if !isNumeric(did) { m.lgr.Infof("[lan] invalid did (non-numeric): %s", did) continue } model, _ := info["model"].(string) if model != "" && m.profileModels[model] { m.lgr.Infof("[lan] model not supported for LAN ctrl: did=%s, model=%s", did, model) continue } if existing, ok := m.devices[did]; ok { existing.UpdateInfo(info) } else { tokenStr, _ := info["token"].(string) if len(tokenStr) != 32 { m.lgr.Infof("[lan] invalid token for device: did=%s", did) continue } ip, _ := info["ip"].(string) dev, err := newLanDevice(m, did, tokenStr, ip) if err != nil { m.lgr.Infof("[lan] create device error: did=%s, err=%v", did, err) continue } m.devices[did] = dev } } } } // DeleteDevices removes devices from LAN management. func (m *MIoTLan) DeleteDevices(dids []string) { if !m.InitDone() { return } m.cmdCh <- func() { m.mu.Lock() defer m.mu.Unlock() for _, did := range dids { if dev, ok := m.devices[did]; ok { dev.OnDelete() delete(m.devices, did) } } } } // SubLanState subscribes to LAN controller state changes (init/deinit). func (m *MIoTLan) SubLanState(key string, handler func(bool)) { m.mu.Lock() defer m.mu.Unlock() m.lanStateSubs[key] = handler } // UnsubLanState unsubscribes from LAN controller state changes. func (m *MIoTLan) UnsubLanState(key string) { m.mu.Lock() defer m.mu.Unlock() delete(m.lanStateSubs, key) } // SubDeviceState subscribes to device online/offline state changes. func (m *MIoTLan) SubDeviceState(key string, handler func(did string, state map[string]interface{}, ctx interface{}), ctx interface{}) { m.cmdCh <- func() { m.mu.Lock() defer m.mu.Unlock() m.deviceStateSubs[key] = &deviceStateCallback{key: key, handler: handler, ctx: ctx} } } // UnsubDeviceState unsubscribes from device state changes. func (m *MIoTLan) UnsubDeviceState(key string) { m.cmdCh <- func() { m.mu.Lock() defer m.mu.Unlock() delete(m.deviceStateSubs, key) } } // SubProp subscribes to property change broadcasts for a device. func (m *MIoTLan) SubProp(did string, handler func(map[string]interface{}, interface{}), siid, piid int, ctx interface{}) { key := fmt.Sprintf("%s/p/#", did) if siid > 0 && piid > 0 { key = fmt.Sprintf("%s/p/%d/%d", did, siid, piid) } m.cmdCh <- func() { m.matcher.Sub(key, handler, ctx) m.lgr.Infof("[lan] sub prop: key=%s", key) } } // UnsubProp unsubscribes from property change broadcasts. func (m *MIoTLan) UnsubProp(did string, siid, piid int) { key := fmt.Sprintf("%s/p/#", did) if siid > 0 && piid > 0 { key = fmt.Sprintf("%s/p/%d/%d", did, siid, piid) } m.cmdCh <- func() { m.matcher.Unsub(key) m.lgr.Infof("[lan] unsub prop: key=%s", key) } } // SubEvent subscribes to event broadcasts for a device. func (m *MIoTLan) SubEvent(did string, handler func(map[string]interface{}, interface{}), siid, eiid int, ctx interface{}) { key := fmt.Sprintf("%s/e/#", did) if siid > 0 && eiid > 0 { key = fmt.Sprintf("%s/e/%d/%d", did, siid, eiid) } m.cmdCh <- func() { m.matcher.Sub(key, handler, ctx) m.lgr.Infof("[lan] sub event: key=%s", key) } } // UnsubEvent unsubscribes from event broadcasts. func (m *MIoTLan) UnsubEvent(did string, siid, eiid int) { key := fmt.Sprintf("%s/e/#", did) if siid > 0 && eiid > 0 { key = fmt.Sprintf("%s/e/%d/%d", did, siid, eiid) } m.cmdCh <- func() { m.matcher.Unsub(key) m.lgr.Infof("[lan] unsub event: key=%s", key) } } // GetDevList returns the list of online devices and their push availability. func (m *MIoTLan) GetDevList() map[string]map[string]interface{} { m.mu.Lock() defer m.mu.Unlock() result := make(map[string]map[string]interface{}) for _, dev := range m.devices { if dev.online { result[dev.DID] = map[string]interface{}{ "online": dev.online, "push_available": dev.subscribed, } } } return result } // GetProps retrieves device properties in batch via LAN. // Sends encrypted UDP requests to each device. func (m *MIoTLan) GetProps(params []PropParam) ([]PropResult, error) { results := make([]PropResult, 0, len(params)) resultCh := make(chan PropResult, len(params)) var wg sync.WaitGroup for _, p := range params { wg.Add(1) go func(param PropParam) { defer wg.Done() resultCh <- m.getPropSingle(param) }(p) } wg.Wait() close(resultCh) for r := range resultCh { results = append(results, r) } return results, nil } // getPropSingle fetches a single property from a LAN device. func (m *MIoTLan) getPropSingle(param PropParam) PropResult { m.mu.Lock() dev, ok := m.devices[param.DID] m.mu.Unlock() if !ok || dev.cipher == nil { return PropResult{ DID: param.DID, Siid: param.Siid, Piid: param.Piid, Code: CodeInternalError, } } msg := map[string]interface{}{ "method": "properties/get", "params": []interface{}{ map[string]interface{}{ "did": param.DID, "siid": param.Siid, "piid": param.Piid, }, }, } resultCh := make(chan PropResult, 1) m.sendToDevice(param.DID, msg, func(resp map[string]interface{}, ctx interface{}) { code, _ := resp["code"].(float64) r := PropResult{ DID: param.DID, Siid: param.Siid, Piid: param.Piid, Code: int(code), } if result, ok := resp["result"].(map[string]interface{}); ok { r.Value = result["value"] } select { case resultCh <- r: default: } }, nil, 5000) select { case r := <-resultCh: return r case <-time.After(6 * time.Second): return PropResult{ DID: param.DID, Siid: param.Siid, Piid: param.Piid, Code: CodeTimeout, } } } // ============================================================================ // Crypto helpers // ============================================================================ // md5Sum computes the MD5 hash of data and returns the 16-byte result. func md5Sum(data []byte) [16]byte { return md5.Sum(data) } // hexDecode decodes a hex string to bytes. func hexDecode(s string) ([]byte, error) { // Use fmt.Sscanf-style approach dst := make([]byte, len(s)/2) for i := 0; i < len(s); i += 2 { var b byte _, err := fmt.Sscanf(s[i:i+2], "%02x", &b) if err != nil { return nil, fmt.Errorf("invalid hex at position %d: %s", i, s[i:i+2]) } dst[i/2] = b } return dst, nil } // ============================================================================ // Math / utility helpers // ============================================================================ func minFloat(a, b float64) float64 { if a < b { return a } return b } func randomizeFloat(value, ratio float64) float64 { return value * (1 - ratio + randFloat()*2*ratio) } func randFloat() float64 { n, _ := rand.Int(rand.Reader, big.NewInt(1<<53)) return float64(n.Int64()) / float64(1<<53) } func mustParseUint64(s string) uint64 { var v uint64 fmt.Sscanf(s, "%d", &v) return v } func toInt(v interface{}) (int, bool) { switch n := v.(type) { case float64: return int(n), true case int: return n, true case int64: return int(n), true case json.Number: i, err := n.Int64() if err != nil { return 0, false } return int(i), true } return 0, false } func isNumeric(s string) bool { for _, c := range s { if c < '0' || c > '9' { return false } } return len(s) > 0 } func trimNull(b []byte) []byte { for len(b) > 0 && b[len(b)-1] == 0 { b = b[:len(b)-1] } return b }