Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion awsimple/__version__.py
Original file line number Diff line number Diff line change
@@ -1,7 +1,7 @@
__application_name__ = "awsimple"
__title__ = __application_name__
__author__ = "abel"
__version__ = "7.1.5"
__version__ = "7.2.0"
__author_email__ = "j@abel.co"
__url__ = "https://github.com/jamesabel/awsimple"
__download_url__ = "https://github.com/jamesabel/awsimple"
Expand Down
39 changes: 19 additions & 20 deletions awsimple/pubsub.py
Original file line number Diff line number Diff line change
Expand Up @@ -4,7 +4,7 @@

import time
from functools import lru_cache
from typing import Any, Dict, List, Callable, Union
from typing import Any, Dict, List, Callable
from datetime import timedelta
from threading import Thread, Event
from queue import Queue
Expand Down Expand Up @@ -33,7 +33,7 @@

@typechecked()
def remove_old_queues(
channel: str, profile_name: Union[str, None] = None, aws_access_key_id: Union[str, None] = None, aws_secret_access_key: Union[str, None] = None, region_name: Union[str, None] = None
channel: str, profile_name: str | None = None, aws_access_key_id: str | None = None, aws_secret_access_key: str | None = None, region_name: str | None = None
) -> list[str]:
"""
Remove old SQS queues that have not been used recently.
Expand Down Expand Up @@ -77,8 +77,7 @@ def _connect_sns_to_sqs(sqs: SQSPollAccess, sns: SNSAccess) -> None:
topic = sns.resource.Topic(topic_arn)

# Subscribe queue to topic
queue_arn = sqs.get_arn()
subscription = topic.subscribe(Protocol="sqs", Endpoint=queue_arn)
subscription = topic.subscribe(Protocol="sqs", Endpoint=sqs_arn)
log.info(f"Subscribed {sqs.queue_name} to topic {topic_arn}. Subscription ARN: {subscription.arn}")

# Update queue policy to allow SNS -> SQS
Expand Down Expand Up @@ -107,8 +106,8 @@ class _SubscriptionThread(Thread):
"""

@typechecked()
def __init__(self, sqs: SQSPollAccess, new_event) -> None:
super().__init__()
def __init__(self, sqs: SQSPollAccess, new_event: Event) -> None:
super().__init__(daemon=True)
self._sqs = sqs
self.sub_queue = Queue() # type: Queue[str]
self._exit_event = Event()
Expand All @@ -118,8 +117,8 @@ def run(self):
while not self._exit_event.is_set():
messages = self._sqs.receive_messages() # long poll
for message in messages:
message = json.loads(message.message)
self.sub_queue.put(message["Message"])
parsed = json.loads(message.message)
self.sub_queue.put(parsed["Message"])
self._new_event.set()

def request_exit(self):
Expand Down Expand Up @@ -149,10 +148,10 @@ def __init__(
node_name: str | None,
sub_callback: Callable | None,
use_sub_queue: bool,
profile_name: Union[str, None],
aws_access_key_id: Union[str, None],
aws_secret_access_key: Union[str, None],
region_name: Union[str, None],
profile_name: str | None,
aws_access_key_id: str | None,
aws_secret_access_key: str | None,
region_name: str | None,
) -> None:
"""
Pub and Sub.
Expand Down Expand Up @@ -305,10 +304,10 @@ def __init__(
self,
channel: str,
node_name: str | None = None,
profile_name: Union[str, None] = None,
aws_access_key_id: Union[str, None] = None,
aws_secret_access_key: Union[str, None] = None,
region_name: Union[str, None] = None,
profile_name: str | None = None,
aws_access_key_id: str | None = None,
aws_secret_access_key: str | None = None,
region_name: str | None = None,
) -> None:
"""
Pub only.
Expand Down Expand Up @@ -336,10 +335,10 @@ def __init__(
channel: str,
node_name: str | None = None,
sub_callback: Callable | None = None,
profile_name: Union[str, None] = None,
aws_access_key_id: Union[str, None] = None,
aws_secret_access_key: Union[str, None] = None,
region_name: Union[str, None] = None,
profile_name: str | None = None,
aws_access_key_id: str | None = None,
aws_secret_access_key: str | None = None,
region_name: str | None = None,
) -> None:
"""
Sub only.
Expand Down
Loading