Coverage for api/utils/check_dead_links/__init__.py: 85%
101 statements
« prev ^ index » next coverage.py v7.15.2, created at 2026-10-07 06:14 +0000
« prev ^ index » next coverage.py v7.15.2, created at 2026-10-07 06:14 +0000
1import asyncio
2import time
4from django.conf import settings
6import aiohttp
7import django_redis
8import structlog
9from asgiref.sync import async_to_sync
10from decouple import config
11from elasticsearch_dsl.response import Hit
12from redis.exceptions import ConnectionError
14from api.utils.aiohttp import get_aiohttp_session
15from api.utils.check_dead_links.provider_status_mappings import provider_status_mappings
16from api.utils.dead_link_mask import get_query_mask, save_query_mask
19logger = structlog.get_logger(__name__)
21CACHE_PREFIX = "valid:"
22HEADERS = {
23 "User-Agent": settings.OUTBOUND_USER_AGENT_TEMPLATE.format(purpose="LinkValidation")
24}
27def _get_cached_statuses(redis, urls):
28 try:
29 cached_statuses = redis.mget([CACHE_PREFIX + url for url in urls])
30 return [
31 int(b.decode("utf-8")) if b is not None else None for b in cached_statuses
32 ]
33 except ConnectionError:
34 logger.warning("Redis connect failed, validating all URLs without cache.")
35 return [None] * len(urls)
38def _get_expiry(status, default):
39 return config(f"LINK_VALIDATION_CACHE_EXPIRY__{status}", default=default, cast=int)
42_timeout = aiohttp.ClientTimeout(total=settings.LINK_VALIDATION_TIMEOUT_SECONDS)
44_ERROR_STATUS = -1
47# Used to filter network errors during liveness checks that we believe
48# we should get errors (rather than warnings) for, e.g., a Sentry issue.
49# As such, this should exclude errors where we don't think we have the
50# ability to control the outcome, whether now or in the future.
51# Some things, like SSL errors and timeouts, may be used for detecting
52# more confident liveness checks or useful catalog information in the future.
53# However, for now, we just need to avoid filling up our error backlog
54# with things we aren't able to actually do anything about right now.
55_NON_ACTIONABLE_NETWORK_EXCEPTIONS = (
56 aiohttp.ClientConnectorCertificateError,
57 aiohttp.ClientConnectorSSLError,
58 aiohttp.ClientConnectorError,
59 aiohttp.ServerDisconnectedError,
60 aiohttp.ClientOSError,
61)
64async def _head(
65 url: str, session: aiohttp.ClientSession, provider: str
66) -> tuple[str, int]:
67 try:
68 async with session.head(
69 url,
70 allow_redirects=False,
71 headers=HEADERS,
72 timeout=_timeout,
73 # do not raise_for_status=True; we want "bad" status codes
74 # so we can cache on them and interpret them per provider
75 # the except block should only handle errors in making or
76 # receiving the response, but not the content of the response
77 # e.g., a 500 from upstream is good information for dead link
78 # validation. An issue with timeouts or being able to actually
79 # access the upstream over the wire are problems we want in
80 # the log and for which we need to fall back to the _ERROR_STATUS
81 # because we won't have a true status from upstream.
82 trace_request_ctx={
83 "timing_event_name": "dead_link_validation_timing",
84 "timing_event_ctx": {
85 "provider": provider,
86 },
87 },
88 ) as response:
89 status = response.status
91 except (aiohttp.ClientError, asyncio.TimeoutError) as exception:
92 # only log non-timeout exceptions. Timeouts are more or less expected
93 # and we have means of visibility into them (e.g., querying for timings
94 # that exceed the timeout); they do _not_ need to be in the error log or go to Sentry
95 # The timeout exception class aiohttp uses is a subclass of `ClientError` and `asyncio.TimeoutError`,
96 # so we have to explicitly check that the error is _not_ an asyncio timeout error, not that it
97 # _is_ an aiohttp error. When it is a timeout error, it will be both `asyncio.TimeoutError` and
98 # `aiohttp.ClientError` because it's a subclass of both
99 if not isinstance(exception, asyncio.TimeoutError):
100 if isinstance(exception, _NON_ACTIONABLE_NETWORK_EXCEPTIONS):
101 log = logger.warning
102 else:
103 # Only error log actionable exceptions
104 log = logger.error
106 log("dead_link_validation_error", exc_info=True, exc=exception)
108 status = _ERROR_STATUS
110 return url, status
113# https://stackoverflow.com/q/55259755
114@async_to_sync
115async def _make_head_requests(
116 urls: dict[str, int], results: list[Hit]
117) -> list[tuple[str, int]]:
118 """
119 Concurrently HEAD request the urls.
121 ``urls`` must map to the index of the corresponding result in ``results``.
123 :param urls: A dictionary with keys of the URLs to request, mapped to the index of that url in ``results``
124 :param results: The ordered list of results, including ones not being validated.
125 """
126 session = await get_aiohttp_session()
127 tasks = [
128 asyncio.ensure_future(_head(url, session, results[idx].provider))
129 for url, idx in urls.items()
130 ]
131 responses = asyncio.gather(*tasks)
132 await responses
133 return responses.result()
136def check_dead_links(query_hash: str, start_slice: int, results: list[Hit]) -> None:
137 """
138 Make sure images exist before we display them.
140 Treat redirects as broken links since most of the time the redirect leads to a
141 generic "not found" placeholder.
143 Results are cached in redis and shared amongst all API servers in the
144 cluster.
145 """
146 if not results:
147 logger.info("link_validation_empty_results")
148 return
150 urls = [result.url for result in results]
152 logger.debug("starting validation")
153 start_time = time.time()
155 # Pull matching images from the cache.
156 redis = django_redis.get_redis_connection("default")
157 cached_statuses = _get_cached_statuses(redis, urls)
158 logger.debug(f"len(cached_statuses)={len(cached_statuses)}")
160 # Anything that isn't in the cache needs to be validated via HEAD request.
161 to_verify = {}
162 for idx, url in enumerate(urls):
163 if cached_statuses[idx] is None:
164 to_verify[url] = idx
165 logger.debug(f"len(to_verify)={len(to_verify)}")
167 verified = _make_head_requests(to_verify, results)
169 # Cache newly verified image statuses.
170 to_cache = {CACHE_PREFIX + url: status for url, status in verified}
172 pipe = redis.pipeline()
173 if len(to_cache) > 0:
174 pipe.mset(to_cache)
176 for key, status in to_cache.items():
177 if status == 200:
178 logger.debug(f"healthy link key={key}")
179 elif status == _ERROR_STATUS: 179 ↛ 180line 179 didn't jump to line 180 because the condition on line 179 was never true
180 logger.debug(f"no response from provider key={key}")
181 else:
182 logger.debug(f"broken link key={key}")
184 expiry = settings.LINK_VALIDATION_CACHE_EXPIRY_CONFIGURATION[status]
185 logger.debug(f"caching status={status} expiry={expiry}")
186 pipe.expire(key, expiry)
188 try:
189 pipe.execute()
190 except ConnectionError:
191 logger.warning("Redis connect failed, cannot cache link liveness.")
193 # Merge newly verified results with cached statuses
194 for idx, url in enumerate(to_verify):
195 cache_idx = to_verify[url]
196 cached_statuses[cache_idx] = verified[idx][1]
198 # Create a new dead link mask
199 new_mask = [1] * len(results)
201 # Delete broken images from the search results response.
202 for idx, _ in enumerate(cached_statuses):
203 del_idx = len(cached_statuses) - idx - 1
204 status = cached_statuses[del_idx]
206 provider = results[del_idx]["provider"]
207 status_mapping = provider_status_mappings[provider]
209 if status in status_mapping.unknown:
210 logger.warning(
211 "Image validation failed due to rate limiting or blocking. "
212 f"url={urls[idx]} "
213 f"status={status} "
214 f"provider={provider} "
215 )
216 elif status not in status_mapping.live:
217 logger.info(
218 "Deleting broken image from results "
219 f"id={results[del_idx]['identifier']} "
220 f"status={status} "
221 f"provider={provider} "
222 )
223 # remove the result, mutating in place
224 del results[del_idx]
225 # update the result's position in the mask to indicate it is dead
226 new_mask[del_idx] = 0
228 # Merge and cache the new mask
229 mask = get_query_mask(query_hash)
230 if mask:
231 # skip the leading part of the mask that represents results that come before
232 # the results we've verified this time around. Overwrite everything after
233 # with our new results validation mask.
234 new_mask = mask[:start_slice] + new_mask
235 save_query_mask(query_hash, new_mask)
237 end_time = time.time()
238 logger.debug(
239 "end validation "
240 f"end_time={end_time} "
241 f"start_time={start_time} "
242 f"delta={end_time - start_time} "
243 )