v0.1.35: query-mode XML-ответы, TOCTOU receive, валидации AWS, purge 60s, сужение локов

This commit is contained in:
“Naeel”
2026-08-14 15:53:08 +04:00
parent 8c40ca592f
commit ec10acc99e
15 changed files with 919 additions and 172 deletions
+22 -4
View File
@@ -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)
}
+18 -3
View File
@@ -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)
+79 -34
View File
@@ -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 {
+24 -2
View File
@@ -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)
// Сохраняем пустые метаданные очереди
+65 -16
View File
@@ -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,
}
+150 -79
View File
@@ -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: без <Message>).
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]
+28 -4
View File
@@ -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)
}
+21 -4
View File
@@ -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)
+4
View File
@@ -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."},
}
}
+41 -9
View File
@@ -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
+63 -15
View File
@@ -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#<Code>",
// "message": "..."}. Раньше для ErrorResponse писался голый ErrorResult
// ({"Type":...,"Code":...,"Message":...}) — SDK не распознавал это как ошибку
// и возвращал Code=None.
// - AwsQueryProtocol: ответ всегда XML (в т.ч. <ErrorResponse> для ошибок).
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
}
+1 -1
View File
@@ -384,7 +384,7 @@ td.msg-expand { padding: 0 !important; border-bottom: 1px solid var(--border); }
<div>Админ-консоль: <a href="/admin">/admin</a> (вход по admin-токену)</div>
<div>Realm: <span style="font-family:monospace">iot-naeel</span> · Persistence: Managed Redis</div>
<div>Версия: <span id="svc-version" style="font-family:monospace">—</span></div>
<div>Образ: <span style="font-family:monospace">naeel/shared-sqs:v0.1.34</span></div>
<div>Образ: <span style="font-family:monospace">naeel/shared-sqs:v0.1.35</span></div>
<div style="margin-top:8px;color:var(--warning)">Статус: тестирование</div>
</div>
</div>