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

1"""Functions for tasks and a few general async tasks.""" 

2 

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 

11 

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 _ 

21 

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 

30 

31from common.settings import get_global_setting, set_global_setting 

32from plugin import registry 

33 

34from .version import isInvenTreeUpToDate 

35 

36logger = structlog.get_logger('inventree') 

37tracer = trace.get_tracer(__name__) 

38 

39 

40def schedule_task(taskname, **kwargs): 

41 """Create a scheduled task. 

42 

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 

48 

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 

54 

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) 

60 

61 Schedule.objects.filter(func=taskname).update(**kwargs) 

62 else: 

63 logger.info("Creating scheduled task '%s'", taskname) 

64 

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 

69 

70 

71def raise_warning(msg): 

72 """Log and raise a warning.""" 

73 logger.warning(msg) 

74 

75 # If testing is running raise a warning that can be asserted 

76 if settings.TESTING: 

77 warnings.warn(msg, stacklevel=2) 

78 

79 

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. 

82 

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) 

86 

87 Returns: 

88 bool: If the task should be run *now*, or wait another day 

89 

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. 

92 

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 

101 

102 attempt_key = f'_{task_name}_ATTEMPT' 

103 success_key = f'_{task_name}_SUCCESS' 

104 

105 # Check for recent success information 

106 last_success = get_global_setting(success_key, '', cache=False) 

107 

108 if last_success: 

109 try: 

110 last_success = datetime.fromisoformat(last_success) 

111 except ValueError: 

112 last_success = None 

113 

114 if last_success: 

115 threshold = datetime.now() - timedelta(days=n_days) 

116 

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 

122 

123 # Check for any information we have about this task 

124 last_attempt = get_global_setting(attempt_key, '', cache=False) 

125 

126 if last_attempt: 

127 try: 

128 last_attempt = datetime.fromisoformat(last_attempt) 

129 except ValueError: 

130 last_attempt = None 

131 

132 if last_attempt: 

133 # Do not attempt if the most recent *attempt* was within 12 hours 

134 threshold = datetime.now() - timedelta(hours=12) 

135 

136 if last_attempt > threshold: 

137 logger.info( 

138 "Last attempt for '%s' was too recent - skipping task", task_name 

139 ) 

140 return False 

141 

142 # Record this attempt 

143 record_task_attempt(task_name) 

144 

145 # No reason *not* to run this task now 

146 return True 

147 

148 

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) 

152 

153 set_global_setting(f'_{task_name}_ATTEMPT', datetime.now().isoformat(), None) 

154 

155 

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) 

159 

160 

161def check_existing_task(taskname, group: str, *args, **kwargs) -> Optional[str]: 

162 """Test if an identical task is already registered with the worker. 

163 

164 This will only return true if the task name, group, args and kwargs all match an existing task. 

165 

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 

171 

172 Returns: 

173 Optional[str]: The ID of the matching task, if found, otherwise None 

174 """ 

175 from django_q.models import OrmQ 

176 

177 task_id = None 

178 

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 

184 

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 

188 

189 if task.args() != args: 

190 # Task args do not match 

191 continue 

192 

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 

196 

197 task_id = task.task_id() 

198 

199 break 

200 

201 return task_id 

202 

203 

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! 

213 

214 If workers are not running or force_sync flag, is set then the task is ran synchronously. 

215 

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 

223 

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 

228 

229 # Extract group information from kwargs 

230 group = kwargs.pop('group', 'inventree') 

231 

232 try: 

233 import importlib 

234 

235 from django_q.tasks import AsyncTask 

236 

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) 

240 

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

248 

249 if force_async: 

250 # Cannot async the task, so return False 

251 return False 

252 else: 

253 force_sync = True 

254 

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 ) 

262 

263 return task_id 

264 

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

270 

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 

295 

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 

304 

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 

312 

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 

321 

322 # Finally, task either completed successfully or was offloaded 

323 return True 

324 

325 

326def get_queued_task(task_id: str): 

327 """Find the task in the queue, if it exists. 

328 

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 

336 

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 

340 

341 task_id = str(task_id) 

342 

343 from django_q.models import OrmQ 

344 

345 while True: 

346 queued_tasks = OrmQ.objects.all().order_by('id')[offset : offset + limit] 

347 if not queued_tasks: 

348 break 

349 

350 for task in queued_tasks: 

351 if task.task_id() == task_id: 

352 return task 

353 

354 offset += limit 

355 

356 # No matching task was discovered 

357 return None 

358 

359 

360@dataclass() 

361class ScheduledTask: 

