702d080a1e
Google BK
398 lines
15 KiB
Python
398 lines
15 KiB
Python
from datetime import datetime, timedelta, date
|
|
from io import IOBase
|
|
from typing import Dict, Generic, List, Optional, Tuple, TypeVar
|
|
|
|
from injector import inject, singleton
|
|
|
|
from .backupscheme import GenerationalScheme, OldestScheme, DeleteAfterUploadScheme
|
|
from backup.config import Config, Setting, CreateOptions
|
|
from backup.exceptions import DeleteMutlipleBackupsError, SimulatedError
|
|
from backup.util import GlobalInfo, Estimator, DataCache
|
|
from .backups import AbstractBackup, Backup
|
|
from .dummybackup import DummyBackup
|
|
from .precache import Precache
|
|
from backup.time import Time
|
|
from backup.worker import Trigger
|
|
from backup.logger import getLogger
|
|
|
|
logger = getLogger(__name__)
|
|
|
|
T = TypeVar('T')
|
|
|
|
|
|
class BackupSource(Trigger, Generic[T]):
|
|
def __init__(self):
|
|
super().__init__()
|
|
pass
|
|
|
|
def name(self) -> str:
|
|
return "Unnamed"
|
|
|
|
def title(self) -> str:
|
|
return "Default"
|
|
|
|
def enabled(self) -> bool:
|
|
return True
|
|
|
|
def needsConfiguration(self) -> bool:
|
|
return not self.enabled()
|
|
|
|
def upload(self) -> bool:
|
|
return True
|
|
|
|
def icon(self) -> str:
|
|
return "sd_card"
|
|
|
|
def freeSpace(self):
|
|
return None
|
|
|
|
async def create(self, options: CreateOptions) -> T:
|
|
pass
|
|
|
|
async def get(self) -> Dict[str, T]:
|
|
pass
|
|
|
|
async def delete(self, backup: T):
|
|
pass
|
|
|
|
async def ignore(self, backup: T, ignore: bool):
|
|
pass
|
|
|
|
async def save(self, backup: AbstractBackup, bytes: IOBase) -> T:
|
|
pass
|
|
|
|
async def read(self, backup: T) -> IOBase:
|
|
pass
|
|
|
|
async def retain(self, backup: T, retain: bool) -> None:
|
|
pass
|
|
|
|
async def note(self, backup, note: str) -> None:
|
|
pass
|
|
|
|
def maxCount(self) -> None:
|
|
return 0
|
|
|
|
def postSync(self) -> None:
|
|
return
|
|
|
|
def detail(self) -> str:
|
|
return ""
|
|
|
|
def isDestination(self) -> bool:
|
|
return False
|
|
|
|
# Gets called after reading state but before any changes are made
|
|
# to check for additional errors.
|
|
def checkBeforeChanges(self) -> None:
|
|
pass
|
|
|
|
|
|
class BackupDestination(BackupSource):
|
|
def isWorking(self):
|
|
return False
|
|
|
|
@property
|
|
def might_be_oob_creds(self) -> bool:
|
|
return False
|
|
|
|
def isDestination(self) -> bool:
|
|
return True
|
|
|
|
|
|
@singleton
|
|
class Model():
|
|
@inject
|
|
def __init__(self, config: Config, time: Time, source: BackupSource, dest: BackupDestination, info: GlobalInfo, estimator: Estimator, data_cache: DataCache):
|
|
self.config: Config = config
|
|
self.time = time
|
|
self.precache: Precache = None
|
|
self.source: BackupSource = source
|
|
self.dest: BackupDestination = dest
|
|
self.reinitialize()
|
|
self.backups: Dict[str, Backup] = {}
|
|
self.firstSync = True
|
|
self.info = info
|
|
self.simulate_error = None
|
|
self.estimator = estimator
|
|
self.waiting_for_startup = False
|
|
self.ignore_startup_delay = False
|
|
self._data_cache = data_cache
|
|
|
|
def enabled(self):
|
|
if self.source.needsConfiguration():
|
|
return False
|
|
if self.dest.needsConfiguration():
|
|
return False
|
|
return True
|
|
|
|
def allSources(self):
|
|
return [self.source, self.dest]
|
|
|
|
def reinitialize(self, precache: Precache = None):
|
|
self.precache = precache
|
|
self._time_of_day: Optional[Tuple[int, int]] = self._parseTimeOfDay()
|
|
|
|
# SOMEDAY: this should be cached in config and regenerated on config updates, not here
|
|
self.generational_config = self.config.getGenerationalConfig()
|
|
|
|
def getTimeOfDay(self):
|
|
return self._time_of_day
|
|
|
|
def _nextBackup(self, now: datetime, last_backup: Optional[datetime]) -> Optional[datetime]:
|
|
timeofDay = self.getTimeOfDay()
|
|
if self.config.get(Setting.DAYS_BETWEEN_BACKUPS) <= 0:
|
|
next = None
|
|
elif self.dest.needsConfiguration():
|
|
next = None
|
|
elif not last_backup:
|
|
# this isn't the cleanest logic, but the idea here is that if there are no backups,
|
|
# then the backups shoudl be made right when the addon starts up.
|
|
next = self.info.start_time
|
|
elif not timeofDay:
|
|
next = last_backup + timedelta(days=self.config.get(Setting.DAYS_BETWEEN_BACKUPS))
|
|
else:
|
|
newest_local: datetime = self.time.toLocal(last_backup)
|
|
time_that_day_local = self.time.localize(datetime(newest_local.year, newest_local.month, newest_local.day, timeofDay[0], timeofDay[1]))
|
|
if newest_local < time_that_day_local:
|
|
# Latest backup is before the backup time for that day
|
|
next = self.time.toUtc(time_that_day_local)
|
|
else:
|
|
# return the next backup after the delta
|
|
next_date = date.fromordinal(int(date(newest_local.year, newest_local.month, newest_local.day).toordinal() + self.config.get(Setting.DAYS_BETWEEN_BACKUPS)))
|
|
next_datetime_local = self.time.localize(datetime(next_date.year, next_date.month, next_date.day, timeofDay[0], timeofDay[1]))
|
|
next = self.time.toUtc(next_datetime_local)
|
|
|
|
if next is None:
|
|
self.waiting_for_startup = False
|
|
return None
|
|
|
|
# Don't backup X minutes after startup, since that can put an unreasonable amount of strain on
|
|
# the system while booting up.
|
|
cooldown_minimum = self.info.backupCooldownTime()
|
|
if next <= now and now < cooldown_minimum and not self.ignore_startup_delay:
|
|
self.waiting_for_startup = True
|
|
return cooldown_minimum
|
|
elif self.ignore_startup_delay:
|
|
self.waiting_for_startup = False
|
|
return next
|
|
elif cooldown_minimum > next:
|
|
self.waiting_for_startup = cooldown_minimum > now
|
|
return cooldown_minimum
|
|
else:
|
|
self.waiting_for_startup = False
|
|
return next
|
|
|
|
def nextBackup(self, now: datetime, include_pending=True):
|
|
latest = max(filter(lambda s: not s.ignore() and (not s.isPending() or include_pending), self.backups.values()),
|
|
default=None, key=lambda s: s.date())
|
|
if latest:
|
|
latest = latest.date()
|
|
return self._nextBackup(now, latest)
|
|
|
|
async def sync(self, now: datetime):
|
|
if self.simulate_error is not None:
|
|
if self.simulate_error.startswith("test"):
|
|
raise Exception(self.simulate_error)
|
|
else:
|
|
raise SimulatedError(self.simulate_error)
|
|
await self._syncBackups([self.source, self.dest], now)
|
|
|
|
self.source.checkBeforeChanges()
|
|
self.dest.checkBeforeChanges()
|
|
|
|
if not self.dest.needsConfiguration():
|
|
if self.source.enabled():
|
|
await self._purge(self.source)
|
|
if self.dest.enabled():
|
|
await self._purge(self.dest)
|
|
|
|
# Delete any "ignored" backups that have expired
|
|
if (self.config.get(Setting.IGNORE_OTHER_BACKUPS) or self.config.get(Setting.IGNORE_UPGRADE_BACKUPS)) and self.config.get(Setting.DELETE_IGNORED_AFTER_DAYS) > 0:
|
|
cutoff = now - timedelta(days=self.config.get(Setting.DELETE_IGNORED_AFTER_DAYS))
|
|
delete = []
|
|
for backup in self.backups.values():
|
|
if backup.ignore() and backup.date() < cutoff:
|
|
delete.append(backup)
|
|
for backup in delete:
|
|
await self.deleteBackup(backup, self.source)
|
|
|
|
self._handleBackupDetails()
|
|
next_backup = self.nextBackup(now)
|
|
if next_backup and now >= next_backup and self.source.enabled() and not self.dest.needsConfiguration():
|
|
if self.config.get(Setting.DELETE_BEFORE_NEW_BACKUP):
|
|
await self._purge(self.source, pre_purge=True)
|
|
await self.createBackup(CreateOptions(now, self.config.get(Setting.BACKUP_NAME)))
|
|
await self._purge(self.source)
|
|
self._handleBackupDetails()
|
|
|
|
if self.dest.enabled() and self.dest.upload():
|
|
# get the backups we should upload
|
|
uploads = []
|
|
for backup in self.backups.values():
|
|
if backup.getSource(self.source.name()) is not None and backup.getSource(self.source.name()).uploadable() and backup.getSource(self.dest.name()) is None and not backup.ignore():
|
|
uploads.append(backup)
|
|
uploads.sort(key=lambda s: s.date())
|
|
uploads.reverse()
|
|
for upload in uploads:
|
|
# only upload if doing so won't result in it being deleted next
|
|
dummy = DummyBackup(
|
|
"", upload.date(), self.dest.name(), "dummy_slug_name")
|
|
proposed = list(self.backups.values())
|
|
proposed.append(dummy)
|
|
if self._nextPurge(self.dest, proposed)[1] != dummy:
|
|
if self.config.get(Setting.DELETE_BEFORE_NEW_BACKUP):
|
|
await self._purge(self.dest, pre_purge=True)
|
|
upload.addSource(await self.dest.save(upload, await self.source.read(upload)))
|
|
await self._purge(self.dest)
|
|
self._handleBackupDetails()
|
|
else:
|
|
break
|
|
if self.config.get(Setting.DELETE_AFTER_UPLOAD):
|
|
await self._purge(self.source)
|
|
self._handleBackupDetails()
|
|
self.source.postSync()
|
|
self.dest.postSync()
|
|
self._data_cache.saveIfDirty()
|
|
|
|
def isWorkingThroughUpload(self):
|
|
return self.dest.isWorking()
|
|
|
|
async def createBackup(self, options):
|
|
if not self.source.enabled():
|
|
return
|
|
|
|
self.estimator.refresh()
|
|
self.estimator.checkSpace(list(self.backups.values()))
|
|
created = await self.source.create(options)
|
|
backup = Backup(created)
|
|
self.backups[backup.slug()] = backup
|
|
|
|
async def deleteBackup(self, backup, source):
|
|
if not backup.getSource(source.name()):
|
|
return
|
|
slug = backup.slug()
|
|
await source.delete(backup)
|
|
backup.removeSource(source.name())
|
|
if backup.isDeleted():
|
|
del self.backups[slug]
|
|
|
|
def getNextPurges(self):
|
|
purges = {}
|
|
for source in [self.source, self.dest]:
|
|
purges[source.name()] = self._nextPurge(
|
|
source, self.backups.values(), findNext=True)[1]
|
|
return purges
|
|
|
|
def _parseTimeOfDay(self) -> Optional[Tuple[int, int]]:
|
|
from_config = self.config.get(Setting.BACKUP_TIME_OF_DAY)
|
|
if len(from_config) == 0:
|
|
return None
|
|
parts = from_config.split(":")
|
|
if len(parts) != 2:
|
|
return None
|
|
try:
|
|
hour: int = int(parts[0])
|
|
minute: int = int(parts[1])
|
|
if hour < 0 or minute < 0 or hour > 23 or minute > 59:
|
|
return None
|
|
return (hour, minute)
|
|
except ValueError:
|
|
# Parse error
|
|
return None
|
|
|
|
async def _syncBackups(self, sources: List[BackupSource], now: datetime):
|
|
for source in sources:
|
|
if source.enabled():
|
|
# check if we have the results from this source precached
|
|
from_source: Dict[str, AbstractBackup] = None
|
|
if self.precache is not None:
|
|
from_source = self.precache.cached(source.name(), now)
|
|
if not from_source:
|
|
from_source = await source.get()
|
|
else:
|
|
from_source: Dict[str, AbstractBackup] = {}
|
|
for backup in from_source.values():
|
|
if backup.slug() not in self.backups:
|
|
self.backups[backup.slug()] = Backup(backup)
|
|
else:
|
|
self.backups[backup.slug()].addSource(backup)
|
|
for backup in list(self.backups.values()):
|
|
if backup.slug() not in from_source:
|
|
slug = backup.slug()
|
|
backup.removeSource(source.name())
|
|
if backup.isDeleted():
|
|
del self.backups[slug]
|
|
self.firstSync = False
|
|
|
|
def _buildDeleteScheme(self, source, findNext=False):
|
|
count = source.maxCount()
|
|
if findNext:
|
|
count -= 1
|
|
if source == self.source and self.config.get(Setting.DELETE_AFTER_UPLOAD):
|
|
return DeleteAfterUploadScheme(source.name(), [self.dest.name()])
|
|
elif self.generational_config:
|
|
return GenerationalScheme(
|
|
self.time, self.generational_config, count=count)
|
|
else:
|
|
return OldestScheme(count=count)
|
|
|
|
def _buildNamingScheme(self):
|
|
source = max(filter(BackupSource.enabled, self.allSources()), key=BackupSource.maxCount)
|
|
return self._buildDeleteScheme(source)
|
|
|
|
def _handleBackupDetails(self):
|
|
self._buildNamingScheme().handleNaming(self.backups.values())
|
|
|
|
def _nextPurge(self, source: BackupSource, backups, findNext=False):
|
|
"""
|
|
Given a list of backups, decides if one should be purged.
|
|
"""
|
|
if not source.enabled() or len(backups) == 0:
|
|
return None, None
|
|
if source.maxCount() == 0 and source.isDestination():
|
|
# When maxCount is zero for a destination, we should never delete from it.
|
|
return None, None
|
|
if source.maxCount() == 0 and not self.config.get(Setting.DELETE_AFTER_UPLOAD):
|
|
return None, None
|
|
|
|
scheme = self._buildDeleteScheme(source, findNext=findNext)
|
|
consider_purging = []
|
|
for backup in backups:
|
|
source_backup = backup.getSource(source.name())
|
|
if source_backup is not None and source_backup.considerForPurge() and not backup.ignore():
|
|
consider_purging.append(backup)
|
|
if len(consider_purging) == 0:
|
|
return None, None
|
|
return scheme.getOldest(consider_purging)
|
|
|
|
async def _purge(self, source: BackupSource, pre_purge=False):
|
|
while True:
|
|
purge = self._getPurgeList(source, pre_purge)
|
|
|
|
reasons = set(map(lambda p: p[1], purge))
|
|
if len(purge) <= 0:
|
|
return
|
|
if len(purge) != len(reasons) and (self.config.get(Setting.CONFIRM_MULTIPLE_DELETES) and not self.info.isPermitMultipleDeletes()):
|
|
raise DeleteMutlipleBackupsError(self._getPurgeStats())
|
|
await self.deleteBackup(purge[0][0], source)
|
|
|
|
def _getPurgeStats(self):
|
|
ret = {}
|
|
for source in [self.source, self.dest]:
|
|
ret[source.name()] = len(self._getPurgeList(source))
|
|
return ret
|
|
|
|
def _getPurgeList(self, source: BackupSource, pre_purge=False):
|
|
if not source.enabled():
|
|
return []
|
|
candidates = list(self.backups.values())
|
|
purges = []
|
|
while True:
|
|
reason, next_purge = self._nextPurge(source, candidates, findNext=pre_purge)
|
|
if next_purge is None:
|
|
return purges
|
|
else:
|
|
purges.append((next_purge, reason))
|
|
candidates.remove(next_purge)
|