From 6fad94afb64cb6a738cbc1353493c4a5be2da839 Mon Sep 17 00:00:00 2001 From: ebuzerdrmz44 Date: Mon, 7 Sep 2026 21:48:49 +0300 Subject: [PATCH] Split VortaScheduler into state and scheduling components and stop holding the lock across job submission --- src/vorta/scheduler.py | 743 ------------------------------ src/vorta/scheduler/__init__.py | 320 +++++++++++++ src/vorta/scheduler/scheduling.py | 352 ++++++++++++++ src/vorta/scheduler/state.py | 246 ++++++++++ tests/unit/test_schedule.py | 5 +- tests/unit/test_scheduler.py | 81 +++- 6 files changed, 1001 insertions(+), 746 deletions(-) delete mode 100644 src/vorta/scheduler.py create mode 100644 src/vorta/scheduler/__init__.py create mode 100644 src/vorta/scheduler/scheduling.py create mode 100644 src/vorta/scheduler/state.py diff --git a/src/vorta/scheduler.py b/src/vorta/scheduler.py deleted file mode 100644 index a4887b0d1..000000000 --- a/src/vorta/scheduler.py +++ /dev/null @@ -1,743 +0,0 @@ -from __future__ import annotations - -import enum -import logging -import threading -from collections.abc import Callable -from datetime import datetime as dt -from datetime import timedelta -from typing import Any, NamedTuple - -import peewee as pw -from packaging import version -from PyQt6 import QtCore, QtDBus -from PyQt6.QtCore import QTimer -from PyQt6.QtWidgets import QApplication - -from vorta import application -from vorta.borg.check import BorgCheckJob -from vorta.borg.compact import BorgCompactJob -from vorta.borg.create import BorgCreateJob -from vorta.borg.list_repo import BorgListRepoJob -from vorta.borg.prune import BorgPruneJob -from vorta.i18n import translate -from vorta.notifications import VortaNotifications -from vorta.store.models import BackupProfileModel, EventLogModel, JobModel, SchedulerPauseModel -from vorta.utils import borg_compat, get_network_status_monitor - -logger = logging.getLogger(__name__) - -RESCHEDULE_INTERVAL_MS = 15 * 60 * 1000 -WAKE_CHECK_INTERVAL_MS = 5 * 60 * 1000 -WAKE_GAP_THRESHOLD = timedelta(minutes=10) - -# A QTimer interval is a C++ int, so a single wait tops out at about 24.8 days. -MAX_TIMER_MS = 2**31 - 1 -# Fire just after the deadline, so the handler always sees it as passed. -TIMER_GRACE_MS = 100 - - -def arm_deadline_timer(deadline: dt, on_expiry: Callable[[], None]) -> QTimer: - """ - Start a timer that calls `on_expiry` once `deadline` has passed. - - Waits longer than `MAX_TIMER_MS` are split into chunks, and every chunk is measured - against the wall clock again, so the deadline holds however far ahead it is. - """ - timer = QTimer() - timer.setSingleShot(True) - - def remaining_ms() -> float: - return (deadline - dt.now()).total_seconds() * 1000 + TIMER_GRACE_MS - - def rearm() -> None: - timer.setInterval(int(max(0, min(remaining_ms(), MAX_TIMER_MS)))) - timer.start() - - def on_timeout() -> None: - if remaining_ms() <= 0: - on_expiry() - else: - rearm() - - timer.timeout.connect(on_timeout) - rearm() - return timer - - -class ScheduleStatusType(enum.Enum): - SCHEDULED = enum.auto() # date provided - UNSCHEDULED = enum.auto() # Unknown - NO_PREVIOUS_BACKUP = enum.auto() # run a manual backup first - PAUSED = enum.auto() # paused after a failed or skipped run, date provided - - -class ScheduleStatus(NamedTuple): - type: ScheduleStatusType - time: dt | None = None - - -#: Timer states that hold a time for a run, and the row status each one shows as. -PENDING_STATUSES = { - ScheduleStatusType.SCHEDULED: JobModel.Status.SCHEDULED.value, - ScheduleStatusType.PAUSED: JobModel.Status.PAUSED.value, -} - - -class PendingJob(NamedTuple): - """A run the scheduler is holding a time for, but that hasn't been recorded yet.""" - - profile_id: int - profile_name: str - repo_url: str | None - scheduled_at: dt - status: str - - -class VortaScheduler(QtCore.QObject): - #: The schedule for a profile changed. - schedule_changed = QtCore.pyqtSignal() - - #: A job outcome was recorded. - jobs_changed = QtCore.pyqtSignal() - - def __init__(self) -> None: - super().__init__() - - #: mapping of profiles to timers - self.timers: dict[int, dict[str, QTimer | dt | ScheduleStatusType | None]] = dict() - - self.app: application.VortaApp = QApplication.instance() - self.lock = threading.Lock() - - # pausing will prevent scheduling for a specified time - self.pauses: dict[int, tuple[dt, QtCore.QTimer]] = dict() - - # Periodic reschedule, in case a run was missed - self.qt_timer = QTimer() - self.qt_timer.timeout.connect(self.reload_all_timers) - self.qt_timer.setInterval(RESCHEDULE_INTERVAL_MS) - self.qt_timer.start() - - # connect signals - self.app.backup_finished_event.connect(lambda res: self.set_timer_for_profile(res['params']['profile_id'])) - - # Connect to network manager to monitor net status - self.net_status = get_network_status_monitor() - self.net_status.network_status_changed.connect(self.networkStatusChanged) - self._net_up = self.net_status.is_network_active() - - # connect to `systemd-logind` to receive sleep/resume events - # The signal `PrepareForSleep` will be emitted before and after hibernation. - service = "org.freedesktop.login1" - path = "/org/freedesktop/login1" - interface = "org.freedesktop.login1.Manager" - name = "PrepareForSleep" - bus = QtDBus.QDBusConnection.systemBus() - if bus.isConnected() and bus.interface().isServiceRegistered(service).value(): - self.bus = bus - self.bus.connect(service, path, interface, name, "b", self.loginSuspendNotify) - else: - logger.info('No systemd-logind to notify us of sleep/resume, watching for clock gaps as well') - - self._last_wake_check = dt.now() - self.wake_timer = QTimer() - self.wake_timer.timeout.connect(self.checkForResume) - self.wake_timer.setInterval(WAKE_CHECK_INTERVAL_MS) - self.wake_timer.start() - - self._restore_pauses() - - @QtCore.pyqtSlot(bool) - def loginSuspendNotify(self, suspend: bool) -> None: - if not suspend: - logger.debug("Got login suspend/resume notification") - self._handle_resume() - - @QtCore.pyqtSlot() - def checkForResume(self) -> None: - now = dt.now() - elapsed = now - self._last_wake_check - self._last_wake_check = now - - if elapsed < WAKE_GAP_THRESHOLD: - return - - logger.debug('Clock jumped %s since the last wake check, assuming the machine slept', elapsed) - self._handle_resume() - - def _handle_resume(self) -> None: - self._last_wake_check = dt.now() - # Defensively refetch in case the network status didn't arrive - self._net_up = self.net_status.is_network_active() - self.reload_all_timers() - - @QtCore.pyqtSlot(bool) - def networkStatusChanged(self, up: bool): - reload = self._net_up != up - self._net_up = up - logger.debug(f"network status up={up}") - if reload: - logger.info("updating schedule due to network status change") - self.reload_all_timers() - - def tr(self, *args: Any, **kwargs: Any) -> str: - scope = self.__class__.__name__ - return translate(scope, *args, **kwargs) - - def pause(self, profile_id: int, until: dt | None = None) -> None: - """ - Call a timeout for scheduling of a given profile. - - If `until` is omitted, a default time for the break is calculated. - - .. warning:: - This method won't work correctly when called from a non-`QThread`. - - - Parameters - ---------- - profile_id : int - The profile to pause the scheduling for. - until : dt | None, optional - The time to end the pause, by default None - """ - profile = BackupProfileModel.get_or_none(id=profile_id) - if profile is None: # profile doesn't exist any more. - return - - if profile.schedule_mode == 'off': - return - - if until is None: - # calculate default timeout - - if profile.schedule_mode == 'interval': - interval = timedelta(**{profile.schedule_interval_unit: profile.schedule_interval_count}) - else: - # fixed - interval = timedelta(days=1) - - timeout = interval // 6 # 60 / 6 = 10 [min] - timeout = max(min(timeout, timedelta(hours=1)), timedelta(minutes=1)) # 1 <= t <= 60 - - until = dt.now().replace(microsecond=0) + timeout - elif until < dt.now(): - return - - # remove existing schedule - self.remove_job(profile_id) - - # set timeout/pause - other_pause = self.pauses.get(profile_id) - if other_pause is not None: - logger.debug(f"Override existing timeout for profile {profile_id}") - - self._set_pause(profile, until) - logger.debug(f"Paused {profile_id} until {until.strftime('%Y-%m-%d %H:%M:%S')}") - - def _set_pause(self, profile: BackupProfileModel, until: dt) -> None: - """Arm the reschedule timer for a pause and store it, so it outlives the process.""" - profile_id = profile.id - replaced = self.pauses.get(profile_id) - if replaced is not None: - replaced[1].stop() - - # setting timer for reschedule is not possible if called - # from a non-QThread - it won't fail but won't work - timer = arm_deadline_timer(until, lambda: self.set_timer_for_profile(profile_id)) - - self.pauses[profile_id] = (until, timer) - - try: - SchedulerPauseModel.replace(profile=profile_id, paused_until=until).execute() - except pw.PeeweeException: - logger.warning('Could not store pause for profile %s.', profile_id, exc_info=True) - - self._mark_paused(profile, until) - - def _mark_paused(self, profile: BackupProfileModel, until: dt) -> None: - """Report the pause as the schedule status, unless the profile has no schedule to block.""" - if profile.repo is None or profile.schedule_mode == 'off': - return - - self.timers[profile.id] = {'type': ScheduleStatusType.PAUSED, 'dt': until} - self.schedule_changed.emit() - - def clear_pause(self, profile_id: int) -> None: - """Drop a pause from memory, from the schedule status and from the database.""" - pause = self.pauses.pop(profile_id, None) - if pause is not None: - pause[1].stop() - - status = self.timers.get(profile_id) - if status is not None and status.get('type') is ScheduleStatusType.PAUSED: - del self.timers[profile_id] - - try: - SchedulerPauseModel.delete().where(SchedulerPauseModel.profile == profile_id).execute() - except pw.PeeweeException: - logger.warning('Could not drop stored pause for profile %s.', profile_id, exc_info=True) - - def _restore_pauses(self) -> None: - """Re-arm the pauses stored by a previous run, dropping the ones that already ran out.""" - now = dt.now() - - for stored in list(SchedulerPauseModel.select()): - profile_id = stored.profile_id - profile = BackupProfileModel.get_or_none(id=profile_id) - - if profile is None or stored.paused_until <= now: - self.clear_pause(profile_id) - continue - - self._set_pause(profile, stored.paused_until) - logger.debug(f"Restored pause for {profile_id} until {stored.paused_until:%Y-%m-%d %H:%M:%S}") - - def unpause(self, profile_id: int) -> None: - """ - Return to scheduling for a profile. - - Parameters - ---------- - profile_id : int - The profile to end the timeout for. - """ - profile = BackupProfileModel.get_or_none(id=profile_id) - if profile is None: # profile doesn't exist any more. - return - - pause = self.pauses.get(profile_id) - if pause is None: # already unpaused - return - - self.clear_pause(profile_id) - - logger.debug(f"Unpaused {profile_id}") - - self.set_timer_for_profile(profile_id) - - def paused(self, profile_id: int) -> bool: - """ - Determine whether scheduling for a profile is paused - - Parameters - ---------- - profile_id : int - - Returns - ------- - bool - """ - return self.pauses.get(profile_id) is not None - - def set_timer_for_profile(self, profile_id: int) -> None: - """ - Set a timer for next scheduled backup run of this profile. - - Removes existing jobs if set to manual only or no repo is assigned. - - Else will look for previous scheduled backups and catch up if - schedule_make_up_missed is enabled. - - Or, if catch-up is not enabled, will add interval to last run to find - next suitable backup time. - """ - profile = BackupProfileModel.get_or_none(id=profile_id) - if profile is None: # profile doesn't exist any more. - return - logger.debug('Profile: %s, %d %d', str(profile), profile.schedule_fixed_hour, profile.schedule_fixed_minute) - - with self.lock: # Acquire lock - self.remove_job(profile_id) # reset schedule - - pause = self.pauses.get(profile_id) - if pause is not None: - pause_end, timer = pause - if dt.now() < pause_end: - logger.debug( - 'Nothing scheduled for profile %s ' + 'because of timeout until %s.', - profile_id, - pause[0].strftime('%Y-%m-%d %H:%M:%S'), - ) - self._mark_paused(profile, pause_end) - return - else: - self.clear_pause(profile_id) - - if profile.repo is None: # No backups without repo set - logger.debug( - 'Nothing scheduled for profile %s because of unset repo.', - profile_id, - ) - # Emit signal so that e.g. the GUI can react to the new schedule - self.schedule_changed.emit() - return - - if profile.schedule_mode == 'off': - logger.debug('Scheduler for profile %s is disabled.', profile_id) - # Emit signal so that e.g. the GUI can react to the new schedule - self.schedule_changed.emit() - return - - logger.info('Setting timer for profile %s', profile_id) - - # determine last backup time - last_run_log = ( - EventLogModel.select() - .where( - EventLogModel.subcommand == 'create', - EventLogModel.category == 'scheduled', - EventLogModel.profile == profile.id, - 0 <= EventLogModel.returncode <= 1, - ) - .order_by(EventLogModel.end_time.desc()) - .first() - ) - - if last_run_log is None: - # look for non scheduled (manual) backup runs - last_run_log = ( - EventLogModel.select() - .where( - EventLogModel.subcommand == 'create', - EventLogModel.profile == profile.id, - 0 <= EventLogModel.returncode <= 1, - ) - .order_by(EventLogModel.end_time.desc()) - .first() - ) - - if last_run_log is None: - logger.info( - f"Nothing scheduled for profile {profile_id} " - + "because it would be the first backup " - + "for this profile." - ) - self.timers[profile_id] = {'type': ScheduleStatusType.NO_PREVIOUS_BACKUP} - # Emit signal so that e.g. the GUI can react to the new schedule - self.schedule_changed.emit() - return - - # calculate next scheduled time - if profile.schedule_mode == 'interval': - last_time: dt = last_run_log.end_time - - interval = {profile.schedule_interval_unit: profile.schedule_interval_count} - next_time = last_time + timedelta(**interval) - - elif profile.schedule_mode == 'fixed': - last_time = last_run_log.end_time - - next_time = last_time.replace( - hour=profile.schedule_fixed_hour, - minute=profile.schedule_fixed_minute, - second=0, - microsecond=0, - ) + timedelta(days=1) - - else: - # unknown schedule mode - raise ValueError("Unknown schedule mode '{}'".format(profile.schedule_mode)) - - logger.debug('Last run time: %s', last_time) - - needs_network = profile.repo is not None and profile.repo.is_remote_repo() - # handle missing of a scheduled time - if next_time <= dt.now(): - if profile.schedule_make_up_missed and (self._net_up or not needs_network): - self.lock.release() - try: - logger.debug( - 'Catching up by running job for %s (%s)', - profile.name, - profile_id, - ) - self.create_backup(profile_id, trigger=JobModel.Trigger.CATCHUP.value) - finally: - self.lock.acquire() # with-statement will try to release - - return # create_backup will lead to a call to this method - elif profile.schedule_make_up_missed and not self._net_up and needs_network: - logger.debug('Skipping catchup %s (%s), the network is not available', profile.name, profile.id) - self._record_skip( - profile, - JobModel.Trigger.CATCHUP.value, - 'Network unavailable for catch-up.', - scheduled_at=next_time, - ) - - # calculate next time from now - if profile.schedule_mode == 'interval': - # next_time % interval should be 0 - # while next_time > now - delta = dt.now() - last_time - next_time = dt.now() - delta % timedelta(**interval) - next_time += timedelta(**interval) - - elif profile.schedule_mode == 'fixed': - # schedule for today - next_time = dt.now().replace( - hour=profile.schedule_fixed_hour, - minute=profile.schedule_fixed_minute, - second=0, - microsecond=0, - ) - - if next_time <= dt.now(): - # time for today has passed, schedule for tomorrow - next_time += timedelta(days=1) - - # start QTimer - logger.debug('Scheduling next run for %s', next_time) - - self.timers[profile_id] = { - 'qtt': arm_deadline_timer(next_time, lambda: self.create_backup(profile_id)), - 'dt': next_time, - 'type': ScheduleStatusType.SCHEDULED, - } - - # Emit signal so that e.g. the GUI can react to the new schedule - self.schedule_changed.emit() - - def reload_all_timers(self) -> None: - logger.debug('Refreshing all scheduler timers') - for profile in BackupProfileModel.select(): - # Only set a timer for the profile if the network is actually up - if profile.repo is None: - logger.debug("nothing scheduled for %s because of unset repo", profile.id) - elif not profile.repo.is_remote_repo() or self._net_up: - logger.debug("scheduling %s", profile.id) - self.set_timer_for_profile(profile.id) - else: - logger.debug("Network is down, not scheduling %s", profile.id) - self.remove_job(profile.id) - - def next_job(self) -> str: - now = dt.now() - - def is_scheduled(timer): - return timer["type"] == ScheduleStatusType.SCHEDULED and timer["qtt"].isActive() and timer["dt"] >= now - - scheduled = {profile_id: timer for profile_id, timer in self.timers.items() if is_scheduled(timer)} - if len(scheduled) == 0: - return self.tr("None scheduled") - - closest_job = min(scheduled.items(), key=lambda item: item[1]["dt"]) - profile_id, timer = closest_job - time = timer["dt"] - profile = BackupProfileModel.get_or_none(id=profile_id) - - time_format = "%H:%M" - if time - now > timedelta(days=1): - time_format = "%b %d, %H:%M" - return f"{time.strftime(time_format)} ({profile.name})" - - def next_job_for_profile(self, profile_id: int) -> ScheduleStatus: - job = self.timers.get(profile_id) - if job is None: - return ScheduleStatus(ScheduleStatusType.UNSCHEDULED) - return ScheduleStatus(job['type'], time=job.get('dt')) # type: ignore[arg-type] - - def pending_jobs(self) -> list[PendingJob]: - """The runs currently on the clock, soonest first.""" - pending = [] - - for profile_id, timer in self.timers.items(): - status = PENDING_STATUSES.get(timer['type']) - if status is None: - continue - - profile = BackupProfileModel.get_or_none(id=profile_id) - if profile is None: - continue - - repo_url = profile.repo.url if profile.repo else None - pending.append(PendingJob(profile_id, profile.name, repo_url, timer['dt'], status)) - - return sorted(pending, key=lambda job: job.scheduled_at) - - def _record_skip( - self, - profile: BackupProfileModel, - trigger: str, - reason: str, - status: str = JobModel.Status.SKIPPED.value, - scheduled_at: dt | None = None, - ) -> None: - """Record a job outcome, deduplicated on the occurrence when one is known.""" - lookup = { - 'profile': str(profile.id), - 'trigger': trigger, - 'status': status, - 'scheduled_at': scheduled_at, - } - details = { - 'profile_name': profile.name, - 'repo_url': profile.repo.url if profile.repo else None, - 'job_type': JobModel.Type.BACKUP.value, - 'reason': reason, - } - - try: - if scheduled_at is None: - JobModel.create(**lookup, **details) - else: - _, recorded = JobModel.get_or_create(**lookup, defaults=details) - if not recorded: - return - except pw.PeeweeException: - logger.warning('Could not record job for profile %s.', profile.id, exc_info=True) - return - - # Two of the call sites hold `self.lock`, and the jobs view reads the table on this signal. - QTimer.singleShot(0, self.jobs_changed.emit) - - def create_backup(self, profile_id: int, trigger: str = JobModel.Trigger.SCHEDULED.value) -> None: - notifier = VortaNotifications.pick() - profile = BackupProfileModel.get_or_none(id=profile_id) - - if profile is None: - logger.info('Profile not found. Maybe deleted?') - return - - # Skip if a job for this profile (repo) is already in progress - if self.app.jobs_manager.is_worker_running(site=profile.repo.id): - logger.debug('A job for repo %s is already active.', profile.repo.id) - self._record_skip(profile, trigger, 'Repository is busy with another job.') - self.pause(profile_id) - return - - with self.lock: - logger.info('Starting background backup for %s', profile.name) - notifier.deliver( - self.tr('Vorta Backup'), - self.tr('Starting background backup for %s.') % profile.name, - level='info', - ) - msg = BorgCreateJob.prepare(profile) - if msg['ok']: - logger.info('Preparation for backup successful.') - msg['category'] = 'scheduled' - job = BorgCreateJob(msg['cmd'], msg, profile.repo.id) - job.result.connect(self.notify) - self.app.jobs_manager.add_job(job) - else: - # Default to 'error': unexpected failures notify. - # Expected skips (WiFi/metered) use 'info' to suppress. - level = msg.get('level', 'error') - if level == 'error': - logger.error('Conditions for backup not met. Aborting.') - logger.error(msg['message']) - notifier.deliver( - self.tr('Vorta Backup'), - translate('messages', msg['message']), - level='error', - ) - status = JobModel.Status.FAILED.value - else: - logger.info('Backup skipped: %s', msg['message']) - status = JobModel.Status.SKIPPED.value - self._record_skip(profile, trigger, msg['message'], status=status) - self.pause(profile_id) - - def notify(self, result: dict[str, Any]) -> None: - notifier = VortaNotifications.pick() - profile_name = result['params']['profile_name'] - profile_id = result['params']['profile'].id - - if result['returncode'] in [0, 1]: - notifier.deliver( - self.tr('Vorta Backup'), - self.tr('Backup successful for %s.') % profile_name, - level='info', - ) - logger.info('Backup creation successful.') - # unpause scheduler - self.unpause(result['params']['profile_id']) - - self.post_backup_tasks(profile_id) - else: - notifier.deliver( - self.tr('Vorta Backup'), - self.tr('Error during backup creation for %s.') % profile_name, - level='error', - ) - logger.error('Error during backup creation.') - # pause scheduler - # if a scheduled backup fails the scheduler should pause - # temporarily. - self.pause(result['params']['profile_id']) - - self.set_timer_for_profile(profile_id) - - def post_backup_tasks(self, profile_id: int) -> None: - """ - Pruning and checking after successful backup. - """ - profile = BackupProfileModel.get(id=profile_id) - notifier = VortaNotifications.pick() - logger.info('Doing post-backup jobs for %s', profile.name) - if profile.prune_on: - msg = BorgPruneJob.prepare(profile) - if msg['ok']: - job = BorgPruneJob(msg['cmd'], msg, profile.repo.id) - self.app.jobs_manager.add_job(job) - - # Refresh archives - msg = BorgListRepoJob.prepare(profile) - if msg['ok']: - job = BorgListRepoJob(msg['cmd'], msg, profile.repo.id) - self.app.jobs_manager.add_job(job) - - validation_cutoff = dt.now() - timedelta(days=7 * profile.validation_weeks) - recent_validations = ( - EventLogModel.select() - .where( - (EventLogModel.subcommand == 'check') - & (EventLogModel.start_time > validation_cutoff) - & (EventLogModel.repo_url == profile.repo.url) - ) - .count() - ) - if profile.validation_on and recent_validations == 0: - msg = BorgCheckJob.prepare(profile) - if msg['ok']: - job = BorgCheckJob(msg['cmd'], msg, profile.repo.id) - self.app.jobs_manager.add_job(job) - - compaction_cutoff = dt.now() - timedelta(days=7 * profile.compaction_weeks) - recent_compactions = ( - EventLogModel.select() - .where( - (EventLogModel.subcommand == '--info') - & (EventLogModel.start_time > compaction_cutoff) - & (EventLogModel.repo_url == profile.repo.url) - ) - .count() - ) - - if ( - profile.compaction_on - and recent_compactions == 0 - and version.parse(borg_compat.version) >= version.parse("1.2") - ): - msg = BorgCompactJob.prepare(profile) - if msg['ok']: - job = BorgCompactJob(msg['cmd'], msg, profile.repo.id) - self.app.jobs_manager.add_job(job) - - logger.info('Finished background task for profile %s', profile.name) - notifier.deliver( - self.tr('Vorta Backup'), - self.tr('Post Backup Tasks successful for %s' % profile.name), - level='info', - ) - - def remove_job(self, profile_id: int) -> None: - if profile_id in self.timers: - qtimer = self.timers[profile_id].get('qtt') - if qtimer is not None: - qtimer.stop() - - del self.timers[profile_id] diff --git a/src/vorta/scheduler/__init__.py b/src/vorta/scheduler/__init__.py new file mode 100644 index 000000000..ae2996a39 --- /dev/null +++ b/src/vorta/scheduler/__init__.py @@ -0,0 +1,320 @@ +from __future__ import annotations + +import logging +from datetime import datetime as dt +from datetime import timedelta +from typing import Any + +from packaging import version +from PyQt6 import QtCore +from PyQt6.QtCore import QTimer +from PyQt6.QtWidgets import QApplication + +from vorta import application +from vorta.borg.check import BorgCheckJob +from vorta.borg.compact import BorgCompactJob +from vorta.borg.create import BorgCreateJob +from vorta.borg.list_repo import BorgListRepoJob +from vorta.borg.prune import BorgPruneJob +from vorta.i18n import translate +from vorta.notifications import VortaNotifications +from vorta.scheduler.scheduling import ( + MAX_TIMER_MS, + PENDING_STATUSES, + RESCHEDULE_INTERVAL_MS, + TIMER_GRACE_MS, + PendingJob, + SchedulerTimers, + ScheduleStatus, + ScheduleStatusType, + arm_deadline_timer, +) +from vorta.scheduler.state import WAKE_CHECK_INTERVAL_MS, WAKE_GAP_THRESHOLD, SchedulerState +from vorta.store.models import BackupProfileModel, EventLogModel, JobModel +from vorta.utils import borg_compat + +logger = logging.getLogger(__name__) + +__all__ = [ + 'MAX_TIMER_MS', + 'PENDING_STATUSES', + 'RESCHEDULE_INTERVAL_MS', + 'TIMER_GRACE_MS', + 'WAKE_CHECK_INTERVAL_MS', + 'WAKE_GAP_THRESHOLD', + 'PendingJob', + 'ScheduleStatus', + 'ScheduleStatusType', + 'VortaScheduler', + 'arm_deadline_timer', +] + + +class VortaScheduler(QtCore.QObject): + #: The schedule for a profile changed. + schedule_changed = QtCore.pyqtSignal() + + #: A job outcome was recorded. + jobs_changed = QtCore.pyqtSignal() + + def __init__(self) -> None: + super().__init__() + + self.app: application.VortaApp = QApplication.instance() + + #: profiles being submitted, so a timer tick cannot submit one twice + self._submitting: set[int] = set() + + # Scheduling is built first: restoring the pauses writes a status into its timers. + self._timers = SchedulerTimers(self) + self._state = SchedulerState(self) + self._state.restore_pauses() + + # connect signals + self.app.backup_finished_event.connect(lambda res: self.set_timer_for_profile(res['params']['profile_id'])) + + @property + def timers(self) -> dict[int, dict[str, QTimer | dt | ScheduleStatusType | None]]: + return self._timers.timers + + @property + def lock(self): + return self._timers.lock + + @property + def qt_timer(self) -> QTimer: + return self._timers.qt_timer + + @property + def net_status(self): + return self._timers.net_status + + @net_status.setter + def net_status(self, monitor) -> None: + self._timers.net_status = monitor + + @property + def _net_up(self) -> bool: + return self._timers._net_up + + @_net_up.setter + def _net_up(self, up: bool) -> None: + self._timers._net_up = up + + @property + def pauses(self) -> dict[int, tuple[dt, QtCore.QTimer]]: + return self._state.pauses + + @property + def wake_timer(self) -> QTimer: + return self._state.wake_timer + + @QtCore.pyqtSlot(bool) + def loginSuspendNotify(self, suspend: bool) -> None: + self._state.loginSuspendNotify(suspend) + + @QtCore.pyqtSlot() + def checkForResume(self) -> None: + self._state.checkForResume() + + @QtCore.pyqtSlot(bool) + def networkStatusChanged(self, up: bool) -> None: + self._timers.networkStatusChanged(up) + + def tr(self, *args: Any, **kwargs: Any) -> str: + scope = self.__class__.__name__ + return translate(scope, *args, **kwargs) + + def pause(self, profile_id: int, until: dt | None = None) -> None: + self._state.pause(profile_id, until) + + def unpause(self, profile_id: int) -> None: + self._state.unpause(profile_id) + + def clear_pause(self, profile_id: int) -> None: + self._state.clear_pause(profile_id) + + def paused(self, profile_id: int) -> bool: + return self._state.paused(profile_id) + + def record_skip( + self, + profile: BackupProfileModel, + trigger: str, + reason: str, + status: str = JobModel.Status.SKIPPED.value, + scheduled_at: dt | None = None, + ) -> None: + self._state.record_skip(profile, trigger, reason, status=status, scheduled_at=scheduled_at) + + def set_timer_for_profile(self, profile_id: int) -> None: + """Set a timer for next scheduled backup run of this profile, and run a missed one.""" + catch_up = self._timers.arm_profile(profile_id) + if catch_up is not None: + self.create_backup(catch_up, trigger=JobModel.Trigger.CATCHUP.value) + + def reload_all_timers(self) -> None: + self._timers.reload_all_timers() + + def remove_job(self, profile_id: int) -> None: + self._timers.remove_job(profile_id) + + def mark_paused(self, profile: BackupProfileModel, until: dt) -> None: + self._timers.mark_paused(profile, until) + + def next_job(self) -> str: + return self._timers.next_job() + + def next_job_for_profile(self, profile_id: int) -> ScheduleStatus: + return self._timers.next_job_for_profile(profile_id) + + def pending_jobs(self) -> list[PendingJob]: + return self._timers.pending_jobs() + + def create_backup(self, profile_id: int, trigger: str = JobModel.Trigger.SCHEDULED.value) -> None: + notifier = VortaNotifications.pick() + profile = BackupProfileModel.get_or_none(id=profile_id) + + if profile is None: + logger.info('Profile not found. Maybe deleted?') + return + + if profile_id in self._submitting: + logger.debug('A run for profile %s is already being submitted.', profile_id) + return + + # Skip if a job for this profile (repo) is already in progress + if self.app.jobs_manager.is_worker_running(site=profile.repo.id): + logger.debug('A job for repo %s is already active.', profile.repo.id) + self.record_skip(profile, trigger, 'Repository is busy with another job.') + self.pause(profile_id) + return + + self._submitting.add(profile_id) + try: + logger.info('Starting background backup for %s', profile.name) + notifier.deliver( + self.tr('Vorta Backup'), + self.tr('Starting background backup for %s.') % profile.name, + level='info', + ) + msg = BorgCreateJob.prepare(profile) + if msg['ok']: + logger.info('Preparation for backup successful.') + msg['category'] = 'scheduled' + job = BorgCreateJob(msg['cmd'], msg, profile.repo.id) + job.result.connect(self.notify) + self.app.jobs_manager.add_job(job) + else: + # Default to 'error': unexpected failures notify. + # Expected skips (WiFi/metered) use 'info' to suppress. + level = msg.get('level', 'error') + if level == 'error': + logger.error('Conditions for backup not met. Aborting.') + logger.error(msg['message']) + notifier.deliver( + self.tr('Vorta Backup'), + translate('messages', msg['message']), + level='error', + ) + status = JobModel.Status.FAILED.value + else: + logger.info('Backup skipped: %s', msg['message']) + status = JobModel.Status.SKIPPED.value + self.record_skip(profile, trigger, msg['message'], status=status) + self.pause(profile_id) + finally: + self._submitting.discard(profile_id) + + def notify(self, result: dict[str, Any]) -> None: + notifier = VortaNotifications.pick() + profile_name = result['params']['profile_name'] + profile_id = result['params']['profile'].id + + if result['returncode'] in [0, 1]: + notifier.deliver( + self.tr('Vorta Backup'), + self.tr('Backup successful for %s.') % profile_name, + level='info', + ) + logger.info('Backup creation successful.') + # unpause scheduler + self.unpause(result['params']['profile_id']) + + self.post_backup_tasks(profile_id) + else: + notifier.deliver( + self.tr('Vorta Backup'), + self.tr('Error during backup creation for %s.') % profile_name, + level='error', + ) + logger.error('Error during backup creation.') + # pause scheduler + # if a scheduled backup fails the scheduler should pause + # temporarily. + self.pause(result['params']['profile_id']) + + self.set_timer_for_profile(profile_id) + + def post_backup_tasks(self, profile_id: int) -> None: + """ + Pruning and checking after successful backup. + """ + profile = BackupProfileModel.get(id=profile_id) + notifier = VortaNotifications.pick() + logger.info('Doing post-backup jobs for %s', profile.name) + if profile.prune_on: + msg = BorgPruneJob.prepare(profile) + if msg['ok']: + job = BorgPruneJob(msg['cmd'], msg, profile.repo.id) + self.app.jobs_manager.add_job(job) + + # Refresh archives + msg = BorgListRepoJob.prepare(profile) + if msg['ok']: + job = BorgListRepoJob(msg['cmd'], msg, profile.repo.id) + self.app.jobs_manager.add_job(job) + + validation_cutoff = dt.now() - timedelta(days=7 * profile.validation_weeks) + recent_validations = ( + EventLogModel.select() + .where( + (EventLogModel.subcommand == 'check') + & (EventLogModel.start_time > validation_cutoff) + & (EventLogModel.repo_url == profile.repo.url) + ) + .count() + ) + if profile.validation_on and recent_validations == 0: + msg = BorgCheckJob.prepare(profile) + if msg['ok']: + job = BorgCheckJob(msg['cmd'], msg, profile.repo.id) + self.app.jobs_manager.add_job(job) + + compaction_cutoff = dt.now() - timedelta(days=7 * profile.compaction_weeks) + recent_compactions = ( + EventLogModel.select() + .where( + (EventLogModel.subcommand == '--info') + & (EventLogModel.start_time > compaction_cutoff) + & (EventLogModel.repo_url == profile.repo.url) + ) + .count() + ) + + if ( + profile.compaction_on + and recent_compactions == 0 + and version.parse(borg_compat.version) >= version.parse("1.2") + ): + msg = BorgCompactJob.prepare(profile) + if msg['ok']: + job = BorgCompactJob(msg['cmd'], msg, profile.repo.id) + self.app.jobs_manager.add_job(job) + + logger.info('Finished background task for profile %s', profile.name) + notifier.deliver( + self.tr('Vorta Backup'), + self.tr('Post Backup Tasks successful for %s' % profile.name), + level='info', + ) diff --git a/src/vorta/scheduler/scheduling.py b/src/vorta/scheduler/scheduling.py new file mode 100644 index 000000000..e6ee77562 --- /dev/null +++ b/src/vorta/scheduler/scheduling.py @@ -0,0 +1,352 @@ +from __future__ import annotations + +import enum +import logging +import threading +from collections.abc import Callable +from datetime import datetime as dt +from datetime import timedelta +from typing import TYPE_CHECKING, NamedTuple + +from PyQt6.QtCore import QTimer + +from vorta.store.models import BackupProfileModel, EventLogModel, JobModel +from vorta.utils import get_network_status_monitor + +if TYPE_CHECKING: + from vorta.scheduler import VortaScheduler + +logger = logging.getLogger(__name__) + +RESCHEDULE_INTERVAL_MS = 15 * 60 * 1000 + +# A QTimer interval is a C++ int, so a single wait tops out at about 24.8 days. +MAX_TIMER_MS = 2**31 - 1 +# Fire just after the deadline, so the handler always sees it as passed. +TIMER_GRACE_MS = 100 + + +def arm_deadline_timer(deadline: dt, on_expiry: Callable[[], None]) -> QTimer: + """ + Start a timer that calls `on_expiry` once `deadline` has passed. + + Waits longer than `MAX_TIMER_MS` are split into chunks, and every chunk is measured + against the wall clock again, so the deadline holds however far ahead it is. + """ + timer = QTimer() + timer.setSingleShot(True) + + def remaining_ms() -> float: + return (deadline - dt.now()).total_seconds() * 1000 + TIMER_GRACE_MS + + def rearm() -> None: + timer.setInterval(int(max(0, min(remaining_ms(), MAX_TIMER_MS)))) + timer.start() + + def on_timeout() -> None: + if remaining_ms() <= 0: + on_expiry() + else: + rearm() + + timer.timeout.connect(on_timeout) + rearm() + return timer + + +class ScheduleStatusType(enum.Enum): + SCHEDULED = enum.auto() # date provided + UNSCHEDULED = enum.auto() # Unknown + NO_PREVIOUS_BACKUP = enum.auto() # run a manual backup first + PAUSED = enum.auto() # paused after a failed or skipped run, date provided + + +class ScheduleStatus(NamedTuple): + type: ScheduleStatusType + time: dt | None = None + + +#: Timer states that hold a time for a run, and the row status each one shows as. +PENDING_STATUSES = { + ScheduleStatusType.SCHEDULED: JobModel.Status.SCHEDULED.value, + ScheduleStatusType.PAUSED: JobModel.Status.PAUSED.value, +} + + +class PendingJob(NamedTuple): + """A run the scheduler is holding a time for, but that hasn't been recorded yet.""" + + profile_id: int + profile_name: str + repo_url: str | None + scheduled_at: dt + status: str + + +class SchedulerTimers: + """The next run time of each profile, the timers holding them and the network gating.""" + + def __init__(self, scheduler: VortaScheduler) -> None: + self.scheduler = scheduler + + #: mapping of profiles to timers + self.timers: dict[int, dict[str, QTimer | dt | ScheduleStatusType | None]] = dict() + + self.lock = threading.Lock() + + # Periodic reschedule, in case a run was missed + self.qt_timer = QTimer() + self.qt_timer.timeout.connect(scheduler.reload_all_timers) + self.qt_timer.setInterval(RESCHEDULE_INTERVAL_MS) + self.qt_timer.start() + + # Connect to network manager to monitor net status + self.net_status = get_network_status_monitor() + self.net_status.network_status_changed.connect(scheduler.networkStatusChanged) + self._net_up = self.net_status.is_network_active() + + def networkStatusChanged(self, up: bool): + reload = self._net_up != up + self._net_up = up + logger.debug(f"network status up={up}") + if reload: + logger.info("updating schedule due to network status change") + self.reload_all_timers() + + def arm_profile(self, profile_id: int) -> int | None: + """ + Set a timer for the next scheduled backup run of this profile. + + Removes existing jobs if set to manual only or no repo is assigned. + + Else will look for previous scheduled backups and catch up if + schedule_make_up_missed is enabled. + + Or, if catch-up is not enabled, will add interval to last run to find + next suitable backup time. + + Returns the profile id whose missed run has to be caught up, if any. + """ + profile = BackupProfileModel.get_or_none(id=profile_id) + if profile is None: # profile doesn't exist any more. + return + logger.debug('Profile: %s, %d %d', str(profile), profile.schedule_fixed_hour, profile.schedule_fixed_minute) + + with self.lock: # Acquire lock + self.remove_job(profile_id) # reset schedule + + pause = self.scheduler.pauses.get(profile_id) + if pause is not None: + pause_end, timer = pause + if dt.now() < pause_end: + logger.debug( + 'Nothing scheduled for profile %s ' + 'because of timeout until %s.', + profile_id, + pause[0].strftime('%Y-%m-%d %H:%M:%S'), + ) + self.mark_paused(profile, pause_end) + return + else: + self.scheduler.clear_pause(profile_id) + + if profile.repo is None: # No backups without repo set + logger.debug( + 'Nothing scheduled for profile %s because of unset repo.', + profile_id, + ) + # Emit signal so that e.g. the GUI can react to the new schedule + self.scheduler.schedule_changed.emit() + return + + if profile.schedule_mode == 'off': + logger.debug('Scheduler for profile %s is disabled.', profile_id) + # Emit signal so that e.g. the GUI can react to the new schedule + self.scheduler.schedule_changed.emit() + return + + logger.info('Setting timer for profile %s', profile_id) + + # determine last backup time + last_run_log = ( + EventLogModel.select() + .where( + EventLogModel.subcommand == 'create', + EventLogModel.category == 'scheduled', + EventLogModel.profile == profile.id, + 0 <= EventLogModel.returncode <= 1, + ) + .order_by(EventLogModel.end_time.desc()) + .first() + ) + + if last_run_log is None: + # look for non scheduled (manual) backup runs + last_run_log = ( + EventLogModel.select() + .where( + EventLogModel.subcommand == 'create', + EventLogModel.profile == profile.id, + 0 <= EventLogModel.returncode <= 1, + ) + .order_by(EventLogModel.end_time.desc()) + .first() + ) + + if last_run_log is None: + logger.info( + f"Nothing scheduled for profile {profile_id} " + + "because it would be the first backup " + + "for this profile." + ) + self.timers[profile_id] = {'type': ScheduleStatusType.NO_PREVIOUS_BACKUP} + # Emit signal so that e.g. the GUI can react to the new schedule + self.scheduler.schedule_changed.emit() + return + + # calculate next scheduled time + if profile.schedule_mode == 'interval': + last_time: dt = last_run_log.end_time + + interval = {profile.schedule_interval_unit: profile.schedule_interval_count} + next_time = last_time + timedelta(**interval) + + elif profile.schedule_mode == 'fixed': + last_time = last_run_log.end_time + + next_time = last_time.replace( + hour=profile.schedule_fixed_hour, + minute=profile.schedule_fixed_minute, + second=0, + microsecond=0, + ) + timedelta(days=1) + + else: + # unknown schedule mode + raise ValueError("Unknown schedule mode '{}'".format(profile.schedule_mode)) + + logger.debug('Last run time: %s', last_time) + + needs_network = profile.repo is not None and profile.repo.is_remote_repo() + # handle missing of a scheduled time + if next_time <= dt.now(): + if profile.schedule_make_up_missed and (self._net_up or not needs_network): + logger.debug( + 'Catching up by running job for %s (%s)', + profile.name, + profile_id, + ) + return profile_id # create_backup will lead to a call to this method + elif profile.schedule_make_up_missed and not self._net_up and needs_network: + logger.debug('Skipping catchup %s (%s), the network is not available', profile.name, profile.id) + self.scheduler.record_skip( + profile, + JobModel.Trigger.CATCHUP.value, + 'Network unavailable for catch-up.', + scheduled_at=next_time, + ) + + # calculate next time from now + if profile.schedule_mode == 'interval': + # next_time % interval should be 0 + # while next_time > now + delta = dt.now() - last_time + next_time = dt.now() - delta % timedelta(**interval) + next_time += timedelta(**interval) + + elif profile.schedule_mode == 'fixed': + # schedule for today + next_time = dt.now().replace( + hour=profile.schedule_fixed_hour, + minute=profile.schedule_fixed_minute, + second=0, + microsecond=0, + ) + + if next_time <= dt.now(): + # time for today has passed, schedule for tomorrow + next_time += timedelta(days=1) + + # start QTimer + logger.debug('Scheduling next run for %s', next_time) + + self.timers[profile_id] = { + 'qtt': arm_deadline_timer(next_time, lambda: self.scheduler.create_backup(profile_id)), + 'dt': next_time, + 'type': ScheduleStatusType.SCHEDULED, + } + + # Emit signal so that e.g. the GUI can react to the new schedule + self.scheduler.schedule_changed.emit() + + def reload_all_timers(self) -> None: + logger.debug('Refreshing all scheduler timers') + for profile in BackupProfileModel.select(): + # Only set a timer for the profile if the network is actually up + if profile.repo is None: + logger.debug("nothing scheduled for %s because of unset repo", profile.id) + elif not profile.repo.is_remote_repo() or self._net_up: + logger.debug("scheduling %s", profile.id) + self.scheduler.set_timer_for_profile(profile.id) + else: + logger.debug("Network is down, not scheduling %s", profile.id) + self.remove_job(profile.id) + + def next_job(self) -> str: + now = dt.now() + + def is_scheduled(timer): + return timer["type"] == ScheduleStatusType.SCHEDULED and timer["qtt"].isActive() and timer["dt"] >= now + + scheduled = {profile_id: timer for profile_id, timer in self.timers.items() if is_scheduled(timer)} + if len(scheduled) == 0: + return self.scheduler.tr("None scheduled") + + closest_job = min(scheduled.items(), key=lambda item: item[1]["dt"]) + profile_id, timer = closest_job + time = timer["dt"] + profile = BackupProfileModel.get_or_none(id=profile_id) + + time_format = "%H:%M" + if time - now > timedelta(days=1): + time_format = "%b %d, %H:%M" + return f"{time.strftime(time_format)} ({profile.name})" + + def next_job_for_profile(self, profile_id: int) -> ScheduleStatus: + job = self.timers.get(profile_id) + if job is None: + return ScheduleStatus(ScheduleStatusType.UNSCHEDULED) + return ScheduleStatus(job['type'], time=job.get('dt')) # type: ignore[arg-type] + + def pending_jobs(self) -> list[PendingJob]: + """The runs currently on the clock, soonest first.""" + pending = [] + + for profile_id, timer in self.timers.items(): + status = PENDING_STATUSES.get(timer['type']) + if status is None: + continue + + profile = BackupProfileModel.get_or_none(id=profile_id) + if profile is None: + continue + + repo_url = profile.repo.url if profile.repo else None + pending.append(PendingJob(profile_id, profile.name, repo_url, timer['dt'], status)) + + return sorted(pending, key=lambda job: job.scheduled_at) + + def mark_paused(self, profile: BackupProfileModel, until: dt) -> None: + """Report the pause as the schedule status, unless the profile has no schedule to block.""" + if profile.repo is None or profile.schedule_mode == 'off': + return + + self.timers[profile.id] = {'type': ScheduleStatusType.PAUSED, 'dt': until} + self.scheduler.schedule_changed.emit() + + def remove_job(self, profile_id: int) -> None: + if profile_id in self.timers: + qtimer = self.timers[profile_id].get('qtt') + if qtimer is not None: + qtimer.stop() + + del self.timers[profile_id] diff --git a/src/vorta/scheduler/state.py b/src/vorta/scheduler/state.py new file mode 100644 index 000000000..2ec2e8b0a --- /dev/null +++ b/src/vorta/scheduler/state.py @@ -0,0 +1,246 @@ +from __future__ import annotations + +import logging +from datetime import datetime as dt +from datetime import timedelta +from typing import TYPE_CHECKING + +import peewee as pw +from PyQt6 import QtCore, QtDBus +from PyQt6.QtCore import QTimer + +from vorta.scheduler.scheduling import ScheduleStatusType, arm_deadline_timer +from vorta.store.models import BackupProfileModel, JobModel, SchedulerPauseModel + +if TYPE_CHECKING: + from vorta.scheduler import VortaScheduler + +logger = logging.getLogger(__name__) + +WAKE_CHECK_INTERVAL_MS = 5 * 60 * 1000 +WAKE_GAP_THRESHOLD = timedelta(minutes=10) + + +class SchedulerState: + """Pauses, recorded job outcomes and resume detection.""" + + def __init__(self, scheduler: VortaScheduler) -> None: + self.scheduler = scheduler + + # pausing will prevent scheduling for a specified time + self.pauses: dict[int, tuple[dt, QtCore.QTimer]] = dict() + + # connect to `systemd-logind` to receive sleep/resume events + # The signal `PrepareForSleep` will be emitted before and after hibernation. + service = "org.freedesktop.login1" + path = "/org/freedesktop/login1" + interface = "org.freedesktop.login1.Manager" + name = "PrepareForSleep" + bus = QtDBus.QDBusConnection.systemBus() + if bus.isConnected() and bus.interface().isServiceRegistered(service).value(): + self.bus = bus + self.bus.connect(service, path, interface, name, "b", scheduler.loginSuspendNotify) + else: + logger.info('No systemd-logind to notify us of sleep/resume, watching for clock gaps as well') + + self._last_wake_check = dt.now() + self.wake_timer = QTimer() + self.wake_timer.timeout.connect(scheduler.checkForResume) + self.wake_timer.setInterval(WAKE_CHECK_INTERVAL_MS) + self.wake_timer.start() + + def loginSuspendNotify(self, suspend: bool) -> None: + if not suspend: + logger.debug("Got login suspend/resume notification") + self._handle_resume() + + def checkForResume(self) -> None: + now = dt.now() + elapsed = now - self._last_wake_check + self._last_wake_check = now + + if elapsed < WAKE_GAP_THRESHOLD: + return + + logger.debug('Clock jumped %s since the last wake check, assuming the machine slept', elapsed) + self._handle_resume() + + def _handle_resume(self) -> None: + self._last_wake_check = dt.now() + # Defensively refetch in case the network status didn't arrive + self.scheduler._net_up = self.scheduler.net_status.is_network_active() + self.scheduler.reload_all_timers() + + def pause(self, profile_id: int, until: dt | None = None) -> None: + """ + Call a timeout for scheduling of a given profile. + + If `until` is omitted, a default time for the break is calculated. + + .. warning:: + This method won't work correctly when called from a non-`QThread`. + + + Parameters + ---------- + profile_id : int + The profile to pause the scheduling for. + until : dt | None, optional + The time to end the pause, by default None + """ + profile = BackupProfileModel.get_or_none(id=profile_id) + if profile is None: # profile doesn't exist any more. + return + + if profile.schedule_mode == 'off': + return + + if until is None: + # calculate default timeout + + if profile.schedule_mode == 'interval': + interval = timedelta(**{profile.schedule_interval_unit: profile.schedule_interval_count}) + else: + # fixed + interval = timedelta(days=1) + + timeout = interval // 6 # 60 / 6 = 10 [min] + timeout = max(min(timeout, timedelta(hours=1)), timedelta(minutes=1)) # 1 <= t <= 60 + + until = dt.now().replace(microsecond=0) + timeout + elif until < dt.now(): + return + + # remove existing schedule + self.scheduler.remove_job(profile_id) + + # set timeout/pause + other_pause = self.pauses.get(profile_id) + if other_pause is not None: + logger.debug(f"Override existing timeout for profile {profile_id}") + + self._set_pause(profile, until) + logger.debug(f"Paused {profile_id} until {until.strftime('%Y-%m-%d %H:%M:%S')}") + + def _set_pause(self, profile: BackupProfileModel, until: dt) -> None: + """Arm the reschedule timer for a pause and store it, so it outlives the process.""" + profile_id = profile.id + replaced = self.pauses.get(profile_id) + if replaced is not None: + replaced[1].stop() + + # setting timer for reschedule is not possible if called + # from a non-QThread - it won't fail but won't work + timer = arm_deadline_timer(until, lambda: self.scheduler.set_timer_for_profile(profile_id)) + + self.pauses[profile_id] = (until, timer) + + try: + SchedulerPauseModel.replace(profile=profile_id, paused_until=until).execute() + except pw.PeeweeException: + logger.warning('Could not store pause for profile %s.', profile_id, exc_info=True) + + self.scheduler.mark_paused(profile, until) + + def clear_pause(self, profile_id: int) -> None: + """Drop a pause from memory, from the schedule status and from the database.""" + pause = self.pauses.pop(profile_id, None) + if pause is not None: + pause[1].stop() + + status = self.scheduler.timers.get(profile_id) + if status is not None and status.get('type') is ScheduleStatusType.PAUSED: + del self.scheduler.timers[profile_id] + + try: + SchedulerPauseModel.delete().where(SchedulerPauseModel.profile == profile_id).execute() + except pw.PeeweeException: + logger.warning('Could not drop stored pause for profile %s.', profile_id, exc_info=True) + + def restore_pauses(self) -> None: + """Re-arm the pauses stored by a previous run, dropping the ones that already ran out.""" + now = dt.now() + + for stored in list(SchedulerPauseModel.select()): + profile_id = stored.profile_id + profile = BackupProfileModel.get_or_none(id=profile_id) + + if profile is None or stored.paused_until <= now: + self.clear_pause(profile_id) + continue + + self._set_pause(profile, stored.paused_until) + logger.debug(f"Restored pause for {profile_id} until {stored.paused_until:%Y-%m-%d %H:%M:%S}") + + def unpause(self, profile_id: int) -> None: + """ + Return to scheduling for a profile. + + Parameters + ---------- + profile_id : int + The profile to end the timeout for. + """ + profile = BackupProfileModel.get_or_none(id=profile_id) + if profile is None: # profile doesn't exist any more. + return + + pause = self.pauses.get(profile_id) + if pause is None: # already unpaused + return + + self.clear_pause(profile_id) + + logger.debug(f"Unpaused {profile_id}") + + self.scheduler.set_timer_for_profile(profile_id) + + def paused(self, profile_id: int) -> bool: + """ + Determine whether scheduling for a profile is paused + + Parameters + ---------- + profile_id : int + + Returns + ------- + bool + """ + return self.pauses.get(profile_id) is not None + + def record_skip( + self, + profile: BackupProfileModel, + trigger: str, + reason: str, + status: str = JobModel.Status.SKIPPED.value, + scheduled_at: dt | None = None, + ) -> None: + """Record a job outcome, deduplicated on the occurrence when one is known.""" + lookup = { + 'profile': str(profile.id), + 'trigger': trigger, + 'status': status, + 'scheduled_at': scheduled_at, + } + details = { + 'profile_name': profile.name, + 'repo_url': profile.repo.url if profile.repo else None, + 'job_type': JobModel.Type.BACKUP.value, + 'reason': reason, + } + + try: + if scheduled_at is None: + JobModel.create(**lookup, **details) + else: + _, recorded = JobModel.get_or_create(**lookup, defaults=details) + if not recorded: + return + except pw.PeeweeException: + logger.warning('Could not record job for profile %s.', profile.id, exc_info=True) + return + + # `arm_profile` records under the scheduler's lock, and the jobs view reads the table here. + QTimer.singleShot(0, self.scheduler.jobs_changed.emit) diff --git a/tests/unit/test_schedule.py b/tests/unit/test_schedule.py index aad602af4..e31660cdc 100644 --- a/tests/unit/test_schedule.py +++ b/tests/unit/test_schedule.py @@ -7,6 +7,8 @@ from PyQt6.QtWidgets import QWidget import vorta.scheduler +import vorta.scheduler.scheduling +import vorta.scheduler.state from vorta.application import VortaApp from vorta.store.models import BackupProfileModel, EventLogModel, JobModel from vorta.views.partials.jobs_table_model import JobsTableModel @@ -18,7 +20,8 @@ @pytest.fixture def clockmock(monkeypatch): datetime_mock = MagicMock(wraps=dt) - monkeypatch.setattr(vorta.scheduler, "dt", datetime_mock) + for module in (vorta.scheduler, vorta.scheduler.scheduling, vorta.scheduler.state): + monkeypatch.setattr(module, "dt", datetime_mock) return datetime_mock diff --git a/tests/unit/test_scheduler.py b/tests/unit/test_scheduler.py index d3e1c9a35..b9ad4c642 100644 --- a/tests/unit/test_scheduler.py +++ b/tests/unit/test_scheduler.py @@ -10,6 +10,8 @@ import vorta.borg import vorta.scheduler +import vorta.scheduler.scheduling +import vorta.scheduler.state from vorta.scheduler import PendingJob, ScheduleStatus, ScheduleStatusType, VortaScheduler from vorta.store.models import BackupProfileModel, EventLogModel, JobModel, SchedulerPauseModel @@ -22,7 +24,8 @@ @pytest.fixture def clockmock(monkeypatch): datetime_mock = MagicMock(wraps=dt) - monkeypatch.setattr(vorta.scheduler, "dt", datetime_mock) + for module in (vorta.scheduler, vorta.scheduler.scheduling, vorta.scheduler.state): + monkeypatch.setattr(module, "dt", datetime_mock) return datetime_mock @@ -288,7 +291,7 @@ def test_deleting_a_paused_profile_clears_the_pause(qapp, qtbot, mocker): prepare_mock = mocker.patch('vorta.scheduler.BorgCreateJob.prepare') mocker.patch.object(QMessageBox, 'question', return_value=QMessageBox.StandardButton.Yes) - mocker.patch.object(qapp.scheduler, '_net_up', True) + mocker.patch.object(qapp.scheduler._timers, '_net_up', True) qtbot.mouseClick(main.profileDeleteButton, QtCore.Qt.MouseButton.LeftButton) # The overdue occurrence must not start a backup on the profile being deleted. @@ -742,6 +745,80 @@ def test_set_timer_records_skip_when_network_down_for_catchup(qtbot, clockmock): assert JobModel.select().count() == jobs_before + 1 +def test_catchup_is_submitted_after_the_lock_is_released(mocker, clockmock): + """The catch-up run used to be submitted inside `lock`, which reentered it in `create_backup`.""" + scheduler = VortaScheduler() + scheduler._net_up = True + + time = dt(2020, 5, 6, 4, 30) + clockmock.now.return_value = time + + profile = BackupProfileModel.get(name=PROFILE_NAME) + profile.schedule_make_up_missed = True + profile.schedule_mode = INTERVAL_SCHEDULE + profile.schedule_interval_unit = 'hours' + profile.schedule_interval_count = 3 + profile.save() + + last_run = time - td(hours=6) + EventLogModel.create( + subcommand='create', + profile=profile.id, + returncode=0, + category='scheduled', + start_time=last_run, + end_time=last_run, + ) + + submitted = [] + mocker.patch.object( + scheduler, + 'create_backup', + side_effect=lambda profile_id, trigger: submitted.append((profile_id, trigger, scheduler.lock.locked())), + ) + + scheduler.set_timer_for_profile(profile.id) + + assert submitted == [(profile.id, JobModel.Trigger.CATCHUP.value, False)] + + +def test_create_backup_does_not_hold_the_lock_while_preparing(qapp, mocker): + """`prepare()` can spin the event loop through a keyring dialog, so the lock has to be free.""" + scheduler = qapp.scheduler + held = [] + + def prepare(profile): + held.append(scheduler.lock.locked()) + return {'ok': False, 'message': 'Current Wifi is not allowed.', 'level': 'info'} + + mocker.patch('vorta.scheduler.BorgCreateJob.prepare', side_effect=prepare) + + scheduler.create_backup(1) + + assert held == [False] + scheduler.unpause(1) + + +def test_create_backup_ignores_a_reentrant_run_for_the_same_profile(qapp, mocker): + """Nothing marks the profile busy until `add_job`, so a tick mid-`prepare()` must not resubmit.""" + scheduler = qapp.scheduler + prepared = [] + + def prepare(profile): + prepared.append(profile.id) + if len(prepared) == 1: + # what a timer tick does while the keyring dialog is open + scheduler.create_backup(profile.id) + return {'ok': False, 'message': 'Current Wifi is not allowed.', 'level': 'info'} + + mocker.patch('vorta.scheduler.BorgCreateJob.prepare', side_effect=prepare) + + scheduler.create_backup(1) + + assert prepared == [1] + scheduler.unpause(1) + + def test_wall_clock_gap_is_treated_as_a_resume(mocker, clockmock): """Without logind, a jump in wall clock time is the only sign that the machine slept.""" clockmock.now.return_value = dt(2020, 5, 6, 4, 0)