Coverage for dataexcept/pipeline_exceptions.py: 100%

83 statements  

« prev     ^ index     » next       coverage.py v7.16.0, created at 2026-09-27 14:44 +0000

1"""Additional exception classes for data pipeline workflows.""" 

2 

3from __future__ import annotations 

4 

5from typing import Any, Optional 

6 

7from ._causes import resolve_cause 

8from .base import DataExceptError 

9from .failure_metadata import FailureMetadata 

10from .redaction import redact_if_url, redact_url 

11 

12 

13class PipelineError(DataExceptError): 

14 """Base exception for pipeline errors.""" 

15 

16 pass 

17 

18 

19class PreprocessingError(PipelineError): 

20 """Raised when a preprocessing step fails.""" 

21 

22 def __init__(self, step_name: str, details: Optional[str] = None) -> None: 

23 default = f"Preprocessing failed at step: '{step_name}'." 

24 message = f"{default} Details: {details}" if details else default 

25 self.step_name = step_name 

26 self.details = details 

27 super().__init__(message) 

28 

29 

30class FeaturePreprocessingError(PreprocessingError): 

31 """Raised when feature engineering fails.""" 

32 

33 def __init__(self, feature: str, reason: Optional[str] = None) -> None: 

34 self.feature = feature 

35 self.reason = reason 

36 super().__init__(step_name=f"feature_{feature}", details=reason) 

37 

38 

39class StorageError(PipelineError): 

40 """Raised when reading from or writing to storage fails.""" 

41 

42 def __init__( 

43 self, 

44 location: str, 

45 operation: str, 

46 message: Optional[str] = None, 

47 *, 

48 cause: Exception | None = None, 

49 ) -> None: 

50 default = f"Storage {operation} failed at location: '{location}'." 

51 self.location = redact_if_url(location) 

52 self.operation = operation 

53 self.cause = resolve_cause(cause=cause) 

54 super().__init__(message or default) 

55 

56 

57class PipelineNotificationError(PipelineError): 

58 """Raised when sending a notification fails.""" 

59 

60 def __init__( 

61 self, 

62 channel: str, 

63 payload: Any, 

64 message: Optional[str] = None, 

65 ) -> None: 

66 default = f"Notification via '{channel}' failed." 

67 self.channel = channel 

68 self.payload = payload 

69 super().__init__(message or default) 

70 

71 

72class RetryLimitExceededError(PipelineError): 

73 """Raised when an operation is retried too many times.""" 

74 

75 def __init__( 

76 self, 

77 operation: str, 

78 retries: int, 

79 message: Optional[str] = None, 

80 ) -> None: 

81 default = ( 

82 "Retry limit exceeded for operation " 

83 f"'{operation}' after {retries} attempts." 

84 ) 

85 self.operation = operation 

86 self.retries = retries 

87 super().__init__(message or default) 

88 

89 

90class ExternalServiceError(PipelineError): 

91 """General failure when calling an external service.""" 

92 

93 def __init__( 

94 self, 

95 service_name: str, 

96 status_code: Optional[int] = None, 

97 response: Optional[Any] = None, 

98 message: Optional[str] = None, 

99 ) -> None: 

100 default = f"Call to external service '{service_name}' failed." 

101 self.service_name = service_name 

102 self.status_code = status_code 

103 self.response = response 

104 super().__init__(message or default) 

105 

106 

107class ServiceAuthenticationError(ExternalServiceError): 

108 """Authentication to an external service failed.""" 

109 

110 _default_failure_metadata = FailureMetadata( 

111 failure_kind="permanent", 

112 retryable=False, 

113 ) 

114 

115 def __init__( 

116 self, 

117 service_name: str, 

118 message: Optional[str] = None, 

119 ) -> None: 

120 default = f"Authentication failed for service '{service_name}'." 

121 super().__init__(service_name=service_name, message=message or default) 

122 

123 

124class ServiceAuthorizationError(ExternalServiceError): 

125 """Authorization was denied by an external service.""" 

126 

127 _default_failure_metadata = FailureMetadata( 

128 failure_kind="permanent", 

129 retryable=False, 

130 ) 

131 

132 def __init__( 

133 self, 

134 service_name: str, 

135 message: Optional[str] = None, 

136 ) -> None: 

137 default = f"Authorization denied for service '{service_name}'." 

138 super().__init__(service_name=service_name, message=message or default) 

139 

140 

141class ServiceTimeoutError(ExternalServiceError): 

142 """A call to an external service exceeded the allotted time.""" 

143 

144 def __init__( 

145 self, 

146 service_name: str, 

147 timeout_seconds: Optional[float] = None, 

148 ) -> None: 

149 default = ( 

150 "Operation timed out after " 

151 f"{timeout_seconds}s on service '{service_name}'." 

152 ) 

153 self.timeout_seconds = timeout_seconds 

154 super().__init__(service_name=service_name, message=default) 

155 

156 

157class ApiError(PipelineError): 

158 """Failure calling a REST API endpoint.""" 

159 

160 def __init__( 

161 self, 

162 endpoint: str, 

163 status_code: Optional[int] = None, 

164 message: Optional[str] = None, 

165 ) -> None: 

166 self.endpoint = redact_url(endpoint) 

167 default = f"API call failed: {self.endpoint}" 

168 if status_code is not None: 

169 default += f" (status {status_code})" 

170 self.status_code = status_code 

171 super().__init__(message or default) 

172 

173 

174class TimeDeltaTooLargeError(PipelineError): 

175 """The time span between records exceeded a threshold.""" 

176 

177 def __init__( 

178 self, 

179 user: str, 

180 delta_minutes: float, 

181 message: Optional[str] = None, 

182 ) -> None: 

183 default = f"Time delta {delta_minutes}m too large for user {user}" 

184 self.user = user 

185 self.delta_minutes = delta_minutes 

186 super().__init__(message or default) 

187 

188 

189class TypeCheckError(PipelineError): 

190 """Invalid type detected during recursive type inspection.""" 

191 

192 

193class DataFetchError(PipelineError): 

194 """Failed to fetch data from a storage backend.""" 

195 

196 def __init__( 

197 self, 

198 source: str, 

199 cid: str, 

200 message: Optional[str] = None, 

201 ) -> None: 

202 default = f"Failed to fetch '{source}' data for cid={cid}" 

203 self.source = redact_if_url(source) 

204 self.cid = cid 

205 super().__init__(message or default) 

206 

207 

208__all__ = [ 

209 "PipelineError", 

210 "PreprocessingError", 

211 "FeaturePreprocessingError", 

212 "StorageError", 

213 "PipelineNotificationError", 

214 "RetryLimitExceededError", 

215 "ExternalServiceError", 

216 "ServiceAuthenticationError", 

217 "ServiceAuthorizationError", 

218 "ServiceTimeoutError", 

219 "ApiError", 

220 "TimeDeltaTooLargeError", 

221 "TypeCheckError", 

222 "DataFetchError", 

223]