Coverage for core/utils.py: 38%

103 statements  

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

1from django.db import DatabaseError, connection 

2from django.http import Http404 

3from django.utils.translation import gettext_lazy as _ 

4from django_rq.queues import get_queue, get_queue_by_index, get_redis_connection 

5from django_rq.settings import get_queues_list, get_queues_map 

6from django_rq.utils import get_jobs, stop_jobs 

7from rq import requeue_job 

8from rq.exceptions import NoSuchJobError 

9from rq.job import Job as RQ_Job 

10from rq.job import JobStatus as RQJobStatus 

11from rq.registry import ( 

12 DeferredJobRegistry, 

13 FailedJobRegistry, 

14 FinishedJobRegistry, 

15 ScheduledJobRegistry, 

16 StartedJobRegistry, 

17) 

18 

19__all__ = ( 

20 'delete_rq_job', 

21 'enqueue_rq_job', 

22 'get_db_schema', 

23 'get_rq_jobs', 

24 'get_rq_jobs_from_status', 

25 'requeue_rq_job', 

26 'stop_rq_job', 

27) 

28 

29 

30def get_rq_jobs(): 

31 """ 

32 Return a list of all RQ jobs. 

33 """ 

34 jobs = set() 

35 

36 for queue in get_queues_list(): 

37 queue = get_queue(queue['name']) 

38 jobs.update(queue.get_jobs()) 

39 

40 return list(jobs) 

41 

42 

43def get_rq_jobs_from_status(queue, status): 

44 """ 

45 Return the RQ jobs with the given status. 

46 """ 

47 jobs = [] 

48 

49 try: 

50 registry_cls = { 

51 RQJobStatus.STARTED: StartedJobRegistry, 

52 RQJobStatus.DEFERRED: DeferredJobRegistry, 

53 RQJobStatus.FINISHED: FinishedJobRegistry, 

54 RQJobStatus.FAILED: FailedJobRegistry, 

55 RQJobStatus.SCHEDULED: ScheduledJobRegistry, 

56 }[status] 

57 except KeyError: 

58 raise Http404 

59 registry = registry_cls(queue.name, queue.connection) 

60 

61 job_ids = registry.get_job_ids() 

62 if status != RQJobStatus.DEFERRED: 

63 jobs = get_jobs(queue, job_ids, registry) 

64 else: 

65 # Deferred jobs require special handling 

66 for job_id in job_ids: 

67 try: 

68 jobs.append(RQ_Job.fetch(job_id, connection=queue.connection, serializer=queue.serializer)) 

69 except NoSuchJobError: 

70 pass 

71 

72 if jobs and status == RQJobStatus.SCHEDULED: 

73 for job in jobs: 

74 job.scheduled_at = registry.get_scheduled_time(job) 

75 

76 return jobs 

77 

78 

79def delete_rq_job(job_id): 

80 """ 

81 Delete the specified RQ job. 

82 """ 

83 config = get_queues_list()[0] 

84 try: 

85 job = RQ_Job.fetch(job_id, connection=get_redis_connection(config['connection_config']),) 

86 except NoSuchJobError: 

87 raise Http404(_("Job {job_id} not found").format(job_id=job_id)) 

88 

89 queue_index = get_queues_map()[job.origin] 

90 queue = get_queue_by_index(queue_index) 

91 

92 # Remove job id from queue and delete the actual job 

93 queue.connection.lrem(queue.key, 0, job.id) 

94 job.delete() 

95 

96 

97def requeue_rq_job(job_id): 

98 """ 

99 Requeue the specified RQ job. 

100 """ 

101 config = get_queues_list()[0] 

102 try: 

103 job = RQ_Job.fetch(job_id, connection=get_redis_connection(config['connection_config']),) 

104 except NoSuchJobError: 

105 raise Http404(_("Job {id} not found.").format(id=job_id)) 

106 

107 queue_index = get_queues_map()[job.origin] 

108 queue = get_queue_by_index(queue_index) 

109 

110 requeue_job(job_id, connection=queue.connection, serializer=queue.serializer) 

111 

112 

113def enqueue_rq_job(job_id): 

114 """ 

115 Enqueue the specified RQ job. 

116 """ 

117 config = get_queues_list()[0] 

118 try: 

119 job = RQ_Job.fetch(job_id, connection=get_redis_connection(config['connection_config']),) 

120 except NoSuchJobError: 

121 raise Http404(_("Job {id} not found.").format(id=job_id)) 

122 

123 queue_index = get_queues_map()[job.origin] 

124 queue = get_queue_by_index(queue_index) 

125 

126 try: 

127 # _enqueue_job is new in RQ 1.14, this is used to enqueue 

128 # job regardless of its dependencies 

129 queue._enqueue_job(job) 

130 except AttributeError: 

131 queue.enqueue_job(job) 

132 

133 # Remove job from correct registry if needed 

134 if job.get_status() == RQJobStatus.DEFERRED: 

135 registry = DeferredJobRegistry(queue.name, queue.connection) 

136 registry.remove(job) 

137 elif job.get_status() == RQJobStatus.FINISHED: 

138 registry = FinishedJobRegistry(queue.name, queue.connection) 

139 registry.remove(job) 

140 elif job.get_status() == RQJobStatus.SCHEDULED: 

141 registry = ScheduledJobRegistry(queue.name, queue.connection) 

142 registry.remove(job) 

143 

144 

145def stop_rq_job(job_id): 

146 """ 

147 Stop the specified RQ job. 

148 """ 

149 config = get_queues_list()[0] 

150 try: 

151 job = RQ_Job.fetch(job_id, connection=get_redis_connection(config['connection_config']),) 

152 except NoSuchJobError: 

153 raise Http404(_("Job {job_id} not found").format(job_id=job_id)) 

154 

155 queue_index = get_queues_map()[job.origin] 

156 queue = get_queue_by_index(queue_index) 

157 

158 return stop_jobs(queue, job_id)[0] 

159 

160 

161def get_db_schema(): 

162 """ 

163 Query the current PostgreSQL schema and return a list of tables, each with its columns and 

164 indexes. Returns an empty list if the database is not accessible. 

165 """ 

166 db_schema = [] 

167 try: 

168 with connection.cursor() as cursor: 

169 cursor.execute(""" 

170 SELECT table_name, column_name, data_type, is_nullable, column_default 

171 FROM information_schema.columns 

172 WHERE table_schema = current_schema() 

173 ORDER BY table_name, ordinal_position 

174 """) 

175 columns_by_table = {} 

176 for table_name, column_name, data_type, is_nullable, column_default in cursor.fetchall(): 

177 columns_by_table.setdefault(table_name, []).append({ 

178 'name': column_name, 

179 'type': data_type, 

180 'nullable': is_nullable == 'YES', 

181 'default': column_default, 

182 }) 

183 

184 cursor.execute(""" 

185 SELECT tablename, indexname, indexdef 

186 FROM pg_indexes 

187 WHERE schemaname = current_schema() 

188 ORDER BY tablename, indexname 

189 """) 

190 indexes_by_table = {} 

191 for table_name, index_name, index_def in cursor.fetchall(): 

192 indexes_by_table.setdefault(table_name, []).append({ 

193 'name': index_name, 

194 'definition': index_def, 

195 }) 

196 

197 for table_name in sorted(columns_by_table.keys()): 

198 db_schema.append({ 

199 'name': table_name, 

200 'columns': columns_by_table[table_name], 

201 'indexes': indexes_by_table.get(table_name, []), 

202 }) 

203 except DatabaseError: 

204 pass 

205 return db_schema