package main import ( "crypto/md5" "encoding/base64" "encoding/hex" "fmt" "io" "os" "path/filepath" "strings" "sync" "time" "lc/common/mqtt" "lc/common/protocol" "lc/common/util" ) const ( deployChunkTimeout = 10 * time.Minute deployTmpDir = "deploy" ) // DeployMgr 部署管理器 type DeployMgr struct { mu sync.Mutex active bool version string fileName string totalChunks int chunkSize int expectedMD5 string chunkBits []bool receivedAt time.Time deployTmp string cancel chan struct{} } var _deployMgrOnce sync.Once var _deployMgr *DeployMgr func GetDeployMgr() *DeployMgr { _deployMgrOnce.Do(func() { _deployMgr = &DeployMgr{} }) return _deployMgr } // HandleTpDeploy 处理来自 MQTT 的部署指令/分片 func HandleTpDeploy(m mqtt.Message) { var obj protocol.Pack_DeployCmd if err := obj.DeCode(string(m.Payload())); err != nil { util.GetTagLog().Errorf("sys", "HandleTpDeploy:DeCode失败,err=%v", err) return } mgr := GetDeployMgr() switch obj.Data.Action { case "start": mgr.start(&obj) case "chunk": mgr.receiveChunk(&obj) case "abort": mgr.abort(&obj) default: util.GetTagLog().Warnf("sys", "HandleTpDeploy:未知Action=%s", obj.Data.Action) } } func (o *DeployMgr) start(cmd *protocol.Pack_DeployCmd) { o.mu.Lock() defer o.mu.Unlock() if o.active { util.GetTagLog().Warnf("sys", "DeployMgr:已有部署进行中,拒绝新部署指令 version=%s", cmd.Data.Version) o.sendAck(cmd.Data.Version, false, "已有部署进行中") return } o.active = true o.version = cmd.Data.Version o.fileName = cmd.Data.FileName o.totalChunks = cmd.Data.TotalChunks o.chunkSize = cmd.Data.ChunkSize o.expectedMD5 = cmd.Data.MD5 o.chunkBits = make([]bool, cmd.Data.TotalChunks) o.receivedAt = time.Now() o.deployTmp = filepath.Join(util.GetPath(4), deployTmpDir) o.cancel = make(chan struct{}, 1) if err := os.MkdirAll(o.deployTmp, os.ModePerm); err != nil { util.GetTagLog().Errorf("sys", "DeployMgr:创建临时目录失败,path=%s,err=%v", o.deployTmp, err) o.cleanup("创建临时目录失败: " + err.Error()) return } util.GetTagLog().Infof("sys", "DeployMgr:开始部署 version=%s totalChunks=%d chunkSize=%d md5=%s", cmd.Data.Version, cmd.Data.TotalChunks, cmd.Data.ChunkSize, cmd.Data.MD5) // 启动超时监控 goroutine(独立于 gopool,避免 worker 饥饿问题) go o.watchTimeout() } func (o *DeployMgr) receiveChunk(cmd *protocol.Pack_DeployCmd) { o.mu.Lock() defer o.mu.Unlock() if !o.active { return } if cmd.Data.Version != o.version { util.GetTagLog().Warnf("sys", "DeployMgr:分片版本不匹配,expect=%s,got=%s", o.version, cmd.Data.Version) return } idx := cmd.Data.ChunkIndex if idx < 0 || idx >= o.totalChunks { util.GetTagLog().Errorf("sys", "DeployMgr:分片索引越界,idx=%d,total=%d", idx, o.totalChunks) return } // base64 decode the chunk data decoded, err := base64.StdEncoding.DecodeString(cmd.Data.Data) if err != nil { util.GetTagLog().Errorf("sys", "DeployMgr:base64解码分片%d失败,err=%v", idx, err) return } chunkFile := filepath.Join(o.deployTmp, fmt.Sprintf("chunk_%d", idx)) if err := os.WriteFile(chunkFile, decoded, os.ModePerm); err != nil { util.GetTagLog().Errorf("sys", "DeployMgr:写分片失败,idx=%d,err=%v", idx, err) return } o.chunkBits[idx] = true o.receivedAt = time.Now() allReceived := true for _, b := range o.chunkBits { if !b { allReceived = false break } } if allReceived { go o.assembleAndDeploy() } } func (o *DeployMgr) abort(cmd *protocol.Pack_DeployCmd) { o.mu.Lock() defer o.mu.Unlock() if !o.active || cmd.Data.Version != o.version { return } util.GetTagLog().Infof("sys", "DeployMgr:收到取消指令 version=%s", o.version) if o.cancel != nil { close(o.cancel) } o.cleanup("已取消") } func (o *DeployMgr) assembleAndDeploy() { o.mu.Lock() defer o.mu.Unlock() if !o.active { return } util.GetTagLog().Infof("sys", "DeployMgr:所有分片已收齐,开始重组 version=%s", o.version) // 1. 重组文件 assembledFile := filepath.Join(o.deployTmp, o.fileName) outFile, err := os.Create(assembledFile) if err != nil { o.cleanup("创建重组文件失败: " + err.Error()) return } defer outFile.Close() for i := 0; i < o.totalChunks; i++ { chunkFile := filepath.Join(o.deployTmp, fmt.Sprintf("chunk_%d", i)) data, err := os.ReadFile(chunkFile) if err != nil { outFile.Close() o.cleanup(fmt.Sprintf("读取分片%d失败: %s", i, err.Error())) return } if _, err := outFile.Write(data); err != nil { outFile.Close() o.cleanup(fmt.Sprintf("写入重组文件失败: %s", err.Error())) return } } // 2. MD5 校验 actualMD5, err := fileMD5(assembledFile) if err != nil { o.cleanup("计算MD5失败: " + err.Error()) return } if !strings.EqualFold(actualMD5, o.expectedMD5) { o.cleanup(fmt.Sprintf("MD5校验失败:期望%s 实际%s", o.expectedMD5, actualMD5)) return } util.GetTagLog().Infof("sys", "DeployMgr:MD5校验通过 md5=%s", actualMD5) // 3. ELF 魔数校验 if !isELF(assembledFile) { o.cleanup("文件格式错误:非Linux可执行文件") return } // 4. 磁盘空间检查(rename 不消耗额外空间,只需 1x 文件大小) cwd, _ := os.Getwd() targetPath := filepath.Join(cwd, o.fileName) fi, _ := os.Stat(assembledFile) needed := fi.Size() + 1*1024*1024 // 文件大小 + 1MB 安全余量 free, err := getFreeSpace(cwd) if err == nil && free < needed { o.cleanup(fmt.Sprintf("磁盘空间不足:需要%dMB 可用%dMB", needed/(1024*1024), free/(1024*1024))) return } // 5. 备份当前二进制,替换新文件 backupPath := targetPath + ".bak" os.Remove(backupPath) if _, err := os.Stat(targetPath); err == nil { if err := os.Rename(targetPath, backupPath); err != nil { o.cleanup("备份旧版失败: " + err.Error()) return } util.GetTagLog().Infof("sys", "DeployMgr:旧版已备份到 %s", backupPath) } if err := os.Rename(assembledFile, targetPath); err != nil { os.Rename(backupPath, targetPath) o.cleanup("替换文件失败: " + err.Error()) return } if err := os.Chmod(targetPath, 0755); err != nil { util.GetTagLog().Warnf("sys", "DeployMgr:chmod失败,err=%v", err) } util.GetTagLog().Infof("sys", "DeployMgr:文件替换成功 version=%s path=%s", o.version, targetPath) // 6. 写部署标记 marker := DeployMarker{ Version: o.version, Timestamp: time.Now().Unix(), Action: "deploy", } markerPath := filepath.Join(util.GetPath(0), "deploy_marker.json") markerContent, _ := json.MarshalToString(marker) if err := os.WriteFile(markerPath, []byte(markerContent), os.ModePerm); err != nil { util.GetTagLog().Errorf("sys", "DeployMgr:写部署标记失败,err=%v", err) } // 7. 发送 ACK o.active = false o.sendAck(o.version, true, "") util.GetTagLog().Infof("sys", "DeployMgr:即将退出进程以完成部署") // 8. 退出让 goforever 拉起新版本 time.Sleep(500 * time.Millisecond) os.Exit(0) } // watchTimeout 分片接收超时监控(独立 goroutine,不走 gopool 避免 worker 饥饿) func (o *DeployMgr) watchTimeout() { ticker := time.NewTicker(30 * time.Second) defer ticker.Stop() for { select { case <-o.cancel: return case <-ticker.C: o.mu.Lock() if !o.active { o.mu.Unlock() return } elapsed := time.Since(o.receivedAt) if elapsed > deployChunkTimeout { received := o.countReceived() o.mu.Unlock() o.cleanup(fmt.Sprintf("分片接收超时:已收 %d/%d", received, o.totalChunks)) return } o.mu.Unlock() } } } func (o *DeployMgr) cleanup(errMsg string) { util.GetTagLog().Errorf("sys", "DeployMgr:部署失败,原因=%s", errMsg) if o.deployTmp != "" { os.RemoveAll(o.deployTmp) } if o.active { o.sendAck(o.version, false, errMsg) } o.active = false o.version = "" o.chunkBits = nil o.cancel = nil } func (o *DeployMgr) sendAck(version string, success bool, errMsg string) { var ack protocol.Pack_DeployAck seq := GetNextUint64() if str, err := ack.EnCode(appConfig.GID, appConfig.GID, seq, version, success, errMsg); err == nil { topic := GetTopic(protocol.DT_GATEWAY, appConfig.GID, protocol.TP_GW_DEPLOY_ACK) GetMQTTMgr().Publish(topic, str, 0, ToCloud) util.GetTagLog().Infof("sys", "DeployMgr:发送ACK version=%s success=%v", version, success) } } func (o *DeployMgr) IsActive() bool { o.mu.Lock() defer o.mu.Unlock() return o.active } func (o *DeployMgr) countReceived() int { n := 0 for _, b := range o.chunkBits { if b { n++ } } return n } // DeployMarker 部署/回滚标记 type DeployMarker struct { Version string `json:"version"` Timestamp int64 `json:"timestamp"` Action string `json:"action"` } // ---- 工具函数 ---- func fileMD5(path string) (string, error) { f, err := os.Open(path) if err != nil { return "", err } defer f.Close() h := md5.New() if _, err := io.Copy(h, f); err != nil { return "", err } return hex.EncodeToString(h.Sum(nil)), nil } func isELF(path string) bool { f, err := os.Open(path) if err != nil { return false } defer f.Close() header := make([]byte, 4) if _, err := io.ReadFull(f, header); err != nil { return false } return header[0] == 0x7f && header[1] == 'E' && header[2] == 'L' && header[3] == 'F' } // init registers the deploy handler func init() { // Deploy handler will be registered in InitCloudMqttSubscribeTopics after appConfig is loaded }