From ebcb20475ad780e8d4099118cd5f36568fa243ad Mon Sep 17 00:00:00 2001 From: Naeel Date: Fri, 10 Apr 2026 19:49:55 +0300 Subject: [PATCH] perf: optimize receive long polling and finalize formatting cleanup --- app/gosqs/delete_message_batch.go | 5 +-- app/gosqs/list_queues.go | 3 +- app/gosqs/receive_message.go | 53 ++++++++++++++++++++----------- 3 files changed, 40 insertions(+), 21 deletions(-) diff --git a/app/gosqs/delete_message_batch.go b/app/gosqs/delete_message_batch.go index df17c68..4217708 100644 --- a/app/gosqs/delete_message_batch.go +++ b/app/gosqs/delete_message_batch.go @@ -6,12 +6,13 @@ import ( "net/http" "strings" - "github.com/gorilla/mux" - log "github.com/sirupsen/logrus" "shared-sqs/app/interfaces" "shared-sqs/app/models" "shared-sqs/app/persistence" "shared-sqs/app/utils" + + "github.com/gorilla/mux" + log "github.com/sirupsen/logrus" ) func DeleteMessageBatchV1(req *http.Request) (int, interfaces.AbstractResponseBody) { diff --git a/app/gosqs/list_queues.go b/app/gosqs/list_queues.go index b6cfa1d..38f5962 100644 --- a/app/gosqs/list_queues.go +++ b/app/gosqs/list_queues.go @@ -7,10 +7,11 @@ import ( "net/http" "strings" - log "github.com/sirupsen/logrus" "shared-sqs/app/interfaces" "shared-sqs/app/models" "shared-sqs/app/utils" + + log "github.com/sirupsen/logrus" ) func ListQueuesV1(req *http.Request) (int, interfaces.AbstractResponseBody) { diff --git a/app/gosqs/receive_message.go b/app/gosqs/receive_message.go index f083d12..6bd03b6 100644 --- a/app/gosqs/receive_message.go +++ b/app/gosqs/receive_message.go @@ -66,33 +66,33 @@ func ReceiveMessageV1(req *http.Request) (int, interfaces.AbstractResponseBody) // Fix #4: clamp WaitTimeSeconds к AWS лимиту 0–20 waitTimeSeconds = ClampInt(waitTimeSeconds, 0, MaxReceiveMessageWaitTimeSeconds) - // Long polling: ждём появления сообщения до waitTimeSeconds*10 итераций по 100ms - loops := waitTimeSeconds * 10 - for loops > 0 { - models.SyncQueues.RLock() - _, queueFound := models.SyncQueues.Queues[key] - if !queueFound { + if waitTimeSeconds > 0 { + deadline := time.Now().Add(time.Duration(waitTimeSeconds) * time.Second) + pollTicker := time.NewTicker(100 * time.Millisecond) + defer pollTicker.Stop() + + for { + models.SyncQueues.RLock() + queue, queueFound := models.SyncQueues.Queues[key] + if !queueFound { + models.SyncQueues.RUnlock() + return utils.CreateErrorResponseV1("QueueNotFound", true) + } + messageFound := queueHasReceivableMessages(queue) models.SyncQueues.RUnlock() - return utils.CreateErrorResponseV1("QueueNotFound", true) - } - messageFound := len(models.SyncQueues.Queues[key].Messages)-numberOfHiddenMessagesInQueue(*models.SyncQueues.Queues[key]) != 0 - models.SyncQueues.RUnlock() - if !messageFound { - continueTimer := time.NewTimer(100 * time.Millisecond) + if messageFound || time.Now().After(deadline) { + break + } + select { case <-req.Context().Done(): - continueTimer.Stop() return http.StatusOK, models.ReceiveMessageResponse{ Xmlns: models.BaseXmlns, Result: models.ReceiveMessageResult{}, Metadata: models.BaseResponseMetadata, } - case <-continueTimer.C: - continueTimer.Stop() + case <-pollTicker.C: } - loops-- - } else { - break } } log.Debugf("Getting Message from Queue:%s (tenant: %s)", queueName, t.ID) @@ -171,3 +171,20 @@ func buildResultMessage(m *models.SqsMessage) *models.ResultMessage { }, } } + +func queueHasReceivableMessages(queue *models.Queue) bool { + for i := range queue.Messages { + msg := &queue.Messages[i] + if msg.ReceiptHandle != "" { + continue + } + if !msg.IsReadyForReceipt() { + continue + } + if queue.IsFIFO && queue.IsLocked(msg.GroupID) { + continue + } + return true + } + return false +}