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