package gosqs import ( "net/url" "time" "shared-sqs/app/models" log "github.com/sirupsen/logrus" ) func init() { models.SyncQueues.Queues = make(map[string]*models.Queue) } // PeriodicTasks — фоновая обработка очередей, выполняется раз в `d` (1 секунда). // // Задачи на каждый такт (для КАЖДОЙ очереди, см. processQueueTick): // 1. Очистка просроченных записей deduplication-окна (FIFO). // 2. Возврат сообщений в видимое состояние по истечении VisibilityTimeout. // 3. Перенос сообщений в DeadLetterQueue при достижении MaxReceiveCount. // // ВАЖНО (фикс производительности): раньше глобальный SyncQueues.Lock удерживался // на время обхода ВСЕХ очередей — при большом числе очередей HTTP-обработчики // блокировались на секунды (наблюдались ReadTimeout и latency до 54 секунд). // Теперь блокировка захватывается ПО ОДНОЙ очереди: между обработкой двух // очередей другие горутины могут вклиниться. func PeriodicTasks(d time.Duration, quit chan bool) { ticker := time.NewTicker(d) for { select { case <-ticker.C: // Снапшот указателей на очереди под коротким RLock. // Указатели остаются валидными после снятия RLock: очереди из map // не удаляются во время итерации, а объекты живут в куче. models.SyncQueues.RLock() queues := make([]*models.Queue, 0, len(models.SyncQueues.Queues)) for _, queue := range models.SyncQueues.Queues { queues = append(queues, queue) } models.SyncQueues.RUnlock() // Каждую очередь обрабатываем под ОТДЕЛЬНЫМ Lock — блокировка одной // очереди не мешает обработке остальных и HTTP-запросам. for _, queue := range queues { models.SyncQueues.Lock() processQueueTick(queue) models.SyncQueues.Unlock() } case <-quit: ticker.Stop() return } } } // processQueueTick — обработка ОДНОЙ очереди за один такт. // // ВЫЗЫВАТЬ ТОЛЬКО под удержанием SyncQueues.Lock. // // Логическая схема: // 1. Deduplication-окно (FIFO): удаляем записи старше DeduplicationPeriod. // 2. Проходим по сообщениям очереди: // - пустой ReceiptHandle (не in-flight) — пропускаем; // - VisibilityTimeout ещё не истёк — пропускаем; // - иначе: возвращаем видимость (сброс ReceiptHandle, разблокировка // FIFO-группы, обновление ReceiptTime), инкрементируем Retry; // - если Retry >= MaxReceiveCount и настроена DLQ — переносим сообщение // в DLQ и удаляем из исходной очереди. func processQueueTick(queue *models.Queue) { // 1. Просроченные записи deduplication-окна. for dedupId, startTime := range queue.Duplicates { if time.Now().After(startTime.Add(models.DeduplicationPeriod)) { log.Debugf("deduplication period for message with deduplicationId [%s] expired", dedupId) delete(queue.Duplicates, dedupId) } } // 2. Возврат видимости истёкших сообщений + перенос в DLQ. log.Debugf("Queue [%s] length [%d]", queue.Name, len(queue.Messages)) for i := 0; i < len(queue.Messages); i++ { msg := &queue.Messages[i] if msg.ReceiptHandle == "" { continue // сообщение не в полёте } if msg.VisibilityTimeout.After(time.Now()) { continue // таймаут видимости ещё действует } log.Debugf("Making message visible again %s", msg.ReceiptHandle) queue.UnlockGroup(msg.GroupID) msg.ReceiptHandle = "" msg.ReceiptTime = time.Now().UTC() msg.Retry++ // 3. Перенос в DLQ при превышении MaxReceiveCount. if queue.MaxReceiveCount > 0 && queue.DeadLetterQueue != nil && msg.Retry >= queue.MaxReceiveCount { queue.DeadLetterQueue.Messages = append(queue.DeadLetterQueue.Messages, *msg) queue.Messages = append(queue.Messages[:i], queue.Messages[i+1:]...) i-- // элемент на позиции i удалён — повторно проверяем новый элемент на ней } } } func numberOfHiddenMessagesInQueue(queue models.Queue) int { num := 0 for _, m := range queue.Messages { if m.ReceiptHandle != "" || m.DelaySecs > 0 && time.Now().Before(m.SentTime.Add(time.Duration(m.DelaySecs)*time.Second)) { num++ } } return num } func getQueueFromPath(formVal string, theUrl string) string { if formVal != "" { return formVal } u, err := url.Parse(theUrl) if err != nil { return "" } return u.Path }