tiphys.queue
Outbound message queue with persistence and retry.
Provides reliable message delivery with automatic retry and dead letter handling.
Example: from tiphys.queue import InMemoryQueue, QueueWorker
# Create queue
queue = InMemoryQueue()
# Enqueue a message
msg_id = await queue.enqueue(
message=OutboundMessage(content="Hello", channel_name="webchat"),
channel_name="webchat",
priority=10,
)
# Create worker
async def deliver(msg):
print(f"Delivering: {msg.message.content}")
worker = QueueWorker(queue, deliver)
await worker.start()
1""" 2Outbound message queue with persistence and retry. 3 4Provides reliable message delivery with automatic retry and dead letter handling. 5 6Example: 7 from tiphys.queue import InMemoryQueue, QueueWorker 8 9 # Create queue 10 queue = InMemoryQueue() 11 12 # Enqueue a message 13 msg_id = await queue.enqueue( 14 message=OutboundMessage(content="Hello", channel_name="webchat"), 15 channel_name="webchat", 16 priority=10, 17 ) 18 19 # Create worker 20 async def deliver(msg): 21 print(f"Delivering: {msg.message.content}") 22 23 worker = QueueWorker(queue, deliver) 24 await worker.start() 25""" 26 27from tiphys.queue.base import ( 28 DeliveryStatus, 29 OutboundQueue, 30 QueuedMessage, 31) 32from tiphys.queue.memory import InMemoryQueue 33from tiphys.queue.worker import ( 34 QueueWorker, 35 QueueWorkerPool, 36 create_channel_delivery_func, 37) 38 39__all__ = [ 40 # Base classes 41 "OutboundQueue", 42 "QueuedMessage", 43 "DeliveryStatus", 44 # Implementations 45 "InMemoryQueue", 46 # Workers 47 "QueueWorker", 48 "QueueWorkerPool", 49 "create_channel_delivery_func", 50]
67class OutboundQueue(ABC): 68 """ 69 Abstract outbound message queue. 70 71 Implementations must handle persistence, ordering, and retry logic. 72 """ 73 74 @abstractmethod 75 async def enqueue( 76 self, 77 message: OutboundMessage, 78 channel_name: str, 79 *, 80 priority: int = 0, 81 delay_seconds: float = 0, 82 max_attempts: int = 3, 83 correlation_id: str | None = None, 84 ) -> str: 85 """ 86 Add a message to the queue. 87 88 Args: 89 message: The outbound message to queue. 90 channel_name: Target channel for delivery. 91 priority: Message priority (higher = more urgent). 92 delay_seconds: Delay before first delivery attempt. 93 max_attempts: Maximum delivery attempts. 94 correlation_id: Optional correlation ID for tracing. 95 96 Returns: 97 Message ID. 98 """ 99 ... 100 101 @abstractmethod 102 async def dequeue( 103 self, 104 channel_name: str | None = None, 105 batch_size: int = 1, 106 ) -> list[QueuedMessage]: 107 """ 108 Get messages ready for delivery. 109 110 Args: 111 channel_name: Optional channel filter. 112 batch_size: Maximum messages to return. 113 114 Returns: 115 List of ready messages (marked as IN_FLIGHT). 116 """ 117 ... 118 119 @abstractmethod 120 async def ack(self, message_id: str) -> bool: 121 """ 122 Acknowledge successful delivery. 123 124 Args: 125 message_id: ID of the delivered message. 126 127 Returns: 128 True if acknowledged, False if not found. 129 """ 130 ... 131 132 @abstractmethod 133 async def nack( 134 self, 135 message_id: str, 136 error: str, 137 *, 138 retry_delay_seconds: float = 60, 139 ) -> bool: 140 """ 141 Mark delivery as failed. 142 143 Will retry if attempts remaining, otherwise moves to dead letter. 144 145 Args: 146 message_id: ID of the failed message. 147 error: Error message. 148 retry_delay_seconds: Delay before retry. 149 150 Returns: 151 True if processed, False if not found. 152 """ 153 ... 154 155 @abstractmethod 156 async def get_pending_count(self, channel_name: str | None = None) -> int: 157 """ 158 Get count of pending messages. 159 160 Args: 161 channel_name: Optional channel filter. 162 163 Returns: 164 Number of pending messages. 165 """ 166 ... 167 168 @abstractmethod 169 async def get_dead_letters( 170 self, 171 channel_name: str | None = None, 172 limit: int = 100, 173 ) -> list[QueuedMessage]: 174 """ 175 Get messages that exceeded max retries. 176 177 Args: 178 channel_name: Optional channel filter. 179 limit: Maximum messages to return. 180 181 Returns: 182 List of dead letter messages. 183 """ 184 ... 185 186 @abstractmethod 187 async def retry_dead_letter(self, message_id: str) -> bool: 188 """ 189 Move a dead letter back to pending. 190 191 Args: 192 message_id: ID of the dead letter message. 193 194 Returns: 195 True if moved, False if not found. 196 """ 197 ... 198 199 @abstractmethod 200 async def delete(self, message_id: str) -> bool: 201 """ 202 Delete a message from the queue. 203 204 Args: 205 message_id: ID of the message to delete. 206 207 Returns: 208 True if deleted, False if not found. 209 """ 210 ... 211 212 @abstractmethod 213 async def clear(self, channel_name: str | None = None) -> int: 214 """ 215 Clear all messages. 216 217 Args: 218 channel_name: Optional channel filter. 219 220 Returns: 221 Number of messages cleared. 222 """ 223 ... 224 225 async def get_stats(self) -> dict[str, int]: 226 """ 227 Get queue statistics. 228 229 Returns: 230 Dict with pending, in_flight, dead_letter counts. 231 """ 232 pending = await self.get_pending_count() 233 dead = len(await self.get_dead_letters(limit=10000)) 234 return { 235 "pending": pending, 236 "dead_letter": dead, 237 }
Abstract outbound message queue.
Implementations must handle persistence, ordering, and retry logic.
74 @abstractmethod 75 async def enqueue( 76 self, 77 message: OutboundMessage, 78 channel_name: str, 79 *, 80 priority: int = 0, 81 delay_seconds: float = 0, 82 max_attempts: int = 3, 83 correlation_id: str | None = None, 84 ) -> str: 85 """ 86 Add a message to the queue. 87 88 Args: 89 message: The outbound message to queue. 90 channel_name: Target channel for delivery. 91 priority: Message priority (higher = more urgent). 92 delay_seconds: Delay before first delivery attempt. 93 max_attempts: Maximum delivery attempts. 94 correlation_id: Optional correlation ID for tracing. 95 96 Returns: 97 Message ID. 98 """ 99 ...
Add a message to the queue.
Args: message: The outbound message to queue. channel_name: Target channel for delivery. priority: Message priority (higher = more urgent). delay_seconds: Delay before first delivery attempt. max_attempts: Maximum delivery attempts. correlation_id: Optional correlation ID for tracing.
Returns: Message ID.
101 @abstractmethod 102 async def dequeue( 103 self, 104 channel_name: str | None = None, 105 batch_size: int = 1, 106 ) -> list[QueuedMessage]: 107 """ 108 Get messages ready for delivery. 109 110 Args: 111 channel_name: Optional channel filter. 112 batch_size: Maximum messages to return. 113 114 Returns: 115 List of ready messages (marked as IN_FLIGHT). 116 """ 117 ...
Get messages ready for delivery.
Args: channel_name: Optional channel filter. batch_size: Maximum messages to return.
Returns: List of ready messages (marked as IN_FLIGHT).
119 @abstractmethod 120 async def ack(self, message_id: str) -> bool: 121 """ 122 Acknowledge successful delivery. 123 124 Args: 125 message_id: ID of the delivered message. 126 127 Returns: 128 True if acknowledged, False if not found. 129 """ 130 ...
Acknowledge successful delivery.
Args: message_id: ID of the delivered message.
Returns: True if acknowledged, False if not found.
132 @abstractmethod 133 async def nack( 134 self, 135 message_id: str, 136 error: str, 137 *, 138 retry_delay_seconds: float = 60, 139 ) -> bool: 140 """ 141 Mark delivery as failed. 142 143 Will retry if attempts remaining, otherwise moves to dead letter. 144 145 Args: 146 message_id: ID of the failed message. 147 error: Error message. 148 retry_delay_seconds: Delay before retry. 149 150 Returns: 151 True if processed, False if not found. 152 """ 153 ...
Mark delivery as failed.
Will retry if attempts remaining, otherwise moves to dead letter.
Args: message_id: ID of the failed message. error: Error message. retry_delay_seconds: Delay before retry.
Returns: True if processed, False if not found.
155 @abstractmethod 156 async def get_pending_count(self, channel_name: str | None = None) -> int: 157 """ 158 Get count of pending messages. 159 160 Args: 161 channel_name: Optional channel filter. 162 163 Returns: 164 Number of pending messages. 165 """ 166 ...
Get count of pending messages.
Args: channel_name: Optional channel filter.
Returns: Number of pending messages.
168 @abstractmethod 169 async def get_dead_letters( 170 self, 171 channel_name: str | None = None, 172 limit: int = 100, 173 ) -> list[QueuedMessage]: 174 """ 175 Get messages that exceeded max retries. 176 177 Args: 178 channel_name: Optional channel filter. 179 limit: Maximum messages to return. 180 181 Returns: 182 List of dead letter messages. 183 """ 184 ...
Get messages that exceeded max retries.
Args: channel_name: Optional channel filter. limit: Maximum messages to return.
Returns: List of dead letter messages.
186 @abstractmethod 187 async def retry_dead_letter(self, message_id: str) -> bool: 188 """ 189 Move a dead letter back to pending. 190 191 Args: 192 message_id: ID of the dead letter message. 193 194 Returns: 195 True if moved, False if not found. 196 """ 197 ...
Move a dead letter back to pending.
Args: message_id: ID of the dead letter message.
Returns: True if moved, False if not found.
199 @abstractmethod 200 async def delete(self, message_id: str) -> bool: 201 """ 202 Delete a message from the queue. 203 204 Args: 205 message_id: ID of the message to delete. 206 207 Returns: 208 True if deleted, False if not found. 209 """ 210 ...
Delete a message from the queue.
Args: message_id: ID of the message to delete.
Returns: True if deleted, False if not found.
212 @abstractmethod 213 async def clear(self, channel_name: str | None = None) -> int: 214 """ 215 Clear all messages. 216 217 Args: 218 channel_name: Optional channel filter. 219 220 Returns: 221 Number of messages cleared. 222 """ 223 ...
Clear all messages.
Args: channel_name: Optional channel filter.
Returns: Number of messages cleared.
225 async def get_stats(self) -> dict[str, int]: 226 """ 227 Get queue statistics. 228 229 Returns: 230 Dict with pending, in_flight, dead_letter counts. 231 """ 232 pending = await self.get_pending_count() 233 dead = len(await self.get_dead_letters(limit=10000)) 234 return { 235 "pending": pending, 236 "dead_letter": dead, 237 }
Get queue statistics.
Returns: Dict with pending, in_flight, dead_letter counts.
31@dataclass 32class QueuedMessage: 33 """ 34 A message in the outbound queue. 35 36 Tracks delivery attempts, status, and scheduling information. 37 """ 38 39 message: OutboundMessage 40 channel_name: str 41 id: str = field(default_factory=lambda: str(uuid4())) 42 peer_id: str | None = None 43 priority: int = 0 # Higher = more urgent 44 status: DeliveryStatus = DeliveryStatus.PENDING 45 attempts: int = 0 46 max_attempts: int = 3 47 created_at: datetime = field(default_factory=datetime.utcnow) 48 scheduled_at: datetime | None = None 49 last_attempt_at: datetime | None = None 50 next_retry_at: datetime | None = None 51 error: str | None = None 52 correlation_id: str | None = None 53 54 def can_retry(self) -> bool: 55 """Check if the message can be retried.""" 56 return self.attempts < self.max_attempts 57 58 def is_ready(self) -> bool: 59 """Check if the message is ready for delivery.""" 60 if self.status != DeliveryStatus.PENDING: 61 return False 62 if self.scheduled_at is not None and self.scheduled_at > datetime.utcnow(): 63 return False 64 return not (self.next_retry_at is not None and self.next_retry_at > datetime.utcnow())
A message in the outbound queue.
Tracks delivery attempts, status, and scheduling information.
54 def can_retry(self) -> bool: 55 """Check if the message can be retried.""" 56 return self.attempts < self.max_attempts
Check if the message can be retried.
58 def is_ready(self) -> bool: 59 """Check if the message is ready for delivery.""" 60 if self.status != DeliveryStatus.PENDING: 61 return False 62 if self.scheduled_at is not None and self.scheduled_at > datetime.utcnow(): 63 return False 64 return not (self.next_retry_at is not None and self.next_retry_at > datetime.utcnow())
Check if the message is ready for delivery.
21class DeliveryStatus(str, Enum): 22 """Message delivery status.""" 23 24 PENDING = "pending" 25 IN_FLIGHT = "in_flight" 26 DELIVERED = "delivered" 27 FAILED = "failed" 28 DEAD_LETTER = "dead_letter"
Message delivery status.
20class InMemoryQueue(OutboundQueue): 21 """ 22 In-memory queue implementation. 23 24 Fast and simple, but messages are lost on restart. 25 Suitable for development, testing, and non-critical workloads. 26 """ 27 28 def __init__(self) -> None: 29 """Initialize an empty queue.""" 30 self._messages: dict[str, QueuedMessage] = {} 31 self._lock = asyncio.Lock() 32 33 async def enqueue( 34 self, 35 message: OutboundMessage, 36 channel_name: str, 37 *, 38 priority: int = 0, 39 delay_seconds: float = 0, 40 max_attempts: int = 3, 41 correlation_id: str | None = None, 42 ) -> str: 43 """Add a message to the queue.""" 44 scheduled_at = None 45 if delay_seconds > 0: 46 scheduled_at = datetime.utcnow() + timedelta(seconds=delay_seconds) 47 48 queued = QueuedMessage( 49 message=message, 50 channel_name=channel_name, 51 peer_id=message.peer_id, 52 priority=priority, 53 max_attempts=max_attempts, 54 scheduled_at=scheduled_at, 55 correlation_id=correlation_id, 56 ) 57 58 async with self._lock: 59 self._messages[queued.id] = queued 60 61 return queued.id 62 63 async def dequeue( 64 self, 65 channel_name: str | None = None, 66 batch_size: int = 1, 67 ) -> list[QueuedMessage]: 68 """Get messages ready for delivery.""" 69 async with self._lock: 70 ready: list[QueuedMessage] = [] 71 72 # Find ready messages sorted by priority (desc) then created_at (asc) 73 candidates = [ 74 msg 75 for msg in self._messages.values() 76 if msg.is_ready() and (channel_name is None or msg.channel_name == channel_name) 77 ] 78 79 # Sort: higher priority first, then older messages first 80 candidates.sort(key=lambda m: (-m.priority, m.created_at)) 81 82 # Take batch_size messages 83 for msg in candidates[:batch_size]: 84 msg.status = DeliveryStatus.IN_FLIGHT 85 msg.last_attempt_at = datetime.utcnow() 86 msg.attempts += 1 87 ready.append(msg) 88 89 return ready 90 91 async def ack(self, message_id: str) -> bool: 92 """Acknowledge successful delivery.""" 93 async with self._lock: 94 if message_id not in self._messages: 95 return False 96 97 msg = self._messages[message_id] 98 msg.status = DeliveryStatus.DELIVERED 99 # Remove delivered messages 100 del self._messages[message_id] 101 return True 102 103 async def nack( 104 self, 105 message_id: str, 106 error: str, 107 *, 108 retry_delay_seconds: float = 60, 109 ) -> bool: 110 """Mark delivery as failed.""" 111 async with self._lock: 112 if message_id not in self._messages: 113 return False 114 115 msg = self._messages[message_id] 116 msg.error = error 117 118 if msg.can_retry(): 119 # Schedule retry 120 msg.status = DeliveryStatus.PENDING 121 msg.next_retry_at = datetime.utcnow() + timedelta(seconds=retry_delay_seconds) 122 else: 123 # Move to dead letter 124 msg.status = DeliveryStatus.DEAD_LETTER 125 126 return True 127 128 async def get_pending_count(self, channel_name: str | None = None) -> int: 129 """Get count of pending messages.""" 130 async with self._lock: 131 return sum( 132 1 133 for msg in self._messages.values() 134 if msg.status == DeliveryStatus.PENDING 135 and (channel_name is None or msg.channel_name == channel_name) 136 ) 137 138 async def get_dead_letters( 139 self, 140 channel_name: str | None = None, 141 limit: int = 100, 142 ) -> list[QueuedMessage]: 143 """Get messages that exceeded max retries.""" 144 async with self._lock: 145 dead = [ 146 msg 147 for msg in self._messages.values() 148 if msg.status == DeliveryStatus.DEAD_LETTER 149 and (channel_name is None or msg.channel_name == channel_name) 150 ] 151 # Sort by created_at (oldest first) 152 dead.sort(key=lambda m: m.created_at) 153 return dead[:limit] 154 155 async def retry_dead_letter(self, message_id: str) -> bool: 156 """Move a dead letter back to pending.""" 157 async with self._lock: 158 if message_id not in self._messages: 159 return False 160 161 msg = self._messages[message_id] 162 if msg.status != DeliveryStatus.DEAD_LETTER: 163 return False 164 165 msg.status = DeliveryStatus.PENDING 166 msg.attempts = 0 167 msg.error = None 168 msg.next_retry_at = None 169 return True 170 171 async def delete(self, message_id: str) -> bool: 172 """Delete a message from the queue.""" 173 async with self._lock: 174 if message_id not in self._messages: 175 return False 176 del self._messages[message_id] 177 return True 178 179 async def clear(self, channel_name: str | None = None) -> int: 180 """Clear all messages.""" 181 async with self._lock: 182 if channel_name is None: 183 count = len(self._messages) 184 self._messages.clear() 185 return count 186 187 to_delete = [ 188 mid for mid, msg in self._messages.items() if msg.channel_name == channel_name 189 ] 190 for mid in to_delete: 191 del self._messages[mid] 192 return len(to_delete) 193 194 async def get_message(self, message_id: str) -> QueuedMessage | None: 195 """ 196 Get a message by ID. 197 198 Args: 199 message_id: The message ID. 200 201 Returns: 202 The message, or None if not found. 203 """ 204 async with self._lock: 205 return self._messages.get(message_id) 206 207 async def get_in_flight_count(self, channel_name: str | None = None) -> int: 208 """ 209 Get count of in-flight messages. 210 211 Args: 212 channel_name: Optional channel filter. 213 214 Returns: 215 Number of in-flight messages. 216 """ 217 async with self._lock: 218 return sum( 219 1 220 for msg in self._messages.values() 221 if msg.status == DeliveryStatus.IN_FLIGHT 222 and (channel_name is None or msg.channel_name == channel_name) 223 ) 224 225 async def get_stats(self) -> dict[str, int]: 226 """Get queue statistics.""" 227 async with self._lock: 228 stats = { 229 "pending": 0, 230 "in_flight": 0, 231 "dead_letter": 0, 232 "total": len(self._messages), 233 } 234 for msg in self._messages.values(): 235 if msg.status == DeliveryStatus.PENDING: 236 stats["pending"] += 1 237 elif msg.status == DeliveryStatus.IN_FLIGHT: 238 stats["in_flight"] += 1 239 elif msg.status == DeliveryStatus.DEAD_LETTER: 240 stats["dead_letter"] += 1 241 return stats
In-memory queue implementation.
Fast and simple, but messages are lost on restart. Suitable for development, testing, and non-critical workloads.
28 def __init__(self) -> None: 29 """Initialize an empty queue.""" 30 self._messages: dict[str, QueuedMessage] = {} 31 self._lock = asyncio.Lock()
Initialize an empty queue.
33 async def enqueue( 34 self, 35 message: OutboundMessage, 36 channel_name: str, 37 *, 38 priority: int = 0, 39 delay_seconds: float = 0, 40 max_attempts: int = 3, 41 correlation_id: str | None = None, 42 ) -> str: 43 """Add a message to the queue.""" 44 scheduled_at = None 45 if delay_seconds > 0: 46 scheduled_at = datetime.utcnow() + timedelta(seconds=delay_seconds) 47 48 queued = QueuedMessage( 49 message=message, 50 channel_name=channel_name, 51 peer_id=message.peer_id, 52 priority=priority, 53 max_attempts=max_attempts, 54 scheduled_at=scheduled_at, 55 correlation_id=correlation_id, 56 ) 57 58 async with self._lock: 59 self._messages[queued.id] = queued 60 61 return queued.id
Add a message to the queue.
63 async def dequeue( 64 self, 65 channel_name: str | None = None, 66 batch_size: int = 1, 67 ) -> list[QueuedMessage]: 68 """Get messages ready for delivery.""" 69 async with self._lock: 70 ready: list[QueuedMessage] = [] 71 72 # Find ready messages sorted by priority (desc) then created_at (asc) 73 candidates = [ 74 msg 75 for msg in self._messages.values() 76 if msg.is_ready() and (channel_name is None or msg.channel_name == channel_name) 77 ] 78 79 # Sort: higher priority first, then older messages first 80 candidates.sort(key=lambda m: (-m.priority, m.created_at)) 81 82 # Take batch_size messages 83 for msg in candidates[:batch_size]: 84 msg.status = DeliveryStatus.IN_FLIGHT 85 msg.last_attempt_at = datetime.utcnow() 86 msg.attempts += 1 87 ready.append(msg) 88 89 return ready
Get messages ready for delivery.
91 async def ack(self, message_id: str) -> bool: 92 """Acknowledge successful delivery.""" 93 async with self._lock: 94 if message_id not in self._messages: 95 return False 96 97 msg = self._messages[message_id] 98 msg.status = DeliveryStatus.DELIVERED 99 # Remove delivered messages 100 del self._messages[message_id] 101 return True
Acknowledge successful delivery.
103 async def nack( 104 self, 105 message_id: str, 106 error: str, 107 *, 108 retry_delay_seconds: float = 60, 109 ) -> bool: 110 """Mark delivery as failed.""" 111 async with self._lock: 112 if message_id not in self._messages: 113 return False 114 115 msg = self._messages[message_id] 116 msg.error = error 117 118 if msg.can_retry(): 119 # Schedule retry 120 msg.status = DeliveryStatus.PENDING 121 msg.next_retry_at = datetime.utcnow() + timedelta(seconds=retry_delay_seconds) 122 else: 123 # Move to dead letter 124 msg.status = DeliveryStatus.DEAD_LETTER 125 126 return True
Mark delivery as failed.
128 async def get_pending_count(self, channel_name: str | None = None) -> int: 129 """Get count of pending messages.""" 130 async with self._lock: 131 return sum( 132 1 133 for msg in self._messages.values() 134 if msg.status == DeliveryStatus.PENDING 135 and (channel_name is None or msg.channel_name == channel_name) 136 )
Get count of pending messages.
138 async def get_dead_letters( 139 self, 140 channel_name: str | None = None, 141 limit: int = 100, 142 ) -> list[QueuedMessage]: 143 """Get messages that exceeded max retries.""" 144 async with self._lock: 145 dead = [ 146 msg 147 for msg in self._messages.values() 148 if msg.status == DeliveryStatus.DEAD_LETTER 149 and (channel_name is None or msg.channel_name == channel_name) 150 ] 151 # Sort by created_at (oldest first) 152 dead.sort(key=lambda m: m.created_at) 153 return dead[:limit]
Get messages that exceeded max retries.
155 async def retry_dead_letter(self, message_id: str) -> bool: 156 """Move a dead letter back to pending.""" 157 async with self._lock: 158 if message_id not in self._messages: 159 return False 160 161 msg = self._messages[message_id] 162 if msg.status != DeliveryStatus.DEAD_LETTER: 163 return False 164 165 msg.status = DeliveryStatus.PENDING 166 msg.attempts = 0 167 msg.error = None 168 msg.next_retry_at = None 169 return True
Move a dead letter back to pending.
171 async def delete(self, message_id: str) -> bool: 172 """Delete a message from the queue.""" 173 async with self._lock: 174 if message_id not in self._messages: 175 return False 176 del self._messages[message_id] 177 return True
Delete a message from the queue.
179 async def clear(self, channel_name: str | None = None) -> int: 180 """Clear all messages.""" 181 async with self._lock: 182 if channel_name is None: 183 count = len(self._messages) 184 self._messages.clear() 185 return count 186 187 to_delete = [ 188 mid for mid, msg in self._messages.items() if msg.channel_name == channel_name 189 ] 190 for mid in to_delete: 191 del self._messages[mid] 192 return len(to_delete)
Clear all messages.
194 async def get_message(self, message_id: str) -> QueuedMessage | None: 195 """ 196 Get a message by ID. 197 198 Args: 199 message_id: The message ID. 200 201 Returns: 202 The message, or None if not found. 203 """ 204 async with self._lock: 205 return self._messages.get(message_id)
Get a message by ID.
Args: message_id: The message ID.
Returns: The message, or None if not found.
207 async def get_in_flight_count(self, channel_name: str | None = None) -> int: 208 """ 209 Get count of in-flight messages. 210 211 Args: 212 channel_name: Optional channel filter. 213 214 Returns: 215 Number of in-flight messages. 216 """ 217 async with self._lock: 218 return sum( 219 1 220 for msg in self._messages.values() 221 if msg.status == DeliveryStatus.IN_FLIGHT 222 and (channel_name is None or msg.channel_name == channel_name) 223 )
Get count of in-flight messages.
Args: channel_name: Optional channel filter.
Returns: Number of in-flight messages.
225 async def get_stats(self) -> dict[str, int]: 226 """Get queue statistics.""" 227 async with self._lock: 228 stats = { 229 "pending": 0, 230 "in_flight": 0, 231 "dead_letter": 0, 232 "total": len(self._messages), 233 } 234 for msg in self._messages.values(): 235 if msg.status == DeliveryStatus.PENDING: 236 stats["pending"] += 1 237 elif msg.status == DeliveryStatus.IN_FLIGHT: 238 stats["in_flight"] += 1 239 elif msg.status == DeliveryStatus.DEAD_LETTER: 240 stats["dead_letter"] += 1 241 return stats
Get queue statistics.
28class QueueWorker: 29 """ 30 Worker that processes messages from the outbound queue. 31 32 Continuously polls the queue for messages and delivers them. 33 Handles retry logic with exponential backoff. 34 """ 35 36 def __init__( 37 self, 38 queue: OutboundQueue, 39 delivery_func: DeliveryFunc, 40 *, 41 channel_name: str | None = None, 42 batch_size: int = 10, 43 poll_interval: float = 0.1, 44 initial_retry_delay: float = 5.0, 45 max_retry_delay: float = 300.0, 46 backoff_multiplier: float = 2.0, 47 ) -> None: 48 """ 49 Initialize the queue worker. 50 51 Args: 52 queue: The queue to process. 53 delivery_func: Async function that delivers a message. 54 channel_name: Optional channel filter. 55 batch_size: Messages to process per batch. 56 poll_interval: Seconds between polls when queue is empty. 57 initial_retry_delay: Initial delay for retries. 58 max_retry_delay: Maximum delay for retries. 59 backoff_multiplier: Multiplier for exponential backoff. 60 """ 61 self.queue = queue 62 self.delivery_func = delivery_func 63 self.channel_name = channel_name 64 self.batch_size = batch_size 65 self.poll_interval = poll_interval 66 self.initial_retry_delay = initial_retry_delay 67 self.max_retry_delay = max_retry_delay 68 self.backoff_multiplier = backoff_multiplier 69 70 self._running = False 71 self._task: asyncio.Task[None] | None = None 72 self._processed_count = 0 73 self._error_count = 0 74 75 @property 76 def is_running(self) -> bool: 77 """Check if the worker is running.""" 78 return self._running 79 80 @property 81 def processed_count(self) -> int: 82 """Number of successfully processed messages.""" 83 return self._processed_count 84 85 @property 86 def error_count(self) -> int: 87 """Number of failed deliveries.""" 88 return self._error_count 89 90 async def start(self) -> None: 91 """Start the worker.""" 92 if self._running: 93 return 94 95 self._running = True 96 self._task = asyncio.create_task(self._run()) 97 98 await logger.ainfo( 99 "Queue worker started", 100 channel=self.channel_name or "all", 101 batch_size=self.batch_size, 102 ) 103 104 async def stop(self, timeout: float = 5.0) -> None: 105 """ 106 Stop the worker gracefully. 107 108 Args: 109 timeout: Maximum seconds to wait for current batch. 110 """ 111 if not self._running: 112 return 113 114 self._running = False 115 116 if self._task: 117 try: 118 await asyncio.wait_for(self._task, timeout=timeout) 119 except TimeoutError: 120 self._task.cancel() 121 with contextlib.suppress(asyncio.CancelledError): 122 await self._task 123 self._task = None 124 125 await logger.ainfo( 126 "Queue worker stopped", 127 processed=self._processed_count, 128 errors=self._error_count, 129 ) 130 131 async def _run(self) -> None: 132 """Main worker loop.""" 133 while self._running: 134 try: 135 # Get batch of messages 136 messages = await self.queue.dequeue( 137 channel_name=self.channel_name, 138 batch_size=self.batch_size, 139 ) 140 141 if not messages: 142 # No messages, wait before polling again 143 await asyncio.sleep(self.poll_interval) 144 continue 145 146 # Process messages concurrently 147 await asyncio.gather( 148 *[self._process_message(msg) for msg in messages], 149 return_exceptions=True, 150 ) 151 152 except asyncio.CancelledError: 153 break 154 except Exception as e: 155 await logger.aerror("Worker loop error", error=str(e)) 156 await asyncio.sleep(1) # Back off on error 157 158 async def _process_message(self, msg: QueuedMessage) -> None: 159 """Process a single message.""" 160 try: 161 await self.delivery_func(msg) 162 await self.queue.ack(msg.id) 163 self._processed_count += 1 164 165 await logger.adebug( 166 "Message delivered", 167 message_id=msg.id, 168 channel=msg.channel_name, 169 correlation_id=msg.correlation_id, 170 ) 171 172 except Exception as e: 173 self._error_count += 1 174 retry_delay = self._calculate_retry_delay(msg.attempts) 175 176 await logger.awarning( 177 "Delivery failed", 178 message_id=msg.id, 179 channel=msg.channel_name, 180 attempt=msg.attempts, 181 retry_delay=retry_delay, 182 error=str(e), 183 ) 184 185 await self.queue.nack( 186 msg.id, 187 str(e), 188 retry_delay_seconds=retry_delay, 189 ) 190 191 def _calculate_retry_delay(self, attempts: int) -> float: 192 """Calculate retry delay with exponential backoff.""" 193 delay = self.initial_retry_delay * (self.backoff_multiplier ** (attempts - 1)) 194 return min(delay, self.max_retry_delay) 195 196 def reset_stats(self) -> None: 197 """Reset processed and error counts.""" 198 self._processed_count = 0 199 self._error_count = 0
Worker that processes messages from the outbound queue.
Continuously polls the queue for messages and delivers them. Handles retry logic with exponential backoff.
36 def __init__( 37 self, 38 queue: OutboundQueue, 39 delivery_func: DeliveryFunc, 40 *, 41 channel_name: str | None = None, 42 batch_size: int = 10, 43 poll_interval: float = 0.1, 44 initial_retry_delay: float = 5.0, 45 max_retry_delay: float = 300.0, 46 backoff_multiplier: float = 2.0, 47 ) -> None: 48 """ 49 Initialize the queue worker. 50 51 Args: 52 queue: The queue to process. 53 delivery_func: Async function that delivers a message. 54 channel_name: Optional channel filter. 55 batch_size: Messages to process per batch. 56 poll_interval: Seconds between polls when queue is empty. 57 initial_retry_delay: Initial delay for retries. 58 max_retry_delay: Maximum delay for retries. 59 backoff_multiplier: Multiplier for exponential backoff. 60 """ 61 self.queue = queue 62 self.delivery_func = delivery_func 63 self.channel_name = channel_name 64 self.batch_size = batch_size 65 self.poll_interval = poll_interval 66 self.initial_retry_delay = initial_retry_delay 67 self.max_retry_delay = max_retry_delay 68 self.backoff_multiplier = backoff_multiplier 69 70 self._running = False 71 self._task: asyncio.Task[None] | None = None 72 self._processed_count = 0 73 self._error_count = 0
Initialize the queue worker.
Args: queue: The queue to process. delivery_func: Async function that delivers a message. channel_name: Optional channel filter. batch_size: Messages to process per batch. poll_interval: Seconds between polls when queue is empty. initial_retry_delay: Initial delay for retries. max_retry_delay: Maximum delay for retries. backoff_multiplier: Multiplier for exponential backoff.
75 @property 76 def is_running(self) -> bool: 77 """Check if the worker is running.""" 78 return self._running
Check if the worker is running.
80 @property 81 def processed_count(self) -> int: 82 """Number of successfully processed messages.""" 83 return self._processed_count
Number of successfully processed messages.
85 @property 86 def error_count(self) -> int: 87 """Number of failed deliveries.""" 88 return self._error_count
Number of failed deliveries.
90 async def start(self) -> None: 91 """Start the worker.""" 92 if self._running: 93 return 94 95 self._running = True 96 self._task = asyncio.create_task(self._run()) 97 98 await logger.ainfo( 99 "Queue worker started", 100 channel=self.channel_name or "all", 101 batch_size=self.batch_size, 102 )
Start the worker.
104 async def stop(self, timeout: float = 5.0) -> None: 105 """ 106 Stop the worker gracefully. 107 108 Args: 109 timeout: Maximum seconds to wait for current batch. 110 """ 111 if not self._running: 112 return 113 114 self._running = False 115 116 if self._task: 117 try: 118 await asyncio.wait_for(self._task, timeout=timeout) 119 except TimeoutError: 120 self._task.cancel() 121 with contextlib.suppress(asyncio.CancelledError): 122 await self._task 123 self._task = None 124 125 await logger.ainfo( 126 "Queue worker stopped", 127 processed=self._processed_count, 128 errors=self._error_count, 129 )
Stop the worker gracefully.
Args: timeout: Maximum seconds to wait for current batch.
202class QueueWorkerPool: 203 """ 204 Pool of queue workers for parallel processing. 205 206 Manages multiple workers for higher throughput. 207 """ 208 209 def __init__( 210 self, 211 queue: OutboundQueue, 212 delivery_func: DeliveryFunc, 213 *, 214 worker_count: int = 3, 215 channel_name: str | None = None, 216 **worker_kwargs: Any, 217 ) -> None: 218 """ 219 Initialize the worker pool. 220 221 Args: 222 queue: The queue to process. 223 delivery_func: Async function that delivers a message. 224 worker_count: Number of workers. 225 channel_name: Optional channel filter. 226 **worker_kwargs: Additional arguments for workers. 227 """ 228 self.workers: list[QueueWorker] = [] 229 for _ in range(worker_count): 230 worker = QueueWorker( 231 queue=queue, 232 delivery_func=delivery_func, 233 channel_name=channel_name, 234 **worker_kwargs, 235 ) 236 self.workers.append(worker) 237 238 async def start(self) -> None: 239 """Start all workers.""" 240 await asyncio.gather(*[w.start() for w in self.workers]) 241 242 async def stop(self, timeout: float = 5.0) -> None: 243 """Stop all workers.""" 244 await asyncio.gather(*[w.stop(timeout) for w in self.workers]) 245 246 @property 247 def total_processed(self) -> int: 248 """Total messages processed across all workers.""" 249 return sum(w.processed_count for w in self.workers) 250 251 @property 252 def total_errors(self) -> int: 253 """Total errors across all workers.""" 254 return sum(w.error_count for w in self.workers)
Pool of queue workers for parallel processing.
Manages multiple workers for higher throughput.
209 def __init__( 210 self, 211 queue: OutboundQueue, 212 delivery_func: DeliveryFunc, 213 *, 214 worker_count: int = 3, 215 channel_name: str | None = None, 216 **worker_kwargs: Any, 217 ) -> None: 218 """ 219 Initialize the worker pool. 220 221 Args: 222 queue: The queue to process. 223 delivery_func: Async function that delivers a message. 224 worker_count: Number of workers. 225 channel_name: Optional channel filter. 226 **worker_kwargs: Additional arguments for workers. 227 """ 228 self.workers: list[QueueWorker] = [] 229 for _ in range(worker_count): 230 worker = QueueWorker( 231 queue=queue, 232 delivery_func=delivery_func, 233 channel_name=channel_name, 234 **worker_kwargs, 235 ) 236 self.workers.append(worker)
Initialize the worker pool.
Args: queue: The queue to process. delivery_func: Async function that delivers a message. worker_count: Number of workers. channel_name: Optional channel filter. **worker_kwargs: Additional arguments for workers.
238 async def start(self) -> None: 239 """Start all workers.""" 240 await asyncio.gather(*[w.start() for w in self.workers])
Start all workers.
242 async def stop(self, timeout: float = 5.0) -> None: 243 """Stop all workers.""" 244 await asyncio.gather(*[w.stop(timeout) for w in self.workers])
Stop all workers.
257def create_channel_delivery_func(context: AppContext) -> DeliveryFunc: 258 """ 259 Create a delivery function that uses channel bridges. 260 261 Args: 262 context: The application context. 263 264 Returns: 265 Async delivery function. 266 """ 267 268 async def deliver(msg: QueuedMessage) -> None: 269 """Deliver message via channel bridge.""" 270 # Get channel from plugin registry 271 channel = context.plugin_registry.get_channel(msg.channel_name) # type: ignore 272 if channel is None: 273 raise ValueError(f"Unknown channel: {msg.channel_name}") 274 275 # Build context for the channel 276 from tiphys.channels.base import ChannelContext 277 278 channel_ctx = ChannelContext( 279 channel_id=msg.peer_id or "", 280 thread_id=msg.message.thread_id, 281 reply_to=msg.message.reply_to, 282 ) 283 284 from tiphys.channels.base import OutboundMessage as ChannelOutboundMessage 285 286 channel_msg = ChannelOutboundMessage( 287 content=msg.message.content, 288 reply_to=msg.message.reply_to, 289 thread_id=msg.message.thread_id, 290 metadata=msg.message.metadata, 291 ) 292 293 # Deliver via channel 294 await channel.send_outbound(channel_msg, channel_ctx) 295 296 return deliver
Create a delivery function that uses channel bridges.
Args: context: The application context.
Returns: Async delivery function.