Coverage for chalicelib/core/jobs.py: 28%
54 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
3from chalicelib.core.sessions import sessions_mobs, sessions_devtool
4from chalicelib.utils import pg_client, helper
6logger = logging.getLogger(__name__)
9class Actions:
10 DELETE_USER_DATA = "delete_user_data"
13class JobStatus:
14 SCHEDULED = "scheduled"
15 COMPLETED = "completed"
16 FAILED = "failed"
17 CANCELLED = "cancelled"
20def update(job_id, job):
21 with pg_client.PostgresClient() as cur:
22 job_data = {
23 "job_id": job_id,
24 "errors": job.get("errors"),
25 **job
26 }
28 query = cur.mogrify(
29 """UPDATE public.jobs
30 SET updated_at = timezone('utc'::text, now()),
31 status = %(status)s,
32 errors = %(errors)s
33 WHERE job_id = %(job_id)s;""",
34 job_data)
36 cur.execute(query=query)
39def __get_session_ids_by_user_ids(project_id, user_ids):
40 with pg_client.PostgresClient() as cur:
41 query = cur.mogrify(
42 """SELECT session_id
43 FROM public.sessions
44 WHERE project_id = %(project_id)s
45 AND user_id IN %(userId)s LIMIT 1000;""",
46 {"project_id": project_id, "userId": tuple(user_ids)})
47 cur.execute(query=query)
48 ids = cur.fetchall()
49 return [s["session_id"] for s in ids]
52def __delete_sessions_by_session_ids(session_ids):
53 with pg_client.PostgresClient(unlimited_query=True) as cur:
54 query = cur.mogrify(
55 """DELETE
56 FROM public.sessions
57 WHERE session_id IN %(session_ids)s""",
58 {"session_ids": tuple(session_ids)}
59 )
60 cur.execute(query=query)
63def __delete_session_mobs_by_session_ids(session_ids, project_id):
64 sessions_mobs.delete_mobs(session_ids=session_ids, project_id=project_id)
65 sessions_devtool.delete_mobs(session_ids=session_ids, project_id=project_id)
68def get_scheduled_jobs():
69 with pg_client.PostgresClient() as cur:
70 query = cur.mogrify(
71 """SELECT *
72 FROM public.jobs
73 WHERE status = %(status)s
74 AND start_at <= (now() at time zone 'utc');""",
75 {"status": JobStatus.SCHEDULED})
76 cur.execute(query=query)
77 data = cur.fetchall()
78 return helper.list_to_camel_case(data)
81def execute_jobs():
82 jobs = get_scheduled_jobs()
83 for job in jobs:
84 logger.info(f"Executing jobId:{job['jobId']}")
85 try:
86 if job["action"] == Actions.DELETE_USER_DATA:
87 session_ids = __get_session_ids_by_user_ids(project_id=job["projectId"],
88 user_ids=[job["referenceId"]])
89 if len(session_ids) > 0:
90 logger.info(f"Deleting {len(session_ids)} sessions")
91 __delete_sessions_by_session_ids(session_ids=session_ids)
92 __delete_session_mobs_by_session_ids(session_ids=session_ids, project_id=job["projectId"])
93 else:
94 raise Exception(f"The action '{job['action']}' not supported.")
96 job["status"] = JobStatus.COMPLETED
97 logger.info(f"Job completed {job['jobId']}")
98 except Exception as e:
99 job["status"] = JobStatus.FAILED
100 job["errors"] = str(e)
101 logger.error(f"Job failed {job['jobId']}")
103 update(job["jobId"], job)