149 lines
6.6 KiB
Go
149 lines
6.6 KiB
Go
package models
|
|
|
|
import (
|
|
"strconv"
|
|
"time"
|
|
|
|
log "github.com/sirupsen/logrus"
|
|
)
|
|
|
|
type MessageStructure string
|
|
type Protocol string
|
|
|
|
type MessageAttribute struct {
|
|
BinaryListValues []string `json:"BinaryListValues,omitempty" xml:"BinaryListValues,omitempty"` // currently unsupported by AWS
|
|
BinaryValue string `json:"BinaryValue,omitempty" xml:"BinaryValue,omitempty"`
|
|
DataType string `json:"DataType,omitempty" xml:"DataType,omitempty"`
|
|
StringListValues []string `json:"StringListValues,omitempty" xml:"StringListValues,omitempty"` // currently unsupported by AWS
|
|
StringValue string `json:"StringValue,omitempty" xml:"StringValue,omitempty"`
|
|
}
|
|
|
|
type SqsMessage struct {
|
|
MessageBody string
|
|
Uuid string
|
|
MD5OfMessageAttributes string
|
|
MD5OfMessageBody string
|
|
ReceiptHandle string
|
|
ReceiptTime time.Time
|
|
VisibilityTimeout time.Time
|
|
NumberOfReceives int
|
|
Retry int
|
|
MessageAttributes map[string]MessageAttribute
|
|
GroupID string
|
|
DeduplicationID string
|
|
SentTime time.Time
|
|
DelaySecs int
|
|
}
|
|
|
|
func (m *SqsMessage) IsReadyForReceipt() bool {
|
|
randomLatency, err := generateRandomLatency()
|
|
if err != nil {
|
|
log.Error(err)
|
|
return true
|
|
}
|
|
showAt := m.SentTime.Add(randomLatency).Add(time.Duration(m.DelaySecs) * time.Second)
|
|
return showAt.Before(time.Now())
|
|
}
|
|
|
|
type Queue struct {
|
|
Name string
|
|
URL string
|
|
Arn string
|
|
VisibilityTimeout int // seconds
|
|
ReceiveMessageWaitTimeSeconds int
|
|
DelaySeconds int
|
|
MaximumMessageSize int
|
|
MessageRetentionPeriod int // seconds // TODO - not used in the code yet
|
|
Messages []SqsMessage
|
|
DeadLetterQueue *Queue
|
|
MaxReceiveCount int
|
|
IsFIFO bool
|
|
FIFOMessages map[string]int
|
|
FIFOSequenceNumbers map[string]int
|
|
EnableDuplicates bool
|
|
Duplicates map[string]time.Time
|
|
Tags map[string]string
|
|
// LastPurgeTime — время последнего успешного PurgeQueue.
|
|
// Используется для запрета повторного purge ранее чем через 60 секунд
|
|
// (требование AWS: PurgeQueueInProgress). Нулевое значение = purge ещё не выполнялся.
|
|
LastPurgeTime time.Time
|
|
}
|
|
|
|
// NextSequenceNumber — возвращает следующий порядковый номер сообщения FIFO-группы.
|
|
//
|
|
// Логическая схема: каждый вызов инкрементирует счётчик группы и возвращает
|
|
// его строковое представление (порядковые номера в FIFO должны расти монотонно).
|
|
//
|
|
// ВАЖНО (фикс): раньше при ПЕРВОМ обращении к группе map пересоздавалась целиком
|
|
// (q.FIFOSequenceNumbers = map[string]int{groupId: 0}) — счётчики ВСЕХ остальных
|
|
// групп терялись, и порядковые номера начинали дублироваться. Теперь map только
|
|
// инициализируется один раз (если nil), а счётчик конкретной группы инкрементируется.
|
|
func (q *Queue) NextSequenceNumber(groupId string) string {
|
|
if q.FIFOSequenceNumbers == nil {
|
|
q.FIFOSequenceNumbers = make(map[string]int)
|
|
}
|
|
q.FIFOSequenceNumbers[groupId]++
|
|
return strconv.Itoa(q.FIFOSequenceNumbers[groupId])
|
|
}
|
|
|
|
// IsLocked — проверяет, заблокирована ли FIFO-группа.
|
|
// Группа блокируется при выдаче её сообщения в ReceiveMessage и остаётся
|
|
// заблокированной, пока сообщение in-flight (до DeleteMessage или таймаута
|
|
// видимости) — это гарантирует порядок обработки внутри группы.
|
|
func (q *Queue) IsLocked(groupId string) bool {
|
|
_, ok := q.FIFOMessages[groupId]
|
|
return ok
|
|
}
|
|
|
|
// LockGroup — блокирует FIFO-группу (помечает как занятую).
|
|
//
|
|
// ВАЖНО (фикс): раньше при ПЕРВОМ обращении map пересоздавалась целиком
|
|
// (q.FIFOMessages = map[string]int{groupId: 0}) — блокировки ВСЕХ остальных
|
|
// групп терялись, и их сообщения могли выдаваться одновременно. Теперь map
|
|
// только инициализируется один раз, а для группы проставляется ключ.
|
|
func (q *Queue) LockGroup(groupId string) {
|
|
if q.FIFOMessages == nil {
|
|
q.FIFOMessages = make(map[string]int)
|
|
}
|
|
q.FIFOMessages[groupId] = 0
|
|
}
|
|
|
|
// UnlockGroup — снимает блокировку FIFO-группы.
|
|
// Вызывается после удаления сообщения или при истечении таймаута видимости
|
|
// (см. PeriodicTasks → processQueueTick).
|
|
func (q *Queue) UnlockGroup(groupId string) {
|
|
if _, ok := q.FIFOMessages[groupId]; ok {
|
|
delete(q.FIFOMessages, groupId)
|
|
}
|
|
}
|
|
|
|
// IsDuplicate — проверяет, было ли сообщение с таким deduplicationId уже
|
|
// добавлено в очередь в пределах deduplication-окна.
|
|
//
|
|
// Логическая схема: дубликаты отслеживаются ТОЛЬКО для FIFO-очередей с
|
|
// включённой дупликацией (EnableDuplicates) и непустым deduplicationId.
|
|
// Для остальных случаев дубликаты не отслеживаются (AWS-семантика).
|
|
func (q *Queue) IsDuplicate(deduplicationId string) bool {
|
|
if !q.EnableDuplicates || !q.IsFIFO || deduplicationId == "" {
|
|
return false
|
|
}
|
|
|
|
_, ok := q.Duplicates[deduplicationId]
|
|
|
|
return ok
|
|
}
|
|
|
|
// InitDuplicatation — регистрирует deduplicationId с текущим временем.
|
|
// Эта отметка — начало deduplication-окна; просроченные записи удаляются
|
|
// фоновой задачей PeriodicTasks (deduplication period expired).
|
|
// Для не-FIFO очередей или пустого deduplicationId — no-op.
|
|
func (q *Queue) InitDuplicatation(deduplicationId string) {
|
|
if !q.EnableDuplicates || !q.IsFIFO || deduplicationId == "" {
|
|
return
|
|
}
|
|
|
|
if _, ok := q.Duplicates[deduplicationId]; !ok {
|
|
q.Duplicates[deduplicationId] = time.Now()
|
|
}
|
|
}
|