Coverage for core/utils.py: 38%
103 statements
« prev ^ index » next coverage.py v7.15.2, created at 2026-10-10 18:35 +0000
« 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)
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)
30def get_rq_jobs():
31 """
32 Return a list of all RQ jobs.
33 """
34 jobs = set()
36 for queue in get_queues_list():
37 queue = get_queue(queue['name'])
38 jobs.update(queue.get_jobs())
40 return list(jobs)
43def get_rq_jobs_from_status(queue, status):
44 """
45 Return the RQ jobs with the given status.
46 """
47 jobs = []
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)
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
72 if jobs and status == RQJobStatus.SCHEDULED:
73 for job in jobs:
74 job.scheduled_at = registry.get_scheduled_time(job)
76 return jobs
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))
89 queue_index = get_queues_map()[job.origin]
90 queue = get_queue_by_index(queue_index)
92 # Remove job id from queue and delete the actual job
93 queue.connection.lrem(queue.key, 0, job.id)
94 job.delete()
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))
107 queue_index = get_queues_map()[job.origin]
108 queue = get_queue_by_index(queue_index)
110 requeue_job(job_id, connection=queue.connection, serializer=queue.serializer)
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))
123 queue_index = get_queues_map()[job.origin]
124 queue = get_queue_by_index(queue_index)
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)
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)
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))
155 queue_index = get_queues_map()[job.origin]
156 queue = get_queue_by_index(queue_index)
158 return stop_jobs(queue, job_id)[0]
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 })
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 })
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