import json
import time
import boto3
REGION = "ap-east-1"
QUEUE_NAME = "dev-order-payment-capture"
DLQ_NAME = "dev-order-payment-capture-dlq"
sqs = boto3.client("sqs", region_name=REGION)
def ensure_queue(queue_name: str, attributes: dict[str, str]) -> str:
sqs.create_queue(QueueName=queue_name, Attributes=attributes)
response = sqs.get_queue_url(QueueName=queue_name)
return response["QueueUrl"]
def get_queue_arn(queue_url: str) -> str:
response = sqs.get_queue_attributes(
QueueUrl=queue_url,
AttributeNames=["QueueArn"],
)
return response["Attributes"]["QueueArn"]
def create_source_and_dlq() -> tuple[str, str]:
dlq_url = ensure_queue(
DLQ_NAME,
{
"MessageRetentionPeriod": "1209600",
"ReceiveMessageWaitTimeSeconds": "20",
"SqsManagedSseEnabled": "true",
},
)
dlq_arn = get_queue_arn(dlq_url)
source_url = ensure_queue(
QUEUE_NAME,
{
"MessageRetentionPeriod": "345600",
"VisibilityTimeout": "120",
"ReceiveMessageWaitTimeSeconds": "20",
"SqsManagedSseEnabled": "true",
"RedrivePolicy": json.dumps(
{
"deadLetterTargetArn": dlq_arn,
"maxReceiveCount": "5",
}
),
},
)
return source_url, dlq_url
def read_queue(queue_url: str) -> dict[str, str]:
response = sqs.get_queue_attributes(
QueueUrl=queue_url,
AttributeNames=[
"QueueArn",
"VisibilityTimeout",
"ReceiveMessageWaitTimeSeconds",
"RedrivePolicy",
],
)
return response["Attributes"]
def update_queue(queue_url: str) -> None:
sqs.set_queue_attributes(
QueueUrl=queue_url,
Attributes={
"VisibilityTimeout": "180",
"ReceiveMessageWaitTimeSeconds": "20",
},
)
def send_message(queue_url: str) -> None:
body = {
"schema_version": 1,
"event_type": "payment.capture.requested",
"request_id": "req-123",
"order_id": "ord-1001",
"amount_cents": 1200,
}
sqs.send_message(
QueueUrl=queue_url,
MessageBody=json.dumps(body),
MessageAttributes={
"trace_id": {
"DataType": "String",
"StringValue": "trace-abc",
},
"schema_version": {
"DataType": "Number",
"StringValue": "1",
},
},
)
def receive_one(queue_url: str) -> dict | None:
response = sqs.receive_message(
QueueUrl=queue_url,
MaxNumberOfMessages=1,
WaitTimeSeconds=20,
AttributeNames=["ApproximateReceiveCount", "SentTimestamp"],
MessageAttributeNames=["All"],
)
messages = response.get("Messages", [])
return messages[0] if messages else None
def extend_visibility(queue_url: str, receipt_handle: str, timeout_seconds: int) -> None:
sqs.change_message_visibility(
QueueUrl=queue_url,
ReceiptHandle=receipt_handle,
VisibilityTimeout=timeout_seconds,
)
def delete_message(queue_url: str, receipt_handle: str) -> None:
sqs.delete_message(QueueUrl=queue_url, ReceiptHandle=receipt_handle)
def purge_queue(queue_url: str) -> None:
sqs.purge_queue(QueueUrl=queue_url)
def delete_queue(queue_url: str) -> None:
sqs.delete_queue(QueueUrl=queue_url)
def main() -> None:
source_url, dlq_url = create_source_and_dlq()
print("source attributes:", read_queue(source_url))
update_queue(source_url)
print("source attributes after update:", read_queue(source_url))
send_message(source_url)
message = receive_one(source_url)
print("received message:", json.dumps(message, indent=2, default=str))
if message:
receipt_handle = message["ReceiptHandle"]
extend_visibility(source_url, receipt_handle, 300)
time.sleep(1)
delete_message(source_url, receipt_handle)
# optional cleanup in dev only
purge_queue(source_url)
purge_queue(dlq_url)
delete_queue(source_url)
delete_queue(dlq_url)
if __name__ == "__main__":
main()