shared-sqs: Этап 4 — изоляция очередей по тенанту

This commit is contained in:
Naeel
2026-04-09 12:53:10 +03:00
parent f4352a17b1
commit 08053ca8e5
16 changed files with 1243 additions and 882 deletions
+114
View File
@@ -0,0 +1,114 @@
// Изменено: 2026-04-09
// Auth middleware для shared-sqs: извлекает AccessKeyId из AWS Authorization header
// и помещает найденного тенанта в context запроса.
package auth
import (
"context"
"encoding/xml"
"net/http"
"strings"
"shared-sqs/app/tenant"
)
// TenantContextKey — ключ для хранения тенанта в request context.
// Тип contextKey предотвращает конфликты с другими пакетами.
type contextKey string
const TenantContextKey contextKey = "tenant"
// AuthMiddleware — middleware: ищет тенанта по AccessKeyId из AWS Authorization header.
// Пропускает /health и /admin/** без tenant-аутентификации.
// Ловушка #4: не ставим короткий таймаут — ReceiveMessage с long polling держит соединение до 20 сек.
func AuthMiddleware(store *tenant.TenantStore) func(http.Handler) http.Handler {
return func(next http.Handler) http.Handler {
return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
// /health — без auth
if r.URL.Path == "/health" {
next.ServeHTTP(w, r)
return
}
// /admin/** — отдельная auth (bearer token, см. admin_handlers.go)
if strings.HasPrefix(r.URL.Path, "/admin/") {
next.ServeHTTP(w, r)
return
}
accessKeyID := extractAccessKeyID(r)
if accessKeyID == "" {
writeSQSAuthError(w, "MissingAuthenticationToken", "Request must contain either AccessKeyId or X-Amz-Credential")
return
}
t, ok := store.GetByAccessKey(accessKeyID)
if !ok || !t.Active {
writeSQSAuthError(w, "InvalidClientTokenId", "The security token included in the request is invalid")
return
}
ctx := context.WithValue(r.Context(), TenantContextKey, t)
next.ServeHTTP(w, r.WithContext(ctx))
})
}
}
// extractAccessKeyID — извлекает AWS AccessKeyId из запроса.
// Поддерживает оба варианта: Authorization header (Signature V4) и X-Amz-Credential query param (presigned URLs).
// Ловушка #3: AWS CLI ВСЕГДА отправляет Signature V4 — нужно парсить, даже не проверяя подпись.
// Ловушка #5: X-Amz-Security-Token (STS) — игнорируем.
func extractAccessKeyID(r *http.Request) string {
// Вариант 1: Authorization header
// Формат: "AWS4-HMAC-SHA256 Credential={AccessKeyId}/{date}/{region}/sqs/aws4_request, ..."
auth := r.Header.Get("Authorization")
if strings.HasPrefix(auth, "AWS4-HMAC-SHA256") {
idx := strings.Index(auth, "Credential=")
if idx >= 0 {
rest := auth[idx+len("Credential="):]
slashIdx := strings.Index(rest, "/")
if slashIdx > 0 {
return rest[:slashIdx]
}
}
}
// Вариант 2: Query parameter (presigned URLs)
// Формат: X-Amz-Credential={AccessKeyId}/{date}/{region}/sqs/aws4_request
if cred := r.URL.Query().Get("X-Amz-Credential"); cred != "" {
parts := strings.SplitN(cred, "/", 2)
if len(parts) > 0 && parts[0] != "" {
return parts[0]
}
}
return ""
}
// sqsAuthError — AWS-совместимый XML ответ об ошибке аутентификации.
type sqsAuthError struct {
XMLName xml.Name `xml:"ErrorResponse"`
Error sqsErrorBody `xml:"Error"`
RequestID string `xml:"RequestId"`
}
type sqsErrorBody struct {
Type string `xml:"Type"`
Code string `xml:"Code"`
Message string `xml:"Message"`
}
// writeSQSAuthError — отвечает AWS-совместимым XML с кодом 403.
func writeSQSAuthError(w http.ResponseWriter, code, message string) {
w.Header().Set("Content-Type", "application/xml")
w.WriteHeader(http.StatusForbidden)
resp := sqsAuthError{
Error: sqsErrorBody{
Type: "Sender",
Code: code,
Message: message,
},
RequestID: "00000000-0000-0000-0000-000000000000",
}
data, _ := xml.Marshal(resp)
w.Write(data)
}
@@ -1,81 +1,87 @@
// Изменено: 2026-04-09
// ChangeMessageVisibilityV1 — меняет visibility timeout сообщения в очереди тенанта.
package gosqs package gosqs
import ( import (
"net/http" "net/http"
"strings" "strings"
"time" "time"
"shared-sqs/app/interfaces" "shared-sqs/app/interfaces"
"shared-sqs/app/models" "shared-sqs/app/models"
"shared-sqs/app/utils" "shared-sqs/app/utils"
"github.com/gorilla/mux" "github.com/gorilla/mux"
log "github.com/sirupsen/logrus" log "github.com/sirupsen/logrus"
) )
func ChangeMessageVisibilityV1(req *http.Request) (int, interfaces.AbstractResponseBody) { func ChangeMessageVisibilityV1(req *http.Request) (int, interfaces.AbstractResponseBody) {
requestBody := models.NewChangeMessageVisibilityRequest() requestBody := models.NewChangeMessageVisibilityRequest()
ok := utils.REQUEST_TRANSFORMER(requestBody, req, false) ok := utils.REQUEST_TRANSFORMER(requestBody, req, false)
if !ok { if !ok {
log.Error("Invalid Request - ChangeMessageVisibilityV1") log.Error("Invalid Request - ChangeMessageVisibilityV1")
return utils.CreateErrorResponseV1("InvalidParameterValue", true) return utils.CreateErrorResponseV1("InvalidParameterValue", true)
} }
vars := mux.Vars(req) t := getTenantFromContext(req)
if t == nil {
queueUrl := requestBody.QueueUrl return utils.CreateErrorResponseV1("InvalidClientTokenId", true)
queueName := "" }
if queueUrl == "" {
queueName = vars["queueName"] vars := mux.Vars(req)
} else { queueUrl := requestBody.QueueUrl
uriSegments := strings.Split(queueUrl, "/") queueName := ""
queueName = uriSegments[len(uriSegments)-1] if queueUrl == "" {
} queueName = vars["queueName"]
} else {
receiptHandle := requestBody.ReceiptHandle uriSegments := strings.Split(queueUrl, "/")
queueName = uriSegments[len(uriSegments)-1]
visibilityTimeout := requestBody.VisibilityTimeout }
if visibilityTimeout > 43200 {
return utils.CreateErrorResponseV1("ValidationError", true) key := tenantQueueKey(t.AccessKey, queueName)
} receiptHandle := requestBody.ReceiptHandle
visibilityTimeout := requestBody.VisibilityTimeout
if _, ok := models.SyncQueues.Queues[queueName]; !ok {
return utils.CreateErrorResponseV1("QueueNotFound", true) if visibilityTimeout > 43200 {
} return utils.CreateErrorResponseV1("ValidationError", true)
}
models.SyncQueues.Lock()
messageFound := false if _, ok := models.SyncQueues.Queues[key]; !ok {
for i := 0; i < len(models.SyncQueues.Queues[queueName].Messages); i++ { return utils.CreateErrorResponseV1("QueueNotFound", true)
queue := models.SyncQueues.Queues[queueName] }
msgs := queue.Messages
if msgs[i].ReceiptHandle == receiptHandle { models.SyncQueues.Lock()
timeout := models.SyncQueues.Queues[queueName].VisibilityTimeout messageFound := false
if visibilityTimeout == 0 { for i := 0; i < len(models.SyncQueues.Queues[key].Messages); i++ {
msgs[i].ReceiptTime = time.Now().UTC() queue := models.SyncQueues.Queues[key]
msgs[i].ReceiptHandle = "" msgs := queue.Messages
msgs[i].VisibilityTimeout = time.Now().Add(time.Duration(timeout) * time.Second) if msgs[i].ReceiptHandle == receiptHandle {
msgs[i].Retry++ timeout := models.SyncQueues.Queues[key].VisibilityTimeout
if queue.MaxReceiveCount > 0 && if visibilityTimeout == 0 {
queue.DeadLetterQueue != nil && msgs[i].ReceiptTime = time.Now().UTC()
msgs[i].Retry >= queue.MaxReceiveCount { msgs[i].ReceiptHandle = ""
queue.DeadLetterQueue.Messages = append(queue.DeadLetterQueue.Messages, msgs[i]) msgs[i].VisibilityTimeout = time.Now().Add(time.Duration(timeout) * time.Second)
queue.Messages = append(queue.Messages[:i], queue.Messages[i+1:]...) msgs[i].Retry++
} if queue.MaxReceiveCount > 0 &&
} else { queue.DeadLetterQueue != nil &&
msgs[i].VisibilityTimeout = time.Now().Add(time.Duration(visibilityTimeout) * time.Second) msgs[i].Retry >= queue.MaxReceiveCount {
} queue.DeadLetterQueue.Messages = append(queue.DeadLetterQueue.Messages, msgs[i])
messageFound = true queue.Messages = append(queue.Messages[:i], queue.Messages[i+1:]...)
break }
} } else {
} msgs[i].VisibilityTimeout = time.Now().Add(time.Duration(visibilityTimeout) * time.Second)
models.SyncQueues.Unlock() }
if !messageFound { messageFound = true
return utils.CreateErrorResponseV1("MessageNotInFlight", true) break
} }
}
respStruct := models.ChangeMessageVisibilityResult{ models.SyncQueues.Unlock()
Xmlns: models.BaseXmlns, if !messageFound {
Metadata: models.BaseResponseMetadata, return utils.CreateErrorResponseV1("MessageNotInFlight", true)
} }
return http.StatusOK, &respStruct respStruct := models.ChangeMessageVisibilityResult{
Xmlns: models.BaseXmlns,
Metadata: models.BaseResponseMetadata,
}
return http.StatusOK, &respStruct
} }
+57 -46
View File
@@ -1,54 +1,65 @@
// Изменено: 2026-04-09
// CreateQueueV1 — создаёт очередь для тенанта из request context.
// Ключ в SyncQueues: "{tenantAccessKey}:{queueName}" для изоляции между тенантами.
package gosqs package gosqs
import ( import (
"net/http" "net/http"
"time" "time"
"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" log "github.com/sirupsen/logrus"
) )
func CreateQueueV1(req *http.Request) (int, interfaces.AbstractResponseBody) { func CreateQueueV1(req *http.Request) (int, interfaces.AbstractResponseBody) {
requestBody := models.NewCreateQueueRequest() requestBody := models.NewCreateQueueRequest()
ok := utils.REQUEST_TRANSFORMER(requestBody, req, false) ok := utils.REQUEST_TRANSFORMER(requestBody, req, false)
if !ok { if !ok {
log.Error("Invalid Request - CreateQueueV1") log.Error("Invalid Request - CreateQueueV1")
return utils.CreateErrorResponseV1("InvalidParameterValue", true) return utils.CreateErrorResponseV1("InvalidParameterValue", true)
} }
queueName := requestBody.QueueName
t := getTenantFromContext(req)
queueUrl := "http://" + models.CurrentEnvironment.Host + ":" + models.CurrentEnvironment.Port + if t == nil {
"/" + models.CurrentEnvironment.AccountID + "/" + queueName return utils.CreateErrorResponseV1("InvalidClientTokenId", true)
if models.CurrentEnvironment.Region != "" { }
queueUrl = "http://" + models.CurrentEnvironment.Region + "." + models.CurrentEnvironment.Host + ":" +
models.CurrentEnvironment.Port + "/" + models.CurrentEnvironment.AccountID + "/" + queueName // Ловушка #8: передаём queueName (не key) в HasFIFOQueueName — иначе .fifo не определится
} queueName := requestBody.QueueName
queueArn := "arn:aws:sqs:" + models.CurrentEnvironment.Region + ":" + models.CurrentEnvironment.AccountID + ":" + queueName key := tenantQueueKey(t.AccessKey, queueName)
queueUrl := tenantQueueURL(t, queueName)
if _, ok := models.SyncQueues.Queues[queueName]; !ok { queueArn := tenantQueueARN(t, queueName)
log.Infof("Creating Queue: %s", queueName)
queue := &models.Queue{ models.SyncQueues.Lock()
Name: queueName, if _, exists := models.SyncQueues.Queues[key]; !exists {
URL: queueUrl, // Проверка лимита очередей тенанта
Arn: queueArn, if t.MaxQueues > 0 && countTenantQueues(t.AccessKey) >= t.MaxQueues {
IsFIFO: utils.HasFIFOQueueName(queueName), models.SyncQueues.Unlock()
EnableDuplicates: models.CurrentEnvironment.EnableDuplicates, return utils.CreateErrorResponseV1("LimitExceeded", true)
Duplicates: make(map[string]time.Time), }
} log.Infof("Creating Queue: %s (tenant: %s)", queueName, t.ID)
if err := setQueueAttributesV1(queue, requestBody.Attributes); err != nil { queue := &models.Queue{
return utils.CreateErrorResponseV1(err.Error(), true) Name: queueName,
} URL: queueUrl,
models.SyncQueues.Lock() Arn: queueArn,
models.SyncQueues.Queues[queueName] = queue IsFIFO: utils.HasFIFOQueueName(queueName),
models.SyncQueues.Unlock() EnableDuplicates: models.CurrentEnvironment.EnableDuplicates,
} Duplicates: make(map[string]time.Time),
}
respStruct := models.CreateQueueResponse{ if err := setQueueAttributesV1(queue, requestBody.Attributes); err != nil {
Xmlns: models.BaseXmlns, models.SyncQueues.Unlock()
Result: models.CreateQueueResult{QueueUrl: queueUrl}, return utils.CreateErrorResponseV1(err.Error(), true)
Metadata: models.BaseResponseMetadata, }
} models.SyncQueues.Queues[key] = queue
return http.StatusOK, respStruct }
models.SyncQueues.Unlock()
respStruct := models.CreateQueueResponse{
Xmlns: models.BaseXmlns,
Result: models.CreateQueueResult{QueueUrl: queueUrl},
Metadata: models.BaseResponseMetadata,
}
return http.StatusOK, respStruct
} }
+56 -57
View File
@@ -1,65 +1,64 @@
// Изменено: 2026-04-09
// DeleteMessageV1 — удаляет сообщение из очереди тенанта по ReceiptHandle.
package gosqs package gosqs
import ( import (
"net/http" "net/http"
"strings" "strings"
"shared-sqs/app/interfaces" "shared-sqs/app/interfaces"
"shared-sqs/app/models" "shared-sqs/app/models"
"shared-sqs/app/utils" "shared-sqs/app/utils"
"github.com/gorilla/mux" "github.com/gorilla/mux"
log "github.com/sirupsen/logrus" log "github.com/sirupsen/logrus"
) )
func DeleteMessageV1(req *http.Request) (int, interfaces.AbstractResponseBody) { func DeleteMessageV1(req *http.Request) (int, interfaces.AbstractResponseBody) {
requestBody := models.NewDeleteMessageRequest() requestBody := models.NewDeleteMessageRequest()
ok := utils.REQUEST_TRANSFORMER(requestBody, req, false) ok := utils.REQUEST_TRANSFORMER(requestBody, req, false)
if !ok { if !ok {
log.Error("Invalid Request - DeleteMessageV1") log.Error("Invalid Request - DeleteMessageV1")
return utils.CreateErrorResponseV1("InvalidParameterValue", true) return utils.CreateErrorResponseV1("InvalidParameterValue", true)
} }
// Retrieve FormValues required t := getTenantFromContext(req)
receiptHandle := requestBody.ReceiptHandle if t == nil {
return utils.CreateErrorResponseV1("InvalidClientTokenId", true)
// Retrieve FormValues required }
queueUrl := requestBody.QueueUrl
queueName := "" receiptHandle := requestBody.ReceiptHandle
if queueUrl == "" { queueUrl := requestBody.QueueUrl
vars := mux.Vars(req) queueName := ""
queueName = vars["queueName"] if queueUrl == "" {
} else { vars := mux.Vars(req)
uriSegments := strings.Split(queueUrl, "/") queueName = vars["queueName"]
queueName = uriSegments[len(uriSegments)-1] } else {
} uriSegments := strings.Split(queueUrl, "/")
queueName = uriSegments[len(uriSegments)-1]
log.Info("Deleting Message, Queue:", queueName, ", ReceiptHandle:", receiptHandle) }
// Find queue/message with the receipt handle and delete key := tenantQueueKey(t.AccessKey, queueName)
models.SyncQueues.Lock() log.Info("Deleting Message, Queue:", queueName, ", ReceiptHandle:", receiptHandle)
defer models.SyncQueues.Unlock()
if _, ok := models.SyncQueues.Queues[queueName]; ok { models.SyncQueues.Lock()
for i, msg := range models.SyncQueues.Queues[queueName].Messages { defer models.SyncQueues.Unlock()
if msg.ReceiptHandle == receiptHandle { if _, ok := models.SyncQueues.Queues[key]; ok {
// Unlock messages for the group for i, msg := range models.SyncQueues.Queues[key].Messages {
log.Debugf("FIFO Queue %s unlocking group %s:", queueName, msg.GroupID) if msg.ReceiptHandle == receiptHandle {
models.SyncQueues.Queues[queueName].UnlockGroup(msg.GroupID) models.SyncQueues.Queues[key].UnlockGroup(msg.GroupID)
//Delete message from Q models.SyncQueues.Queues[key].Messages = append(models.SyncQueues.Queues[key].Messages[:i], models.SyncQueues.Queues[key].Messages[i+1:]...)
models.SyncQueues.Queues[queueName].Messages = append(models.SyncQueues.Queues[queueName].Messages[:i], models.SyncQueues.Queues[queueName].Messages[i+1:]...) delete(models.SyncQueues.Queues[key].Duplicates, msg.DeduplicationID)
delete(models.SyncQueues.Queues[queueName].Duplicates, msg.DeduplicationID) respStruct := models.DeleteMessageResponse{
Xmlns: models.BaseXmlns,
// Create, encode/xml and send response Metadata: models.BaseResponseMetadata,
respStruct := models.DeleteMessageResponse{ }
Xmlns: models.BaseXmlns, return 200, &respStruct
Metadata: models.BaseResponseMetadata, }
} }
return 200, &respStruct log.Warning("Receipt Handle not found")
} } else {
} log.Warning("Queue not found")
log.Warning("Receipt Handle not found") }
} else {
log.Warning("Queue not found") return utils.CreateErrorResponseV1("MessageDoesNotExist", true)
}
return utils.CreateErrorResponseV1("MessageDoesNotExist", true)
} }
+93 -92
View File
@@ -1,118 +1,119 @@
// Изменено: 2026-04-09
// DeleteMessageBatchV1 — пакетное удаление сообщений из очереди тенанта.
package gosqs package gosqs
import ( import (
"net/http" "net/http"
"strings" "strings"
"shared-sqs/app/interfaces" "shared-sqs/app/interfaces"
"shared-sqs/app/models" "shared-sqs/app/models"
"shared-sqs/app/utils" "shared-sqs/app/utils"
"github.com/gorilla/mux" "github.com/gorilla/mux"
log "github.com/sirupsen/logrus" log "github.com/sirupsen/logrus"
) )
func DeleteMessageBatchV1(req *http.Request) (int, interfaces.AbstractResponseBody) { func DeleteMessageBatchV1(req *http.Request) (int, interfaces.AbstractResponseBody) {
requestBody := models.NewDeleteMessageBatchRequest() requestBody := models.NewDeleteMessageBatchRequest()
ok := utils.REQUEST_TRANSFORMER(requestBody, req, false) ok := utils.REQUEST_TRANSFORMER(requestBody, req, false)
if !ok { if !ok {
log.Error("Invalid Request - DeleteMessageBatchV1") log.Error("Invalid Request - DeleteMessageBatchV1")
return utils.CreateErrorResponseV1("InvalidParameterValue", true) return utils.CreateErrorResponseV1("InvalidParameterValue", true)
} }
queueUrl := requestBody.QueueUrl t := getTenantFromContext(req)
if t == nil {
return utils.CreateErrorResponseV1("InvalidClientTokenId", true)
}
queueName := "" queueUrl := requestBody.QueueUrl
if queueUrl == "" { queueName := ""
vars := mux.Vars(req) if queueUrl == "" {
queueName = vars["queueName"] vars := mux.Vars(req)
} else { queueName = vars["queueName"]
uriSegments := strings.Split(queueUrl, "/") } else {
queueName = uriSegments[len(uriSegments)-1] uriSegments := strings.Split(queueUrl, "/")
} queueName = uriSegments[len(uriSegments)-1]
}
if _, ok := models.SyncQueues.Queues[queueName]; !ok { key := tenantQueueKey(t.AccessKey, queueName)
return utils.CreateErrorResponseV1("QueueNotFound", true)
}
if len(requestBody.Entries) == 0 { if _, ok := models.SyncQueues.Queues[key]; !ok {
return utils.CreateErrorResponseV1("EmptyBatchRequest", true) return utils.CreateErrorResponseV1("QueueNotFound", true)
} }
if len(requestBody.Entries) > 10 { if len(requestBody.Entries) == 0 {
return utils.CreateErrorResponseV1("TooManyEntriesInBatchRequest", true) return utils.CreateErrorResponseV1("EmptyBatchRequest", true)
} }
ids := map[string]bool{} if len(requestBody.Entries) > 10 {
for _, v := range requestBody.Entries { return utils.CreateErrorResponseV1("TooManyEntriesInBatchRequest", true)
if _, found := ids[v.Id]; found { }
return utils.CreateErrorResponseV1("BatchEntryIdsNotDistinct", true)
}
ids[v.Id] = true
}
models.SyncQueues.Lock() ids := map[string]bool{}
defer models.SyncQueues.Unlock() for _, v := range requestBody.Entries {
if _, found := ids[v.Id]; found {
return utils.CreateErrorResponseV1("BatchEntryIdsNotDistinct", true)
}
ids[v.Id] = true
}
// create deleteMessageMap models.SyncQueues.Lock()
deleteMessageMap := make(map[string]*deleteEntry) defer models.SyncQueues.Unlock()
for _, entry := range requestBody.Entries {
deleteMessageMap[entry.ReceiptHandle] = &deleteEntry{
Id: entry.Id,
ReceiptHandle: entry.ReceiptHandle,
Deleted: false,
}
}
deletedEntries := make([]models.DeleteMessageBatchResultEntry, 0) deleteMessageMap := make(map[string]*deleteEntry)
// create a slice to hold messages that are not deleted for _, entry := range requestBody.Entries {
remainingMessages := make([]models.SqsMessage, 0, len(models.SyncQueues.Queues[queueName].Messages)) deleteMessageMap[entry.ReceiptHandle] = &deleteEntry{
Id: entry.Id,
ReceiptHandle: entry.ReceiptHandle,
Deleted: false,
}
}
// delete message from queue deletedEntries := make([]models.DeleteMessageBatchResultEntry, 0)
for _, message := range models.SyncQueues.Queues[queueName].Messages { remainingMessages := make([]models.SqsMessage, 0, len(models.SyncQueues.Queues[key].Messages))
if deleteEntry, found := deleteMessageMap[message.ReceiptHandle]; found {
// Unlock messages for the group
log.Debugf("FIFO Queue %s unlocking group %s:", queueName, message.GroupID)
models.SyncQueues.Queues[queueName].UnlockGroup(message.GroupID)
delete(models.SyncQueues.Queues[queueName].Duplicates, message.DeduplicationID)
deleteEntry.Deleted = true
deletedEntries = append(deletedEntries, models.DeleteMessageBatchResultEntry{Id: deleteEntry.Id})
} else {
remainingMessages = append(remainingMessages, message)
}
}
// Update the queue with the remaining mesages for _, message := range models.SyncQueues.Queues[key].Messages {
models.SyncQueues.Queues[queueName].Messages = remainingMessages if de, found := deleteMessageMap[message.ReceiptHandle]; found {
log.Debugf("FIFO Queue %s unlocking group %s:", queueName, message.GroupID)
models.SyncQueues.Queues[key].UnlockGroup(message.GroupID)
delete(models.SyncQueues.Queues[key].Duplicates, message.DeduplicationID)
de.Deleted = true
deletedEntries = append(deletedEntries, models.DeleteMessageBatchResultEntry{Id: de.Id})
} else {
remainingMessages = append(remainingMessages, message)
}
}
// Process not found entries models.SyncQueues.Queues[key].Messages = remainingMessages
notFoundEntries := make([]models.BatchResultErrorEntry, 0)
for _, deleteEntry := range deleteMessageMap {
if !deleteEntry.Deleted {
notFoundEntries = append(notFoundEntries, models.BatchResultErrorEntry{
Code: "1",
Id: deleteEntry.Id,
Message: "Message not found",
SenderFault: true,
})
}
}
respStruct := models.DeleteMessageBatchResponse{ notFoundEntries := make([]models.BatchResultErrorEntry, 0)
Xmlns: models.BaseXmlns, for _, de := range deleteMessageMap {
Result: models.DeleteMessageBatchResult{ if !de.Deleted {
Successful: deletedEntries, notFoundEntries = append(notFoundEntries, models.BatchResultErrorEntry{
Failed: notFoundEntries, Code: "1",
}, Id: de.Id,
Metadata: models.BaseResponseMetadata, Message: "Message not found",
} SenderFault: true,
})
}
}
return http.StatusOK, respStruct respStruct := models.DeleteMessageBatchResponse{
Xmlns: models.BaseXmlns,
Result: models.DeleteMessageBatchResult{
Successful: deletedEntries,
Failed: notFoundEntries,
},
Metadata: models.BaseResponseMetadata,
}
return http.StatusOK, respStruct
} }
type deleteEntry struct { type deleteEntry struct {
Id string Id string
ReceiptHandle string ReceiptHandle string
Error string Error string
Deleted bool Deleted bool
} }
+35 -29
View File
@@ -1,37 +1,43 @@
// Изменено: 2026-04-09
// DeleteQueueV1 — удаляет очередь тенанта по tenant-scoped ключу.
package gosqs package gosqs
import ( import (
"net/http" "net/http"
"strings" "strings"
"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"
log "github.com/sirupsen/logrus"
) )
func DeleteQueueV1(req *http.Request) (int, interfaces.AbstractResponseBody) { func DeleteQueueV1(req *http.Request) (int, interfaces.AbstractResponseBody) {
requestBody := models.NewDeleteQueueRequest() requestBody := models.NewDeleteQueueRequest()
ok := utils.REQUEST_TRANSFORMER(requestBody, req, false) ok := utils.REQUEST_TRANSFORMER(requestBody, req, false)
if !ok { if !ok {
log.Error("Invalid Request - DeleteQueueV1") log.Error("Invalid Request - DeleteQueueV1")
return utils.CreateErrorResponseV1("InvalidParameterValue", true) return utils.CreateErrorResponseV1("InvalidParameterValue", true)
} }
uriSegments := strings.Split(requestBody.QueueUrl, "/") t := getTenantFromContext(req)
queueName := uriSegments[len(uriSegments)-1] if t == nil {
return utils.CreateErrorResponseV1("InvalidClientTokenId", true)
log.Infof("Deleting Queue: %s", queueName) }
models.SyncQueues.Lock() uriSegments := strings.Split(requestBody.QueueUrl, "/")
delete(models.SyncQueues.Queues, queueName) queueName := uriSegments[len(uriSegments)-1]
models.SyncQueues.Unlock() key := tenantQueueKey(t.AccessKey, queueName)
respStruct := models.DeleteQueueResponse{ log.Infof("Deleting Queue: %s (tenant: %s)", queueName, t.ID)
Xmlns: models.BaseXmlns,
Metadata: models.BaseResponseMetadata, models.SyncQueues.Lock()
} delete(models.SyncQueues.Queues, key)
return http.StatusOK, respStruct models.SyncQueues.Unlock()
respStruct := models.DeleteQueueResponse{
Xmlns: models.BaseXmlns,
Metadata: models.BaseResponseMetadata,
}
return http.StatusOK, respStruct
} }
+110 -127
View File
@@ -1,135 +1,118 @@
// Изменено: 2026-04-09
// GetQueueAttributesV1 — возвращает атрибуты очереди тенанта.
package gosqs package gosqs
import ( import (
"fmt" "fmt"
"net/http" "net/http"
"strconv" "strconv"
"strings" "strings"
"shared-sqs/app/models" "shared-sqs/app/interfaces"
"shared-sqs/app/utils" "shared-sqs/app/models"
"github.com/mitchellh/copystructure" "shared-sqs/app/utils"
"github.com/mitchellh/copystructure"
"shared-sqs/app/interfaces" log "github.com/sirupsen/logrus"
log "github.com/sirupsen/logrus"
) )
func GetQueueAttributesV1(req *http.Request) (int, interfaces.AbstractResponseBody) { func GetQueueAttributesV1(req *http.Request) (int, interfaces.AbstractResponseBody) {
requestBody := models.NewGetQueueAttributesRequest() requestBody := models.NewGetQueueAttributesRequest()
ok := utils.REQUEST_TRANSFORMER(requestBody, req, false) ok := utils.REQUEST_TRANSFORMER(requestBody, req, false)
if !ok { if !ok {
log.Error("Invalid Request - GetQueueAttributesV1") log.Error("Invalid Request - GetQueueAttributesV1")
return utils.CreateErrorResponseV1("InvalidParameterValue", true) return utils.CreateErrorResponseV1("InvalidParameterValue", true)
} }
if requestBody.QueueUrl == "" { if requestBody.QueueUrl == "" {
log.Error("Missing QueueUrl - GetQueueAttributesV1") log.Error("Missing QueueUrl - GetQueueAttributesV1")
return utils.CreateErrorResponseV1("InvalidParameterValue", true) return utils.CreateErrorResponseV1("InvalidParameterValue", true)
} }
requestedAttributes := func() map[string]bool { t := getTenantFromContext(req)
attrs := map[string]bool{} if t == nil {
if len(requestBody.AttributeNames) == 0 { return utils.CreateErrorResponseV1("InvalidClientTokenId", true)
return map[string]bool{"All": true} }
}
for _, attr := range requestBody.AttributeNames { requestedAttributes := func() map[string]bool {
if "All" == attr { attrs := map[string]bool{}
return map[string]bool{"All": true} if len(requestBody.AttributeNames) == 0 {
} return map[string]bool{"All": true}
attrs[attr] = true }
} for _, attr := range requestBody.AttributeNames {
return attrs if "All" == attr {
}() return map[string]bool{"All": true}
}
dupe, _ := copystructure.Copy(models.AvailableQueueAttributes) attrs[attr] = true
includedAttributes, _ := dupe.(map[string]bool) }
_, ok = requestedAttributes["All"] return attrs
if !ok { }()
for attr, _ := range includedAttributes {
_, ok := requestedAttributes[attr] dupe, _ := copystructure.Copy(models.AvailableQueueAttributes)
if !ok { includedAttributes, _ := dupe.(map[string]bool)
delete(includedAttributes, attr) _, ok = requestedAttributes["All"]
} if !ok {
} for attr := range includedAttributes {
} if _, ok := requestedAttributes[attr]; !ok {
delete(includedAttributes, attr)
uriSegments := strings.Split(requestBody.QueueUrl, "/") }
queueName := uriSegments[len(uriSegments)-1] }
}
log.Infof("Get Queue QueueAttributes: %s", queueName)
queueAttributes := make([]models.Attribute, 0, 0) uriSegments := strings.Split(requestBody.QueueUrl, "/")
queueName := uriSegments[len(uriSegments)-1]
models.SyncQueues.RLock() key := tenantQueueKey(t.AccessKey, queueName)
defer models.SyncQueues.RUnlock()
queue, ok := models.SyncQueues.Queues[queueName] log.Infof("Get Queue Attributes: %s (tenant: %s)", queueName, t.ID)
if !ok { queueAttributes := make([]models.Attribute, 0)
log.Errorf("Get Queue URL: %s queue does not exist!!!", queueName)
return utils.CreateErrorResponseV1("InvalidParameterValue", true) models.SyncQueues.RLock()
} defer models.SyncQueues.RUnlock()
queue, ok := models.SyncQueues.Queues[key]
if _, ok := includedAttributes["DelaySeconds"]; ok { if !ok {
attr := models.Attribute{Name: "DelaySeconds", Value: strconv.Itoa(queue.DelaySeconds)} log.Errorf("Get Queue Attributes: %s queue does not exist for tenant %s", queueName, t.ID)
queueAttributes = append(queueAttributes, attr) return utils.CreateErrorResponseV1("QueueNotFound", true)
} }
if _, ok := includedAttributes["MaximumMessageSize"]; ok {
attr := models.Attribute{Name: "MaximumMessageSize", Value: strconv.Itoa(queue.MaximumMessageSize)} if _, ok := includedAttributes["DelaySeconds"]; ok {
queueAttributes = append(queueAttributes, attr) queueAttributes = append(queueAttributes, models.Attribute{Name: "DelaySeconds", Value: strconv.Itoa(queue.DelaySeconds)})
} }
if _, ok := includedAttributes["MessageRetentionPeriod"]; ok { if _, ok := includedAttributes["MaximumMessageSize"]; ok {
attr := models.Attribute{Name: "MessageRetentionPeriod", Value: strconv.Itoa(queue.MessageRetentionPeriod)} queueAttributes = append(queueAttributes, models.Attribute{Name: "MaximumMessageSize", Value: strconv.Itoa(queue.MaximumMessageSize)})
queueAttributes = append(queueAttributes, attr) }
} if _, ok := includedAttributes["MessageRetentionPeriod"]; ok {
if _, ok := includedAttributes["ReceiveMessageWaitTimeSeconds"]; ok { queueAttributes = append(queueAttributes, models.Attribute{Name: "MessageRetentionPeriod", Value: strconv.Itoa(queue.MessageRetentionPeriod)})
attr := models.Attribute{Name: "ReceiveMessageWaitTimeSeconds", Value: strconv.Itoa(queue.ReceiveMessageWaitTimeSeconds)} }
queueAttributes = append(queueAttributes, attr) if _, ok := includedAttributes["ReceiveMessageWaitTimeSeconds"]; ok {
} queueAttributes = append(queueAttributes, models.Attribute{Name: "ReceiveMessageWaitTimeSeconds", Value: strconv.Itoa(queue.ReceiveMessageWaitTimeSeconds)})
if _, ok := includedAttributes["VisibilityTimeout"]; ok { }
attr := models.Attribute{Name: "VisibilityTimeout", Value: strconv.Itoa(queue.VisibilityTimeout)} if _, ok := includedAttributes["VisibilityTimeout"]; ok {
queueAttributes = append(queueAttributes, attr) queueAttributes = append(queueAttributes, models.Attribute{Name: "VisibilityTimeout", Value: strconv.Itoa(queue.VisibilityTimeout)})
} }
if _, ok := includedAttributes["ApproximateNumberOfMessages"]; ok { if _, ok := includedAttributes["ApproximateNumberOfMessages"]; ok {
attr := models.Attribute{Name: "ApproximateNumberOfMessages", Value: strconv.Itoa(len(queue.Messages))} queueAttributes = append(queueAttributes, models.Attribute{Name: "ApproximateNumberOfMessages", Value: strconv.Itoa(len(queue.Messages))})
queueAttributes = append(queueAttributes, attr) }
} if _, ok := includedAttributes["ApproximateNumberOfMessagesNotVisible"]; ok {
// TODO - implement queueAttributes = append(queueAttributes, models.Attribute{Name: "ApproximateNumberOfMessagesNotVisible", Value: strconv.Itoa(numberOfHiddenMessagesInQueue(*queue))})
//if _, ok := includedAttributes["ApproximateNumberOfMessagesDelayed"]; ok { }
// attr := models.Attribute{Name: "ApproximateNumberOfMessagesDelayed", Value: strconv.Itoa(len(queue.Messages))} if _, ok := includedAttributes["CreatedTimestamp"]; ok {
// queueAttributes = append(queueAttributes, attr) queueAttributes = append(queueAttributes, models.Attribute{Name: "CreatedTimestamp", Value: "0000000000"})
//} }
if _, ok := includedAttributes["ApproximateNumberOfMessagesNotVisible"]; ok { if _, ok := includedAttributes["LastModifiedTimestamp"]; ok {
attr := models.Attribute{Name: "ApproximateNumberOfMessagesNotVisible", Value: strconv.Itoa(numberOfHiddenMessagesInQueue(*queue))} queueAttributes = append(queueAttributes, models.Attribute{Name: "LastModifiedTimestamp", Value: "0000000000"})
queueAttributes = append(queueAttributes, attr) }
} if _, ok := includedAttributes["QueueArn"]; ok {
if _, ok := includedAttributes["CreatedTimestamp"]; ok { queueAttributes = append(queueAttributes, models.Attribute{Name: "QueueArn", Value: queue.Arn})
attr := models.Attribute{Name: "CreatedTimestamp", Value: "0000000000"} }
queueAttributes = append(queueAttributes, attr) if _, ok := includedAttributes["RedrivePolicy"]; ok && queue.DeadLetterQueue != nil {
} queueAttributes = append(queueAttributes, models.Attribute{
if _, ok := includedAttributes["LastModifiedTimestamp"]; ok { Name: "RedrivePolicy",
attr := models.Attribute{Name: "LastModifiedTimestamp", Value: "0000000000"} Value: fmt.Sprintf(`{"maxReceiveCount":"%d", "deadLetterTargetArn":"%s"}`, queue.MaxReceiveCount, queue.DeadLetterQueue.Arn),
queueAttributes = append(queueAttributes, attr) })
} }
if _, ok := includedAttributes["QueueArn"]; ok {
attr := models.Attribute{Name: "QueueArn", Value: queue.Arn} respStruct := models.GetQueueAttributesResponse{
queueAttributes = append(queueAttributes, attr) Xmlns: models.BaseXmlns,
} Result: models.GetQueueAttributesResult{Attrs: queueAttributes},
// TODO - implement Metadata: models.BaseResponseMetadata,
//if _, ok := includedAttributes["Policy"]; ok { }
// attr := models.Attribute{Name: "Policy", Value: ""} return http.StatusOK, respStruct
// queueAttributes = append(queueAttributes, attr)
//}
//if _, ok := includedAttributes["RedriveAllowPolicy"]; ok {
// attr := models.Attribute{Name: "RedriveAllowPolicy", Value: ""}
// queueAttributes = append(queueAttributes, attr)
//}
if _, ok := includedAttributes["RedrivePolicy"]; ok && queue.DeadLetterQueue != nil {
attr := models.Attribute{Name: "RedrivePolicy", Value: fmt.Sprintf(`{"maxReceiveCount":"%d", "deadLetterTargetArn":"%s"}`, queue.MaxReceiveCount, queue.DeadLetterQueue.Arn)}
queueAttributes = append(queueAttributes, attr)
}
respStruct := models.GetQueueAttributesResponse{
Xmlns: models.BaseXmlns,
Result: models.GetQueueAttributesResult{Attrs: queueAttributes},
Metadata: models.BaseResponseMetadata,
}
return http.StatusOK, respStruct
} }
+36 -28
View File
@@ -1,36 +1,44 @@
// Изменено: 2026-04-09
// GetQueueUrlV1 — возвращает URL очереди тенанта по имени.
package gosqs package gosqs
import ( import (
"net/http" "net/http"
"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" log "github.com/sirupsen/logrus"
) )
func GetQueueUrlV1(req *http.Request) (int, interfaces.AbstractResponseBody) { func GetQueueUrlV1(req *http.Request) (int, interfaces.AbstractResponseBody) {
requestBody := models.NewGetQueueUrlRequest() requestBody := models.NewGetQueueUrlRequest()
ok := utils.REQUEST_TRANSFORMER(requestBody, req, false) ok := utils.REQUEST_TRANSFORMER(requestBody, req, false)
if !ok { if !ok {
log.Error("Invalid Request - GetQueueUrlV1") log.Error("Invalid Request - GetQueueUrlV1")
return utils.CreateErrorResponseV1("InvalidParameterValue", true) return utils.CreateErrorResponseV1("InvalidParameterValue", true)
} }
queueName := requestBody.QueueName t := getTenantFromContext(req)
if _, ok := models.SyncQueues.Queues[queueName]; !ok { if t == nil {
log.Error("Get Queue URL:", queueName, ", queue does not exist!!!") return utils.CreateErrorResponseV1("InvalidClientTokenId", true)
return utils.CreateErrorResponseV1("QueueNotFound", true) }
}
queueName := requestBody.QueueName
queue := models.SyncQueues.Queues[queueName] key := tenantQueueKey(t.AccessKey, queueName)
log.Debug("Get Queue URL:", queue.Name)
if _, ok := models.SyncQueues.Queues[key]; !ok {
result := models.GetQueueUrlResult{QueueUrl: queue.URL} log.Errorf("Get Queue URL: %s, queue does not exist for tenant %s", queueName, t.ID)
respStruct := models.GetQueueUrlResponse{ return utils.CreateErrorResponseV1("QueueNotFound", true)
Xmlns: models.BaseXmlns, }
Result: result,
Metadata: models.BaseResponseMetadata, queue := models.SyncQueues.Queues[key]
} log.Debug("Get Queue URL:", queue.Name)
return http.StatusOK, respStruct
respStruct := models.GetQueueUrlResponse{
Xmlns: models.BaseXmlns,
Result: models.GetQueueUrlResult{QueueUrl: queue.URL},
Metadata: models.BaseResponseMetadata,
}
return http.StatusOK, respStruct
} }
+43 -37
View File
@@ -1,45 +1,51 @@
// Изменено: 2026-04-09
// ListQueuesV1 — возвращает только очереди текущего тенанта.
// Изоляция: фильтруем SyncQueues по префиксу "{tenantAccessKey}:".
package gosqs package gosqs
import ( import (
"net/http" "net/http"
"strings" "strings"
"shared-sqs/app/utils" "shared-sqs/app/interfaces"
"shared-sqs/app/models"
"shared-sqs/app/models" "shared-sqs/app/utils"
log "github.com/sirupsen/logrus"
"shared-sqs/app/interfaces"
log "github.com/sirupsen/logrus"
) )
// TODO - set up MaxResults, NextToken request params
//
// https://docs.aws.amazon.com/AWSSimpleQueueService/latest/APIReference/API_ListQueues.html
func ListQueuesV1(req *http.Request) (int, interfaces.AbstractResponseBody) { func ListQueuesV1(req *http.Request) (int, interfaces.AbstractResponseBody) {
requestBody := models.NewListQueuesRequest() requestBody := models.NewListQueuesRequest()
ok := utils.REQUEST_TRANSFORMER(requestBody, req, true) ok := utils.REQUEST_TRANSFORMER(requestBody, req, true)
if !ok { if !ok {
log.Error("Invalid Request - ListQueuesV1") log.Error("Invalid Request - ListQueuesV1")
return utils.CreateErrorResponseV1("InvalidParameterValue", true) return utils.CreateErrorResponseV1("InvalidParameterValue", true)
} }
log.Info("Listing Queues") t := getTenantFromContext(req)
queueUrls := make([]string, 0) if t == nil {
models.SyncQueues.Lock() return utils.CreateErrorResponseV1("InvalidClientTokenId", true)
for _, queue := range models.SyncQueues.Queues { }
if strings.HasPrefix(queue.Name, requestBody.QueueNamePrefix) {
queueUrls = append(queueUrls, queue.URL) log.Infof("Listing Queues for tenant: %s", t.ID)
} queueUrls := make([]string, 0)
} prefix := t.AccessKey + ":"
models.SyncQueues.Unlock()
models.SyncQueues.Lock()
respStruct := models.ListQueuesResponse{ for key, queue := range models.SyncQueues.Queues {
Xmlns: models.BaseXmlns, // Показываем только очереди этого тенанта
Metadata: models.BaseResponseMetadata, if strings.HasPrefix(key, prefix) {
Result: models.ListQueuesResult{ if strings.HasPrefix(queue.Name, requestBody.QueueNamePrefix) {
QueueUrls: queueUrls, queueUrls = append(queueUrls, queue.URL)
}, }
} }
}
return http.StatusOK, respStruct models.SyncQueues.Unlock()
respStruct := models.ListQueuesResponse{
Xmlns: models.BaseXmlns,
Metadata: models.BaseResponseMetadata,
Result: models.ListQueuesResult{QueueUrls: queueUrls},
}
return http.StatusOK, respStruct
} }
+41 -34
View File
@@ -1,42 +1,49 @@
// Изменено: 2026-04-09
// PurgeQueueV1 — очищает все сообщения в очереди тенанта.
package gosqs package gosqs
import ( import (
"net/http" "net/http"
"strings" "strings"
"time" "time"
"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"
log "github.com/sirupsen/logrus"
) )
func PurgeQueueV1(req *http.Request) (int, interfaces.AbstractResponseBody) { func PurgeQueueV1(req *http.Request) (int, interfaces.AbstractResponseBody) {
requestBody := models.NewPurgeQueueRequest() requestBody := models.NewPurgeQueueRequest()
ok := utils.REQUEST_TRANSFORMER(requestBody, req, false) ok := utils.REQUEST_TRANSFORMER(requestBody, req, false)
if !ok { if !ok {
log.Error("Invalid Request - PurgeQueueV1") log.Error("Invalid Request - PurgeQueueV1")
return utils.CreateErrorResponseV1("InvalidParameterValue", true) return utils.CreateErrorResponseV1("InvalidParameterValue", true)
} }
uriSegments := strings.Split(requestBody.QueueUrl, "/") t := getTenantFromContext(req)
queueName := uriSegments[len(uriSegments)-1] if t == nil {
return utils.CreateErrorResponseV1("InvalidClientTokenId", true)
models.SyncQueues.Lock() }
defer models.SyncQueues.Unlock()
if _, ok := models.SyncQueues.Queues[queueName]; !ok { uriSegments := strings.Split(requestBody.QueueUrl, "/")
log.Errorf("Purge Queue: %s, queue does not exist!!!", queueName) queueName := uriSegments[len(uriSegments)-1]
return utils.CreateErrorResponseV1("QueueNotFound", true) key := tenantQueueKey(t.AccessKey, queueName)
}
models.SyncQueues.Lock()
log.Infof("Purging Queue: %s", queueName) defer models.SyncQueues.Unlock()
models.SyncQueues.Queues[queueName].Messages = nil if _, ok := models.SyncQueues.Queues[key]; !ok {
models.SyncQueues.Queues[queueName].Duplicates = make(map[string]time.Time) log.Errorf("Purge Queue: %s, queue does not exist for tenant %s", queueName, t.ID)
return utils.CreateErrorResponseV1("QueueNotFound", true)
respStruct := models.PurgeQueueResponse{ }
Xmlns: models.BaseXmlns,
Metadata: models.BaseResponseMetadata, log.Infof("Purging Queue: %s (tenant: %s)", queueName, t.ID)
} models.SyncQueues.Queues[key].Messages = nil
return http.StatusOK, respStruct models.SyncQueues.Queues[key].Duplicates = make(map[string]time.Time)
respStruct := models.PurgeQueueResponse{
Xmlns: models.BaseXmlns,
Metadata: models.BaseResponseMetadata,
}
return http.StatusOK, respStruct
} }
+138 -134
View File
@@ -1,164 +1,168 @@
// Изменено: 2026-04-09
// ReceiveMessageV1 — получает сообщения из очереди тенанта с поддержкой long polling.
// Ловушка #4: long polling держит соединение до 20 сек — не прерываем принудительно.
package gosqs package gosqs
import ( import (
"fmt" "fmt"
"net/http" "net/http"
"strings" "strings"
"time" "time"
"github.com/google/uuid" "github.com/google/uuid"
"shared-sqs/app/interfaces" "shared-sqs/app/interfaces"
"shared-sqs/app/models" "shared-sqs/app/models"
"shared-sqs/app/utils" "shared-sqs/app/utils"
"github.com/gorilla/mux" "github.com/gorilla/mux"
log "github.com/sirupsen/logrus" log "github.com/sirupsen/logrus"
) )
// TODO - Admiral-Piett - could we refactor the way we hide messages? Change data structure to a queue
// organized by "reveal time" or a map with the key being a timestamp of when it could be shown?
// Ordered Map - https://github.com/elliotchance/orderedmap
func ReceiveMessageV1(req *http.Request) (int, interfaces.AbstractResponseBody) { func ReceiveMessageV1(req *http.Request) (int, interfaces.AbstractResponseBody) {
requestBody := models.NewReceiveMessageRequest() requestBody := models.NewReceiveMessageRequest()
ok := utils.REQUEST_TRANSFORMER(requestBody, req, false) ok := utils.REQUEST_TRANSFORMER(requestBody, req, false)
if !ok { if !ok {
log.Error("Invalid Request - ReceiveMessageV1") log.Error("Invalid Request - ReceiveMessageV1")
return utils.CreateErrorResponseV1("InvalidParameterValue", true) return utils.CreateErrorResponseV1("InvalidParameterValue", true)
} }
maxNumberOfMessages := requestBody.MaxNumberOfMessages t := getTenantFromContext(req)
if maxNumberOfMessages == 0 { if t == nil {
maxNumberOfMessages = 1 return utils.CreateErrorResponseV1("InvalidClientTokenId", true)
} }
queueName := "" maxNumberOfMessages := requestBody.MaxNumberOfMessages
if requestBody.QueueUrl == "" { if maxNumberOfMessages == 0 {
vars := mux.Vars(req) maxNumberOfMessages = 1
queueName = vars["queueName"] }
} else {
uriSegments := strings.Split(requestBody.QueueUrl, "/")
queueName = uriSegments[len(uriSegments)-1]
}
if _, ok := models.SyncQueues.Queues[queueName]; !ok { queueName := ""
return utils.CreateErrorResponseV1("QueueNotFound", true) if requestBody.QueueUrl == "" {
} vars := mux.Vars(req)
queueName = vars["queueName"]
} else {
uriSegments := strings.Split(requestBody.QueueUrl, "/")
queueName = uriSegments[len(uriSegments)-1]
}
var messages []*models.ResultMessage key := tenantQueueKey(t.AccessKey, queueName)
respStruct := models.ReceiveMessageResponse{}
waitTimeSeconds := requestBody.WaitTimeSeconds if _, ok := models.SyncQueues.Queues[key]; !ok {
if waitTimeSeconds == 0 { return utils.CreateErrorResponseV1("QueueNotFound", true)
models.SyncQueues.RLock() }
waitTimeSeconds = models.SyncQueues.Queues[queueName].ReceiveMessageWaitTimeSeconds
models.SyncQueues.RUnlock()
}
loops := waitTimeSeconds * 10 var messages []*models.ResultMessage
for loops > 0 { respStruct := models.ReceiveMessageResponse{}
models.SyncQueues.RLock()
_, queueFound := models.SyncQueues.Queues[queueName]
if !queueFound {
models.SyncQueues.RUnlock()
return utils.CreateErrorResponseV1("QueueNotFound", true)
}
messageFound := len(models.SyncQueues.Queues[queueName].Messages)-numberOfHiddenMessagesInQueue(*models.SyncQueues.Queues[queueName]) != 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
}
} waitTimeSeconds := requestBody.WaitTimeSeconds
log.Debugf("Getting Message from Queue:%s", queueName) if waitTimeSeconds == 0 {
models.SyncQueues.RLock()
waitTimeSeconds = models.SyncQueues.Queues[key].ReceiveMessageWaitTimeSeconds
models.SyncQueues.RUnlock()
}
models.SyncQueues.Lock() // Lock the Queues // Long polling: ждём появления сообщения до waitTimeSeconds*10 итераций по 100ms
defer models.SyncQueues.Unlock() // Unlock the Queues 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)
if len(models.SyncQueues.Queues[queueName].Messages) > 0 { models.SyncQueues.Lock()
numMsg := 0 defer models.SyncQueues.Unlock()
messages = make([]*models.ResultMessage, 0)
for i := range models.SyncQueues.Queues[queueName].Messages {
if numMsg >= maxNumberOfMessages {
break
}
if models.SyncQueues.Queues[queueName].Messages[i].ReceiptHandle != "" { if len(models.SyncQueues.Queues[key].Messages) > 0 {
continue numMsg := 0
} messages = make([]*models.ResultMessage, 0)
for i := range models.SyncQueues.Queues[key].Messages {
if numMsg >= maxNumberOfMessages {
break
}
msg := &models.SyncQueues.Queues[queueName].Messages[i] if models.SyncQueues.Queues[key].Messages[i].ReceiptHandle != "" {
if !msg.IsReadyForReceipt() { continue
continue }
}
if models.SyncQueues.Queues[queueName].IsFIFO { msg := &models.SyncQueues.Queues[key].Messages[i]
// If we got messages here it means we have not processed it yet, so get next if !msg.IsReadyForReceipt() {
if models.SyncQueues.Queues[queueName].IsLocked(msg.GroupID) { continue
continue }
}
// Otherwise lock messages for group ID
models.SyncQueues.Queues[queueName].LockGroup(msg.GroupID)
}
randomId := uuid.NewString() if models.SyncQueues.Queues[key].IsFIFO {
msg.ReceiptHandle = msg.Uuid + "#" + randomId if models.SyncQueues.Queues[key].IsLocked(msg.GroupID) {
msg.ReceiptTime = time.Now().UTC() continue
}
models.SyncQueues.Queues[key].LockGroup(msg.GroupID)
}
if requestBody.VisibilityTimeout != 0 { randomId := uuid.NewString()
msg.VisibilityTimeout = time.Now().Add(time.Duration(requestBody.VisibilityTimeout) * time.Second) msg.ReceiptHandle = msg.Uuid + "#" + randomId
} else { msg.ReceiptTime = time.Now().UTC()
msg.VisibilityTimeout = time.Now().Add(time.Duration(models.SyncQueues.Queues[queueName].VisibilityTimeout) * time.Second)
}
messages = append(messages, buildResultMessage(msg)) 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)
}
numMsg++ messages = append(messages, buildResultMessage(msg))
} numMsg++
}
respStruct = models.ReceiveMessageResponse{ respStruct = models.ReceiveMessageResponse{
"http://queue.amazonaws.com/doc/2012-11-05/", "http://queue.amazonaws.com/doc/2012-11-05/",
models.ReceiveMessageResult{ models.ReceiveMessageResult{Messages: messages},
Messages: messages, models.ResponseMetadata{RequestId: "00000000-0000-0000-0000-000000000000"},
}, }
models.ResponseMetadata{ } else {
RequestId: "00000000-0000-0000-0000-000000000000", log.Warning("No messages in Queue:", queueName)
}, respStruct = models.ReceiveMessageResponse{
} Xmlns: "http://queue.amazonaws.com/doc/2012-11-05/",
} else { Result: models.ReceiveMessageResult{},
log.Warning("No messages in Queue:", queueName) Metadata: models.ResponseMetadata{RequestId: "00000000-0000-0000-0000-000000000000"},
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 return http.StatusOK, respStruct
} }
func buildResultMessage(m *models.SqsMessage) *models.ResultMessage { func buildResultMessage(m *models.SqsMessage) *models.ResultMessage {
return &models.ResultMessage{ return &models.ResultMessage{
MessageId: m.Uuid, MessageId: m.Uuid,
Body: m.MessageBody, Body: m.MessageBody,
ReceiptHandle: m.ReceiptHandle, ReceiptHandle: m.ReceiptHandle,
MD5OfBody: utils.GetMD5Hash(m.MessageBody), MD5OfBody: utils.GetMD5Hash(m.MessageBody),
MD5OfMessageAttributes: m.MD5OfMessageAttributes, MD5OfMessageAttributes: m.MD5OfMessageAttributes,
MessageAttributes: m.MessageAttributes, MessageAttributes: m.MessageAttributes,
Attributes: map[string]string{ Attributes: map[string]string{
"ApproximateFirstReceiveTimestamp": fmt.Sprintf("%d", m.ReceiptTime.UnixNano()/int64(time.Millisecond)), "ApproximateFirstReceiveTimestamp": fmt.Sprintf("%d", m.ReceiptTime.UnixNano()/int64(time.Millisecond)),
"SenderId": models.CurrentEnvironment.AccountID, "SenderId": models.CurrentEnvironment.AccountID,
"ApproximateReceiveCount": fmt.Sprintf("%d", m.NumberOfReceives+1), "ApproximateReceiveCount": fmt.Sprintf("%d", m.NumberOfReceives+1),
"SentTimestamp": fmt.Sprintf("%d", time.Now().UTC().UnixNano()/int64(time.Millisecond)), "SentTimestamp": fmt.Sprintf("%d", time.Now().UTC().UnixNano()/int64(time.Millisecond)),
}, },
} }
} }
+97 -89
View File
@@ -1,100 +1,108 @@
// Изменено: 2026-04-09
// SendMessageV1 — добавляет сообщение в очередь тенанта.
// Ловушка #6: queueName извлекается как ПОСЛЕДНИЙ сегмент URL — при URL вида
// http://host/tenantID/queueName последний сегмент = queueName (правильно).
package gosqs package gosqs
import ( import (
"net/http" "net/http"
"strings" "strings"
"time" "time"
"github.com/google/uuid" "github.com/google/uuid"
"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"
log "github.com/sirupsen/logrus" "github.com/gorilla/mux"
"github.com/gorilla/mux"
) )
func SendMessageV1(req *http.Request) (int, interfaces.AbstractResponseBody) { func SendMessageV1(req *http.Request) (int, interfaces.AbstractResponseBody) {
requestBody := models.NewSendMessageRequest() requestBody := models.NewSendMessageRequest()
ok := utils.REQUEST_TRANSFORMER(requestBody, req, false) ok := utils.REQUEST_TRANSFORMER(requestBody, req, false)
if !ok { if !ok {
log.Error("Invalid Request - SendMessageV1") log.Error("Invalid Request - SendMessageV1")
return utils.CreateErrorResponseV1("InvalidParameterValue", true) return utils.CreateErrorResponseV1("InvalidParameterValue", true)
} }
messageBody := requestBody.MessageBody
messageGroupID := requestBody.MessageGroupId t := getTenantFromContext(req)
messageDeduplicationID := requestBody.MessageDeduplicationId if t == nil {
return utils.CreateErrorResponseV1("InvalidClientTokenId", true)
queueUrl := getQueueFromPath(requestBody.QueueUrl, req.URL.String()) }
queueName := "" messageBody := requestBody.MessageBody
if queueUrl == "" { messageGroupID := requestBody.MessageGroupId
// TODO: Remove this query param logic if it's not still valid or something messageDeduplicationID := requestBody.MessageDeduplicationId
vars := mux.Vars(req)
queueName = vars["queueName"] queueUrl := getQueueFromPath(requestBody.QueueUrl, req.URL.String())
} else { queueName := ""
uriSegments := strings.Split(queueUrl, "/") if queueUrl == "" {
queueName = uriSegments[len(uriSegments)-1] vars := mux.Vars(req)
} queueName = vars["queueName"]
} else {
if _, ok := models.SyncQueues.Queues[queueName]; !ok { // Ловушка #6: берём последний сегмент — это queueName, не tenantID
// Queue does not exist uriSegments := strings.Split(queueUrl, "/")
return utils.CreateErrorResponseV1("QueueNotFound", true) queueName = uriSegments[len(uriSegments)-1]
} }
if models.SyncQueues.Queues[queueName].MaximumMessageSize > 0 && key := tenantQueueKey(t.AccessKey, queueName)
len(messageBody) > models.SyncQueues.Queues[queueName].MaximumMessageSize {
// Message size is too big if _, ok := models.SyncQueues.Queues[key]; !ok {
return utils.CreateErrorResponseV1("MessageTooBig", true) return utils.CreateErrorResponseV1("QueueNotFound", true)
} }
delaySecs := models.SyncQueues.Queues[queueName].DelaySeconds if models.SyncQueues.Queues[key].MaximumMessageSize > 0 &&
if requestBody.DelaySeconds != 0 { len(messageBody) > models.SyncQueues.Queues[key].MaximumMessageSize {
delaySecs = requestBody.DelaySeconds return utils.CreateErrorResponseV1("MessageTooBig", true)
} }
log.Debugf("Putting Message in Queue: [%s]", queueName) delaySecs := models.SyncQueues.Queues[key].DelaySeconds
msg := models.SqsMessage{MessageBody: messageBody} if requestBody.DelaySeconds != 0 {
if len(requestBody.MessageAttributes) > 0 { delaySecs = requestBody.DelaySeconds
msg.MessageAttributes = requestBody.MessageAttributes }
msg.MD5OfMessageAttributes = utils.HashAttributes(requestBody.MessageAttributes)
} log.Debugf("Putting Message in Queue: [%s] tenant: [%s]", queueName, t.ID)
msg.MD5OfMessageBody = utils.GetMD5Hash(messageBody) msg := models.SqsMessage{MessageBody: messageBody}
msg.Uuid = uuid.NewString() if len(requestBody.MessageAttributes) > 0 {
msg.GroupID = messageGroupID msg.MessageAttributes = requestBody.MessageAttributes
msg.DeduplicationID = messageDeduplicationID msg.MD5OfMessageAttributes = utils.HashAttributes(requestBody.MessageAttributes)
msg.SentTime = time.Now() }
msg.DelaySecs = delaySecs msg.MD5OfMessageBody = utils.GetMD5Hash(messageBody)
msg.Uuid = uuid.NewString()
models.SyncQueues.Lock() msg.GroupID = messageGroupID
fifoSeqNumber := "" msg.DeduplicationID = messageDeduplicationID
if models.SyncQueues.Queues[queueName].IsFIFO { msg.SentTime = time.Now()
fifoSeqNumber = models.SyncQueues.Queues[queueName].NextSequenceNumber(messageGroupID) msg.DelaySecs = delaySecs
}
models.SyncQueues.Lock()
if !models.SyncQueues.Queues[queueName].IsDuplicate(messageDeduplicationID) { fifoSeqNumber := ""
models.SyncQueues.Queues[queueName].Messages = append(models.SyncQueues.Queues[queueName].Messages, msg) if models.SyncQueues.Queues[key].IsFIFO {
} else { fifoSeqNumber = models.SyncQueues.Queues[key].NextSequenceNumber(messageGroupID)
log.Debugf("Message with deduplicationId [%s] in queue [%s] is duplicate ", messageDeduplicationID, queueName) }
}
if !models.SyncQueues.Queues[key].IsDuplicate(messageDeduplicationID) {
models.SyncQueues.Queues[queueName].InitDuplicatation(messageDeduplicationID) models.SyncQueues.Queues[key].Messages = append(models.SyncQueues.Queues[key].Messages, msg)
models.SyncQueues.Unlock() } else {
log.Infof("%s: Queue: %s, Message: %s\n", time.Now().Format("2006-01-02 15:04:05"), queueName, msg.MessageBody) log.Debugf("Duplicate message deduplicationId [%s] in queue [%s]", messageDeduplicationID, queueName)
}
respStruct := models.SendMessageResponse{
Xmlns: models.BaseXmlns, models.SyncQueues.Queues[key].InitDuplicatation(messageDeduplicationID)
Result: models.SendMessageResult{ models.SyncQueues.Unlock()
MD5OfMessageAttributes: msg.MD5OfMessageAttributes, log.Infof("%s: Queue: %s, Message: %s\n", time.Now().Format("2006-01-02 15:04:05"), queueName, msg.MessageBody)
MD5OfMessageBody: msg.MD5OfMessageBody,
MessageId: msg.Uuid, respStruct := models.SendMessageResponse{
SequenceNumber: fifoSeqNumber, Xmlns: models.BaseXmlns,
}, Result: models.SendMessageResult{
Metadata: models.BaseResponseMetadata, MD5OfMessageAttributes: msg.MD5OfMessageAttributes,
} MD5OfMessageBody: msg.MD5OfMessageBody,
MessageId: msg.Uuid,
return http.StatusOK, respStruct SequenceNumber: fifoSeqNumber,
},
Metadata: models.BaseResponseMetadata,
}
return http.StatusOK, respStruct
} }
+100 -96
View File
@@ -1,105 +1,109 @@
// Изменено: 2026-04-09
// SendMessageBatchV1 — пакетная отправка сообщений в очередь тенанта.
package gosqs package gosqs
import ( import (
"net/http" "net/http"
"strings" "strings"
"time" "time"
"github.com/google/uuid" "github.com/google/uuid"
"shared-sqs/app/interfaces" "shared-sqs/app/interfaces"
"shared-sqs/app/models" "shared-sqs/app/models"
"shared-sqs/app/utils" "shared-sqs/app/utils"
"github.com/gorilla/mux" "github.com/gorilla/mux"
log "github.com/sirupsen/logrus" log "github.com/sirupsen/logrus"
) )
func SendMessageBatchV1(req *http.Request) (int, interfaces.AbstractResponseBody) { func SendMessageBatchV1(req *http.Request) (int, interfaces.AbstractResponseBody) {
requestBody := models.NewSendMessageBatchRequest() requestBody := models.NewSendMessageBatchRequest()
ok := utils.REQUEST_TRANSFORMER(requestBody, req, false) ok := utils.REQUEST_TRANSFORMER(requestBody, req, false)
if !ok { if !ok {
log.Error("Invalid Request - SendMessageBatchV1") log.Error("Invalid Request - SendMessageBatchV1")
return utils.CreateErrorResponseV1("InvalidParameterValue", true) return utils.CreateErrorResponseV1("InvalidParameterValue", true)
} }
queueUrl := requestBody.QueueUrl t := getTenantFromContext(req)
if t == nil {
// TODO: Remove this query param logic if it's not still valid or something return utils.CreateErrorResponseV1("InvalidClientTokenId", true)
queueName := "" }
if queueUrl == "" {
vars := mux.Vars(req) queueUrl := requestBody.QueueUrl
queueName = vars["queueName"] queueName := ""
} else { if queueUrl == "" {
uriSegments := strings.Split(queueUrl, "/") vars := mux.Vars(req)
queueName = uriSegments[len(uriSegments)-1] queueName = vars["queueName"]
} } else {
uriSegments := strings.Split(queueUrl, "/")
if _, ok := models.SyncQueues.Queues[queueName]; !ok { queueName = uriSegments[len(uriSegments)-1]
return utils.CreateErrorResponseV1("QueueNotFound", true) }
}
key := tenantQueueKey(t.AccessKey, queueName)
sendEntries := requestBody.Entries
if _, ok := models.SyncQueues.Queues[key]; !ok {
if len(sendEntries) == 0 { return utils.CreateErrorResponseV1("QueueNotFound", true)
return utils.CreateErrorResponseV1("EmptyBatchRequest", true) }
}
sendEntries := requestBody.Entries
if len(sendEntries) > 10 {
return utils.CreateErrorResponseV1("TooManyEntriesInBatchRequest", true) if len(sendEntries) == 0 {
} return utils.CreateErrorResponseV1("EmptyBatchRequest", true)
ids := map[string]struct{}{} }
for _, v := range sendEntries {
if _, ok := ids[v.Id]; ok { if len(sendEntries) > 10 {
return utils.CreateErrorResponseV1("BatchEntryIdsNotDistinct", true) return utils.CreateErrorResponseV1("TooManyEntriesInBatchRequest", true)
} }
ids[v.Id] = struct{}{} ids := map[string]struct{}{}
} for _, v := range sendEntries {
if _, ok := ids[v.Id]; ok {
sentEntries := make([]models.SendMessageBatchResultEntry, 0) return utils.CreateErrorResponseV1("BatchEntryIdsNotDistinct", true)
log.Debug("Putting Message in Queue:", queueName) }
for _, sendEntry := range sendEntries { ids[v.Id] = struct{}{}
msg := models.SqsMessage{MessageBody: sendEntry.MessageBody} }
if len(sendEntry.MessageAttributes) > 0 {
msg.MessageAttributes = sendEntry.MessageAttributes sentEntries := make([]models.SendMessageBatchResultEntry, 0)
msg.MD5OfMessageAttributes = utils.HashAttributes(sendEntry.MessageAttributes) log.Debugf("Batch sending to Queue: %s (tenant: %s)", queueName, t.ID)
} for _, sendEntry := range sendEntries {
msg.MD5OfMessageBody = utils.GetMD5Hash(sendEntry.MessageBody) msg := models.SqsMessage{MessageBody: sendEntry.MessageBody}
msg.GroupID = sendEntry.MessageGroupId if len(sendEntry.MessageAttributes) > 0 {
msg.DeduplicationID = sendEntry.MessageDeduplicationId msg.MessageAttributes = sendEntry.MessageAttributes
msg.Uuid = uuid.NewString() msg.MD5OfMessageAttributes = utils.HashAttributes(sendEntry.MessageAttributes)
msg.SentTime = time.Now() }
models.SyncQueues.Lock() msg.MD5OfMessageBody = utils.GetMD5Hash(sendEntry.MessageBody)
fifoSeqNumber := "" msg.GroupID = sendEntry.MessageGroupId
if models.SyncQueues.Queues[queueName].IsFIFO { msg.DeduplicationID = sendEntry.MessageDeduplicationId
fifoSeqNumber = models.SyncQueues.Queues[queueName].NextSequenceNumber(sendEntry.MessageGroupId) msg.Uuid = uuid.NewString()
} msg.SentTime = time.Now()
if !models.SyncQueues.Queues[queueName].IsDuplicate(sendEntry.MessageDeduplicationId) { models.SyncQueues.Lock()
models.SyncQueues.Queues[queueName].Messages = append(models.SyncQueues.Queues[queueName].Messages, msg) fifoSeqNumber := ""
} else { if models.SyncQueues.Queues[key].IsFIFO {
log.Debugf("Message with deduplicationId [%s] in queue [%s] is duplicate ", sendEntry.MessageDeduplicationId, queueName) fifoSeqNumber = models.SyncQueues.Queues[key].NextSequenceNumber(sendEntry.MessageGroupId)
} }
if !models.SyncQueues.Queues[key].IsDuplicate(sendEntry.MessageDeduplicationId) {
models.SyncQueues.Queues[queueName].InitDuplicatation(sendEntry.MessageDeduplicationId) models.SyncQueues.Queues[key].Messages = append(models.SyncQueues.Queues[key].Messages, msg)
} else {
models.SyncQueues.Unlock() log.Debugf("Duplicate deduplicationId [%s] in queue [%s]", sendEntry.MessageDeduplicationId, queueName)
se := models.SendMessageBatchResultEntry{ }
Id: sendEntry.Id, models.SyncQueues.Queues[key].InitDuplicatation(sendEntry.MessageDeduplicationId)
MessageId: msg.Uuid, models.SyncQueues.Unlock()
MD5OfMessageBody: msg.MD5OfMessageBody,
MD5OfMessageAttributes: msg.MD5OfMessageAttributes, sentEntries = append(sentEntries, models.SendMessageBatchResultEntry{
SequenceNumber: fifoSeqNumber, Id: sendEntry.Id,
} MessageId: msg.Uuid,
sentEntries = append(sentEntries, se) MD5OfMessageBody: msg.MD5OfMessageBody,
log.Infof("%s: Queue: %s, Message: %s\n", time.Now().Format("2006-01-02 15:04:05"), queueName, msg.MessageBody) MD5OfMessageAttributes: msg.MD5OfMessageAttributes,
} SequenceNumber: fifoSeqNumber,
})
respStruct := models.SendMessageBatchResponse{ log.Infof("%s: Queue: %s, Message: %s", time.Now().Format("2006-01-02 15:04:05"), queueName, msg.MessageBody)
Xmlns: models.BaseXmlns, }
Result: models.SendMessageBatchResult{Entry: sentEntries},
Metadata: models.BaseResponseMetadata, respStruct := models.SendMessageBatchResponse{
} Xmlns: models.BaseXmlns,
Result: models.SendMessageBatchResult{Entry: sentEntries},
return http.StatusOK, respStruct Metadata: models.BaseResponseMetadata,
}
return http.StatusOK, respStruct
} }
+46 -40
View File
@@ -1,48 +1,54 @@
// Изменено: 2026-04-09
// SetQueueAttributesV1 — устанавливает атрибуты очереди тенанта.
// Ловушка #9: при RedrivePolicy парсим ARN DLQ и DLQ тоже должна принадлежать тому же тенанту.
package gosqs package gosqs
import ( import (
"net/http" "net/http"
"strings" "strings"
"shared-sqs/app/models" "shared-sqs/app/interfaces"
"shared-sqs/app/utils" "shared-sqs/app/models"
"shared-sqs/app/utils"
"shared-sqs/app/interfaces" log "github.com/sirupsen/logrus"
log "github.com/sirupsen/logrus"
) )
func SetQueueAttributesV1(req *http.Request) (int, interfaces.AbstractResponseBody) { func SetQueueAttributesV1(req *http.Request) (int, interfaces.AbstractResponseBody) {
requestBody := models.NewSetQueueAttributesRequest() requestBody := models.NewSetQueueAttributesRequest()
ok := utils.REQUEST_TRANSFORMER(requestBody, req, false) ok := utils.REQUEST_TRANSFORMER(requestBody, req, false)
if !ok { if !ok {
log.Error("Invalid Request - GetQueueAttributesV1") log.Error("Invalid Request - SetQueueAttributesV1")
return utils.CreateErrorResponseV1("InvalidParameterValue", true) return utils.CreateErrorResponseV1("InvalidParameterValue", true)
} }
if requestBody.QueueUrl == "" { if requestBody.QueueUrl == "" {
log.Error("Missing QueueUrl - GetQueueAttributesV1") log.Error("Missing QueueUrl - SetQueueAttributesV1")
return utils.CreateErrorResponseV1("InvalidParameterValue", true) return utils.CreateErrorResponseV1("InvalidParameterValue", true)
} }
// NOTE: I tore out the handling for devining the url from a param. I can't find documentation that t := getTenantFromContext(req)
// that is valid any longer. if t == nil {
uriSegments := strings.Split(requestBody.QueueUrl, "/") return utils.CreateErrorResponseV1("InvalidClientTokenId", true)
queueName := uriSegments[len(uriSegments)-1] }
log.Infof("Set Queue QueueAttributes: %s", queueName) uriSegments := strings.Split(requestBody.QueueUrl, "/")
models.SyncQueues.Lock() queueName := uriSegments[len(uriSegments)-1]
defer models.SyncQueues.Unlock() key := tenantQueueKey(t.AccessKey, queueName)
queue, ok := models.SyncQueues.Queues[queueName]
if !ok { log.Infof("Set Queue Attributes: %s (tenant: %s)", queueName, t.ID)
log.Warningf("Get Queue URL: %s, queue does not exist!!!", queueName) models.SyncQueues.Lock()
return utils.CreateErrorResponseV1("QueueNotFound", true) defer models.SyncQueues.Unlock()
} queue, ok := models.SyncQueues.Queues[key]
if err := setQueueAttributesV1(queue, requestBody.Attributes); err != nil { if !ok {
return utils.CreateErrorResponseV1(err.Error(), true) log.Warningf("Set Queue Attributes: %s, queue does not exist for tenant %s", queueName, t.ID)
} return utils.CreateErrorResponseV1("QueueNotFound", true)
}
respStruct := models.SetQueueAttributesResponse{ if err := setQueueAttributesV1(queue, requestBody.Attributes); err != nil {
Xmlns: models.BaseXmlns, return utils.CreateErrorResponseV1(err.Error(), true)
Metadata: models.BaseResponseMetadata, }
}
return http.StatusOK, respStruct respStruct := models.SetQueueAttributesResponse{
Xmlns: models.BaseXmlns,
Metadata: models.BaseResponseMetadata,
}
return http.StatusOK, respStruct
} }
+56
View File
@@ -0,0 +1,56 @@
// Изменено: 2026-04-09
// Helper-функции для tenant-scoped операций с очередями.
// Используются всеми SQS handlers для изоляции очередей между тенантами.
package gosqs
import (
"net/http"
"shared-sqs/app/auth"
"shared-sqs/app/models"
"shared-sqs/app/tenant"
)
// tenantQueueKey — внутренний ключ очереди в SyncQueues в формате "{accessKey}:{queueName}".
// Такой формат гарантирует изоляцию: тенант видит только очереди с префиксом своего accessKey.
func tenantQueueKey(tenantAccessKey, queueName string) string {
return tenantAccessKey + ":" + queueName
}
// getTenantFromContext — извлекает тенанта из request context.
// Возвращает nil если тенант не найден (не должно быть — auth middleware должен это поймать раньше).
func getTenantFromContext(r *http.Request) *tenant.Tenant {
t, _ := r.Context().Value(auth.TenantContextKey).(*tenant.Tenant)
return t
}
// tenantQueueURL — формирует URL очереди для тенанта.
// Ловушка #10: QueueUrl ОБЯЗАН содержать tenantID в пути, иначе AWS SDK не сможет send/receive.
func tenantQueueURL(t *tenant.Tenant, queueName string) string {
host := models.CurrentEnvironment.Host
port := models.CurrentEnvironment.Port
region := models.CurrentEnvironment.Region
if region != "" {
return "http://" + region + "." + host + ":" + port + "/" + t.ID + "/" + queueName
}
return "http://" + host + ":" + port + "/" + t.ID + "/" + queueName
}
// tenantQueueARN — формирует ARN очереди для тенанта.
func tenantQueueARN(t *tenant.Tenant, queueName string) string {
return "arn:aws:sqs:" + models.CurrentEnvironment.Region + ":" + t.ID + ":" + queueName
}
// countTenantQueues — считает количество очередей тенанта в SyncQueues.
// Используется для проверки лимита MaxQueues.
// Вызывать под SyncQueues.RLock().
func countTenantQueues(tenantAccessKey string) int {
prefix := tenantAccessKey + ":"
count := 0
for key := range models.SyncQueues.Queues {
if len(key) > len(prefix) && key[:len(prefix)] == prefix {
count++
}
}
return count
}
+142
View File
@@ -0,0 +1,142 @@
// Изменено: 2026-04-09
// Tenant model и in-memory хранилище тенантов для shared-sqs.
package tenant
import (
"crypto/rand"
"encoding/hex"
"fmt"
"sync"
"time"
)
// Tenant — модель тенанта shared-sqs.
// AccessKey используется как идентификатор в AWS Authorization header.
type Tenant struct {
ID string // уникальный идентификатор тенанта (t-<hex>)
Name string // человекочитаемое имя
AccessKey string // аналог AWS AccessKeyId (SSAK-<hex>)
SecretKey string // аналог AWS SecretAccessKey (64 hex chars)
MaxQueues int // лимит очередей (0 = безлимит)
CreatedAt time.Time
Active bool
}
// TenantStore — потокобезопасное in-memory хранилище тенантов.
// Два индекса позволяют быстро искать как по ID (admin API), так и по AccessKey (auth middleware).
type TenantStore struct {
mu sync.RWMutex
byID map[string]*Tenant
byAccessKey map[string]*Tenant
}
// NewTenantStore — создаёт пустое хранилище тенантов.
func NewTenantStore() *TenantStore {
return &TenantStore{
byID: make(map[string]*Tenant),
byAccessKey: make(map[string]*Tenant),
}
}
// Create — создаёт нового тенанта, генерирует ключи, сохраняет в оба индекса.
func (s *TenantStore) Create(name string, maxQueues int) (*Tenant, error) {
id, err := generateTenantID()
if err != nil {
return nil, fmt.Errorf("generate tenant id: %w", err)
}
accessKey, err := generateAccessKey()
if err != nil {
return nil, fmt.Errorf("generate access key: %w", err)
}
secretKey, err := generateSecretKey()
if err != nil {
return nil, fmt.Errorf("generate secret key: %w", err)
}
t := &Tenant{
ID: id,
Name: name,
AccessKey: accessKey,
SecretKey: secretKey,
MaxQueues: maxQueues,
CreatedAt: time.Now().UTC(),
Active: true,
}
s.mu.Lock()
s.byID[t.ID] = t
s.byAccessKey[t.AccessKey] = t
s.mu.Unlock()
return t, nil
}
// GetByAccessKey — поиск тенанта по AccessKeyId (используется в auth middleware).
func (s *TenantStore) GetByAccessKey(accessKey string) (*Tenant, bool) {
s.mu.RLock()
t, ok := s.byAccessKey[accessKey]
s.mu.RUnlock()
return t, ok
}
// GetByID — поиск тенанта по ID (используется в admin API).
func (s *TenantStore) GetByID(id string) (*Tenant, bool) {
s.mu.RLock()
t, ok := s.byID[id]
s.mu.RUnlock()
return t, ok
}
// Delete — удаляет тенанта из ОБОИХ индексов.
// Ловушка #2: если удалить только из одного индекса — orphaned данные и memory leak.
func (s *TenantStore) Delete(id string) bool {
s.mu.Lock()
defer s.mu.Unlock()
t, ok := s.byID[id]
if !ok {
return false
}
delete(s.byID, t.ID)
delete(s.byAccessKey, t.AccessKey)
return true
}
// List — список всех тенантов (для admin GET /tenants).
func (s *TenantStore) List() []*Tenant {
s.mu.RLock()
result := make([]*Tenant, 0, len(s.byID))
for _, t := range s.byID {
result = append(result, t)
}
s.mu.RUnlock()
return result
}
// generateTenantID — генерирует уникальный ID тенанта в формате t-<12 hex bytes>.
func generateTenantID() (string, error) {
b := make([]byte, 8)
if _, err := rand.Read(b); err != nil {
return "", err
}
return "t-" + hex.EncodeToString(b), nil
}
// generateAccessKey — генерирует AccessKey в формате SSAK-<12 hex bytes>.
// SSAK = Shared SQS Access Key. Используем crypto/rand (ловушка #1: не math/rand).
func generateAccessKey() (string, error) {
b := make([]byte, 12)
if _, err := rand.Read(b); err != nil {
return "", err
}
return "SSAK-" + hex.EncodeToString(b), nil
}
// generateSecretKey — генерирует SecretKey как 64 hex символа (32 random bytes).
// Используем crypto/rand (ловушка #1).
func generateSecretKey() (string, error) {
b := make([]byte, 32)
if _, err := rand.Read(b); err != nil {
return "", err
}
return hex.EncodeToString(b), nil
}