HISTORY: FIFO/DLQ e2e — FIFO PASS, найден баг RedrivePolicy (tenant-scoped ключ)

This commit is contained in:
“Naeel”
2026-08-15 08:43:10 +04:00
parent 63261d1238
commit 26c77ff9e1
2 changed files with 136 additions and 0 deletions
+122
View File
@@ -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-<id>/<name>
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())