132 lines
4.2 KiB
Go
132 lines
4.2 KiB
Go
// Создано: 2026-04-11, Изменено: 2026-04-11 — добавлена Redis persistence
|
|
// ChangeMessageVisibilityBatchV1 — пакетная смена таймаута видимости (до 10 сообщений).
|
|
// Паттерн аналогичен DeleteMessageBatchV1: валидация Id, цикл по записям, partial success.
|
|
package gosqs
|
|
|
|
import (
|
|
"net/http"
|
|
"strings"
|
|
"time"
|
|
|
|
"shared-sqs/app/interfaces"
|
|
"shared-sqs/app/models"
|
|
"shared-sqs/app/persistence"
|
|
"shared-sqs/app/utils"
|
|
|
|
"github.com/gorilla/mux"
|
|
log "github.com/sirupsen/logrus"
|
|
)
|
|
|
|
func ChangeMessageVisibilityBatchV1(req *http.Request) (int, interfaces.AbstractResponseBody) {
|
|
requestBody := models.NewChangeMessageVisibilityBatchRequest()
|
|
ok := utils.REQUEST_TRANSFORMER(requestBody, req, false)
|
|
if !ok {
|
|
log.Error("Invalid Request - ChangeMessageVisibilityBatchV1")
|
|
return utils.CreateErrorResponseV1("InvalidParameterValue", true)
|
|
}
|
|
|
|
t := getTenantFromContext(req)
|
|
if t == nil {
|
|
return utils.CreateErrorResponseV1("InvalidClientTokenId", true)
|
|
}
|
|
|
|
queueUrl := requestBody.QueueUrl
|
|
queueName := ""
|
|
if queueUrl == "" {
|
|
vars := mux.Vars(req)
|
|
queueName = vars["queueName"]
|
|
} else {
|
|
uriSegments := strings.Split(queueUrl, "/")
|
|
queueName = uriSegments[len(uriSegments)-1]
|
|
}
|
|
|
|
key := tenantQueueKey(t.AccessKey, queueName)
|
|
|
|
if _, ok := models.SyncQueues.Queues[key]; !ok {
|
|
return utils.CreateErrorResponseV1("QueueNotFound", true)
|
|
}
|
|
|
|
if len(requestBody.Entries) == 0 {
|
|
return utils.CreateErrorResponseV1("EmptyBatchRequest", true)
|
|
}
|
|
|
|
if len(requestBody.Entries) > 10 {
|
|
return utils.CreateErrorResponseV1("TooManyEntriesInBatchRequest", true)
|
|
}
|
|
|
|
// Проверка уникальности Id в пределах запроса
|
|
ids := map[string]bool{}
|
|
for _, entry := range requestBody.Entries {
|
|
if _, found := ids[entry.Id]; found {
|
|
return utils.CreateErrorResponseV1("BatchEntryIdsNotDistinct", true)
|
|
}
|
|
ids[entry.Id] = true
|
|
}
|
|
|
|
models.SyncQueues.Lock()
|
|
defer models.SyncQueues.Unlock()
|
|
|
|
successEntries := make([]models.ChangeMessageVisibilityBatchResultEntry, 0)
|
|
failedEntries := make([]models.BatchResultErrorEntry, 0)
|
|
|
|
for _, entry := range requestBody.Entries {
|
|
if entry.VisibilityTimeout > 43200 {
|
|
failedEntries = append(failedEntries, models.BatchResultErrorEntry{
|
|
Code: "InvalidParameterValue",
|
|
Id: entry.Id,
|
|
Message: "VisibilityTimeout must be between 0 and 43200",
|
|
SenderFault: true,
|
|
})
|
|
continue
|
|
}
|
|
|
|
messageFound := false
|
|
queue := models.SyncQueues.Queues[key]
|
|
for i := 0; i < len(queue.Messages); i++ {
|
|
if queue.Messages[i].ReceiptHandle == entry.ReceiptHandle {
|
|
if entry.VisibilityTimeout == 0 {
|
|
// Сброс: сообщение снова видимо, счётчик retry++
|
|
queue.Messages[i].ReceiptTime = time.Now().UTC()
|
|
queue.Messages[i].ReceiptHandle = ""
|
|
queue.Messages[i].VisibilityTimeout = time.Now().Add(time.Duration(queue.VisibilityTimeout) * time.Second)
|
|
queue.Messages[i].Retry++
|
|
} else {
|
|
queue.Messages[i].VisibilityTimeout = time.Now().Add(time.Duration(entry.VisibilityTimeout) * time.Second)
|
|
}
|
|
// Персистим изменённое сообщение отдельно
|
|
persistence.SaveMessage(key, &queue.Messages[i])
|
|
messageFound = true
|
|
break
|
|
}
|
|
}
|
|
|
|
if messageFound {
|
|
successEntries = append(successEntries, models.ChangeMessageVisibilityBatchResultEntry{Id: entry.Id})
|
|
} else {
|
|
failedEntries = append(failedEntries, models.BatchResultErrorEntry{
|
|
Code: "ReceiptHandleIsInvalid",
|
|
Id: entry.Id,
|
|
Message: "Message not found",
|
|
SenderFault: true,
|
|
})
|
|
}
|
|
}
|
|
|
|
// Персистим изменения в Redis под Lock
|
|
if len(successEntries) > 0 {
|
|
persistence.SaveQueue(key, models.SyncQueues.Queues[key])
|
|
}
|
|
|
|
respStruct := models.ChangeMessageVisibilityBatchResponse{
|
|
Xmlns: models.BaseXmlns,
|
|
Result: models.ChangeMessageVisibilityBatchResult{
|
|
Successful: successEntries,
|
|
Failed: failedEntries,
|
|
},
|
|
Metadata: models.BaseResponseMetadata,
|
|
}
|
|
|
|
log.Debugf("ChangeMessageVisibilityBatch: %s — %d ok, %d failed", queueName, len(successEntries), len(failedEntries))
|
|
return http.StatusOK, respStruct
|
|
}
|