Coverage for src/backend/InvenTree/InvenTree/tasks.py: 40%
457 statements
« prev ^ index » next coverage.py v7.15.2, created at 2026-10-07 17:47 +0000
« prev ^ index » next coverage.py v7.15.2, created at 2026-10-07 17:47 +0000
1"""Functions for tasks and a few general async tasks."""
3import json
4import os
5import re
6import warnings
7from collections.abc import Callable
8from dataclasses import dataclass
9from datetime import datetime, timedelta
10from typing import Optional
12from django.conf import settings
13from django.contrib.auth import get_user_model
14from django.core.exceptions import AppRegistryNotReady, ValidationError
15from django.core.management import call_command
16from django.db import DEFAULT_DB_ALIAS, connections
17from django.db.migrations.executor import MigrationExecutor
18from django.db.utils import NotSupportedError, OperationalError, ProgrammingError
19from django.utils import timezone
20from django.utils.translation import gettext_lazy as _
22import requests
23import structlog
24from maintenance_mode.core import (
25 get_maintenance_mode,
26 maintenance_mode_on,
27 set_maintenance_mode,
28)
29from opentelemetry import trace
31from common.settings import get_global_setting, set_global_setting
32from plugin import registry
34from .version import isInvenTreeUpToDate
36logger = structlog.get_logger('inventree')
37tracer = trace.get_tracer(__name__)
40def schedule_task(taskname, **kwargs):
41 """Create a scheduled task.
43 If the task has already been scheduled, ignore!
44 """
45 # If unspecified, repeat indefinitely
46 repeats = kwargs.pop('repeats', -1)
47 kwargs['repeats'] = repeats
49 try:
50 from django_q.models import Schedule
51 except AppRegistryNotReady: # pragma: no cover
52 logger.info('Could not start background tasks - App registry not ready')
53 return
55 try:
56 # If this task is already scheduled, don't schedule it again
57 # Instead, update the scheduling parameters
58 if Schedule.objects.filter(func=taskname).exists():
59 logger.debug("Scheduled task '%s' already exists - updating!", taskname)
61 Schedule.objects.filter(func=taskname).update(**kwargs)
62 else:
63 logger.info("Creating scheduled task '%s'", taskname)
65 Schedule.objects.create(name=taskname, func=taskname, **kwargs)
66 except (OperationalError, ProgrammingError): # pragma: no cover
67 # Required if the DB is not ready yet
68 pass
71def raise_warning(msg):
72 """Log and raise a warning."""
73 logger.warning(msg)
75 # If testing is running raise a warning that can be asserted
76 if settings.TESTING:
77 warnings.warn(msg, stacklevel=2)
80def check_daily_holdoff(task_name: str, n_days: int = 1) -> bool:
81 """Check if a periodic task should be run, based on the provided setting name.
83 Arguments:
84 task_name (str): The name of the task being run, e.g. 'dummy_task'
85 n_days (int): The number of days between task runs (default = 1)
87 Returns:
88 bool: If the task should be run *now*, or wait another day
90 This function will determine if the task should be run *today*,
91 based on when it was last run, or if we have a record of it running at all.
93 Note that this function creates some *hidden* global settings (designated with the _ prefix),
94 which are used to keep a running track of when the particular task was was last run.
95 """
96 if n_days <= 0:
97 logger.info(
98 "Specified interval for task '%s' < 1 - task will not run", task_name
99 )
100 return False
102 attempt_key = f'_{task_name}_ATTEMPT'
103 success_key = f'_{task_name}_SUCCESS'
105 # Check for recent success information
106 last_success = get_global_setting(success_key, '', cache=False)
108 if last_success:
109 try:
110 last_success = datetime.fromisoformat(last_success)
111 except ValueError:
112 last_success = None
114 if last_success:
115 threshold = datetime.now() - timedelta(days=n_days)
117 if last_success.date() > threshold.date():
118 logger.info(
119 "Last successful run for '%s' was too recent - skipping task", task_name
120 )
121 return False
123 # Check for any information we have about this task
124 last_attempt = get_global_setting(attempt_key, '', cache=False)
126 if last_attempt:
127 try:
128 last_attempt = datetime.fromisoformat(last_attempt)
129 except ValueError:
130 last_attempt = None
132 if last_attempt:
133 # Do not attempt if the most recent *attempt* was within 12 hours
134 threshold = datetime.now() - timedelta(hours=12)
136 if last_attempt > threshold:
137 logger.info(
138 "Last attempt for '%s' was too recent - skipping task", task_name
139 )
140 return False
142 # Record this attempt
143 record_task_attempt(task_name)
145 # No reason *not* to run this task now
146 return True
149def record_task_attempt(task_name: str):
150 """Record that a multi-day task has been attempted *now*."""
151 logger.info("Logging task attempt for '%s'", task_name)
153 set_global_setting(f'_{task_name}_ATTEMPT', datetime.now().isoformat(), None)
156def record_task_success(task_name: str):
157 """Record that a multi-day task was successful *now*."""
158 set_global_setting(f'_{task_name}_SUCCESS', datetime.now().isoformat(), None)
161def check_existing_task(taskname, group: str, *args, **kwargs) -> Optional[str]:
162 """Test if an identical task is already registered with the worker.
164 This will only return true if the task name, group, args and kwargs all match an existing task.
166 Arguments:
167 taskname: The name of the task to check for, in the format 'app.module.function'
168 group: The group that the task belongs to
169 *args: Positional arguments to match
170 **kwargs: Keyword arguments to match
172 Returns:
173 Optional[str]: The ID of the matching task, if found, otherwise None
174 """
175 from django_q.models import OrmQ
177 task_id = None
179 # Iterate through all available tasks, with the most recent first
180 for task in OrmQ.objects.all().order_by('-id'):
181 if task.func() != taskname and task.task.get('func') != taskname:
182 # Task does not match
183 continue
185 if task.group() != group: 185 ↛ 187line 185 didn't jump to line 187 because the condition on line 185 was never true
186 # Group does not match
187 continue
189 if task.args() != args:
190 # Task args do not match
191 continue
193 if task.kwargs() != kwargs: 193 ↛ 195line 193 didn't jump to line 195 because the condition on line 193 was never true
194 # Task kwargs do not match
195 continue
197 task_id = task.task_id()
199 break
201 return task_id
204def offload_task(
205 taskname,
206 *args,
207 force_async: bool = False,
208 force_sync: bool = False,
209 check_duplicates: bool = True,
210 **kwargs,
211) -> str | bool:
212 """Create an AsyncTask if workers are running. This is different to a 'scheduled' task, in that it only runs once!
214 If workers are not running or force_sync flag, is set then the task is ran synchronously.
216 Arguments:
217 taskname: The name of the task to be run, in the format 'app.module.function'
218 *args: Positional arguments to be passed to the task function
219 force_async: If True, force the task to be offloaded (even if workers are not running)
220 force_sync: If True, force the task to be run synchronously (even if workers are running)
221 check_duplicates: If True, check for existing identical tasks before offloading
222 **kwargs: Keyword arguments to be passed to the task function
224 Returns:
225 str | bool: Task ID if the task was offloaded, True if ran synchronously, False otherwise
226 """
227 from InvenTree.exceptions import log_error
229 # Extract group information from kwargs
230 group = kwargs.pop('group', 'inventree')
232 try:
233 import importlib
235 from django_q.tasks import AsyncTask
237 from InvenTree.status import is_worker_running
238 except AppRegistryNotReady: # pragma: no cover
239 logger.warning("Could not offload task '%s' - app registry not ready", taskname)
241 if force_async:
242 # Cannot async the task, so return False
243 return False
244 else:
245 force_sync = True
246 except (OperationalError, ProgrammingError): # pragma: no cover
247 raise_warning(f"Could not offload task '{taskname}' - database not ready")
249 if force_async:
250 # Cannot async the task, so return False
251 return False
252 else:
253 force_sync = True
255 if force_async or (is_worker_running() and not force_sync):
256 # Before offloading, check if a duplicate task exists
257 if not force_sync and check_duplicates: 257 ↛ 266line 257 didn't jump to line 266 because the condition on line 257 was always true
258 if task_id := check_existing_task(taskname, group, *args, **kwargs):
259 logger.debug(
260 "Skipping duplicate task '%s' with ID '%s'", taskname, task_id
261 )
263 return task_id
265 # Running as asynchronous task
266 try:
267 task = AsyncTask(taskname, *args, group=group, **kwargs)
268 with tracer.start_as_current_span(f'async worker: {taskname}'):
269 task.run()
271 # Return the ID of the offloaded task, so that it can be tracked if needed
272 return task.id
273 except ImportError:
274 raise_warning(f"WARNING: '{taskname}' not offloaded - Function not found")
275 return False
276 except Exception as exc:
277 raise_warning(f"WARNING: '{taskname}' not offloaded due to {exc!s}")
278 log_error('offload_task', scope='worker')
279 return False
280 else:
281 if callable(taskname): 281 ↛ 288line 281 didn't jump to line 288 because the condition on line 281 was always true
282 # function was passed - use that
283 _func = taskname
284 else:
285 # Split on the last dot: everything before is the module path,
286 # everything after is the function name. rsplit handles any depth
287 # (e.g. 'app.module.func' or 'app.sub.module.func').
288 try:
289 module_path, func_name = taskname.rsplit('.', 1)
290 except ValueError:
291 raise_warning(
292 f"WARNING: '{taskname}' not started - Malformed function path"
293 )
294 return False
296 try:
297 _mod = importlib.import_module(module_path)
298 except ModuleNotFoundError:
299 log_error('offload_task', scope='worker')
300 raise_warning(
301 f"WARNING: '{taskname}' not started - No module named '{module_path}'"
302 )
303 return False
305 _func = getattr(_mod, func_name, None)
306 if _func is None:
307 log_error('offload_task', scope='worker')
308 raise_warning(
309 f"WARNING: '{taskname}' not started - No function named '{func_name}'"
310 )
311 return False
313 # Workers are not running: run it as synchronous task
314 try:
315 with tracer.start_as_current_span(f'sync worker: {taskname}'):
316 _func(*args, **kwargs)
317 except Exception as exc:
318 log_error('offload_task', scope='worker')
319 raise_warning(f"WARNING: '{taskname}' failed due to {exc!s}")
320 raise exc
322 # Finally, task either completed successfully or was offloaded
323 return True
326def get_queued_task(task_id: str):
327 """Find the task in the queue, if it exists.
329 Note that the OrmQ table does NOT keep the task ID as a database field,
330 it is instead stored in the payload data.
331 If there are a large number of pending tasks, this query may be inefficient,
332 but there is no other way to find a queued task by ID.
333 """
334 offset = 0
335 limit = 500
337 if not task_id: 337 ↛ 339line 337 didn't jump to line 339 because the condition on line 337 was never true
338 # Return early if no task ID was provided
339 return None
341 task_id = str(task_id)
343 from django_q.models import OrmQ
345 while True:
346 queued_tasks = OrmQ.objects.all().order_by('id')[offset : offset + limit]
347 if not queued_tasks:
348 break
350 for task in queued_tasks:
351 if task.task_id() == task_id:
352 return task
354 offset += limit
356 # No matching task was discovered
357 return None
360@dataclass()
361class ScheduledTask:
362 """A scheduled task.
364 - interval: The interval at which the task should be run
365 - minutes: The number of minutes between task runs
366 - func: The function to be run
367 """
369 func: Callable
370 interval: str
371 minutes: Optional[int] = None
373 MINUTES: str = 'I'
374 HOURLY: str = 'H'
375 DAILY: str = 'D'
376 WEEKLY: str = 'W'
377 MONTHLY: str = 'M'
378 QUARTERLY: str = 'Q'
379 YEARLY: str = 'Y'
381 TYPE: tuple[str] = (MINUTES, HOURLY, DAILY, WEEKLY, MONTHLY, QUARTERLY, YEARLY)
384class TaskRegister:
385 """Registry for periodic tasks."""
387 task_list: list[ScheduledTask] = []
389 def register(self, task, schedule, minutes: Optional[int] = None):
390 """Register a task with the que."""
391 self.task_list.append(ScheduledTask(task, schedule, minutes))
394tasks = TaskRegister()
397def scheduled_task(
398 interval: str,
399 minutes: Optional[int] = None,
400 tasklist: Optional[TaskRegister] = None,
401):
402 """Register the given task as a scheduled task.
404 Example:
405 ```python
406 @scheduled_task(ScheduledTask.DAILY)
407 def my_custom_function():
408 # Perform a custom function once per day
409 ...
410 ```
412 Args:
413 interval (str): The interval at which the task should be run
414 minutes (int, optional): The number of minutes between task runs. Defaults to None.
415 tasklist (TaskRegister, optional): The list the tasks should be registered to. Defaults to None.
417 Returns:
418 _type_: _description_
420 Raises:
421 ValueError: If decorated object is not callable
422 ValueError: If interval is not valid
423 """
425 def _task_wrapper(admin_class):
426 if not isinstance(admin_class, Callable): 426 ↛ 427line 426 didn't jump to line 427 because the condition on line 426 was never true
427 raise ValueError('Wrapped object must be a function')
429 if interval not in ScheduledTask.TYPE: 429 ↛ 430line 429 didn't jump to line 430 because the condition on line 429 was never true
430 raise ValueError(f'Invalid interval. Must be one of {ScheduledTask.TYPE}')
432 _tasks = tasklist if tasklist else tasks
433 _tasks.register(admin_class, interval, minutes=minutes)
435 return admin_class
437 return _task_wrapper
440@tracer.start_as_current_span('rebuild_model_tree')
441def rebuild_model_tree(model: str, tree_id: int) -> None:
442 """Rebuild the tree structure (and pathstring values) for a tree model.
444 This task is offloaded to the background worker whenever nodes are
445 restructured (e.g. re-parented), to avoid expensive tree rebuild
446 operations blocking the calling thread.
448 Arguments:
449 model: Label of the model class to rebuild, e.g. 'stock.stocklocation'
450 tree_id: ID of the tree to rebuild
451 """
452 from django.apps import apps
454 import InvenTree.models
456 try:
457 model_class = apps.get_model(model)
458 except (LookupError, ValueError):
459 logger.warning("rebuild_model_tree: Model '%s' does not exist", model)
460 return
462 if not issubclass(model_class, InvenTree.models.InvenTreeTree):
463 logger.warning("rebuild_model_tree: Model '%s' is not a tree model", model)
464 return
466 # Rebuild the tree structure, based on the parent-child relationships
467 model_class.rebuild_trees([tree_id])
469 # Rebuild the 'pathstring' values for the entire tree (if applicable)
470 if issubclass(model_class, InvenTree.models.PathStringMixin):
471 model_class.rebuild_tree_pathstring_values([tree_id])
474@tracer.start_as_current_span('heartbeat')
475@scheduled_task(ScheduledTask.MINUTES, 1)
476def heartbeat():
477 """Simple task which runs at 1 minute intervals, so we can determine that the background worker is actually running."""
478 try:
479 from django_q.models import OrmQ, Success
480 except AppRegistryNotReady: # pragma: no cover
481 logger.info('Could not perform heartbeat task - App registry not ready')
482 return
484 # Write a timestamp file so that health checks can verify worker liveness
485 # without needing to start a full Django process.
486 import tempfile
487 from pathlib import Path
489 try:
490 Path(tempfile.gettempdir()).joinpath('inventree_worker_heartbeat').write_text(
491 str(timezone.now().timestamp())
492 )
493 except Exception:
494 pass
496 threshold = timezone.now() - timedelta(minutes=15)
498 # Delete heartbeat results more than 15 minutes old,
499 # otherwise they just create extra noise
500 heartbeats = Success.objects.filter(
501 func='InvenTree.tasks.heartbeat', started__lte=threshold
502 )
504 heartbeats.delete()
506 # Clear out any other pending heartbeat tasks
507 for task in OrmQ.objects.all():
508 if task.func() == 'InvenTree.tasks.heartbeat':
509 task.delete()
512@tracer.start_as_current_span('delete_successful_tasks')
513@scheduled_task(ScheduledTask.DAILY)
514def delete_successful_tasks():
515 """Delete successful task logs which are older than a specified period."""
516 try:
517 from django_q.models import Success
519 days = get_global_setting('INVENTREE_DELETE_TASKS_DAYS', 30)
520 threshold = timezone.now() - timedelta(days=days)
522 # Delete successful tasks
523 results = Success.objects.filter(started__lte=threshold)
525 if results.count() > 0:
526 logger.info('Deleting %s successful task records', results.count())
527 results.delete()
529 except AppRegistryNotReady: # pragma: no cover
530 logger.info(
531 "Could not perform 'delete_successful_tasks' - App registry not ready"
532 )
535@tracer.start_as_current_span('delete_failed_tasks')
536@scheduled_task(ScheduledTask.DAILY)
537def delete_failed_tasks():
538 """Delete failed task logs which are older than a specified period."""
539 try:
540 from django_q.models import Failure
542 days = get_global_setting('INVENTREE_DELETE_TASKS_DAYS', 30)
543 threshold = timezone.now() - timedelta(days=days)
545 # Delete failed tasks
546 results = Failure.objects.filter(started__lte=threshold)
548 if results.count() > 0:
549 logger.info('Deleting %s failed task records', results.count())
550 results.delete()
552 except AppRegistryNotReady: # pragma: no cover
553 logger.info("Could not perform 'delete_failed_tasks' - App registry not ready")
556@tracer.start_as_current_span('delete_old_error_logs')
557@scheduled_task(ScheduledTask.DAILY)
558def delete_old_error_logs():
559 """Delete old error logs from the server."""
560 try:
561 from error_report.models import Error
563 days = get_global_setting('INVENTREE_DELETE_ERRORS_DAYS', 30)
564 threshold = timezone.now() - timedelta(days=days)
566 errors = Error.objects.filter(when__lte=threshold)
568 if errors.count() > 0:
569 logger.info('Deleting %s old error logs', errors.count())
570 errors.delete()
572 except AppRegistryNotReady: # pragma: no cover
573 # Apps not yet loaded
574 logger.info(
575 "Could not perform 'delete_old_error_logs' - App registry not ready"
576 )
579@tracer.start_as_current_span('delete_old_notifications')
580@scheduled_task(ScheduledTask.DAILY)
581def delete_old_notifications():
582 """Delete old notification logs."""
583 try:
584 from common.models import NotificationEntry, NotificationMessage
586 days = get_global_setting('INVENTREE_DELETE_NOTIFICATIONS_DAYS', 30)
587 threshold = timezone.now() - timedelta(days=days)
589 items = NotificationEntry.objects.filter(updated__lte=threshold)
591 if items.count() > 0:
592 logger.info('Deleted %s old notification entries', items.count())
593 items.delete()
595 items = NotificationMessage.objects.filter(creation__lte=threshold)
597 if items.count() > 0:
598 logger.info('Deleted %s old notification messages', items.count())
599 items.delete()
601 except AppRegistryNotReady:
602 logger.info(
603 "Could not perform 'delete_old_notifications' - App registry not ready"
604 )
607@tracer.start_as_current_span('delete_old_emails')
608@scheduled_task(ScheduledTask.DAILY)
609def delete_old_emails():
610 """Delete old email messages."""
611 try:
612 from common.models import EmailMessage
614 days = get_global_setting('INVENTREE_DELETE_EMAIL_DAYS', 30)
615 threshold = timezone.now() - timedelta(days=days)
617 emails = EmailMessage.objects.filter(timestamp__lte=threshold)
619 if emails.count() > 0:
620 try:
621 emails.delete()
622 logger.info('Deleted %s old email messages', emails.count())
623 except ValidationError:
624 logger.info(
625 'Did not delete %s old email messages because of a validation error',
626 emails.count(),
627 )
629 except AppRegistryNotReady: # pragma: no cover
630 logger.info("Could not perform 'delete_old_emails' - App registry not ready")
633@tracer.start_as_current_span('check_for_updates')
634@scheduled_task(ScheduledTask.DAILY)
635def check_for_updates():
636 """Check if there is an update for InvenTree."""
637 try:
638 from common.notifications import trigger_notification
639 from plugin.builtin.integration.core_notifications import (
640 InvenTreeUINotifications,
641 )
642 except AppRegistryNotReady: # pragma: no cover
643 # Apps not yet loaded!
644 logger.info("Could not perform 'check_for_updates' - App registry not ready")
645 return
647 interval = int(
648 get_global_setting('INVENTREE_UPDATE_CHECK_INTERVAL', 7, cache=False)
649 )
651 # Check if we should check for updates *today*
652 if not check_daily_holdoff('check_for_updates', interval):
653 return
655 logger.info('Checking for InvenTree software updates')
657 headers = {}
659 # If running within github actions, use authentication token
660 if settings.TESTING:
661 token = os.getenv('GITHUB_TOKEN', None)
663 if token:
664 headers['Authorization'] = f'Bearer {token}'
666 response = requests.get(
667 'https://api.github.com/repos/inventree/inventree/releases/latest',
668 headers=headers,
669 )
671 if response.status_code != 200:
672 raise ValueError(
673 f'Unexpected status code from GitHub API: {response.status_code}'
674 ) # pragma: no cover
676 data = json.loads(response.text)
678 tag = data.get('tag_name', None)
680 if not tag:
681 raise ValueError("'tag_name' missing from GitHub response") # pragma: no cover
683 match = re.match(r'^.*(\d+)\.(\d+)\.(\d+).*$', tag)
685 if not match or len(match.groups()) != 3: # pragma: no cover
686 logger.warning("Version '%s' did not match expected pattern", tag)
687 return
689 latest_version = [int(x) for x in match.groups()]
691 if len(latest_version) != 3:
692 raise ValueError(f"Version '{tag}' is not correct format") # pragma: no cover
694 logger.info("Latest InvenTree version: '%s'", tag)
696 # Save the version to the database
697 set_global_setting('_INVENTREE_LATEST_VERSION', tag, None)
699 # Record that this task was successful
700 record_task_success('check_for_updates')
702 # Send notification if there is a new version
703 if not isInvenTreeUpToDate():
704 # Send notification to superusers
705 trigger_notification(
706 None,
707 'update_available',
708 targets=get_user_model().objects.filter(is_superuser=True),
709 delivery_methods={InvenTreeUINotifications},
710 context={
711 'name': _('Update Available'),
712 'message': _('An update for InvenTree is available'),
713 },
714 )
717@tracer.start_as_current_span('update_exchange_rates')
718@scheduled_task(ScheduledTask.DAILY)
719def update_exchange_rates(force: bool = False):
720 """Update currency exchange rates.
722 Arguments:
723 force: If True, force the update to run regardless of the last update time
724 """
725 from InvenTree.ready import canAppAccessDatabase, isRunningMigrations
727 if isRunningMigrations(): 727 ↛ 728line 727 didn't jump to line 728 because the condition on line 727 was never true
728 return
730 # Do not update exchange rates if we cannot access the database
731 if not canAppAccessDatabase(allow_test=True, allow_shell=True): 731 ↛ 732line 731 didn't jump to line 732 because the condition on line 731 was never true
732 return
734 try:
735 from djmoney.contrib.exchange.models import Rate
737 from common.currency import currency_code_default, currency_codes
738 from InvenTree.exchange import InvenTreeExchange
739 except AppRegistryNotReady: # pragma: no cover
740 # Apps not yet loaded!
741 logger.info(
742 "Could not perform 'update_exchange_rates' - App registry not ready"
743 )
744 return
745 except Exception as exc: # pragma: no cover
746 logger.info("Could not perform 'update_exchange_rates' - %s", exc)
747 return
749 if not force: 749 ↛ 750line 749 didn't jump to line 750 because the condition on line 749 was never true
750 interval = int(get_global_setting('CURRENCY_UPDATE_INTERVAL', 1, cache=False))
752 if not check_daily_holdoff('update_exchange_rates', interval):
753 logger.info('Skipping exchange rate update (interval not reached)')
754 return
756 backend = InvenTreeExchange()
757 base = currency_code_default()
758 logger.info("Updating exchange rates using base currency '%s'", base)
760 try:
761 backend.update_rates(base_currency=base)
763 # Remove any exchange rates which are not in the provided currencies
764 Rate.objects.filter(backend='InvenTreeExchange').exclude(
765 currency__in=currency_codes()
766 ).delete()
768 # Record successful task execution
769 record_task_success('update_exchange_rates')
771 except (AppRegistryNotReady, OperationalError, ProgrammingError):
772 logger.warning('Could not update exchange rates - database not ready')
773 except Exception as e: # pragma: no cover
774 logger.exception('Error updating exchange rates: %s', type(e))
777@tracer.start_as_current_span('run_backup')
778@scheduled_task(ScheduledTask.DAILY)
779def run_backup():
780 """Run the backup command."""
781 if not get_global_setting('INVENTREE_BACKUP_ENABLE', False, cache=False):
782 # Backups are not enabled - exit early
783 return
785 interval = int(get_global_setting('INVENTREE_BACKUP_DAYS', 1, cache=False))
787 # Check if should run this task *today*
788 if not check_daily_holdoff('run_backup', interval):
789 return
791 logger.info('Performing automated database backup task')
793 call_command('dbbackup', noinput=True, clean=True, compress=True, interactive=False)
794 call_command(
795 'mediabackup', noinput=True, clean=True, compress=True, interactive=False
796 )
798 # Record that this task was successful
799 record_task_success('run_backup')
802def get_migration_plan():
803 """Returns a list of migrations which are needed to be run."""
804 executor = MigrationExecutor(connections[DEFAULT_DB_ALIAS])
805 plan = executor.migration_plan(executor.loader.graph.leaf_nodes())
806 return plan
809def get_migration_count():
810 """Returns the number of all detected migrations."""
811 executor = MigrationExecutor(connections[DEFAULT_DB_ALIAS])
812 return executor.loader.applied_migrations
815@tracer.start_as_current_span('check_for_migrations')
816@scheduled_task(ScheduledTask.DAILY)
817def check_for_migrations(force: bool = False, reload_registry: bool = True) -> bool:
818 """Checks if migrations are needed.
820 If the setting auto_update is enabled we will start updating.
822 Returns bool indicating if migrations are up to date
823 """
824 from . import ready
826 if ready.isRunningMigrations() or ready.isRunningBackup(): 826 ↛ 828line 826 didn't jump to line 828 because the condition on line 826 was never true
827 # Migrations are already running!
828 return False
830 def set_pending_migrations(n: int):
831 """Helper function to inform the user about pending migrations."""
832 logger.info('There are %s pending migrations', n)
834 try:
835 set_global_setting('_PENDING_MIGRATIONS', n, None)
836 except Exception:
837 # If the setting cannot be set, we just log a warning
838 logger.error('Could not clear _PENDING_MIGRATIONS flag')
840 logger.info('Checking for pending database migrations')
842 if reload_registry: 842 ↛ 846line 842 didn't jump to line 846 because the condition on line 842 was always true
843 # Force plugin registry reload
844 registry.check_reload()
846 plan = get_migration_plan()
848 n = len(plan)
850 # Check if there are any open migrations
851 if not plan:
852 set_pending_migrations(0)
853 return True
855 set_pending_migrations(n)
857 # Test if auto-updates are enabled
858 if not force and not settings.AUTO_UPDATE: 858 ↛ 859line 858 didn't jump to line 859 because the condition on line 858 was never true
859 logger.info('Auto-update is disabled - skipping migrations')
860 return False
862 # Log open migrations
863 for migration in plan:
864 logger.info('- %s', migration[0])
866 # Set the application to maintenance mode - no access from now on.
867 set_maintenance_mode(True)
869 # To be sure we are in maintenance this is wrapped
870 with maintenance_mode_on():
871 logger.info('Starting migration process...')
873 try:
874 call_command('migrate', interactive=False)
875 except NotSupportedError as e: # pragma: no cover
876 if settings.DATABASES['default']['ENGINE'] != 'django.db.backends.sqlite3':
877 raise e
878 logger.exception('Error during migrations: %s', e)
879 except Exception as e: # pragma: no cover
880 logger.exception('Error during migrations: %s', e)
881 else:
882 set_pending_migrations(0)
884 logger.info('Completed %s migrations', n)
886 # Make sure we are out of maintenance mode
887 if get_maintenance_mode(): 887 ↛ 888line 887 didn't jump to line 888 because the condition on line 887 was never true
888 logger.warning('Maintenance mode was not disabled - forcing it now')
889 set_maintenance_mode(False)
890 logger.info('Manually released maintenance mode')
892 if reload_registry: 892 ↛ 897line 892 didn't jump to line 897 because the condition on line 892 was always true
893 # We should be current now - triggering full reload to make sure all models
894 # are loaded fully in their new state.
895 registry.reload_plugins(full_reload=True, force_reload=True, collect=True)
897 return True
900def email_user(user_id: int, subject: str, message: str) -> None:
901 """Send a message to a user."""
902 try:
903 user = get_user_model().objects.get(pk=user_id)
904 except Exception:
905 logger.warning('User <%s> not found - cannot send welcome message', user_id)
906 return
908 from InvenTree.helpers_email import get_email_for_user, send_email
910 if email := get_email_for_user(user):
911 send_email(subject, message, [email])
914@tracer.start_as_current_span('run_oauth_maintenance')
915@scheduled_task(ScheduledTask.DAILY)
916def run_oauth_maintenance():
917 """Run the OAuth maintenance task(s)."""
918 from oauth2_provider.models import clear_expired
920 logger.info('Starting OAuth maintenance task')
921 clear_expired()
922 logger.info('Completed OAuth maintenance task')