package lc import ( "sync" "time" ) // CommandQueue 命令队列,用于串行处理Send操作(带150ms间隔) type CommandQueue struct { tasks chan []byte // 存储待发送的数据 running bool // 队列运行状态 mu sync.Mutex // 保护running状态的互斥锁 wg sync.WaitGroup // 用于等待队列结束 interval time.Duration // 任务执行间隔(毫秒) } // NewCommandQueue 创建新的命令队列(指定间隔,默认150ms) func NewCommandQueue(bufferSize int, interval ...time.Duration) *CommandQueue { // 默认间隔150毫秒,支持自定义传入 intervalMs := 200 * time.Millisecond if len(interval) > 0 { intervalMs = interval[0] } return &CommandQueue{ tasks: make(chan []byte, bufferSize), interval: intervalMs, } } // Start 启动队列处理循环 func (q *CommandQueue) Start(handle func([]byte)) { q.mu.Lock() defer q.mu.Unlock() if q.running { return } q.running = true q.wg.Add(1) go func() { defer q.wg.Done() // 循环处理队列中的发送任务 for data := range q.tasks { handle(data) // 执行实际发送操作 time.Sleep(q.interval) } }() } // AddTask 添加发送任务到队列 func (q *CommandQueue) AddTask(data []byte) { q.mu.Lock() defer q.mu.Unlock() if !q.running { return } q.tasks <- data } // Stop 停止队列,等待所有任务处理完毕 func (q *CommandQueue) Stop() { q.mu.Lock() defer q.mu.Unlock() if !q.running { return } close(q.tasks) // 关闭通道,让处理goroutine退出 q.wg.Wait() // 等待所有任务处理完毕 q.running = false }