Files
SQS-service/app/gosqs/send_message.go
T

138 lines
4.7 KiB
Go

// Изменено: 2026-04-10 — добавлена Redis persistence
// SendMessageV1 — добавляет сообщение в очередь тенанта.
// Ловушка #6: queueName извлекается как ПОСЛЕДНИЙ сегмент URL — при URL вида
// http://host/tenantID/queueName последний сегмент = queueName (правильно).
package gosqs
import (
"net/http"
"strings"
"time"
"github.com/google/uuid"
"shared-sqs/app/interfaces"
"shared-sqs/app/models"
"shared-sqs/app/persistence"
"shared-sqs/app/utils"
log "github.com/sirupsen/logrus"
"github.com/gorilla/mux"
)
func SendMessageV1(req *http.Request) (int, interfaces.AbstractResponseBody) {
requestBody := models.NewSendMessageRequest()
ok := utils.REQUEST_TRANSFORMER(requestBody, req, false)
if !ok {
log.Error("Invalid Request - SendMessageV1")
return utils.CreateErrorResponseV1("InvalidParameterValue", true)
}
t := getTenantFromContext(req)
if t == nil {
return utils.CreateErrorResponseV1("InvalidClientTokenId", true)
}
messageBody := requestBody.MessageBody
if messageBody == "" {
return utils.CreateErrorResponseV1("MissingParameter", true)
}
messageGroupID := requestBody.MessageGroupId
messageDeduplicationID := requestBody.MessageDeduplicationId
// Валидация DeduplicationID и GroupID (макс 128 chars по AWS)
if err := ValidateDeduplicationID(messageDeduplicationID); err != nil {
return utils.CreateErrorResponseV1("InvalidParameterValue", true)
}
if err := ValidateGroupID(messageGroupID); err != nil {
return utils.CreateErrorResponseV1("InvalidParameterValue", true)
}
// Валидация количества message attributes (макс 10)
if len(requestBody.MessageAttributes) > MaxMessageAttributes {
return utils.CreateErrorResponseV1("InvalidParameterValue", true)
}
queueUrl := getQueueFromPath(requestBody.QueueUrl, req.URL.String())
queueName := ""
if queueUrl == "" {
vars := mux.Vars(req)
queueName = vars["queueName"]
} else {
// Ловушка #6: берём последний сегмент — это queueName, не tenantID
uriSegments := strings.Split(queueUrl, "/")
queueName = uriSegments[len(uriSegments)-1]
}
key := tenantQueueKey(t.AccessKey, queueName)
models.SyncQueues.RLock()
queue, exists := models.SyncQueues.Queues[key]
if !exists {
models.SyncQueues.RUnlock()
return utils.CreateErrorResponseV1("QueueNotFound", true)
}
maxMessageSize := queue.MaximumMessageSize
currentMsgCount := len(queue.Messages)
queueIsFIFO := queue.IsFIFO
delaySecs := queue.DelaySeconds
models.SyncQueues.RUnlock()
if maxMessageSize > 0 && len(messageBody) > maxMessageSize {
return utils.CreateErrorResponseV1("MessageTooBig", true)
}
// Fix #11: лимит сообщений в очереди — защита от OOM
if currentMsgCount >= MaxMessagesForQueue(queueIsFIFO) {
return utils.CreateErrorResponseV1("OverLimit", true)
}
if requestBody.DelaySeconds != 0 {
delaySecs = requestBody.DelaySeconds
}
log.Debugf("Putting Message in Queue: [%s] tenant: [%s]", queueName, t.ID)
msg := models.SqsMessage{MessageBody: messageBody}
if len(requestBody.MessageAttributes) > 0 {
msg.MessageAttributes = requestBody.MessageAttributes
msg.MD5OfMessageAttributes = utils.HashAttributes(requestBody.MessageAttributes)
}
msg.MD5OfMessageBody = utils.GetMD5Hash(messageBody)
msg.Uuid = uuid.NewString()
msg.GroupID = messageGroupID
msg.DeduplicationID = messageDeduplicationID
msg.SentTime = time.Now()
msg.DelaySecs = delaySecs
models.SyncQueues.Lock()
fifoSeqNumber := ""
if models.SyncQueues.Queues[key].IsFIFO {
fifoSeqNumber = models.SyncQueues.Queues[key].NextSequenceNumber(messageGroupID)
}
if !models.SyncQueues.Queues[key].IsDuplicate(messageDeduplicationID) {
models.SyncQueues.Queues[key].Messages = append(models.SyncQueues.Queues[key].Messages, msg)
} else {
log.Debugf("Duplicate message deduplicationId [%s] in queue [%s]", messageDeduplicationID, queueName)
}
models.SyncQueues.Queues[key].InitDuplicatation(messageDeduplicationID)
// Сохраняем очередь в Redis пока держим Lock
persistence.SaveQueue(key, models.SyncQueues.Queues[key])
models.SyncQueues.Unlock()
log.Infof("%s: Queue: %s, Message: %s\n", time.Now().Format("2006-01-02 15:04:05"), queueName, msg.MessageBody)
respStruct := models.SendMessageResponse{
Xmlns: models.BaseXmlns,
Result: models.SendMessageResult{
MD5OfMessageAttributes: msg.MD5OfMessageAttributes,
MD5OfMessageBody: msg.MD5OfMessageBody,
MessageId: msg.Uuid,
SequenceNumber: fifoSeqNumber,
},
Metadata: models.BaseResponseMetadata,
}
return http.StatusOK, respStruct
}