267 lines
13 KiB
Go
267 lines
13 KiB
Go
// Изменено: 2026-04-11 — фикс: SentTimestamp из m.SentTime, MD5 из кэша
|
||
// ReceiveMessageV1 — получает сообщения из очереди тенанта с поддержкой long polling.
|
||
// Ловушка #4: long polling держит соединение до 20 сек — не прерываем принудительно.
|
||
package gosqs
|
||
|
||
import (
|
||
"fmt"
|
||
"net/http"
|
||
"strings"
|
||
"time"
|
||
|
||
"github.com/google/uuid"
|
||
|
||
"shared-sqs/app/interfaces"
|
||
"shared-sqs/app/models"
|
||
"shared-sqs/app/utils"
|
||
|
||
"github.com/gorilla/mux"
|
||
log "github.com/sirupsen/logrus"
|
||
)
|
||
|
||
// 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 {
|
||
log.Error("Invalid Request - ReceiveMessageV1")
|
||
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 = MinNumberOfMessagesLimit
|
||
}
|
||
if maxNumberOfMessages < MinNumberOfMessagesLimit || maxNumberOfMessages > MaxNumberOfMessagesLimit {
|
||
return utils.CreateErrorResponseV1("InvalidParameterValue", true)
|
||
}
|
||
|
||
// 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)
|
||
queueName = vars["queueName"]
|
||
} else {
|
||
uriSegments := strings.Split(requestBody.QueueUrl, "/")
|
||
queueName = uriSegments[len(uriSegments)-1]
|
||
}
|
||
|
||
key := tenantQueueKey(t.AccessKey, queueName)
|
||
|
||
// ОДИН RLock на проверку существования И чтение атрибута очереди:
|
||
// раньше было два отдельных RLock — между ними очередь могла быть удалена
|
||
// (nil-deref при чтении ReceiveMessageWaitTimeSeconds).
|
||
models.SyncQueues.RLock()
|
||
queue, queueExists := models.SyncQueues.Queues[key]
|
||
if !queueExists {
|
||
models.SyncQueues.RUnlock()
|
||
return utils.CreateErrorResponseV1("QueueNotFound", true)
|
||
}
|
||
waitTimeSeconds := requestBody.WaitTimeSeconds
|
||
if waitTimeSeconds == 0 {
|
||
waitTimeSeconds = queue.ReceiveMessageWaitTimeSeconds
|
||
}
|
||
models.SyncQueues.RUnlock()
|
||
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 {
|
||
// Атомарная проверка+выборка: между проверкой и выдачей другой
|
||
// получатель не может «съесть» сообщения — они помечаются
|
||
// ReceiptHandle в том же критическом участке (см. receiveMessagesUnderLock).
|
||
messages, found := receiveMessagesUnderLock(key, requestBody.VisibilityTimeout, maxNumberOfMessages)
|
||
if found {
|
||
return http.StatusOK, buildReceiveResponse(messages)
|
||
}
|
||
if time.Now().After(deadline) {
|
||
return http.StatusOK, buildReceiveResponse(nil)
|
||
}
|
||
|
||
select {
|
||
case <-req.Context().Done():
|
||
// Клиент разорвал соединение — завершаем без ошибки.
|
||
return http.StatusOK, buildReceiveResponse(nil)
|
||
case <-pollTicker.C:
|
||
}
|
||
}
|
||
}
|
||
|
||
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()
|
||
|
||
queue, queueFound := models.SyncQueues.Queues[key]
|
||
if !queueFound || len(queue.Messages) == 0 {
|
||
return nil, false
|
||
}
|
||
// Быстрая предпроверка: есть ли вообще хотя бы одно готовое сообщение.
|
||
if !queueHasReceivableMessages(queue) {
|
||
return nil, false
|
||
}
|
||
|
||
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,
|
||
Body: m.MessageBody,
|
||
ReceiptHandle: m.ReceiptHandle,
|
||
MD5OfBody: m.MD5OfMessageBody, // Используем кэшированный MD5 вместо пересчёта
|
||
MD5OfMessageAttributes: m.MD5OfMessageAttributes,
|
||
MessageAttributes: m.MessageAttributes,
|
||
Attributes: map[string]string{
|
||
"ApproximateFirstReceiveTimestamp": fmt.Sprintf("%d", m.ReceiptTime.UnixNano()/int64(time.Millisecond)),
|
||
"SenderId": models.CurrentEnvironment.AccountID,
|
||
"ApproximateReceiveCount": fmt.Sprintf("%d", m.NumberOfReceives+1),
|
||
"SentTimestamp": fmt.Sprintf("%d", m.SentTime.UnixNano()/int64(time.Millisecond)), // Фикс: реальное время отправки
|
||
},
|
||
}
|
||
}
|
||
|
||
// queueHasReceivableMessages — проверяет, есть ли в очереди хотя бы ОДНО готовое
|
||
// к выдаче сообщение: не in-flight, задержка доставки истекла, FIFO-группа не
|
||
// заблокирована. Используется как быстрая предпроверка перед выборкой.
|
||
//
|
||
// ВАЖНО: вызывается ТОЛЬКО под удержанием блокировки (Lock или RLock) —
|
||
// читает общее состояние очереди.
|
||
func queueHasReceivableMessages(queue *models.Queue) bool {
|
||
for i := range queue.Messages {
|
||
msg := &queue.Messages[i]
|
||
if msg.ReceiptHandle != "" {
|
||
continue
|
||
}
|
||
if !msg.IsReadyForReceipt() {
|
||
continue
|
||
}
|
||
if queue.IsFIFO && queue.IsLocked(msg.GroupID) {
|
||
continue
|
||
}
|
||
return true
|
||
}
|
||
return false
|
||
}
|