perf: optimize receive long polling and finalize formatting cleanup
This commit is contained in:
@@ -6,12 +6,13 @@ import (
|
|||||||
"net/http"
|
"net/http"
|
||||||
"strings"
|
"strings"
|
||||||
|
|
||||||
"github.com/gorilla/mux"
|
|
||||||
log "github.com/sirupsen/logrus"
|
|
||||||
"shared-sqs/app/interfaces"
|
"shared-sqs/app/interfaces"
|
||||||
"shared-sqs/app/models"
|
"shared-sqs/app/models"
|
||||||
"shared-sqs/app/persistence"
|
"shared-sqs/app/persistence"
|
||||||
"shared-sqs/app/utils"
|
"shared-sqs/app/utils"
|
||||||
|
|
||||||
|
"github.com/gorilla/mux"
|
||||||
|
log "github.com/sirupsen/logrus"
|
||||||
)
|
)
|
||||||
|
|
||||||
func DeleteMessageBatchV1(req *http.Request) (int, interfaces.AbstractResponseBody) {
|
func DeleteMessageBatchV1(req *http.Request) (int, interfaces.AbstractResponseBody) {
|
||||||
|
|||||||
@@ -7,10 +7,11 @@ import (
|
|||||||
"net/http"
|
"net/http"
|
||||||
"strings"
|
"strings"
|
||||||
|
|
||||||
log "github.com/sirupsen/logrus"
|
|
||||||
"shared-sqs/app/interfaces"
|
"shared-sqs/app/interfaces"
|
||||||
"shared-sqs/app/models"
|
"shared-sqs/app/models"
|
||||||
"shared-sqs/app/utils"
|
"shared-sqs/app/utils"
|
||||||
|
|
||||||
|
log "github.com/sirupsen/logrus"
|
||||||
)
|
)
|
||||||
|
|
||||||
func ListQueuesV1(req *http.Request) (int, interfaces.AbstractResponseBody) {
|
func ListQueuesV1(req *http.Request) (int, interfaces.AbstractResponseBody) {
|
||||||
|
|||||||
@@ -66,33 +66,33 @@ func ReceiveMessageV1(req *http.Request) (int, interfaces.AbstractResponseBody)
|
|||||||
// Fix #4: clamp WaitTimeSeconds к AWS лимиту 0–20
|
// Fix #4: clamp WaitTimeSeconds к AWS лимиту 0–20
|
||||||
waitTimeSeconds = ClampInt(waitTimeSeconds, 0, MaxReceiveMessageWaitTimeSeconds)
|
waitTimeSeconds = ClampInt(waitTimeSeconds, 0, MaxReceiveMessageWaitTimeSeconds)
|
||||||
|
|
||||||
// Long polling: ждём появления сообщения до waitTimeSeconds*10 итераций по 100ms
|
if waitTimeSeconds > 0 {
|
||||||
loops := waitTimeSeconds * 10
|
deadline := time.Now().Add(time.Duration(waitTimeSeconds) * time.Second)
|
||||||
for loops > 0 {
|
pollTicker := time.NewTicker(100 * time.Millisecond)
|
||||||
models.SyncQueues.RLock()
|
defer pollTicker.Stop()
|
||||||
_, queueFound := models.SyncQueues.Queues[key]
|
|
||||||
if !queueFound {
|
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()
|
models.SyncQueues.RUnlock()
|
||||||
return utils.CreateErrorResponseV1("QueueNotFound", true)
|
if messageFound || time.Now().After(deadline) {
|
||||||
}
|
break
|
||||||
messageFound := len(models.SyncQueues.Queues[key].Messages)-numberOfHiddenMessagesInQueue(*models.SyncQueues.Queues[key]) != 0
|
}
|
||||||
models.SyncQueues.RUnlock()
|
|
||||||
if !messageFound {
|
|
||||||
continueTimer := time.NewTimer(100 * time.Millisecond)
|
|
||||||
select {
|
select {
|
||||||
case <-req.Context().Done():
|
case <-req.Context().Done():
|
||||||
continueTimer.Stop()
|
|
||||||
return http.StatusOK, models.ReceiveMessageResponse{
|
return http.StatusOK, models.ReceiveMessageResponse{
|
||||||
Xmlns: models.BaseXmlns,
|
Xmlns: models.BaseXmlns,
|
||||||
Result: models.ReceiveMessageResult{},
|
Result: models.ReceiveMessageResult{},
|
||||||
Metadata: models.BaseResponseMetadata,
|
Metadata: models.BaseResponseMetadata,
|
||||||
}
|
}
|
||||||
case <-continueTimer.C:
|
case <-pollTicker.C:
|
||||||
continueTimer.Stop()
|
|
||||||
}
|
}
|
||||||
loops--
|
|
||||||
} else {
|
|
||||||
break
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
log.Debugf("Getting Message from Queue:%s (tenant: %s)", queueName, t.ID)
|
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
|
||||||
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user