Coverage for dataexcept/broker_exceptions.py: 100%
78 statements
« prev ^ index » next coverage.py v7.16.0, created at 2026-09-27 14:44 +0000
« prev ^ index » next coverage.py v7.16.0, created at 2026-09-27 14:44 +0000
1"""Custom exceptions for message-broker operations.
3A broker is a data-engineering boundary like a database or an object store, and
4it fails in ways that are specific to it: a publish is rejected, a consumer
5cannot reach its group, an offset is never committed. Mapping those onto
6``ServiceConnectionError`` and ``OperationTimeoutError`` loses which of them
7happened, and with it the topic, partition, offset and consumer group that say
8where.
10The hierarchy is deliberately about the *operation* rather than the product.
11Kafka, RabbitMQ, Pulsar and NATS disagree about almost everything else, but all
12four connect, publish, consume and acknowledge, so a pipeline can catch these
13without knowing which broker or client library is underneath -- and DataExcept
14depends on none of them.
16Failure metadata stays at the conservative default. A broker refusing a publish
17may be a leader election that resolves in a second or a topic that does not
18exist, and the exception cannot tell which. An integration that *does* know --
19because it read the broker's own error code -- attaches that with
20``with_failure_metadata`` or ``wrap(..., failure_metadata=...)``.
21"""
23from __future__ import annotations
25from typing import Optional
27from ._causes import resolve_cause
28from .base import DataExceptError
29from .redaction import redact_if_url
32def _describe_position(
33 topic: str,
34 partition: Optional[int] = None,
35 offset: Optional[int] = None,
36 consumer_group: Optional[str] = None,
37) -> str:
38 """Render as much of a message's coordinates as the caller supplied."""
39 parts = [f"topic '{topic}'"]
40 if partition is not None:
41 parts.append(f"partition {partition}")
42 if offset is not None:
43 parts.append(f"offset {offset}")
44 if consumer_group is not None:
45 parts.append(f"group '{consumer_group}'")
46 return ", ".join(parts)
49def _on_broker(broker: Optional[str]) -> str:
50 return f" on broker '{broker}'" if broker else ""
53class MessageBrokerError(DataExceptError):
54 """Base exception for message-broker failures.
56 Catching this catches every broker failure the library raises, without
57 catching a database or HTTP one.
58 """
60 pass
63class BrokerConnectionError(MessageBrokerError):
64 """Raised when a connection to the broker cannot be established."""
66 def __init__(
67 self,
68 broker: str,
69 message: Optional[str] = None,
70 *,
71 cause: Optional[Exception] = None,
72 ) -> None:
73 """Initialize BrokerConnectionError.
75 Args:
76 broker: The broker being connected to. A bootstrap address, a
77 cluster name, or a URL -- redacted when it is a URL, since one
78 routinely carries credentials.
79 message: Optional custom error message.
80 cause: Optional underlying exception, also set as ``__cause__``.
81 """
82 self.broker = redact_if_url(broker)
83 self.cause = resolve_cause(cause=cause)
84 default = f"Failed to connect to message broker '{self.broker}'"
85 if self.cause:
86 default += f": {self.cause}"
87 super().__init__(message or default)
90class BrokerTimeoutError(MessageBrokerError):
91 """Raised when a broker operation exceeds its time limit."""
93 def __init__(
94 self,
95 broker: str,
96 operation: Optional[str] = None,
97 timeout_seconds: Optional[float] = None,
98 message: Optional[str] = None,
99 *,
100 cause: Optional[Exception] = None,
101 ) -> None:
102 """Initialize BrokerTimeoutError.
104 Args:
105 broker: The broker the operation was waiting on.
106 operation: What was being attempted, such as ``"publish"``.
107 timeout_seconds: The limit that was exceeded.
108 message: Optional custom error message.
109 cause: Optional underlying exception, also set as ``__cause__``.
110 """
111 self.broker = redact_if_url(broker)
112 self.operation = operation
113 self.timeout_seconds = timeout_seconds
114 self.cause = resolve_cause(cause=cause)
115 attempted = (
116 f"Broker operation '{operation}'" if operation else "Broker operation"
117 )
118 default = f"{attempted} timed out"
119 if timeout_seconds is not None:
120 default += f" after {timeout_seconds}s"
121 default += _on_broker(self.broker)
122 if self.cause:
123 default += f": {self.cause}"
124 super().__init__(message or default)
127class MessagePublishError(MessageBrokerError):
128 """Raised when publishing a message fails."""
130 def __init__(
131 self,
132 topic: str,
133 broker: Optional[str] = None,
134 partition: Optional[int] = None,
135 message: Optional[str] = None,
136 *,
137 cause: Optional[Exception] = None,
138 ) -> None:
139 """Initialize MessagePublishError.
141 Args:
142 topic: The topic, queue or subject published to.
143 broker: Optional broker the publish was addressed to.
144 partition: Optional partition the message was keyed to.
145 message: Optional custom error message.
146 cause: Optional underlying exception, also set as ``__cause__``.
147 """
148 self.topic = topic
149 self.broker = redact_if_url(broker)
150 self.partition = partition
151 self.cause = resolve_cause(cause=cause)
152 position = _describe_position(topic, partition)
153 default = f"Failed to publish to {position}{_on_broker(self.broker)}"
154 if self.cause:
155 default += f": {self.cause}"
156 super().__init__(message or default)
159class MessageConsumeError(MessageBrokerError):
160 """Raised when consuming a message fails."""
162 def __init__(
163 self,
164 topic: str,
165 broker: Optional[str] = None,
166 partition: Optional[int] = None,
167 offset: Optional[int] = None,
168 consumer_group: Optional[str] = None,
169 message: Optional[str] = None,
170 *,
171 cause: Optional[Exception] = None,
172 ) -> None:
173 """Initialize MessageConsumeError.
175 Args:
176 topic: The topic, queue or subject consumed from.
177 broker: Optional broker the consumer was reading from.
178 partition: Optional partition being read.
179 offset: Optional offset the failure happened at.
180 consumer_group: Optional consumer group the reader belongs to.
181 message: Optional custom error message.
182 cause: Optional underlying exception, also set as ``__cause__``.
183 """
184 self.topic = topic
185 self.broker = redact_if_url(broker)
186 self.partition = partition
187 self.offset = offset
188 self.consumer_group = consumer_group
189 self.cause = resolve_cause(cause=cause)
190 position = _describe_position(topic, partition, offset, consumer_group)
191 default = f"Failed to consume from {position}{_on_broker(self.broker)}"
192 if self.cause:
193 default += f": {self.cause}"
194 super().__init__(message or default)
197class MessageAcknowledgementError(MessageBrokerError):
198 """Raised when acknowledging or committing a message fails.
200 Distinct from a consume failure on purpose: the message was read and
201 processed, and it is the record of that which did not stick -- so it will
202 be delivered again.
203 """
205 def __init__(
206 self,
207 topic: str,
208 broker: Optional[str] = None,
209 partition: Optional[int] = None,
210 offset: Optional[int] = None,
211 consumer_group: Optional[str] = None,
212 message: Optional[str] = None,
213 *,
214 cause: Optional[Exception] = None,
215 ) -> None:
216 """Initialize MessageAcknowledgementError.
218 Args:
219 topic: The topic, queue or subject the message came from.
220 broker: Optional broker the acknowledgement was sent to.
221 partition: Optional partition the message was read from.
222 offset: Optional offset that failed to commit.
223 consumer_group: Optional consumer group committing the offset.
224 message: Optional custom error message.
225 cause: Optional underlying exception, also set as ``__cause__``.
226 """
227 self.topic = topic
228 self.broker = redact_if_url(broker)
229 self.partition = partition
230 self.offset = offset
231 self.consumer_group = consumer_group
232 self.cause = resolve_cause(cause=cause)
233 position = _describe_position(topic, partition, offset, consumer_group)
234 default = f"Failed to acknowledge {position}{_on_broker(self.broker)}"
235 if self.cause:
236 default += f": {self.cause}"
237 super().__init__(message or default)
240__all__ = [
241 "MessageBrokerError",
242 "BrokerConnectionError",
243 "BrokerTimeoutError",
244 "MessagePublishError",
245 "MessageConsumeError",
246 "MessageAcknowledgementError",
247]