diff --git a/doc/thinking/2026-04-10.md b/doc/thinking/2026-04-10.md new file mode 100644 index 0000000..b359b11 --- /dev/null +++ b/doc/thinking/2026-04-10.md @@ -0,0 +1,52 @@ +# Thinking Log — 2026-04-10 +# Agent: GitHub Copilot (Claude Sonnet 4.6) + +--- + +## Сессия 1 + +### Задача +1. Задокументировать итоги работы над shared-sqs (v0.1.11–v0.1.14) +2. Закоммитить и запушить все изменения +3. Найти тесты харбора и прогнать нагрузочно после апгрейда ресурсов + +### Контекст (из предыдущих сессий) + +#### Что было сделано над shared-sqs: +- **v0.1.11** — Redis write-through persistence (очереди и сообщения сохраняются при рестарте) +- **v0.1.12** — промежуточный билд +- **v0.1.13** — КРИТИЧЕСКИЙ фикс дедлока в `create_queue.go`: `SyncQueues.Lock()` захватывался без `Unlock()` в happy path, из-за чего после первого успешного CreateQueue сервис замирал навсегда +- **v0.1.14** — фикс UI: JS читал поле `m.sent`, API отдавал `m.sent_at` → даты сообщений всегда показывались как `—` + +#### Статус тестирования: +- 23/23 PASS — суровые тесты с ВМ (наeel@5.172.178.213) +- 6/6 PASS — quick_test.sh из публичной gitea репы Nail/shared-SQS + +#### Важный вывод о продукте: +Аналогов нет. GitHub search `multi-tenant sqs compatible` → 0 результатов. +Ближайшее: ElasticMQ (single-tenant, local dev only) и GoAws (то же самое). +shared-sqs занимает нишу "SQS-as-a-Service для private cloud" — её в open source нет. + +### Изменённые файлы в текущем коммите: +- `app/gosqs/create_queue.go` — фикс дедлока (Unlock перед return в happy path) +- `app/gosqs/delete_queue.go` — рефакторинг под новую модель с Redis +- `app/gosqs/purge_queue.go` — то же +- `app/gosqs/send_message.go` — то же +- `app/gosqs/set_queue_attributes.go` — то же +- `app/router/router.go` — маршруты +- `app/ui/index.html` — фикс `m.sent` → `m.sent_at` +- `deployments/k8s/deployment.yaml` — образ v0.1.14 +- `deployments/k8s/ingress.yaml` — TLS endpoint qu.kube5s.ru +- `deployments/k8s/redis.yaml` — новый: деплой Redis в кластере + +### Исправленная ошибка агента +Агент пытался выполнять команды (git, bash) локально через терминал. +**ПРАВИЛО**: `/home/naeel/remote_dev/sless` — это sshfs-mount. +Все файлы физически на ВМ `naeel@5.172.178.213:/home/naeel/terra/sless`. +Все команды — ТОЛЬКО через SSH на ВМ. + +### План на сессию +1. ✅ Написать thinking log +2. Закоммитить изменения shared-sqs на ВМ +3. Найти `test_harbor_load.sh` в корне проекта, изучить +4. Прогнать нагрузочный тест харбора с ВМ, сравнить с предыдущими результатами diff --git a/shared-sqs/app/gosqs/create_queue.go b/shared-sqs/app/gosqs/create_queue.go index 45aa23f..357d27b 100644 --- a/shared-sqs/app/gosqs/create_queue.go +++ b/shared-sqs/app/gosqs/create_queue.go @@ -4,63 +4,65 @@ package gosqs import ( -"net/http" -"time" + "net/http" + "time" -"shared-sqs/app/interfaces" -"shared-sqs/app/models" -"shared-sqs/app/persistence" -"shared-sqs/app/utils" -log "github.com/sirupsen/logrus" + "shared-sqs/app/interfaces" + "shared-sqs/app/models" + "shared-sqs/app/persistence" + "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) -} + 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) -} + 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) + // Ловушка #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.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 + } // Сохраняем очередь в Redis пока держим Lock — консистентный снапшот persistence.SaveQueue(key, models.SyncQueues.Queues[key]) -respStruct := models.CreateQueueResponse{ -Xmlns: models.BaseXmlns, -Result: models.CreateQueueResult{QueueUrl: queueUrl}, -Metadata: models.BaseResponseMetadata, -} -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 } diff --git a/shared-sqs/app/gosqs/delete_queue.go b/shared-sqs/app/gosqs/delete_queue.go index 03b5283..674181f 100644 --- a/shared-sqs/app/gosqs/delete_queue.go +++ b/shared-sqs/app/gosqs/delete_queue.go @@ -3,45 +3,46 @@ package gosqs import ( -"net/http" -"strings" + "net/http" + "strings" -"shared-sqs/app/interfaces" -"shared-sqs/app/models" -"shared-sqs/app/persistence" -"shared-sqs/app/utils" -log "github.com/sirupsen/logrus" + "shared-sqs/app/interfaces" + "shared-sqs/app/models" + "shared-sqs/app/persistence" + "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) -} - -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() - -// Удаляем из Redis асинхронно -persistence.DeleteQueue(key) - -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() + + // Удаляем из Redis асинхронно + persistence.DeleteQueue(key) + + respStruct := models.DeleteQueueResponse{ + Xmlns: models.BaseXmlns, + Metadata: models.BaseResponseMetadata, + } + return http.StatusOK, respStruct } diff --git a/shared-sqs/app/gosqs/purge_queue.go b/shared-sqs/app/gosqs/purge_queue.go index 0a56a01..e661640 100644 --- a/shared-sqs/app/gosqs/purge_queue.go +++ b/shared-sqs/app/gosqs/purge_queue.go @@ -3,50 +3,51 @@ package gosqs import ( -"net/http" -"strings" -"time" + "net/http" + "strings" + "time" -"shared-sqs/app/interfaces" -"shared-sqs/app/models" -"shared-sqs/app/persistence" -"shared-sqs/app/utils" -log "github.com/sirupsen/logrus" + "shared-sqs/app/interfaces" + "shared-sqs/app/models" + "shared-sqs/app/persistence" + "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) -} + 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) -} + 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) + 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) -} + 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) -// Сохраняем пустую очередь в Redis пока держим Lock -persistence.SaveQueue(key, models.SyncQueues.Queues[key]) + 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) + // Сохраняем пустую очередь в Redis пока держим Lock + persistence.SaveQueue(key, models.SyncQueues.Queues[key]) -respStruct := models.PurgeQueueResponse{ -Xmlns: models.BaseXmlns, -Metadata: models.BaseResponseMetadata, -} -return http.StatusOK, respStruct + respStruct := models.PurgeQueueResponse{ + Xmlns: models.BaseXmlns, + Metadata: models.BaseResponseMetadata, + } + return http.StatusOK, respStruct } diff --git a/shared-sqs/app/gosqs/send_message.go b/shared-sqs/app/gosqs/send_message.go index 2b40fc1..d0ea550 100644 --- a/shared-sqs/app/gosqs/send_message.go +++ b/shared-sqs/app/gosqs/send_message.go @@ -5,107 +5,107 @@ 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/persistence" -"shared-sqs/app/utils" + "shared-sqs/app/interfaces" + "shared-sqs/app/models" + "shared-sqs/app/persistence" + "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) -} - -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) -// Сохраняем очередь в Redis пока держим Lock -persistence.SaveQueue(key, models.SyncQueues.Queues[key]) -models.SyncQueues.Unlock() -log.Infof("%s: Queue: %s, Message: %s\n", time.Now().Format("2006-01-02 15:04:05"), queueName, msg.MessageBody) - -respStruct := models.SendMessageResponse{ -Xmlns: models.BaseXmlns, -Result: models.SendMessageResult{ -MD5OfMessageAttributes: msg.MD5OfMessageAttributes, -MD5OfMessageBody: msg.MD5OfMessageBody, -MessageId: msg.Uuid, -SequenceNumber: fifoSeqNumber, -}, -Metadata: models.BaseResponseMetadata, -} - -return http.StatusOK, respStruct + 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) + // Сохраняем очередь в Redis пока держим Lock + persistence.SaveQueue(key, models.SyncQueues.Queues[key]) + models.SyncQueues.Unlock() + log.Infof("%s: Queue: %s, Message: %s\n", time.Now().Format("2006-01-02 15:04:05"), queueName, msg.MessageBody) + + respStruct := models.SendMessageResponse{ + Xmlns: models.BaseXmlns, + Result: models.SendMessageResult{ + MD5OfMessageAttributes: msg.MD5OfMessageAttributes, + MD5OfMessageBody: msg.MD5OfMessageBody, + MessageId: msg.Uuid, + SequenceNumber: fifoSeqNumber, + }, + Metadata: models.BaseResponseMetadata, + } + + return http.StatusOK, respStruct } diff --git a/shared-sqs/app/gosqs/set_queue_attributes.go b/shared-sqs/app/gosqs/set_queue_attributes.go index d150103..fb6420f 100644 --- a/shared-sqs/app/gosqs/set_queue_attributes.go +++ b/shared-sqs/app/gosqs/set_queue_attributes.go @@ -4,54 +4,55 @@ package gosqs import ( -"net/http" -"strings" + "net/http" + "strings" -"shared-sqs/app/interfaces" -"shared-sqs/app/models" -"shared-sqs/app/persistence" -"shared-sqs/app/utils" -log "github.com/sirupsen/logrus" + "shared-sqs/app/interfaces" + "shared-sqs/app/models" + "shared-sqs/app/persistence" + "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 - SetQueueAttributesV1") -return utils.CreateErrorResponseV1("InvalidParameterValue", true) -} -if requestBody.QueueUrl == "" { -log.Error("Missing QueueUrl - SetQueueAttributesV1") -return utils.CreateErrorResponseV1("InvalidParameterValue", true) -} + 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) -} + 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) + 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) -} -// Сохраняем атрибуты в Redis пока держим Lock (через defer) -persistence.SaveQueue(key, queue) + 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) + } + // Сохраняем атрибуты в Redis пока держим Lock (через defer) + persistence.SaveQueue(key, queue) -respStruct := models.SetQueueAttributesResponse{ -Xmlns: models.BaseXmlns, -Metadata: models.BaseResponseMetadata, -} -return http.StatusOK, respStruct + respStruct := models.SetQueueAttributesResponse{ + Xmlns: models.BaseXmlns, + Metadata: models.BaseResponseMetadata, + } + return http.StatusOK, respStruct } diff --git a/shared-sqs/app/router/router.go b/shared-sqs/app/router/router.go index be57cc5..7c9b2b5 100644 --- a/shared-sqs/app/router/router.go +++ b/shared-sqs/app/router/router.go @@ -39,13 +39,14 @@ func New(tenantStore *tenant.TenantStore, adminToken string) http.Handler { // UI console — встроенный SPA, публичный доступ r.PathPrefix("/ui").Handler(http.StripPrefix("/ui", ui.Handler())) - // SQS API — tenant auth middleware - sqsRouter := r.NewRoute().Subrouter() - sqsRouter.Use(auth.AuthMiddleware(tenantStore)) - sqsRouter.HandleFunc("/", actionHandler).Methods("GET", "POST") - sqsRouter.HandleFunc("/{account}", actionHandler).Methods("GET", "POST") - sqsRouter.HandleFunc("/queue/{queueName}", actionHandler).Methods("GET", "POST") - sqsRouter.HandleFunc("/{account}/{queueName}", actionHandler).Methods("GET", "POST") + // SQS API — tenant auth middleware оборачивает каждый handler отдельно. + // r.NewRoute().Subrouter() с Use() некорректно работает в gorilla/mux v1.8.0 + // при пустом prefix — ответы теряются. Поэтому используем явную обёртку. + sqsAuth := auth.AuthMiddleware(tenantStore) + r.Handle("/", sqsAuth(http.HandlerFunc(actionHandler))).Methods("GET", "POST") + r.Handle("/{account}", sqsAuth(http.HandlerFunc(actionHandler))).Methods("GET", "POST") + r.Handle("/queue/{queueName}", sqsAuth(http.HandlerFunc(actionHandler))).Methods("GET", "POST") + r.Handle("/{account}/{queueName}", sqsAuth(http.HandlerFunc(actionHandler))).Methods("GET", "POST") return r } diff --git a/shared-sqs/app/ui/index.html b/shared-sqs/app/ui/index.html index de81902..d45228c 100644 --- a/shared-sqs/app/ui/index.html +++ b/shared-sqs/app/ui/index.html @@ -857,7 +857,7 @@ function renderQueueMessages(queueName, msgs) { onmouseover="this.style.background='rgba(26,127,212,0.08)'" onmouseout="this.style.background=''">