Coverage for chalicelib/core/health.py: 68%

168 statements  

« prev     ^ index     » next       coverage.py v7.15.2, created at 2026-10-10 12:56 +0000

1import logging 

2 

3import redis 

4import requests 

5from decouple import config 

6 

7from chalicelib.utils import pg_client 

8from chalicelib.utils.TimeUTC import TimeUTC 

9from chalicelib.utils.log import sanitize 

10 

11logger = logging.getLogger(__name__) 

12 

13 

14def app_connection_string(name, port, path): 

15 namespace = config("POD_NAMESPACE", default="app") 

16 conn_string = config("CLUSTER_URL", default="svc.cluster.local") 

17 return f"http://{name}.{namespace}.{conn_string}:{port}/{path}" 

18 

19 

20HEALTH_ENDPOINTS = { 

21 "alerts": app_connection_string("alerts-openreplay", 8888, "health"), 

22 "assets": app_connection_string("assets-openreplay", 8888, "metrics"), 

23 "assist": app_connection_string("assist-openreplay", 8888, "health"), 

24 "chalice": app_connection_string("chalice-openreplay", 8888, "metrics"), 

25 "db": app_connection_string("db-openreplay", 8888, "metrics"), 

26 "ender": app_connection_string("ender-openreplay", 8888, "metrics"), 

27 "http": app_connection_string("http-openreplay", 8888, "metrics"), 

28 "ingress-nginx": app_connection_string("ingress-nginx-openreplay", 80, "healthz"), 

29 "sink": app_connection_string("sink-openreplay", 8888, "metrics"), 

30 "sourcemapreader": app_connection_string( 

31 "sourcemapreader-openreplay", 8888, "health" 

32 ), 

33 "storage": app_connection_string("storage-openreplay", 8888, "metrics"), 

34} 

35 

36 

37def __check_database_pg(*_): 

38 logger.info("__check_database_pg: start") 

39 fail_response = { 

40 "health": False, 

41 "details": {"errors": ["Postgres health-check failed"]}, 

42 } 

43 with pg_client.PostgresClient() as cur: 

44 try: 

45 cur.execute("SHOW server_version;") 

46 # server_version = cur.fetchone() 

47 except Exception as e: 

48 logger.error("!! health failed: postgres not responding") 

49 logger.exception(e) 

50 logger.info("__check_database_pg: end --failed") 

51 return fail_response 

52 try: 

53 cur.execute("SELECT openreplay_version() AS version;") 

54 # schema_version = cur.fetchone() 

55 except Exception as e: 

56 logger.error("!! health failed: openreplay_version not defined") 

57 logger.exception(e) 

58 logger.info("__check_database_pg: end --failed") 

59 return fail_response 

60 logger.info("__check_database_pg: end") 

61 return { 

62 "health": True, 

63 "details": { 

64 # "version": server_version["server_version"], 

65 # "schema": schema_version["version"] 

66 }, 

67 } 

68 

69 

70def __always_healthy(*_): 

71 return {"health": True, "details": {}} 

72 

73 

74def __check_be_service(service_name): 

75 def fn(*_): 

76 logger.info("__check_be_service(%s): start", service_name) 

77 fail_response = { 

78 "health": False, 

79 "details": {"errors": ["server health-check failed"]}, 

80 } 

81 try: 

82 results = requests.get(HEALTH_ENDPOINTS.get(service_name), timeout=2) 

83 if results.status_code != 200: 

84 logger.error( 

85 f"!! issue with the {service_name}-health code:{results.status_code}" 

86 ) 

87 logger.error(sanitize(results.text)) 

88 # fail_response["details"]["errors"].append(results.text) 

89 logger.info("__check_be_service(%s): end", service_name) 

90 return fail_response 

91 except requests.exceptions.Timeout: 

92 logger.error(f"!! Timeout getting {service_name}-health") 

93 # fail_response["details"]["errors"].append("timeout") 

94 logger.info("__check_be_service(%s): end --failed", service_name) 

95 return fail_response 

96 except Exception as e: 

97 logger.error(f"!! Issue getting {service_name}-health response") 

98 logger.exception(e) 

99 try: 

100 logger.error(sanitize(results.text)) 

101 # fail_response["details"]["errors"].append(results.text) 

102 except Exception: 

103 logger.error("couldn't get response") 

104 # fail_response["details"]["errors"].append(str(e)) 

105 logger.info("__check_be_service(%s): end --failed", service_name) 

106 return fail_response 

107 logger.info("__check_be_service(%s): end", service_name) 

108 return {"health": True, "details": {}} 

109 

110 return fn 

111 

112 

113def __check_redis(*_): 

114 logger.info("__check_redis: start") 

115 fail_response = { 

116 "health": False, 

117 "details": {"errors": ["server health-check failed"]}, 

118 } 

