Files
SQS-service/app/models/models.go
T

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()
}
}