package service import ( "context" "encoding/json" "lone-services/pkg/logRedis" "lone-services/pkg/utils" "lone-services/services/lonelog/internal/dao" "lone-services/services/lonelog/internal/model" "sync" "time" "github.com/zeromicro/go-zero/core/logx" ) // LogWatcher 队列监听器,内部直接使用 logredis.LogRedisClient 包全局变量,不再传入redis参数 type LogWatcher struct { ctx context.Context cancel context.CancelFunc wg sync.WaitGroup } // GlobalLogWatcher 全局单例,main初始化之后其它包直接 service.GlobalLogWatcher 访问 var GlobalLogWatcher *LogWatcher // NewLogWatcher 无入参,直接使用包全局 logredis.LogRedisClient func NewLogWatcher() *LogWatcher { ctx, cancel := context.WithCancel(context.Background()) return &LogWatcher{ ctx: ctx, cancel: cancel, } } // StartWatch 启动两个消费协程 func (w *LogWatcher) StartWatch() { w.wg.Add(2) go w.loopActionLog() go w.loopLoginLog() logx.Infof("LogWatcher start, watch %s , %s", utils.LogActionQueueKey, utils.LogLoginQueueKey) } // loopActionLog 消费action日志队列,带recover自动重启 func (w *LogWatcher) loopActionLog() { defer func() { if r := recover(); r != nil { logx.Errorf("loopActionLog panic recovered, err=%v", r) time.Sleep(2 * time.Second) w.wg.Add(1) go w.loopActionLog() } w.wg.Done() }() for { select { case <-w.ctx.Done(): logx.Info("loopActionLog received exit signal, quit") return default: } res, err := logRedis.LogRedisClient.BRPop(w.ctx, 5*time.Second, utils.LogActionQueueKey).Result() if err != nil { if err.Error() == "redis: nil" { continue } logx.Errorf("BRPop %s redis err: %v", utils.LogActionQueueKey, err) time.Sleep(1 * time.Second) continue } if len(res) < 2 { continue } payload := res[1] var action utils.ActionLog if err := json.Unmarshal([]byte(payload), &action); err != nil { logx.Errorf("unmarshal actionlog failed payload=%s err=%v", payload, err) continue } ctxBg := context.Background() logTime, _ := utils.ParseCustomTime(action.CreateTime) edit := dao.ActionLog{ Content: action.Content, AdminId: action.AdminId, AdminName: action.AdminName, Type: uint8(action.Type), ModuleName: action.ModuleName, Ip: action.Ip, BrowserInfo: action.BrowserInfo, BrowserName: action.BrowserName, BrowserVersion: action.BrowserVersion, Reason: action.Reason, CreateTime: logTime, } err = model.ActionModel{}.Init().Create(&edit) if err != nil { logx.Errorf("insert actionlog mysql err=%v payload=%s", err, payload) _ = logRedis.LogRedisClient.LPush(ctxBg, utils.LogActionQueueKey, payload).Err() time.Sleep(500 * time.Millisecond) } } } // loopLoginLog 消费登录日志队列,带recover自动重启 func (w *LogWatcher) loopLoginLog() { defer func() { if r := recover(); r != nil { logx.Errorf("loopLoginLog panic recovered, err=%v", r) time.Sleep(2 * time.Second) w.wg.Add(1) go w.loopLoginLog() } w.wg.Done() }() for { select { case <-w.ctx.Done(): logx.Info("loopLoginLog received exit signal, quit") return default: } res, err := logRedis.LogRedisClient.BRPop(w.ctx, 5*time.Second, utils.LogLoginQueueKey).Result() if err != nil { if err.Error() == "redis: nil" { continue } logx.Errorf("BRPop %s redis err: %v", utils.LogLoginQueueKey, err) time.Sleep(1 * time.Second) continue } if len(res) < 2 { continue } payload := res[1] var login utils.LoginLog if err := json.Unmarshal([]byte(payload), &login); err != nil { logx.Errorf("unmarshal loginlog failed payload=%s err=%v", payload, err) continue } ctxBg := context.Background() logTime, _ := utils.ParseCustomTime(login.CreateTime) edit := dao.LoginLog{ AdminId: login.AdminId, AdminName: login.Name, Ip: login.Ip, BrowserInfo: login.BrowserInfo, BrowserName: login.BrowserName, BrowserVersion: login.BrowserVersion, Reason: login.Reason, CreateTime: logTime, Status: uint8(login.Status), } err = model.LoginModel{}.Init().Create(&edit) if err != nil { logx.Errorf("insert loginlog mysql err=%v payload=%s", err, payload) _ = logRedis.LogRedisClient.LPush(ctxBg, utils.LogLoginQueueKey, payload).Err() time.Sleep(500 * time.Millisecond) } } } // StopWatch 优雅停止消费协程 func (w *LogWatcher) StopWatch() { w.cancel() w.wg.Wait() logx.Info("LogWatcher stopped complete") }