119 if config("REDIS_STRING", default=None) is None: 119 ↛ 121line 119 didn't jump to line 121 because the condition on line 119 was never true

120 # fail_response["details"]["errors"].append("REDIS_STRING not defined in env-vars") 

121 logger.info("__check_redis: end --failed") 

122 return fail_response 

123 

124 try: 

125 r = redis.from_url(config("REDIS_STRING"), socket_timeout=2) 

126 r.ping() 

127 except Exception as e: 

128 logger.error("!! Issue getting redis-health response") 

129 logger.exception(e) 

130 # fail_response["details"]["errors"].append(str(e)) 

131 logger.info("__check_redis: end --failed") 

132 return fail_response 

133 

134 logger.info("__check_redis: end") 

135 return { 

136 "health": True, 

137 "details": { 

138 # "version": r.execute_command('INFO')['redis_version'] 

139 }, 

140 } 

141 

142 

143def __check_SSL(*_): 

144 logger.info("__check_SSL: start") 

145 fail_response = { 

146 "health": False, 

147 "details": {"errors": ["SSL Certificate health-check failed"]}, 

148 } 

149 try: 

150 timeout = config("SSL_VERIFICATION_TIMEOUT_S", cast=int, default=10) 

151 requests.get(config("SITE_URL"), verify=True, allow_redirects=True, timeout=timeout) 

152 except Exception as e: 

153 logger.error("!! health failed: SSL Certificate") 

154 logger.exception(e) 

155 logger.info("__check_SSL: end --failed") 

156 return fail_response 

157 logger.info("__check_SSL: end") 

158 return {"health": True, "details": {}} 

159 

160 

161def __get_sessions_stats(*_): 

162 logger.info("__get_sessions_stats: start") 

163 with pg_client.PostgresClient() as cur: 

164 constraints = ["projects.deleted_at IS NULL"] 

165 query = cur.mogrify( 

166 f"""SELECT COALESCE(SUM(sessions_count),0) AS s_c, 

167 COALESCE(SUM(events_count),0) AS e_c 

168 FROM public.projects_stats 

169 INNER JOIN public.projects USING(project_id) 

170 WHERE {" AND ".join(constraints)};""" 

171 ) 

172 cur.execute(query) 

173 row = cur.fetchone() 

174 logger.info("__get_sessions_stats: end") 

175 return {"numberOfSessionsCaptured": row["s_c"], "numberOfEventCaptured": row["e_c"]} 

176 

177 

178def get_health(tenant_id=None): 

179 logger.info("get_health: start") 

180 health_map = { 

181 "databases": {"postgres": __check_database_pg}, 

182 "ingestionPipeline": {"redis": __check_redis}, 

183 "backendServices": { 

184 "alerts": __check_be_service("alerts"), 

185 "assets": __check_be_service("assets"), 

186 "assist": __check_be_service("assist"), 

187 "chalice": __always_healthy, 

188 "db": __check_be_service("db"), 

189 "ender": __check_be_service("ender"), 

190 "frontend": __always_healthy, 

191 "http": __check_be_service("http"), 

192 "ingress-nginx": __always_healthy, 

193 "sink": __check_be_service("sink"), 

194 "sourcemapreader": __check_be_service("sourcemapreader"), 

195 "storage": __check_be_service("storage"), 

196 }, 

197 "details": __get_sessions_stats, 

198 "ssl": __check_SSL, 

199 } 

200 result = __process_health(health_map=health_map) 

201 logger.info("get_health: end") 

202 return result 

203 

204 

205def __process_health(health_map): 

206 logger.info("__process_health: start") 

207 response = dict(health_map) 

208 for parent_key in health_map.keys(): 

209 if config(f"SKIP_H_{parent_key.upper()}", cast=bool, default=False): 209 ↛ 210line 209 didn't jump to line 210 because the condition on line 209 was never true

210 response.pop(parent_key) 

211 elif isinstance(health_map[parent_key], dict): 

212 for element_key in health_map[parent_key]: 

213 if config( 213 ↛ 218line 213 didn't jump to line 218 because the condition on line 213 was never true

214 f"SKIP_H_{parent_key.upper()}_{element_key.upper()}", 

215 cast=bool, 

216 default=False, 

217 ): 

218 response[parent_key].pop(element_key) 

219 else: 

220 response[parent_key][element_key] = health_map[parent_key][ 

221 element_key 

222 ]() 

223 else: 

224 response[parent_key] = health_map[parent_key]() 

225 logger.info("__process_health: end") 

226 return response 

227 

228 

229def cron(): 

230 logger.info("cron: start") 

231 with pg_client.PostgresClient() as cur: 

