fix(shared-sqs): deadlock in create_queue, UI sent_at bug, Redis deploy, TLS ingress (v0.1.13-v0.1.14)
- fix: add SyncQueues.Unlock() before return in create_queue.go happy path (deadlock after first CreateQueue) - fix: UI m.sent -> m.sent_at (message dates always showed as dash) - feat: add deployments/k8s/redis.yaml (Redis persistence) - chore: update deployment image to naeel/shared-sqs:v0.1.14 - chore: update ingress.yaml (TLS, qu.kube5s.ru) - docs: add thinking log 2026-04-10
This commit is contained in:
@@ -4,63 +4,65 @@
|
||||
package gosqs
|
||||
|
||||
import (
|
||||
"net/http"
|
||||
"time"
|
||||
"net/http"
|
||||
"time"
|
||||
|
||||
"shared-sqs/app/interfaces"
|
||||
"shared-sqs/app/models"
|
||||
"shared-sqs/app/persistence"
|
||||
"shared-sqs/app/utils"
|
||||
log "github.com/sirupsen/logrus"
|
||||
"shared-sqs/app/interfaces"
|
||||
"shared-sqs/app/models"
|
||||
"shared-sqs/app/persistence"
|
||||
"shared-sqs/app/utils"
|
||||
|
||||
log "github.com/sirupsen/logrus"
|
||||
)
|
||||
|
||||
func CreateQueueV1(req *http.Request) (int, interfaces.AbstractResponseBody) {
|
||||
requestBody := models.NewCreateQueueRequest()
|
||||
ok := utils.REQUEST_TRANSFORMER(requestBody, req, false)
|
||||
if !ok {
|
||||
log.Error("Invalid Request - CreateQueueV1")
|
||||
return utils.CreateErrorResponseV1("InvalidParameterValue", true)
|
||||
}
|
||||
requestBody := models.NewCreateQueueRequest()
|
||||
ok := utils.REQUEST_TRANSFORMER(requestBody, req, false)
|
||||
if !ok {
|
||||
log.Error("Invalid Request - CreateQueueV1")
|
||||
return utils.CreateErrorResponseV1("InvalidParameterValue", true)
|
||||
}
|
||||
|
||||
t := getTenantFromContext(req)
|
||||
if t == nil {
|
||||
return utils.CreateErrorResponseV1("InvalidClientTokenId", true)
|
||||
}
|
||||
t := getTenantFromContext(req)
|
||||
if t == nil {
|
||||
return utils.CreateErrorResponseV1("InvalidClientTokenId", true)
|
||||
}
|
||||
|
||||
// Ловушка #8: передаём queueName (не key) в HasFIFOQueueName — иначе .fifo не определится
|
||||
queueName := requestBody.QueueName
|
||||
key := tenantQueueKey(t.AccessKey, queueName)
|
||||
queueUrl := tenantQueueURL(t, queueName)
|
||||
queueArn := tenantQueueARN(t, queueName)
|
||||
// Ловушка #8: передаём queueName (не key) в HasFIFOQueueName — иначе .fifo не определится
|
||||
queueName := requestBody.QueueName
|
||||
key := tenantQueueKey(t.AccessKey, queueName)
|
||||
queueUrl := tenantQueueURL(t, queueName)
|
||||
queueArn := tenantQueueARN(t, queueName)
|
||||
|
||||
models.SyncQueues.Lock()
|
||||
if _, exists := models.SyncQueues.Queues[key]; !exists {
|
||||
// Проверка лимита очередей тенанта
|
||||
if t.MaxQueues > 0 && countTenantQueues(t.AccessKey) >= t.MaxQueues {
|
||||
models.SyncQueues.Unlock()
|
||||
return utils.CreateErrorResponseV1("LimitExceeded", true)
|
||||
}
|
||||
log.Infof("Creating Queue: %s (tenant: %s)", queueName, t.ID)
|
||||
queue := &models.Queue{
|
||||
Name: queueName,
|
||||
URL: queueUrl,
|
||||
Arn: queueArn,
|
||||
IsFIFO: utils.HasFIFOQueueName(queueName),
|
||||
EnableDuplicates: models.CurrentEnvironment.EnableDuplicates,
|
||||
Duplicates: make(map[string]time.Time),
|
||||
}
|
||||
if err := setQueueAttributesV1(queue, requestBody.Attributes); err != nil {
|
||||
models.SyncQueues.Unlock()
|
||||
return utils.CreateErrorResponseV1(err.Error(), true)
|
||||
}
|
||||
models.SyncQueues.Queues[key] = queue
|
||||
}
|
||||
models.SyncQueues.Lock()
|
||||
if _, exists := models.SyncQueues.Queues[key]; !exists {
|
||||
// Проверка лимита очередей тенанта
|
||||
if t.MaxQueues > 0 && countTenantQueues(t.AccessKey) >= t.MaxQueues {
|
||||
models.SyncQueues.Unlock()
|
||||
return utils.CreateErrorResponseV1("LimitExceeded", true)
|
||||
}
|
||||
log.Infof("Creating Queue: %s (tenant: %s)", queueName, t.ID)
|
||||
queue := &models.Queue{
|
||||
Name: queueName,
|
||||
URL: queueUrl,
|
||||
Arn: queueArn,
|
||||
IsFIFO: utils.HasFIFOQueueName(queueName),
|
||||
EnableDuplicates: models.CurrentEnvironment.EnableDuplicates,
|
||||
Duplicates: make(map[string]time.Time),
|
||||
}
|
||||
if err := setQueueAttributesV1(queue, requestBody.Attributes); err != nil {
|
||||
models.SyncQueues.Unlock()
|
||||
return utils.CreateErrorResponseV1(err.Error(), true)
|
||||
}
|
||||
models.SyncQueues.Queues[key] = queue
|
||||
}
|
||||
// Сохраняем очередь в Redis пока держим Lock — консистентный снапшот
|
||||
persistence.SaveQueue(key, models.SyncQueues.Queues[key])
|
||||
respStruct := models.CreateQueueResponse{
|
||||
Xmlns: models.BaseXmlns,
|
||||
Result: models.CreateQueueResult{QueueUrl: queueUrl},
|
||||
Metadata: models.BaseResponseMetadata,
|
||||
}
|
||||
return http.StatusOK, respStruct
|
||||
models.SyncQueues.Unlock()
|
||||
respStruct := models.CreateQueueResponse{
|
||||
Xmlns: models.BaseXmlns,
|
||||
Result: models.CreateQueueResult{QueueUrl: queueUrl},
|
||||
Metadata: models.BaseResponseMetadata,
|
||||
}
|
||||
return http.StatusOK, respStruct
|
||||
}
|
||||
|
||||
@@ -3,45 +3,46 @@
|
||||
package gosqs
|
||||
|
||||
import (
|
||||
"net/http"
|
||||
"strings"
|
||||
"net/http"
|
||||
"strings"
|
||||
|
||||
"shared-sqs/app/interfaces"
|
||||
"shared-sqs/app/models"
|
||||
"shared-sqs/app/persistence"
|
||||
"shared-sqs/app/utils"
|
||||
log "github.com/sirupsen/logrus"
|
||||
"shared-sqs/app/interfaces"
|
||||
"shared-sqs/app/models"
|
||||
"shared-sqs/app/persistence"
|
||||
"shared-sqs/app/utils"
|
||||
|
||||
log "github.com/sirupsen/logrus"
|
||||
)
|
||||
|
||||
func DeleteQueueV1(req *http.Request) (int, interfaces.AbstractResponseBody) {
|
||||
requestBody := models.NewDeleteQueueRequest()
|
||||
ok := utils.REQUEST_TRANSFORMER(requestBody, req, false)
|
||||
if !ok {
|
||||
log.Error("Invalid Request - DeleteQueueV1")
|
||||
return utils.CreateErrorResponseV1("InvalidParameterValue", true)
|
||||
}
|
||||
|
||||
t := getTenantFromContext(req)
|
||||
if t == nil {
|
||||
return utils.CreateErrorResponseV1("InvalidClientTokenId", true)
|
||||
}
|
||||
|
||||
uriSegments := strings.Split(requestBody.QueueUrl, "/")
|
||||
queueName := uriSegments[len(uriSegments)-1]
|
||||
key := tenantQueueKey(t.AccessKey, queueName)
|
||||
|
||||
log.Infof("Deleting Queue: %s (tenant: %s)", queueName, t.ID)
|
||||
|
||||
models.SyncQueues.Lock()
|
||||
delete(models.SyncQueues.Queues, key)
|
||||
models.SyncQueues.Unlock()
|
||||
|
||||
// Удаляем из Redis асинхронно
|
||||
persistence.DeleteQueue(key)
|
||||
|
||||
respStruct := models.DeleteQueueResponse{
|
||||
Xmlns: models.BaseXmlns,
|
||||
Metadata: models.BaseResponseMetadata,
|
||||
}
|
||||
return http.StatusOK, respStruct
|
||||
requestBody := models.NewDeleteQueueRequest()
|
||||
ok := utils.REQUEST_TRANSFORMER(requestBody, req, false)
|
||||
if !ok {
|
||||
log.Error("Invalid Request - DeleteQueueV1")
|
||||
return utils.CreateErrorResponseV1("InvalidParameterValue", true)
|
||||
}
|
||||
|
||||
t := getTenantFromContext(req)
|
||||
if t == nil {
|
||||
return utils.CreateErrorResponseV1("InvalidClientTokenId", true)
|
||||
}
|
||||
|
||||
uriSegments := strings.Split(requestBody.QueueUrl, "/")
|
||||
queueName := uriSegments[len(uriSegments)-1]
|
||||
key := tenantQueueKey(t.AccessKey, queueName)
|
||||
|
||||
log.Infof("Deleting Queue: %s (tenant: %s)", queueName, t.ID)
|
||||
|
||||
models.SyncQueues.Lock()
|
||||
delete(models.SyncQueues.Queues, key)
|
||||
models.SyncQueues.Unlock()
|
||||
|
||||
// Удаляем из Redis асинхронно
|
||||
persistence.DeleteQueue(key)
|
||||
|
||||
respStruct := models.DeleteQueueResponse{
|
||||
Xmlns: models.BaseXmlns,
|
||||
Metadata: models.BaseResponseMetadata,
|
||||
}
|
||||
return http.StatusOK, respStruct
|
||||
}
|
||||
|
||||
@@ -3,50 +3,51 @@
|
||||
package gosqs
|
||||
|
||||
import (
|
||||
"net/http"
|
||||
"strings"
|
||||
"time"
|
||||
"net/http"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
"shared-sqs/app/interfaces"
|
||||
"shared-sqs/app/models"
|
||||
"shared-sqs/app/persistence"
|
||||
"shared-sqs/app/utils"
|
||||
log "github.com/sirupsen/logrus"
|
||||
"shared-sqs/app/interfaces"
|
||||
"shared-sqs/app/models"
|
||||
"shared-sqs/app/persistence"
|
||||
"shared-sqs/app/utils"
|
||||
|
||||
log "github.com/sirupsen/logrus"
|
||||
)
|
||||
|
||||
func PurgeQueueV1(req *http.Request) (int, interfaces.AbstractResponseBody) {
|
||||
requestBody := models.NewPurgeQueueRequest()
|
||||
ok := utils.REQUEST_TRANSFORMER(requestBody, req, false)
|
||||
if !ok {
|
||||
log.Error("Invalid Request - PurgeQueueV1")
|
||||
return utils.CreateErrorResponseV1("InvalidParameterValue", true)
|
||||
}
|
||||
requestBody := models.NewPurgeQueueRequest()
|
||||
ok := utils.REQUEST_TRANSFORMER(requestBody, req, false)
|
||||
if !ok {
|
||||
log.Error("Invalid Request - PurgeQueueV1")
|
||||
return utils.CreateErrorResponseV1("InvalidParameterValue", true)
|
||||
}
|
||||
|
||||
t := getTenantFromContext(req)
|
||||
if t == nil {
|
||||
return utils.CreateErrorResponseV1("InvalidClientTokenId", true)
|
||||
}
|
||||
t := getTenantFromContext(req)
|
||||
if t == nil {
|
||||
return utils.CreateErrorResponseV1("InvalidClientTokenId", true)
|
||||
}
|
||||
|
||||
uriSegments := strings.Split(requestBody.QueueUrl, "/")
|
||||
queueName := uriSegments[len(uriSegments)-1]
|
||||
key := tenantQueueKey(t.AccessKey, queueName)
|
||||
uriSegments := strings.Split(requestBody.QueueUrl, "/")
|
||||
queueName := uriSegments[len(uriSegments)-1]
|
||||
key := tenantQueueKey(t.AccessKey, queueName)
|
||||
|
||||
models.SyncQueues.Lock()
|
||||
defer models.SyncQueues.Unlock()
|
||||
if _, ok := models.SyncQueues.Queues[key]; !ok {
|
||||
log.Errorf("Purge Queue: %s, queue does not exist for tenant %s", queueName, t.ID)
|
||||
return utils.CreateErrorResponseV1("QueueNotFound", true)
|
||||
}
|
||||
models.SyncQueues.Lock()
|
||||
defer models.SyncQueues.Unlock()
|
||||
if _, ok := models.SyncQueues.Queues[key]; !ok {
|
||||
log.Errorf("Purge Queue: %s, queue does not exist for tenant %s", queueName, t.ID)
|
||||
return utils.CreateErrorResponseV1("QueueNotFound", 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)
|
||||
// Сохраняем пустую очередь в Redis пока держим Lock
|
||||
persistence.SaveQueue(key, models.SyncQueues.Queues[key])
|
||||
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)
|
||||
// Сохраняем пустую очередь в Redis пока держим Lock
|
||||
persistence.SaveQueue(key, models.SyncQueues.Queues[key])
|
||||
|
||||
respStruct := models.PurgeQueueResponse{
|
||||
Xmlns: models.BaseXmlns,
|
||||
Metadata: models.BaseResponseMetadata,
|
||||
}
|
||||
return http.StatusOK, respStruct
|
||||
respStruct := models.PurgeQueueResponse{
|
||||
Xmlns: models.BaseXmlns,
|
||||
Metadata: models.BaseResponseMetadata,
|
||||
}
|
||||
return http.StatusOK, respStruct
|
||||
}
|
||||
|
||||
@@ -5,107 +5,107 @@
|
||||
package gosqs
|
||||
|
||||
import (
|
||||
"net/http"
|
||||
"strings"
|
||||
"time"
|
||||
"net/http"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
"github.com/google/uuid"
|
||||
"github.com/google/uuid"
|
||||
|
||||
"shared-sqs/app/interfaces"
|
||||
"shared-sqs/app/models"
|
||||
"shared-sqs/app/persistence"
|
||||
"shared-sqs/app/utils"
|
||||
"shared-sqs/app/interfaces"
|
||||
"shared-sqs/app/models"
|
||||
"shared-sqs/app/persistence"
|
||||
"shared-sqs/app/utils"
|
||||
|
||||
log "github.com/sirupsen/logrus"
|
||||
log "github.com/sirupsen/logrus"
|
||||
|
||||
"github.com/gorilla/mux"
|
||||
"github.com/gorilla/mux"
|
||||
)
|
||||
|
||||
func SendMessageV1(req *http.Request) (int, interfaces.AbstractResponseBody) {
|
||||
requestBody := models.NewSendMessageRequest()
|
||||
ok := utils.REQUEST_TRANSFORMER(requestBody, req, false)
|
||||
if !ok {
|
||||
log.Error("Invalid Request - SendMessageV1")
|
||||
return utils.CreateErrorResponseV1("InvalidParameterValue", true)
|
||||
}
|
||||
|
||||
t := getTenantFromContext(req)
|
||||
if t == nil {
|
||||
return utils.CreateErrorResponseV1("InvalidClientTokenId", true)
|
||||
}
|
||||
|
||||
messageBody := requestBody.MessageBody
|
||||
messageGroupID := requestBody.MessageGroupId
|
||||
messageDeduplicationID := requestBody.MessageDeduplicationId
|
||||
|
||||
queueUrl := getQueueFromPath(requestBody.QueueUrl, req.URL.String())
|
||||
queueName := ""
|
||||
if queueUrl == "" {
|
||||
vars := mux.Vars(req)
|
||||
queueName = vars["queueName"]
|
||||
} else {
|
||||
// Ловушка #6: берём последний сегмент — это queueName, не tenantID
|
||||
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 models.SyncQueues.Queues[key].MaximumMessageSize > 0 &&
|
||||
len(messageBody) > models.SyncQueues.Queues[key].MaximumMessageSize {
|
||||
return utils.CreateErrorResponseV1("MessageTooBig", true)
|
||||
}
|
||||
|
||||
delaySecs := models.SyncQueues.Queues[key].DelaySeconds
|
||||
if requestBody.DelaySeconds != 0 {
|
||||
delaySecs = requestBody.DelaySeconds
|
||||
}
|
||||
|
||||
log.Debugf("Putting Message in Queue: [%s] tenant: [%s]", queueName, t.ID)
|
||||
msg := models.SqsMessage{MessageBody: messageBody}
|
||||
if len(requestBody.MessageAttributes) > 0 {
|
||||
msg.MessageAttributes = requestBody.MessageAttributes
|
||||
msg.MD5OfMessageAttributes = utils.HashAttributes(requestBody.MessageAttributes)
|
||||
}
|
||||
msg.MD5OfMessageBody = utils.GetMD5Hash(messageBody)
|
||||
msg.Uuid = uuid.NewString()
|
||||
msg.GroupID = messageGroupID
|
||||
msg.DeduplicationID = messageDeduplicationID
|
||||
msg.SentTime = time.Now()
|
||||
msg.DelaySecs = delaySecs
|
||||
|
||||
models.SyncQueues.Lock()
|
||||
fifoSeqNumber := ""
|
||||
if models.SyncQueues.Queues[key].IsFIFO {
|
||||
fifoSeqNumber = models.SyncQueues.Queues[key].NextSequenceNumber(messageGroupID)
|
||||
}
|
||||
|
||||
if !models.SyncQueues.Queues[key].IsDuplicate(messageDeduplicationID) {
|
||||
models.SyncQueues.Queues[key].Messages = append(models.SyncQueues.Queues[key].Messages, msg)
|
||||
} else {
|
||||
log.Debugf("Duplicate message deduplicationId [%s] in queue [%s]", messageDeduplicationID, queueName)
|
||||
}
|
||||
|
||||
models.SyncQueues.Queues[key].InitDuplicatation(messageDeduplicationID)
|
||||
// Сохраняем очередь в Redis пока держим Lock
|
||||
persistence.SaveQueue(key, models.SyncQueues.Queues[key])
|
||||
models.SyncQueues.Unlock()
|
||||
log.Infof("%s: Queue: %s, Message: %s\n", time.Now().Format("2006-01-02 15:04:05"), queueName, msg.MessageBody)
|
||||
|
||||
respStruct := models.SendMessageResponse{
|
||||
Xmlns: models.BaseXmlns,
|
||||
Result: models.SendMessageResult{
|
||||
MD5OfMessageAttributes: msg.MD5OfMessageAttributes,
|
||||
MD5OfMessageBody: msg.MD5OfMessageBody,
|
||||
MessageId: msg.Uuid,
|
||||
SequenceNumber: fifoSeqNumber,
|
||||
},
|
||||
Metadata: models.BaseResponseMetadata,
|
||||
}
|
||||
|
||||
return http.StatusOK, respStruct
|
||||
requestBody := models.NewSendMessageRequest()
|
||||
ok := utils.REQUEST_TRANSFORMER(requestBody, req, false)
|
||||
if !ok {
|
||||
log.Error("Invalid Request - SendMessageV1")
|
||||
return utils.CreateErrorResponseV1("InvalidParameterValue", true)
|
||||
}
|
||||
|
||||
t := getTenantFromContext(req)
|
||||
if t == nil {
|
||||
return utils.CreateErrorResponseV1("InvalidClientTokenId", true)
|
||||
}
|
||||
|
||||
messageBody := requestBody.MessageBody
|
||||
messageGroupID := requestBody.MessageGroupId
|
||||
messageDeduplicationID := requestBody.MessageDeduplicationId
|
||||
|
||||
queueUrl := getQueueFromPath(requestBody.QueueUrl, req.URL.String())
|
||||
queueName := ""
|
||||
if queueUrl == "" {
|
||||
vars := mux.Vars(req)
|
||||
queueName = vars["queueName"]
|
||||
} else {
|
||||
// Ловушка #6: берём последний сегмент — это queueName, не tenantID
|
||||
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 models.SyncQueues.Queues[key].MaximumMessageSize > 0 &&
|
||||
len(messageBody) > models.SyncQueues.Queues[key].MaximumMessageSize {
|
||||
return utils.CreateErrorResponseV1("MessageTooBig", true)
|
||||
}
|
||||
|
||||
delaySecs := models.SyncQueues.Queues[key].DelaySeconds
|
||||
if requestBody.DelaySeconds != 0 {
|
||||
delaySecs = requestBody.DelaySeconds
|
||||
}
|
||||
|
||||
log.Debugf("Putting Message in Queue: [%s] tenant: [%s]", queueName, t.ID)
|
||||
msg := models.SqsMessage{MessageBody: messageBody}
|
||||
if len(requestBody.MessageAttributes) > 0 {
|
||||
msg.MessageAttributes = requestBody.MessageAttributes
|
||||
msg.MD5OfMessageAttributes = utils.HashAttributes(requestBody.MessageAttributes)
|
||||
}
|
||||
msg.MD5OfMessageBody = utils.GetMD5Hash(messageBody)
|
||||
msg.Uuid = uuid.NewString()
|
||||
msg.GroupID = messageGroupID
|
||||
msg.DeduplicationID = messageDeduplicationID
|
||||
msg.SentTime = time.Now()
|
||||
msg.DelaySecs = delaySecs
|
||||
|
||||
models.SyncQueues.Lock()
|
||||
fifoSeqNumber := ""
|
||||
if models.SyncQueues.Queues[key].IsFIFO {
|
||||
fifoSeqNumber = models.SyncQueues.Queues[key].NextSequenceNumber(messageGroupID)
|
||||
}
|
||||
|
||||
if !models.SyncQueues.Queues[key].IsDuplicate(messageDeduplicationID) {
|
||||
models.SyncQueues.Queues[key].Messages = append(models.SyncQueues.Queues[key].Messages, msg)
|
||||
} else {
|
||||
log.Debugf("Duplicate message deduplicationId [%s] in queue [%s]", messageDeduplicationID, queueName)
|
||||
}
|
||||
|
||||
models.SyncQueues.Queues[key].InitDuplicatation(messageDeduplicationID)
|
||||
// Сохраняем очередь в Redis пока держим Lock
|
||||
persistence.SaveQueue(key, models.SyncQueues.Queues[key])
|
||||
models.SyncQueues.Unlock()
|
||||
log.Infof("%s: Queue: %s, Message: %s\n", time.Now().Format("2006-01-02 15:04:05"), queueName, msg.MessageBody)
|
||||
|
||||
respStruct := models.SendMessageResponse{
|
||||
Xmlns: models.BaseXmlns,
|
||||
Result: models.SendMessageResult{
|
||||
MD5OfMessageAttributes: msg.MD5OfMessageAttributes,
|
||||
MD5OfMessageBody: msg.MD5OfMessageBody,
|
||||
MessageId: msg.Uuid,
|
||||
SequenceNumber: fifoSeqNumber,
|
||||
},
|
||||
Metadata: models.BaseResponseMetadata,
|
||||
}
|
||||
|
||||
return http.StatusOK, respStruct
|
||||
}
|
||||
|
||||
@@ -4,54 +4,55 @@
|
||||
package gosqs
|
||||
|
||||
import (
|
||||
"net/http"
|
||||
"strings"
|
||||
"net/http"
|
||||
"strings"
|
||||
|
||||
"shared-sqs/app/interfaces"
|
||||
"shared-sqs/app/models"
|
||||
"shared-sqs/app/persistence"
|
||||
"shared-sqs/app/utils"
|
||||
log "github.com/sirupsen/logrus"
|
||||
"shared-sqs/app/interfaces"
|
||||
"shared-sqs/app/models"
|
||||
"shared-sqs/app/persistence"
|
||||
"shared-sqs/app/utils"
|
||||
|
||||
log "github.com/sirupsen/logrus"
|
||||
)
|
||||
|
||||
func SetQueueAttributesV1(req *http.Request) (int, interfaces.AbstractResponseBody) {
|
||||
requestBody := models.NewSetQueueAttributesRequest()
|
||||
ok := utils.REQUEST_TRANSFORMER(requestBody, req, false)
|
||||
if !ok {
|
||||
log.Error("Invalid Request - SetQueueAttributesV1")
|
||||
return utils.CreateErrorResponseV1("InvalidParameterValue", true)
|
||||
}
|
||||
if requestBody.QueueUrl == "" {
|
||||
log.Error("Missing QueueUrl - SetQueueAttributesV1")
|
||||
return utils.CreateErrorResponseV1("InvalidParameterValue", true)
|
||||
}
|
||||
requestBody := models.NewSetQueueAttributesRequest()
|
||||
ok := utils.REQUEST_TRANSFORMER(requestBody, req, false)
|
||||
if !ok {
|
||||
log.Error("Invalid Request - SetQueueAttributesV1")
|
||||
return utils.CreateErrorResponseV1("InvalidParameterValue", true)
|
||||
}
|
||||
if requestBody.QueueUrl == "" {
|
||||
log.Error("Missing QueueUrl - SetQueueAttributesV1")
|
||||
return utils.CreateErrorResponseV1("InvalidParameterValue", true)
|
||||
}
|
||||
|
||||
t := getTenantFromContext(req)
|
||||
if t == nil {
|
||||
return utils.CreateErrorResponseV1("InvalidClientTokenId", true)
|
||||
}
|
||||
t := getTenantFromContext(req)
|
||||
if t == nil {
|
||||
return utils.CreateErrorResponseV1("InvalidClientTokenId", true)
|
||||
}
|
||||
|
||||
uriSegments := strings.Split(requestBody.QueueUrl, "/")
|
||||
queueName := uriSegments[len(uriSegments)-1]
|
||||
key := tenantQueueKey(t.AccessKey, queueName)
|
||||
uriSegments := strings.Split(requestBody.QueueUrl, "/")
|
||||
queueName := uriSegments[len(uriSegments)-1]
|
||||
key := tenantQueueKey(t.AccessKey, queueName)
|
||||
|
||||
log.Infof("Set Queue Attributes: %s (tenant: %s)", queueName, t.ID)
|
||||
models.SyncQueues.Lock()
|
||||
defer models.SyncQueues.Unlock()
|
||||
queue, ok := models.SyncQueues.Queues[key]
|
||||
if !ok {
|
||||
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 {
|
||||
return utils.CreateErrorResponseV1(err.Error(), true)
|
||||
}
|
||||
// Сохраняем атрибуты в Redis пока держим Lock (через defer)
|
||||
persistence.SaveQueue(key, queue)
|
||||
log.Infof("Set Queue Attributes: %s (tenant: %s)", queueName, t.ID)
|
||||
models.SyncQueues.Lock()
|
||||
defer models.SyncQueues.Unlock()
|
||||
queue, ok := models.SyncQueues.Queues[key]
|
||||
if !ok {
|
||||
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 {
|
||||
return utils.CreateErrorResponseV1(err.Error(), true)
|
||||
}
|
||||
// Сохраняем атрибуты в Redis пока держим Lock (через defer)
|
||||
persistence.SaveQueue(key, queue)
|
||||
|
||||
respStruct := models.SetQueueAttributesResponse{
|
||||
Xmlns: models.BaseXmlns,
|
||||
Metadata: models.BaseResponseMetadata,
|
||||
}
|
||||
return http.StatusOK, respStruct
|
||||
respStruct := models.SetQueueAttributesResponse{
|
||||
Xmlns: models.BaseXmlns,
|
||||
Metadata: models.BaseResponseMetadata,
|
||||
}
|
||||
return http.StatusOK, respStruct
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user