180 lines
8.1 KiB
Python
180 lines
8.1 KiB
Python
"""Explicit operator restore with a durable, non-destructive recovery boundary."""
|
|
|
|
|
|
def tree_metadata(folder):
|
|
return {'schemaVersion': 1, 'application': 'werkjournal', 'createdAt': int(time.time()),
|
|
'release': manifest(folder / 'release'),
|
|
'files': {p.relative_to(folder).as_posix(): checksum(p) for p in sorted(folder.rglob('*'))
|
|
if p.is_file() and p.name != 'backup.json'}}
|
|
|
|
|
|
def verify_tree(folder):
|
|
expected = json.loads((folder / 'backup.json').read_text())
|
|
if any(p.is_symlink() for p in folder.rglob('*')):
|
|
raise ValueError('Restore tree must not contain symlinks')
|
|
actual = tree_metadata(folder)
|
|
if expected['release'] != actual['release'] or expected['files'] != actual['files']:
|
|
raise ValueError('Restore tree checksum mismatch')
|
|
return expected
|
|
|
|
|
|
class OperatorRestore:
|
|
def __init__(self, deployment):
|
|
self.d = deployment
|
|
self.pending = deployment.app / 'pending-restore.json'
|
|
|
|
def record(self, state, phase):
|
|
state['phase'] = phase
|
|
atomic(self.pending, json.dumps(state))
|
|
|
|
def bundle(self, state):
|
|
if (not re.fullmatch('[0-9a-f]{32}', state.get('id', ''))
|
|
or any(not re.fullmatch('[0-9a-f]{40}', state.get(key, {}).get('commit', '')) for key in ('old', 'new'))
|
|
or any(type(state.get(key)) is not bool for key in ('wasRunning', 'hadMaintenance'))):
|
|
raise ValueError('Invalid restore identifier')
|
|
return self.d.app / 'restores' / state['id']
|
|
|
|
def install_data(self, source, bundle):
|
|
replace_database(source, self.d.app / 'data', bundle)
|
|
|
|
def activate(self, release):
|
|
self.d.start()
|
|
self.d.ready(release)
|
|
|
|
def finish(self, state):
|
|
# Terminal intent is written before opening the site. Never reapply a DB
|
|
# snapshot when resuming after this boundary: new writes may exist already.
|
|
if not state['hadMaintenance'] and self.d.flag.exists():
|
|
self.d.flag.unlink()
|
|
sync_directory(self.d.site)
|
|
self.pending.unlink()
|
|
sync_directory(self.d.app)
|
|
|
|
def recover(self):
|
|
self.d.preflight(check_data=False, allow_restore=True)
|
|
if self.d.pending.exists():
|
|
raise RuntimeError('A deployment journal also exists; operator inspection required')
|
|
state = json.loads(self.pending.read_text())
|
|
bundle = self.bundle(state)
|
|
if state['phase'] in {'published', 'reverted'}:
|
|
expected = state['new'] if state['phase'] == 'published' else state['old']
|
|
if (manifest(self.d.app / 'current') != expected
|
|
or json.loads((self.d.app / 'deployed.json').read_text()) != expected):
|
|
raise RuntimeError('Published restore state differs from current release')
|
|
if state['wasRunning']:
|
|
self.activate(expected)
|
|
else:
|
|
self.d.stop()
|
|
self.finish(state)
|
|
return
|
|
if state['phase'] not in {'prepared', 'backed_up', 'installed'}:
|
|
raise ValueError('Unknown restore phase')
|
|
if not self.d.flag.exists():
|
|
self.d.maintenance()
|
|
self.d.stop()
|
|
if state['phase'] != 'prepared':
|
|
prior = bundle / 'prior'
|
|
verify_tree(prior)
|
|
self.install_data(prior / 'data', bundle)
|
|
old = self.d.app / 'releases' / state['old']['commit']
|
|
if manifest(old) != state['old']:
|
|
raise ValueError('Previous release failed verification')
|
|
link(self.d.app / 'current', old.relative_to(self.d.app))
|
|
atomic(self.d.app / 'deployed.json', json.dumps(state['old']))
|
|
if state['wasRunning']:
|
|
self.activate(state['old'])
|
|
self.record(state, 'reverted')
|
|
self.finish(state)
|
|
|
|
def restore(self, archive):
|
|
self.d.preflight()
|
|
if self.d.pending.exists():
|
|
raise RuntimeError('Resolve pending deployment before restoring')
|
|
metadata = verify(archive)
|
|
old = manifest(self.d.app / 'current')
|
|
if old != json.loads((self.d.app / 'deployed.json').read_text()):
|
|
raise RuntimeError('Current release differs from the verified deployment')
|
|
active = self.d.command('systemctl', '--user', 'show', self.d.service, '--property=ActiveState', '--value')
|
|
if active not in {'active', 'inactive', 'failed'}:
|
|
raise RuntimeError('Service is transitioning')
|
|
state = {'id': uuid.uuid4().hex, 'phase': 'prepared', 'old': old,
|
|
'new': metadata['release'], 'wasRunning': active == 'active',
|
|
'hadMaintenance': self.d.flag.exists()}
|
|
bundle = self.bundle(state)
|
|
source = bundle / 'source'
|
|
source.mkdir(parents=True, mode=0o700)
|
|
# verify() has already rejected paths, links, duplicates and bad checksums.
|
|
with tarfile.open(archive, 'r:gz') as content:
|
|
for entry in content:
|
|
target = source / entry.name
|
|
target.parent.mkdir(parents=True, exist_ok=True, mode=0o700)
|
|
with content.extractfile(entry) as incoming, target.open('xb') as output:
|
|
shutil.copyfileobj(incoming, output)
|
|
verify_tree(source)
|
|
release = self.d.app / 'releases' / state['new']['commit']
|
|
if release.exists():
|
|
if manifest(release) != state['new']:
|
|
raise ValueError('Release ID already has different contents')
|
|
else:
|
|
staged_release = Path(tempfile.mkdtemp(prefix='restore-release-', dir=self.d.app))
|
|
try:
|
|
shutil.copytree(source / 'release', staged_release, dirs_exist_ok=True)
|
|
durable_tree(staged_release)
|
|
staged_release.rename(release)
|
|
sync_directory(release.parent)
|
|
finally:
|
|
if staged_release.exists():
|
|
shutil.rmtree(staged_release)
|
|
durable_tree(source)
|
|
sync_directory(bundle)
|
|
sync_directory(bundle.parent)
|
|
self.record(state, 'prepared')
|
|
try:
|
|
if not self.d.flag.exists():
|
|
self.d.maintenance()
|
|
self.d.stop()
|
|
prior = bundle / 'prior'
|
|
shutil.copytree(self.d.app / 'data', prior / 'data')
|
|
shutil.copytree(self.d.app / 'current', prior / 'release')
|
|
atomic(prior / 'backup.json', json.dumps(tree_metadata(prior)))
|
|
durable_tree(prior)
|
|
sync_directory(bundle)
|
|
self.record(state, 'backed_up')
|
|
self.install_data(source / 'data', bundle)
|
|
link(self.d.app / 'current', release.relative_to(self.d.app))
|
|
self.record(state, 'installed')
|
|
self.activate(state['new'])
|
|
if not state['wasRunning']:
|
|
self.d.stop()
|
|
atomic(self.d.app / 'deployed.json', json.dumps(state['new']))
|
|
self.record(state, 'published')
|
|
self.finish(state)
|
|
except Exception:
|
|
# Preserve the original exception even after a successful rollback.
|
|
try:
|
|
self.recover()
|
|
except Exception:
|
|
try:
|
|
self.d.stop()
|
|
except Exception:
|
|
pass
|
|
print('Restore recovery failed; maintenance and prior snapshot retained. Run restore --recover after inspection.', file=sys.stderr)
|
|
raise
|
|
return {'result': 'RESTORED', 'commit': state['new']['commit'], 'priorSnapshot': str(bundle / 'prior')}
|
|
|
|
|
|
if __name__ == '__main__':
|
|
os.umask(0o077)
|
|
deployment = Deployment(Path.home())
|
|
operation = OperatorRestore(deployment)
|
|
with (deployment.app / '.deploy.lock').open('a') as lock:
|
|
fcntl.flock(lock, fcntl.LOCK_EX | fcntl.LOCK_NB)
|
|
if '--recover' in sys.argv:
|
|
operation.recover()
|
|
print(json.dumps({'result': 'RECOVERED'}))
|
|
else:
|
|
with tempfile.NamedTemporaryFile(prefix='restore-upload-', dir=deployment.app) as archive:
|
|
shutil.copyfileobj(sys.stdin.buffer, archive)
|
|
archive.flush()
|
|
print(json.dumps(operation.restore(archive.name)))
|