// Изменено: 2026-04-09 // 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" ) func ReceiveMessageV1(req *http.Request) (int, interfaces.AbstractResponseBody) { requestBody := models.NewReceiveMessageRequest() ok := utils.REQUEST_TRANSFORMER(requestBody, req, false) if !ok { log.Error("Invalid Request - ReceiveMessageV1") return utils.CreateErrorResponseV1("InvalidParameterValue", true) } t := getTenantFromContext(req) if t == nil { return utils.CreateErrorResponseV1("InvalidClientTokenId", true) } maxNumberOfMessages := requestBody.MaxNumberOfMessages if maxNumberOfMessages == 0 { maxNumberOfMessages = 1 } 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) if _, ok := models.SyncQueues.Queues[key]; !ok { return utils.CreateErrorResponseV1("QueueNotFound", true) } var messages []*models.ResultMessage respStruct := models.ReceiveMessageResponse{} waitTimeSeconds := requestBody.WaitTimeSeconds if waitTimeSeconds == 0 { models.SyncQueues.RLock() waitTimeSeconds = models.SyncQueues.Queues[key].ReceiveMessageWaitTimeSeconds models.SyncQueues.RUnlock() } // Long polling: ждём появления сообщения до waitTimeSeconds*10 итераций по 100ms loops := waitTimeSeconds * 10 for loops > 0 { models.SyncQueues.RLock() _, queueFound := models.SyncQueues.Queues[key] if !queueFound { models.SyncQueues.RUnlock() return utils.CreateErrorResponseV1("QueueNotFound", true) } messageFound := len(models.SyncQueues.Queues[key].Messages)-numberOfHiddenMessagesInQueue(*models.SyncQueues.Queues[key]) != 0 models.SyncQueues.RUnlock() if !messageFound { continueTimer := time.NewTimer(100 * time.Millisecond) select { case <-req.Context().Done(): continueTimer.Stop() return http.StatusOK, models.ReceiveMessageResponse{ Xmlns: models.BaseXmlns, Result: models.ReceiveMessageResult{}, Metadata: models.BaseResponseMetadata, } case <-continueTimer.C: continueTimer.Stop() } loops-- } else { break } } log.Debugf("Getting Message from Queue:%s (tenant: %s)", queueName, t.ID) 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"}, } } return http.StatusOK, respStruct } func buildResultMessage(m *models.SqsMessage) *models.ResultMessage { return &models.ResultMessage{ MessageId: m.Uuid, Body: m.MessageBody, ReceiptHandle: m.ReceiptHandle, MD5OfBody: utils.GetMD5Hash(m.MessageBody), MD5OfMessageAttributes: m.MD5OfMessageAttributes, MessageAttributes: m.MessageAttributes, Attributes: map[string]string{ "ApproximateFirstReceiveTimestamp": fmt.Sprintf("%d", m.ReceiptTime.UnixNano()/int64(time.Millisecond)), "SenderId": models.CurrentEnvironment.AccountID, "ApproximateReceiveCount": fmt.Sprintf("%d", m.NumberOfReceives+1), "SentTimestamp": fmt.Sprintf("%d", time.Now().UTC().UnixNano()/int64(time.Millisecond)), }, } }