362 """A scheduled task. 

363 

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

368 

369 func: Callable 

370 interval: str 

371 minutes: Optional[int] = None 

372 

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' 

380 

381 TYPE: tuple[str] = (MINUTES, HOURLY, DAILY, WEEKLY, MONTHLY, QUARTERLY, YEARLY) 

382 

383 

384class TaskRegister: 

385 """Registry for periodic tasks.""" 

386 

387 task_list: list[ScheduledTask] = [] 

388 

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

392 

393 

394tasks = TaskRegister() 

395 

396 

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. 

403 

404 Example: 

405 ```python 

406 @scheduled_task(ScheduledTask.DAILY) 

407 def my_custom_function(): 

408 # Perform a custom function once per day 

409 ... 

410 ``` 

411 

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. 

416 

417 Returns: 

418 _type_: _description_ 

419 

420 Raises: 

421 ValueError: If decorated object is not callable 

422 ValueError: If interval is not valid 

423 """ 

424 

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

428 

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

431 

432 _tasks = tasklist if tasklist else tasks 

433 _tasks.register(admin_class, interval, minutes=minutes) 

434 

435 return admin_class 

436 

437 return _task_wrapper 

438 

439 

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. 

443 

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. 

447 

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 

453 

454 import InvenTree.models 

455 

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 

461 

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 

465 

466 # Rebuild the tree structure, based on the parent-child relationships 

467 model_class.rebuild_trees([tree_id]) 

468 

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

472 

473 

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 

483 

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 

488 

489 try: 

490 Path(tempfile.gettempdir()).joinpath('inventree_worker_heartbeat').write_text( 

491 str(timezone.now().timestamp()) 

492 ) 

493 except Exception: 

494 pass 

495 

496 threshold = timezone.now() - timedelta(minutes=15) 

497 

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 ) 

503 

504 heartbeats.delete() 

505 

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

510 

511 

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 

518 

519 days = get_global_setting('INVENTREE_DELETE_TASKS_DAYS', 30) 

520 threshold = timezone.now() - timedelta(days=days) 

521 

522 # Delete successful tasks 

523 results = Success.objects.filter(started__lte=threshold) 

524 

525 if results.count() > 0: 

526 logger.info('Deleting %s successful task records', results.count()) 

527 results.delete() 

528 

529 except AppRegistryNotReady: # pragma: no cover 

530 logger.info( 

531 "Could not perform 'delete_successful_tasks' - App registry not ready" 

532 ) 

533 

534 

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 

541 

542 days = get_global_setting('INVENTREE_DELETE_TASKS_DAYS', 30) 

543 threshold = timezone.now() - timedelta(days=days) 

544 

545 # Delete failed tasks 

546 results = Failure.objects.filter(started__lte=threshold) 

547 

548 if results.count() > 0: 

549 logger.info('Deleting %s failed task records', results.count()) 

550 results.delete() 

551 

552 except AppRegistryNotReady: # pragma: no cover 

553 logger.info("Could not perform 'delete_failed_tasks' - App registry not ready") 

554 

555 

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 

562 

563 days = get_global_setting('INVENTREE_DELETE_ERRORS_DAYS', 30) 

564 threshold = timezone.now() - timedelta(days=days) 

565 

566 errors = Error.objects.filter(when__lte=threshold) 

567 

568 if errors.count() > 0: 

569 logger.info('Deleting %s old error logs', errors.count()) 

570 errors.delete() 

571 

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 ) 

577 

578 

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 

585 

586 days = get_global_setting('INVENTREE_DELETE_NOTIFICATIONS_DAYS', 30) 

587 threshold = timezone.now() - timedelta(days=days) 

588 

589 items = NotificationEntry.objects.filter(updated__lte=threshold) 

590 

591 if items.count() > 0: 

592 logger.info('Deleted %s old notification entries', items.count()) 

593 items.delete() 

594 

595 items = NotificationMessage.objects.filter(creation__lte=threshold) 

596 

597 if items.count() > 0: 

598 logger.info('Deleted %s old notification messages', items.count()) 

599 items.delete() 

600 

601 except AppRegistryNotReady: 

602 logger.info( 

603 "Could not perform 'delete_old_notifications' - App registry not ready" 

604 ) 

605 

606 

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 

613 

614 days = get_global_setting('INVENTREE_DELETE_EMAIL_DAYS', 30) 

615 threshold = timezone.now() - timedelta(days=days) 

616 

617 emails = EmailMessage.objects.filter(timestamp__lte=threshold) 

618 

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 ) 

628 

629 except AppRegistryNotReady: # pragma: no cover 

