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

1import logging 

2import uuid 

3from dataclasses import asdict 

4from functools import partial 

5 

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 

18 

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 

29 

30__all__ = ( 

31 'Job', 

32) 

33 

34 

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 ) 

142 

143 objects = RestrictedQuerySet.as_manager() 

144 

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') 

153 

154 def __str__(self): 

155 return self.name 

156 

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]) 

165 

166 def get_status_color(self): 

167 return JobStatusChoices.colors.get(self.status) 

168 

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) 

175 

176 def clean(self): 

177 super().clean() 

178 

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 ) 

184 

185 @property 

186 def duration(self): 

187 if not self.completed: 

188 return None 

189 

190 start_time = self.started or self.created 

191 

192 if not start_time: 

193 return None 

194 

195 duration = self.completed - start_time 

196 minutes, seconds = divmod(duration.total_seconds(), 60) 

197 

198 return f"{int(minutes)} minutes, {seconds:.2f} seconds" 

199 

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) 

204 

205 super().delete(*args, **kwargs) 

206 

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) 

210 

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 

217 

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 

224 

225 # Start the job 

226 self.started = timezone.now() 

227 self.status = JobStatusChoices.STATUS_RUNNING 

228 self.save() 

229 

230 # Send signal 

231 job_start.send(self) 

232 start.alters_data = True 

233 

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 ) 

244 

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() 

253 

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() 

265 

266 # Send signal 

267 job_end.send(self) 

268 terminate.alters_data = True 

269 

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

276 

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 

293 

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

307 

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() 

330 

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) 

335 

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) 

340 

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) 

345 

346 return job