From 4a8357549128d1b1817f62ee13edd23d2162b2f6 Mon Sep 17 00:00:00 2001 From: Naeel Date: Sun, 12 Apr 2026 08:44:46 +0300 Subject: [PATCH] feat: persist messages separately and update ingress docs --- app/admin/admin.go | 4 + app/gosqs/change_message_visibility.go | 21 +- app/gosqs/change_message_visibility_batch.go | 2 + app/gosqs/delete_message.go | 5 +- app/gosqs/delete_message_batch.go | 6 +- app/gosqs/get_queue_attributes.go | 157 ++++++------ app/gosqs/purge_queue.go | 4 +- app/gosqs/send_message.go | 4 +- app/gosqs/send_message_batch.go | 6 +- app/persistence/redis.go | 245 +++++++++++++++++-- deployments/k8s/ingress.yaml | 9 +- 11 files changed, 359 insertions(+), 104 deletions(-) diff --git a/app/admin/admin.go b/app/admin/admin.go index 921eaa5..dc699d9 100644 --- a/app/admin/admin.go +++ b/app/admin/admin.go @@ -575,6 +575,8 @@ func (h *Handler) sendMessageToQueue(w http.ResponseWriter, r *http.Request) { } models.SyncQueues.Lock() models.SyncQueues.Queues[key].Messages = append(models.SyncQueues.Queues[key].Messages, msg) + // Персистим одно сообщение отдельно + persistence.SaveMessage(key, &msg) persistence.SaveQueue(key, models.SyncQueues.Queues[key]) models.SyncQueues.Unlock() w.Header().Set("Content-Type", "application/json") @@ -599,6 +601,8 @@ func (h *Handler) purgeQueue(w http.ResponseWriter, r *http.Request) { } models.SyncQueues.Lock() models.SyncQueues.Queues[key].Messages = models.SyncQueues.Queues[key].Messages[:0] + // Удаляем все сообщения из Redis одной командой + persistence.PurgeMessagesPersist(key) persistence.SaveQueue(key, models.SyncQueues.Queues[key]) models.SyncQueues.Unlock() log.Infof("admin: purged queue %s for tenant %s", queueName, t.ID) diff --git a/app/gosqs/change_message_visibility.go b/app/gosqs/change_message_visibility.go index 084f8d9..c2dbb3c 100644 --- a/app/gosqs/change_message_visibility.go +++ b/app/gosqs/change_message_visibility.go @@ -53,6 +53,8 @@ func ChangeMessageVisibilityV1(req *http.Request) (int, interfaces.AbstractRespo models.SyncQueues.Lock() messageFound := false + var changedMsgUuid string + var msgRemoved bool for i := 0; i < len(models.SyncQueues.Queues[key].Messages); i++ { queue := models.SyncQueues.Queues[key] msgs := queue.Messages @@ -66,18 +68,35 @@ func ChangeMessageVisibilityV1(req *http.Request) (int, interfaces.AbstractRespo if queue.MaxReceiveCount > 0 && queue.DeadLetterQueue != nil && msgs[i].Retry >= queue.MaxReceiveCount { + changedMsgUuid = msgs[i].Uuid queue.DeadLetterQueue.Messages = append(queue.DeadLetterQueue.Messages, msgs[i]) queue.Messages = append(queue.Messages[:i], queue.Messages[i+1:]...) + msgRemoved = true + } else { + changedMsgUuid = msgs[i].Uuid } } else { msgs[i].VisibilityTimeout = time.Now().Add(time.Duration(visibilityTimeout) * time.Second) + changedMsgUuid = msgs[i].Uuid } messageFound = true break } } - // Персистим изменение в Redis под Lock + // Персистим изменения в Redis под Lock if messageFound { + if msgRemoved { + // Сообщение удалено (перемещено в DLQ) — удаляем из Redis + persistence.DeleteMessagePersist(key, changedMsgUuid) + } else { + // Сообщение обновлено — перезаписываем в Redis + for i := range models.SyncQueues.Queues[key].Messages { + if models.SyncQueues.Queues[key].Messages[i].Uuid == changedMsgUuid { + persistence.SaveMessage(key, &models.SyncQueues.Queues[key].Messages[i]) + break + } + } + } persistence.SaveQueue(key, models.SyncQueues.Queues[key]) } models.SyncQueues.Unlock() diff --git a/app/gosqs/change_message_visibility_batch.go b/app/gosqs/change_message_visibility_batch.go index ad7c00c..182963c 100644 --- a/app/gosqs/change_message_visibility_batch.go +++ b/app/gosqs/change_message_visibility_batch.go @@ -93,6 +93,8 @@ func ChangeMessageVisibilityBatchV1(req *http.Request) (int, interfaces.Abstract } else { queue.Messages[i].VisibilityTimeout = time.Now().Add(time.Duration(entry.VisibilityTimeout) * time.Second) } + // Персистим изменённое сообщение отдельно + persistence.SaveMessage(key, &queue.Messages[i]) messageFound = true break } diff --git a/app/gosqs/delete_message.go b/app/gosqs/delete_message.go index 0f5e2f1..e31dd0a 100644 --- a/app/gosqs/delete_message.go +++ b/app/gosqs/delete_message.go @@ -47,10 +47,13 @@ func DeleteMessageV1(req *http.Request) (int, interfaces.AbstractResponseBody) { if _, ok := models.SyncQueues.Queues[key]; ok { for i, msg := range models.SyncQueues.Queues[key].Messages { if msg.ReceiptHandle == receiptHandle { + msgUuid := msg.Uuid models.SyncQueues.Queues[key].UnlockGroup(msg.GroupID) models.SyncQueues.Queues[key].Messages = append(models.SyncQueues.Queues[key].Messages[:i], models.SyncQueues.Queues[key].Messages[i+1:]...) delete(models.SyncQueues.Queues[key].Duplicates, msg.DeduplicationID) - // Сохраняем очередь в Redis пока держим Lock + // Удаляем одно сообщение из Redis — O(1) + persistence.DeleteMessagePersist(key, msgUuid) + // Метаданные очереди (duplicates, FIFO state) persistence.SaveQueue(key, models.SyncQueues.Queues[key]) respStruct := models.DeleteMessageResponse{ Xmlns: models.BaseXmlns, diff --git a/app/gosqs/delete_message_batch.go b/app/gosqs/delete_message_batch.go index 4217708..dcaa50e 100644 --- a/app/gosqs/delete_message_batch.go +++ b/app/gosqs/delete_message_batch.go @@ -73,6 +73,7 @@ func DeleteMessageBatchV1(req *http.Request) (int, interfaces.AbstractResponseBo } deletedEntries := make([]models.DeleteMessageBatchResultEntry, 0) + deletedUuids := make([]string, 0) remainingMessages := make([]models.SqsMessage, 0, len(models.SyncQueues.Queues[key].Messages)) for _, message := range models.SyncQueues.Queues[key].Messages { @@ -82,13 +83,16 @@ func DeleteMessageBatchV1(req *http.Request) (int, interfaces.AbstractResponseBo delete(models.SyncQueues.Queues[key].Duplicates, message.DeduplicationID) de.Deleted = true deletedEntries = append(deletedEntries, models.DeleteMessageBatchResultEntry{Id: de.Id}) + deletedUuids = append(deletedUuids, message.Uuid) } else { remainingMessages = append(remainingMessages, message) } } models.SyncQueues.Queues[key].Messages = remainingMessages - // Персистим обновлённое состояние очереди, чтобы не терять batch-delete после рестарта. + // Удаляем сообщения из Redis отдельно — O(batch_size) + persistence.DeleteMessagesPersist(key, deletedUuids) + // Метаданные очереди (duplicates, FIFO state) persistence.SaveQueue(key, models.SyncQueues.Queues[key]) notFoundEntries := make([]models.BatchResultErrorEntry, 0) diff --git a/app/gosqs/get_queue_attributes.go b/app/gosqs/get_queue_attributes.go index b5f2d9b..8e5c68e 100644 --- a/app/gosqs/get_queue_attributes.go +++ b/app/gosqs/get_queue_attributes.go @@ -3,35 +3,36 @@ package gosqs import ( -"fmt" -"net/http" -"strconv" -"strings" + "fmt" + "net/http" + "strconv" + "strings" -"shared-sqs/app/interfaces" -"shared-sqs/app/models" -"shared-sqs/app/utils" -log "github.com/sirupsen/logrus" + "shared-sqs/app/interfaces" + "shared-sqs/app/models" + "shared-sqs/app/utils" + + log "github.com/sirupsen/logrus" ) func GetQueueAttributesV1(req *http.Request) (int, interfaces.AbstractResponseBody) { -requestBody := models.NewGetQueueAttributesRequest() -ok := utils.REQUEST_TRANSFORMER(requestBody, req, false) -if !ok { -log.Error("Invalid Request - GetQueueAttributesV1") -return utils.CreateErrorResponseV1("InvalidParameterValue", true) -} -if requestBody.QueueUrl == "" { -log.Error("Missing QueueUrl - GetQueueAttributesV1") -return utils.CreateErrorResponseV1("InvalidParameterValue", true) -} + requestBody := models.NewGetQueueAttributesRequest() + ok := utils.REQUEST_TRANSFORMER(requestBody, req, false) + if !ok { + log.Error("Invalid Request - GetQueueAttributesV1") + return utils.CreateErrorResponseV1("InvalidParameterValue", true) + } + if requestBody.QueueUrl == "" { + log.Error("Missing QueueUrl - GetQueueAttributesV1") + return utils.CreateErrorResponseV1("InvalidParameterValue", true) + } -t := getTenantFromContext(req) -if t == nil { -return utils.CreateErrorResponseV1("InvalidClientTokenId", true) -} + t := getTenantFromContext(req) + if t == nil { + return utils.CreateErrorResponseV1("InvalidClientTokenId", true) + } -// Определяем набор запрошенных атрибутов (или All) + // Определяем набор запрошенных атрибутов (или All) requestedAttributes := func() map[string]bool { attrs := map[string]bool{} if len(requestBody.AttributeNames) == 0 { @@ -54,64 +55,64 @@ return utils.CreateErrorResponseV1("InvalidClientTokenId", true) } _, ok := requestedAttributes[attr] return ok -} + } -uriSegments := strings.Split(requestBody.QueueUrl, "/") -queueName := uriSegments[len(uriSegments)-1] -key := tenantQueueKey(t.AccessKey, queueName) + uriSegments := strings.Split(requestBody.QueueUrl, "/") + queueName := uriSegments[len(uriSegments)-1] + key := tenantQueueKey(t.AccessKey, queueName) -log.Infof("Get Queue Attributes: %s (tenant: %s)", queueName, t.ID) -queueAttributes := make([]models.Attribute, 0) + log.Infof("Get Queue Attributes: %s (tenant: %s)", queueName, t.ID) + queueAttributes := make([]models.Attribute, 0) -models.SyncQueues.RLock() -defer models.SyncQueues.RUnlock() -queue, ok := models.SyncQueues.Queues[key] -if !ok { -log.Errorf("Get Queue Attributes: %s queue does not exist for tenant %s", queueName, t.ID) -return utils.CreateErrorResponseV1("QueueNotFound", true) -} + models.SyncQueues.RLock() + defer models.SyncQueues.RUnlock() + queue, ok := models.SyncQueues.Queues[key] + if !ok { + log.Errorf("Get Queue Attributes: %s queue does not exist for tenant %s", queueName, t.ID) + return utils.CreateErrorResponseV1("QueueNotFound", true) + } -if shouldInclude("DelaySeconds") { -queueAttributes = append(queueAttributes, models.Attribute{Name: "DelaySeconds", Value: strconv.Itoa(queue.DelaySeconds)}) -} -if shouldInclude("MaximumMessageSize") { -queueAttributes = append(queueAttributes, models.Attribute{Name: "MaximumMessageSize", Value: strconv.Itoa(queue.MaximumMessageSize)}) -} -if shouldInclude("MessageRetentionPeriod") { -queueAttributes = append(queueAttributes, models.Attribute{Name: "MessageRetentionPeriod", Value: strconv.Itoa(queue.MessageRetentionPeriod)}) -} -if shouldInclude("ReceiveMessageWaitTimeSeconds") { -queueAttributes = append(queueAttributes, models.Attribute{Name: "ReceiveMessageWaitTimeSeconds", Value: strconv.Itoa(queue.ReceiveMessageWaitTimeSeconds)}) -} -if shouldInclude("VisibilityTimeout") { -queueAttributes = append(queueAttributes, models.Attribute{Name: "VisibilityTimeout", Value: strconv.Itoa(queue.VisibilityTimeout)}) -} -if shouldInclude("ApproximateNumberOfMessages") { -queueAttributes = append(queueAttributes, models.Attribute{Name: "ApproximateNumberOfMessages", Value: strconv.Itoa(len(queue.Messages))}) -} -if shouldInclude("ApproximateNumberOfMessagesNotVisible") { -queueAttributes = append(queueAttributes, models.Attribute{Name: "ApproximateNumberOfMessagesNotVisible", Value: strconv.Itoa(numberOfHiddenMessagesInQueue(*queue))}) -} -if shouldInclude("CreatedTimestamp") { -queueAttributes = append(queueAttributes, models.Attribute{Name: "CreatedTimestamp", Value: "0000000000"}) -} -if shouldInclude("LastModifiedTimestamp") { -queueAttributes = append(queueAttributes, models.Attribute{Name: "LastModifiedTimestamp", Value: "0000000000"}) -} -if shouldInclude("QueueArn") { -queueAttributes = append(queueAttributes, models.Attribute{Name: "QueueArn", Value: queue.Arn}) -} -if shouldInclude("RedrivePolicy") && queue.DeadLetterQueue != nil { -queueAttributes = append(queueAttributes, models.Attribute{ -Name: "RedrivePolicy", -Value: fmt.Sprintf(`{"maxReceiveCount":"%d", "deadLetterTargetArn":"%s"}`, queue.MaxReceiveCount, queue.DeadLetterQueue.Arn), -}) -} + if shouldInclude("DelaySeconds") { + queueAttributes = append(queueAttributes, models.Attribute{Name: "DelaySeconds", Value: strconv.Itoa(queue.DelaySeconds)}) + } + if shouldInclude("MaximumMessageSize") { + queueAttributes = append(queueAttributes, models.Attribute{Name: "MaximumMessageSize", Value: strconv.Itoa(queue.MaximumMessageSize)}) + } + if shouldInclude("MessageRetentionPeriod") { + queueAttributes = append(queueAttributes, models.Attribute{Name: "MessageRetentionPeriod", Value: strconv.Itoa(queue.MessageRetentionPeriod)}) + } + if shouldInclude("ReceiveMessageWaitTimeSeconds") { + queueAttributes = append(queueAttributes, models.Attribute{Name: "ReceiveMessageWaitTimeSeconds", Value: strconv.Itoa(queue.ReceiveMessageWaitTimeSeconds)}) + } + if shouldInclude("VisibilityTimeout") { + queueAttributes = append(queueAttributes, models.Attribute{Name: "VisibilityTimeout", Value: strconv.Itoa(queue.VisibilityTimeout)}) + } + if shouldInclude("ApproximateNumberOfMessages") { + queueAttributes = append(queueAttributes, models.Attribute{Name: "ApproximateNumberOfMessages", Value: strconv.Itoa(len(queue.Messages))}) + } + if shouldInclude("ApproximateNumberOfMessagesNotVisible") { + queueAttributes = append(queueAttributes, models.Attribute{Name: "ApproximateNumberOfMessagesNotVisible", Value: strconv.Itoa(numberOfHiddenMessagesInQueue(*queue))}) + } + if shouldInclude("CreatedTimestamp") { + queueAttributes = append(queueAttributes, models.Attribute{Name: "CreatedTimestamp", Value: "0000000000"}) + } + if shouldInclude("LastModifiedTimestamp") { + queueAttributes = append(queueAttributes, models.Attribute{Name: "LastModifiedTimestamp", Value: "0000000000"}) + } + if shouldInclude("QueueArn") { + queueAttributes = append(queueAttributes, models.Attribute{Name: "QueueArn", Value: queue.Arn}) + } + if shouldInclude("RedrivePolicy") && queue.DeadLetterQueue != nil { + queueAttributes = append(queueAttributes, models.Attribute{ + Name: "RedrivePolicy", + Value: fmt.Sprintf(`{"maxReceiveCount":"%d", "deadLetterTargetArn":"%s"}`, queue.MaxReceiveCount, queue.DeadLetterQueue.Arn), + }) + } -respStruct := models.GetQueueAttributesResponse{ -Xmlns: models.BaseXmlns, -Result: models.GetQueueAttributesResult{Attrs: queueAttributes}, -Metadata: models.BaseResponseMetadata, -} -return http.StatusOK, respStruct + respStruct := models.GetQueueAttributesResponse{ + Xmlns: models.BaseXmlns, + Result: models.GetQueueAttributesResult{Attrs: queueAttributes}, + Metadata: models.BaseResponseMetadata, + } + return http.StatusOK, respStruct } diff --git a/app/gosqs/purge_queue.go b/app/gosqs/purge_queue.go index e661640..0ae1466 100644 --- a/app/gosqs/purge_queue.go +++ b/app/gosqs/purge_queue.go @@ -42,7 +42,9 @@ func PurgeQueueV1(req *http.Request) (int, interfaces.AbstractResponseBody) { log.Infof("Purging Queue: %s (tenant: %s)", queueName, t.ID) models.SyncQueues.Queues[key].Messages = nil models.SyncQueues.Queues[key].Duplicates = make(map[string]time.Time) - // Сохраняем пустую очередь в Redis пока держим Lock + // Удаляем все сообщения из Redis одной командой DEL + persistence.PurgeMessagesPersist(key) + // Сохраняем пустые метаданные очереди persistence.SaveQueue(key, models.SyncQueues.Queues[key]) respStruct := models.PurgeQueueResponse{ diff --git a/app/gosqs/send_message.go b/app/gosqs/send_message.go index e1130c0..41cacef 100644 --- a/app/gosqs/send_message.go +++ b/app/gosqs/send_message.go @@ -112,12 +112,14 @@ func SendMessageV1(req *http.Request) (int, interfaces.AbstractResponseBody) { if !models.SyncQueues.Queues[key].IsDuplicate(messageDeduplicationID) { models.SyncQueues.Queues[key].Messages = append(models.SyncQueues.Queues[key].Messages, msg) + // Персистим одно сообщение отдельно — O(msg_size) вместо O(N*msg_size) + persistence.SaveMessage(key, &msg) } else { log.Debugf("Duplicate message deduplicationId [%s] in queue [%s]", messageDeduplicationID, queueName) } models.SyncQueues.Queues[key].InitDuplicatation(messageDeduplicationID) - // Сохраняем очередь в Redis пока держим Lock + // Сохраняем метаданные очереди (FIFO state, duplicates) — без Messages persistence.SaveQueue(key, models.SyncQueues.Queues[key]) models.SyncQueues.Unlock() // Логируем только метаданные — тело сообщения не логируется (perf + security) diff --git a/app/gosqs/send_message_batch.go b/app/gosqs/send_message_batch.go index b5e8c53..5cf24b4 100644 --- a/app/gosqs/send_message_batch.go +++ b/app/gosqs/send_message_batch.go @@ -106,6 +106,7 @@ func SendMessageBatchV1(req *http.Request) (int, interfaces.AbstractResponseBody models.SyncQueues.Lock() queue = models.SyncQueues.Queues[key] + newMsgs := make([]models.SqsMessage, 0, len(sendEntries)) for _, sendEntry := range sendEntries { msg := models.SqsMessage{MessageBody: sendEntry.MessageBody} if len(sendEntry.MessageAttributes) > 0 { @@ -136,10 +137,13 @@ func SendMessageBatchV1(req *http.Request) (int, interfaces.AbstractResponseBody MD5OfMessageAttributes: msg.MD5OfMessageAttributes, SequenceNumber: fifoSeqNumber, }) + newMsgs = append(newMsgs, msg) // Логируем только метаданные — тело сообщения не логируется (perf + security) log.Debugf("Queue: %s, MessageId: %s, Size: %d bytes", queueName, msg.Uuid, len(sendEntry.MessageBody)) } - // Персистим батч-изменение одним снапшотом под lock. + // Персистим новые сообщения отдельно — O(batch_size*msg_size) вместо O(N*msg_size) + persistence.SaveMessages(key, newMsgs) + // Метаданные очереди (FIFO state, duplicates) — без Messages persistence.SaveQueue(key, queue) models.SyncQueues.Unlock() diff --git a/app/persistence/redis.go b/app/persistence/redis.go index ef7871f..4e1cd42 100644 --- a/app/persistence/redis.go +++ b/app/persistence/redis.go @@ -4,7 +4,13 @@ // Redis — источник правды для восстановления после рестарта. // Все записи в Redis асинхронны (горутина) — не блокируют SQS-операции. // Сериализация (json.Marshal) происходит синхронно пока вызывающий держит мьютекс — консистентный снапшот. -// Created: 2026-04-10 +// +// Схема хранения v2 (2026-04-11): +// ssq:queues (HASH) — queueKey → JSON(метаданные очереди без Messages) +// ssq:msg:{queueKey} (HASH) — uuid → JSON(SqsMessage) +// Это позволяет O(1) на каждое сообщение вместо O(N×msg_size) при SaveQueue. +// Миграция со старого формата (Messages внутри ssq:queues) — автоматическая при LoadAllQueues. +// Created: 2026-04-10 | Modified: 2026-04-11 package persistence @@ -34,8 +40,10 @@ var ( const ( // redisHashTenants — HASH: tenantID → JSON тенанта redisHashTenants = "ssq:tenants" - // redisHashQueues — HASH: queueKey → JSON очереди (включая сообщения) + // redisHashQueues — HASH: queueKey → JSON метаданных очереди (без Messages с v2) redisHashQueues = "ssq:queues" + // redisMsgHashPrefix — префикс для per-message HASH: ssq:msg:{queueKey} → uuid → JSON(SqsMessage) + redisMsgHashPrefix = "ssq:msg:" ) // Connect — подключается к Redis и проверяет ping. @@ -74,22 +82,37 @@ func asyncWrite(fn func()) { }() } -// SaveQueue — сохраняет очередь (с сообщениями) в Redis асинхронно. -// ВАЖНО: вызывать пока вызывающий держит SyncQueues.Lock() — тогда json.Marshal +// marshalQueueMeta — сериализует очередь БЕЗ Messages (и без Messages в DLQ). +// DeadLetterQueue сохраняется как ссылка (Name/URL/Arn), но без сообщений. +// ВАЖНО: вызывать под SyncQueues.Lock() — временно nil'ит Messages. +func marshalQueueMeta(queue *models.Queue) ([]byte, error) { + savedMsgs := queue.Messages + queue.Messages = nil + var savedDLQMsgs []models.SqsMessage + if queue.DeadLetterQueue != nil { + savedDLQMsgs = queue.DeadLetterQueue.Messages + queue.DeadLetterQueue.Messages = nil + } + data, err := json.Marshal(queue) + queue.Messages = savedMsgs + if queue.DeadLetterQueue != nil { + queue.DeadLetterQueue.Messages = savedDLQMsgs + } + return data, err +} + +// SaveQueue — сохраняет МЕТАДАННЫЕ очереди (без сообщений) в Redis асинхронно. +// ВАЖНО: вызывать пока вызывающий держит SyncQueues.Lock() — тогда marshalQueueMeta // создаёт консистентный снапшот. Горутина только делает сетевой вызов. +// Для сохранения сообщений используй SaveMessage/SaveMessages. func SaveQueue(key string, queue *models.Queue) { if Client == nil { return } - // Сериализуем синхронно под мьютексом вызывающего → консистентный снапшот - data, err := json.Marshal(queue) + // Сериализуем метаданные синхронно под мьютексом вызывающего → консистентный снапшот + data, err := marshalQueueMeta(queue) if err != nil { - log.Errorf("persistence: marshal queue %q: %v", key, err) - return - } - // Fix #4.3: Redis size guard — не сохранять если > 50MB (OOM protection) - if len(data) > 50*1024*1024 { - log.Warnf("persistence: queue %q too large for Redis (%d bytes), skipping", key, len(data)) + log.Errorf("persistence: marshal queue meta %q: %v", key, err) return } writeSeq := nextQueueWriteSeq(key) @@ -100,42 +123,162 @@ func SaveQueue(key string, queue *models.Queue) { ctx, cancel := context.WithTimeout(context.Background(), 3*time.Second) defer cancel() if err := Client.HSet(ctx, redisHashQueues, key, string(data)).Err(); err != nil { - log.Errorf("persistence: HSet queue %q: %v", key, err) + log.Errorf("persistence: HSet queue meta %q: %v", key, err) } }) } -// DeleteQueue — удаляет очередь из Redis асинхронно. +// SaveMessage — сохраняет ОДНО сообщение в Redis асинхронно. +// Ключ: ssq:msg:{queueKey}, поле: msg.Uuid, значение: JSON(SqsMessage). +// Вызывать под Lock — json.Marshal делает консистентный снапшот сообщения. +func SaveMessage(queueKey string, msg *models.SqsMessage) { + if Client == nil { + return + } + data, err := json.Marshal(msg) + if err != nil { + log.Errorf("persistence: marshal message %s/%s: %v", queueKey, msg.Uuid, err) + return + } + hashKey := redisMsgHashPrefix + queueKey + asyncWrite(func() { + ctx, cancel := context.WithTimeout(context.Background(), 3*time.Second) + defer cancel() + if err := Client.HSet(ctx, hashKey, msg.Uuid, string(data)).Err(); err != nil { + log.Errorf("persistence: HSet message %s/%s: %v", queueKey, msg.Uuid, err) + } + }) +} + +// SaveMessages — сохраняет несколько сообщений в Redis одним pipeline. +// Вызывать под Lock. +func SaveMessages(queueKey string, msgs []models.SqsMessage) { + if Client == nil || len(msgs) == 0 { + return + } + // Сериализуем все сообщения синхронно под Lock + fields := make(map[string]string, len(msgs)) + for i := range msgs { + data, err := json.Marshal(&msgs[i]) + if err != nil { + log.Errorf("persistence: marshal message %s/%s: %v", queueKey, msgs[i].Uuid, err) + continue + } + fields[msgs[i].Uuid] = string(data) + } + if len(fields) == 0 { + return + } + hashKey := redisMsgHashPrefix + queueKey + asyncWrite(func() { + ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second) + defer cancel() + // Конвертируем map[string]string → []interface{} для HSet + args := make([]interface{}, 0, len(fields)*2) + for k, v := range fields { + args = append(args, k, v) + } + if err := Client.HSet(ctx, hashKey, args...).Err(); err != nil { + log.Errorf("persistence: HSet messages %s (%d msgs): %v", queueKey, len(fields), err) + } + }) +} + +// DeleteMessagePersist — удаляет одно сообщение из Redis асинхронно. +// Вызывать при DeleteMessage (после удаления из in-memory). +func DeleteMessagePersist(queueKey string, msgUuid string) { + if Client == nil || msgUuid == "" { + return + } + hashKey := redisMsgHashPrefix + queueKey + asyncWrite(func() { + ctx, cancel := context.WithTimeout(context.Background(), 3*time.Second) + defer cancel() + if err := Client.HDel(ctx, hashKey, msgUuid).Err(); err != nil { + log.Errorf("persistence: HDel message %s/%s: %v", queueKey, msgUuid, err) + } + }) +} + +// DeleteMessagesPersist — удаляет несколько сообщений из Redis. +// Вызывать при DeleteMessageBatch. +func DeleteMessagesPersist(queueKey string, uuids []string) { + if Client == nil || len(uuids) == 0 { + return + } + hashKey := redisMsgHashPrefix + queueKey + // Копируем uuids — вызывающий может переиспользовать slice + ids := make([]string, len(uuids)) + copy(ids, uuids) + asyncWrite(func() { + ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second) + defer cancel() + if err := Client.HDel(ctx, hashKey, ids...).Err(); err != nil { + log.Errorf("persistence: HDel messages %s (%d): %v", queueKey, len(ids), err) + } + }) +} + +// PurgeMessagesPersist — удаляет ВСЕ сообщения очереди из Redis (для PurgeQueue). +// Удаляет весь HASH ssq:msg:{queueKey}. +func PurgeMessagesPersist(queueKey string) { + if Client == nil { + return + } + hashKey := redisMsgHashPrefix + queueKey + asyncWrite(func() { + ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second) + defer cancel() + if err := Client.Del(ctx, hashKey).Err(); err != nil { + log.Errorf("persistence: DEL messages hash %s: %v", queueKey, err) + } + }) +} + +// DeleteQueue — удаляет метаданные очереди И все её сообщения из Redis асинхронно. func DeleteQueue(key string) { if Client == nil { return } writeSeq := nextQueueWriteSeq(key) + msgHashKey := redisMsgHashPrefix + key asyncWrite(func() { if !isLatestQueueWriteSeq(key, writeSeq) { return } - ctx, cancel := context.WithTimeout(context.Background(), 3*time.Second) + ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second) defer cancel() + // Удаляем метаданные из общего HASH if err := Client.HDel(ctx, redisHashQueues, key).Err(); err != nil { - log.Errorf("persistence: HDel queue %q: %v", key, err) + log.Errorf("persistence: HDel queue meta %q: %v", key, err) + } + // Удаляем весь HASH с сообщениями + if err := Client.Del(ctx, msgHashKey).Err(); err != nil { + log.Errorf("persistence: DEL messages hash %q: %v", key, err) } }) } // LoadAllQueues — загружает все очереди из Redis в память при старте сервиса. -// Инициализирует nil-maps чтобы избежать panic при deduplication/FIFO операциях. +// Схема v2: метаданные из ssq:queues, сообщения из ssq:msg:{key}. +// Миграция v1→v2: если в ssq:queues лежит JSON со встроенными Messages, +// они извлекаются в отдельный HASH и метаданные пересохраняются без Messages. func LoadAllQueues() (map[string]*models.Queue, error) { if Client == nil { return nil, nil } - ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second) + ctx, cancel := context.WithTimeout(context.Background(), 30*time.Second) defer cancel() + + // 1. Загружаем метаданные очередей raw, err := Client.HGetAll(ctx, redisHashQueues).Result() if err != nil { return nil, fmt.Errorf("redis HGetAll queues: %w", err) } + queues := make(map[string]*models.Queue, len(raw)) + migrateKeys := make([]string, 0) // ключи для миграции v1→v2 + for k, v := range raw { var q models.Queue if err := json.Unmarshal([]byte(v), &q); err != nil { @@ -152,9 +295,73 @@ func LoadAllQueues() (map[string]*models.Queue, error) { if q.FIFOSequenceNumbers == nil { q.FIFOSequenceNumbers = make(map[string]int) } + + // Миграция v1→v2: если в JSON есть Messages — это старый формат + if len(q.Messages) > 0 { + migrateKeys = append(migrateKeys, k) + log.Infof("persistence: миграция v1→v2 для %q (%d сообщений)", k, len(q.Messages)) + } + queues[k] = &q } - log.Infof("persistence: загружено %d очередей из Redis", len(queues)) + + // 2. Миграция: выносим Messages из старого формата в отдельные HASH'ы + for _, k := range migrateKeys { + q := queues[k] + hashKey := redisMsgHashPrefix + k + // Сохраняем каждое сообщение отдельно + if len(q.Messages) > 0 { + args := make([]interface{}, 0, len(q.Messages)*2) + for i := range q.Messages { + data, err := json.Marshal(&q.Messages[i]) + if err != nil { + log.Errorf("persistence: migrate marshal msg %s/%s: %v", k, q.Messages[i].Uuid, err) + continue + } + args = append(args, q.Messages[i].Uuid, string(data)) + } + if len(args) > 0 { + if err := Client.HSet(ctx, hashKey, args...).Err(); err != nil { + log.Errorf("persistence: migrate HSet messages %q: %v", k, err) + } + } + } + // Пересохраняем метаданные без Messages + metaData, err := marshalQueueMeta(q) + if err == nil { + if err := Client.HSet(ctx, redisHashQueues, k, string(metaData)).Err(); err != nil { + log.Errorf("persistence: migrate HSet meta %q: %v", k, err) + } + } + log.Infof("persistence: миграция %q завершена — %d сообщений вынесены в %s", k, len(q.Messages), hashKey) + } + + // 3. Загружаем сообщения из отдельных HASH'ов для каждой очереди + for k, q := range queues { + hashKey := redisMsgHashPrefix + k + msgRaw, err := Client.HGetAll(ctx, hashKey).Result() + if err != nil { + log.Errorf("persistence: HGetAll messages %q: %v", k, err) + continue + } + if len(msgRaw) > 0 { + msgs := make([]models.SqsMessage, 0, len(msgRaw)) + for uuid, v := range msgRaw { + var m models.SqsMessage + if err := json.Unmarshal([]byte(v), &m); err != nil { + log.Errorf("persistence: unmarshal message %s/%s: %v", k, uuid, err) + continue + } + msgs = append(msgs, m) + } + q.Messages = msgs + } else if len(q.Messages) == 0 { + // Пустая очередь — Messages уже nil, оставляем + q.Messages = nil + } + } + + log.Infof("persistence: загружено %d очередей из Redis (миграций: %d)", len(queues), len(migrateKeys)) return queues, nil } diff --git a/deployments/k8s/ingress.yaml b/deployments/k8s/ingress.yaml index d99a082..483f7ef 100644 --- a/deployments/k8s/ingress.yaml +++ b/deployments/k8s/ingress.yaml @@ -1,6 +1,6 @@ # deployments/k8s/ingress.yaml # Ingress для shared-sqs на домене qu.kube5s.ru -# Created: 2026-04-09, Updated: 2026-04-10 +# Created: 2026-04-09, Updated: 2026-04-11 apiVersion: networking.k8s.io/v1 kind: Ingress metadata: @@ -11,6 +11,13 @@ metadata: nginx.ingress.kubernetes.io/proxy-body-size: 10m nginx.ingress.kubernetes.io/proxy-read-timeout: "30" nginx.ingress.kubernetes.io/proxy-send-timeout: "30" + # Fix: 64KB+ payload (boto3/urllib3 2.0 шлёт headers и body отдельными TCP write, + # botocore убирает TCP_NODELAY → последние ~16KB body не доходят за client_body_timeout). + # Решение: proxy-request-buffering=off стримит body напрямую в backend без ожидания полного тела. + nginx.ingress.kubernetes.io/client-body-buffer-size: "512k" + nginx.ingress.kubernetes.io/proxy-request-buffering: "off" + nginx.ingress.kubernetes.io/server-snippet: | + client_body_timeout 120s; spec: ingressClassName: nginx rules: