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
« 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
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 _
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
23from ..choices import *
24from ..exceptions import SyncError
26__all__ = (
27 'AutoSyncRecord',
28 'DataFile',
29 'DataSource',
30)
32logger = logging.getLogger('netbox.core.data')
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 )
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 ]
94 def __str__(self):
95 return f'{self.name}'
97 @property
98 def docs_url(self):
99 return f'{settings.STATIC_URL}docs/models/{self._meta.app_label}/{self._meta.model_name}/'
101 def get_type_display(self):
102 if backend := registry['data_backends'].get(self.type):
103 return backend.label
104 return None
106 def get_status_color(self):
107 return DataSourceStatusChoices.colors.get(self.status)
109 @property
110 def url_scheme(self):
111 return urlparse(self.source_url).scheme.lower()
113 @property
114 def backend_class(self):
115 return registry['data_backends'].get(self.type)
117 @property
118 def ready_for_sync(self):
119 return self.enabled and self.status not in (
120 DataSourceStatusChoices.QUEUED,
121 DataSourceStatusChoices.SYNCING
122 )
124 def clean(self):
125 super().clean()
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 })
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 })
141 def save(self, *args, **kwargs):
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
152 super().save(*args, **kwargs)
154 def to_objectchange(self, action):
155 objectchange = super().to_objectchange(action)
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
174 return objectchange
176 def get_backend(self):
177 backend_params = self.parameters or {}
178 return self.backend_class(self.source_url, **backend_params)
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
186 if self.status == DataSourceStatusChoices.SYNCING:
187 raise SyncError(_("Cannot initiate sync; syncing already in progress."))
189 # Emit the pre_sync signal
190 pre_sync.send(sender=self.__class__, instance=self)
192 self.status = DataSourceStatusChoices.SYNCING
193 DataSource.objects.filter(pk=self.pk).update(status=self.status)
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:
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')
209 # Check for any updated/deleted files
210 updated_files = []
211 deleted_file_ids = []
212 for datafile in data_files:
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
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")
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")
232 # Walk the local replication to find new files
233 new_paths = self._walk(local_path) - known_paths
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")
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)
250 # Emit the post_sync signal
251 post_sync.send(sender=self.__class__, instance=self)
252 sync.alters_data = True
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()
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)
270 logger.debug(f"Found {len(paths)} files")
271 return paths
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
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()
326 objects = RestrictedQuerySet.as_manager()
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')
339 def __str__(self):
340 return self.path
342 def get_absolute_url(self):
343 return reverse('core:datafile', args=[self.pk])
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
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)
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()
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()
381 return is_modified
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 )
404 _netbox_private = True
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')