| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382 |
- 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
- }
|