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
« prev ^ index » next coverage.py v7.16.0, created at 2026-09-27 14:44 +0000
1"""Additional exception classes for data pipeline workflows."""
3from __future__ import annotations
5from typing import Any, Optional
7from ._causes import resolve_cause
8from .base import DataExceptError
9from .failure_metadata import FailureMetadata
10from .redaction import redact_if_url, redact_url
13class PipelineError(DataExceptError):
14 """Base exception for pipeline errors."""
16 pass
19class PreprocessingError(PipelineError):
20 """Raised when a preprocessing step fails."""
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)
30class FeaturePreprocessingError(PreprocessingError):
31 """Raised when feature engineering fails."""
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)
39class StorageError(PipelineError):
40 """Raised when reading from or writing to storage fails."""
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)
57class PipelineNotificationError(PipelineError):
58 """Raised when sending a notification fails."""
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)
72class RetryLimitExceededError(PipelineError):
73 """Raised when an operation is retried too many times."""
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)
90class ExternalServiceError(PipelineError):
91 """General failure when calling an external service."""
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)
107class ServiceAuthenticationError(ExternalServiceError):
108 """Authentication to an external service failed."""
110 _default_failure_metadata = FailureMetadata(
111 failure_kind="permanent",
112 retryable=False,
113 )
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)
124class ServiceAuthorizationError(ExternalServiceError):
125 """Authorization was denied by an external service."""
127 _default_failure_metadata = FailureMetadata(
128 failure_kind="permanent",
129 retryable=False,
130 )
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)
141class ServiceTimeoutError(ExternalServiceError):
142 """A call to an external service exceeded the allotted time."""
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)
157class ApiError(PipelineError):
158 """Failure calling a REST API endpoint."""
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)
174class TimeDeltaTooLargeError(PipelineError):
175 """The time span between records exceeded a threshold."""
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)
189class TypeCheckError(PipelineError):
190 """Invalid type detected during recursive type inspection."""
193class DataFetchError(PipelineError):
194 """Failed to fetch data from a storage backend."""
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)
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]