From fbecde72eb4b49ba7374a0dcedec8fcd7f35e8fb Mon Sep 17 00:00:00 2001 From: Naeel Date: Sat, 11 Apr 2026 07:51:17 +0300 Subject: [PATCH] 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 --- app/gosqs/change_message_visibility_batch.go | 123 +++ app/gosqs/list_queue_tags.go | 65 ++ app/gosqs/tag_queue.go | 70 ++ app/gosqs/untag_queue.go | 66 ++ app/models/models.go | 1 + app/models/requests.go | 90 +++ app/models/responses.go | 71 ++ app/router/router.go | 4 + doc/api/yandex-message-queue-api-reference.md | 708 ++++++++++++++++++ doc/thinking/2026-04-11.md | 52 ++ 10 files changed, 1250 insertions(+) create mode 100644 app/gosqs/change_message_visibility_batch.go create mode 100644 app/gosqs/list_queue_tags.go create mode 100644 app/gosqs/tag_queue.go create mode 100644 app/gosqs/untag_queue.go create mode 100644 doc/api/yandex-message-queue-api-reference.md create mode 100644 doc/thinking/2026-04-11.md diff --git a/app/gosqs/change_message_visibility_batch.go b/app/gosqs/change_message_visibility_batch.go new file mode 100644 index 0000000..d0e625a --- /dev/null +++ b/app/gosqs/change_message_visibility_batch.go @@ -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 +} diff --git a/app/gosqs/list_queue_tags.go b/app/gosqs/list_queue_tags.go new file mode 100644 index 0000000..75c6e60 --- /dev/null +++ b/app/gosqs/list_queue_tags.go @@ -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 +} diff --git a/app/gosqs/tag_queue.go b/app/gosqs/tag_queue.go new file mode 100644 index 0000000..a1af472 --- /dev/null +++ b/app/gosqs/tag_queue.go @@ -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 +} diff --git a/app/gosqs/untag_queue.go b/app/gosqs/untag_queue.go new file mode 100644 index 0000000..694e4cc --- /dev/null +++ b/app/gosqs/untag_queue.go @@ -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 +} diff --git a/app/models/models.go b/app/models/models.go index 637a4fd..366813b 100644 --- a/app/models/models.go +++ b/app/models/models.go @@ -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 { diff --git a/app/models/requests.go b/app/models/requests.go index 606bf21..3f58632 100644 --- a/app/models/requests.go +++ b/app/models/requests.go @@ -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) {} diff --git a/app/models/responses.go b/app/models/responses.go index a69e652..d6ee5b5 100644 --- a/app/models/responses.go +++ b/app/models/responses.go @@ -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 +} diff --git a/app/router/router.go b/app/router/router.go index 504ee8c..5983c6b 100644 --- a/app/router/router.go +++ b/app/router/router.go @@ -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) { diff --git a/doc/api/yandex-message-queue-api-reference.md b/doc/api/yandex-message-queue-api-reference.md new file mode 100644 index 0000000..09502f7 --- /dev/null +++ b/doc/api/yandex-message-queue-api-reference.md @@ -0,0 +1,708 @@ +# Yandex Message Queue API Reference + +**Дата документации:** 11 апреля 2026 + +## Обзор + +Yandex Message Queue предоставляет HTTP API, частично совместимый с Amazon SQS API. + +### Базовые параметры всех запросов + +**Адрес:** `POST https://message-queue.api.cloud.yandex.net/` + +**Заголовки:** +- `Content-Type: application/x-www-form-urlencoded` +- `Authorization: Authorization string (AWS Signature Version 4)` + +**Параметры запроса:** +- `Action` — название вызываемого метода API +- `Version` — всегда `2012-11-05` + +### Формат передачи массивов параметров + +Элементы массивов передаются с индексами начиная с 1: +``` +Attribute.1.Name=VisibilityTimeout +Attribute.1.Value=40 +Attribute.2.Name=MessageRetentionPeriod +Attribute.2.Value=1000 +``` + +### Формат ответов + +**Успешный ответ:** +```xml + + + + + + + UUID + + +``` + +**Ошибочный ответ:** +```xml + + + Sender|Receiver + ошибка + Описание + + UUID + +``` + +--- + +## API Команды + +### Управление очередями + +#### CreateQueue +Создание новой стандартной или FIFO очереди. + +**Параметры запроса:** +| Параметр | Тип | Обязательный | Описание | +|---|---|---|---| +| QueueName | string | Да | Имя очереди (макс 80 символов). Для FIFO должно оканчиваться на .fifo | +| Attributes.N.* | список | Нет | Атрибуты очереди | +| Tags.N.* | список | Нет | Метки очереди | + +**Атрибуты очереди:** +| Атрибут | Тип | Описание | +|---|---|---| +| DelaySeconds | integer | 0-900 сек. Default: 0 | +| MaximumMessageSize | integer | 1024-262144 байт. Default: 262144 | +| MessageRetentionPeriod | integer | 60-1209600 сек. Default: 345600 | +| ReceiveMessageWaitTimeSeconds | integer | 0-20 сек. Default: 0 | +| RedrivePolicy | string | JSON с deadLetterTargetArn и maxReceiveCount | +| VisibilityTimeout | integer | 0-43000 сек. Default: 30 | +| FifoQueue | boolean | true/false — создание FIFO очереди | +| ContentBasedDeduplication | boolean | true/false — дедупликация (FIFO) | + +**Выходные параметры:** +| Параметр | Тип | Описание | +|---|---|---| +| QueueUrl | string | URL созданной очереди | + +**Ошибки:** +- 400 `QueueDeletedRecently` — очередь удалена недавно, ждать 60 сек +- 400 `QueueAlreadyExists` — очередь уже существует + +--- + +#### DeleteQueue +Удаление очереди. Процесс занимает до 60 секунд. + +**Параметры запроса:** +| Параметр | Тип | Обязательный | Описание | +|---|---|---|---| +| QueueUrl | string | Да | URL очереди (чувствителен к регистру) | + +**Выходные параметры:** +- Нет полей + +**Ошибки:** +- Только стандартные + +--- + +#### GetQueueAttributes +Получение атрибутов очереди. + +**Параметры запроса:** +| Параметр | Тип | Обязательный | Описание | +|---|---|---|---| +| QueueUrl | string | Да | URL очереди | +| AttributeNames.N | array | Нет | Список запрашиваемых атрибутов | + +**Значения AttributeNames:** +- `All` — все атрибуты +- `ApproximateNumberOfMessages` — ориентировочное количество готовых сообщений +- `ApproximateNumberOfMessagesDelayed` — отложенные сообщения +- `ApproximateNumberOfMessagesNotVisible` — сообщения в процессе передачи +- `CreatedTimestamp` — время создания (epoch time) +- `DelaySeconds` +- `LastModifiedTimestamp` — время последнего изменения (epoch time) +- `MaximumMessageSize` +- `MessageRetentionPeriod` +- `QueueArn` — ARN очереди +- `ReceiveMessageWaitTimeSeconds` +- `RedrivePolicy` — политика DLQ +- `VisibilityTimeout` +- `FifoQueue` — флаг FIFO +- `ContentBasedDeduplication` — флаг дедупликации + +**Выходные параметры:** +| Параметр | Тип | Описание | +|---|---|---| +| Attributes.N.* | array | Массив атрибутов (Name, Value) | + +**Ошибки:** +- 400 `InvalidAttributeName` — неверное имя атрибута + +--- + +#### GetQueueUrl +Получение URL очереди по имени. + +**Параметры запроса:** +| Параметр | Тип | Обязательный | Описание | +|---|---|---|---| +| QueueName | string | Да | Имя очереди (макс 80 символов, чувствителен к регистру) | +| QueueOwnerAWSAccountId | string | Нет | Параметр игнорируется | + +**Выходные параметры:** +| Параметр | Тип | Описание | +|---|---|---| +| QueueUrl | string | URL очереди | + +**Ошибки:** +- 400 `NonExistentQueue` — очередь не существует + +--- + +#### ListQueues +Получение списка очередей в каталоге (макс 1000). + +**Параметры запроса:** +| Параметр | Тип | Обязательный | Описание | +|---|---|---|---| +| QueueNamePrefix | string | Нет | Префикс для фильтрации (чувствителен к регистру) | + +**Выходные параметры:** +| Параметр | Тип | Описание | +|---|---|---| +| QueueUrl.N | array | Массив URL очередей (до 1000) | + +**Ошибки:** +- Только стандартные + +--- + +#### PurgeQueue +Очистка очереди (удаление всех сообщений). + +**Параметры запроса:** +| Параметр | Тип | Обязательный | Описание | +|---|---|---|---| +| QueueUrl | string | Да | URL очереди (чувствителен к регистру) | + +**Выходные параметры:** +- Нет полей + +**Ошибки:** +- 400 `NonExistentQueue` — очередь не существует +- 403 `PurgeQueueInProgress` — PurgeQueue уже вызывали за последние 60 сек + +--- + +#### SetQueueAttributes +Изменение атрибутов очереди (может занять до 60 сек). + +**Параметры запроса:** +| Параметр | Тип | Обязательный | Описание | +|---|---|---|---| +| QueueUrl | string | Да | URL очереди | +| Attributes.N.* | список | Да | Список атрибутов | + +**Поддерживаемые атрибуты:** +- DelaySeconds +- MaximumMessageSize +- MessageRetentionPeriod +- ReceiveMessageWaitTimeSeconds +- RedrivePolicy +- VisibilityTimeout +- ContentBasedDeduplication (только FIFO) + +**Выходные параметры:** +- Нет полей + +**Ошибки:** +- 400 `InvalidAttributeName` — неверное имя атрибута + +--- + +#### TagQueue +Добавление/изменение меток очереди (может занять до 60 сек). + +**Параметры запроса:** +| Параметр | Тип | Обязательный | Описание | +|---|---|---|---| +| QueueUrl | string | Да | URL очереди | +| Tags.N.* | список | Да | Список меток (Tag.N.Key, Tag.N.Value) | + +**Выходные параметры:** +- Нет полей + +**Ошибки:** +- Только стандартные + +--- + +#### UntagQueue +Удаление меток очереди (может занять до 60 сек). + +**Параметры запроса:** +| Параметр | Тип | Обязательный | Описание | +|---|---|---|---| +| QueueUrl | string | Да | URL очереди | +| TagKeys.N | array | Да | Список ключей для удаления (TagKey.N) | + +**Выходные параметры:** +- Нет полей + +**Ошибки:** +- Только стандартные + +--- + +### Управление сообщениями + +#### SendMessage +Отправка одного сообщения в очередь. + +**Параметры запроса:** +| Параметр | Тип | Обязательный | Описание | +|---|---|---|---| +| QueueUrl | string | Да | URL очереди | +| MessageBody | string | Да | Тело сообщения (макс 256 КБ). XML/JSON/text | +| DelaySeconds | integer | Нет | 0-900 сек отложки | +| MessageAttributeName.N / MessageAttributeValue.N | array | Нет | Пользовательские атрибуты | +| MessageDeduplicationId | string | Да (FIFO) | Макс 128 символов | +| MessageGroupId | string | Да (FIFO) | Макс 128 символов | + +**Выходные параметры:** +| Параметр | Тип | Описание | +|---|---|---| +| MD5OfMessageBody | string | MD5 хэш тела | +| MD5OfMessageAttributes | string | MD5 хэш атрибутов | +| MessageId | string | Id отправленного сообщения | +| SequenceNumber | string | Номер в FIFO (только FIFO) | + +**Ошибки:** +- 400 `UnsupportedOperation` — неподдерживаемая операция +- 400 `InvalidMessageContents` — запрещённые символы + +--- + +#### SendMessageBatch +Отправка до 10 сообщений одновременно (макс 256 КБ общий размер). + +**Параметры запроса:** +| Параметр | Тип | Обязательный | Описание | +|---|---|---|---| +| QueueUrl | string | Да | URL очереди | +| SendMessageBatchRequestEntry.N | array | Да | Массив до 10 сообщений | + +**Выходные параметры:** +| Параметр | Тип | Описание | +|---|---|---| +| BatchResultErrorEntry.N | array | Ошибки отправки | +| SendMessageBatchResultEntry.N | array | Id, MD5, MessageId успешных | + +**Ошибки:** +- 400 `BatchEntryIdsNotDistinct` — одинаковые Id +- 400 `BatchRequestTooLong` — общая длина превышена +- 400 `EmptyBatchRequest` — нет сообщений +- 400 `InvalidBatchEntryId` — неправильный Id +- 400 `TooManyEntriesInBatchRequest` — >10 сообщений +- 400 `UnsupportedOperation` — неподдерживаемая операция + +--- + +#### ReceiveMessage +Получение 1-10 сообщений из очереди. + +**Параметры запроса:** +| Параметр | Тип | Обязательный | Описание | +|---|---|---|---| +| QueueUrl | string | Да | URL очереди | +| MaxNumberOfMessages | string | Нет | 1-10 сообщений. Default: 1 | +| MessageAttributeName.N | array | Нет | Имена атрибутов (макс 256 символов) | +| ReceiveRequestAttemptId | string | Нет | Id для повтора попытки (FIFO) | +| VisibilityTimeout | string | Нет | Таймаут видимости | +| WaitTimeSeconds | string | Нет | Long-polling ожидание (0-20 сек) | + +**Атрибуты сообщения:** +- `All` — все атрибуты +- `ApproximateFirstReceiveTimestamp` — время первого получения +- `ApproximateReceiveCount` — количество получений без удаления +- `SenderId` — Id отправителя (IAM) +- `SentTimestamp` — время отправки +- `MessageDeduplicationId` — Id дедупликации (FIFO) +- `MessageGroupId` — Id группы (FIFO) +- `SequenceNumber` — номер в группе (FIFO) + +**Выходные параметры:** +| Параметр | Тип | Описание | +|---|---|---| +| Message | array | Массив сообщений | + +**Ошибки:** +- 403 `OverLimit` — превышен один из установленных лимитов + +--- + +#### DeleteMessage +Удаление одного сообщения из очереди. + +**Параметры запроса:** +| Параметр | Тип | Обязательный | Описание | +|---|---|---|---| +| QueueUrl | string | Да | URL очереди | +| ReceiptHandle | string | Да | Id получения из ReceiveMessage | + +**Выходные параметры:** +- Нет полей + +**Ошибки:** +- 400 `InvalidIdFormat` — некорректный формат ReceiptHandle +- 400 `ReceiptHandleIsInvalid` — неверный ReceiptHandle + +--- + +#### DeleteMessageBatch +Удаление до 10 сообщений одновременно. + +**Параметры запроса:** +| Параметр | Тип | Обязательный | Описание | +|---|---|---|---| +| QueueUrl | string | Да | URL очереди | +| DeleteMessageBatchRequestEntry.N | array | Да | Массив до 10 сообщений (Id, ReceiptHandle) | + +**Выходные параметры:** +| Параметр | Тип | Описание | +|---|---|---| +| BatchResultErrorEntry.N | array | Ошибки удаления | +| DeleteMessageBatchResultEntry.N | array | Id успешно удаленных | + +**Ошибки:** +- 400 `BatchEntryIdsNotDistinct` — одинаковые Id +- 400 `EmptyBatchRequest` — нет сообщений +- 400 `InvalidBatchEntryId` — неправильный Id +- 400 `TooManyEntriesInBatchRequest` — >10 сообщений + +--- + +#### ChangeMessageVisibility +Изменение таймаута видимости сообщения в обработке. + +**Параметры запроса:** +| Параметр | Тип | Обязательный | Описание | +|---|---|---|---| +| QueueUrl | string | Да | URL очереди | +| ReceiptHandle | string | Да | Id получения | +| VisibilityTimeout | integer | Да | 0-43200 сек новый таймаут | + +**Выходные параметры:** +- Нет полей + +**Ошибки:** +- 400 `MessageNotInflight` — сообщение не в обработке +- 400 `ReceiptHandleIsInvalid` — неверный ReceiptHandle + +**Примечание:** Суммарная длительность таймаута не может быть более 12 часов. + +--- + +#### ChangeMessageVisibilityBatch +Изменение таймаута видимости до 10 сообщений одновременно. + +**Параметры запроса:** +| Параметр | Тип | Обязательный | Описание | +|---|---|---|---| +| QueueUrl | string | Да | URL очереди (чувствителен к регистру) | +| ChangeMessageVisibilityBatchRequestEntry.N | array | Да | Массив до 10 сообщений | + +**Выходные параметры:** +| Параметр | Тип | Описание | +|---|---|---| +| BatchResultErrorEntry.N | array | Ошибки | +| ChangeMessageVisibilityBatchResultEntry.N | array | Id успешно измененных | + +**Ошибки:** +- 400 `BatchEntryIdsNotDistinct` — одинаковые Id +- 400 `EmptyBatchRequest` — нет сообщений +- 400 `InvalidBatchEntryId` — неправильный Id +- 400 `TooManyEntriesInBatchRequest` — >10 сообщений + +--- + +## Типы данных + +### BatchResultErrorEntry +Описание ошибки выполнения действия из группы. + +| Параметр | Тип | Обязательный | Описание | +|---|---|---|---| +| Code | string | Да | Код ошибки | +| Id | string | Да | Id сообщения в группе | +| Message | string | Нет | Описание ошибки | +| SenderFault | boolean | Да | Ошибка на стороне отправителя | + +--- + +### ChangeMessageVisibilityBatchRequestEntry +Элемент массива для ChangeMessageVisibilityBatch. + +| Параметр | Тип | Обязательный | Описание | +|---|---|---|---| +| Id | string | Да | Id ReceiptHandle, уникален в пределе запроса | +| ReceiptHandle | string | Да | Id получения сообщения | +| VisibilityTimeout | boolean | Нет | Новый таймаут в секундах | + +--- + +### ChangeMessageVisibilityBatchResultEntry +Результат для одного сообщения в ChangeMessageVisibilityBatch. + +| Параметр | Тип | Описание | +|---|---|---| +| Id | string | Id сообщения с измененным таймаутом | + +--- + +### DeleteMessageBatchRequestEntry +Элемент массива для DeleteMessageBatch. + +| Параметр | Тип | Обязательный | Описание | +|---|---|---|---| +| Id | string | Да | Id, уникален в пределе запроса | +| ReceiptHandle | string | Да | Id получения сообщения | + +--- + +### DeleteMessageBatchResultEntry +Результат для одного сообщения в DeleteMessageBatch. + +| Параметр | Тип | Описание | +|---|---|---| +| Id | string | Id удаленного сообщения | + +--- + +### Message +Сообщение из ReceiveMessage. + +| Параметр | Тип | Описание | +|---|---|---| +| Attribute.N | array | Системные атрибуты (ApproximateReceiveCount, ApproximateFirstReceiveTimestamp, MessageDeduplicationId, MessageGroupId, SenderId, SentTimestamp, SequenceNumber) | +| Body | string | Тело сообщения | +| MD5OfBody | string | MD5 хэш тела | +| MD5OfMessageAttributes | string | MD5 хэш атрибутов | +| MessageAttribute | array | Пользовательские атрибуты (MessageAttributeValue) | +| MessageId | string | Уникальный Id | +| ReceiptHandle | string | Id получения (новый при каждом получении) | + +--- + +### MessageAttributeValue +Значение пользовательского атрибута сообщения. + +| Параметр | Тип | Обязательный | Описание | +|---|---|---|---| +| BinaryListValue.N | array | Нет | Не реализовано | +| BinaryValue | base64 | Нет | Двоичные данные | +| DataType | string | Да | String, Number или Binary (для Number используется StringValue) | +| StringListValue.N | string | Нет | Не реализовано | +| StringValue | string | Нет | Строка UTF-8 | + +**Примечание:** Имя, тип, значение и тело не могут быть пустыми. Суммарный размер всех частей ≤ 256 КБ. + +--- + +### SendMessageBatchRequestEntry +Элемент массива для SendMessageBatch. + +| Параметр | Тип | Обязательный | Описание | +|---|---|---|---| +| DelaySeconds | integer | Нет | Отложка в секундах | +| Id | string | Да | Id сообщения в списке | +| MessageAttribute | string | Нет | Атрибуты (имя, тип, значение) | +| MessageBody | string | Нет | Тело сообщения | +| MessageDeduplicationId | string | Нет | Id дедупликации | +| MessageGroupId | string | Нет | Id группы (только FIFO) | + +--- + +### SendMessageBatchResultEntry +Результат для одного сообщения в SendMessageBatch. + +| Параметр | Тип | Обязательный | Описание | +|---|---|---|---| +| Id | string | Да | Id сообщения в группе | +| MD5OfMessageAttributes | string | Нет | MD5 хэш атрибутов | +| MD5OfMessageBody | string | Да | MD5 хэш тела | +| MessageId | string | Да | Id сообщения | +| SequenceNumber | string | Нет | Номер (128 бит, только FIFO) | + +--- + +## Стандартные ошибки + +Ошибки, возвращаемые всеми методами: + +| HTTP | Код ошибки | Описание | +|---|---|---| +| 400 | AccessDeniedException | Недостаточно прав для выполнения действия | +| 400 | IncompleteSignature | Подпись запроса не соответствует стандартам AWS | +| 500 | InternalFailure | Неизвестная ошибка | +| 400 | InvalidAction | Неизвестное значение параметра Action | +| 403 | InvalidClientTokenId | Неверный ключ сервисного аккаунта | +| 400 | InvalidParameterCombination | Одновременно используются несовместимые параметры | +| 400 | InvalidParameterValue | Параметр задан неверно или вне диапазона | +| 400 | InvalidQueryParameter | Используется несуществующий параметр | +| 404 | MalformedQueryString | Синтаксическая ошибка в запросе | +| 400 | MissingAction | Не указан параметр Action | +| 403 | MissingAuthenticationToken | Не указан Id ключа сервисного аккаунта | +| 400 | MissingParameter | Отсутствует обязательный параметр | +| 400 | OptInRequired | Id ключа требует подписки на сервис | +| 400 | RequestExpired | Запрос получен через >15 минут от времени в запросе | +| 503 | ServiceUnavailable | Сервис недоступен | +| 403 | ThrottlingException | Ограничение по числу запросов | +| 400 | ValidationError | Значения не соответствуют ограничениям | + +--- + +## Примечания + +1. **FIFO очереди:** + - Имя должно оканчиваться на `.fifo` + - Требуют `MessageGroupId` для SendMessage + - Требуют `MessageDeduplicationId` для дедупликации + - При ReceiveMessage из одной группы за вызов получается только одно сообщение + - Сообщения обрабатываются в порядке отправления + +2. **Таймауты:** + - VisibilityTimeout: 0-43000 сек, при смене максимум до 12 часов суммарно + - ReceiveMessage с WaitTimeSeconds — long-polling (0-20 сек) + +3. **Batch операции:** + - Максимум 10 элементов в batch + - SendMessageBatch: макс 256 КБ общий размер + - Результаты проверяются индивидуально (могут быть успехи и ошибки) + +4. **Удаление очередей:** + - Процесс может занять до 60 секунд + - Новую очередь с таким же именем можно создать через 60 сек после удаления + +5. **ARN очереди:** + - Используется в RedrivePolicy + - Может быть получен через GetQueueAttributes с параметром QueueArn + +--- + +## Audit Trails (Аудитные логи) + +В Audit Trails для Yandex Message Queue поддерживается отслеживание событий уровня конфигурации (Control Plane). + +### Формат event_type + +``` +yandex.cloud.audit.ymq.<имя_события> +``` + +### События уровня конфигурации + +| Имя события | Описание | +|---|---| +| CreateMessageQueue | Создание очереди сообщений | +| DeleteMessageQueue | Удаление очереди сообщений | +| UpdateMessageQueue | Изменение очереди сообщений | + +--- + +## Monitoring (Метрики) + +В Yandex Monitoring поддерживаются метрики сервиса Message Queue. + +**Общая метка для всех метрик:** `service=message-queue` + +### Метрики HTTP API + +| Имя метрики | Тип, единицы измерения | Описание | Метка | +|---|---|---|---| +| api.http.errors_count_per_second | DGAUGE, ошибки/с | Количество ошибок выполнения запросов в секунду | method — метод API | +| api.http.request_duration_milliseconds | DGAUGE, миллисекунды | Продолжительность выполнения запросов | method — метод API | +| api.http.requests_count_per_second | DGAUGE, запросы/с | Количество обработанных запросов в секунду | method — метод API | + +### Метрики сервиса + +| Имя метрики | Тип, единицы измерения | Описание | +|---|---|---| +| queue.messages.client_processing_duration_milliseconds | DGAUGE, миллисекунды | Время обработки сообщений получателем | +| queue.messages.deduplicated_count_per_second | DGAUGE, сообщения/с | Частота дедупликации сообщений | +| queue.messages.deleted_count_per_second | DGAUGE, сообщения/с | Частота удаления сообщений из очереди | +| queue.messages.empty_receive_attempts_count_per_second | DGAUGE, попытки/с | Кол-во попыток получения пустого сообщения в сек | +| queue.messages.inflight_count | DGAUGE, штуки | Кол-во сообщений в обработке (активных) | +| queue.messages.oldest_age_milliseconds | DGAUGE, секунды | Время хранения наиболее раннего сообщения в очереди | +| queue.messages.purged_count_per_second | DGAUGE, сообщения/с | Частота удаления сообщений методом PurgeQueue | +| queue.messages.receive_attempts_count_rate | DGAUGE, штуки | Количество попыток получения сообщений из очереди | +| queue.messages.received_bytes_per_second | DGAUGE, байты/с | Общий размер полученных сообщений в сек | +| queue.messages.received_count_per_second | DGAUGE, сообщения/с | Количество полученных сообщений в сек | +| queue.messages.request_timeouts_count_per_second | DGAUGE, ошибки/с | Кол-во ошибок выполнения запросов ReceiveMessage | +| queue.messages.reside_duration_milliseconds | DGAUGE, миллисекунды | Время обработки сообщений в очереди | +| queue.messages.sent_bytes_per_second | DGAUGE, байты/с | Общий размер отправленных сообщений в сек | +| queue.messages.sent_count_per_second | DGAUGE, сообщения/с | Количество отправленных сообщений в сек | +| queue.messages.stored_count | DGAUGE, штуки | Количество сообщений в очереди в текущий момент | + +--- + +## Часто задаваемые вопросы + +### Ограничения и лимиты + +**Q: Что означает ошибка «Cannot create queue: Too many queues»?** + +A: Достигнут лимит на максимальное количество очередей. Для увеличения лимита обратитесь в техническую поддержку с указанием: +- Идентификателя облака +- Нужного количества очередей +- Назначения (зачем требуется такое количество) + +**Q: Какой максимальный размер сообщения?** + +A: Максимальный размер сообщения — 256 КБ. О других ограничениях см. в разделе Квоты и лимиты. + +### Доступ и аутентификация + +**Q: Мне нужна регистрация в Amazon для использования AWS CLI с Message Queue?** + +A: Нет, AWS CLI можно использовать без регистрации и ключей AWS. Подробнее см. раздел Инструменты. + +**Q: Какой ключ доступа требуется для работы с Message Queue?** + +A: Требуется статический ключ доступа. Создайте его в консоли облака. + +**Q: Что вводить в поле Default output format при настройке AWS CLI?** + +A: Оставьте это поле пустым. Укажите только идентификатор ключа и секретный ключ. + +### Операция и обслуживание + +**Q: Какой SLA для Message Queue?** + +A: Message Queue имеет SLA 99,90% на доступность сервиса. + +**Q: Могу ли я мониторить Message Queue через Prometheus?** + +A: Да, можно экспортировать метрики в Prometheus и получить список метрик для различных объектов. + +**Q: Почему все мои сообщения висят в очереди со статусом «В обработке» продолжительное время?** + +A: Это может быть связано с большим таймаутом видимости в настройках очереди. VisibilityTimeout — это время, на которое сообщение скрывается из очереди после чтения. + +### Обработка сообщений + +**Q: Как удалить сообщение и его дубль из очереди?** + +A: Дубликаты сообщений не записываются в течение 5 минут. Если прошло 5 минут, дубликат удаляется по его собственному ReceiptHandle. Дубли при чтении указывают на то, что сервис не успел удалить сообщение по истечении taймаута видимости — продлите таймаут. + +### Диагностикаи логи + +**Q: Могу ли я получить логи своей работы в сервисах?** + +A: Да, обратитесь в техническую поддержку для получения информации о работе с вашими ресурсами из логов сервисов Yandex Cloud. diff --git a/doc/thinking/2026-04-11.md b/doc/thinking/2026-04-11.md new file mode 100644 index 0000000..70a0dea --- /dev/null +++ b/doc/thinking/2026-04-11.md @@ -0,0 +1,52 @@ +# Thinking Log — 2026-04-11 +# Agent: GitHub Copilot (Claude Opus 4.6) + +--- + +## Задача: Добавить 4 недостающие API команды для совместимости с Yandex/AWS SQS + +### Контекст +Пользователь скопировал всю документацию Yandex Message Queue API (16 команд). +Сравнение показало, что у нас реализовано 13 из 16. Не хватает: +- ChangeMessageVisibilityBatch +- TagQueue +- UntagQueue +- ListQueueTags + +### Анализ +1. **ChangeMessageVisibilityBatch** — паттерн полностью аналогичен DeleteMessageBatch: + - Валидация: пустой batch, >10 entries, дублирование Id + - Partial success: отдельно Successful и Failed массивы + - Логика: цикл по ChangeMessageVisibility для каждого Entry + +2. **Tag-операции** — требуют добавления `Tags map[string]string` в Queue struct: + - TagQueue: merge tags (новый ключ перезаписывает старый) + - UntagQueue: delete по списку ключей + - ListQueueTags: read-only, RLock достаточно + - Persistence: SaveQueue после изменения tags (TagQueue, UntagQueue) + +3. **Юридический вопрос** — пользователь спросил про авторские права. + Ответ: API интерфейс не защищён (Oracle v. Google 2021). Yandex сам реализует AWS SQS API. + Десятки компаний делают то же самое (ElasticMQ, LocalStack, MinIO). + +### Решения +- Добавил поле `Tags map[string]string` в Queue struct — минимально инвазивное изменение +- Для старых очередей (без Tags) — nil-safe: проверка `if queue.Tags == nil` перед операциями +- ListQueueTags использует RLock (не Lock) — read-only операция +- ChangeMessageVisibilityBatch НЕ персистит в Redis (аналогично одиночному ChangeMessageVisibility) +- TagQueue/UntagQueue персистят через SaveQueue (tags — часть конфигурации очереди) + +### Что создано +- `app/gosqs/change_message_visibility_batch.go` — ~120 строк +- `app/gosqs/tag_queue.go` — ~70 строк +- `app/gosqs/untag_queue.go` — ~65 строк +- `app/gosqs/list_queue_tags.go` — ~65 строк +- Модели request/response в `app/models/requests.go` и `app/models/responses.go` +- Routing в `app/router/router.go` +- Поле `Tags` в `app/models/models.go` Queue struct + +### Итого API команд: 17 +Полный список: CreateQueue, DeleteQueue, GetQueueAttributes, GetQueueUrl, ListQueues, +PurgeQueue, SetQueueAttributes, SendMessage, SendMessageBatch, ReceiveMessage, +DeleteMessage, DeleteMessageBatch, ChangeMessageVisibility, ChangeMessageVisibilityBatch, +TagQueue, UntagQueue, ListQueueTags