diff --git a/HISTORY/2026-08-14-session-log.md b/HISTORY/2026-08-14-session-log.md index 05e2daf..c3bde8b 100644 --- a/HISTORY/2026-08-14-session-log.md +++ b/HISTORY/2026-08-14-session-log.md @@ -791,3 +791,17 @@ SIGKILL/рестарте пода запись не успевает → оче `LoadAllQueues` восстанавливает её. Удаления должны быть синхронными (хотя бы DEL очереди) — НЕ ИСПРАВЛЕНО, задокументировано. +## P2.6 FIFO/DLQ e2e — результаты (15.08.2026) + +`tests/fifo_dlq_probe.py`: +- **FIFO порядок**: PASS — 10 сообщений группы выданы строго m0..m9. +- **FIFO dedup** (дубль с тем же MessageDeduplicationId внутри окна): PASS — + в очереди одно сообщение («first»). +- **DLQ (RedrivePolicy)**: **FAIL — найден баг**: `setQueueAttributesV1` ищет DLQ + по ГОЛОМУ имени `models.SyncQueues.Queues[queueName]`, а ключи в map + tenant-scoped "{accessKey}:{queueName}" → DLQ не находится → CreateQueue + с RedrivePolicy всегда InvalidAttributeValue. RedrivePolicy не работает. + НЕ ИСПРАВЛЕНО, задокументировано. + +**Истечение dedup-окна (5 мин) не проверено** — тест быстрый, только «внутри окна». + diff --git a/tests/fifo_dlq_probe.py b/tests/fifo_dlq_probe.py new file mode 100644 index 0000000..fdd7359 --- /dev/null +++ b/tests/fifo_dlq_probe.py @@ -0,0 +1,122 @@ +#!/usr/bin/env python3 +"""tests/fifo_dlq_probe.py — e2e: FIFO порядок + dedup внутри окна + DLQ (план Соннета P2). + +Проверки: + 1. FIFO: 10 сообщений в одну группу выдаются строго по порядку. + 2. FIFO dedup: два send с одинаковым MessageDeduplicationId внутри окна + (5 мин) → в очереди оказывается ОДНО сообщение. + 3. DLQ: сообщение после MaxReceiveCount=2 попыток (VisibilityTimeout=1с) + переносится в dead-letter очередь. +""" +import json +import os +import sys +import time +import uuid + +import boto3 +from botocore.config import Config + +ENDPOINT = os.environ.get("ENDPOINT_URL", "https://sqs.containerk8s.dev.nubes.ru") +REGION = os.environ.get("REGION", "us-east-1") + +sqs = boto3.client("sqs", endpoint_url=ENDPOINT, region_name=REGION, + config=Config(connect_timeout=10, read_timeout=30, retries={"max_attempts": 3})) + +uid = str(uuid.uuid4())[:8] + + +def drain(url, n=10): + msgs = [] + while True: + b = sqs.receive_message(QueueUrl=url, MaxNumberOfMessages=n).get("Messages", []) + if not b: + break + msgs.extend(b) + return msgs + + +results = {} + + +def check_fifo_order(): + q = "fifo-order-%s.fifo" % uid + u = sqs.create_queue(QueueName=q, Attributes={"FifoQueue": "true"})["QueueUrl"] + for i in range(10): + sqs.send_message(QueueUrl=u, MessageBody="m%d" % i, MessageGroupId="g", MessageDeduplicationId="d%d" % (i,)) + got = [] + while True: + b = sqs.receive_message(QueueUrl=u, MaxNumberOfMessages=10).get("Messages", []) + if not b: + break + got.extend(b) + for m in b: + sqs.delete_message(QueueUrl=u, ReceiptHandle=m["ReceiptHandle"]) + bodies = [m["Body"] for m in got] + want = ["m%d" % i for i in range(10)] + results["fifo_order"] = bodies == want + print("fifo_order: got=%s want=%s %s" % (bodies, want, "PASS" if bodies == want else "FAIL"), flush=True) + sqs.delete_queue(QueueUrl=u) + + +def check_fifo_dedup(): + q = "fifo-dedup-%s.fifo" % uid + u = sqs.create_queue(QueueName=q, Attributes={"FifoQueue": "true"})["QueueUrl"] + sqs.send_message(QueueUrl=u, MessageBody="first", MessageGroupId="g", MessageDeduplicationId="SAME") + sqs.send_message(QueueUrl=u, MessageBody="dup", MessageGroupId="g", MessageDeduplicationId="SAME") + msgs = drain(u) + ok = len(msgs) == 1 and msgs[0]["Body"] == "first" + results["fifo_dedup"] = ok + print("fifo_dedup: count=%d body=%s %s" % (len(msgs), msgs[0]["Body"] if msgs else None, + "PASS" if ok else "FAIL"), flush=True) + sqs.delete_queue(QueueUrl=u) + + +def check_dlq(): + dlq_name = "fifo-dlq-%s" % uid + main_name = "fifo-main-%s" % uid + dlq_url = sqs.create_queue(QueueName=dlq_name)["QueueUrl"] + # tenantID из URL DLQ: http://us-east-1.goaws.com:4100/t-/ + tenant = dlq_url.split("/")[-2] + arn = "arn:aws:sqs:%s:%s:%s" % (REGION, tenant, dlq_name) + policy = json.dumps({"deadLetterTargetArn": arn, "maxReceiveCount": "2"}) + main_url = sqs.create_queue( + QueueName=main_name, + Attributes={"VisibilityTimeout": "1", "RedrivePolicy": policy}, + )["QueueUrl"] + sqs.send_message(QueueUrl=main_url, MessageBody="to-dlq") + for attempt in range(3): + b = sqs.receive_message(QueueUrl=main_url, MaxNumberOfMessages=1).get("Messages", []) + print("attempt %d: %s" % (attempt + 1, "received" if b else "empty"), flush=True) + time.sleep(1.5) + dlq_msgs = drain(dlq_url) + ok = len(dlq_msgs) == 1 and dlq_msgs[0]["Body"] == "to-dlq" + results["dlq"] = ok + print("dlq: count=%d body=%s %s" % (len(dlq_msgs), dlq_msgs[0]["Body"] if dlq_msgs else None, + "PASS" if ok else "FAIL"), flush=True) + sqs.delete_queue(QueueUrl=main_url) + sqs.delete_queue(QueueUrl=dlq_url) + + +def main(): + try: + check_fifo_order() + except Exception as e: + results["fifo_order"] = False + print("fifo_order: EXC %s %s" % (type(e).__name__, str(e)[:120]), flush=True) + try: + check_fifo_dedup() + except Exception as e: + results["fifo_dedup"] = False + print("fifo_dedup: EXC %s %s" % (type(e).__name__, str(e)[:120]), flush=True) + try: + check_dlq() + except Exception as e: + results["dlq"] = False + print("dlq: EXC %s %s" % (type(e).__name__, str(e)[:120]), flush=True) + print("SUMMARY: %s" % json.dumps(results), flush=True) + return 0 if all(results.values()) else 1 + + +if __name__ == "__main__": + sys.exit(main())