diff --git a/shared-sqs/app/auth/auth_middleware.go b/shared-sqs/app/auth/auth_middleware.go new file mode 100644 index 0000000..997363d --- /dev/null +++ b/shared-sqs/app/auth/auth_middleware.go @@ -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) +} diff --git a/shared-sqs/app/gosqs/change_message_visibility.go b/shared-sqs/app/gosqs/change_message_visibility.go index 56427c2..f34e9e5 100644 --- a/shared-sqs/app/gosqs/change_message_visibility.go +++ b/shared-sqs/app/gosqs/change_message_visibility.go @@ -1,81 +1,87 @@ +// Изменено: 2026-04-09 +// ChangeMessageVisibilityV1 — меняет visibility timeout сообщения в очереди тенанта. package gosqs import ( - "net/http" - "strings" - "time" +"net/http" +"strings" +"time" - "shared-sqs/app/interfaces" - "shared-sqs/app/models" - "shared-sqs/app/utils" - "github.com/gorilla/mux" - log "github.com/sirupsen/logrus" +"shared-sqs/app/interfaces" +"shared-sqs/app/models" +"shared-sqs/app/utils" +"github.com/gorilla/mux" +log "github.com/sirupsen/logrus" ) func ChangeMessageVisibilityV1(req *http.Request) (int, interfaces.AbstractResponseBody) { - requestBody := models.NewChangeMessageVisibilityRequest() - ok := utils.REQUEST_TRANSFORMER(requestBody, req, false) - if !ok { - log.Error("Invalid Request - ChangeMessageVisibilityV1") - return utils.CreateErrorResponseV1("InvalidParameterValue", true) - } - - vars := mux.Vars(req) - - queueUrl := requestBody.QueueUrl - queueName := "" - if queueUrl == "" { - queueName = vars["queueName"] - } else { - uriSegments := strings.Split(queueUrl, "/") - queueName = uriSegments[len(uriSegments)-1] - } - - receiptHandle := requestBody.ReceiptHandle - - visibilityTimeout := requestBody.VisibilityTimeout - if visibilityTimeout > 43200 { - return utils.CreateErrorResponseV1("ValidationError", true) - } - - if _, ok := models.SyncQueues.Queues[queueName]; !ok { - return utils.CreateErrorResponseV1("QueueNotFound", true) - } - - models.SyncQueues.Lock() - messageFound := false - for i := 0; i < len(models.SyncQueues.Queues[queueName].Messages); i++ { - queue := models.SyncQueues.Queues[queueName] - msgs := queue.Messages - if msgs[i].ReceiptHandle == receiptHandle { - timeout := models.SyncQueues.Queues[queueName].VisibilityTimeout - if visibilityTimeout == 0 { - msgs[i].ReceiptTime = time.Now().UTC() - msgs[i].ReceiptHandle = "" - msgs[i].VisibilityTimeout = time.Now().Add(time.Duration(timeout) * time.Second) - msgs[i].Retry++ - if queue.MaxReceiveCount > 0 && - queue.DeadLetterQueue != nil && - msgs[i].Retry >= queue.MaxReceiveCount { - queue.DeadLetterQueue.Messages = append(queue.DeadLetterQueue.Messages, msgs[i]) - queue.Messages = append(queue.Messages[:i], queue.Messages[i+1:]...) - } - } else { - msgs[i].VisibilityTimeout = time.Now().Add(time.Duration(visibilityTimeout) * time.Second) - } - messageFound = true - break - } - } - models.SyncQueues.Unlock() - if !messageFound { - return utils.CreateErrorResponseV1("MessageNotInFlight", true) - } - - respStruct := models.ChangeMessageVisibilityResult{ - Xmlns: models.BaseXmlns, - Metadata: models.BaseResponseMetadata, - } - - return http.StatusOK, &respStruct +requestBody := models.NewChangeMessageVisibilityRequest() +ok := utils.REQUEST_TRANSFORMER(requestBody, req, false) +if !ok { +log.Error("Invalid Request - ChangeMessageVisibilityV1") +return utils.CreateErrorResponseV1("InvalidParameterValue", true) +} + +t := getTenantFromContext(req) +if t == nil { +return utils.CreateErrorResponseV1("InvalidClientTokenId", true) +} + +vars := mux.Vars(req) +queueUrl := requestBody.QueueUrl +queueName := "" +if queueUrl == "" { +queueName = vars["queueName"] +} else { +uriSegments := strings.Split(queueUrl, "/") +queueName = uriSegments[len(uriSegments)-1] +} + +key := tenantQueueKey(t.AccessKey, queueName) +receiptHandle := requestBody.ReceiptHandle +visibilityTimeout := requestBody.VisibilityTimeout + +if visibilityTimeout > 43200 { +return utils.CreateErrorResponseV1("ValidationError", true) +} + +if _, ok := models.SyncQueues.Queues[key]; !ok { +return utils.CreateErrorResponseV1("QueueNotFound", true) +} + +models.SyncQueues.Lock() +messageFound := false +for i := 0; i < len(models.SyncQueues.Queues[key].Messages); i++ { +queue := models.SyncQueues.Queues[key] +msgs := queue.Messages +if msgs[i].ReceiptHandle == receiptHandle { +timeout := models.SyncQueues.Queues[key].VisibilityTimeout +if visibilityTimeout == 0 { +msgs[i].ReceiptTime = time.Now().UTC() +msgs[i].ReceiptHandle = "" +msgs[i].VisibilityTimeout = time.Now().Add(time.Duration(timeout) * time.Second) +msgs[i].Retry++ +if queue.MaxReceiveCount > 0 && +queue.DeadLetterQueue != nil && +msgs[i].Retry >= queue.MaxReceiveCount { +queue.DeadLetterQueue.Messages = append(queue.DeadLetterQueue.Messages, msgs[i]) +queue.Messages = append(queue.Messages[:i], queue.Messages[i+1:]...) +} +} else { +msgs[i].VisibilityTimeout = time.Now().Add(time.Duration(visibilityTimeout) * time.Second) +} +messageFound = true +break +} +} +models.SyncQueues.Unlock() +if !messageFound { +return utils.CreateErrorResponseV1("MessageNotInFlight", true) +} + +respStruct := models.ChangeMessageVisibilityResult{ +Xmlns: models.BaseXmlns, +Metadata: models.BaseResponseMetadata, +} +return http.StatusOK, &respStruct } diff --git a/shared-sqs/app/gosqs/create_queue.go b/shared-sqs/app/gosqs/create_queue.go index 2c84faa..2e42e5d 100644 --- a/shared-sqs/app/gosqs/create_queue.go +++ b/shared-sqs/app/gosqs/create_queue.go @@ -1,54 +1,65 @@ +// Изменено: 2026-04-09 +// CreateQueueV1 — создаёт очередь для тенанта из request context. +// Ключ в SyncQueues: "{tenantAccessKey}:{queueName}" для изоляции между тенантами. package gosqs import ( - "net/http" - "time" +"net/http" +"time" - "shared-sqs/app/interfaces" - "shared-sqs/app/models" - "shared-sqs/app/utils" - log "github.com/sirupsen/logrus" +"shared-sqs/app/interfaces" +"shared-sqs/app/models" +"shared-sqs/app/utils" +log "github.com/sirupsen/logrus" ) func CreateQueueV1(req *http.Request) (int, interfaces.AbstractResponseBody) { - requestBody := models.NewCreateQueueRequest() - ok := utils.REQUEST_TRANSFORMER(requestBody, req, false) - if !ok { - log.Error("Invalid Request - CreateQueueV1") - return utils.CreateErrorResponseV1("InvalidParameterValue", true) - } - queueName := requestBody.QueueName - - queueUrl := "http://" + models.CurrentEnvironment.Host + ":" + models.CurrentEnvironment.Port + - "/" + models.CurrentEnvironment.AccountID + "/" + queueName - if models.CurrentEnvironment.Region != "" { - queueUrl = "http://" + models.CurrentEnvironment.Region + "." + models.CurrentEnvironment.Host + ":" + - models.CurrentEnvironment.Port + "/" + models.CurrentEnvironment.AccountID + "/" + queueName - } - queueArn := "arn:aws:sqs:" + models.CurrentEnvironment.Region + ":" + models.CurrentEnvironment.AccountID + ":" + queueName - - if _, ok := models.SyncQueues.Queues[queueName]; !ok { - log.Infof("Creating Queue: %s", queueName) - queue := &models.Queue{ - Name: queueName, - URL: queueUrl, - Arn: queueArn, - IsFIFO: utils.HasFIFOQueueName(queueName), - EnableDuplicates: models.CurrentEnvironment.EnableDuplicates, - Duplicates: make(map[string]time.Time), - } - if err := setQueueAttributesV1(queue, requestBody.Attributes); err != nil { - return utils.CreateErrorResponseV1(err.Error(), true) - } - models.SyncQueues.Lock() - models.SyncQueues.Queues[queueName] = queue - models.SyncQueues.Unlock() - } - - respStruct := models.CreateQueueResponse{ - Xmlns: models.BaseXmlns, - Result: models.CreateQueueResult{QueueUrl: queueUrl}, - Metadata: models.BaseResponseMetadata, - } - return http.StatusOK, respStruct +requestBody := models.NewCreateQueueRequest() +ok := utils.REQUEST_TRANSFORMER(requestBody, req, false) +if !ok { +log.Error("Invalid Request - CreateQueueV1") +return utils.CreateErrorResponseV1("InvalidParameterValue", true) +} + +t := getTenantFromContext(req) +if t == nil { +return utils.CreateErrorResponseV1("InvalidClientTokenId", true) +} + +// Ловушка #8: передаём queueName (не key) в HasFIFOQueueName — иначе .fifo не определится +queueName := requestBody.QueueName +key := tenantQueueKey(t.AccessKey, queueName) +queueUrl := tenantQueueURL(t, queueName) +queueArn := tenantQueueARN(t, queueName) + +models.SyncQueues.Lock() +if _, exists := models.SyncQueues.Queues[key]; !exists { +// Проверка лимита очередей тенанта +if t.MaxQueues > 0 && countTenantQueues(t.AccessKey) >= t.MaxQueues { +models.SyncQueues.Unlock() +return utils.CreateErrorResponseV1("LimitExceeded", true) +} +log.Infof("Creating Queue: %s (tenant: %s)", queueName, t.ID) +queue := &models.Queue{ +Name: queueName, +URL: queueUrl, +Arn: queueArn, +IsFIFO: utils.HasFIFOQueueName(queueName), +EnableDuplicates: models.CurrentEnvironment.EnableDuplicates, +Duplicates: make(map[string]time.Time), +} +if err := setQueueAttributesV1(queue, requestBody.Attributes); err != nil { +models.SyncQueues.Unlock() +return utils.CreateErrorResponseV1(err.Error(), true) +} +models.SyncQueues.Queues[key] = queue +} +models.SyncQueues.Unlock() + +respStruct := models.CreateQueueResponse{ +Xmlns: models.BaseXmlns, +Result: models.CreateQueueResult{QueueUrl: queueUrl}, +Metadata: models.BaseResponseMetadata, +} +return http.StatusOK, respStruct } diff --git a/shared-sqs/app/gosqs/delete_message.go b/shared-sqs/app/gosqs/delete_message.go index 0f7080a..1e2edbf 100644 --- a/shared-sqs/app/gosqs/delete_message.go +++ b/shared-sqs/app/gosqs/delete_message.go @@ -1,65 +1,64 @@ +// Изменено: 2026-04-09 +// DeleteMessageV1 — удаляет сообщение из очереди тенанта по ReceiptHandle. package gosqs import ( - "net/http" - "strings" +"net/http" +"strings" - "shared-sqs/app/interfaces" - "shared-sqs/app/models" - "shared-sqs/app/utils" - "github.com/gorilla/mux" - log "github.com/sirupsen/logrus" +"shared-sqs/app/interfaces" +"shared-sqs/app/models" +"shared-sqs/app/utils" +"github.com/gorilla/mux" +log "github.com/sirupsen/logrus" ) func DeleteMessageV1(req *http.Request) (int, interfaces.AbstractResponseBody) { - requestBody := models.NewDeleteMessageRequest() - ok := utils.REQUEST_TRANSFORMER(requestBody, req, false) - if !ok { - log.Error("Invalid Request - DeleteMessageV1") - return utils.CreateErrorResponseV1("InvalidParameterValue", true) - } - - // Retrieve FormValues required - receiptHandle := requestBody.ReceiptHandle - - // Retrieve FormValues required - queueUrl := requestBody.QueueUrl - queueName := "" - if queueUrl == "" { - vars := mux.Vars(req) - queueName = vars["queueName"] - } 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 - models.SyncQueues.Lock() - defer models.SyncQueues.Unlock() - if _, ok := models.SyncQueues.Queues[queueName]; ok { - for i, msg := range models.SyncQueues.Queues[queueName].Messages { - if msg.ReceiptHandle == receiptHandle { - // Unlock messages for the group - log.Debugf("FIFO Queue %s unlocking group %s:", queueName, msg.GroupID) - models.SyncQueues.Queues[queueName].UnlockGroup(msg.GroupID) - //Delete message from Q - models.SyncQueues.Queues[queueName].Messages = append(models.SyncQueues.Queues[queueName].Messages[:i], models.SyncQueues.Queues[queueName].Messages[i+1:]...) - delete(models.SyncQueues.Queues[queueName].Duplicates, msg.DeduplicationID) - - // Create, encode/xml and send response - respStruct := models.DeleteMessageResponse{ - Xmlns: models.BaseXmlns, - Metadata: models.BaseResponseMetadata, - } - return 200, &respStruct - } - } - log.Warning("Receipt Handle not found") - } else { - log.Warning("Queue not found") - } - - return utils.CreateErrorResponseV1("MessageDoesNotExist", true) +requestBody := models.NewDeleteMessageRequest() +ok := utils.REQUEST_TRANSFORMER(requestBody, req, false) +if !ok { +log.Error("Invalid Request - DeleteMessageV1") +return utils.CreateErrorResponseV1("InvalidParameterValue", true) +} + +t := getTenantFromContext(req) +if t == nil { +return utils.CreateErrorResponseV1("InvalidClientTokenId", true) +} + +receiptHandle := requestBody.ReceiptHandle +queueUrl := requestBody.QueueUrl +queueName := "" +if queueUrl == "" { +vars := mux.Vars(req) +queueName = vars["queueName"] +} else { +uriSegments := strings.Split(queueUrl, "/") +queueName = uriSegments[len(uriSegments)-1] +} + +key := tenantQueueKey(t.AccessKey, queueName) +log.Info("Deleting Message, Queue:", queueName, ", ReceiptHandle:", receiptHandle) + +models.SyncQueues.Lock() +defer models.SyncQueues.Unlock() +if _, ok := models.SyncQueues.Queues[key]; ok { +for i, msg := range models.SyncQueues.Queues[key].Messages { +if msg.ReceiptHandle == receiptHandle { +models.SyncQueues.Queues[key].UnlockGroup(msg.GroupID) +models.SyncQueues.Queues[key].Messages = append(models.SyncQueues.Queues[key].Messages[:i], models.SyncQueues.Queues[key].Messages[i+1:]...) +delete(models.SyncQueues.Queues[key].Duplicates, msg.DeduplicationID) +respStruct := models.DeleteMessageResponse{ +Xmlns: models.BaseXmlns, +Metadata: models.BaseResponseMetadata, +} +return 200, &respStruct +} +} +log.Warning("Receipt Handle not found") +} else { +log.Warning("Queue not found") +} + +return utils.CreateErrorResponseV1("MessageDoesNotExist", true) } diff --git a/shared-sqs/app/gosqs/delete_message_batch.go b/shared-sqs/app/gosqs/delete_message_batch.go index 01f1a23..ebafc99 100644 --- a/shared-sqs/app/gosqs/delete_message_batch.go +++ b/shared-sqs/app/gosqs/delete_message_batch.go @@ -1,118 +1,119 @@ +// Изменено: 2026-04-09 +// DeleteMessageBatchV1 — пакетное удаление сообщений из очереди тенанта. package gosqs import ( - "net/http" - "strings" +"net/http" +"strings" - "shared-sqs/app/interfaces" - "shared-sqs/app/models" - "shared-sqs/app/utils" - "github.com/gorilla/mux" - log "github.com/sirupsen/logrus" +"shared-sqs/app/interfaces" +"shared-sqs/app/models" +"shared-sqs/app/utils" +"github.com/gorilla/mux" +log "github.com/sirupsen/logrus" ) func DeleteMessageBatchV1(req *http.Request) (int, interfaces.AbstractResponseBody) { - requestBody := models.NewDeleteMessageBatchRequest() - ok := utils.REQUEST_TRANSFORMER(requestBody, req, false) - if !ok { - log.Error("Invalid Request - DeleteMessageBatchV1") - return utils.CreateErrorResponseV1("InvalidParameterValue", true) - } +requestBody := models.NewDeleteMessageBatchRequest() +ok := utils.REQUEST_TRANSFORMER(requestBody, req, false) +if !ok { +log.Error("Invalid Request - DeleteMessageBatchV1") +return utils.CreateErrorResponseV1("InvalidParameterValue", true) +} - queueUrl := requestBody.QueueUrl +t := getTenantFromContext(req) +if t == nil { +return utils.CreateErrorResponseV1("InvalidClientTokenId", true) +} - queueName := "" - if queueUrl == "" { - vars := mux.Vars(req) - queueName = vars["queueName"] - } else { - uriSegments := strings.Split(queueUrl, "/") - queueName = uriSegments[len(uriSegments)-1] - } +queueUrl := requestBody.QueueUrl +queueName := "" +if queueUrl == "" { +vars := mux.Vars(req) +queueName = vars["queueName"] +} else { +uriSegments := strings.Split(queueUrl, "/") +queueName = uriSegments[len(uriSegments)-1] +} - if _, ok := models.SyncQueues.Queues[queueName]; !ok { - return utils.CreateErrorResponseV1("QueueNotFound", true) - } +key := tenantQueueKey(t.AccessKey, queueName) - if len(requestBody.Entries) == 0 { - return utils.CreateErrorResponseV1("EmptyBatchRequest", true) - } +if _, ok := models.SyncQueues.Queues[key]; !ok { +return utils.CreateErrorResponseV1("QueueNotFound", true) +} - if len(requestBody.Entries) > 10 { - return utils.CreateErrorResponseV1("TooManyEntriesInBatchRequest", true) - } +if len(requestBody.Entries) == 0 { +return utils.CreateErrorResponseV1("EmptyBatchRequest", true) +} - ids := map[string]bool{} - for _, v := range requestBody.Entries { - if _, found := ids[v.Id]; found { - return utils.CreateErrorResponseV1("BatchEntryIdsNotDistinct", true) - } - ids[v.Id] = true - } +if len(requestBody.Entries) > 10 { +return utils.CreateErrorResponseV1("TooManyEntriesInBatchRequest", true) +} - models.SyncQueues.Lock() - defer models.SyncQueues.Unlock() +ids := map[string]bool{} +for _, v := range requestBody.Entries { +if _, found := ids[v.Id]; found { +return utils.CreateErrorResponseV1("BatchEntryIdsNotDistinct", true) +} +ids[v.Id] = true +} - // create deleteMessageMap - deleteMessageMap := make(map[string]*deleteEntry) - for _, entry := range requestBody.Entries { - deleteMessageMap[entry.ReceiptHandle] = &deleteEntry{ - Id: entry.Id, - ReceiptHandle: entry.ReceiptHandle, - Deleted: false, - } - } +models.SyncQueues.Lock() +defer models.SyncQueues.Unlock() - deletedEntries := make([]models.DeleteMessageBatchResultEntry, 0) - // create a slice to hold messages that are not deleted - remainingMessages := make([]models.SqsMessage, 0, len(models.SyncQueues.Queues[queueName].Messages)) +deleteMessageMap := make(map[string]*deleteEntry) +for _, entry := range requestBody.Entries { +deleteMessageMap[entry.ReceiptHandle] = &deleteEntry{ +Id: entry.Id, +ReceiptHandle: entry.ReceiptHandle, +Deleted: false, +} +} - // delete message from queue - for _, message := range models.SyncQueues.Queues[queueName].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) - } - } +deletedEntries := make([]models.DeleteMessageBatchResultEntry, 0) +remainingMessages := make([]models.SqsMessage, 0, len(models.SyncQueues.Queues[key].Messages)) - // Update the queue with the remaining mesages - models.SyncQueues.Queues[queueName].Messages = remainingMessages +for _, message := range models.SyncQueues.Queues[key].Messages { +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 - 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, - }) - } - } +models.SyncQueues.Queues[key].Messages = remainingMessages - respStruct := models.DeleteMessageBatchResponse{ - Xmlns: models.BaseXmlns, - Result: models.DeleteMessageBatchResult{ - Successful: deletedEntries, - Failed: notFoundEntries, - }, - Metadata: models.BaseResponseMetadata, - } +notFoundEntries := make([]models.BatchResultErrorEntry, 0) +for _, de := range deleteMessageMap { +if !de.Deleted { +notFoundEntries = append(notFoundEntries, models.BatchResultErrorEntry{ +Code: "1", +Id: de.Id, +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 { - Id string - ReceiptHandle string - Error string - Deleted bool +Id string +ReceiptHandle string +Error string +Deleted bool } diff --git a/shared-sqs/app/gosqs/delete_queue.go b/shared-sqs/app/gosqs/delete_queue.go index 6f21d80..7ca36e3 100644 --- a/shared-sqs/app/gosqs/delete_queue.go +++ b/shared-sqs/app/gosqs/delete_queue.go @@ -1,37 +1,43 @@ +// Изменено: 2026-04-09 +// DeleteQueueV1 — удаляет очередь тенанта по tenant-scoped ключу. package gosqs import ( - "net/http" - "strings" +"net/http" +"strings" - "shared-sqs/app/interfaces" - - "shared-sqs/app/models" - "shared-sqs/app/utils" - - log "github.com/sirupsen/logrus" +"shared-sqs/app/interfaces" +"shared-sqs/app/models" +"shared-sqs/app/utils" +log "github.com/sirupsen/logrus" ) func DeleteQueueV1(req *http.Request) (int, interfaces.AbstractResponseBody) { - requestBody := models.NewDeleteQueueRequest() - ok := utils.REQUEST_TRANSFORMER(requestBody, req, false) - if !ok { - log.Error("Invalid Request - DeleteQueueV1") - return utils.CreateErrorResponseV1("InvalidParameterValue", true) - } - - uriSegments := strings.Split(requestBody.QueueUrl, "/") - queueName := uriSegments[len(uriSegments)-1] - - log.Infof("Deleting Queue: %s", queueName) - - models.SyncQueues.Lock() - delete(models.SyncQueues.Queues, queueName) - models.SyncQueues.Unlock() - - respStruct := models.DeleteQueueResponse{ - Xmlns: models.BaseXmlns, - Metadata: models.BaseResponseMetadata, - } - return http.StatusOK, respStruct +requestBody := models.NewDeleteQueueRequest() +ok := utils.REQUEST_TRANSFORMER(requestBody, req, false) +if !ok { +log.Error("Invalid Request - DeleteQueueV1") +return utils.CreateErrorResponseV1("InvalidParameterValue", true) +} + +t := getTenantFromContext(req) +if t == nil { +return utils.CreateErrorResponseV1("InvalidClientTokenId", true) +} + +uriSegments := strings.Split(requestBody.QueueUrl, "/") +queueName := uriSegments[len(uriSegments)-1] +key := tenantQueueKey(t.AccessKey, queueName) + +log.Infof("Deleting Queue: %s (tenant: %s)", queueName, t.ID) + +models.SyncQueues.Lock() +delete(models.SyncQueues.Queues, key) +models.SyncQueues.Unlock() + +respStruct := models.DeleteQueueResponse{ +Xmlns: models.BaseXmlns, +Metadata: models.BaseResponseMetadata, +} +return http.StatusOK, respStruct } diff --git a/shared-sqs/app/gosqs/get_queue_attributes.go b/shared-sqs/app/gosqs/get_queue_attributes.go index 2e0d95d..c0f9cba 100644 --- a/shared-sqs/app/gosqs/get_queue_attributes.go +++ b/shared-sqs/app/gosqs/get_queue_attributes.go @@ -1,135 +1,118 @@ +// Изменено: 2026-04-09 +// GetQueueAttributesV1 — возвращает атрибуты очереди тенанта. package gosqs import ( - "fmt" - "net/http" - "strconv" - "strings" +"fmt" +"net/http" +"strconv" +"strings" - "shared-sqs/app/models" - "shared-sqs/app/utils" - "github.com/mitchellh/copystructure" - - "shared-sqs/app/interfaces" - - log "github.com/sirupsen/logrus" +"shared-sqs/app/interfaces" +"shared-sqs/app/models" +"shared-sqs/app/utils" +"github.com/mitchellh/copystructure" +log "github.com/sirupsen/logrus" ) func GetQueueAttributesV1(req *http.Request) (int, interfaces.AbstractResponseBody) { - requestBody := models.NewGetQueueAttributesRequest() - ok := utils.REQUEST_TRANSFORMER(requestBody, req, false) - if !ok { - log.Error("Invalid Request - GetQueueAttributesV1") - return utils.CreateErrorResponseV1("InvalidParameterValue", true) - } - if requestBody.QueueUrl == "" { - log.Error("Missing QueueUrl - GetQueueAttributesV1") - return utils.CreateErrorResponseV1("InvalidParameterValue", true) - } - - requestedAttributes := func() map[string]bool { - attrs := map[string]bool{} - if len(requestBody.AttributeNames) == 0 { - return map[string]bool{"All": true} - } - for _, attr := range requestBody.AttributeNames { - if "All" == attr { - return map[string]bool{"All": true} - } - attrs[attr] = true - } - return attrs - }() - - dupe, _ := copystructure.Copy(models.AvailableQueueAttributes) - includedAttributes, _ := dupe.(map[string]bool) - _, ok = requestedAttributes["All"] - if !ok { - for attr, _ := range includedAttributes { - _, ok := requestedAttributes[attr] - if !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) - - models.SyncQueues.RLock() - defer models.SyncQueues.RUnlock() - queue, ok := models.SyncQueues.Queues[queueName] - if !ok { - log.Errorf("Get Queue URL: %s queue does not exist!!!", queueName) - return utils.CreateErrorResponseV1("InvalidParameterValue", true) - } - - if _, ok := includedAttributes["DelaySeconds"]; ok { - attr := models.Attribute{Name: "DelaySeconds", Value: strconv.Itoa(queue.DelaySeconds)} - queueAttributes = append(queueAttributes, attr) - } - if _, ok := includedAttributes["MaximumMessageSize"]; ok { - attr := models.Attribute{Name: "MaximumMessageSize", Value: strconv.Itoa(queue.MaximumMessageSize)} - queueAttributes = append(queueAttributes, attr) - } - if _, ok := includedAttributes["MessageRetentionPeriod"]; ok { - attr := models.Attribute{Name: "MessageRetentionPeriod", Value: strconv.Itoa(queue.MessageRetentionPeriod)} - queueAttributes = append(queueAttributes, attr) - } - if _, ok := includedAttributes["ReceiveMessageWaitTimeSeconds"]; ok { - attr := models.Attribute{Name: "ReceiveMessageWaitTimeSeconds", Value: strconv.Itoa(queue.ReceiveMessageWaitTimeSeconds)} - queueAttributes = append(queueAttributes, attr) - } - if _, ok := includedAttributes["VisibilityTimeout"]; ok { - attr := models.Attribute{Name: "VisibilityTimeout", Value: strconv.Itoa(queue.VisibilityTimeout)} - queueAttributes = append(queueAttributes, attr) - } - if _, ok := includedAttributes["ApproximateNumberOfMessages"]; ok { - attr := models.Attribute{Name: "ApproximateNumberOfMessages", Value: strconv.Itoa(len(queue.Messages))} - queueAttributes = append(queueAttributes, attr) - } - // TODO - implement - //if _, ok := includedAttributes["ApproximateNumberOfMessagesDelayed"]; ok { - // attr := models.Attribute{Name: "ApproximateNumberOfMessagesDelayed", Value: strconv.Itoa(len(queue.Messages))} - // queueAttributes = append(queueAttributes, attr) - //} - if _, ok := includedAttributes["ApproximateNumberOfMessagesNotVisible"]; ok { - attr := models.Attribute{Name: "ApproximateNumberOfMessagesNotVisible", Value: strconv.Itoa(numberOfHiddenMessagesInQueue(*queue))} - queueAttributes = append(queueAttributes, attr) - } - if _, ok := includedAttributes["CreatedTimestamp"]; ok { - attr := models.Attribute{Name: "CreatedTimestamp", Value: "0000000000"} - queueAttributes = append(queueAttributes, attr) - } - if _, ok := includedAttributes["LastModifiedTimestamp"]; ok { - attr := models.Attribute{Name: "LastModifiedTimestamp", Value: "0000000000"} - queueAttributes = append(queueAttributes, attr) - } - if _, ok := includedAttributes["QueueArn"]; ok { - attr := models.Attribute{Name: "QueueArn", Value: queue.Arn} - queueAttributes = append(queueAttributes, attr) - } - // TODO - implement - //if _, ok := includedAttributes["Policy"]; ok { - // attr := models.Attribute{Name: "Policy", Value: ""} - // 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 +requestBody := models.NewGetQueueAttributesRequest() +ok := utils.REQUEST_TRANSFORMER(requestBody, req, false) +if !ok { +log.Error("Invalid Request - GetQueueAttributesV1") +return utils.CreateErrorResponseV1("InvalidParameterValue", true) +} +if requestBody.QueueUrl == "" { +log.Error("Missing QueueUrl - GetQueueAttributesV1") +return utils.CreateErrorResponseV1("InvalidParameterValue", true) +} + +t := getTenantFromContext(req) +if t == nil { +return utils.CreateErrorResponseV1("InvalidClientTokenId", true) +} + +requestedAttributes := func() map[string]bool { +attrs := map[string]bool{} +if len(requestBody.AttributeNames) == 0 { +return map[string]bool{"All": true} +} +for _, attr := range requestBody.AttributeNames { +if "All" == attr { +return map[string]bool{"All": true} +} +attrs[attr] = true +} +return attrs +}() + +dupe, _ := copystructure.Copy(models.AvailableQueueAttributes) +includedAttributes, _ := dupe.(map[string]bool) +_, 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] +key := tenantQueueKey(t.AccessKey, queueName) + +log.Infof("Get Queue Attributes: %s (tenant: %s)", queueName, t.ID) +queueAttributes := make([]models.Attribute, 0) + +models.SyncQueues.RLock() +defer models.SyncQueues.RUnlock() +queue, ok := models.SyncQueues.Queues[key] +if !ok { +log.Errorf("Get Queue Attributes: %s queue does not exist for tenant %s", queueName, t.ID) +return utils.CreateErrorResponseV1("QueueNotFound", true) +} + +if _, ok := includedAttributes["DelaySeconds"]; ok { +queueAttributes = append(queueAttributes, models.Attribute{Name: "DelaySeconds", Value: strconv.Itoa(queue.DelaySeconds)}) +} +if _, ok := includedAttributes["MaximumMessageSize"]; ok { +queueAttributes = append(queueAttributes, models.Attribute{Name: "MaximumMessageSize", Value: strconv.Itoa(queue.MaximumMessageSize)}) +} +if _, ok := includedAttributes["MessageRetentionPeriod"]; ok { +queueAttributes = append(queueAttributes, models.Attribute{Name: "MessageRetentionPeriod", Value: strconv.Itoa(queue.MessageRetentionPeriod)}) +} +if _, ok := includedAttributes["ReceiveMessageWaitTimeSeconds"]; ok { +queueAttributes = append(queueAttributes, models.Attribute{Name: "ReceiveMessageWaitTimeSeconds", Value: strconv.Itoa(queue.ReceiveMessageWaitTimeSeconds)}) +} +if _, ok := includedAttributes["VisibilityTimeout"]; ok { +queueAttributes = append(queueAttributes, models.Attribute{Name: "VisibilityTimeout", Value: strconv.Itoa(queue.VisibilityTimeout)}) +} +if _, ok := includedAttributes["ApproximateNumberOfMessages"]; ok { +queueAttributes = append(queueAttributes, models.Attribute{Name: "ApproximateNumberOfMessages", Value: strconv.Itoa(len(queue.Messages))}) +} +if _, ok := includedAttributes["ApproximateNumberOfMessagesNotVisible"]; ok { +queueAttributes = append(queueAttributes, models.Attribute{Name: "ApproximateNumberOfMessagesNotVisible", Value: strconv.Itoa(numberOfHiddenMessagesInQueue(*queue))}) +} +if _, ok := includedAttributes["CreatedTimestamp"]; ok { +queueAttributes = append(queueAttributes, models.Attribute{Name: "CreatedTimestamp", Value: "0000000000"}) +} +if _, ok := includedAttributes["LastModifiedTimestamp"]; ok { +queueAttributes = append(queueAttributes, models.Attribute{Name: "LastModifiedTimestamp", Value: "0000000000"}) +} +if _, ok := includedAttributes["QueueArn"]; ok { +queueAttributes = append(queueAttributes, models.Attribute{Name: "QueueArn", Value: queue.Arn}) +} +if _, ok := includedAttributes["RedrivePolicy"]; ok && queue.DeadLetterQueue != nil { +queueAttributes = append(queueAttributes, models.Attribute{ +Name: "RedrivePolicy", +Value: fmt.Sprintf(`{"maxReceiveCount":"%d", "deadLetterTargetArn":"%s"}`, queue.MaxReceiveCount, queue.DeadLetterQueue.Arn), +}) +} + +respStruct := models.GetQueueAttributesResponse{ +Xmlns: models.BaseXmlns, +Result: models.GetQueueAttributesResult{Attrs: queueAttributes}, +Metadata: models.BaseResponseMetadata, +} +return http.StatusOK, respStruct } diff --git a/shared-sqs/app/gosqs/get_queue_url.go b/shared-sqs/app/gosqs/get_queue_url.go index 258724b..6a243f0 100644 --- a/shared-sqs/app/gosqs/get_queue_url.go +++ b/shared-sqs/app/gosqs/get_queue_url.go @@ -1,36 +1,44 @@ +// Изменено: 2026-04-09 +// GetQueueUrlV1 — возвращает URL очереди тенанта по имени. package gosqs import ( - "net/http" +"net/http" - "shared-sqs/app/interfaces" - "shared-sqs/app/models" - "shared-sqs/app/utils" - log "github.com/sirupsen/logrus" +"shared-sqs/app/interfaces" +"shared-sqs/app/models" +"shared-sqs/app/utils" +log "github.com/sirupsen/logrus" ) func GetQueueUrlV1(req *http.Request) (int, interfaces.AbstractResponseBody) { - requestBody := models.NewGetQueueUrlRequest() - ok := utils.REQUEST_TRANSFORMER(requestBody, req, false) - if !ok { - log.Error("Invalid Request - GetQueueUrlV1") - return utils.CreateErrorResponseV1("InvalidParameterValue", true) - } - - queueName := requestBody.QueueName - if _, ok := models.SyncQueues.Queues[queueName]; !ok { - log.Error("Get Queue URL:", queueName, ", queue does not exist!!!") - return utils.CreateErrorResponseV1("QueueNotFound", true) - } - - queue := models.SyncQueues.Queues[queueName] - log.Debug("Get Queue URL:", queue.Name) - - result := models.GetQueueUrlResult{QueueUrl: queue.URL} - respStruct := models.GetQueueUrlResponse{ - Xmlns: models.BaseXmlns, - Result: result, - Metadata: models.BaseResponseMetadata, - } - return http.StatusOK, respStruct +requestBody := models.NewGetQueueUrlRequest() +ok := utils.REQUEST_TRANSFORMER(requestBody, req, false) +if !ok { +log.Error("Invalid Request - GetQueueUrlV1") +return utils.CreateErrorResponseV1("InvalidParameterValue", true) +} + +t := getTenantFromContext(req) +if t == nil { +return utils.CreateErrorResponseV1("InvalidClientTokenId", true) +} + +queueName := requestBody.QueueName +key := tenantQueueKey(t.AccessKey, queueName) + +if _, ok := models.SyncQueues.Queues[key]; !ok { +log.Errorf("Get Queue URL: %s, queue does not exist for tenant %s", queueName, t.ID) +return utils.CreateErrorResponseV1("QueueNotFound", true) +} + +queue := models.SyncQueues.Queues[key] +log.Debug("Get Queue URL:", queue.Name) + +respStruct := models.GetQueueUrlResponse{ +Xmlns: models.BaseXmlns, +Result: models.GetQueueUrlResult{QueueUrl: queue.URL}, +Metadata: models.BaseResponseMetadata, +} +return http.StatusOK, respStruct } diff --git a/shared-sqs/app/gosqs/list_queues.go b/shared-sqs/app/gosqs/list_queues.go index ecf613f..8bb098f 100644 --- a/shared-sqs/app/gosqs/list_queues.go +++ b/shared-sqs/app/gosqs/list_queues.go @@ -1,45 +1,51 @@ +// Изменено: 2026-04-09 +// ListQueuesV1 — возвращает только очереди текущего тенанта. +// Изоляция: фильтруем SyncQueues по префиксу "{tenantAccessKey}:". package gosqs import ( - "net/http" - "strings" +"net/http" +"strings" - "shared-sqs/app/utils" - - "shared-sqs/app/models" - - "shared-sqs/app/interfaces" - log "github.com/sirupsen/logrus" +"shared-sqs/app/interfaces" +"shared-sqs/app/models" +"shared-sqs/app/utils" +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) { - requestBody := models.NewListQueuesRequest() - ok := utils.REQUEST_TRANSFORMER(requestBody, req, true) - if !ok { - log.Error("Invalid Request - ListQueuesV1") - return utils.CreateErrorResponseV1("InvalidParameterValue", true) - } - - log.Info("Listing Queues") - queueUrls := make([]string, 0) - models.SyncQueues.Lock() - for _, queue := range models.SyncQueues.Queues { - if strings.HasPrefix(queue.Name, requestBody.QueueNamePrefix) { - queueUrls = append(queueUrls, queue.URL) - } - } - models.SyncQueues.Unlock() - - respStruct := models.ListQueuesResponse{ - Xmlns: models.BaseXmlns, - Metadata: models.BaseResponseMetadata, - Result: models.ListQueuesResult{ - QueueUrls: queueUrls, - }, - } - - return http.StatusOK, respStruct +requestBody := models.NewListQueuesRequest() +ok := utils.REQUEST_TRANSFORMER(requestBody, req, true) +if !ok { +log.Error("Invalid Request - ListQueuesV1") +return utils.CreateErrorResponseV1("InvalidParameterValue", true) +} + +t := getTenantFromContext(req) +if t == nil { +return utils.CreateErrorResponseV1("InvalidClientTokenId", true) +} + +log.Infof("Listing Queues for tenant: %s", t.ID) +queueUrls := make([]string, 0) +prefix := t.AccessKey + ":" + +models.SyncQueues.Lock() +for key, queue := range models.SyncQueues.Queues { +// Показываем только очереди этого тенанта +if strings.HasPrefix(key, prefix) { +if strings.HasPrefix(queue.Name, requestBody.QueueNamePrefix) { +queueUrls = append(queueUrls, queue.URL) +} +} +} +models.SyncQueues.Unlock() + +respStruct := models.ListQueuesResponse{ +Xmlns: models.BaseXmlns, +Metadata: models.BaseResponseMetadata, +Result: models.ListQueuesResult{QueueUrls: queueUrls}, +} + +return http.StatusOK, respStruct } diff --git a/shared-sqs/app/gosqs/purge_queue.go b/shared-sqs/app/gosqs/purge_queue.go index 31b80a9..9f6e018 100644 --- a/shared-sqs/app/gosqs/purge_queue.go +++ b/shared-sqs/app/gosqs/purge_queue.go @@ -1,42 +1,49 @@ +// Изменено: 2026-04-09 +// PurgeQueueV1 — очищает все сообщения в очереди тенанта. package gosqs import ( - "net/http" - "strings" - "time" +"net/http" +"strings" +"time" - "shared-sqs/app/interfaces" - "shared-sqs/app/models" - "shared-sqs/app/utils" - - log "github.com/sirupsen/logrus" +"shared-sqs/app/interfaces" +"shared-sqs/app/models" +"shared-sqs/app/utils" +log "github.com/sirupsen/logrus" ) func PurgeQueueV1(req *http.Request) (int, interfaces.AbstractResponseBody) { - requestBody := models.NewPurgeQueueRequest() - ok := utils.REQUEST_TRANSFORMER(requestBody, req, false) - if !ok { - log.Error("Invalid Request - PurgeQueueV1") - return utils.CreateErrorResponseV1("InvalidParameterValue", true) - } - - uriSegments := strings.Split(requestBody.QueueUrl, "/") - queueName := uriSegments[len(uriSegments)-1] - - models.SyncQueues.Lock() - defer models.SyncQueues.Unlock() - if _, ok := models.SyncQueues.Queues[queueName]; !ok { - log.Errorf("Purge Queue: %s, queue does not exist!!!", queueName) - return utils.CreateErrorResponseV1("QueueNotFound", true) - } - - log.Infof("Purging Queue: %s", queueName) - models.SyncQueues.Queues[queueName].Messages = nil - models.SyncQueues.Queues[queueName].Duplicates = make(map[string]time.Time) - - respStruct := models.PurgeQueueResponse{ - Xmlns: models.BaseXmlns, - Metadata: models.BaseResponseMetadata, - } - return http.StatusOK, respStruct +requestBody := models.NewPurgeQueueRequest() +ok := utils.REQUEST_TRANSFORMER(requestBody, req, false) +if !ok { +log.Error("Invalid Request - PurgeQueueV1") +return utils.CreateErrorResponseV1("InvalidParameterValue", true) +} + +t := getTenantFromContext(req) +if t == nil { +return utils.CreateErrorResponseV1("InvalidClientTokenId", true) +} + +uriSegments := strings.Split(requestBody.QueueUrl, "/") +queueName := uriSegments[len(uriSegments)-1] +key := tenantQueueKey(t.AccessKey, queueName) + +models.SyncQueues.Lock() +defer models.SyncQueues.Unlock() +if _, ok := models.SyncQueues.Queues[key]; !ok { +log.Errorf("Purge Queue: %s, queue does not exist for tenant %s", queueName, t.ID) +return utils.CreateErrorResponseV1("QueueNotFound", true) +} + +log.Infof("Purging Queue: %s (tenant: %s)", queueName, t.ID) +models.SyncQueues.Queues[key].Messages = nil +models.SyncQueues.Queues[key].Duplicates = make(map[string]time.Time) + +respStruct := models.PurgeQueueResponse{ +Xmlns: models.BaseXmlns, +Metadata: models.BaseResponseMetadata, +} +return http.StatusOK, respStruct } diff --git a/shared-sqs/app/gosqs/receive_message.go b/shared-sqs/app/gosqs/receive_message.go index 99d4bba..792aa81 100644 --- a/shared-sqs/app/gosqs/receive_message.go +++ b/shared-sqs/app/gosqs/receive_message.go @@ -1,164 +1,168 @@ +// Изменено: 2026-04-09 +// ReceiveMessageV1 — получает сообщения из очереди тенанта с поддержкой long polling. +// Ловушка #4: long polling держит соединение до 20 сек — не прерываем принудительно. package gosqs import ( - "fmt" - "net/http" - "strings" - "time" +"fmt" +"net/http" +"strings" +"time" - "github.com/google/uuid" +"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" +"shared-sqs/app/interfaces" +"shared-sqs/app/models" +"shared-sqs/app/utils" +"github.com/gorilla/mux" +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) { - requestBody := models.NewReceiveMessageRequest() - ok := utils.REQUEST_TRANSFORMER(requestBody, req, false) - if !ok { - log.Error("Invalid Request - ReceiveMessageV1") - return utils.CreateErrorResponseV1("InvalidParameterValue", true) - } +requestBody := models.NewReceiveMessageRequest() +ok := utils.REQUEST_TRANSFORMER(requestBody, req, false) +if !ok { +log.Error("Invalid Request - ReceiveMessageV1") +return utils.CreateErrorResponseV1("InvalidParameterValue", true) +} - maxNumberOfMessages := requestBody.MaxNumberOfMessages - if maxNumberOfMessages == 0 { - maxNumberOfMessages = 1 - } +t := getTenantFromContext(req) +if t == nil { +return utils.CreateErrorResponseV1("InvalidClientTokenId", true) +} - queueName := "" - if requestBody.QueueUrl == "" { - vars := mux.Vars(req) - queueName = vars["queueName"] - } else { - uriSegments := strings.Split(requestBody.QueueUrl, "/") - queueName = uriSegments[len(uriSegments)-1] - } +maxNumberOfMessages := requestBody.MaxNumberOfMessages +if maxNumberOfMessages == 0 { +maxNumberOfMessages = 1 +} - if _, ok := models.SyncQueues.Queues[queueName]; !ok { - return utils.CreateErrorResponseV1("QueueNotFound", true) - } +queueName := "" +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 - respStruct := models.ReceiveMessageResponse{} +key := tenantQueueKey(t.AccessKey, queueName) - waitTimeSeconds := requestBody.WaitTimeSeconds - if waitTimeSeconds == 0 { - models.SyncQueues.RLock() - waitTimeSeconds = models.SyncQueues.Queues[queueName].ReceiveMessageWaitTimeSeconds - models.SyncQueues.RUnlock() - } +if _, ok := models.SyncQueues.Queues[key]; !ok { +return utils.CreateErrorResponseV1("QueueNotFound", true) +} - loops := waitTimeSeconds * 10 - for loops > 0 { - 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 - } +var messages []*models.ResultMessage +respStruct := models.ReceiveMessageResponse{} - } - log.Debugf("Getting Message from Queue:%s", queueName) +waitTimeSeconds := requestBody.WaitTimeSeconds +if waitTimeSeconds == 0 { +models.SyncQueues.RLock() +waitTimeSeconds = models.SyncQueues.Queues[key].ReceiveMessageWaitTimeSeconds +models.SyncQueues.RUnlock() +} - models.SyncQueues.Lock() // Lock the Queues - defer models.SyncQueues.Unlock() // Unlock the Queues +// 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) - if len(models.SyncQueues.Queues[queueName].Messages) > 0 { - numMsg := 0 - messages = make([]*models.ResultMessage, 0) - for i := range models.SyncQueues.Queues[queueName].Messages { - if numMsg >= maxNumberOfMessages { - break - } +models.SyncQueues.Lock() +defer models.SyncQueues.Unlock() - if models.SyncQueues.Queues[queueName].Messages[i].ReceiptHandle != "" { - continue - } +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 +} - msg := &models.SyncQueues.Queues[queueName].Messages[i] - if !msg.IsReadyForReceipt() { - continue - } +if models.SyncQueues.Queues[key].Messages[i].ReceiptHandle != "" { +continue +} - if models.SyncQueues.Queues[queueName].IsFIFO { - // If we got messages here it means we have not processed it yet, so get next - if models.SyncQueues.Queues[queueName].IsLocked(msg.GroupID) { - continue - } - // Otherwise lock messages for group ID - models.SyncQueues.Queues[queueName].LockGroup(msg.GroupID) - } +msg := &models.SyncQueues.Queues[key].Messages[i] +if !msg.IsReadyForReceipt() { +continue +} - randomId := uuid.NewString() - msg.ReceiptHandle = msg.Uuid + "#" + randomId - msg.ReceiptTime = time.Now().UTC() +if models.SyncQueues.Queues[key].IsFIFO { +if models.SyncQueues.Queues[key].IsLocked(msg.GroupID) { +continue +} +models.SyncQueues.Queues[key].LockGroup(msg.GroupID) +} - 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[queueName].VisibilityTimeout) * time.Second) - } +randomId := uuid.NewString() +msg.ReceiptHandle = msg.Uuid + "#" + randomId +msg.ReceiptTime = time.Now().UTC() - 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{ - "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"}} - } +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 +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)), - }, - } +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)), +}, +} } diff --git a/shared-sqs/app/gosqs/send_message.go b/shared-sqs/app/gosqs/send_message.go index 7673c1f..2bf4e4b 100644 --- a/shared-sqs/app/gosqs/send_message.go +++ b/shared-sqs/app/gosqs/send_message.go @@ -1,100 +1,108 @@ +// Изменено: 2026-04-09 +// SendMessageV1 — добавляет сообщение в очередь тенанта. +// Ловушка #6: queueName извлекается как ПОСЛЕДНИЙ сегмент URL — при URL вида +// http://host/tenantID/queueName последний сегмент = queueName (правильно). package gosqs import ( - "net/http" - "strings" - "time" +"net/http" +"strings" +"time" - "github.com/google/uuid" +"github.com/google/uuid" - "shared-sqs/app/interfaces" - "shared-sqs/app/models" +"shared-sqs/app/interfaces" +"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) { - requestBody := models.NewSendMessageRequest() - ok := utils.REQUEST_TRANSFORMER(requestBody, req, false) - if !ok { - log.Error("Invalid Request - SendMessageV1") - return utils.CreateErrorResponseV1("InvalidParameterValue", true) - } - messageBody := requestBody.MessageBody - messageGroupID := requestBody.MessageGroupId - messageDeduplicationID := requestBody.MessageDeduplicationId - - queueUrl := getQueueFromPath(requestBody.QueueUrl, req.URL.String()) - - queueName := "" - if queueUrl == "" { - // TODO: Remove this query param logic if it's not still valid or something - vars := mux.Vars(req) - queueName = vars["queueName"] - } else { - uriSegments := strings.Split(queueUrl, "/") - queueName = uriSegments[len(uriSegments)-1] - } - - if _, ok := models.SyncQueues.Queues[queueName]; !ok { - // Queue does not exist - return utils.CreateErrorResponseV1("QueueNotFound", true) - } - - if models.SyncQueues.Queues[queueName].MaximumMessageSize > 0 && - len(messageBody) > models.SyncQueues.Queues[queueName].MaximumMessageSize { - // Message size is too big - return utils.CreateErrorResponseV1("MessageTooBig", true) - } - - delaySecs := models.SyncQueues.Queues[queueName].DelaySeconds - if requestBody.DelaySeconds != 0 { - delaySecs = requestBody.DelaySeconds - } - - log.Debugf("Putting Message in Queue: [%s]", queueName) - 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[queueName].IsFIFO { - fifoSeqNumber = models.SyncQueues.Queues[queueName].NextSequenceNumber(messageGroupID) - } - - if !models.SyncQueues.Queues[queueName].IsDuplicate(messageDeduplicationID) { - models.SyncQueues.Queues[queueName].Messages = append(models.SyncQueues.Queues[queueName].Messages, msg) - } else { - log.Debugf("Message with deduplicationId [%s] in queue [%s] is duplicate ", messageDeduplicationID, queueName) - } - - models.SyncQueues.Queues[queueName].InitDuplicatation(messageDeduplicationID) - 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 +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 + +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) +} + +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) +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 } diff --git a/shared-sqs/app/gosqs/send_message_batch.go b/shared-sqs/app/gosqs/send_message_batch.go index b70901a..e6a3348 100644 --- a/shared-sqs/app/gosqs/send_message_batch.go +++ b/shared-sqs/app/gosqs/send_message_batch.go @@ -1,105 +1,109 @@ +// Изменено: 2026-04-09 +// SendMessageBatchV1 — пакетная отправка сообщений в очередь тенанта. package gosqs import ( - "net/http" - "strings" - "time" +"net/http" +"strings" +"time" - "github.com/google/uuid" +"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" +"shared-sqs/app/interfaces" +"shared-sqs/app/models" +"shared-sqs/app/utils" +"github.com/gorilla/mux" +log "github.com/sirupsen/logrus" ) func SendMessageBatchV1(req *http.Request) (int, interfaces.AbstractResponseBody) { - requestBody := models.NewSendMessageBatchRequest() - ok := utils.REQUEST_TRANSFORMER(requestBody, req, false) - if !ok { - log.Error("Invalid Request - SendMessageBatchV1") - return utils.CreateErrorResponseV1("InvalidParameterValue", true) - } - - queueUrl := requestBody.QueueUrl - - // TODO: Remove this query param logic if it's not still valid or something - queueName := "" - if queueUrl == "" { - vars := mux.Vars(req) - queueName = vars["queueName"] - } else { - uriSegments := strings.Split(queueUrl, "/") - queueName = uriSegments[len(uriSegments)-1] - } - - if _, ok := models.SyncQueues.Queues[queueName]; !ok { - return utils.CreateErrorResponseV1("QueueNotFound", true) - } - - sendEntries := requestBody.Entries - - if len(sendEntries) == 0 { - return utils.CreateErrorResponseV1("EmptyBatchRequest", true) - } - - if len(sendEntries) > 10 { - return utils.CreateErrorResponseV1("TooManyEntriesInBatchRequest", true) - } - ids := map[string]struct{}{} - for _, v := range sendEntries { - if _, ok := ids[v.Id]; ok { - return utils.CreateErrorResponseV1("BatchEntryIdsNotDistinct", true) - } - ids[v.Id] = struct{}{} - } - - sentEntries := make([]models.SendMessageBatchResultEntry, 0) - log.Debug("Putting Message in Queue:", queueName) - for _, sendEntry := range sendEntries { - msg := models.SqsMessage{MessageBody: sendEntry.MessageBody} - if len(sendEntry.MessageAttributes) > 0 { - msg.MessageAttributes = sendEntry.MessageAttributes - msg.MD5OfMessageAttributes = utils.HashAttributes(sendEntry.MessageAttributes) - } - msg.MD5OfMessageBody = utils.GetMD5Hash(sendEntry.MessageBody) - msg.GroupID = sendEntry.MessageGroupId - msg.DeduplicationID = sendEntry.MessageDeduplicationId - msg.Uuid = uuid.NewString() - msg.SentTime = time.Now() - models.SyncQueues.Lock() - fifoSeqNumber := "" - if models.SyncQueues.Queues[queueName].IsFIFO { - fifoSeqNumber = models.SyncQueues.Queues[queueName].NextSequenceNumber(sendEntry.MessageGroupId) - } - - if !models.SyncQueues.Queues[queueName].IsDuplicate(sendEntry.MessageDeduplicationId) { - models.SyncQueues.Queues[queueName].Messages = append(models.SyncQueues.Queues[queueName].Messages, msg) - } else { - log.Debugf("Message with deduplicationId [%s] in queue [%s] is duplicate ", sendEntry.MessageDeduplicationId, queueName) - } - - models.SyncQueues.Queues[queueName].InitDuplicatation(sendEntry.MessageDeduplicationId) - - models.SyncQueues.Unlock() - se := models.SendMessageBatchResultEntry{ - Id: sendEntry.Id, - MessageId: msg.Uuid, - MD5OfMessageBody: msg.MD5OfMessageBody, - MD5OfMessageAttributes: msg.MD5OfMessageAttributes, - SequenceNumber: fifoSeqNumber, - } - sentEntries = append(sentEntries, se) - log.Infof("%s: Queue: %s, Message: %s\n", time.Now().Format("2006-01-02 15:04:05"), queueName, msg.MessageBody) - } - - respStruct := models.SendMessageBatchResponse{ - Xmlns: models.BaseXmlns, - Result: models.SendMessageBatchResult{Entry: sentEntries}, - Metadata: models.BaseResponseMetadata, - } - - return http.StatusOK, respStruct - +requestBody := models.NewSendMessageBatchRequest() +ok := utils.REQUEST_TRANSFORMER(requestBody, req, false) +if !ok { +log.Error("Invalid Request - SendMessageBatchV1") +return utils.CreateErrorResponseV1("InvalidParameterValue", true) +} + +t := getTenantFromContext(req) +if t == nil { +return utils.CreateErrorResponseV1("InvalidClientTokenId", true) +} + +queueUrl := requestBody.QueueUrl +queueName := "" +if queueUrl == "" { +vars := mux.Vars(req) +queueName = vars["queueName"] +} else { +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) +} + +sendEntries := requestBody.Entries + +if len(sendEntries) == 0 { +return utils.CreateErrorResponseV1("EmptyBatchRequest", true) +} + +if len(sendEntries) > 10 { +return utils.CreateErrorResponseV1("TooManyEntriesInBatchRequest", true) +} +ids := map[string]struct{}{} +for _, v := range sendEntries { +if _, ok := ids[v.Id]; ok { +return utils.CreateErrorResponseV1("BatchEntryIdsNotDistinct", true) +} +ids[v.Id] = struct{}{} +} + +sentEntries := make([]models.SendMessageBatchResultEntry, 0) +log.Debugf("Batch sending to Queue: %s (tenant: %s)", queueName, t.ID) +for _, sendEntry := range sendEntries { +msg := models.SqsMessage{MessageBody: sendEntry.MessageBody} +if len(sendEntry.MessageAttributes) > 0 { +msg.MessageAttributes = sendEntry.MessageAttributes +msg.MD5OfMessageAttributes = utils.HashAttributes(sendEntry.MessageAttributes) +} +msg.MD5OfMessageBody = utils.GetMD5Hash(sendEntry.MessageBody) +msg.GroupID = sendEntry.MessageGroupId +msg.DeduplicationID = sendEntry.MessageDeduplicationId +msg.Uuid = uuid.NewString() +msg.SentTime = time.Now() + +models.SyncQueues.Lock() +fifoSeqNumber := "" +if models.SyncQueues.Queues[key].IsFIFO { +fifoSeqNumber = models.SyncQueues.Queues[key].NextSequenceNumber(sendEntry.MessageGroupId) +} +if !models.SyncQueues.Queues[key].IsDuplicate(sendEntry.MessageDeduplicationId) { +models.SyncQueues.Queues[key].Messages = append(models.SyncQueues.Queues[key].Messages, msg) +} else { +log.Debugf("Duplicate deduplicationId [%s] in queue [%s]", sendEntry.MessageDeduplicationId, queueName) +} +models.SyncQueues.Queues[key].InitDuplicatation(sendEntry.MessageDeduplicationId) +models.SyncQueues.Unlock() + +sentEntries = append(sentEntries, models.SendMessageBatchResultEntry{ +Id: sendEntry.Id, +MessageId: msg.Uuid, +MD5OfMessageBody: msg.MD5OfMessageBody, +MD5OfMessageAttributes: msg.MD5OfMessageAttributes, +SequenceNumber: fifoSeqNumber, +}) +log.Infof("%s: Queue: %s, Message: %s", time.Now().Format("2006-01-02 15:04:05"), queueName, msg.MessageBody) +} + +respStruct := models.SendMessageBatchResponse{ +Xmlns: models.BaseXmlns, +Result: models.SendMessageBatchResult{Entry: sentEntries}, +Metadata: models.BaseResponseMetadata, +} + +return http.StatusOK, respStruct } diff --git a/shared-sqs/app/gosqs/set_queue_attributes.go b/shared-sqs/app/gosqs/set_queue_attributes.go index c9de081..f691585 100644 --- a/shared-sqs/app/gosqs/set_queue_attributes.go +++ b/shared-sqs/app/gosqs/set_queue_attributes.go @@ -1,48 +1,54 @@ +// Изменено: 2026-04-09 +// SetQueueAttributesV1 — устанавливает атрибуты очереди тенанта. +// Ловушка #9: при RedrivePolicy парсим ARN DLQ и DLQ тоже должна принадлежать тому же тенанту. package gosqs import ( - "net/http" - "strings" +"net/http" +"strings" - "shared-sqs/app/models" - "shared-sqs/app/utils" - - "shared-sqs/app/interfaces" - log "github.com/sirupsen/logrus" +"shared-sqs/app/interfaces" +"shared-sqs/app/models" +"shared-sqs/app/utils" +log "github.com/sirupsen/logrus" ) func SetQueueAttributesV1(req *http.Request) (int, interfaces.AbstractResponseBody) { - requestBody := models.NewSetQueueAttributesRequest() - ok := utils.REQUEST_TRANSFORMER(requestBody, req, false) - if !ok { - log.Error("Invalid Request - GetQueueAttributesV1") - return utils.CreateErrorResponseV1("InvalidParameterValue", true) - } - if requestBody.QueueUrl == "" { - log.Error("Missing QueueUrl - GetQueueAttributesV1") - return utils.CreateErrorResponseV1("InvalidParameterValue", true) - } - - // NOTE: I tore out the handling for devining the url from a param. I can't find documentation that - // that is valid any longer. - uriSegments := strings.Split(requestBody.QueueUrl, "/") - queueName := uriSegments[len(uriSegments)-1] - - log.Infof("Set Queue QueueAttributes: %s", queueName) - models.SyncQueues.Lock() - defer models.SyncQueues.Unlock() - queue, ok := models.SyncQueues.Queues[queueName] - if !ok { - log.Warningf("Get Queue URL: %s, queue does not exist!!!", queueName) - return utils.CreateErrorResponseV1("QueueNotFound", true) - } - if err := setQueueAttributesV1(queue, requestBody.Attributes); err != nil { - return utils.CreateErrorResponseV1(err.Error(), true) - } - - respStruct := models.SetQueueAttributesResponse{ - Xmlns: models.BaseXmlns, - Metadata: models.BaseResponseMetadata, - } - return http.StatusOK, respStruct +requestBody := models.NewSetQueueAttributesRequest() +ok := utils.REQUEST_TRANSFORMER(requestBody, req, false) +if !ok { +log.Error("Invalid Request - SetQueueAttributesV1") +return utils.CreateErrorResponseV1("InvalidParameterValue", true) +} +if requestBody.QueueUrl == "" { +log.Error("Missing QueueUrl - SetQueueAttributesV1") +return utils.CreateErrorResponseV1("InvalidParameterValue", true) +} + +t := getTenantFromContext(req) +if t == nil { +return utils.CreateErrorResponseV1("InvalidClientTokenId", true) +} + +uriSegments := strings.Split(requestBody.QueueUrl, "/") +queueName := uriSegments[len(uriSegments)-1] +key := tenantQueueKey(t.AccessKey, queueName) + +log.Infof("Set Queue Attributes: %s (tenant: %s)", queueName, t.ID) +models.SyncQueues.Lock() +defer models.SyncQueues.Unlock() +queue, ok := models.SyncQueues.Queues[key] +if !ok { +log.Warningf("Set Queue Attributes: %s, queue does not exist for tenant %s", queueName, t.ID) +return utils.CreateErrorResponseV1("QueueNotFound", true) +} +if err := setQueueAttributesV1(queue, requestBody.Attributes); err != nil { +return utils.CreateErrorResponseV1(err.Error(), true) +} + +respStruct := models.SetQueueAttributesResponse{ +Xmlns: models.BaseXmlns, +Metadata: models.BaseResponseMetadata, +} +return http.StatusOK, respStruct } diff --git a/shared-sqs/app/gosqs/tenant_helpers.go b/shared-sqs/app/gosqs/tenant_helpers.go new file mode 100644 index 0000000..f1f6337 --- /dev/null +++ b/shared-sqs/app/gosqs/tenant_helpers.go @@ -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 +} diff --git a/shared-sqs/app/tenant/tenant_store.go b/shared-sqs/app/tenant/tenant_store.go new file mode 100644 index 0000000..3ce75b1 --- /dev/null +++ b/shared-sqs/app/tenant/tenant_store.go @@ -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-) +Name string // человекочитаемое имя +AccessKey string // аналог AWS AccessKeyId (SSAK-) +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 +}