128 lines
5.2 KiB
Go
128 lines
5.2 KiB
Go
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
|
|
}
|