From ec10acc99eff0b1312446cb1f1775826d36df115 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E2=80=9CNaeel=E2=80=9D?= Date: Fri, 14 Aug 2026 15:53:08 +0400 Subject: [PATCH] =?UTF-8?q?v0.1.35:=20query-mode=20XML-=D0=BE=D1=82=D0=B2?= =?UTF-8?q?=D0=B5=D1=82=D1=8B,=20TOCTOU=20receive,=20=D0=B2=D0=B0=D0=BB?= =?UTF-8?q?=D0=B8=D0=B4=D0=B0=D1=86=D0=B8=D0=B8=20AWS,=20purge=2060s,=20?= =?UTF-8?q?=D1=81=D1=83=D0=B6=D0=B5=D0=BD=D0=B8=D0=B5=20=D0=BB=D0=BE=D0=BA?= =?UTF-8?q?=D0=BE=D0=B2?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- HISTORY/2026-08-14-session-log.md | 54 +++++ Makefile | 2 +- app/gosqs/create_queue.go | 26 ++- app/gosqs/delete_queue.go | 21 +- app/gosqs/gosqs.go | 113 +++++++--- app/gosqs/purge_queue.go | 26 ++- app/gosqs/queue_attributes.go | 81 +++++-- app/gosqs/receive_message.go | 229 +++++++++++++------- app/gosqs/send_message.go | 32 ++- app/gosqs/set_queue_attributes.go | 25 ++- app/models/errors.go | 4 + app/models/models.go | 50 ++++- app/router/router.go | 78 +++++-- app/ui/index.html | 2 +- tests/load_test.py | 348 ++++++++++++++++++++++++++++++ 15 files changed, 919 insertions(+), 172 deletions(-) create mode 100644 tests/load_test.py diff --git a/HISTORY/2026-08-14-session-log.md b/HISTORY/2026-08-14-session-log.md index b6d5174..e0194c0 100644 --- a/HISTORY/2026-08-14-session-log.md +++ b/HISTORY/2026-08-14-session-log.md @@ -465,3 +465,57 @@ DeleteQueue всегда успешен; NextSequenceNumber пересоздаё PurgeQueueInProgress (LastPurgeTime); FIFO MissingParameter; SendMessage пустое тело/размер — AWS-коды. 10. NumberOfReceives инкремент при выдаче. После: 25-цикловая проверка целостности, smoke, повторный 30-мин нагрузочный тест. + +--- + +## v0.1.35 — реализация фиксов по итогам нагрузочного теста (14.08.2026) + +**Изменённые файлы** (правки с подробными комментариями каждой функции и логики): + +1. `app/router/router.go`: + - `resolveProtocol`: заголовок `x-amzn-query-mode: true` (AWS CLI v2/boto3) → + принудительно Query-протокол (XML-ответы). Раньше query-mode запросы получали + JSON-ошибки → botocore Code=None («сырые 400/404»). + - `extractAction`: Action берётся сначала из `X-Amz-Target` (JSON и query-mode + клиенты), затем из form `Action` — для query-mode раньше Action не находился. + - `encodeResponse` (JSON): ошибки — в формате AWS-JSON + `{"__type":"com.amazonaws.sqs#","message":"..."}` (раньше голый ErrorResult). + - `actionHandler`: неизвестный Action → GeneralError в формате протокола клиента + (раньше текст "Bad Request"). Убран импорт `io`. + +2. `app/gosqs/receive_message.go`: + - TOCTOU устранён: проверка готовности И пометка in-flight — под одним Lock + (`receiveMessagesUnderLock`). Раньше два параллельных receive видели одно + сообщение, второй получал пусто → «потери». + - MaxNumberOfMessages: 1–10, выход → InvalidParameterValue (был clamp). + - Проверка существования очереди — под RLock (была гонка чтения map). + - `msg.NumberOfReceives++` при выдаче (счётчик доставок). + +3. `app/gosqs/send_message.go`: + - FIFO без MessageGroupId → MissingParameter (было: принималось). + - Док-схема функции; Redis-записи асинхронны, под Lock только marshal. + +4. `app/gosqs/gosqs.go`: + - `PeriodicTasks`: снапшот указателей под RLock + обработка каждой очереди под + отдельным Lock (раньше глобальный Lock на весь обход всех очередей). + +5. `app/models/models.go`: + - `LockGroup`/`NextSequenceNumber`: map больше НЕ пересоздаётся (терялись + блокировки/счётчики остальных групп). + - Новое поле `Queue.LastPurgeTime`. + +6. `app/gosqs/queue_attributes.go` + `set_queue_attributes.go` + `create_queue.go`: + - Применяются только явно переданные атрибуты (!= 0) — раньше zero-value + обнулял непереданные атрибуты. + - Bounds-check → InvalidParameterValue (был clamp). + - Whitelist имён → InvalidAttributeName (Query-протокол). + +7. `app/gosqs/purge_queue.go`: повторный purge < 60 с → PurgeQueueInProgress. + +8. `app/gosqs/delete_queue.go`: удаление несуществующей очереди → QueueNotFound. + +9. `app/models/errors.go`: добавлены `InvalidAttributeName`, `PurgeQueueInProgress`. + +**Проверки**: go build OK (v0.1.35), go vet чистый (кроме известной некомпиляции +тестов gosqs — пакет fixtures), go test app/models OK. + diff --git a/Makefile b/Makefile index d306b77..9f0e0b1 100644 --- a/Makefile +++ b/Makefile @@ -3,7 +3,7 @@ # Registry: Docker Hub (naeel/shared-sqs) IMAGE_REPO=naeel/shared-sqs -VERSION=v0.1.34 +VERSION=v0.1.35 LDFLAGS=-X shared-sqs/app/models.Version=$(VERSION) BINARY=shared-sqs diff --git a/app/gosqs/create_queue.go b/app/gosqs/create_queue.go index 0506591..8e294db 100644 --- a/app/gosqs/create_queue.go +++ b/app/gosqs/create_queue.go @@ -1,6 +1,3 @@ -// Изменено: 2026-04-10 — добавлена Redis persistence -// CreateQueueV1 — создаёт очередь для тенанта из request context. -// Ключ в SyncQueues: "{tenantAccessKey}:{queueName}" для изоляции между тенантами. package gosqs import ( @@ -15,6 +12,21 @@ import ( log "github.com/sirupsen/logrus" ) +// CreateQueueV1 — создаёт очередь для тенанта из request context. +// +// Логическая схема: +// 1. Разбор тела запроса (Query/JSON) → CreateQueueRequest. +// 2. Тенант из контекста (SigV4 → AccessKey). +// 3. Валидация имени очереди по AWS-правилам: до 80 символов, [a-zA-Z0-9_-], +// опциональный суффикс .fifo (ValidateQueueName). +// 4. Под Lock: +// - если очередь уже есть — просто возвращаем её URL (идемпотентность AWS); +// - проверка лимита очередей тенанта (MaxQueues) → LimitExceeded; +// - создание очереди с дефолтами и применение переданных атрибутов +// (setQueueAttributesV1 с валидацией диапазонов); +// - вставка в map с tenant-scoped ключом "{tenantAccessKey}:{queueName}". +// 5. Сохранение очереди в Redis (асинхронно, маршалинг под Lock). +// 6. Ответ: QueueUrl. func CreateQueueV1(req *http.Request) (int, interfaces.AbstractResponseBody) { requestBody := models.NewCreateQueueRequest() ok := utils.REQUEST_TRANSFORMER(requestBody, req, false) @@ -53,7 +65,13 @@ func CreateQueueV1(req *http.Request) (int, interfaces.AbstractResponseBody) { EnableDuplicates: models.CurrentEnvironment.EnableDuplicates, Duplicates: make(map[string]time.Time), } - if err := setQueueAttributesV1(queue, requestBody.Attributes); err != nil { + // Список имён переданных атрибутов — для whitelist-валидации + // (Query-протокол: Attribute.N.Name; JSON: nil — проверка пропускается). + var provided map[string]string + if req.Header.Get("Content-Type") != "application/x-amz-json-1.0" { + provided = utils.ExtractQueueAttributes(req.PostForm) + } + if err := setQueueAttributesV1(queue, requestBody.Attributes, provided); err != nil { models.SyncQueues.Unlock() return utils.CreateErrorResponseV1(err.Error(), true) } diff --git a/app/gosqs/delete_queue.go b/app/gosqs/delete_queue.go index 674181f..ec29c1a 100644 --- a/app/gosqs/delete_queue.go +++ b/app/gosqs/delete_queue.go @@ -1,5 +1,3 @@ -// Изменено: 2026-04-10 — добавлена Redis persistence -// DeleteQueueV1 — удаляет очередь тенанта по tenant-scoped ключу. package gosqs import ( @@ -14,6 +12,17 @@ import ( log "github.com/sirupsen/logrus" ) +// DeleteQueueV1 — удаляет очередь тенанта по tenant-scoped ключу. +// +// Логическая схема: +// 1. Разбор тела запроса → DeleteQueueRequest. +// 2. Тенант из контекста (SigV4 → AccessKey). +// 3. Имя очереди из QueueUrl (последний сегмент) + tenant-scoped ключ. +// 4. Под Lock: проверка существования очереди → QueueNotFound (фикс: раньше +// удаление несуществующей очереди возвращало 200, как будто она удалена); +// удаление очереди из in-memory map. +// 5. Асинхронное удаление метаданных и всех сообщений очереди из Redis +// (persistence.DeleteQueue — защищено write-seq от гонки со SaveQueue). func DeleteQueueV1(req *http.Request) (int, interfaces.AbstractResponseBody) { requestBody := models.NewDeleteQueueRequest() ok := utils.REQUEST_TRANSFORMER(requestBody, req, false) @@ -34,8 +43,14 @@ func DeleteQueueV1(req *http.Request) (int, interfaces.AbstractResponseBody) { log.Infof("Deleting Queue: %s (tenant: %s)", queueName, t.ID) models.SyncQueues.Lock() + defer models.SyncQueues.Unlock() + // Проверка существования: AWS возвращает NonExistentQueue, а не 200 + // (раньше удаление несуществующей очереди молча «успевало»). + if _, ok := models.SyncQueues.Queues[key]; !ok { + log.Warnf("Delete Queue: %s does not exist (tenant: %s)", queueName, t.ID) + return utils.CreateErrorResponseV1("QueueNotFound", true) + } delete(models.SyncQueues.Queues, key) - models.SyncQueues.Unlock() // Удаляем из Redis асинхронно persistence.DeleteQueue(key) diff --git a/app/gosqs/gosqs.go b/app/gosqs/gosqs.go index 6a45977..e0f1fbf 100644 --- a/app/gosqs/gosqs.go +++ b/app/gosqs/gosqs.go @@ -13,46 +13,40 @@ func init() { models.SyncQueues.Queues = make(map[string]*models.Queue) } +// PeriodicTasks — фоновая обработка очередей, выполняется раз в `d` (1 секунда). +// +// Задачи на каждый такт (для КАЖДОЙ очереди, см. processQueueTick): +// 1. Очистка просроченных записей deduplication-окна (FIFO). +// 2. Возврат сообщений в видимое состояние по истечении VisibilityTimeout. +// 3. Перенос сообщений в DeadLetterQueue при достижении MaxReceiveCount. +// +// ВАЖНО (фикс производительности): раньше глобальный SyncQueues.Lock удерживался +// на время обхода ВСЕХ очередей — при большом числе очередей HTTP-обработчики +// блокировались на секунды (наблюдались ReadTimeout и latency до 54 секунд). +// Теперь блокировка захватывается ПО ОДНОЙ очереди: между обработкой двух +// очередей другие горутины могут вклиниться. func PeriodicTasks(d time.Duration, quit chan bool) { ticker := time.NewTicker(d) for { select { case <-ticker.C: - models.SyncQueues.Lock() - for qName := range models.SyncQueues.Queues { - queue := models.SyncQueues.Queues[qName] - - // Reset deduplication period - for dedupId, startTime := range queue.Duplicates { - if time.Now().After(startTime.Add(models.DeduplicationPeriod)) { - log.Debugf("deduplication period for message with deduplicationId [%s] expired", dedupId) - delete(queue.Duplicates, dedupId) - } - } - - log.Debugf("Queue [%s] length [%d]", queue.Name, len(queue.Messages)) - for i := 0; i < len(queue.Messages); i++ { - msg := &queue.Messages[i] - - if msg.ReceiptHandle != "" { - if msg.VisibilityTimeout.Before(time.Now()) { - log.Debugf("Making message visible again %s", msg.ReceiptHandle) - queue.UnlockGroup(msg.GroupID) - msg.ReceiptHandle = "" - msg.ReceiptTime = time.Now().UTC() - msg.Retry++ - if queue.MaxReceiveCount > 0 && - queue.DeadLetterQueue != nil && - msg.Retry >= queue.MaxReceiveCount { - queue.DeadLetterQueue.Messages = append(queue.DeadLetterQueue.Messages, *msg) - queue.Messages = append(queue.Messages[:i], queue.Messages[i+1:]...) - i-- - } - } - } - } + // Снапшот указателей на очереди под коротким RLock. + // Указатели остаются валидными после снятия RLock: очереди из map + // не удаляются во время итерации, а объекты живут в куче. + models.SyncQueues.RLock() + queues := make([]*models.Queue, 0, len(models.SyncQueues.Queues)) + for _, queue := range models.SyncQueues.Queues { + queues = append(queues, queue) + } + models.SyncQueues.RUnlock() + + // Каждую очередь обрабатываем под ОТДЕЛЬНЫМ Lock — блокировка одной + // очереди не мешает обработке остальных и HTTP-запросам. + for _, queue := range queues { + models.SyncQueues.Lock() + processQueueTick(queue) + models.SyncQueues.Unlock() } - models.SyncQueues.Unlock() case <-quit: ticker.Stop() return @@ -60,6 +54,57 @@ func PeriodicTasks(d time.Duration, quit chan bool) { } } +// processQueueTick — обработка ОДНОЙ очереди за один такт. +// +// ВЫЗЫВАТЬ ТОЛЬКО под удержанием SyncQueues.Lock. +// +// Логическая схема: +// 1. Deduplication-окно (FIFO): удаляем записи старше DeduplicationPeriod. +// 2. Проходим по сообщениям очереди: +// - пустой ReceiptHandle (не in-flight) — пропускаем; +// - VisibilityTimeout ещё не истёк — пропускаем; +// - иначе: возвращаем видимость (сброс ReceiptHandle, разблокировка +// FIFO-группы, обновление ReceiptTime), инкрементируем Retry; +// - если Retry >= MaxReceiveCount и настроена DLQ — переносим сообщение +// в DLQ и удаляем из исходной очереди. +func processQueueTick(queue *models.Queue) { + // 1. Просроченные записи deduplication-окна. + for dedupId, startTime := range queue.Duplicates { + if time.Now().After(startTime.Add(models.DeduplicationPeriod)) { + log.Debugf("deduplication period for message with deduplicationId [%s] expired", dedupId) + delete(queue.Duplicates, dedupId) + } + } + + // 2. Возврат видимости истёкших сообщений + перенос в DLQ. + log.Debugf("Queue [%s] length [%d]", queue.Name, len(queue.Messages)) + for i := 0; i < len(queue.Messages); i++ { + msg := &queue.Messages[i] + + if msg.ReceiptHandle == "" { + continue // сообщение не в полёте + } + if msg.VisibilityTimeout.After(time.Now()) { + continue // таймаут видимости ещё действует + } + + log.Debugf("Making message visible again %s", msg.ReceiptHandle) + queue.UnlockGroup(msg.GroupID) + msg.ReceiptHandle = "" + msg.ReceiptTime = time.Now().UTC() + msg.Retry++ + + // 3. Перенос в DLQ при превышении MaxReceiveCount. + if queue.MaxReceiveCount > 0 && + queue.DeadLetterQueue != nil && + msg.Retry >= queue.MaxReceiveCount { + queue.DeadLetterQueue.Messages = append(queue.DeadLetterQueue.Messages, *msg) + queue.Messages = append(queue.Messages[:i], queue.Messages[i+1:]...) + i-- // элемент на позиции i удалён — повторно проверяем новый элемент на ней + } + } +} + func numberOfHiddenMessagesInQueue(queue models.Queue) int { num := 0 for _, m := range queue.Messages { diff --git a/app/gosqs/purge_queue.go b/app/gosqs/purge_queue.go index 0ae1466..cb371d0 100644 --- a/app/gosqs/purge_queue.go +++ b/app/gosqs/purge_queue.go @@ -1,5 +1,3 @@ -// Изменено: 2026-04-10 — добавлена Redis persistence -// PurgeQueueV1 — очищает все сообщения в очереди тенанта. package gosqs import ( @@ -15,6 +13,21 @@ import ( log "github.com/sirupsen/logrus" ) +// PurgeQueueV1 — очищает ВСЕ сообщения очереди тенанта. +// +// Логическая схема: +// 1. Разбор тела запроса → PurgeQueueRequest. +// 2. Тенант из контекста (SigV4 → AccessKey). +// 3. Имя очереди из QueueUrl (последний сегмент). +// 4. Под Lock: +// - проверка существования очереди → QueueNotFound; +// - защита от повторного purge (фикс): повторный вызов в течение 60 секунд +// → PurgeQueueInProgress (требование AWS). Раньше purge можно было +// вызывать без ограничений; +// - очистка in-memory сообщений и dedup-записей; +// - фиксация LastPurgeTime; +// - асинхронная очистка Redis (DEL хэша сообщений) + SaveQueue. +// 5. Ответ: пустой PurgeQueueResponse. func PurgeQueueV1(req *http.Request) (int, interfaces.AbstractResponseBody) { requestBody := models.NewPurgeQueueRequest() ok := utils.REQUEST_TRANSFORMER(requestBody, req, false) @@ -39,9 +52,18 @@ func PurgeQueueV1(req *http.Request) (int, interfaces.AbstractResponseBody) { return utils.CreateErrorResponseV1("QueueNotFound", true) } + // AWS: PurgeQueue допускается не чаще одного раза в 60 секунд на очередь. + // Нулевой LastPurgeTime (purge ещё не выполнялся) → запрет не срабатывает. + if time.Since(models.SyncQueues.Queues[key].LastPurgeTime) < 60*time.Second { + log.Warnf("Purge Queue: %s rejected — purge in progress (last: %s)", queueName, models.SyncQueues.Queues[key].LastPurgeTime.Format(time.RFC3339)) + return utils.CreateErrorResponseV1("PurgeQueueInProgress", 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) + // Фиксируем время purge — защита от повторного вызова в течение 60 секунд. + models.SyncQueues.Queues[key].LastPurgeTime = time.Now() // Удаляем все сообщения из Redis одной командой DEL persistence.PurgeMessagesPersist(key) // Сохраняем пустые метаданные очереди diff --git a/app/gosqs/queue_attributes.go b/app/gosqs/queue_attributes.go index 28fc14d..8f35a90 100644 --- a/app/gosqs/queue_attributes.go +++ b/app/gosqs/queue_attributes.go @@ -9,27 +9,65 @@ import ( "shared-sqs/app/models" ) -// TODO - Support: -// - attr.MessageRetentionPeriod -// - attr.Policy -// - attr.RedriveAllowPolicy -func setQueueAttributesV1(q *models.Queue, attr models.QueueAttributes) error { - // AWS-совместимые лимиты: clamp значений к допустимым диапазонам - if attr.DelaySeconds >= 0 { - q.DelaySeconds = ClampInt(attr.DelaySeconds.Int(), 0, MaxDelaySeconds) +// setQueueAttributesV1 — применяет атрибуты очереди с валидацией по AWS-спецификации. +// +// Логическая схема: +// 1. Whitelist имён: неизвестное имя атрибута → InvalidAttributeName. +// Список имён (provided) доступен для Query-протокола (form-параметры +// Attribute.N.Name); для JSON-протокола provided == nil и проверка пропускается +// (структура уже разобрана из тела, неизвестные поля отброшены). +// 2. Применяются ТОЛЬКО явно переданные атрибуты (значение != 0). +// ВАЖНО (фикс): раньше условие `attr.X >= 0` срабатывало и для НЕ переданных +// полей (zero-value = 0) — вызов SetQueueAttributes с одним атрибутом +// обнулял все остальные атрибуты очереди. +// 3. Значения вне диапазонов AWS → ошибка InvalidParameterValue +// (раньше значения молча клэмпились в допустимый диапазон). +// 4. RedrivePolicy: ARN DLQ разбирается, DLQ должна существовать, иначе +// InvalidAttributeValue. +func setQueueAttributesV1(q *models.Queue, attr models.QueueAttributes, provided map[string]string) error { + // Шаг 1: whitelist имён атрибутов. + for name := range provided { + if !attrNameWhitelist[name] { + return fmt.Errorf("InvalidAttributeName") + } } - if attr.MaximumMessageSize >= 0 { - q.MaximumMessageSize = ClampInt(attr.MaximumMessageSize.Int(), MinMessageSizeLimit, MaxMessageSizeLimit) + + // Шаг 2: применяем только явно переданные атрибуты с проверкой диапазонов. + // Диапазоны AWS: DelaySeconds 0–900, MaximumMessageSize 1024–262144, + // MessageRetentionPeriod 60–1209600, ReceiveMessageWaitTimeSeconds 0–20, + // VisibilityTimeout 0–43200. + if attr.DelaySeconds != 0 { + if attr.DelaySeconds.Int() < 0 || attr.DelaySeconds.Int() > MaxDelaySeconds { + return fmt.Errorf("InvalidParameterValue") + } + q.DelaySeconds = attr.DelaySeconds.Int() } - if attr.MessageRetentionPeriod > 0 { - q.MessageRetentionPeriod = ClampInt(attr.MessageRetentionPeriod.Int(), MinMessageRetentionPeriod, MaxMessageRetentionPeriod) + if attr.MaximumMessageSize != 0 { + if attr.MaximumMessageSize.Int() < MinMessageSizeLimit || attr.MaximumMessageSize.Int() > MaxMessageSizeLimit { + return fmt.Errorf("InvalidParameterValue") + } + q.MaximumMessageSize = attr.MaximumMessageSize.Int() } - if attr.ReceiveMessageWaitTimeSeconds > 0 { - q.ReceiveMessageWaitTimeSeconds = ClampInt(attr.ReceiveMessageWaitTimeSeconds.Int(), 0, MaxReceiveMessageWaitTimeSeconds) + if attr.MessageRetentionPeriod != 0 { + if attr.MessageRetentionPeriod.Int() < MinMessageRetentionPeriod || attr.MessageRetentionPeriod.Int() > MaxMessageRetentionPeriod { + return fmt.Errorf("InvalidParameterValue") + } + q.MessageRetentionPeriod = attr.MessageRetentionPeriod.Int() } - if attr.VisibilityTimeout >= 0 { - q.VisibilityTimeout = ClampInt(attr.VisibilityTimeout.Int(), 0, MaxVisibilityTimeout) + if attr.ReceiveMessageWaitTimeSeconds != 0 { + if attr.ReceiveMessageWaitTimeSeconds.Int() < 0 || attr.ReceiveMessageWaitTimeSeconds.Int() > MaxReceiveMessageWaitTimeSeconds { + return fmt.Errorf("InvalidParameterValue") + } + q.ReceiveMessageWaitTimeSeconds = attr.ReceiveMessageWaitTimeSeconds.Int() } + if attr.VisibilityTimeout != 0 { + if attr.VisibilityTimeout.Int() < 0 || attr.VisibilityTimeout.Int() > MaxVisibilityTimeout { + return fmt.Errorf("InvalidParameterValue") + } + q.VisibilityTimeout = attr.VisibilityTimeout.Int() + } + + // Шаг 3: RedrivePolicy (dead-letter queue). if attr.RedrivePolicy != (models.RedrivePolicy{}) { arnArray := strings.Split(attr.RedrivePolicy.DeadLetterTargetArn, ":") queueName := arnArray[len(arnArray)-1] @@ -43,3 +81,14 @@ func setQueueAttributesV1(q *models.Queue, attr models.QueueAttributes) error { } return nil } + +// attrNameWhitelist — имена атрибутов, которые SetQueueAttributes может изменять. +// Всё, чего нет в списке, отклоняется с ошибкой InvalidAttributeName. +var attrNameWhitelist = map[string]bool{ + "DelaySeconds": true, + "MaximumMessageSize": true, + "MessageRetentionPeriod": true, + "ReceiveMessageWaitTimeSeconds": true, + "VisibilityTimeout": true, + "RedrivePolicy": true, +} diff --git a/app/gosqs/receive_message.go b/app/gosqs/receive_message.go index ca9d160..22955bc 100644 --- a/app/gosqs/receive_message.go +++ b/app/gosqs/receive_message.go @@ -19,7 +19,27 @@ import ( log "github.com/sirupsen/logrus" ) +// ReceiveMessageV1 — получает до MaxNumberOfMessages сообщений из очереди тенанта. +// +// Логическая схема: +// 1. Разбор тела запроса (Query-форма или JSON) → ReceiveMessageRequest. +// 2. Аутентифицированный тенант из контекста (SigV4 → AccessKey). +// 3. Валидация входных параметров (AWS-совместимая): +// - MaxNumberOfMessages: 1–10; 0 (не задан) = 1. Выход за диапазон — ошибка +// InvalidParameterValue (раньше значение молча клэмпилось). +// - VisibilityTimeout: 0–43200 секунд, иначе InvalidParameterValue. +// - WaitTimeSeconds: 0–20 секунд; 0 = брать атрибут очереди +// ReceiveMessageWaitTimeSeconds. +// 4. Имя очереди: из QueueUrl (последний сегмент) или из пути /{account}/{queueName}. +// 5. Выборка сообщений — АТОМАРНО под одним Lock (см. receiveMessagesUnderLock): +// устранена гонка TOCTOU, когда два параллельных receive видели одни и те же +// сообщения, но один из них получал пустой ответ («потери» сообщений). +// 6. Long polling: при WaitTimeSeconds > 0 ждём сообщений тиками по 100 мс до +// истечения deadline; прерываемся по отмене HTTP-запроса (клиент отвалился). +// +// Возвращает HTTP 200 и ReceiveMessageResponse (пустой Result, если сообщений нет). func ReceiveMessageV1(req *http.Request) (int, interfaces.AbstractResponseBody) { + // --- Шаг 1: разбор тела запроса --- requestBody := models.NewReceiveMessageRequest() ok := utils.REQUEST_TRANSFORMER(requestBody, req, false) if !ok { @@ -27,18 +47,31 @@ func ReceiveMessageV1(req *http.Request) (int, interfaces.AbstractResponseBody) return utils.CreateErrorResponseV1("InvalidParameterValue", true) } + // --- Шаг 2: тенант из контекста --- t := getTenantFromContext(req) if t == nil { return utils.CreateErrorResponseV1("InvalidClientTokenId", true) } + // --- Шаг 3: валидация параметров --- + // MaxNumberOfMessages: AWS допускает 1–10; 0 (параметр не задан) означает 1. + // Выход за диапазон — ошибка, а не молчаливый clamp (так делает AWS). maxNumberOfMessages := requestBody.MaxNumberOfMessages if maxNumberOfMessages == 0 { - maxNumberOfMessages = 1 + maxNumberOfMessages = MinNumberOfMessagesLimit + } + if maxNumberOfMessages < MinNumberOfMessagesLimit || maxNumberOfMessages > MaxNumberOfMessagesLimit { + return utils.CreateErrorResponseV1("InvalidParameterValue", true) } - // Fix #8: clamp MaxNumberOfMessages к AWS лимиту 1–10 - maxNumberOfMessages = ClampInt(maxNumberOfMessages, MinNumberOfMessagesLimit, MaxNumberOfMessagesLimit) + // VisibilityTimeout: AWS допускает 0–43200. + if requestBody.VisibilityTimeout < 0 || requestBody.VisibilityTimeout > MaxVisibilityTimeout { + return utils.CreateErrorResponseV1("InvalidParameterValue", true) + } + + // --- Шаг 4: имя очереди --- + // QueueUrl не задан → берём из пути /{account}/{queueName} (старый стиль клиентов). + // QueueUrl задан → последний сегмент URL = имя очереди. queueName := "" if requestBody.QueueUrl == "" { vars := mux.Vars(req) @@ -50,118 +83,150 @@ func ReceiveMessageV1(req *http.Request) (int, interfaces.AbstractResponseBody) key := tenantQueueKey(t.AccessKey, queueName) - if _, ok := models.SyncQueues.Queues[key]; !ok { + // Проверка существования очереди — под RLock (защита от параллельной записи + // в map; раньше map читалась без блокировки — гонка данных). + models.SyncQueues.RLock() + _, queueExists := models.SyncQueues.Queues[key] + models.SyncQueues.RUnlock() + if !queueExists { return utils.CreateErrorResponseV1("QueueNotFound", true) } - var messages []*models.ResultMessage - respStruct := models.ReceiveMessageResponse{} - - // Валидация VisibilityTimeout: AWS SQS допускает 0–43200 - if requestBody.VisibilityTimeout < 0 || requestBody.VisibilityTimeout > MaxVisibilityTimeout { - return utils.CreateErrorResponseV1("InvalidParameterValue", true) - } - + // --- Шаг 5: WaitTimeSeconds (long polling) --- + // 0 (не задан) → атрибут очереди ReceiveMessageWaitTimeSeconds (default 0). waitTimeSeconds := requestBody.WaitTimeSeconds if waitTimeSeconds == 0 { models.SyncQueues.RLock() waitTimeSeconds = models.SyncQueues.Queues[key].ReceiveMessageWaitTimeSeconds models.SyncQueues.RUnlock() } - // Валидация WaitTimeSeconds: AWS SQS допускает 0–20, иначе ошибка if waitTimeSeconds < 0 || waitTimeSeconds > MaxReceiveMessageWaitTimeSeconds { return utils.CreateErrorResponseV1("InvalidParameterValue", true) } + // --- Шаги 6–7: выборка сообщений (атомарно) + long polling --- if waitTimeSeconds > 0 { deadline := time.Now().Add(time.Duration(waitTimeSeconds) * time.Second) pollTicker := time.NewTicker(100 * time.Millisecond) defer pollTicker.Stop() for { - models.SyncQueues.RLock() - queue, queueFound := models.SyncQueues.Queues[key] - if !queueFound { - models.SyncQueues.RUnlock() - return utils.CreateErrorResponseV1("QueueNotFound", true) + // Атомарная проверка+выборка: между проверкой и выдачей другой + // получатель не может «съесть» сообщения — они помечаются + // ReceiptHandle в том же критическом участке (см. receiveMessagesUnderLock). + messages, found := receiveMessagesUnderLock(key, requestBody.VisibilityTimeout, maxNumberOfMessages) + if found { + return http.StatusOK, buildReceiveResponse(messages) } - messageFound := queueHasReceivableMessages(queue) - models.SyncQueues.RUnlock() - if messageFound || time.Now().After(deadline) { - break + if time.Now().After(deadline) { + return http.StatusOK, buildReceiveResponse(nil) } select { case <-req.Context().Done(): - return http.StatusOK, models.ReceiveMessageResponse{ - Xmlns: models.BaseXmlns, - Result: models.ReceiveMessageResult{}, - Metadata: models.BaseResponseMetadata, - } + // Клиент разорвал соединение — завершаем без ошибки. + return http.StatusOK, buildReceiveResponse(nil) case <-pollTicker.C: } } } - log.Debugf("Getting Message from Queue:%s (tenant: %s)", queueName, t.ID) + messages, _ := receiveMessagesUnderLock(key, requestBody.VisibilityTimeout, maxNumberOfMessages) + return http.StatusOK, buildReceiveResponse(messages) +} + +// receiveMessagesUnderLock — АТОМАРНАЯ проверка и выборка сообщений очереди. +// +// ВЫЗЫВАТЬ БЕЗ удержания SyncQueues-блокировок. +// Возвращает: +// - messages — готовые к ответу ResultMessage (уже помечены in-flight); +// - found=false — пригодных сообщений нет (очередь пуста, всё in-flight, +// всё на задержке доставки или все FIFO-группы заблокированы). +// +// Почему весь цикл под одним Lock (фикс гонки TOCTOU): +// раньше проверка queueHasReceivableMessages выполнялась под RLock, а выборка — +// под Lock отдельно. Два параллельных receive оба проходили проверку, затем +// первый забирал все сообщения, а второй получал пустой ответ — сообщения +// «терялись» для второго клиента. Теперь проверка и пометка ReceiptHandle +// выполняются в одном критическом участке. +// +// Для каждого выбранного сообщения: +// - генерируется ReceiptHandle = {Uuid}#{случайный id} — уникален для каждой +// доставки (два разных receive не имеют одинакового хендла); +// - проставляется VisibilityTimeout (параметр запроса или атрибут очереди); +// - инкрементируется NumberOfReceives — счётчик доставок +// (ApproximateReceiveCount в ответе); +// - для FIFO-очереди блокируется группа (LockGroup) до DeleteMessage/таймаута. +func receiveMessagesUnderLock(key string, visibilityTimeout, maxNumberOfMessages int) ([]*models.ResultMessage, bool) { models.SyncQueues.Lock() defer models.SyncQueues.Unlock() - 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 - } - - if models.SyncQueues.Queues[key].Messages[i].ReceiptHandle != "" { - continue - } - - msg := &models.SyncQueues.Queues[key].Messages[i] - if !msg.IsReadyForReceipt() { - continue - } - - if models.SyncQueues.Queues[key].IsFIFO { - if models.SyncQueues.Queues[key].IsLocked(msg.GroupID) { - continue - } - models.SyncQueues.Queues[key].LockGroup(msg.GroupID) - } - - randomId := uuid.NewString() - msg.ReceiptHandle = msg.Uuid + "#" + randomId - msg.ReceiptTime = time.Now().UTC() - - 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) - } - - 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"}, - } + queue, queueFound := models.SyncQueues.Queues[key] + if !queueFound || len(queue.Messages) == 0 { + return nil, false + } + // Быстрая предпроверка: есть ли вообще хотя бы одно готовое сообщение. + if !queueHasReceivableMessages(queue) { + return nil, false } - return http.StatusOK, respStruct + messages := make([]*models.ResultMessage, 0) + numMsg := 0 + for i := range queue.Messages { + if numMsg >= maxNumberOfMessages { + break + } + msg := &queue.Messages[i] + if msg.ReceiptHandle != "" { + continue // сообщение уже в полёте у другого получателя + } + if !msg.IsReadyForReceipt() { + continue // задержка доставки (DelaySeconds / random latency) ещё не истекла + } + if queue.IsFIFO { + if queue.IsLocked(msg.GroupID) { + continue // FIFO-группа заблокирована предыдущим in-flight сообщением + } + queue.LockGroup(msg.GroupID) + } + + // Пометка in-flight: уникальный ReceiptHandle + таймаут видимости. + msg.ReceiptHandle = msg.Uuid + "#" + uuid.NewString() + msg.ReceiptTime = time.Now().UTC() + if visibilityTimeout != 0 { + msg.VisibilityTimeout = time.Now().Add(time.Duration(visibilityTimeout) * time.Second) + } else { + msg.VisibilityTimeout = time.Now().Add(time.Duration(queue.VisibilityTimeout) * time.Second) + } + msg.NumberOfReceives++ + + messages = append(messages, buildResultMessage(msg)) + numMsg++ + } + + if numMsg == 0 { + return nil, false + } + return messages, true } +// buildReceiveResponse — формирует успешный ответ ReceiveMessage. +// При messages == nil Result содержит пустой список (как AWS: без ). +func buildReceiveResponse(messages []*models.ResultMessage) models.ReceiveMessageResponse { + return models.ReceiveMessageResponse{ + Xmlns: models.BaseXmlns, + Result: models.ReceiveMessageResult{Messages: messages}, + Metadata: models.BaseResponseMetadata, + } +} + +// buildResultMessage — конвертирует SqsMessage в ResultMessage (XML/JSON ответ). +// +// Атрибуты AWS-совместимы: +// - ApproximateFirstReceiveTimestamp — время первой выдачи (ReceiptTime); +// - ApproximateReceiveCount — число доставок (NumberOfReceives + 1); +// - SentTimestamp — реальное время отправки (SentTime), НЕ время выдачи. +// MD5 тела берётся из кэша (посчитан при SendMessage) — без пересчёта. func buildResultMessage(m *models.SqsMessage) *models.ResultMessage { return &models.ResultMessage{ MessageId: m.Uuid, @@ -179,6 +244,12 @@ func buildResultMessage(m *models.SqsMessage) *models.ResultMessage { } } +// queueHasReceivableMessages — проверяет, есть ли в очереди хотя бы ОДНО готовое +// к выдаче сообщение: не in-flight, задержка доставки истекла, FIFO-группа не +// заблокирована. Используется как быстрая предпроверка перед выборкой. +// +// ВАЖНО: вызывается ТОЛЬКО под удержанием блокировки (Lock или RLock) — +// читает общее состояние очереди. func queueHasReceivableMessages(queue *models.Queue) bool { for i := range queue.Messages { msg := &queue.Messages[i] diff --git a/app/gosqs/send_message.go b/app/gosqs/send_message.go index 41cacef..2a209ef 100644 --- a/app/gosqs/send_message.go +++ b/app/gosqs/send_message.go @@ -1,7 +1,3 @@ -// Изменено: 2026-04-11 — убрано логирование тела сообщения (perf + security) -// SendMessageV1 — добавляет сообщение в очередь тенанта. -// Ловушка #6: queueName извлекается как ПОСЛЕДНИЙ сегмент URL — при URL вида -// http://host/tenantID/queueName последний сегмент = queueName (правильно). package gosqs import ( @@ -21,6 +17,27 @@ import ( "github.com/gorilla/mux" ) +// SendMessageV1 — добавляет сообщение в очередь тенанта. +// +// Логическая схема: +// 1. Разбор тела запроса (Query/JSON) → SendMessageRequest. +// 2. Аутентифицированный тенант из контекста (SigV4 → AccessKey). +// 3. Валидации: +// - MessageBody обязателен (пустой → MissingParameter); +// - DeduplicationId и GroupId ≤ 128 символов; +// - MessageAttributes ≤ 10; +// - FIFO-очередь БЕЗ MessageGroupId → MissingParameter (AWS требует GroupId +// для каждой отправки в FIFO); +// - размер тела ≤ MaximumMessageSize очереди → MessageTooBig; +// - число сообщений < лимита очереди (OOM-защита) → OverLimit. +// 4. Построение SqsMessage: MD5 тела и атрибутов, UUID, метки времени, задержка. +// 5. Атомарно под Lock: +// - FIFO: порядковый номер группы (NextSequenceNumber); +// - проверка дубликата по deduplicationId (дубль НЕ добавляется); +// - добавление в очередь + SaveMessage/SaveQueue в Redis. +// Примечание: сетевые Redis-записи асинхронны (asyncWrite), под Lock +// выполняется только json.Marshal — блокировка не удерживается на сетевом I/O. +// 6. Ответ: MessageId, MD5OfMessageBody/MD5OfMessageAttributes, SequenceNumber (FIFO). func SendMessageV1(req *http.Request) (int, interfaces.AbstractResponseBody) { requestBody := models.NewSendMessageRequest() ok := utils.REQUEST_TRANSFORMER(requestBody, req, false) @@ -78,6 +95,13 @@ func SendMessageV1(req *http.Request) (int, interfaces.AbstractResponseBody) { delaySecs := queue.DelaySeconds models.SyncQueues.RUnlock() + // FIFO-очередь БЕЗ MessageGroupId — ошибка: AWS требует GroupId для каждой + // отправки в FIFO (иначе невозможно гарантировать порядок внутри группы). + // Раньше сообщение просто попадало в группу с пустым именем. + if queueIsFIFO && messageGroupID == "" { + return utils.CreateErrorResponseV1("MissingParameter", true) + } + if maxMessageSize > 0 && len(messageBody) > maxMessageSize { return utils.CreateErrorResponseV1("MessageTooBig", true) } diff --git a/app/gosqs/set_queue_attributes.go b/app/gosqs/set_queue_attributes.go index fb6420f..cd9ab76 100644 --- a/app/gosqs/set_queue_attributes.go +++ b/app/gosqs/set_queue_attributes.go @@ -1,6 +1,3 @@ -// Изменено: 2026-04-10 — добавлена Redis persistence -// SetQueueAttributesV1 — устанавливает атрибуты очереди тенанта. -// Ловушка #9: при RedrivePolicy парсим ARN DLQ и DLQ тоже должна принадлежать тому же тенанту. package gosqs import ( @@ -15,6 +12,17 @@ import ( log "github.com/sirupsen/logrus" ) +// SetQueueAttributesV1 — устанавливает атрибуты очереди тенанта. +// +// Логическая схема: +// 1. Разбор тела запроса (Query/JSON) → SetQueueAttributesRequest. +// 2. QueueUrl обязателен → InvalidParameterValue. +// 3. Тенант из контекста (SigV4 → AccessKey). +// 4. Имя очереди из QueueUrl (последний сегмент) + tenant-scoped ключ. +// 5. Под Lock: проверка существования очереди → QueueNotFound; +// применение атрибутов с валидацией (setQueueAttributesV1); +// сохранение очереди в Redis (асинхронно, маршалинг под Lock). +// 6. Ответ: пустой SetQueueAttributesResponse (AWS так и отвечает). func SetQueueAttributesV1(req *http.Request) (int, interfaces.AbstractResponseBody) { requestBody := models.NewSetQueueAttributesRequest() ok := utils.REQUEST_TRANSFORMER(requestBody, req, false) @@ -44,7 +52,16 @@ func SetQueueAttributesV1(req *http.Request) (int, interfaces.AbstractResponseBo 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 { + // Список имён переданных атрибутов — для whitelist-валидации. + // Query-протокол: имена видны в form (Attribute.N.Name). + // JSON-протокол: структура уже разобрана из тела, список имён недоступен → nil + // (whitelist-проверка пропускается, неизвестные поля отброшены парсером). + var provided map[string]string + if req.Header.Get("Content-Type") != "application/x-amz-json-1.0" { + provided = utils.ExtractQueueAttributes(req.PostForm) + } + + if err := setQueueAttributesV1(queue, requestBody.Attributes, provided); err != nil { return utils.CreateErrorResponseV1(err.Error(), true) } // Сохраняем атрибуты в Redis пока держим Lock (через defer) diff --git a/app/models/errors.go b/app/models/errors.go index 4aca2da..c36eebd 100644 --- a/app/models/errors.go +++ b/app/models/errors.go @@ -24,6 +24,10 @@ func init() { "LimitExceeded": {HttpError: http.StatusBadRequest, Type: "LimitExceeded", Code: "AWS.SimpleQueueService.LimitExceeded", Message: "You've reached the limit on the number of queues."}, // MissingParameter — обязательный параметр отсутствует (например, пустой MessageBody) "MissingParameter": {HttpError: http.StatusBadRequest, Type: "MissingParameter", Code: "MissingParameter", Message: "The request must contain the parameter MessageBody."}, + // InvalidAttributeName — передан атрибут с неизвестным именем (SetQueueAttributes) + "InvalidAttributeName": {HttpError: http.StatusBadRequest, Type: "InvalidAttributeName", Code: "AWS.SimpleQueueService.InvalidAttributeName", Message: "Unknown Attribute."}, + // PurgeQueueInProgress — повторный PurgeQueue на очереди ранее чем через 60 секунд + "PurgeQueueInProgress": {HttpError: http.StatusForbidden, Type: "Sender", Code: "AWS.SimpleQueueService.PurgeQueueInProgress", Message: "Only one PurgeQueue operation on a queue is allowed every 60 seconds."}, } } diff --git a/app/models/models.go b/app/models/models.go index 5529afe..332e172 100644 --- a/app/models/models.go +++ b/app/models/models.go @@ -63,38 +63,66 @@ type Queue struct { EnableDuplicates bool Duplicates map[string]time.Time Tags map[string]string + // LastPurgeTime — время последнего успешного PurgeQueue. + // Используется для запрета повторного purge ранее чем через 60 секунд + // (требование AWS: PurgeQueueInProgress). Нулевое значение = purge ещё не выполнялся. + LastPurgeTime time.Time } +// NextSequenceNumber — возвращает следующий порядковый номер сообщения FIFO-группы. +// +// Логическая схема: каждый вызов инкрементирует счётчик группы и возвращает +// его строковое представление (порядковые номера в FIFO должны расти монотонно). +// +// ВАЖНО (фикс): раньше при ПЕРВОМ обращении к группе map пересоздавалась целиком +// (q.FIFOSequenceNumbers = map[string]int{groupId: 0}) — счётчики ВСЕХ остальных +// групп терялись, и порядковые номера начинали дублироваться. Теперь map только +// инициализируется один раз (если nil), а счётчик конкретной группы инкрементируется. func (q *Queue) NextSequenceNumber(groupId string) string { - if _, ok := q.FIFOSequenceNumbers[groupId]; !ok { - q.FIFOSequenceNumbers = map[string]int{ - groupId: 0, - } + if q.FIFOSequenceNumbers == nil { + q.FIFOSequenceNumbers = make(map[string]int) } - q.FIFOSequenceNumbers[groupId]++ return strconv.Itoa(q.FIFOSequenceNumbers[groupId]) } +// IsLocked — проверяет, заблокирована ли FIFO-группа. +// Группа блокируется при выдаче её сообщения в ReceiveMessage и остаётся +// заблокированной, пока сообщение in-flight (до DeleteMessage или таймаута +// видимости) — это гарантирует порядок обработки внутри группы. func (q *Queue) IsLocked(groupId string) bool { _, ok := q.FIFOMessages[groupId] return ok } +// LockGroup — блокирует FIFO-группу (помечает как занятую). +// +// ВАЖНО (фикс): раньше при ПЕРВОМ обращении map пересоздавалась целиком +// (q.FIFOMessages = map[string]int{groupId: 0}) — блокировки ВСЕХ остальных +// групп терялись, и их сообщения могли выдаваться одновременно. Теперь map +// только инициализируется один раз, а для группы проставляется ключ. func (q *Queue) LockGroup(groupId string) { - if _, ok := q.FIFOMessages[groupId]; !ok { - q.FIFOMessages = map[string]int{ - groupId: 0, - } + if q.FIFOMessages == nil { + q.FIFOMessages = make(map[string]int) } + q.FIFOMessages[groupId] = 0 } +// UnlockGroup — снимает блокировку FIFO-группы. +// Вызывается после удаления сообщения или при истечении таймаута видимости +// (см. PeriodicTasks → processQueueTick). func (q *Queue) UnlockGroup(groupId string) { if _, ok := q.FIFOMessages[groupId]; ok { delete(q.FIFOMessages, groupId) } } +// IsDuplicate — проверяет, было ли сообщение с таким deduplicationId уже +// добавлено в очередь в пределах deduplication-окна. +// +// Логическая схема: дубликаты отслеживаются ТОЛЬКО для FIFO-очередей с +// включённой дупликацией (EnableDuplicates) и непустым deduplicationId. +// Для остальных случаев дубликаты не отслеживаются (AWS-семантика). func (q *Queue) IsDuplicate(deduplicationId string) bool { if !q.EnableDuplicates || !q.IsFIFO || deduplicationId == "" { return false @@ -105,6 +133,10 @@ func (q *Queue) IsDuplicate(deduplicationId string) bool { return ok } +// InitDuplicatation — регистрирует deduplicationId с текущим временем. +// Эта отметка — начало deduplication-окна; просроченные записи удаляются +// фоновой задачей PeriodicTasks (deduplication period expired). +// Для не-FIFO очередей или пустого deduplicationId — no-op. func (q *Queue) InitDuplicatation(deduplicationId string) { if !q.EnableDuplicates || !q.IsFIFO || deduplicationId == "" { return diff --git a/app/router/router.go b/app/router/router.go index cbf4046..29d4e80 100644 --- a/app/router/router.go +++ b/app/router/router.go @@ -7,7 +7,6 @@ import ( "encoding/json" "encoding/xml" "fmt" - "io" "net/http" "strings" @@ -103,10 +102,35 @@ func requestBodyLimitMiddleware(next http.Handler) http.Handler { }) } +// encodeResponse — сериализует ответ в формате протокола клиента. +// +// Логическая схема: +// - AwsJsonProtocol (чистый JSON-протокол): успешный ответ — JSON тела; +// ОШИБКА — спец-формат AWS-JSON {"__type": "com.amazonaws.sqs#", +// "message": "..."}. Раньше для ErrorResponse писался голый ErrorResult +// ({"Type":...,"Code":...,"Message":...}) — SDK не распознавал это как ошибку +// и возвращал Code=None. +// - AwsQueryProtocol: ответ всегда XML (в т.ч. для ошибок). func encodeResponse(w http.ResponseWriter, req *http.Request, statusCode int, body interfaces.AbstractResponseBody) { protocol := resolveProtocol(req) switch protocol { case AwsJsonProtocol: + // Ошибка в JSON-протоколе: __type + message — единственный формат, + // который botocore распознаёт как исключение с кодом. + if errResp, ok := body.(models.ErrorResponse); ok { + w.Header().Set("x-amzn-RequestId", errResp.RequestId) + w.Header().Set("Content-Type", "application/x-amz-json-1.0") + w.WriteHeader(statusCode) + err := json.NewEncoder(w).Encode(map[string]string{ + "__type": "com.amazonaws.sqs#" + errResp.Result.Code, + "message": errResp.Result.Message, + }) + if err != nil { + log.Errorf("Response Encoding Error: %v\nResponse: %+v", err, body) + } + return + } + // Успешный ответ: JSON тела (Result). w.Header().Set("x-amzn-RequestId", body.GetRequestId()) w.Header().Set("Content-Type", "application/x-amz-json-1.0") w.WriteHeader(statusCode) @@ -196,9 +220,16 @@ func actionHandler(w http.ResponseWriter, req *http.Request) { } return } + // Неизвестный Action — отвечаем ошибкой в ФОРМАТЕ ПРОТОКОЛА клиента + // (XML ErrorResponse для Query, AWS-JSON для JSON-протокола). + // Раньше возвращался голый текст "Bad Request" (200/400 без AWS-XML) — + // SDK получал ClientError с Code=None и не мог понять причину. log.Warnf("Bad Request - Action: %s", action) - w.WriteHeader(http.StatusBadRequest) - io.WriteString(w, "Bad Request") + errType := models.SqsErrors["GeneralError"] + encodeResponse(w, req, errType.StatusCode(), models.ErrorResponse{ + Result: errType.Response(), + RequestId: "00000000-0000-0000-0000-000000000000", + }) } type AwsProtocol int @@ -208,25 +239,42 @@ const ( AwsQueryProtocol AwsProtocol = iota ) -// extractAction — извлекает Action из запроса (Query Protocol или JSON Protocol) +// extractAction — извлекает имя операции SQS из запроса. +// +// Логическая схема (не зависит от resolveProtocol): +// 1. Заголовок X-Amz-Target ("AmazonSQS.SendMessage") — используют клиенты +// JSON-протокола И AWS CLI v2 в query-mode. Извлекаем часть после точки. +// 2. Иначе — form-параметр Action (классический query-протокол). +// +// ВАЖНО: определение по X-Amz-Target в первую очередь обязательно — AWS CLI v2 +// шлёт Content-Type: application/x-amz-json-1.0 + x-amzn-query-mode: true, но +// Action у него в заголовке, а НЕ в form-параметрах. func extractAction(req *http.Request) string { - protocol := resolveProtocol(req) - switch protocol { - case AwsJsonProtocol: - action := req.Header.Get("X-Amz-Target") + if action := req.Header.Get("X-Amz-Target"); action != "" { parts := strings.SplitN(action, ".", 2) - if len(parts) != 2 || parts[1] == "" { - return "" + if len(parts) == 2 && parts[1] != "" { + return parts[1] } - return parts[1] - case AwsQueryProtocol: - return req.FormValue("Action") } - return "" + return req.FormValue("Action") } -// resolveProtocol — определяет протокол по Content-Type +// resolveProtocol — определяет протокол ОТВЕТА по заголовкам запроса. +// +// Логическая схема: +// 1. Заголовок x-amzn-query-mode: true (шлют AWS CLI v2 и boto3) означает +// «query-compatible JSON»: тело запроса сериализовано как JSON, но клиент +// ЖДЁТ ответ в XML (query-протокол). Если на такой запрос ответить JSON — +// botocore не распарсит тело ошибки (ClientError с Code=None — именно это +// наблюдалось: «сырые 400/404» без AWS-XML). Поэтому для query-mode +// принудительно отвечаем XML. +// 2. Content-Type: application/x-amz-json-1.0 БЕЗ query-mode — чистый +// JSON-протокол (SDK): и запрос, и ответ в JSON. +// 3. Всё остальное (form-urlencoded от старых клиентов) — query/XML. func resolveProtocol(req *http.Request) AwsProtocol { + if req.Header.Get("x-amzn-query-mode") == "true" { + return AwsQueryProtocol + } if req.Header.Get("Content-Type") == "application/x-amz-json-1.0" { return AwsJsonProtocol } diff --git a/app/ui/index.html b/app/ui/index.html index 0932152..5ee4511 100644 --- a/app/ui/index.html +++ b/app/ui/index.html @@ -384,7 +384,7 @@ td.msg-expand { padding: 0 !important; border-bottom: 1px solid var(--border); }
Админ-консоль: /admin (вход по admin-токену)
Realm: iot-naeel · Persistence: Managed Redis
Версия:
-
Образ: naeel/shared-sqs:v0.1.34
+
Образ: naeel/shared-sqs:v0.1.35
Статус: тестирование
diff --git a/tests/load_test.py b/tests/load_test.py new file mode 100644 index 0000000..ab02a61 --- /dev/null +++ b/tests/load_test.py @@ -0,0 +1,348 @@ +#!/usr/bin/env python3 +"""tests/load_test.py — длительная нагрузочная проверка устойчивости shared-sqs. + +Режимы работы в каждом цикле: +1) Корректные операции: create/get-url/attributes/send/batch/receive(long poll)/ + change-visibility/delete/batch-delete/tags/purge/delete-queue + FIFO. +2) Крайние условия: сообщения до 256КБ, юникод, пустые тела, пакеты по 10, + сообщение > лимита (300КБ — ожидается отказ). +3) Некорректные условия (ожидаемые отказы с известными кодами): + несуществующая очередь, битый ReceiptHandle, некорректные атрибуты, + повторный purge, >10 записей в batch, дубли Id в batch, пустой batch, + FIFO без MessageGroupId, некорректные имена очередей. + +Контроль: отдельный поток каждые 10с проверяет GET /health. +Каждые 30с — строка прогресса. В конце — сводка по кодам ошибок и латентности. + +Переменные окружения: + ENDPOINT_URL (по умолчанию https://sqs.containerk8s.dev.nubes.ru) + REGION (us-east-1) + DURATION_SECONDS (1800 = 30 минут) + WORKERS (4) + AWS_ACCESS_KEY_ID / AWS_SECRET_ACCESS_KEY — обязательны +""" +import json +import os +import random +import string +import sys +import threading +import time +import urllib.request + +import boto3 +from botocore.config import Config +from botocore.exceptions import ClientError + +ENDPOINT = os.environ.get("ENDPOINT_URL", "https://sqs.containerk8s.dev.nubes.ru") +REGION = os.environ.get("REGION", "us-east-1") +DURATION = int(os.environ.get("DURATION_SECONDS", "1800")) +WORKERS = int(os.environ.get("WORKERS", "4")) +HEALTH_EVERY = 10 + +# Коды ошибок, которые ожидаемы для некорректных условий. +EXPECTED_CODES = { + "NonExistentQueue", + "QueueDeletedRecently", + "PurgeQueueInProgress", + "ReceiptHandleIsInvalid", + "InvalidParameterValue", + "InvalidAttributeValue", + "InvalidAttributeName", + "MissingParameter", + "ReadCountOutOfRange", + "TooManyEntriesInBatchRequest", + "BatchEntryIdsNotDistinct", + "EmptyBatchRequest", + "InvalidBatchEntryId", + "MessageTooLong", + "InvalidMessageContents", + "AWS.SimpleQueueService.NonExistentQueue", +} + +stats = { + "valid_ok": 0, + "expected_err": 0, + "unexpected_fail": 0, + "integrity_fail": 0, + "soft_warn": 0, # некорректный запрос неожиданно прошёл + "health_fail": 0, + "ops_total": 0, + "err_by_code": {}, +} +stats_lock = threading.Lock() +start_ts = time.time() +stop = threading.Event() + +LATENCY = [] +LATENCY_LOCK = threading.Lock() + + +def add_latency(sec): + with LATENCY_LOCK: + LATENCY.append(sec) + if len(LATENCY) > 100000: + del LATENCY[:50000] + + +def bump(key, n=1, code=None): + with stats_lock: + stats[key] += n + stats["ops_total"] += n + if code: + stats["err_by_code"][code] = stats["err_by_code"].get(code, 0) + 1 + + +def record_expected(e: ClientError): + code = e.response.get("Error", {}).get("Code", "Unknown") + bump("expected_err", code=code) + + +def record_unexpected(e): + code = getattr(e, "response", None) and e.response.get("Error", {}).get("Code", "Unknown") or type(e).__name__ + bump("unexpected_fail", code=code) + + +def rand_body(max_size=250 * 1024): + size = random.randint(0, max_size) + return "".join(random.choices(string.ascii_letters + "абвгд\x00\xff😀", k=size)) + + +def make_client(): + cfg = Config(connect_timeout=10, read_timeout=25, retries={"max_attempts": 2}) + return boto3.client("sqs", endpoint_url=ENDPOINT, region_name=REGION, config=cfg) + + +def health_monitor(): + while not stop.is_set(): + try: + with urllib.request.urlopen( + urllib.request.Request(ENDPOINT + "/health", headers={"User-Agent": "load-test"}), + timeout=10, + ) as r: + body = r.read().decode("utf-8", "replace") + if r.status != 200 or '"status":"ok"' not in body: + bump("health_fail") + except Exception: + bump("health_fail") + stop.wait(HEALTH_EVERY) + + +def valid_cycle(sqs, wid, it): + q = f"load-{wid}-{it}-{random.randint(0, 10**9)}" + url = sqs.create_queue(QueueName=q)["QueueUrl"] + try: + assert sqs.get_queue_url(QueueName=q)["QueueUrl"] == url + + # сообщение произвольного размера, включая 256КБ и юникод + body = rand_body() + sqs.send_message(QueueUrl=url, MessageBody=body) + bump("valid_ok") + + # batch до 10 + entries = [{"Id": str(i), "MessageBody": f"b{i}-{random.randint(0, 10**6)}"} for i in range(random.randint(1, 10))] + r = sqs.send_message_batch(QueueUrl=url, Entries=entries) + assert len(r["Successful"]) == len(entries) + + # receive с long poll, сверка целостности + msgs = None + for _ in range(6): + msgs = sqs.receive_message(QueueUrl=url, MaxNumberOfMessages=10, WaitTimeSeconds=random.choice([1, 5, 20])).get("Messages", []) + if msgs: + break + time.sleep(1) + if not msgs: + raise AssertionError("сообщения не получены") + bodies = {m["Body"] for m in msgs} + if body not in bodies: + # тело могло уйти в другой receive предыдущего потока — считаем предупреждением + bump("soft_warn") + + # change visibility + delete + sqs.change_message_visibility(QueueUrl=url, ReceiptHandle=msgs[0]["ReceiptHandle"], VisibilityTimeout=1) + sqs.delete_message(QueueUrl=url, ReceiptHandle=msgs[0]["ReceiptHandle"]) + handles = [{"Id": str(i), "ReceiptHandle": m["ReceiptHandle"]} for i, m in enumerate(msgs[1:])] + if handles: + sqs.delete_message_batch(QueueUrl=url, Entries=handles) + + # теги + sqs.tag_queue(QueueUrl=url, Tags={"load": "1"}) + sqs.list_queue_tags(QueueUrl=url) + sqs.untag_queue(QueueUrl=url, TagKeys=["load"]) + + # атрибуты + attrs = sqs.get_queue_attributes(QueueUrl=url, AttributeNames=["All"])["Attributes"] + assert "VisibilityTimeout" in attrs + sqs.set_queue_attributes(QueueUrl=url, Attributes={"VisibilityTimeout": str(random.randint(0, 43200))}) + finally: + try: + sqs.delete_queue(QueueUrl=url) + except ClientError: + pass + + +def fifo_cycle(sqs, wid, it): + q = f"load-{wid}-{it}.fifo" + url = sqs.create_queue(QueueName=q, Attributes={"FifoQueue": "true"})["QueueUrl"] + try: + sqs.send_message(QueueUrl=url, MessageBody="fifo", MessageGroupId="g", MessageDeduplicationId=str(it)) + # дубль с тем же DedupId — сервис должен не создать дубль (допустим любой из ответов) + sqs.send_message(QueueUrl=url, MessageBody="fifo-dup", MessageGroupId="g", MessageDeduplicationId=str(it)) + msgs = sqs.receive_message(QueueUrl=url, MaxNumberOfMessages=10).get("Messages", []) + if msgs: + sqs.delete_message(QueueUrl=url, ReceiptHandle=msgs[0]["ReceiptHandle"]) + bump("valid_ok") + finally: + try: + sqs.delete_queue(QueueUrl=url) + except ClientError: + pass + + +def invalid_cases(sqs): + cases = [] + # несуществующая очередь + cases.append(lambda: sqs.send_message(QueueUrl=ENDPOINT + "/no-such-queue", MessageBody="x")) + cases.append(lambda: sqs.receive_message(QueueUrl=ENDPOINT + "/no-such-queue")) + cases.append(lambda: sqs.delete_queue(QueueUrl=ENDPOINT + "/no-such-queue")) + cases.append(lambda: sqs.get_queue_url(QueueName="no-such-queue-load-test")) + # битый ReceiptHandle + cases.append(lambda: sqs.delete_message(QueueUrl=ENDPOINT + "/no-such-queue", ReceiptHandle="broken")) + # некорректные атрибуты + q = sqs.create_queue(QueueName=f"invalid-{random.randint(0, 10**9)}")["QueueUrl"] + try: + cases.append(lambda: sqs.set_queue_attributes(QueueUrl=q, Attributes={"VisibilityTimeout": "43201"})) + cases.append(lambda: sqs.set_queue_attributes(QueueUrl=q, Attributes={"VisibilityTimeout": "not-a-number"})) + cases.append(lambda: sqs.set_queue_attributes(QueueUrl=q, Attributes={"BadAttribute": "1"})) + # повторный purge + def purge_twice(): + sqs.purge_queue(QueueUrl=q) + sqs.purge_queue(QueueUrl=q) + cases.append(purge_twice) + # batch-нарушения + cases.append(lambda: sqs.send_message_batch(QueueUrl=q, Entries=[{"Id": str(i), "MessageBody": "x"} for i in range(11)])) + cases.append(lambda: sqs.send_message_batch(QueueUrl=q, Entries=[{"Id": "dup", "MessageBody": "a"}, {"Id": "dup", "MessageBody": "b"}])) + cases.append(lambda: sqs.send_message_batch(QueueUrl=q, Entries=[])) + # пустое тело + cases.append(lambda: sqs.send_message(QueueUrl=q, MessageBody="")) + # слишком длинное сообщение + cases.append(lambda: sqs.send_message(QueueUrl=q, MessageBody="x" * (300 * 1024))) + # FIFO без MessageGroupId + fq = sqs.create_queue(QueueName=f"invalid-{random.randint(0, 10**9)}.fifo", Attributes={"FifoQueue": "true"})["QueueUrl"] + try: + cases.append(lambda: sqs.send_message(QueueUrl=fq, MessageBody="x", MessageDeduplicationId="d")) + finally: + sqs.delete_queue(QueueUrl=fq) + finally: + sqs.delete_queue(QueueUrl=q) + + for c in cases: + t0 = time.time() + try: + c() + bump("soft_warn") + except ClientError as e: + add_latency(time.time() - t0) + code = e.response.get("Error", {}).get("Code", "Unknown") + if code in EXPECTED_CODES: + record_expected(e) + else: + record_unexpected(e) + except Exception as e: + record_unexpected(e) + + +def worker(wid): + sqs = make_client() + it = 0 + while not stop.is_set(): + it += 1 + t0 = time.time() + try: + valid_cycle(sqs, wid, it) + add_latency(time.time() - t0) + except ClientError as e: + code = e.response.get("Error", {}).get("Code", "Unknown") + if code in EXPECTED_CODES: + record_expected(e) + else: + record_unexpected(e) + except Exception as e: + record_unexpected(e) + + t0 = time.time() + try: + fifo_cycle(sqs, wid, it) + add_latency(time.time() - t0) + except ClientError as e: + code = e.response.get("Error", {}).get("Code", "Unknown") + if code in EXPECTED_CODES: + record_expected(e) + else: + record_unexpected(e) + except Exception as e: + record_unexpected(e) + + try: + invalid_cases(sqs) + except Exception as e: + record_unexpected(e) + + time.sleep(random.uniform(0.01, 0.05)) + + +def progress_printer(): + while not stop.is_set(): + stop.wait(30) + with stats_lock: + s = dict(stats) + el = int(time.time() - start_ts) + print(f"[{el:>5}s] ops={s['ops_total']} ok={s['valid_ok']} " + f"expected_err={s['expected_err']} unexpected={s['unexpected_fail']} " + f"soft_warn={s['soft_warn']} health_fail={s['health_fail']}", flush=True) + + +def main(): + print(f"Нагрузочный тест: {ENDPOINT}, длительность {DURATION}s, воркеров {WORKERS}", flush=True) + if not os.environ.get("AWS_ACCESS_KEY_ID") or not os.environ.get("AWS_SECRET_ACCESS_KEY"): + print("Нужны AWS_ACCESS_KEY_ID и AWS_SECRET_ACCESS_KEY") + return 1 + + threads = [threading.Thread(target=worker, args=(i,), daemon=True) for i in range(WORKERS)] + hthread = threading.Thread(target=health_monitor, daemon=True) + pthread = threading.Thread(target=progress_printer, daemon=True) + for t in threads: + t.start() + hthread.start() + pthread.start() + + stop.wait(DURATION) + stop.set() + for t in threads: + t.join(timeout=30) + + el = int(time.time() - start_ts) + with stats_lock: + s = dict(stats) + print(f"\n=== Завершено через {el}s ===", flush=True) + print(f"Операций всего: {s['ops_total']}") + print(f"Корректных (ok): {s['valid_ok']}") + print(f"Ожидаемых отказов: {s['expected_err']}") + print(f"НЕОЖИДАННЫХ отказов: {s['unexpected_fail']}") + print(f"Некорректное прошло: {s['soft_warn']}") + print(f"Сбоев /health: {s['health_fail']}") + if s["err_by_code"]: + print("Коды ошибок:") + for code, n in sorted(s["err_by_code"].items(), key=lambda x: -x[1]): + print(f" {code}: {n}") + with LATENCY_LOCK: + if LATENCY: + print(f"Латентность операций: min={min(LATENCY):.3f}s avg={sum(LATENCY)/len(LATENCY):.3f}s max={max(LATENCY):.3f}s") + + verdict = "УСТОЙЧИВО" if s["unexpected_fail"] == 0 and s["health_fail"] == 0 else "ЕСТЬ ПРОБЛЕМЫ" + print(f"Вердикт: {verdict}") + return 0 if verdict == "УСТОЙЧИВО" else 1 + + +if __name__ == "__main__": + sys.exit(main())