feat: add 4 missing SQS API commands for Yandex/AWS compatibility
- ChangeMessageVisibilityBatch: batch visibility timeout (up to 10 msgs) - TagQueue: add/update queue tags - UntagQueue: remove queue tags by keys - ListQueueTags: list all queue tags - Added Tags field to Queue struct - Request/Response models for all 4 commands - Registered in router (17 total API commands now) - Yandex Message Queue API reference doc
This commit is contained in:
@@ -0,0 +1,123 @@
|
||||
// Создано: 2026-04-11
|
||||
// ChangeMessageVisibilityBatchV1 — пакетная смена таймаута видимости (до 10 сообщений).
|
||||
// Паттерн аналогичен DeleteMessageBatchV1: валидация Id, цикл по записям, partial success.
|
||||
package gosqs
|
||||
|
||||
import (
|
||||
"net/http"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
"shared-sqs/app/interfaces"
|
||||
"shared-sqs/app/models"
|
||||
"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)
|
||||
}
|
||||
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,
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
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
|
||||
}
|
||||
@@ -0,0 +1,65 @@
|
||||
// Создано: 2026-04-11
|
||||
// ListQueueTagsV1 — возвращает все метки (tags) очереди тенанта.
|
||||
// Совместимость с AWS SQS ListQueueTags.
|
||||
package gosqs
|
||||
|
||||
import (
|
||||
"net/http"
|
||||
"strings"
|
||||
|
||||
"shared-sqs/app/interfaces"
|
||||
"shared-sqs/app/models"
|
||||
"shared-sqs/app/utils"
|
||||
|
||||
"github.com/gorilla/mux"
|
||||
log "github.com/sirupsen/logrus"
|
||||
)
|
||||
|
||||
func ListQueueTagsV1(req *http.Request) (int, interfaces.AbstractResponseBody) {
|
||||
requestBody := models.NewListQueueTagsRequest()
|
||||
ok := utils.REQUEST_TRANSFORMER(requestBody, req, false)
|
||||
if !ok {
|
||||
log.Error("Invalid Request - ListQueueTagsV1")
|
||||
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)
|
||||
|
||||
models.SyncQueues.RLock()
|
||||
defer models.SyncQueues.RUnlock()
|
||||
|
||||
queue, exists := models.SyncQueues.Queues[key]
|
||||
if !exists {
|
||||
return utils.CreateErrorResponseV1("QueueNotFound", true)
|
||||
}
|
||||
|
||||
tags := queue.Tags
|
||||
if tags == nil {
|
||||
tags = make(map[string]string)
|
||||
}
|
||||
|
||||
log.Debugf("ListQueueTags: %s — %d tags", queueName, len(tags))
|
||||
respStruct := models.ListQueueTagsResponse{
|
||||
Xmlns: models.BaseXmlns,
|
||||
Result: models.ListQueueTagsResult{
|
||||
Tags: tags,
|
||||
},
|
||||
Metadata: models.BaseResponseMetadata,
|
||||
}
|
||||
return http.StatusOK, respStruct
|
||||
}
|
||||
@@ -0,0 +1,70 @@
|
||||
// Создано: 2026-04-11
|
||||
// TagQueueV1 — добавляет или обновляет метки (tags) очереди тенанта.
|
||||
// Совместимость с AWS SQS TagQueue: Tag.N.Key / Tag.N.Value.
|
||||
package gosqs
|
||||
|
||||
import (
|
||||
"net/http"
|
||||
"strings"
|
||||
|
||||
"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 TagQueueV1(req *http.Request) (int, interfaces.AbstractResponseBody) {
|
||||
requestBody := models.NewTagQueueRequest()
|
||||
ok := utils.REQUEST_TRANSFORMER(requestBody, req, false)
|
||||
if !ok {
|
||||
log.Error("Invalid Request - TagQueueV1")
|
||||
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)
|
||||
|
||||
models.SyncQueues.Lock()
|
||||
defer models.SyncQueues.Unlock()
|
||||
|
||||
queue, exists := models.SyncQueues.Queues[key]
|
||||
if !exists {
|
||||
return utils.CreateErrorResponseV1("QueueNotFound", true)
|
||||
}
|
||||
|
||||
// Инициализируем Tags если nil (для старых очередей, созданных без Tags)
|
||||
if queue.Tags == nil {
|
||||
queue.Tags = make(map[string]string)
|
||||
}
|
||||
|
||||
// Merge: новая метка с совпадающим ключом заменяет существующую
|
||||
for k, v := range requestBody.Tags {
|
||||
queue.Tags[k] = v
|
||||
}
|
||||
|
||||
persistence.SaveQueue(key, queue)
|
||||
|
||||
log.Debugf("TagQueue: %s — added %d tags", queueName, len(requestBody.Tags))
|
||||
respStruct := models.TagQueueResponse{
|
||||
Xmlns: models.BaseXmlns,
|
||||
Metadata: models.BaseResponseMetadata,
|
||||
}
|
||||
return http.StatusOK, respStruct
|
||||
}
|
||||
@@ -0,0 +1,66 @@
|
||||
// Создано: 2026-04-11
|
||||
// UntagQueueV1 — удаляет метки (tags) очереди тенанта по ключам.
|
||||
// Совместимость с AWS SQS UntagQueue: TagKey.N.
|
||||
package gosqs
|
||||
|
||||
import (
|
||||
"net/http"
|
||||
"strings"
|
||||
|
||||
"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 UntagQueueV1(req *http.Request) (int, interfaces.AbstractResponseBody) {
|
||||
requestBody := models.NewUntagQueueRequest()
|
||||
ok := utils.REQUEST_TRANSFORMER(requestBody, req, false)
|
||||
if !ok {
|
||||
log.Error("Invalid Request - UntagQueueV1")
|
||||
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)
|
||||
|
||||
models.SyncQueues.Lock()
|
||||
defer models.SyncQueues.Unlock()
|
||||
|
||||
queue, exists := models.SyncQueues.Queues[key]
|
||||
if !exists {
|
||||
return utils.CreateErrorResponseV1("QueueNotFound", true)
|
||||
}
|
||||
|
||||
if queue.Tags != nil {
|
||||
for _, tagKey := range requestBody.TagKeys {
|
||||
delete(queue.Tags, tagKey)
|
||||
}
|
||||
}
|
||||
|
||||
persistence.SaveQueue(key, queue)
|
||||
|
||||
log.Debugf("UntagQueue: %s — removed %d tag keys", queueName, len(requestBody.TagKeys))
|
||||
respStruct := models.UntagQueueResponse{
|
||||
Xmlns: models.BaseXmlns,
|
||||
Metadata: models.BaseResponseMetadata,
|
||||
}
|
||||
return http.StatusOK, respStruct
|
||||
}
|
||||
@@ -63,6 +63,7 @@ type Queue struct {
|
||||
FIFOSequenceNumbers map[string]int
|
||||
EnableDuplicates bool
|
||||
Duplicates map[string]time.Time
|
||||
Tags map[string]string
|
||||
}
|
||||
|
||||
func (q *Queue) NextSequenceNumber(groupId string) string {
|
||||
|
||||
@@ -530,3 +530,93 @@ func (r *DeleteMessageBatchRequest) SetAttributesFromForm(values url.Values) {
|
||||
r.Entries = entries
|
||||
}
|
||||
}
|
||||
|
||||
/*** ChangeMessageVisibilityBatch Request ***/
|
||||
type ChangeMessageVisibilityBatchRequestEntry struct {
|
||||
Id string `json:"Id" schema:"Id"`
|
||||
ReceiptHandle string `json:"ReceiptHandle" schema:"ReceiptHandle"`
|
||||
VisibilityTimeout int `json:"VisibilityTimeout" schema:"VisibilityTimeout"`
|
||||
}
|
||||
|
||||
type ChangeMessageVisibilityBatchRequest struct {
|
||||
Entries []ChangeMessageVisibilityBatchRequestEntry `json:"Entries"`
|
||||
QueueUrl string `json:"QueueUrl" schema:"QueueUrl"`
|
||||
}
|
||||
|
||||
func NewChangeMessageVisibilityBatchRequest() *ChangeMessageVisibilityBatchRequest {
|
||||
return &ChangeMessageVisibilityBatchRequest{}
|
||||
}
|
||||
|
||||
// SetAttributesFromForm — парсит AWS Query Protocol: ChangeMessageVisibilityBatchRequestEntry.N.*
|
||||
func (r *ChangeMessageVisibilityBatchRequest) SetAttributesFromForm(values url.Values) {
|
||||
for i := 1; ; i++ {
|
||||
id := values.Get(fmt.Sprintf("ChangeMessageVisibilityBatchRequestEntry.%d.Id", i))
|
||||
receiptHandle := values.Get(fmt.Sprintf("ChangeMessageVisibilityBatchRequestEntry.%d.ReceiptHandle", i))
|
||||
if id == "" || receiptHandle == "" {
|
||||
break
|
||||
}
|
||||
entry := ChangeMessageVisibilityBatchRequestEntry{
|
||||
Id: id,
|
||||
ReceiptHandle: receiptHandle,
|
||||
}
|
||||
vt := values.Get(fmt.Sprintf("ChangeMessageVisibilityBatchRequestEntry.%d.VisibilityTimeout", i))
|
||||
if vt != "" {
|
||||
entry.VisibilityTimeout, _ = strconv.Atoi(vt)
|
||||
}
|
||||
r.Entries = append(r.Entries, entry)
|
||||
}
|
||||
}
|
||||
|
||||
/*** TagQueue Request ***/
|
||||
type TagQueueRequest struct {
|
||||
QueueUrl string `json:"QueueUrl" schema:"QueueUrl"`
|
||||
Tags map[string]string `json:"Tags"`
|
||||
}
|
||||
|
||||
func NewTagQueueRequest() *TagQueueRequest {
|
||||
return &TagQueueRequest{Tags: make(map[string]string)}
|
||||
}
|
||||
|
||||
// SetAttributesFromForm — парсит AWS Query Protocol: Tag.N.Key / Tag.N.Value
|
||||
func (r *TagQueueRequest) SetAttributesFromForm(values url.Values) {
|
||||
for i := 1; ; i++ {
|
||||
key := values.Get(fmt.Sprintf("Tag.%d.Key", i))
|
||||
value := values.Get(fmt.Sprintf("Tag.%d.Value", i))
|
||||
if key == "" {
|
||||
break
|
||||
}
|
||||
r.Tags[key] = value
|
||||
}
|
||||
}
|
||||
|
||||
/*** UntagQueue Request ***/
|
||||
type UntagQueueRequest struct {
|
||||
QueueUrl string `json:"QueueUrl" schema:"QueueUrl"`
|
||||
TagKeys []string `json:"TagKeys"`
|
||||
}
|
||||
|
||||
func NewUntagQueueRequest() *UntagQueueRequest {
|
||||
return &UntagQueueRequest{}
|
||||
}
|
||||
|
||||
// SetAttributesFromForm — парсит AWS Query Protocol: TagKey.N
|
||||
func (r *UntagQueueRequest) SetAttributesFromForm(values url.Values) {
|
||||
for i := 1; ; i++ {
|
||||
key := values.Get(fmt.Sprintf("TagKey.%d", i))
|
||||
if key == "" {
|
||||
break
|
||||
}
|
||||
r.TagKeys = append(r.TagKeys, key)
|
||||
}
|
||||
}
|
||||
|
||||
/*** ListQueueTags Request ***/
|
||||
type ListQueueTagsRequest struct {
|
||||
QueueUrl string `json:"QueueUrl" schema:"QueueUrl"`
|
||||
}
|
||||
|
||||
func NewListQueueTagsRequest() *ListQueueTagsRequest {
|
||||
return &ListQueueTagsRequest{}
|
||||
}
|
||||
|
||||
func (r *ListQueueTagsRequest) SetAttributesFromForm(values url.Values) {}
|
||||
|
||||
@@ -338,3 +338,74 @@ func (r DeleteMessageBatchResponse) GetResult() interface{} {
|
||||
func (r DeleteMessageBatchResponse) GetRequestId() string {
|
||||
return r.Metadata.RequestId
|
||||
}
|
||||
|
||||
/*** ChangeMessageVisibilityBatch Response ***/
|
||||
type ChangeMessageVisibilityBatchResultEntry struct {
|
||||
Id string `json:"Id" xml:"Id"`
|
||||
}
|
||||
|
||||
type ChangeMessageVisibilityBatchResult struct {
|
||||
Successful []ChangeMessageVisibilityBatchResultEntry `json:"Successful" xml:"ChangeMessageVisibilityBatchResultEntry"`
|
||||
Failed []BatchResultErrorEntry `json:"Failed,omitempty" xml:"BatchResultErrorEntry,omitempty"`
|
||||
}
|
||||
|
||||
type ChangeMessageVisibilityBatchResponse struct {
|
||||
Xmlns string `json:"Xmlns" xml:"xmlns,attr"`
|
||||
Result ChangeMessageVisibilityBatchResult `json:"ChangeMessageVisibilityBatchResult" xml:"ChangeMessageVisibilityBatchResult"`
|
||||
Metadata ResponseMetadata `json:"ResponseMetadata" xml:"ResponseMetadata"`
|
||||
}
|
||||
|
||||
func (r ChangeMessageVisibilityBatchResponse) GetResult() interface{} {
|
||||
return r.Result
|
||||
}
|
||||
|
||||
func (r ChangeMessageVisibilityBatchResponse) GetRequestId() string {
|
||||
return r.Metadata.RequestId
|
||||
}
|
||||
|
||||
/*** TagQueue Response ***/
|
||||
type TagQueueResponse struct {
|
||||
Xmlns string `json:"Xmlns" xml:"xmlns,attr"`
|
||||
Metadata ResponseMetadata `json:"ResponseMetadata" xml:"ResponseMetadata"`
|
||||
}
|
||||
|
||||
func (r TagQueueResponse) GetResult() interface{} {
|
||||
return nil
|
||||
}
|
||||
|
||||
func (r TagQueueResponse) GetRequestId() string {
|
||||
return r.Metadata.RequestId
|
||||
}
|
||||
|
||||
/*** UntagQueue Response ***/
|
||||
type UntagQueueResponse struct {
|
||||
Xmlns string `json:"Xmlns" xml:"xmlns,attr"`
|
||||
Metadata ResponseMetadata `json:"ResponseMetadata" xml:"ResponseMetadata"`
|
||||
}
|
||||
|
||||
func (r UntagQueueResponse) GetResult() interface{} {
|
||||
return nil
|
||||
}
|
||||
|
||||
func (r UntagQueueResponse) GetRequestId() string {
|
||||
return r.Metadata.RequestId
|
||||
}
|
||||
|
||||
/*** ListQueueTags Response ***/
|
||||
type ListQueueTagsResult struct {
|
||||
Tags map[string]string `json:"Tags" xml:"Tag"`
|
||||
}
|
||||
|
||||
type ListQueueTagsResponse struct {
|
||||
Xmlns string `json:"Xmlns" xml:"xmlns,attr"`
|
||||
Result ListQueueTagsResult `json:"ListQueueTagsResult" xml:"ListQueueTagsResult"`
|
||||
Metadata ResponseMetadata `json:"ResponseMetadata" xml:"ResponseMetadata"`
|
||||
}
|
||||
|
||||
func (r ListQueueTagsResponse) GetResult() interface{} {
|
||||
return r.Result
|
||||
}
|
||||
|
||||
func (r ListQueueTagsResponse) GetRequestId() string {
|
||||
return r.Metadata.RequestId
|
||||
}
|
||||
|
||||
@@ -110,6 +110,10 @@ var routingTableV1 = map[string]func(r *http.Request) (int, interfaces.AbstractR
|
||||
"DeleteQueue": sqs.DeleteQueueV1,
|
||||
"SendMessageBatch": sqs.SendMessageBatchV1,
|
||||
"DeleteMessageBatch": sqs.DeleteMessageBatchV1,
|
||||
"ChangeMessageVisibilityBatch": sqs.ChangeMessageVisibilityBatchV1,
|
||||
"TagQueue": sqs.TagQueueV1,
|
||||
"UntagQueue": sqs.UntagQueueV1,
|
||||
"ListQueueTags": sqs.ListQueueTagsV1,
|
||||
}
|
||||
|
||||
func health(w http.ResponseWriter, req *http.Request) {
|
||||
|
||||
Reference in New Issue
Block a user