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

1import logging 

2 

3from chalicelib.core.sessions import sessions_mobs, sessions_devtool 

4from chalicelib.utils import pg_client, helper 

5 

6logger = logging.getLogger(__name__) 

7 

8 

9class Actions: 

10 DELETE_USER_DATA = "delete_user_data" 

11 

12 

13class JobStatus: 

14 SCHEDULED = "scheduled" 

15 COMPLETED = "completed" 

16 FAILED = "failed" 

17 CANCELLED = "cancelled" 

18 

19 

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 } 

27 

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) 

35 

36 cur.execute(query=query) 

37 

38 

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] 

50 

51 

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) 

61 

62 

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) 

66 

67 

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) 

79 

80 

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

95 

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

102 

103 update(job["jobId"], job)