Coverage for src/ai_lls_lib/providers/external.py: 95%
161 statements
« prev ^ index » next coverage.py v7.15.2, created at 2026-08-24 12:44 +0000
« prev ^ index » next coverage.py v7.15.2, created at 2026-08-24 12:44 +0000
1"""
2External API provider for production phone verification
3"""
5import os
6import random
7import threading
8import time
9from collections.abc import Callable
11import httpx
12from aws_lambda_powertools import Logger
14from ..core.models import LineType, VerificationResult
15from .exceptions import (
16 ProviderAuthError,
17 ProviderError,
18 ProviderQuotaError,
19 ProviderUnavailableError,
20)
22logger = Logger()
24RETRYABLE_STATUSES = frozenset({408, 429, 500, 502, 503, 504})
25RETRY_AFTER_CAP_SECONDS = 60.0
27# Best-effort classification of the vendor's HTTP 400 `errorMsg`. The vendor returns 400 for
28# a missing key, an invalid key and an exhausted balance alike; these keyword sets are the
29# only signal available and must be kept in step with observed real responses.
30QUOTA_KEYWORDS = ("balance", "credit", "quota", "insufficient", "exhaust", "limit reached")
31AUTH_KEYWORDS = ("api key", "apikey", "api_key", "unauthori", "invalid key")
34def _env_float(name: str, default: float) -> float:
35 raw = os.environ.get(name)
36 return float(raw) if raw not in (None, "") else default
39def _env_int(name: str, default: int) -> int:
40 raw = os.environ.get(name)
41 return int(raw) if raw not in (None, "") else default
44class ExternalAPIProvider:
45 """
46 Production provider that calls external verification APIs.
47 Uses landlineremover.com which returns both line type and DNC status in a single call.
48 """
50 def __init__(
51 self,
52 timeout: float = 10.0,
53 *,
54 max_retries: int | None = None,
55 backoff_base: float = 0.5,
56 backoff_max: float = 8.0,
57 rate_limit_per_second: float | None = None,
58 breaker_failure_threshold: int | None = None,
59 breaker_reset_seconds: float = 30.0,
60 pool_size: int = 20,
61 transport: httpx.BaseTransport | None = None,
62 ):
63 """
64 Initialize external API provider.
66 Resilience settings default from the environment so deployments can tune them
67 without a code change:
68 LLS_PROVIDER_MAX_RETRIES (default 3): retries on 429/5xx/transport errors
69 LLS_PROVIDER_RATE_LIMIT_PER_SECOND (default 0 = off): client-side ceiling
70 LLS_PROVIDER_BREAKER_THRESHOLD (default 5): consecutive failed calls (after
71 retries) that open the circuit for breaker_reset_seconds
73 The limiter and breaker are process-local: they bound one Lambda invocation or
74 worker pool, not the fleet.
76 Args:
77 timeout: HTTP request timeout in seconds
78 max_retries: bounded retry count; backoff is exponential with full jitter
79 backoff_base: first backoff ceiling in seconds; doubles per attempt
80 backoff_max: cap on a single backoff wait
81 rate_limit_per_second: token-bucket rate; burst equals the rate
82 breaker_failure_threshold: consecutive failures that open the circuit
83 breaker_reset_seconds: how long the circuit stays open before one probe call
84 pool_size: httpx connection pool size, sized for parallel workers
85 transport: httpx transport override for tests (never call the vendor in tests)
86 """
87 self.api_key = os.environ.get("LANDLINE_REMOVER_API_KEY", "")
88 if not self.api_key:
89 logger.warning("LANDLINE_REMOVER_API_KEY not set")
91 self.api_url = "https://app.landlineremover.com/api/check-number"
92 self.timeout = timeout
93 self.max_retries = (
94 max_retries if max_retries is not None else _env_int("LLS_PROVIDER_MAX_RETRIES", 3)
95 )
96 self.backoff_base = backoff_base
97 self.backoff_max = backoff_max
98 self.rate_limit_per_second = (
99 rate_limit_per_second
100 if rate_limit_per_second is not None
101 else _env_float("LLS_PROVIDER_RATE_LIMIT_PER_SECOND", 0.0)
102 )
103 self.breaker_failure_threshold = (
104 breaker_failure_threshold
105 if breaker_failure_threshold is not None
106 else _env_int("LLS_PROVIDER_BREAKER_THRESHOLD", 5)
107 )
108 self.breaker_reset_seconds = breaker_reset_seconds
109 self.pool_limits = httpx.Limits(
110 max_connections=pool_size, max_keepalive_connections=pool_size
111 )
112 self.http_client = httpx.Client(
113 timeout=timeout, limits=self.pool_limits, transport=transport
114 )
116 # Injectable for tests; never patch time globally
117 self.sleep: Callable[[float], None] = time.sleep
118 self.clock: Callable[[], float] = time.monotonic
120 self._lock = threading.Lock()
121 self._tokens = float(self.rate_limit_per_second)
122 self._last_refill = self.clock()
123 self._consecutive_failures = 0
124 self._breaker_opened_at: float | None = None
125 self._breaker_probe_in_flight = False
127 # ---- rate limiter -------------------------------------------------------------
129 def _acquire_token(self) -> None:
130 """Token bucket: wait until a request token is available (no-op when disabled)."""
131 if self.rate_limit_per_second <= 0:
132 return
133 while True:
134 with self._lock:
135 now = self.clock()
136 elapsed = max(0.0, now - self._last_refill)
137 self._tokens = min(
138 float(self.rate_limit_per_second),
139 self._tokens + elapsed * self.rate_limit_per_second,
140 )
141 self._last_refill = now
142 if self._tokens >= 1.0:
143 self._tokens -= 1.0
144 return
145 wait = (1.0 - self._tokens) / self.rate_limit_per_second
146 self.sleep(wait)
148 # ---- circuit breaker ----------------------------------------------------------
150 def _check_breaker(self) -> None:
151 """Fail fast while the circuit is open; allow a single probe once the reset elapses."""
152 with self._lock:
153 if self._breaker_opened_at is None:
154 return
155 elapsed = self.clock() - self._breaker_opened_at
156 if elapsed >= self.breaker_reset_seconds and not self._breaker_probe_in_flight:
157 self._breaker_probe_in_flight = True # half-open: one call may try
158 return
159 raise ProviderUnavailableError(
160 f"Provider circuit open after {self._consecutive_failures} consecutive failures; "
161 f"retry in {max(0.0, self.breaker_reset_seconds - elapsed):.0f}s"
162 )
164 def _record_success(self) -> None:
165 with self._lock:
166 self._consecutive_failures = 0
167 self._breaker_opened_at = None
168 self._breaker_probe_in_flight = False
170 def _record_failure(self) -> None:
171 with self._lock:
172 self._consecutive_failures += 1
173 self._breaker_probe_in_flight = False
174 if self._consecutive_failures >= self.breaker_failure_threshold:
175 if self._breaker_opened_at is None:
176 logger.error(
177 f"Provider circuit opened after {self._consecutive_failures} "
178 f"consecutive failures"
179 )
180 self._breaker_opened_at = self.clock()
182 # ---- retry --------------------------------------------------------------------
184 def _backoff(self, attempt: int, response: httpx.Response | None) -> float:
185 """Seconds to wait before retry `attempt` (0-based); honours Retry-After."""
186 if response is not None:
187 retry_after = response.headers.get("Retry-After")
188 if retry_after:
189 try:
190 return min(float(retry_after), RETRY_AFTER_CAP_SECONDS)
191 except ValueError:
192 pass # HTTP-date form: fall through to computed backoff
193 ceiling = min(self.backoff_base * (2**attempt), self.backoff_max)
194 return random.uniform(0, ceiling) # noqa: S311 - jitter, not security
196 def _classify_400(self, response: httpx.Response) -> ProviderError:
197 """Turn the vendor's catch-all 400 into a specific exception where the body allows."""
198 message = ""
199 try:
200 body = response.json()
201 if isinstance(body, dict):
202 message = str(body.get("errorMsg") or body.get("error") or "")
203 except ValueError:
204 message = response.text
205 lowered = message.lower()
206 if any(k in lowered for k in QUOTA_KEYWORDS):
207 return ProviderQuotaError(f"Provider quota exhausted: {message}", status_code=400)
208 if any(k in lowered for k in AUTH_KEYWORDS):
209 return ProviderAuthError(f"Provider rejected API key: {message}", status_code=400)
210 return ProviderError(f"API request failed with status 400: {message}", status_code=400)
212 def _request(self, phone: str) -> httpx.Response:
213 """GET with bounded retry on 429/5xx/transport errors. Raises ProviderError."""
214 attempt = 0
215 while True:
216 self._acquire_token()
217 response: httpx.Response | None = None
218 error: ProviderError
219 try:
220 response = self.http_client.get(
221 self.api_url,
222 params={"apikey": self.api_key, "number": phone},
223 follow_redirects=True,
224 )
225 if response.status_code < 400:
226 return response
227 if response.status_code == 400:
228 raise self._classify_400(response)
229 if response.status_code not in RETRYABLE_STATUSES:
230 logger.error(f"API error: {response.status_code} - {response.text}")
231 raise ProviderError(
232 f"API request failed with status {response.status_code}",
233 status_code=response.status_code,
234 )
235 error = ProviderError(
236 f"API request failed with status {response.status_code}",
237 status_code=response.status_code,
238 )
239 except httpx.RequestError as e:
240 error = ProviderError(f"Network error during API call: {str(e)}")
242 if attempt >= self.max_retries:
243 logger.error(f"Provider call failed after {attempt + 1} attempts: {error}")
244 raise error
245 wait = self._backoff(attempt, response)
246 logger.warning(
247 f"Provider call failed ({error}); retry {attempt + 1}/{self.max_retries} "
248 f"in {wait:.2f}s"
249 )
250 self.sleep(wait)
251 attempt += 1
253 def verify_phone(self, phone: str) -> VerificationResult:
254 """
255 Verify phone using landlineremover.com API.
257 This API returns line type, DNC status, and litigator status in a
258 single call, which is more efficient than making separate API calls.
260 Args:
261 phone: E.164 formatted phone number
263 Returns:
264 VerificationResult of (line_type, dnc, known_litigator)
266 Raises:
267 httpx.HTTPError: For API communication errors
268 ValueError: For invalid responses
269 """
270 logger.debug(f"Verifying phone {phone[:6]}*** via external API")
272 if not self.api_key:
273 # A provider configuration failure, not a data error: callers must not
274 # classify it as an invalid customer phone number.
275 raise ProviderError("API key not configured")
277 self._check_breaker()
279 try:
280 # Single call returns line type, DNC and litigator status; retried on
281 # transient failures inside _request
282 response = self._request(phone)
284 # Parse response
285 json_response = response.json()
287 # Extract data from response wrapper
288 if "data" in json_response:
289 data = json_response["data"]
290 else:
291 data = json_response
293 # Map line type from API response
294 line_type = self._map_line_type(data)
296 # Map DNC status - API uses "DNCType" field
297 # Values can be "dnc", "clean", "litigator", etc.
298 dnc_type = data.get("DNCType", data.get("dnc_type", "")).lower()
299 is_dnc = dnc_type != "clean" and dnc_type != ""
301 # Litigator status is signaled via DNCType rather than a separate
302 # field; substring match tolerates combined variants like "dnc_litigator"
303 known_litigator = "litigator" in dnc_type
305 logger.debug(
306 f"Verification complete for {phone[:6]}***",
307 extra={
308 "line_type": line_type.value,
309 "is_dnc": is_dnc,
310 "dnc_type": dnc_type,
311 "known_litigator": known_litigator,
312 },
313 )
315 self._record_success()
316 return VerificationResult(line_type, is_dnc, known_litigator)
318 except (ProviderQuotaError, ProviderAuthError):
319 # The vendor answered; the account is the problem. Not an outage signal.
320 raise
321 except ProviderError:
322 self._record_failure()
323 raise
324 except Exception as e:
325 logger.error(f"Unexpected error during verification: {str(e)}")
326 self._record_failure()
327 raise
329 def _map_line_type(self, data: dict) -> LineType:
330 """
331 Map API response to LineType enum.
333 Args:
334 data: API response dictionary
336 Returns:
337 LineType enum value
338 """
339 # API uses "LineType" (capitalized) field
340 line_type_str = data.get("LineType", data.get("line_type", "")).lower()
342 # Map common line types
343 line_type_map = {
344 "mobile": LineType.MOBILE,
345 "landline": LineType.LANDLINE,
346 "voip": LineType.VOIP,
347 "wireless": LineType.MOBILE, # Some APIs return "wireless" for mobile
348 "fixed": LineType.LANDLINE, # Some APIs return "fixed" for landline
349 }
351 return line_type_map.get(line_type_str, LineType.UNKNOWN)
353 def __del__(self) -> None:
354 """Cleanup HTTP client"""
355 if hasattr(self, "http_client"):
356 self.http_client.close()