Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 2 additions & 0 deletions src/vorta/borg/borg_job.py
Original file line number Diff line number Diff line change
Expand Up @@ -343,6 +343,8 @@ def read_async(fd):
log_entry.returncode = p.returncode
log_entry.repo_url = self.params.get('repo_url', None)
log_entry.end_time = dt.now()
result['log_entry_id'] = log_entry.id

with db_lock:
log_entry.save()
self.process_result(result)
Expand Down
172 changes: 16 additions & 156 deletions src/vorta/scheduler/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -2,22 +2,15 @@

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.execution import SchedulerExecution
from vorta.scheduler.scheduling import (
MAX_TIMER_MS,
PENDING_STATUSES,
Expand All @@ -30,8 +23,7 @@
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
from vorta.store.models import BackupProfileModel, JobModel

logger = logging.getLogger(__name__)

Expand Down Expand Up @@ -62,10 +54,9 @@ def __init__(self) -> None:

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.
# Execution first: it starts nothing, and the other two arm timers that route back into it.
self._execution = SchedulerExecution(self)
# Scheduling before State: restoring the pauses writes a status into its timers.
self._timers = SchedulerTimers(self)
self._state = SchedulerState(self)
self._state.restore_pauses()
Expand Down Expand Up @@ -147,6 +138,14 @@ def record_skip(
) -> None:
self._state.record_skip(profile, trigger, reason, status=status, scheduled_at=scheduled_at)

def record_start(self, profile: BackupProfileModel, trigger: str) -> int | None:
return self._state.record_start(profile, trigger)

def record_finish(
self, record_id: int, status: str, log_entry_id: int | None = None, reason: str | None = None
) -> None:
self._state.record_finish(record_id, status, log_entry_id, reason)

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)
Expand All @@ -172,149 +171,10 @@ 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)
self._execution.create_backup(profile_id, trigger)

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)
self._execution.notify(result)

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',
)
self._execution.post_backup_tasks(profile_id)
Loading
Loading