Coverage for core/models/data.py: 34%

205 statements  

« prev     ^ index     » next       coverage.py v7.15.2, created at 2026-10-10 18:35 +0000

1import hashlib 

2import logging 

3import os 

4from fnmatch import fnmatchcase 

5from urllib.parse import urlparse 

6 

7import yaml 

8from django.conf import settings 

9from django.core.exceptions import ValidationError 

10from django.core.validators import RegexValidator 

11from django.db import models 

12from django.urls import reverse 

13from django.utils import timezone 

14from django.utils.translation import gettext_lazy as _ 

15 

16from netbox.constants import CENSOR_TOKEN, CENSOR_TOKEN_CHANGED 

17from netbox.models import PrimaryModel 

18from netbox.models.features import JobsMixin 

19from netbox.registry import registry 

20from utilities.fields import RestrictedGenericForeignKey 

21from utilities.querysets import RestrictedQuerySet 

22 

23from ..choices import * 

24from ..exceptions import SyncError 

25 

26__all__ = ( 

27 'AutoSyncRecord', 

28 'DataFile', 

29 'DataSource', 

30) 

31 

32logger = logging.getLogger('netbox.core.data') 

33 

34 

35class DataSource(JobsMixin, PrimaryModel): 

36 """ 

37 A remote source, such as a git repository, from which DataFiles are synchronized. 

38 """ 

39 name = models.CharField( 

40 verbose_name=_('name'), 

41 max_length=100, 

42 unique=True 

43 ) 

44 type = models.CharField( 

45 verbose_name=_('type'), 

46 max_length=50 

47 ) 

48 source_url = models.CharField( 

49 max_length=200, 

50 verbose_name=_('URL') 

51 ) 

52 status = models.CharField( 

53 verbose_name=_('status'), 

54 max_length=50, 

55 choices=DataSourceStatusChoices, 

56 default=DataSourceStatusChoices.NEW, 

57 editable=False 

58 ) 

59 enabled = models.BooleanField( 

60 verbose_name=_('enabled'), 

61 default=True 

62 ) 

63 sync_interval = models.PositiveSmallIntegerField( 

64 verbose_name=_('sync interval'), 

65 choices=JobIntervalChoices, 

66 blank=True, 

67 null=True 

68 ) 

69 ignore_rules = models.TextField( 

70 verbose_name=_('ignore rules'), 

71 blank=True, 

72 help_text=_("Patterns (one per line) matching files or paths to ignore when syncing") 

73 ) 

74 parameters = models.JSONField( 

75 verbose_name=_('parameters'), 

76 blank=True, 

77 null=True 

78 ) 

79 last_synced = models.DateTimeField( 

80 verbose_name=_('last synced'), 

81 blank=True, 

82 null=True, 

83 editable=False 

84 ) 

85 

86 class Meta: 

87 ordering = ('name',) 

88 verbose_name = _('data source') 

89 verbose_name_plural = _('data sources') 

90 permissions = [ 

91 ('sync', 'Synchronize data from remote source'), 

92 ] 

93 

94 def __str__(self): 

95 return f'{self.name}' 

96 

97 @property 

98 def docs_url(self): 

99 return f'{settings.STATIC_URL}docs/models/{self._meta.app_label}/{self._meta.model_name}/' 

100 

101 def get_type_display(self): 

102 if backend := registry['data_backends'].get(self.type): 

103 return backend.label 

104 return None 

105 

106 def get_status_color(self): 

107 return DataSourceStatusChoices.colors.get(self.status) 

108 

109 @property 

110 def url_scheme(self): 

111 return urlparse(self.source_url).scheme.lower() 

112 

113 @property 

114 def backend_class(self): 

115 return registry['data_backends'].get(self.type) 

116 

117 @property 

118 def ready_for_sync(self): 

119 return self.enabled and self.status not in ( 

120 DataSourceStatusChoices.QUEUED, 

121 DataSourceStatusChoices.SYNCING 

122 ) 

123 

124 def clean(self): 

125 super().clean() 

126 

127 # Validate data backend type 

128 if self.type and self.type not in registry['data_backends']: 

129 raise ValidationError({ 

130 'type': _("Unknown backend type: {type}".format(type=self.type)) 

131 }) 

132 

133 # Ensure URL scheme matches selected type 

134 if self.backend_class.is_local and self.url_scheme not in ('file', ''): 

135 raise ValidationError({ 

136 'source_url': _("URLs for local sources must start with {scheme} (or specify no scheme)").format( 

137 scheme='file://' 

138 ) 

139 }) 

140 

141 def save(self, *args, **kwargs): 

142 

143 # If recurring sync is disabled for an existing DataSource, clear any pending sync jobs for it and reset its 

144 # "queued" status 

145 if not self._state.adding and not self.sync_interval: 

146 self.jobs.filter(status=JobStatusChoices.STATUS_PENDING).delete() 