232 query = cur.mogrify( 

233 """SELECT projects.project_id, 

234 projects.created_at, 

235 projects.sessions_last_check_at, 

236 projects.first_recorded_session_at, 

237 projects_stats.last_update_at 

238 FROM public.projects 

239 LEFT JOIN public.projects_stats USING (project_id) 

240 WHERE projects.deleted_at IS NULL 

241 ORDER BY project_id;""" 

242 ) 

243 cur.execute(query) 

244 rows = cur.fetchall() 

245 for r in rows: 

246 insert = False 

247 if r["last_update_at"] is None: 

248 # never counted before, must insert 

249 insert = True 

250 if r["first_recorded_session_at"] is None: 250 ↛ 256line 250 didn't jump to line 256 because the condition on line 250 was always true

251 if r["sessions_last_check_at"] is None: 251 ↛ 254line 251 didn't jump to line 254 because the condition on line 251 was always true

252 count_start_from = r["created_at"] 

253 else: 

254 count_start_from = r["sessions_last_check_at"] 

255 else: 

256 count_start_from = r["first_recorded_session_at"] 

257 

258 else: 

259 # counted before, must update 

260 count_start_from = r["last_update_at"] 

261 

262 count_start_from = TimeUTC.datetime_to_timestamp(count_start_from) 

263 params = { 

264 "project_id": r["project_id"], 

265 "start_ts": count_start_from, 

266 "end_ts": TimeUTC.now(), 

267 "sessions_count": 0, 

268 "events_count": 0, 

269 } 

270 

271 query = cur.mogrify( 

272 """SELECT COUNT(1) AS sessions_count, 

273 COALESCE(SUM(events_count),0) AS events_count 

274 FROM public.sessions 

275 WHERE project_id=%(project_id)s 

276 AND start_ts>=%(start_ts)s 

277 AND start_ts<=%(end_ts)s 

278 AND duration IS NOT NULL;""", 

279 params, 

280 ) 

281 cur.execute(query) 

282 row = cur.fetchone() 

283 if row is not None: 283 ↛ 287line 283 didn't jump to line 287 because the condition on line 283 was always true

284 params["sessions_count"] = row["sessions_count"] 

285 params["events_count"] = row["events_count"] 

286 

287 if insert: 

288 query = cur.mogrify( 

289 """INSERT INTO public.projects_stats(project_id, sessions_count, events_count, last_update_at) 

290 VALUES (%(project_id)s, %(sessions_count)s, %(events_count)s, (now() AT TIME ZONE 'utc'::text));""", 

291 params, 

292 ) 

293 else: 

294 query = cur.mogrify( 

295 """UPDATE public.projects_stats 

296 SET sessions_count=sessions_count+%(sessions_count)s, 

297 events_count=events_count+%(events_count)s, 

298 last_update_at=(now() AT TIME ZONE 'utc'::text) 

299 WHERE project_id=%(project_id)s;""", 

300 params, 

301 ) 

302 cur.execute(query) 

303 logger.info("cron: end") 

304 

305 

306# this cron is used to correct the sessions&events count every week 

307def weekly_cron(): 

308 logger.info("weekly_cron: start") 

309 with pg_client.PostgresClient(long_query=True) as cur: 

310 query = cur.mogrify( 

311 """SELECT project_id, 

312 projects_stats.last_update_at 

313 FROM public.projects 

314 LEFT JOIN public.projects_stats USING (project_id) 

315 WHERE projects.deleted_at IS NULL 

316 ORDER BY project_id;""" 

317 ) 

318 cur.execute(query) 

319 rows = cur.fetchall() 

320 for r in rows: 

321 if r["last_update_at"] is None: 

322 continue 

323 

324 params = { 

325 "project_id": r["project_id"], 

326 "end_ts": TimeUTC.now(), 

327 "sessions_count": 0, 

328 "events_count": 0, 

329 } 

330 

331 query = cur.mogrify( 

332 """SELECT COUNT(1) AS sessions_count, 

333 COALESCE(SUM(events_count),0) AS events_count 

334 FROM public.sessions 

335 WHERE project_id=%(project_id)s 

336 AND start_ts<=%(end_ts)s 

337 AND duration IS NOT NULL;""", 

338 params, 

339 ) 

340 cur.execute(query) 

341 row = cur.fetchone() 

342 if row is not None: 

343 params["sessions_count"] = row["sessions_count"] 

344 params["events_count"] = row["events_count"] 

345 

346 query = cur.mogrify( 

347 """UPDATE public.projects_stats 

348 SET sessions_count=%(sessions_count)s, 

349 events_count=%(events_count)s, 

350 last_update_at=(now() AT TIME ZONE 'utc'::text) 

351 WHERE project_id=%(project_id)s;""", 

352 params, 

353 ) 

354 cur.execute(query) 

355 logger.info("weekly_cron: end")