Coverage for core/jobs.py: 20%

131 statements  

« 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 

4 

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 

11 

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 

17 

18from .choices import DataSourceStatusChoices, JobIntervalChoices, ObjectChangeActionChoices 

19from .models import DataSource 

20 

21 

22class SyncDataSourceJob(JobRunner): 

23 """ 

24 Call sync() on a DataSource. 

25 """ 

26 

27 class Meta: 

28 name = 'Synchronization' 

29 

30 @classmethod 

31 def enqueue(cls, *args, **kwargs): 

32 job = super().enqueue(*args, **kwargs) 

33 

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) 

38 

39 return job 

40 

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

44 

45 try: 

46 self.logger.info(f"Syncing data source {datasource}") 

47 datasource.sync() 

48 

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

52 

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 

57 

58 self.logger.info("Syncing completed successfully") 

59 

60 

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" 

68 

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 

77 

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

83 

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 

95 

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 

110 

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 ) 

125 

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 

135 

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

138 

139 expired_qs = ObjectChange.objects.filter(time__lt=cutoff) 

140 

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

145 

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 ) 

151 

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 ) 

159 

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 ) 

169 

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

172 

173 count = expired_qs.delete()[0] 

174 self.logger.info(f'Deleted {count} expired changelog records') 

175 

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 

185 

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 ) 

190 

191 count = Job.objects.filter(created__lt=cutoff).delete()[0] 

192 self.logger.info(f"Deleted {count} expired jobs") 

193 

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 

205 

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 

218 

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

231 

232 # Cache the most recent release 

233 cache.set('latest_release', latest_release, None)