v0.1.35: resolveProtocol — query-mode клиенты парсят JSON-ответы (эмпирика), ошибки AWS-JSON __type
This commit is contained in:
@@ -23,7 +23,7 @@ import (
|
|||||||
// - если очередь уже есть — просто возвращаем её URL (идемпотентность AWS);
|
// - если очередь уже есть — просто возвращаем её URL (идемпотентность AWS);
|
||||||
// - проверка лимита очередей тенанта (MaxQueues) → LimitExceeded;
|
// - проверка лимита очередей тенанта (MaxQueues) → LimitExceeded;
|
||||||
// - создание очереди с дефолтами и применение переданных атрибутов
|
// - создание очереди с дефолтами и применение переданных атрибутов
|
||||||
// (setQueueAttributesV1 с валидацией диапазонов);
|
// (setQueueAttributesV1 с валидацией диапазонов);
|
||||||
// - вставка в map с tenant-scoped ключом "{tenantAccessKey}:{queueName}".
|
// - вставка в map с tenant-scoped ключом "{tenantAccessKey}:{queueName}".
|
||||||
// 5. Сохранение очереди в Redis (асинхронно, маршалинг под Lock).
|
// 5. Сохранение очереди в Redis (асинхронно, маршалинг под Lock).
|
||||||
// 6. Ответ: QueueUrl.
|
// 6. Ответ: QueueUrl.
|
||||||
|
|||||||
+2
-2
@@ -64,9 +64,9 @@ func PeriodicTasks(d time.Duration, quit chan bool) {
|
|||||||
// - пустой ReceiptHandle (не in-flight) — пропускаем;
|
// - пустой ReceiptHandle (не in-flight) — пропускаем;
|
||||||
// - VisibilityTimeout ещё не истёк — пропускаем;
|
// - VisibilityTimeout ещё не истёк — пропускаем;
|
||||||
// - иначе: возвращаем видимость (сброс ReceiptHandle, разблокировка
|
// - иначе: возвращаем видимость (сброс ReceiptHandle, разблокировка
|
||||||
// FIFO-группы, обновление ReceiptTime), инкрементируем Retry;
|
// FIFO-группы, обновление ReceiptTime), инкрементируем Retry;
|
||||||
// - если Retry >= MaxReceiveCount и настроена DLQ — переносим сообщение
|
// - если Retry >= MaxReceiveCount и настроена DLQ — переносим сообщение
|
||||||
// в DLQ и удаляем из исходной очереди.
|
// в DLQ и удаляем из исходной очереди.
|
||||||
func processQueueTick(queue *models.Queue) {
|
func processQueueTick(queue *models.Queue) {
|
||||||
// 1. Просроченные записи deduplication-окна.
|
// 1. Просроченные записи deduplication-окна.
|
||||||
for dedupId, startTime := range queue.Duplicates {
|
for dedupId, startTime := range queue.Duplicates {
|
||||||
|
|||||||
@@ -22,8 +22,8 @@ import (
|
|||||||
// 4. Под Lock:
|
// 4. Под Lock:
|
||||||
// - проверка существования очереди → QueueNotFound;
|
// - проверка существования очереди → QueueNotFound;
|
||||||
// - защита от повторного purge (фикс): повторный вызов в течение 60 секунд
|
// - защита от повторного purge (фикс): повторный вызов в течение 60 секунд
|
||||||
// → PurgeQueueInProgress (требование AWS). Раньше purge можно было
|
// → PurgeQueueInProgress (требование AWS). Раньше purge можно было
|
||||||
// вызывать без ограничений;
|
// вызывать без ограничений;
|
||||||
// - очистка in-memory сообщений и dedup-записей;
|
// - очистка in-memory сообщений и dedup-записей;
|
||||||
// - фиксация LastPurgeTime;
|
// - фиксация LastPurgeTime;
|
||||||
// - асинхронная очистка Redis (DEL хэша сообщений) + SaveQueue.
|
// - асинхронная очистка Redis (DEL хэша сообщений) + SaveQueue.
|
||||||
|
|||||||
@@ -26,10 +26,10 @@ import (
|
|||||||
// 2. Аутентифицированный тенант из контекста (SigV4 → AccessKey).
|
// 2. Аутентифицированный тенант из контекста (SigV4 → AccessKey).
|
||||||
// 3. Валидация входных параметров (AWS-совместимая):
|
// 3. Валидация входных параметров (AWS-совместимая):
|
||||||
// - MaxNumberOfMessages: 1–10; 0 (не задан) = 1. Выход за диапазон — ошибка
|
// - MaxNumberOfMessages: 1–10; 0 (не задан) = 1. Выход за диапазон — ошибка
|
||||||
// InvalidParameterValue (раньше значение молча клэмпилось).
|
// InvalidParameterValue (раньше значение молча клэмпилось).
|
||||||
// - VisibilityTimeout: 0–43200 секунд, иначе InvalidParameterValue.
|
// - VisibilityTimeout: 0–43200 секунд, иначе InvalidParameterValue.
|
||||||
// - WaitTimeSeconds: 0–20 секунд; 0 = брать атрибут очереди
|
// - WaitTimeSeconds: 0–20 секунд; 0 = брать атрибут очереди
|
||||||
// ReceiveMessageWaitTimeSeconds.
|
// ReceiveMessageWaitTimeSeconds.
|
||||||
// 4. Имя очереди: из QueueUrl (последний сегмент) или из пути /{account}/{queueName}.
|
// 4. Имя очереди: из QueueUrl (последний сегмент) или из пути /{account}/{queueName}.
|
||||||
// 5. Выборка сообщений — АТОМАРНО под одним Lock (см. receiveMessagesUnderLock):
|
// 5. Выборка сообщений — АТОМАРНО под одним Lock (см. receiveMessagesUnderLock):
|
||||||
// устранена гонка TOCTOU, когда два параллельных receive видели одни и те же
|
// устранена гонка TOCTOU, когда два параллельных receive видели одни и те же
|
||||||
@@ -226,6 +226,7 @@ func buildReceiveResponse(messages []*models.ResultMessage) models.ReceiveMessag
|
|||||||
// - ApproximateFirstReceiveTimestamp — время первой выдачи (ReceiptTime);
|
// - ApproximateFirstReceiveTimestamp — время первой выдачи (ReceiptTime);
|
||||||
// - ApproximateReceiveCount — число доставок (NumberOfReceives + 1);
|
// - ApproximateReceiveCount — число доставок (NumberOfReceives + 1);
|
||||||
// - SentTimestamp — реальное время отправки (SentTime), НЕ время выдачи.
|
// - SentTimestamp — реальное время отправки (SentTime), НЕ время выдачи.
|
||||||
|
//
|
||||||
// MD5 тела берётся из кэша (посчитан при SendMessage) — без пересчёта.
|
// MD5 тела берётся из кэша (посчитан при SendMessage) — без пересчёта.
|
||||||
func buildResultMessage(m *models.SqsMessage) *models.ResultMessage {
|
func buildResultMessage(m *models.SqsMessage) *models.ResultMessage {
|
||||||
return &models.ResultMessage{
|
return &models.ResultMessage{
|
||||||
|
|||||||
@@ -27,7 +27,7 @@ import (
|
|||||||
// - DeduplicationId и GroupId ≤ 128 символов;
|
// - DeduplicationId и GroupId ≤ 128 символов;
|
||||||
// - MessageAttributes ≤ 10;
|
// - MessageAttributes ≤ 10;
|
||||||
// - FIFO-очередь БЕЗ MessageGroupId → MissingParameter (AWS требует GroupId
|
// - FIFO-очередь БЕЗ MessageGroupId → MissingParameter (AWS требует GroupId
|
||||||
// для каждой отправки в FIFO);
|
// для каждой отправки в FIFO);
|
||||||
// - размер тела ≤ MaximumMessageSize очереди → MessageTooBig;
|
// - размер тела ≤ MaximumMessageSize очереди → MessageTooBig;
|
||||||
// - число сообщений < лимита очереди (OOM-защита) → OverLimit.
|
// - число сообщений < лимита очереди (OOM-защита) → OverLimit.
|
||||||
// 4. Построение SqsMessage: MD5 тела и атрибутов, UUID, метки времени, задержка.
|
// 4. Построение SqsMessage: MD5 тела и атрибутов, UUID, метки времени, задержка.
|
||||||
|
|||||||
+10
-12
@@ -262,19 +262,17 @@ func extractAction(req *http.Request) string {
|
|||||||
// resolveProtocol — определяет протокол ОТВЕТА по заголовкам запроса.
|
// resolveProtocol — определяет протокол ОТВЕТА по заголовкам запроса.
|
||||||
//
|
//
|
||||||
// Логическая схема:
|
// Логическая схема:
|
||||||
// 1. Заголовок x-amzn-query-mode: true (шлют AWS CLI v2 и boto3) означает
|
// 1. Content-Type: application/x-amz-json-1.0 → JSON-протокол (ответ JSON).
|
||||||
// «query-compatible JSON»: тело запроса сериализовано как JSON, но клиент
|
// Сюда входят: чистые JSON-клиенты (SDK) И AWS CLI v2/boto3 в режиме
|
||||||
// ЖДЁТ ответ в XML (query-протокол). Если на такой запрос ответить JSON —
|
// «query-compatible JSON» (x-amzn-query-mode: true). ВАЖНО (эмпирика):
|
||||||
// botocore не распарсит тело ошибки (ClientError с Code=None — именно это
|
// query-mode клиент парсит именно JSON-ответы — при XML-ответе список
|
||||||
// наблюдалось: «сырые 400/404» без AWS-XML). Поэтому для query-mode
|
// очередей распознаётся как пустой (None), при JSON — корректно.
|
||||||
// принудительно отвечаем XML.
|
// 2. Всё остальное (form-urlencoded от классических клиентов) — query/XML.
|
||||||
// 2. Content-Type: application/x-amz-json-1.0 БЕЗ query-mode — чистый
|
//
|
||||||
// JSON-протокол (SDK): и запрос, и ответ в JSON.
|
// Ошибки в JSON-протоколе оформляются как {"__type":"...","message":"..."}
|
||||||
// 3. Всё остальное (form-urlencoded от старых клиентов) — query/XML.
|
// (см. encodeResponse) — только такой формат botocore распознаёт как ошибку
|
||||||
|
// с кодом (раньше отдавался голый ErrorResult → Code=None).
|
||||||
func resolveProtocol(req *http.Request) AwsProtocol {
|
func resolveProtocol(req *http.Request) AwsProtocol {
|
||||||
if req.Header.Get("x-amzn-query-mode") == "true" {
|
|
||||||
return AwsQueryProtocol
|
|
||||||
}
|
|
||||||
if req.Header.Get("Content-Type") == "application/x-amz-json-1.0" {
|
if req.Header.Get("Content-Type") == "application/x-amz-json-1.0" {
|
||||||
return AwsJsonProtocol
|
return AwsJsonProtocol
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user