Coverage for core/jobs.py: 20%
131 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 sys
2from datetime import timedelta
3from importlib import import_module
5import requests
6from django.conf import settings
7from django.core.cache import cache
8from django.db.models import Exists, OuterRef, Subquery
9from django.utils import timezone
10from packaging import version
12from core.models import Job, ObjectChange
13from netbox.config import Config
14from netbox.jobs import JobRunner, system_job
15from netbox.search.backends import search_backend
16from utilities.proxy import resolve_proxies
18from .choices import DataSourceStatusChoices, JobIntervalChoices, ObjectChangeActionChoices
19from .models import DataSource
22class SyncDataSourceJob(JobRunner):
23 """
24 Call sync() on a DataSource.
25 """
27 class Meta:
28 name = 'Synchronization'
30 @classmethod
31 def enqueue(cls, *args, **kwargs):
32 job = super().enqueue(*args, **kwargs)
34 # Update the DataSource's synchronization status to queued
35 if datasource := job.object:
36 datasource.status = DataSourceStatusChoices.QUEUED
37 DataSource.objects.filter(pk=datasource.pk).update(status=datasource.status)
39 return job
41 def run(self, *args, **kwargs):
42 datasource = DataSource.objects.get(pk=self.job.object_id)
43 self.logger.debug(f"Found DataSource ID {datasource.pk}")
45 try:
46 self.logger.info(f"Syncing data source {datasource}")
47 datasource.sync()
49 # Update the search cache for DataFiles belonging to this source
50 self.logger.debug("Updating search cache for data files")
51 search_backend.cache(datasource.datafiles.iterator())
53 except Exception as e:
54 self.logger.error(f"Error syncing data source: {e}")
55 DataSource.objects.filter(pk=datasource.pk).update(status=DataSourceStatusChoices.FAILED)
56 raise e
58 self.logger.info("Syncing completed successfully")
61@system_job(interval=JobIntervalChoices.INTERVAL_DAILY)
62class SystemHousekeepingJob(JobRunner):
63 """
64 Perform daily system housekeeping functions.
65 """
66 class Meta:
67 name = "System Housekeeping"
69 def run(self, *args, **kwargs):
70 # Skip if running in development or test mode
71 if settings.DEBUG:
72 self.logger.warning("Aborting execution: Debug is enabled")
73 return
74 if 'test' in sys.argv:
75 self.logger.warning("Aborting execution: Tests are running")
76 return
78 self.send_census_report()
79 self.clear_expired_sessions()
80 self.prune_changelog()
81 self.delete_expired_jobs()
82 self.check_for_new_releases()
84 def send_census_report(self):
85 """
86 Send a census report (if enabled).
87 """
88 self.logger.info("Reporting census data...")
89 if settings.ISOLATED_DEPLOYMENT:
90 self.logger.info("ISOLATED_DEPLOYMENT is enabled; skipping")
91 return
92 if not settings.CENSUS_REPORTING_ENABLED:
93 self.logger.info("CENSUS_REPORTING_ENABLED is disabled; skipping")
94 return
96 census_data = {
97 'version': settings.RELEASE.full_version,
98 'python_version': sys.version.split()[0],
99 'deployment_id': settings.DEPLOYMENT_ID,
100 }
101 try:
102 requests.get(
103 url=settings.CENSUS_URL,
104 params=census_data,
105 timeout=3,
106 proxies=resolve_proxies(url=settings.CENSUS_URL)
107 )
108 except requests.exceptions.RequestException:
109 pass
111 def clear_expired_sessions(self):
112 """
113 Clear any expired sessions from the database.
114 """
115 self.logger.info("Clearing expired sessions...")
116 engine = import_module(settings.SESSION_ENGINE)
117 try:
118 engine.SessionStore.clear_expired()
119 self.logger.info("Sessions cleared.")
120 except NotImplementedError:
121 self.logger.warning(
122 f"The configured session engine ({settings.SESSION_ENGINE}) does not support "
123 f"clearing sessions; skipping."
124 )
126 def prune_changelog(self):
127 """
128 Delete any ObjectChange records older than the configured changelog retention time (if any).
129 """
130 self.logger.info('Pruning old changelog entries...')
131 config = Config()
132 if not config.CHANGELOG_RETENTION:
133 self.logger.info('No retention period specified; skipping.')
134 return
136 cutoff = timezone.now() - timedelta(days=config.CHANGELOG_RETENTION)
137 self.logger.debug(f'Changelog retention period: {config.CHANGELOG_RETENTION} days ({cutoff:%Y-%m-%d %H:%M:%S})')
139 expired_qs = ObjectChange.objects.filter(time__lt=cutoff)
141 # When enabled, retain each object's original create record and most recent update record while pruning expired
142 # changelog entries. This applies only to objects without a delete record.
143 if config.CHANGELOG_RETAIN_CREATE_LAST_UPDATE:
144 self.logger.debug('Retaining changelog create records and last update records (excluding deleted objects)')
146 deleted_exists = ObjectChange.objects.filter(
147 action=ObjectChangeActionChoices.ACTION_DELETE,
148 changed_object_type_id=OuterRef('changed_object_type_id'),
149 changed_object_id=OuterRef('changed_object_id'),
150 )
152 # Keep create records only where no delete exists for that object
153 create_pks_to_keep = (
154 ObjectChange.objects.filter(action=ObjectChangeActionChoices.ACTION_CREATE)
155 .annotate(has_delete=Exists(deleted_exists))
156 .filter(has_delete=False)
157 .values('pk')
158 )
160 # Keep the most recent update per object only where no delete exists for the object
161 latest_update_pks_to_keep = (
162 ObjectChange.objects.filter(action=ObjectChangeActionChoices.ACTION_UPDATE)
163 .annotate(has_delete=Exists(deleted_exists))
164 .filter(has_delete=False)
165 .order_by('changed_object_type_id', 'changed_object_id', '-time', '-pk')
166 .distinct('changed_object_type_id', 'changed_object_id')
167 .values('pk')
168 )
170 expired_qs = expired_qs.exclude(pk__in=Subquery(create_pks_to_keep))
171 expired_qs = expired_qs.exclude(pk__in=Subquery(latest_update_pks_to_keep))
173 count = expired_qs.delete()[0]
174 self.logger.info(f'Deleted {count} expired changelog records')
176 def delete_expired_jobs(self):
177 """
178 Delete any jobs older than the configured retention period (if any).
179 """
180 self.logger.info("Deleting expired jobs...")
181 config = Config()
182 if not config.JOB_RETENTION:
183 self.logger.info("No retention period specified; skipping.")
184 return
186 cutoff = timezone.now() - timedelta(days=config.JOB_RETENTION)
187 self.logger.debug(
188 f"Job retention period: {config.JOB_RETENTION} days ({cutoff:%Y-%m-%d %H:%M:%S})"
189 )
191 count = Job.objects.filter(created__lt=cutoff).delete()[0]
192 self.logger.info(f"Deleted {count} expired jobs")
194 def check_for_new_releases(self):
195 """
196 Check for new releases and cache the latest release.
197 """
198 self.logger.info("Checking for new releases...")
199 if settings.ISOLATED_DEPLOYMENT:
200 self.logger.info("ISOLATED_DEPLOYMENT is enabled; skipping")
201 return
202 if not settings.RELEASE_CHECK_URL:
203 self.logger.info("RELEASE_CHECK_URL is not set; skipping")
204 return
206 # Fetch the latest releases
207 self.logger.debug(f"Release check URL: {settings.RELEASE_CHECK_URL}")
208 try:
209 response = requests.get(
210 url=settings.RELEASE_CHECK_URL,
211 headers={'Accept': 'application/vnd.github.v3+json'},
212 proxies=resolve_proxies(url=settings.RELEASE_CHECK_URL)
213 )
214 response.raise_for_status()
215 except requests.exceptions.RequestException as exc:
216 self.logger.error(f"Error fetching release: {exc}")
217 return
219 # Determine the most recent stable release
220 releases = []
221 for release in response.json():
222 if 'tag_name' not in release or release.get('devrelease') or release.get('prerelease'):
223 continue
224 releases.append((version.parse(release['tag_name']), release.get('html_url')))
225 self.logger.debug(f"Found {len(response.json())} releases; {len(releases)} usable")
226 if not releases:
227 self.logger.info("No usable releases found; skipping")
228 return
229 latest_release = max(releases)
230 self.logger.info(f"Latest release: {latest_release[0]}")
232 # Cache the most recent release
233 cache.set('latest_release', latest_release, None)