147 if self.status == DataSourceStatusChoices.QUEUED and self.last_synced: 

148 self.status = DataSourceStatusChoices.COMPLETED 

149 elif self.status == DataSourceStatusChoices.QUEUED: 

150 self.status = DataSourceStatusChoices.NEW 

151 

152 super().save(*args, **kwargs) 

153 

154 def to_objectchange(self, action): 

155 objectchange = super().to_objectchange(action) 

156 

157 # Censor any backend parameters marked as sensitive in the serialized data 

158 pre_change_params = {} 

159 post_change_params = {} 

160 if objectchange.prechange_data: 

161 pre_change_params = objectchange.prechange_data.get('parameters') or {} # parameters may be None 

162 if objectchange.postchange_data: 

163 post_change_params = objectchange.postchange_data.get('parameters') or {} 

164 for param in self.backend_class.sensitive_parameters: 

165 if post_change_params.get(param): 

166 if post_change_params[param] != pre_change_params.get(param): 

167 # Set the "changed" token if the parameter's value has been modified 

168 post_change_params[param] = CENSOR_TOKEN_CHANGED 

169 else: 

170 post_change_params[param] = CENSOR_TOKEN 

171 if pre_change_params.get(param): 

172 pre_change_params[param] = CENSOR_TOKEN 

173 

174 return objectchange 

175 

176 def get_backend(self): 

177 backend_params = self.parameters or {} 

178 return self.backend_class(self.source_url, **backend_params) 

179 

180 def sync(self): 

181 """ 

182 Create/update/delete child DataFiles as necessary to synchronize with the remote source. 

183 """ 

184 from core.signals import post_sync, pre_sync 

185 

186 if self.status == DataSourceStatusChoices.SYNCING: 

187 raise SyncError(_("Cannot initiate sync; syncing already in progress.")) 

188 

189 # Emit the pre_sync signal 

190 pre_sync.send(sender=self.__class__, instance=self) 

191 

192 self.status = DataSourceStatusChoices.SYNCING 

193 DataSource.objects.filter(pk=self.pk).update(status=self.status) 

194 

195 # Replicate source data locally 

196 try: 

197 backend = self.get_backend() 

198 except ModuleNotFoundError as e: 

199 raise SyncError( 

200 _("There was an error initializing the backend. A dependency needs to be installed: ") + str(e) 

201 ) 

202 with backend.fetch() as local_path: 

203 

204 logger.debug(f'Syncing files from source root {local_path}') 

205 data_files = self.datafiles.all() 

206 known_paths = {df.path for df in data_files} 

207 logger.debug(f'Starting with {len(known_paths)} known files') 

208 

209 # Check for any updated/deleted files 

210 updated_files = [] 

211 deleted_file_ids = [] 

212 for datafile in data_files: 

213 

214 try: 

215 if datafile.refresh_from_disk(source_root=local_path): 

216 updated_files.append(datafile) 

217 except FileNotFoundError: 

218 # File no longer exists 

219 deleted_file_ids.append(datafile.pk) 

220 continue 

221 

222 # Bulk update modified files 

223 updated_count = DataFile.objects.bulk_update( 

224 updated_files, ('last_updated', 'size', 'hash', 'data'), batch_size=settings.BULK_UPDATE_CHUNK_SIZE 

225 ) 

226 logger.debug(f"Updated {updated_count} files") 

227 

228 # Bulk delete deleted files 

229 deleted_count, __ = DataFile.objects.filter(pk__in=deleted_file_ids).delete() 

230 logger.debug(f"Deleted {deleted_count} files") 

231 

232 # Walk the local replication to find new files 

233 new_paths = self._walk(local_path) - known_paths 

234 

235 # Bulk create new files 

236 new_datafiles = [] 

237 for path in new_paths: 

238 datafile = DataFile(source=self, path=path) 

239 datafile.refresh_from_disk(source_root=local_path) 

240 datafile.full_clean() 

241 new_datafiles.append(datafile) 

242 created_count = len(DataFile.objects.bulk_create(new_datafiles, batch_size=100)) 

243 logger.debug(f"Created {created_count} data files") 

244 

245 # Update status & last_synced time 

246 self.status = DataSourceStatusChoices.COMPLETED 

247 self.last_synced = timezone.now() 

248 DataSource.objects.filter(pk=self.pk).update(status=self.status, last_synced=self.last_synced) 

249 

250 # Emit the post_sync signal 

251 post_sync.send(sender=self.__class__, instance=self) 

252 sync.alters_data = True 

253 

254 def _walk(self, root): 

255 """ 

256 Return a set of all non-excluded files within the root path. 

257 """ 

258 logger.debug(f"Walking {root}...") 

259 paths = set() 

260 

261 for path, dir_names, file_names in os.walk(root): 

262 path = path.split(root)[1].lstrip('/') # Strip root path 

263 if path.startswith('.'): 

264 continue 

265 for file_name in file_names: 

