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

1""" 

2External API provider for production phone verification 

3""" 

4 

5import os 

6import random 

7import threading 

8import time 

9from collections.abc import Callable 

10 

11import httpx 

12from aws_lambda_powertools import Logger 

13 

14from ..core.models import LineType, VerificationResult 

15from .exceptions import ( 

16 ProviderAuthError, 

17 ProviderError, 

18 ProviderQuotaError, 

19 ProviderUnavailableError, 

20) 

21 

22logger = Logger() 

23 

24RETRYABLE_STATUSES = frozenset({408, 429, 500, 502, 503, 504}) 

25RETRY_AFTER_CAP_SECONDS = 60.0 

26 

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

32 

33 

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 

37 

38 

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 

42 

43 

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

49 

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. 

65 

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 

72 

73 The limiter and breaker are process-local: they bound one Lambda invocation or 

74 worker pool, not the fleet. 

75 

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

90 

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 ) 

115 

116 # Injectable for tests; never patch time globally 

117 self.sleep: Callable[[float], None] = time.sleep 

118 self.clock: Callable[[], float] = time.monotonic 

119 

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 

126 

127 # ---- rate limiter ------------------------------------------------------------- 

128 

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) 

147 

148 # ---- circuit breaker ---------------------------------------------------------- 

149 

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 ) 

163 

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 

169 

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() 

181 

182 # ---- retry -------------------------------------------------------------------- 

183 

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 

195 

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) 

211 

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

241 

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 

252 

253 def verify_phone(self, phone: str) -> VerificationResult: 

254 """ 

255 Verify phone using landlineremover.com API. 

256 

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. 

259 

260 Args: 

261 phone: E.164 formatted phone number 

262 

263 Returns: 

264 VerificationResult of (line_type, dnc, known_litigator) 

265 

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

271 

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

276 

277 self._check_breaker() 

278 

279 try: 

280 # Single call returns line type, DNC and litigator status; retried on 

281 # transient failures inside _request 

282 response = self._request(phone) 

283 

284 # Parse response 

285 json_response = response.json() 

286 

287 # Extract data from response wrapper 

288 if "data" in json_response: 

289 data = json_response["data"] 

290 else: 

291 data = json_response 

292 

293 # Map line type from API response 

294 line_type = self._map_line_type(data) 

295 

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

300 

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 

304 

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 ) 

314 

315 self._record_success() 

316 return VerificationResult(line_type, is_dnc, known_litigator) 

317 

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 

328 

329 def _map_line_type(self, data: dict) -> LineType: 

330 """ 

331 Map API response to LineType enum. 

332 

333 Args: 

334 data: API response dictionary 

335 

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() 

341 

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 } 

350 

351 return line_type_map.get(line_type_str, LineType.UNKNOWN) 

352 

353 def __del__(self) -> None: 

354 """Cleanup HTTP client""" 

355 if hasattr(self, "http_client"): 

356 self.http_client.close()