diff --git a/app/gosqs/change_message_visibility.go b/app/gosqs/change_message_visibility.go index c2dbb3c..9836875 100644 --- a/app/gosqs/change_message_visibility.go +++ b/app/gosqs/change_message_visibility.go @@ -47,19 +47,26 @@ func ChangeMessageVisibilityV1(req *http.Request) (int, interfaces.AbstractRespo return utils.CreateErrorResponseV1("InvalidParameterValue", true) } - if _, ok := models.SyncQueues.Queues[key]; !ok { + // ВСЯ работа с очередью — под одним Lock. + // Раньше проверка существования шла БЕЗ блокировки (чтение map) — + // concurrent map read and map write при параллельном DeleteQueue + // (fatal error, под падал под нагрузкой). + // defer Unlock: при любой панике в секции мьютекс не останется захваченным. + models.SyncQueues.Lock() + defer models.SyncQueues.Unlock() + + queue, exists := models.SyncQueues.Queues[key] + if !exists { return utils.CreateErrorResponseV1("QueueNotFound", true) } - 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] + for i := 0; i < len(queue.Messages); i++ { msgs := queue.Messages if msgs[i].ReceiptHandle == receiptHandle { - timeout := models.SyncQueues.Queues[key].VisibilityTimeout + timeout := queue.VisibilityTimeout if visibilityTimeout == 0 { msgs[i].ReceiptTime = time.Now().UTC() msgs[i].ReceiptHandle = "" @@ -83,23 +90,22 @@ func ChangeMessageVisibilityV1(req *http.Request) (int, interfaces.AbstractRespo 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]) + for i := range queue.Messages { + if queue.Messages[i].Uuid == changedMsgUuid { + persistence.SaveMessage(key, &queue.Messages[i]) break } } } - persistence.SaveQueue(key, models.SyncQueues.Queues[key]) + persistence.SaveQueue(key, queue) } - models.SyncQueues.Unlock() if !messageFound { return utils.CreateErrorResponseV1("MessageNotInFlight", true) } diff --git a/app/gosqs/change_message_visibility_batch.go b/app/gosqs/change_message_visibility_batch.go index 182963c..ad86ad1 100644 --- a/app/gosqs/change_message_visibility_batch.go +++ b/app/gosqs/change_message_visibility_batch.go @@ -42,10 +42,6 @@ func ChangeMessageVisibilityBatchV1(req *http.Request) (int, interfaces.Abstract key := tenantQueueKey(t.AccessKey, queueName) - if _, ok := models.SyncQueues.Queues[key]; !ok { - return utils.CreateErrorResponseV1("QueueNotFound", true) - } - if len(requestBody.Entries) == 0 { return utils.CreateErrorResponseV1("EmptyBatchRequest", true) } @@ -66,6 +62,13 @@ func ChangeMessageVisibilityBatchV1(req *http.Request) (int, interfaces.Abstract models.SyncQueues.Lock() defer models.SyncQueues.Unlock() + // Проверка существования очереди ВНУТРИ Lock: раньше была БЕЗ блокировки + // (чтение map) — concurrent map read/write при параллельном DeleteQueue. + queue, exists := models.SyncQueues.Queues[key] + if !exists { + return utils.CreateErrorResponseV1("QueueNotFound", true) + } + successEntries := make([]models.ChangeMessageVisibilityBatchResultEntry, 0) failedEntries := make([]models.BatchResultErrorEntry, 0) @@ -81,7 +84,6 @@ func ChangeMessageVisibilityBatchV1(req *http.Request) (int, interfaces.Abstract } messageFound := false - queue := models.SyncQueues.Queues[key] for i := 0; i < len(queue.Messages); i++ { if queue.Messages[i].ReceiptHandle == entry.ReceiptHandle { if entry.VisibilityTimeout == 0 { @@ -114,7 +116,7 @@ func ChangeMessageVisibilityBatchV1(req *http.Request) (int, interfaces.Abstract // Персистим изменения в Redis под Lock if len(successEntries) > 0 { - persistence.SaveQueue(key, models.SyncQueues.Queues[key]) + persistence.SaveQueue(key, queue) } respStruct := models.ChangeMessageVisibilityBatchResponse{ diff --git a/app/gosqs/delete_message_batch.go b/app/gosqs/delete_message_batch.go index dcaa50e..40d36b8 100644 --- a/app/gosqs/delete_message_batch.go +++ b/app/gosqs/delete_message_batch.go @@ -40,10 +40,6 @@ func DeleteMessageBatchV1(req *http.Request) (int, interfaces.AbstractResponseBo key := tenantQueueKey(t.AccessKey, queueName) - if _, ok := models.SyncQueues.Queues[key]; !ok { - return utils.CreateErrorResponseV1("QueueNotFound", true) - } - if len(requestBody.Entries) == 0 { return utils.CreateErrorResponseV1("EmptyBatchRequest", true) } @@ -63,6 +59,13 @@ func DeleteMessageBatchV1(req *http.Request) (int, interfaces.AbstractResponseBo models.SyncQueues.Lock() defer models.SyncQueues.Unlock() + // Проверка существования очереди ВНУТРИ Lock: раньше была БЕЗ блокировки + // (чтение map) — concurrent map read/write при параллельном DeleteQueue. + queue, exists := models.SyncQueues.Queues[key] + if !exists { + return utils.CreateErrorResponseV1("QueueNotFound", true) + } + deleteMessageMap := make(map[string]*deleteEntry) for _, entry := range requestBody.Entries { deleteMessageMap[entry.ReceiptHandle] = &deleteEntry{ @@ -74,13 +77,13 @@ 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)) + remainingMessages := make([]models.SqsMessage, 0, len(queue.Messages)) - for _, message := range models.SyncQueues.Queues[key].Messages { + for _, message := range queue.Messages { if de, found := deleteMessageMap[message.ReceiptHandle]; found { log.Debugf("FIFO Queue %s unlocking group %s:", queueName, message.GroupID) - models.SyncQueues.Queues[key].UnlockGroup(message.GroupID) - delete(models.SyncQueues.Queues[key].Duplicates, message.DeduplicationID) + queue.UnlockGroup(message.GroupID) + delete(queue.Duplicates, message.DeduplicationID) de.Deleted = true deletedEntries = append(deletedEntries, models.DeleteMessageBatchResultEntry{Id: de.Id}) deletedUuids = append(deletedUuids, message.Uuid) @@ -89,11 +92,11 @@ func DeleteMessageBatchV1(req *http.Request) (int, interfaces.AbstractResponseBo } } - models.SyncQueues.Queues[key].Messages = remainingMessages + queue.Messages = remainingMessages // Удаляем сообщения из Redis отдельно — O(batch_size) persistence.DeleteMessagesPersist(key, deletedUuids) // Метаданные очереди (duplicates, FIFO state) - persistence.SaveQueue(key, models.SyncQueues.Queues[key]) + persistence.SaveQueue(key, queue) notFoundEntries := make([]models.BatchResultErrorEntry, 0) for _, de := range deleteMessageMap { diff --git a/app/gosqs/receive_message.go b/app/gosqs/receive_message.go index 799c1a8..ff071a0 100644 --- a/app/gosqs/receive_message.go +++ b/app/gosqs/receive_message.go @@ -83,23 +83,20 @@ func ReceiveMessageV1(req *http.Request) (int, interfaces.AbstractResponseBody) key := tenantQueueKey(t.AccessKey, queueName) - // Проверка существования очереди — под RLock (защита от параллельной записи - // в map; раньше map читалась без блокировки — гонка данных). + // ОДИН RLock на проверку существования И чтение атрибута очереди: + // раньше было два отдельных RLock — между ними очередь могла быть удалена + // (nil-deref при чтении ReceiveMessageWaitTimeSeconds). models.SyncQueues.RLock() - _, queueExists := models.SyncQueues.Queues[key] - models.SyncQueues.RUnlock() + queue, queueExists := models.SyncQueues.Queues[key] if !queueExists { + models.SyncQueues.RUnlock() return utils.CreateErrorResponseV1("QueueNotFound", true) } - - // --- Шаг 5: WaitTimeSeconds (long polling) --- - // 0 (не задан) → атрибут очереди ReceiveMessageWaitTimeSeconds (default 0). waitTimeSeconds := requestBody.WaitTimeSeconds if waitTimeSeconds == 0 { - models.SyncQueues.RLock() - waitTimeSeconds = models.SyncQueues.Queues[key].ReceiveMessageWaitTimeSeconds - models.SyncQueues.RUnlock() + waitTimeSeconds = queue.ReceiveMessageWaitTimeSeconds } + models.SyncQueues.RUnlock() if waitTimeSeconds < 0 || waitTimeSeconds > MaxReceiveMessageWaitTimeSeconds { return utils.CreateErrorResponseV1("InvalidParameterValue", true) } diff --git a/app/gosqs/send_message.go b/app/gosqs/send_message.go index a21b4af..f764845 100644 --- a/app/gosqs/send_message.go +++ b/app/gosqs/send_message.go @@ -129,22 +129,36 @@ func SendMessageV1(req *http.Request) (int, interfaces.AbstractResponseBody) { msg.DelaySecs = delaySecs models.SyncQueues.Lock() - fifoSeqNumber := "" - if models.SyncQueues.Queues[key].IsFIFO { - fifoSeqNumber = models.SyncQueues.Queues[key].NextSequenceNumber(messageGroupID) + // Перепроверка под Lock: очередь могла быть удалена между RLock-проверкой + // выше и захватом Lock (иначе nil-pointer deref). + queue = models.SyncQueues.Queues[key] + if queue == nil { + models.SyncQueues.Unlock() + return utils.CreateErrorResponseV1("QueueNotFound", true) + } + // Повторная проверка лимита под Lock: между RLock-снимком и захватом Lock + // другие горутины могли добавить сообщения (мягкий OOM-лимит). + if len(queue.Messages) >= MaxMessagesForQueue(queue.IsFIFO) { + models.SyncQueues.Unlock() + return utils.CreateErrorResponseV1("OverLimit", true) } - if !models.SyncQueues.Queues[key].IsDuplicate(messageDeduplicationID) { - models.SyncQueues.Queues[key].Messages = append(models.SyncQueues.Queues[key].Messages, msg) + fifoSeqNumber := "" + if queue.IsFIFO { + fifoSeqNumber = queue.NextSequenceNumber(messageGroupID) + } + + if !queue.IsDuplicate(messageDeduplicationID) { + queue.Messages = append(queue.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) + queue.InitDuplicatation(messageDeduplicationID) // Сохраняем метаданные очереди (FIFO state, duplicates) — без Messages - persistence.SaveQueue(key, models.SyncQueues.Queues[key]) + persistence.SaveQueue(key, queue) models.SyncQueues.Unlock() // Логируем только метаданные — тело сообщения не логируется (perf + security) log.Infof("Queue: %s, MessageId: %s, Size: %d bytes", queueName, msg.Uuid, len(messageBody)) diff --git a/app/gosqs/send_message_batch.go b/app/gosqs/send_message_batch.go index 5cf24b4..9a46896 100644 --- a/app/gosqs/send_message_batch.go +++ b/app/gosqs/send_message_batch.go @@ -105,7 +105,13 @@ func SendMessageBatchV1(req *http.Request) (int, interfaces.AbstractResponseBody } models.SyncQueues.Lock() + // Перепроверка под Lock: очередь могла быть удалена между RLock-проверкой + // выше и захватом Lock (иначе nil-pointer deref). queue = models.SyncQueues.Queues[key] + if queue == nil { + models.SyncQueues.Unlock() + return utils.CreateErrorResponseV1("QueueNotFound", true) + } newMsgs := make([]models.SqsMessage, 0, len(sendEntries)) for _, sendEntry := range sendEntries { msg := models.SqsMessage{MessageBody: sendEntry.MessageBody} diff --git a/app/models/errors.go b/app/models/errors.go index c36eebd..d52b344 100644 --- a/app/models/errors.go +++ b/app/models/errors.go @@ -28,6 +28,8 @@ func init() { "InvalidAttributeName": {HttpError: http.StatusBadRequest, Type: "InvalidAttributeName", Code: "AWS.SimpleQueueService.InvalidAttributeName", Message: "Unknown Attribute."}, // PurgeQueueInProgress — повторный PurgeQueue на очереди ранее чем через 60 секунд "PurgeQueueInProgress": {HttpError: http.StatusForbidden, Type: "Sender", Code: "AWS.SimpleQueueService.PurgeQueueInProgress", Message: "Only one PurgeQueue operation on a queue is allowed every 60 seconds."}, + // OverLimit — превышен лимит сообщений в очереди (OOM-защита) + "OverLimit": {HttpError: http.StatusBadRequest, Type: "OverLimit", Code: "OverLimit", Message: "The queue contains the maximum number of messages."}, } } diff --git a/tests/__pycache__/load_test.cpython-312.pyc b/tests/__pycache__/load_test.cpython-312.pyc new file mode 100644 index 0000000..bd9c3ea Binary files /dev/null and b/tests/__pycache__/load_test.cpython-312.pyc differ