Phase 1 (Critical): - #1 JWT auth (done in v0.1.16) - #2 Batch message size validation in send_message_batch.go - #10 RLock in GetQueueUrlV1 (data race fix) Phase 2 (AWS-compatible limits): - #3 QueueName validation: max 80 chars, [a-zA-Z0-9_-](.fifo)? - #4 WaitTimeSeconds clamped to 0-20 - #5 ReceiveMessageWaitTimeSeconds clamped to 0-20 - #6 DelaySeconds clamped to 0-900 - #7 VisibilityTimeout clamped to 0-43200 - #8 MaxNumberOfMessages clamped to 1-10 - #9 Message attributes limited to 10 per message - #15 BatchEntryId length validated (max 80) - #16 DeduplicationID length validated (max 128) - #17 GroupID length validated (max 128) Phase 3 (Per-tenant resource limits): - #11 Max messages per queue (120K standard, 20K FIFO) - #12 Global tenant limit (1000) - #3.5 HTTP request body size limit (1MB via MaxBytesReader) Phase 4 (Stability): - #14 Duplicates map cleanup (already in PeriodicTasks) - #13 FIFO group lock timeout (already in visibility timeout reset) - #18 Redis size guard: skip save if >50MB Skipped (Low, no real risk): - #19 {account} URL param (informational only, not used for access) - #20 ReceiptHandle format (self-validating UUID#UUID) New file: app/gosqs/validation.go — centralized AWS SQS limits and validators
133 lines
4.6 KiB
Go
133 lines
4.6 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
|
|
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)
|
|
|
|
if _, ok := models.SyncQueues.Queues[key]; !ok {
|
|
return utils.CreateErrorResponseV1("QueueNotFound", true)
|
|
}
|
|
|
|
if models.SyncQueues.Queues[key].MaximumMessageSize > 0 &&
|
|
len(messageBody) > models.SyncQueues.Queues[key].MaximumMessageSize {
|
|
return utils.CreateErrorResponseV1("MessageTooBig", true)
|
|
}
|
|
|
|
// Fix #11: лимит сообщений в очереди — защита от OOM
|
|
models.SyncQueues.RLock()
|
|
currentMsgCount := len(models.SyncQueues.Queues[key].Messages)
|
|
queueIsFIFO := models.SyncQueues.Queues[key].IsFIFO
|
|
models.SyncQueues.RUnlock()
|
|
if currentMsgCount >= MaxMessagesForQueue(queueIsFIFO) {
|
|
return utils.CreateErrorResponseV1("OverLimit", true)
|
|
}
|
|
|
|
delaySecs := models.SyncQueues.Queues[key].DelaySeconds
|
|
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
|
|
}
|