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

1import asyncio 

2import time 

3 

4from django.conf import settings 

5 

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 

13 

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 

17 

18 

19logger = structlog.get_logger(__name__) 

20 

21CACHE_PREFIX = "valid:" 

22HEADERS = { 

23 "User-Agent": settings.OUTBOUND_USER_AGENT_TEMPLATE.format(purpose="LinkValidation") 

24} 

25 

26 

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) 

36 

37 

38def _get_expiry(status, default): 

39 return config(f"LINK_VALIDATION_CACHE_EXPIRY__{status}", default=default, cast=int) 

40 

41 

42_timeout = aiohttp.ClientTimeout(total=settings.LINK_VALIDATION_TIMEOUT_SECONDS) 

43 

44_ERROR_STATUS = -1 

45 

46 

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) 

62 

63 

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 

90 

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 

105 

106 log("dead_link_validation_error", exc_info=True, exc=exception) 

107 

108 status = _ERROR_STATUS 

109 

110 return url, status 

111 

112 

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. 

120 

121 ``urls`` must map to the index of the corresponding result in ``results``. 

122 

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

134 

135 

136def check_dead_links(query_hash: str, start_slice: int, results: list[Hit]) -> None: 

137 """ 

138 Make sure images exist before we display them. 

139 

140 Treat redirects as broken links since most of the time the redirect leads to a 

141 generic "not found" placeholder. 

142 

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 

149 

150 urls = [result.url for result in results] 

151 

152 logger.debug("starting validation") 

153 start_time = time.time() 

154 

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

159 

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

166 

167 verified = _make_head_requests(to_verify, results) 

168 

169 # Cache newly verified image statuses. 

170 to_cache = {CACHE_PREFIX + url: status for url, status in verified} 

171 

172 pipe = redis.pipeline() 

173 if len(to_cache) > 0: 

174 pipe.mset(to_cache) 

175 

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

183 

184 expiry = settings.LINK_VALIDATION_CACHE_EXPIRY_CONFIGURATION[status] 

185 logger.debug(f"caching status={status} expiry={expiry}") 

186 pipe.expire(key, expiry) 

187 

188 try: 

189 pipe.execute() 

190 except ConnectionError: 

191 logger.warning("Redis connect failed, cannot cache link liveness.") 

192 

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] 

197 

198 # Create a new dead link mask 

199 new_mask = [1] * len(results) 

200 

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] 

205 

206 provider = results[del_idx]["provider"] 

207 status_mapping = provider_status_mappings[provider] 

208 

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 

227 

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) 

236 

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 )