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
« prev ^ index » next coverage.py v7.15.2, created at 2026-10-10 12:56 +0000
1import logging
3import redis
4import requests
5from decouple import config
7from chalicelib.utils import pg_client
8from chalicelib.utils.TimeUTC import TimeUTC
9from chalicelib.utils.log import sanitize
11logger = logging.getLogger(__name__)
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}"
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}
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 }
70def __always_healthy(*_):
71 return {"health": True, "details": {}}
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": {}}
110 return fn
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
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
134 logger.info("__check_redis: end")
135 return {
136 "health": True,
137 "details": {
138 # "version": r.execute_command('INFO')['redis_version']
139 },
140 }
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": {}}
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"]}
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
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
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"]
258 else:
259 # counted before, must update
260 count_start_from = r["last_update_at"]
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 }
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"]
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")
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
324 params = {
325 "project_id": r["project_id"],
326 "end_ts": TimeUTC.now(),
327 "sessions_count": 0,
328 "events_count": 0,
329 }
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"]
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")