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]
class OutboundQueue(abc.ABC):
 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.

@abstractmethod
async def enqueue( self, message: tiphys.pipeline.OutboundMessage, channel_name: str, *, priority: int = 0, delay_seconds: float = 0, max_attempts: int = 3, correlation_id: str | None = None) -> str:
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.

@abstractmethod
async def dequeue( self, channel_name: str | None = None, batch_size: int = 1) -> list[QueuedMessage]:
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).

@abstractmethod
async def ack(self, message_id: str) -> bool:
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.

@abstractmethod
async def nack( self, message_id: str, error: str, *, retry_delay_seconds: float = 60) -> bool:
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.

@abstractmethod
async def get_pending_count(self, channel_name: str | None = None) -> int:
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.

@abstractmethod
async def get_dead_letters( self, channel_name: str | None = None, limit: int = 100) -> list[QueuedMessage]:
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.

@abstractmethod
async def retry_dead_letter(self, message_id: str) -> bool:
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.

@abstractmethod
async def delete(self, message_id: str) -> bool:
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.

@abstractmethod
async def clear(self, channel_name: str | None = None) -> int:
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.

async def get_stats(self) -> dict[str, int]:
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.

@dataclass
class QueuedMessage:
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.

QueuedMessage( message: tiphys.pipeline.OutboundMessage, channel_name: str, id: str = <factory>, peer_id: str | None = None, priority: int = 0, status: DeliveryStatus = <DeliveryStatus.PENDING: 'pending'>, attempts: int = 0, max_attempts: int = 3, created_at: datetime.datetime = <factory>, scheduled_at: datetime.datetime | None = None, last_attempt_at: datetime.datetime | None = None, next_retry_at: datetime.datetime | None = None, error: str | None = None, correlation_id: str | None = None)
channel_name: str
id: str
peer_id: str | None = None
priority: int = 0
status: DeliveryStatus = <DeliveryStatus.PENDING: 'pending'>
attempts: int = 0
max_attempts: int = 3
created_at: datetime.datetime
scheduled_at: datetime.datetime | None = None
last_attempt_at: datetime.datetime | None = None
next_retry_at: datetime.datetime | None = None
error: str | None = None
correlation_id: str | None = None
def can_retry(self) -> bool:
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.

def is_ready(self) -> bool:
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.

class DeliveryStatus(builtins.str, enum.Enum):
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.

PENDING = <DeliveryStatus.PENDING: 'pending'>
IN_FLIGHT = <DeliveryStatus.IN_FLIGHT: 'in_flight'>
DELIVERED = <DeliveryStatus.DELIVERED: 'delivered'>
FAILED = <DeliveryStatus.FAILED: 'failed'>
DEAD_LETTER = <DeliveryStatus.DEAD_LETTER: 'dead_letter'>
class InMemoryQueue(tiphys.queue.OutboundQueue):
 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.

InMemoryQueue()
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.

async def enqueue( self, message: tiphys.pipeline.OutboundMessage, channel_name: str, *, priority: int = 0, delay_seconds: float = 0, max_attempts: int = 3, correlation_id: str | None = None) -> str:
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.

async def dequeue( self, channel_name: str | None = None, batch_size: int = 1) -> list[QueuedMessage]:
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.

async def ack(self, message_id: str) -> bool:
 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.

async def nack( self, message_id: str, error: str, *, retry_delay_seconds: float = 60) -> bool:
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.

async def get_pending_count(self, channel_name: str | None = None) -> int:
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.

async def get_dead_letters( self, channel_name: str | None = None, limit: int = 100) -> list[QueuedMessage]:
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.

async def retry_dead_letter(self, message_id: str) -> bool:
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.

async def delete(self, message_id: str) -> bool:
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.

async def clear(self, channel_name: str | None = None) -> int:
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.

async def get_message(self, message_id: str) -> QueuedMessage | None:
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.

async def get_in_flight_count(self, channel_name: str | None = None) -> int:
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.

async def get_stats(self) -> dict[str, int]:
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.

class QueueWorker:
 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.

QueueWorker( queue: OutboundQueue, delivery_func: Callable[[QueuedMessage], Awaitable[None]], *, channel_name: str | None = None, batch_size: int = 10, poll_interval: float = 0.1, initial_retry_delay: float = 5.0, max_retry_delay: float = 300.0, backoff_multiplier: float = 2.0)
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.

queue
delivery_func
channel_name
batch_size
poll_interval
initial_retry_delay
max_retry_delay
backoff_multiplier
is_running: bool
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.

processed_count: int
80    @property
81    def processed_count(self) -> int:
82        """Number of successfully processed messages."""
83        return self._processed_count

Number of successfully processed messages.

error_count: int
85    @property
86    def error_count(self) -> int:
87        """Number of failed deliveries."""
88        return self._error_count

Number of failed deliveries.

async def start(self) -> None:
 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.

async def stop(self, timeout: float = 5.0) -> None:
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.

def reset_stats(self) -> None:
196    def reset_stats(self) -> None:
197        """Reset processed and error counts."""
198        self._processed_count = 0
199        self._error_count = 0

Reset processed and error counts.

class QueueWorkerPool:
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.

QueueWorkerPool( queue: OutboundQueue, delivery_func: Callable[[QueuedMessage], Awaitable[None]], *, worker_count: int = 3, channel_name: str | None = None, **worker_kwargs: Any)
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.

workers: list[QueueWorker]
async def start(self) -> None:
238    async def start(self) -> None:
239        """Start all workers."""
240        await asyncio.gather(*[w.start() for w in self.workers])

Start all workers.

async def stop(self, timeout: float = 5.0) -> None:
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.

total_processed: int
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)

Total messages processed across all workers.

total_errors: int
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)

Total errors across all workers.

def create_channel_delivery_func( context: tiphys.context.AppContext) -> Callable[[QueuedMessage], Awaitable[None]]:
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.