Coverage for core/models/jobs.py: 52%
138 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
1import logging
2import uuid
3from dataclasses import asdict
4from functools import partial
6import django_rq
7from django.conf import settings
8from django.contrib.contenttypes.fields import GenericForeignKey
9from django.contrib.postgres.fields import ArrayField
10from django.core.exceptions import ValidationError
11from django.core.serializers.json import DjangoJSONEncoder
12from django.core.validators import MinValueValidator
13from django.db import models, transaction
14from django.urls import reverse
15from django.utils import timezone
16from django.utils.translation import gettext as _
17from rq.exceptions import InvalidJobOperation
19from core.choices import JobNotificationChoices, JobStatusChoices
20from core.dataclasses import JobLogEntry
21from core.events import JOB_COMPLETED, JOB_ERRORED, JOB_FAILED
22from core.models import ObjectType
23from core.signals import job_end, job_start
24from extras.models import Notification
25from netbox.models.features import has_feature
26from utilities.json import JobLogDecoder
27from utilities.querysets import RestrictedQuerySet
28from utilities.rqworker import get_queue_for_model
30__all__ = (
31 'Job',
32)
35class Job(models.Model):
36 """
37 Tracks the lifecycle of a job which represents a background task (e.g. the execution of a custom script).
38 """
39 object_type = models.ForeignKey(
40 to='contenttypes.ContentType',
41 related_name='jobs',
42 on_delete=models.CASCADE,
43 blank=True,
44 null=True
45 )
46 object_id = models.PositiveBigIntegerField(
47 blank=True,
48 null=True
49 )
50 object = GenericForeignKey(
51 ct_field='object_type',
52 fk_field='object_id',
53 for_concrete_model=False
54 )
55 name = models.CharField(
56 verbose_name=_('name'),
57 max_length=200
58 )
59 created = models.DateTimeField(
60 verbose_name=_('created'),
61 auto_now_add=True
62 )
63 scheduled = models.DateTimeField(
64 verbose_name=_('scheduled'),
65 null=True,
66 blank=True
67 )
68 interval = models.PositiveIntegerField(
69 verbose_name=_('interval'),
70 blank=True,
71 null=True,
72 validators=(
73 MinValueValidator(1),
74 ),
75 help_text=_('Recurrence interval (in minutes)')
76 )
77 started = models.DateTimeField(
78 verbose_name=_('started'),
79 null=True,
80 blank=True
81 )
82 completed = models.DateTimeField(
83 verbose_name=_('completed'),
84 null=True,
85 blank=True
86 )
87 execution_time = models.DurationField(
88 verbose_name=_('execution time'),
89 null=True,
90 blank=True,
91 editable=False
92 )
93 user = models.ForeignKey(
94 to=settings.AUTH_USER_MODEL,
95 on_delete=models.SET_NULL,
96 related_name='+',
97 blank=True,
98 null=True
99 )
100 status = models.CharField(
101 verbose_name=_('status'),
102 max_length=30,
103 choices=JobStatusChoices,
104 default=JobStatusChoices.STATUS_PENDING
105 )
106 data = models.JSONField(
107 verbose_name=_('data'),
108 encoder=DjangoJSONEncoder,
109 null=True,
110 blank=True,
111 )
112 error = models.TextField(
113 verbose_name=_('error'),
114 editable=False,
115 blank=True
116 )
117 job_id = models.UUIDField(
118 verbose_name=_('job ID'),
119 unique=True
120 )
121 queue_name = models.CharField(
122 verbose_name=_('queue name'),
123 max_length=100,
124 blank=True,
125 help_text=_('Name of the queue in which this job was enqueued')
126 )
127 notifications = models.CharField(
128 verbose_name=_('notifications'),
129 max_length=30,
130 choices=JobNotificationChoices,
131 default=JobNotificationChoices.NOTIFICATION_ALWAYS
132 )
133 log_entries = ArrayField(
134 verbose_name=_('log entries'),
135 base_field=models.JSONField(
136 encoder=DjangoJSONEncoder,
137 decoder=JobLogDecoder,
138 ),
139 blank=True,
140 default=list,
141 )
143 objects = RestrictedQuerySet.as_manager()
145 class Meta:
146 ordering = ['-created']
147 indexes = (
148 models.Index(fields=('-created',)), # Default ordering
149 models.Index(fields=('object_type', 'object_id')),
150 )
151 verbose_name = _('job')
152 verbose_name_plural = _('jobs')
154 def __str__(self):
155 return self.name
157 def get_absolute_url(self):
158 # TODO: Employ dynamic registration
159 if self.object_type:
160 if self.object_type.model == 'reportmodule':
161 return reverse('extras:report_result', kwargs={'job_pk': self.pk})
162 if self.object_type.model == 'scriptmodule':
163 return reverse('extras:script_result', kwargs={'job_pk': self.pk})
164 return reverse('core:job', args=[self.pk])
166 def get_status_color(self):
167 return JobStatusChoices.colors.get(self.status)
169 def get_event_type(self):
170 return {
171 JobStatusChoices.STATUS_COMPLETED: JOB_COMPLETED,
172 JobStatusChoices.STATUS_FAILED: JOB_FAILED,
173 JobStatusChoices.STATUS_ERRORED: JOB_ERRORED,
174 }.get(self.status)
176 def clean(self):
177 super().clean()
179 # Validate the assigned object type
180 if self.object_type and not has_feature(self.object_type, 'jobs'): 180 ↛ 181line 180 didn't jump to line 181 because the condition on line 180 was never true
181 raise ValidationError(
182 _("Jobs cannot be assigned to this object type ({type}).").format(type=self.object_type)
183 )
185 @property
186 def duration(self):
187 if not self.completed:
188 return None
190 start_time = self.started or self.created
192 if not start_time:
193 return None
195 duration = self.completed - start_time
196 minutes, seconds = divmod(duration.total_seconds(), 60)
198 return f"{int(minutes)} minutes, {seconds:.2f} seconds"
200 def delete(self, *args, **kwargs):
201 # Use the stored queue name, or fall back to get_queue_for_model for legacy jobs
202 rq_queue_name = self.queue_name or get_queue_for_model(self.object_type.model if self.object_type else None)
203 rq_job_id = str(self.job_id)
205 super().delete(*args, **kwargs)
207 # Cancel the RQ job using the stored queue name
208 queue = django_rq.get_queue(rq_queue_name)
209 job = queue.fetch_job(rq_job_id)
211 if job:
212 try:
213 job.cancel()
214 except InvalidJobOperation:
215 # Job may raise this exception from get_status() if missing from Redis
216 pass
218 def start(self):
219 """
220 Record the job's start time and update its status to "running."
221 """
222 if self.started is not None:
223 return
225 # Start the job
226 self.started = timezone.now()
227 self.status = JobStatusChoices.STATUS_RUNNING
228 self.save()
230 # Send signal
231 job_start.send(self)
232 start.alters_data = True
234 def terminate(self, status=JobStatusChoices.STATUS_COMPLETED, error=None):
235 """
236 Mark the job as completed, optionally specifying a particular termination status.
237 """
238 if status not in JobStatusChoices.TERMINAL_STATE_CHOICES:
239 raise ValueError(
240 _("Invalid status for job termination. Choices are: {choices}").format(
241 choices=', '.join(JobStatusChoices.TERMINAL_STATE_CHOICES)
242 )
243 )
245 # Set the job's status and completion time
246 self.status = status
247 if error:
248 self.error = error
249 self.completed = timezone.now()
250 if self.started:
251 self.execution_time = self.completed - self.started
252 self.save()
254 # Notify the user (if any) of completion
255 if self.user and self.notifications != JobNotificationChoices.NOTIFICATION_NEVER:
256 if (
257 self.notifications == JobNotificationChoices.NOTIFICATION_ALWAYS or
258 status != JobStatusChoices.STATUS_COMPLETED
259 ):
260 Notification(
261 user=self.user,
262 object=self,
263 event_type=self.get_event_type(),
264 ).save()
266 # Send signal
267 job_end.send(self)
268 terminate.alters_data = True
270 def log(self, record: logging.LogRecord):
271 """
272 Record a LogRecord from Python's native logging in the job's log.
273 """
274 entry = JobLogEntry.from_logrecord(record)
275 self.log_entries.append(asdict(entry))
277 @classmethod
278 def enqueue(
279 cls,
280 func,
281 instance=None,
282 name='',
283 user=None,
284 schedule_at=None,
285 interval=None,
286 immediate=False,
287 queue_name=None,
288 notifications=None,
289 **kwargs
290 ):
291 """
292 Create a Job instance and enqueue a job using the given callable
294 Args:
295 func: The callable object to be enqueued for execution
296 instance: The NetBox object to which this job pertains (optional)
297 name: Name for the job (optional)
298 user: The user responsible for running the job
299 schedule_at: Schedule the job to be executed at the passed date and time
300 interval: Recurrence interval (in minutes)
301 immediate: Run the job immediately without scheduling it in the background. Should be used for interactive
302 management commands only.
303 notifications: Notification behavior on job completion (always, on_failure, or never)
304 """
305 if schedule_at and immediate: 305 ↛ 306line 305 didn't jump to line 306 because the condition on line 305 was never true
306 raise ValueError(_("enqueue() cannot be called with values for both schedule_at and immediate."))
308 if instance: 308 ↛ 309line 308 didn't jump to line 309 because the condition on line 308 was never true
309 object_type = ObjectType.objects.get_for_model(instance, for_concrete_model=False)
310 object_id = instance.pk
311 else:
312 object_type = object_id = None
313 rq_queue_name = queue_name if queue_name else get_queue_for_model(object_type.model if object_type else None)
314 queue = django_rq.get_queue(rq_queue_name)
315 status = JobStatusChoices.STATUS_SCHEDULED if schedule_at else JobStatusChoices.STATUS_PENDING
316 job = Job(
317 object_type=object_type,
318 object_id=object_id,
319 name=name,
320 status=status,
321 scheduled=schedule_at,
322 interval=interval,
323 user=user,
324 job_id=uuid.uuid4(),
325 queue_name=rq_queue_name,
326 notifications=notifications if notifications is not None else JobNotificationChoices.NOTIFICATION_ALWAYS
327 )
328 job.full_clean()
329 job.save()
331 # Run the job immediately, rather than enqueuing it as a background task. Note that this is a synchronous
332 # (blocking) operation, and execution will pause until the job completes.
333 if immediate: 333 ↛ 334line 333 didn't jump to line 334 because the condition on line 333 was never true
334 func(job_id=str(job.job_id), job=job, **kwargs)
336 # Schedule the job to run at a specific date & time.
337 elif schedule_at: 337 ↛ 338line 337 didn't jump to line 338 because the condition on line 337 was never true
338 callback = partial(queue.enqueue_at, schedule_at, func, job_id=str(job.job_id), job=job, **kwargs)
339 transaction.on_commit(callback)
341 # Schedule the job to run asynchronously at this first available opportunity.
342 else:
343 callback = partial(queue.enqueue, func, job_id=str(job.job_id), job=job, **kwargs)
344 transaction.on_commit(callback)
346 return job