1import asyncio
2from datetime import datetime
3from typing import Any
4
5import httpx
6from apify import Actor
7
8
9BASE_URL = "https://api.nhtsa.gov/recalls"
10
11
12def _as_int(value: Any) -> int | None:
13 try:
14 return int(value) if value not in (None, "") else None
15 except (TypeError, ValueError):
16 return None
17
18
19def _date(value: Any) -> str | None:
20 if not value:
21 return None
22 text = str(value).strip()
23 for pattern in ("%d/%m/%Y", "%m/%d/%Y", "%Y-%m-%d"):
24 try:
25 return datetime.strptime(text, pattern).date().isoformat()
26 except ValueError:
27 pass
28 return text
29
30
31async def _get(client: httpx.AsyncClient, url: str, params: dict[str, Any]) -> httpx.Response:
32 for attempt in range(4):
33 try:
34 response = await client.get(url, params=params)
35 except httpx.TransportError:
36 if attempt == 3:
37 raise
38 else:
39 if response.status_code not in {429, 500, 502, 503, 504}:
40 break
41 if attempt == 3:
42 break
43 await asyncio.sleep(2**attempt)
44 if response.status_code >= 400:
45 raise RuntimeError(
46 f"NHTSA API request failed with HTTP {response.status_code}: "
47 f"{response.text[:500]}"
48 )
49 return response
50
51
52def flatten(record: dict[str, Any], search_mode: str, total: int, api_url: str) -> dict[str, Any]:
53 return {
54 "manufacturer": record.get("Manufacturer"),
55 "campaign_number": record.get("NHTSACampaignNumber"),
56 "report_received_date": _date(record.get("ReportReceivedDate")),
57 "component": record.get("Component"),
58 "potential_units_affected": _as_int(record.get("PotentialNumberofUnitsAffected")),
59 "summary": record.get("Summary"),
60 "consequence": record.get("Consequence"),
61 "remedy": record.get("Remedy"),
62 "notes": record.get("Notes"),
63 "model_year": _as_int(record.get("ModelYear")),
64 "make": record.get("Make"),
65 "model": record.get("Model"),
66 "park_it": record.get("parkIt"),
67 "park_outside": record.get("parkOutSide"),
68 "over_the_air_update": record.get("overTheAirUpdate"),
69 "search_mode": search_mode,
70 "total_matches": total,
71 "api_url": api_url,
72 }
73
74
75async def main() -> None:
76 async with Actor:
77 actor_input = await Actor.get_input() or {}
78 max_items = int(actor_input.get("max_items", 10))
79 campaign_number = str(actor_input.get("campaign_number") or "").strip()
80 if campaign_number:
81 search_mode = "campaign"
82 url = f"{BASE_URL}/campaignNumber"
83 params: dict[str, Any] = {"campaignNumber": campaign_number}
84 else:
85 make = str(actor_input.get("make") or "").strip()
86 model = str(actor_input.get("model") or "").strip()
87 model_year = actor_input.get("model_year")
88 if not make or not model or not model_year:
89 raise ValueError(
90 "Provide a campaign_number, or provide make, model, and model_year together."
91 )
92 search_mode = "vehicle"
93 url = f"{BASE_URL}/recallsByVehicle"
94 params = {"make": make, "model": model, "modelYear": int(model_year)}
95
96 async with httpx.AsyncClient(
97 timeout=httpx.Timeout(45.0),
98 headers={"User-Agent": "Apify NHTSA Vehicle Recall Search"},
99 ) as client:
100 response = await _get(client, url, params)
101 payload = response.json()
102 records = payload.get("results")
103 if not isinstance(records, list):
104 raise RuntimeError("NHTSA returned an unexpected results payload.")
105 total = int(payload.get("Count") or len(records))
106 rows = [
107 flatten(record, search_mode, total, str(response.request.url))
108 for record in records[:max_items]
109 if isinstance(record, dict)
110 ]
111 if rows:
112 await Actor.push_data(rows)
113 Actor.log.info("Saved %s of %s matching NHTSA recall rows.", len(rows), total)
114
115
116if __name__ == "__main__":
117 asyncio.run(main())