feat: persist messages separately and update ingress docs
This commit is contained in:
@@ -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)
|
||||
|
||||
@@ -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()
|
||||
|
||||
@@ -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
|
||||
}
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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
|
||||
}
|
||||
|
||||
@@ -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{
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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()
|
||||
|
||||
|
||||
+226
-19
@@ -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
|
||||
}
|
||||
|
||||
|
||||
@@ -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:
|
||||
|
||||
Reference in New Issue
Block a user