v0.1.35: фиксы — синхронные удаления в Redis (воскрешение очередей) + RedrivePolicy tenant-scoped DLQ
This commit is contained in:
@@ -71,7 +71,7 @@ func CreateQueueV1(req *http.Request) (int, interfaces.AbstractResponseBody) {
|
|||||||
if req.Header.Get("Content-Type") != "application/x-amz-json-1.0" {
|
if req.Header.Get("Content-Type") != "application/x-amz-json-1.0" {
|
||||||
provided = utils.ExtractQueueAttributes(req.PostForm)
|
provided = utils.ExtractQueueAttributes(req.PostForm)
|
||||||
}
|
}
|
||||||
if err := setQueueAttributesV1(queue, requestBody.Attributes, provided); err != nil {
|
if err := setQueueAttributesV1(queue, requestBody.Attributes, provided, t.AccessKey); err != nil {
|
||||||
models.SyncQueues.Unlock()
|
models.SyncQueues.Unlock()
|
||||||
return utils.CreateErrorResponseV1(err.Error(), true)
|
return utils.CreateErrorResponseV1(err.Error(), true)
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -22,9 +22,11 @@ import (
|
|||||||
// обнулял все остальные атрибуты очереди.
|
// обнулял все остальные атрибуты очереди.
|
||||||
// 3. Значения вне диапазонов AWS → ошибка InvalidParameterValue
|
// 3. Значения вне диапазонов AWS → ошибка InvalidParameterValue
|
||||||
// (раньше значения молча клэмпились в допустимый диапазон).
|
// (раньше значения молча клэмпились в допустимый диапазон).
|
||||||
// 4. RedrivePolicy: ARN DLQ разбирается, DLQ должна существовать, иначе
|
// 4. RedrivePolicy: ARN DLQ разбирается, DLQ ищется по TENANT-SCOPED ключу
|
||||||
// InvalidAttributeValue.
|
// "{accessKey}:{queueName}" (раньше — по голому имени, а ключи в map
|
||||||
func setQueueAttributesV1(q *models.Queue, attr models.QueueAttributes, provided map[string]string) error {
|
// tenant-scoped → DLQ никогда не находилась, RedrivePolicy не работал).
|
||||||
|
// DLQ должна принадлежать тому же тенанту.
|
||||||
|
func setQueueAttributesV1(q *models.Queue, attr models.QueueAttributes, provided map[string]string, tenantAccessKey string) error {
|
||||||
// Шаг 1: whitelist имён атрибутов (единый список — models.AttrNameWhitelist).
|
// Шаг 1: whitelist имён атрибутов (единый список — models.AttrNameWhitelist).
|
||||||
for name := range provided {
|
for name := range provided {
|
||||||
if !models.AttrNameWhitelist[name] {
|
if !models.AttrNameWhitelist[name] {
|
||||||
@@ -71,7 +73,9 @@ func setQueueAttributesV1(q *models.Queue, attr models.QueueAttributes, provided
|
|||||||
if attr.RedrivePolicy != (models.RedrivePolicy{}) {
|
if attr.RedrivePolicy != (models.RedrivePolicy{}) {
|
||||||
arnArray := strings.Split(attr.RedrivePolicy.DeadLetterTargetArn, ":")
|
arnArray := strings.Split(attr.RedrivePolicy.DeadLetterTargetArn, ":")
|
||||||
queueName := arnArray[len(arnArray)-1]
|
queueName := arnArray[len(arnArray)-1]
|
||||||
deadLetterQueue, ok := models.SyncQueues.Queues[queueName]
|
// DLQ ищем по tenant-scoped ключу — как хранятся все очереди.
|
||||||
|
dlqKey := tenantAccessKey + ":" + queueName
|
||||||
|
deadLetterQueue, ok := models.SyncQueues.Queues[dlqKey]
|
||||||
if !ok {
|
if !ok {
|
||||||
log.Error("Invalid RedrivePolicy Attribute")
|
log.Error("Invalid RedrivePolicy Attribute")
|
||||||
return fmt.Errorf("InvalidAttributeValue")
|
return fmt.Errorf("InvalidAttributeValue")
|
||||||
|
|||||||
@@ -61,7 +61,7 @@ func SetQueueAttributesV1(req *http.Request) (int, interfaces.AbstractResponseBo
|
|||||||
provided = utils.ExtractQueueAttributes(req.PostForm)
|
provided = utils.ExtractQueueAttributes(req.PostForm)
|
||||||
}
|
}
|
||||||
|
|
||||||
if err := setQueueAttributesV1(queue, requestBody.Attributes, provided); err != nil {
|
if err := setQueueAttributesV1(queue, requestBody.Attributes, provided, t.AccessKey); err != nil {
|
||||||
return utils.CreateErrorResponseV1(err.Error(), true)
|
return utils.CreateErrorResponseV1(err.Error(), true)
|
||||||
}
|
}
|
||||||
// Сохраняем атрибуты в Redis пока держим Lock (через defer)
|
// Сохраняем атрибуты в Redis пока держим Lock (через defer)
|
||||||
|
|||||||
+16
-21
@@ -184,55 +184,49 @@ func SaveMessages(queueKey string, msgs []models.SqsMessage) {
|
|||||||
})
|
})
|
||||||
}
|
}
|
||||||
|
|
||||||
// DeleteMessagePersist — удаляет одно сообщение из Redis асинхронно.
|
// DeleteMessagePersist — удаляет одно сообщение из Redis СИНХРОННО.
|
||||||
// Вызывать при DeleteMessage (после удаления из in-memory).
|
// Синхронность обязательна: если удаление потеряется при рестарте пода,
|
||||||
|
// сообщение «воскреснет» из Redis (факт: удалённые очереди возвращались
|
||||||
|
// после рестарта, когда удаления были асинхронными).
|
||||||
func DeleteMessagePersist(queueKey string, msgUuid string) {
|
func DeleteMessagePersist(queueKey string, msgUuid string) {
|
||||||
if Client == nil || msgUuid == "" {
|
if Client == nil || msgUuid == "" {
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
hashKey := redisMsgHashPrefix + queueKey
|
hashKey := redisMsgHashPrefix + queueKey
|
||||||
asyncWrite(func() {
|
|
||||||
ctx, cancel := context.WithTimeout(context.Background(), 3*time.Second)
|
ctx, cancel := context.WithTimeout(context.Background(), 3*time.Second)
|
||||||
defer cancel()
|
defer cancel()
|
||||||
if err := Client.HDel(ctx, hashKey, msgUuid).Err(); err != nil {
|
if err := Client.HDel(ctx, hashKey, msgUuid).Err(); err != nil {
|
||||||
log.Errorf("persistence: HDel message %s/%s: %v", queueKey, msgUuid, err)
|
log.Errorf("persistence: HDel message %s/%s: %v", queueKey, msgUuid, err)
|
||||||
}
|
}
|
||||||
})
|
|
||||||
}
|
}
|
||||||
|
|
||||||
// DeleteMessagesPersist — удаляет несколько сообщений из Redis.
|
// DeleteMessagesPersist — удаляет несколько сообщений из Redis СИНХРОННО
|
||||||
// Вызывать при DeleteMessageBatch.
|
// (иначе удалённые сообщения воскреснут после рестарта пода).
|
||||||
func DeleteMessagesPersist(queueKey string, uuids []string) {
|
func DeleteMessagesPersist(queueKey string, uuids []string) {
|
||||||
if Client == nil || len(uuids) == 0 {
|
if Client == nil || len(uuids) == 0 {
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
hashKey := redisMsgHashPrefix + queueKey
|
hashKey := redisMsgHashPrefix + queueKey
|
||||||
// Копируем uuids — вызывающий может переиспользовать slice
|
|
||||||
ids := make([]string, len(uuids))
|
|
||||||
copy(ids, uuids)
|
|
||||||
asyncWrite(func() {
|
|
||||||
ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
|
ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
|
||||||
defer cancel()
|
defer cancel()
|
||||||
if err := Client.HDel(ctx, hashKey, ids...).Err(); err != nil {
|
if err := Client.HDel(ctx, hashKey, uuids...).Err(); err != nil {
|
||||||
log.Errorf("persistence: HDel messages %s (%d): %v", queueKey, len(ids), err)
|
log.Errorf("persistence: HDel messages %s (%d): %v", queueKey, len(uuids), err)
|
||||||
}
|
}
|
||||||
})
|
|
||||||
}
|
}
|
||||||
|
|
||||||
// PurgeMessagesPersist — удаляет ВСЕ сообщения очереди из Redis (для PurgeQueue).
|
// PurgeMessagesPersist — удаляет ВСЕ сообщения очереди из Redis СИНХРОННО.
|
||||||
// Удаляет весь HASH ssq:msg:{queueKey}.
|
// Синхронность обязательна: при рестарте пода удалённые сообщения не должны
|
||||||
|
// воскреснуть (факт: асинхронное удаление очередей приводило к их возврату).
|
||||||
func PurgeMessagesPersist(queueKey string) {
|
func PurgeMessagesPersist(queueKey string) {
|
||||||
if Client == nil {
|
if Client == nil {
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
hashKey := redisMsgHashPrefix + queueKey
|
hashKey := redisMsgHashPrefix + queueKey
|
||||||
asyncWrite(func() {
|
|
||||||
ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
|
ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
|
||||||
defer cancel()
|
defer cancel()
|
||||||
if err := Client.Del(ctx, hashKey).Err(); err != nil {
|
if err := Client.Del(ctx, hashKey).Err(); err != nil {
|
||||||
log.Errorf("persistence: DEL messages hash %s: %v", queueKey, err)
|
log.Errorf("persistence: DEL messages hash %s: %v", queueKey, err)
|
||||||
}
|
}
|
||||||
})
|
|
||||||
}
|
}
|
||||||
|
|
||||||
// DeleteQueue — удаляет метаданные очереди И все её сообщения из Redis асинхронно.
|
// DeleteQueue — удаляет метаданные очереди И все её сообщения из Redis асинхронно.
|
||||||
@@ -242,12 +236,14 @@ func DeleteQueue(key string) {
|
|||||||
}
|
}
|
||||||
writeSeq := nextQueueWriteSeq(key)
|
writeSeq := nextQueueWriteSeq(key)
|
||||||
msgHashKey := redisMsgHashPrefix + key
|
msgHashKey := redisMsgHashPrefix + key
|
||||||
asyncWrite(func() {
|
// СИНХРОННО: удаление очереди обязано попасть в Redis до ответа клиенту.
|
||||||
|
// Асинхронное удаление приводило к «воскрешению» удалённых очередей после
|
||||||
|
// рестарта пода (факт 2026-08-15: 50 мусорных очередей вернулись из Redis).
|
||||||
|
ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
|
||||||
|
defer cancel()
|
||||||
if !isLatestQueueWriteSeq(key, writeSeq) {
|
if !isLatestQueueWriteSeq(key, writeSeq) {
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
|
|
||||||
defer cancel()
|
|
||||||
// Удаляем метаданные из общего HASH
|
// Удаляем метаданные из общего HASH
|
||||||
if err := Client.HDel(ctx, redisHashQueues, key).Err(); err != nil {
|
if err := Client.HDel(ctx, redisHashQueues, key).Err(); err != nil {
|
||||||
log.Errorf("persistence: HDel queue meta %q: %v", key, err)
|
log.Errorf("persistence: HDel queue meta %q: %v", key, err)
|
||||||
@@ -256,7 +252,6 @@ func DeleteQueue(key string) {
|
|||||||
if err := Client.Del(ctx, msgHashKey).Err(); err != nil {
|
if err := Client.Del(ctx, msgHashKey).Err(); err != nil {
|
||||||
log.Errorf("persistence: DEL messages hash %q: %v", key, err)
|
log.Errorf("persistence: DEL messages hash %q: %v", key, err)
|
||||||
}
|
}
|
||||||
})
|
|
||||||
}
|
}
|
||||||
|
|
||||||
// LoadAllQueues — загружает все очереди из Redis в память при старте сервиса.
|
// LoadAllQueues — загружает все очереди из Redis в память при старте сервиса.
|
||||||
|
|||||||
Reference in New Issue
Block a user