169 lines
4.9 KiB
Go
169 lines
4.9 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
|
|
}
|
|
|
|
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()
|
|
}
|
|
|
|
// Long polling: ждём появления сообщения до waitTimeSeconds*10 итераций по 100ms
|
|
loops := waitTimeSeconds * 10
|
|
for loops > 0 {
|
|
models.SyncQueues.RLock()
|
|
_, queueFound := models.SyncQueues.Queues[key]
|
|
if !queueFound {
|
|
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)
|
|
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()
|
|
}
|
|
loops--
|
|
} else {
|
|
break
|
|
}
|
|
}
|
|
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)),
|
|
},
|
|
}
|
|
}
|