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

1"""Custom exceptions for message-broker operations. 

2 

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. 

9 

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. 

15 

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""" 

22 

23from __future__ import annotations 

24 

25from typing import Optional 

26 

27from ._causes import resolve_cause 

28from .base import DataExceptError 

29from .redaction import redact_if_url 

30 

31 

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) 

47 

48 

49def _on_broker(broker: Optional[str]) -> str: 

50 return f" on broker '{broker}'" if broker else "" 

51 

52 

53class MessageBrokerError(DataExceptError): 

54 """Base exception for message-broker failures. 

55 

56 Catching this catches every broker failure the library raises, without 

57 catching a database or HTTP one. 

58 """ 

59 

60 pass 

61 

62 

63class BrokerConnectionError(MessageBrokerError): 

64 """Raised when a connection to the broker cannot be established.""" 

65 

66 def __init__( 

67 self, 

68 broker: str, 

69 message: Optional[str] = None, 

70 *, 

71 cause: Optional[Exception] = None, 

72 ) -> None: 

73 """Initialize BrokerConnectionError. 

74 

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) 

88 

89 

90class BrokerTimeoutError(MessageBrokerError): 

91 """Raised when a broker operation exceeds its time limit.""" 

92 

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. 

103 

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) 

125 

126 

127class MessagePublishError(MessageBrokerError): 

128 """Raised when publishing a message fails.""" 

129 

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. 

140 

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) 

157 

158 

159class MessageConsumeError(MessageBrokerError): 

160 """Raised when consuming a message fails.""" 

161 

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. 

174 

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) 

195 

196 

197class MessageAcknowledgementError(MessageBrokerError): 

198 """Raised when acknowledging or committing a message fails. 

199 

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 """ 

204 

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. 

217 

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) 

238 

239 

240__all__ = [ 

241 "MessageBrokerError", 

242 "BrokerConnectionError", 

243 "BrokerTimeoutError", 

244 "MessagePublishError", 

245 "MessageConsumeError", 

246 "MessageAcknowledgementError", 

247]