Files
4566704 2b7de69403 refactor(logger): 日志改用 logrus 全局包级调用,移除 Logger 注入接口
为什么:并入宿主项目后须与宿主日志体系一致——宿主统一用 pkg/logger 初始化
logrus 全局实例并直接 logrus.Xxx 包级调用;包一层 Logger 接口会让 logrus 的
caller 定位(报告调用文件:行号)失准。

改动:
- 删除 logger 包(Logger 接口 / Default / SetDefault,默认本就是 logrus.StandardLogger)
- 11 个结构体移除 lgr logger.Logger 字段、SetLogger 方法及构造中的 lgr 初始化
- 192 处 lgr.Xxxf 调用改为 logrus.Xxxf,日志消息文本保持不变
- xiaomi.Client 移除 Logger() 访问器,air_conditioner 回读日志去掉判空调用直连 logrus
- miot_i18n.go 两处标准库 log.Printf 改为 logrus.Errorf,统一日志出口
- go.mod:logrus 从 indirect 提升为直接依赖
- goimports/goformat 全树规范化:此前部分文件未 gofmt(单行 if、对齐),本次顺带
  格式化,纯空白/换行变化,无逻辑改动

验证:go build / go vet / go test ./... 全部通过(bridge、miot、xiaomi、devices、specs)
2026-09-15 20:16:26 +08:00

1707 lines
42 KiB
Go

// 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"
"github.com/sirupsen/logrus"
)
// ============================================================================
// 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)
})
logrus.Debugf("[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
logrus.Debugf("[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 {
logrus.Warnf("[lan] subscribe error: no result, did=%s, msg=%v", d.DID, msg)
return
}
code, _ := result["code"].(float64)
if code != 0 {
logrus.Warnf("[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,
})
logrus.Debugf("[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 {
logrus.Warnf("[lan] unsubscribe error: code=%v, did=%s, msg=%v", code, d.DID, msg)
return
}
}
logrus.Debugf("[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
}
logrus.Debugf("[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 {
logrus.Warnf("[lan] update token cipher error: %v, did=%s", err, d.DID)
return
}
d.cipher = block
d.aesIV = aesIV[:]
logrus.Debugf("[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 {
logrus.Debugf("[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 == "" {
logrus.Debugf("[lan] if_name not set for device: did=%s", d.DID)
return
}
if d.IP == "" {
logrus.Debugf("[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) {
logrus.Debugf("[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 {
logrus.Debugf("[lan] unstable device detected: did=%s", d.DID)
d.unstableTimer = time.AfterFunc(time.Duration(networkUnstableResumeTh*float64(time.Second)), func() {
logrus.Debugf("[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
}
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),
}
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 {
logrus.Infof("[lan] no central hub gateway service, scheduling init")
go func() {
m.Init()
}()
}
return m, nil
}
// 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()
logrus.Infof("[lan] already initialized")
return
}
if len(m.netIFs) == 0 {
m.mu.Unlock()
logrus.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()
logrus.Infof("[lan] no vote for lan ctrl")
return
}
if len(m.mipsService.GetServices("")) > 0 {
m.mu.Unlock()
logrus.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()
logrus.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()
logrus.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)
}
logrus.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()
logrus.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)
}
logrus.Infof("[lan] deinitialized")
}
// ============================================================================
// Internal loop (runs in a dedicated goroutine)
// ============================================================================
func (m *MIoTLan) internalLoop() {
logrus.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()
logrus.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 == "" {
logrus.Infof("[lan] no IP for interface: %s", ifName)
return
}
conn, err := createBoundSocket(info.IP, m.localPort)
if err != nil {
logrus.Warnf("[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)
logrus.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()
logrus.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 {
logrus.Debugf("[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 {
logrus.Debugf("[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 {
logrus.Debugf("[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
}
logrus.Debugf("[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)
logrus.Debugf("[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) {
logrus.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()
logrus.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{}) {
logrus.Infof("[lan] mips service change: groupID=%s, state=%s, data=%v", groupID, state, data)
if len(m.mipsService.GetServices("")) > 0 {
logrus.Infof("[lan] central hub gateway found, deinit LAN")
go m.Deinit()
} else {
logrus.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) {
logrus.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) {
logrus.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) {
logrus.Infof("[lan] invalid did (non-numeric): %s", did)
continue
}
model, _ := info["model"].(string)
if model != "" && m.profileModels[model] {
logrus.Debugf("[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 {
logrus.Debugf("[lan] invalid token for device: did=%s", did)
continue
}
ip, _ := info["ip"].(string)
dev, err := newLanDevice(m, did, tokenStr, ip)
if err != nil {
logrus.Warnf("[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)
logrus.Debugf("[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)
logrus.Debugf("[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)
logrus.Debugf("[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)
logrus.Debugf("[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
}