191 lines
5.7 KiB
Go
191 lines
5.7 KiB
Go
// Изменено: 2026-04-09
|
||
// ReceiveMessageV1 — получает сообщения из очереди тенанта с поддержкой long polling.
|
||
// Ловушка #4: long polling держит соединение до 20 сек — не прерываем принудительно.
|
||
package gosqs
|
||
|
||
import (
|
||
"fmt"
|
||
"net/http"
|
||
"strings"
|
||
"time"
|
||
|
||
"github.com/google/uuid"
|
||
|
||
"shared-sqs/app/interfaces"
|
||
"shared-sqs/app/models"
|
||
"shared-sqs/app/utils"
|
||
|
||
"github.com/gorilla/mux"
|
||
log "github.com/sirupsen/logrus"
|
||
)
|
||
|
||
func ReceiveMessageV1(req *http.Request) (int, interfaces.AbstractResponseBody) {
|
||
requestBody := models.NewReceiveMessageRequest()
|
||
ok := utils.REQUEST_TRANSFORMER(requestBody, req, false)
|
||
if !ok {
|
||
log.Error("Invalid Request - ReceiveMessageV1")
|
||
return utils.CreateErrorResponseV1("InvalidParameterValue", true)
|
||
}
|
||
|
||
t := getTenantFromContext(req)
|
||
if t == nil {
|
||
return utils.CreateErrorResponseV1("InvalidClientTokenId", true)
|
||
}
|
||
|
||
maxNumberOfMessages := requestBody.MaxNumberOfMessages
|
||
if maxNumberOfMessages == 0 {
|
||
maxNumberOfMessages = 1
|
||
}
|
||
// Fix #8: clamp MaxNumberOfMessages к AWS лимиту 1–10
|
||
maxNumberOfMessages = ClampInt(maxNumberOfMessages, MinNumberOfMessagesLimit, MaxNumberOfMessagesLimit)
|
||
|
||
queueName := ""
|
||
if requestBody.QueueUrl == "" {
|
||
vars := mux.Vars(req)
|
||
queueName = vars["queueName"]
|
||
} else {
|
||
uriSegments := strings.Split(requestBody.QueueUrl, "/")
|
||
queueName = uriSegments[len(uriSegments)-1]
|
||
}
|
||
|
||
key := tenantQueueKey(t.AccessKey, queueName)
|
||
|
||
if _, ok := models.SyncQueues.Queues[key]; !ok {
|
||
return utils.CreateErrorResponseV1("QueueNotFound", true)
|
||
}
|
||
|
||
var messages []*models.ResultMessage
|
||
respStruct := models.ReceiveMessageResponse{}
|
||
|
||
waitTimeSeconds := requestBody.WaitTimeSeconds
|
||
if waitTimeSeconds == 0 {
|
||
models.SyncQueues.RLock()
|
||
waitTimeSeconds = models.SyncQueues.Queues[key].ReceiveMessageWaitTimeSeconds
|
||
models.SyncQueues.RUnlock()
|
||
}
|
||
// Fix #4: clamp WaitTimeSeconds к AWS лимиту 0–20
|
||
waitTimeSeconds = ClampInt(waitTimeSeconds, 0, MaxReceiveMessageWaitTimeSeconds)
|
||
|
||
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()
|
||
if messageFound || time.Now().After(deadline) {
|
||
break
|
||
}
|
||
|
||
select {
|
||
case <-req.Context().Done():
|
||
return http.StatusOK, models.ReceiveMessageResponse{
|
||
Xmlns: models.BaseXmlns,
|
||
Result: models.ReceiveMessageResult{},
|
||
Metadata: models.BaseResponseMetadata,
|
||
}
|
||
case <-pollTicker.C:
|
||
}
|
||
}
|
||
}
|
||
log.Debugf("Getting Message from Queue:%s (tenant: %s)", queueName, t.ID)
|
||
|
||
models.SyncQueues.Lock()
|
||
defer models.SyncQueues.Unlock()
|
||
|
||
if len(models.SyncQueues.Queues[key].Messages) > 0 {
|
||
numMsg := 0
|
||
messages = make([]*models.ResultMessage, 0)
|
||
for i := range models.SyncQueues.Queues[key].Messages {
|
||
if numMsg >= maxNumberOfMessages {
|
||
break
|
||
}
|
||
|
||
if models.SyncQueues.Queues[key].Messages[i].ReceiptHandle != "" {
|
||
continue
|
||
}
|
||
|
||
msg := &models.SyncQueues.Queues[key].Messages[i]
|
||
if !msg.IsReadyForReceipt() {
|
||
continue
|
||
}
|
||
|
||
if models.SyncQueues.Queues[key].IsFIFO {
|
||
if models.SyncQueues.Queues[key].IsLocked(msg.GroupID) {
|
||
continue
|
||
}
|
||
models.SyncQueues.Queues[key].LockGroup(msg.GroupID)
|
||
}
|
||
|
||
randomId := uuid.NewString()
|
||
msg.ReceiptHandle = msg.Uuid + "#" + randomId
|
||
msg.ReceiptTime = time.Now().UTC()
|
||
|
||
if requestBody.VisibilityTimeout != 0 {
|
||
msg.VisibilityTimeout = time.Now().Add(time.Duration(requestBody.VisibilityTimeout) * time.Second)
|
||
} else {
|
||
msg.VisibilityTimeout = time.Now().Add(time.Duration(models.SyncQueues.Queues[key].VisibilityTimeout) * time.Second)
|
||
}
|
||
|
||
messages = append(messages, buildResultMessage(msg))
|
||
numMsg++
|
||
}
|
||
|
||
respStruct = models.ReceiveMessageResponse{
|
||
"http://queue.amazonaws.com/doc/2012-11-05/",
|
||
models.ReceiveMessageResult{Messages: messages},
|
||
models.ResponseMetadata{RequestId: "00000000-0000-0000-0000-000000000000"},
|
||
}
|
||
} else {
|
||
log.Warning("No messages in Queue:", queueName)
|
||
respStruct = models.ReceiveMessageResponse{
|
||
Xmlns: "http://queue.amazonaws.com/doc/2012-11-05/",
|
||
Result: models.ReceiveMessageResult{},
|
||
Metadata: models.ResponseMetadata{RequestId: "00000000-0000-0000-0000-000000000000"},
|
||
}
|
||
}
|
||
|
||
return http.StatusOK, respStruct
|
||
}
|
||
|
||
func buildResultMessage(m *models.SqsMessage) *models.ResultMessage {
|
||
return &models.ResultMessage{
|
||
MessageId: m.Uuid,
|
||
Body: m.MessageBody,
|
||
ReceiptHandle: m.ReceiptHandle,
|
||
MD5OfBody: utils.GetMD5Hash(m.MessageBody),
|
||
MD5OfMessageAttributes: m.MD5OfMessageAttributes,
|
||
MessageAttributes: m.MessageAttributes,
|
||
Attributes: map[string]string{
|
||
"ApproximateFirstReceiveTimestamp": fmt.Sprintf("%d", m.ReceiptTime.UnixNano()/int64(time.Millisecond)),
|
||
"SenderId": models.CurrentEnvironment.AccountID,
|
||
"ApproximateReceiveCount": fmt.Sprintf("%d", m.NumberOfReceives+1),
|
||
"SentTimestamp": fmt.Sprintf("%d", time.Now().UTC().UnixNano()/int64(time.Millisecond)),
|
||
},
|
||
}
|
||
}
|
||
|
||
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
|
||
}
|