630 logger.info("Could not perform 'delete_old_emails' - App registry not ready") 

631 

632 

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 

646 

647 interval = int( 

648 get_global_setting('INVENTREE_UPDATE_CHECK_INTERVAL', 7, cache=False) 

649 ) 

650 

651 # Check if we should check for updates *today* 

652 if not check_daily_holdoff('check_for_updates', interval): 

653 return 

654 

655 logger.info('Checking for InvenTree software updates') 

656 

657 headers = {} 

658 

659 # If running within github actions, use authentication token 

660 if settings.TESTING: 

661 token = os.getenv('GITHUB_TOKEN', None) 

662 

663 if token: 

664 headers['Authorization'] = f'Bearer {token}' 

665 

666 response = requests.get( 

667 'https://api.github.com/repos/inventree/inventree/releases/latest', 

668 headers=headers, 

669 ) 

670 

671 if response.status_code != 200: 

672 raise ValueError( 

673 f'Unexpected status code from GitHub API: {response.status_code}' 

674 ) # pragma: no cover 

675 

676 data = json.loads(response.text) 

677 

678 tag = data.get('tag_name', None) 

679 

680 if not tag: 

681 raise ValueError("'tag_name' missing from GitHub response") # pragma: no cover 

682 

683 match = re.match(r'^.*(\d+)\.(\d+)\.(\d+).*$', tag) 

684 

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 

688 

689 latest_version = [int(x) for x in match.groups()] 

690 

691 if len(latest_version) != 3: 

692 raise ValueError(f"Version '{tag}' is not correct format") # pragma: no cover 

693 

694 logger.info("Latest InvenTree version: '%s'", tag) 

695 

696 # Save the version to the database 

697 set_global_setting('_INVENTREE_LATEST_VERSION', tag, None) 

698 

699 # Record that this task was successful 

700 record_task_success('check_for_updates') 

701 

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 ) 

715 

716 

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. 

721 

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 

726 

727 if isRunningMigrations(): 727 ↛ 728line 727 didn't jump to line 728 because the condition on line 727 was never true

728 return 

729 

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 

733 

734 try: 

735 from djmoney.contrib.exchange.models import Rate 

736 

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 

748 

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

751 

752 if not check_daily_holdoff('update_exchange_rates', interval): 

753 logger.info('Skipping exchange rate update (interval not reached)') 

754 return 

755 

756 backend = InvenTreeExchange() 

757 base = currency_code_default() 

758 logger.info("Updating exchange rates using base currency '%s'", base) 

759 

760 try: 

761 backend.update_rates(base_currency=base) 

762 

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

767 

768 # Record successful task execution 

769 record_task_success('update_exchange_rates') 

770 

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

775 

776 

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 

784 

785 interval = int(get_global_setting('INVENTREE_BACKUP_DAYS', 1, cache=False)) 

786 

787 # Check if should run this task *today* 

788 if not check_daily_holdoff('run_backup', interval): 

789 return 

790 

791 logger.info('Performing automated database backup task') 

792 

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 ) 

797 

798 # Record that this task was successful 

799 record_task_success('run_backup') 

800 

801 

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 

807 

808 

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 

813 

814 

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. 

819 

820 If the setting auto_update is enabled we will start updating. 

821 

822 Returns bool indicating if migrations are up to date 

823 """ 

824 from . import ready 

825 

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 

829 

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) 

833 

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

839 

840 logger.info('Checking for pending database migrations') 

841 

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

845 

846 plan = get_migration_plan() 

847 

848 n = len(plan) 

849 

850 # Check if there are any open migrations 

851 if not plan: 

852 set_pending_migrations(0) 

853 return True 

854 

855 set_pending_migrations(n) 

856 

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 

861 

862 # Log open migrations 

863 for migration in plan: 

864 logger.info('- %s', migration[0]) 

865 

866 # Set the application to maintenance mode - no access from now on. 

867 set_maintenance_mode(True) 

868 

869 # To be sure we are in maintenance this is wrapped 

870 with maintenance_mode_on(): 

871 logger.info('Starting migration process...') 

872 

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) 

883 

884 logger.info('Completed %s migrations', n) 

885 

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

891 

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) 

896 

897 return True 

898 

899 

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 

907 

908 from InvenTree.helpers_email import get_email_for_user, send_email 

909 

910 if email := get_email_for_user(user): 

911 send_email(subject, message, [email]) 

912 

913 

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 

919 

920 logger.info('Starting OAuth maintenance task') 

921 clear_expired() 

922 logger.info('Completed OAuth maintenance task')