266 file_path = os.path.join(path, file_name) 

267 if not self._ignore(file_path): 

268 paths.add(file_path) 

269 

270 logger.debug(f"Found {len(paths)} files") 

271 return paths 

272 

273 def _ignore(self, file_path): 

274 """ 

275 Returns a boolean indicating whether the file should be ignored per the DataSource's configured 

276 ignore rules. file_path is the full relative path (e.g. "subdir/file.txt"). 

277 """ 

278 if os.path.basename(file_path).startswith('.'): 

279 return True 

280 for rule in self.ignore_rules.splitlines(): 

281 if fnmatchcase(file_path, rule) or fnmatchcase(os.path.basename(file_path), rule): 

282 return True 

283 return False 

284 

285 

286class DataFile(models.Model): 

287 """ 

288 The database representation of a remote file fetched from a remote DataSource. DataFile instances should be created, 

289 updated, or deleted only by calling DataSource.sync(). 

290 """ 

291 created = models.DateTimeField( 

292 verbose_name=_('created'), 

293 auto_now_add=True 

294 ) 

295 last_updated = models.DateTimeField( 

296 verbose_name=_('last updated'), 

297 editable=False 

298 ) 

299 source = models.ForeignKey( 

300 to='core.DataSource', 

301 on_delete=models.CASCADE, 

302 related_name='datafiles', 

303 editable=False 

304 ) 

305 path = models.CharField( 

306 verbose_name=_('path'), 

307 max_length=1000, 

308 editable=False, 

309 help_text=_("File path relative to the data source's root") 

310 ) 

311 size = models.PositiveIntegerField( 

312 editable=False, 

313 verbose_name=_('size') 

314 ) 

315 hash = models.CharField( 

316 verbose_name=_('hash'), 

317 max_length=64, 

318 editable=False, 

319 validators=[ 

320 RegexValidator(regex='^[0-9a-f]{64}$', message=_("Length must be 64 hexadecimal characters.")) 

321 ], 

322 help_text=_('SHA256 hash of the file data') 

323 ) 

324 data = models.BinaryField() 

325 

326 objects = RestrictedQuerySet.as_manager() 

327 

328 class Meta: 

329 ordering = ('source', 'path') 

330 constraints = ( 

331 models.UniqueConstraint( 

332 fields=('source', 'path'), 

333 name='%(app_label)s_%(class)s_unique_source_path' 

334 ), 

335 ) 

336 verbose_name = _('data file') 

337 verbose_name_plural = _('data files') 

338 

339 def __str__(self): 

340 return self.path 

341 

342 def get_absolute_url(self): 

343 return reverse('core:datafile', args=[self.pk]) 

344 

345 @property 

346 def data_as_string(self): 

347 if not self.data: 

348 return None 

349 try: 

350 return self.data.decode('utf-8') 

351 except UnicodeDecodeError: 

352 return None 

353 

354 def get_data(self): 

355 """ 

356 Attempt to read the file data as JSON/YAML and return a native Python object. Returns None if the file 

357 content cannot be decoded. 

358 """ 

359 # TODO: Something more robust 

360 if (data := self.data_as_string) is None: 

361 return None 

362 return yaml.safe_load(data) 

363 

364 def refresh_from_disk(self, source_root): 

365 """ 

366 Update instance attributes from the file on disk. Returns True if any attribute 

367 has changed. 

368 """ 

369 file_path = os.path.join(source_root, self.path) 

370 with open(file_path, 'rb') as f: 

371 file_hash = hashlib.sha256(f.read()).hexdigest() 

372 

373 # Update instance file attributes & data 

374 if is_modified := file_hash != self.hash: 

375 self.last_updated = timezone.now() 

376 self.size = os.path.getsize(file_path) 

377 self.hash = file_hash 

378 with open(file_path, 'rb') as f: 

379 self.data = f.read() 

380 

381 return is_modified 

382 

383 

384class AutoSyncRecord(models.Model): 

385 """ 

386 Maps a DataFile to a synced object for efficient automatic updating. 

387 """ 

388 datafile = models.ForeignKey( 

389 to=DataFile, 

390 on_delete=models.CASCADE, 

391 related_name='+' 

392 ) 

393 object_type = models.ForeignKey( 

394 to='contenttypes.ContentType', 

395 on_delete=models.CASCADE, 

396 related_name='+' 

397 ) 

398 object_id = models.PositiveBigIntegerField() 

399 object = RestrictedGenericForeignKey( 

400 ct_field='object_type', 

401 fk_field='object_id' 

402 ) 

403 

404 _netbox_private = True 

405 

406 class Meta: 

407 constraints = ( 

408 models.UniqueConstraint( 

409 fields=('object_type', 'object_id'), 

410 name='%(app_label)s_%(class)s_object' 

411 ), 

412 ) 

413 verbose_name = _('auto sync record') 

414 verbose_name_plural = _('auto sync records')