360 lines
18 KiB
Python
360 lines
18 KiB
Python
"""Exercise real release/database file transitions with a controlled service."""
|
|
import importlib.util
|
|
import hashlib
|
|
import io
|
|
import json
|
|
from pathlib import Path
|
|
import tarfile
|
|
import tempfile
|
|
import unittest
|
|
from unittest.mock import patch
|
|
|
|
spec = importlib.util.spec_from_file_location('update', Path(__file__).resolve().parents[2] / 'scripts/deploy-update-remote.py')
|
|
update = importlib.util.module_from_spec(spec)
|
|
spec.loader.exec_module(update)
|
|
|
|
|
|
def release(commit):
|
|
payload = commit.encode()
|
|
digest = hashlib.sha256(payload).hexdigest()
|
|
return {'schemaVersion': 1, 'runtime': 'java-25', 'version': 'test', 'commit': commit,
|
|
'sha256': digest, 'files': {'werkjournal.jar': digest}, 'ancestors': [c * 40 for c in 'abcd' if c < commit[0]]}, payload
|
|
|
|
|
|
def archive(commit):
|
|
value, payload = release(commit)
|
|
stream = io.BytesIO()
|
|
with tarfile.open(fileobj=stream, mode='w') as tar:
|
|
for name, data in [('manifest.json', json.dumps(value).encode()), ('werkjournal.jar', payload)]:
|
|
entry = tarfile.TarInfo(name); entry.size = len(data)
|
|
tar.addfile(entry, io.BytesIO(data))
|
|
stream.seek(0)
|
|
return stream
|
|
|
|
|
|
class Controlled(update.Deployment):
|
|
def __init__(self, home):
|
|
super().__init__(home)
|
|
self.running = True
|
|
self.fail = set()
|
|
self.fail_restore = False
|
|
self.closing_bytes = None
|
|
|
|
def preflight(self, check_data=True):
|
|
pass
|
|
|
|
def stop(self):
|
|
assert self.flag.exists()
|
|
if self.closing_bytes is not None:
|
|
(self.app / 'data/werkjournal.mv.db').write_bytes(self.closing_bytes)
|
|
self.closing_bytes = None
|
|
self.running = False
|
|
|
|
def start(self):
|
|
assert self.flag.exists()
|
|
self.running = True
|
|
if (self.app / 'current').resolve().name in self.fail:
|
|
(self.app / 'data/werkjournal.mv.db').write_bytes(b'partially migrated')
|
|
|
|
def ready(self, expected):
|
|
assert self.running
|
|
if expected['commit'] in self.fail or self.fail_restore:
|
|
raise RuntimeError('Injected startup failure')
|
|
|
|
def maintenance(self):
|
|
update.atomic(self.flag, 'Wartungsmodus')
|
|
|
|
|
|
class UpgradeTest(unittest.TestCase):
|
|
def setUp(self):
|
|
self.temp = tempfile.TemporaryDirectory()
|
|
self.addCleanup(self.temp.cleanup)
|
|
self.d = Controlled(Path(self.temp.name))
|
|
self.d.site.mkdir(parents=True)
|
|
self.old = 'a' * 40
|
|
folder = self.d.app / 'releases' / self.old
|
|
folder.mkdir(parents=True)
|
|
value, payload = release(self.old)
|
|
(folder / 'manifest.json').write_text(json.dumps(value))
|
|
(folder / 'werkjournal.jar').write_bytes(payload)
|
|
(self.d.app / 'current').symlink_to(Path('releases') / self.old)
|
|
(self.d.app / 'deployed.json').write_text(json.dumps(value))
|
|
(self.d.app / 'data').mkdir()
|
|
self.db = self.d.app / 'data/werkjournal.mv.db'
|
|
self.db.write_bytes(b'bookings before first upgrade')
|
|
|
|
def test_stopped_java_exit_143_is_closed_but_any_live_pid_is_rejected(self):
|
|
for state in ['inactive', 'failed']:
|
|
with patch.object(self.d, 'command', side_effect=lambda *args: state if '--property=ActiveState' in args else '0' if '--property=MainPID' in args else ''):
|
|
update.Deployment.stop(self.d)
|
|
|
|
with patch.object(self.d, 'command', side_effect=lambda *args: 'failed' if '--property=ActiveState' in args else '123' if '--property=MainPID' in args else ''):
|
|
with self.assertRaisesRegex(RuntimeError, 'process'):
|
|
update.Deployment.stop(self.d)
|
|
|
|
def test_upload_is_repeatable_and_never_changes_running_release_or_data(self):
|
|
self.d.flag.write_text('Existing operator maintenance')
|
|
before = (self.db.read_bytes(), (self.d.app / 'deployed.json').read_bytes(),
|
|
(self.d.app / 'current').readlink(), self.d.flag.read_bytes())
|
|
with patch.object(self.d, 'stop', side_effect=AssertionError('Upload must not stop')), \
|
|
patch.object(self.d, 'start', side_effect=AssertionError('Upload must not start')):
|
|
for _ in range(2):
|
|
result = self.d.upload(archive('b' * 40))
|
|
self.assertEqual('STAGED', result['result'])
|
|
self.assertEqual('b' * 40, result['commit'])
|
|
after = (self.db.read_bytes(), (self.d.app / 'deployed.json').read_bytes(),
|
|
(self.d.app / 'current').readlink(), self.d.flag.read_bytes())
|
|
self.assertEqual(before, after)
|
|
self.assertTrue(self.d.running)
|
|
self.assertTrue((self.d.app / 'releases' / ('b' * 40) / 'werkjournal.jar').exists())
|
|
self.assertFalse(self.d.pending.exists())
|
|
self.assertEqual([], list(self.d.app.glob('candidate-*')))
|
|
|
|
def test_upload_rejects_pending_deployment_without_staging(self):
|
|
self.d.pending.write_text('{"phase":"prepared"}')
|
|
with self.assertRaisesRegex(RuntimeError, 'pending'):
|
|
self.d.upload(archive('b' * 40))
|
|
self.assertFalse((self.d.app / 'releases' / ('b' * 40)).exists())
|
|
self.assertEqual('{"phase":"prepared"}', self.d.pending.read_text())
|
|
|
|
def test_uploaded_candidate_requires_explicit_deploy_and_can_then_be_activated(self):
|
|
self.d.upload(archive('b' * 40))
|
|
self.assertEqual(self.old, (self.d.app / 'current').resolve().name)
|
|
self.assertEqual('PASS', self.d.deploy(archive('b' * 40))['result'])
|
|
self.assertEqual('b' * 40, (self.d.app / 'current').resolve().name)
|
|
self.assertEqual(b'bookings before first upgrade', self.db.read_bytes())
|
|
self.assertEqual(self.old, (self.d.app / 'previous').resolve().name)
|
|
|
|
def test_upload_rejects_bad_checksum_and_cleans_temporary_candidate(self):
|
|
value, payload = release('b' * 40)
|
|
stream = io.BytesIO()
|
|
with tarfile.open(fileobj=stream, mode='w') as tar:
|
|
for name, data in [('manifest.json', json.dumps(value).encode()), ('werkjournal.jar', payload + b'corrupt')]:
|
|
entry = tarfile.TarInfo(name)
|
|
entry.size = len(data)
|
|
tar.addfile(entry, io.BytesIO(data))
|
|
stream.seek(0)
|
|
with self.assertRaises(ValueError):
|
|
self.d.upload(stream)
|
|
self.assertFalse((self.d.app / 'releases' / ('b' * 40)).exists())
|
|
self.assertEqual([], list(self.d.app.glob('candidate-*')))
|
|
self.assertTrue(self.d.running)
|
|
|
|
def test_snapshot_includes_transaction_flushed_during_stop(self):
|
|
self.d.closing_bytes = b'last committed transaction on shutdown'
|
|
self.d.deploy(archive('b' * 40))
|
|
self.assertEqual(b'last committed transaction on shutdown',
|
|
((self.d.app / 'rollback').resolve() / 'data/werkjournal.mv.db').read_bytes())
|
|
|
|
def test_shared_restore_copy_failure_does_not_displace_live_database(self):
|
|
source = self.d.app / 'snapshot'
|
|
source.mkdir()
|
|
(source / 'werkjournal.mv.db').write_bytes(b'snapshot')
|
|
displaced = self.d.app / 'displaced'
|
|
displaced.mkdir()
|
|
with patch.object(update, 'durable_tree', side_effect=OSError('flush failed')):
|
|
with self.assertRaisesRegex(OSError, 'flush failed'):
|
|
update.replace_database(source, self.d.app / 'data', displaced)
|
|
self.assertEqual(b'bookings before first upgrade', self.db.read_bytes())
|
|
self.assertEqual([], list(displaced.iterdir()))
|
|
self.assertEqual([], list(self.d.app.glob('restore-data-*')))
|
|
|
|
def test_runtime_preflight_rejects_external_database_and_inline_overrides(self):
|
|
unit = self.d.home / '.config/systemd/user/werkjournal-backend.service'
|
|
unit.parent.mkdir(parents=True)
|
|
text = f"""[Service]
|
|
WorkingDirectory={self.d.app}
|
|
ExecStart={self.d.home}/opt/jdk25/bin/java -Xms32m -Xmx384m -XX:MaxMetaspaceSize=192m -XX:ActiveProcessorCount=2 -Djava.awt.headless=true -jar {self.d.app}/current/werkjournal.jar
|
|
Environment=BACKEND_PORT=18090
|
|
Environment=SERVER_ADDRESS=127.0.0.1
|
|
Environment=SERVER_FORWARD_HEADERS_STRATEGY=native
|
|
Environment=SERVER_TOMCAT_REMOTEIP_HOST_HEADER=x-forwarded-host
|
|
EnvironmentFile=-%h/.config/werkjournal/environment
|
|
"""
|
|
unit.write_text(text)
|
|
with patch.object(self.d, 'command', side_effect=lambda *args: 'no' if '--property=NeedDaemonReload' in args else ''):
|
|
update.Deployment.preflight(self.d)
|
|
env = self.d.home / '.config/werkjournal/environment'
|
|
env.parent.mkdir(parents=True)
|
|
env.write_text('WERKJOURNAL_OIDC_CLIENT_SECRET="test secret"\n')
|
|
env.chmod(0o600)
|
|
update.Deployment.preflight(self.d)
|
|
env.write_text('WERKJOURNAL_DATABASE_URL=jdbc:h2:file:/outside/database\n')
|
|
with self.assertRaisesRegex(RuntimeError, 'outside'):
|
|
update.Deployment.preflight(self.d)
|
|
unit.write_text(text + 'Environment=SPRING_DATASOURCE_URL=jdbc:h2:file:/outside/database\n')
|
|
with self.assertRaisesRegex(RuntimeError, 'inline'):
|
|
update.Deployment.preflight(self.d)
|
|
|
|
def test_corrupt_backup_is_not_restored_and_service_remains_stopped(self):
|
|
self.d.fail.add('b' * 40)
|
|
self.d.fail_restore = True
|
|
with self.assertRaises(RuntimeError):
|
|
self.d.deploy(archive('b' * 40))
|
|
state = json.loads(self.d.pending.read_text())
|
|
(self.d.app / 'rollbacks' / state['backup'] / 'data/werkjournal.mv.db').write_bytes(b'corrupted')
|
|
self.d.fail_restore = False
|
|
with self.assertRaisesRegex(RuntimeError, 'checksum'):
|
|
self.d.recover()
|
|
self.assertTrue(self.d.flag.exists())
|
|
self.assertFalse(self.d.running)
|
|
|
|
def test_success_retains_matching_application_and_closed_database(self):
|
|
self.d.deploy(archive('b' * 40))
|
|
self.assertEqual('b' * 40, (self.d.app / 'current').resolve().name)
|
|
self.assertEqual(self.old, (self.d.app / 'previous').resolve().name)
|
|
backup = (self.d.app / 'rollback').resolve()
|
|
self.assertEqual(self.old, json.loads((backup / 'manifest.json').read_text())['commit'])
|
|
self.assertEqual(b'bookings before first upgrade', (backup / 'data/werkjournal.mv.db').read_bytes())
|
|
self.assertFalse(self.d.flag.exists())
|
|
self.assertFalse(self.d.pending.exists())
|
|
|
|
def test_two_failed_upgrades_preserve_recent_bookings_and_last_successful_fallback(self):
|
|
self.d.deploy(archive('b' * 40))
|
|
fallback = (self.d.app / 'rollback').resolve()
|
|
for candidate in ['c' * 40, 'd' * 40]:
|
|
self.db.write_bytes(candidate.encode() + b' latest bookings')
|
|
self.d.fail.add(candidate)
|
|
with self.assertRaisesRegex(RuntimeError, 'Injected'):
|
|
self.d.deploy(archive(candidate))
|
|
self.assertEqual(candidate.encode() + b' latest bookings', self.db.read_bytes())
|
|
self.assertEqual('b' * 40, (self.d.app / 'current').resolve().name)
|
|
self.assertEqual(self.old, (self.d.app / 'previous').resolve().name)
|
|
self.assertEqual(fallback, (self.d.app / 'rollback').resolve())
|
|
self.assertFalse(self.d.flag.exists())
|
|
|
|
def test_failed_rollback_leaves_maintenance_and_recoverable_journal(self):
|
|
self.d.fail.add('b' * 40)
|
|
self.d.fail_restore = True
|
|
with self.assertRaises(RuntimeError):
|
|
self.d.deploy(archive('b' * 40))
|
|
self.assertTrue(self.d.flag.exists())
|
|
self.assertTrue(self.d.pending.exists())
|
|
self.assertFalse(self.d.running)
|
|
self.d.fail_restore = False
|
|
self.d.recover()
|
|
self.assertFalse(self.d.flag.exists())
|
|
self.assertEqual(b'bookings before first upgrade', self.db.read_bytes())
|
|
|
|
def test_crash_after_publication_does_not_restore_over_new_bookings(self):
|
|
original = self.d.finalize
|
|
self.d.finalize = lambda state: (_ for _ in ()).throw(RuntimeError('crash after publication'))
|
|
with self.assertRaisesRegex(RuntimeError, 'crash'):
|
|
self.d.deploy(archive('b' * 40))
|
|
self.db.write_bytes(b'new public bookings')
|
|
self.d.finalize = original
|
|
self.d.recover()
|
|
self.assertEqual(b'new public bookings', self.db.read_bytes())
|
|
self.assertEqual('b' * 40, (self.d.app / 'current').resolve().name)
|
|
self.assertFalse(self.d.pending.exists())
|
|
|
|
def test_publication_intent_prevents_rollback_if_flag_reappears_after_power_loss(self):
|
|
original = self.d.finalize
|
|
self.d.finalize = lambda state: (_ for _ in ()).throw(RuntimeError('power loss'))
|
|
with self.assertRaises(RuntimeError):
|
|
self.d.deploy(archive('b' * 40))
|
|
self.db.write_bytes(b'accepted after publication')
|
|
self.d.flag.write_text('simulated non-durable unlink')
|
|
self.d.finalize = original
|
|
self.d.recover()
|
|
self.assertEqual(b'accepted after publication', self.db.read_bytes())
|
|
self.assertFalse(self.d.flag.exists())
|
|
|
|
def test_old_candidate_cannot_replace_newer_successful_release(self):
|
|
self.d.deploy(archive('b' * 40))
|
|
with self.assertRaisesRegex(RuntimeError, 'descend'):
|
|
self.d.deploy(archive(self.old))
|
|
self.assertEqual('b' * 40, (self.d.app / 'current').resolve().name)
|
|
self.assertFalse(self.d.flag.exists())
|
|
|
|
def test_recovery_after_interruption_before_publication_restores_snapshot(self):
|
|
original = self.d.ready
|
|
def interrupted(expected):
|
|
if expected['commit'] == 'b' * 40:
|
|
self.db.write_bytes(b'candidate migration')
|
|
raise KeyboardInterrupt()
|
|
original(expected)
|
|
self.d.ready = interrupted
|
|
with self.assertRaises(KeyboardInterrupt):
|
|
self.d.deploy(archive('b' * 40))
|
|
self.assertTrue(self.d.flag.exists())
|
|
self.d.ready = original
|
|
self.d.recover()
|
|
self.assertEqual(b'bookings before first upgrade', self.db.read_bytes())
|
|
self.assertEqual(self.old, (self.d.app / 'current').resolve().name)
|
|
|
|
def test_changed_artifact_is_rejected_before_stopping(self):
|
|
broken = archive('b' * 40)
|
|
data = broken.getvalue().replace(b'b' * 40, b'c' * 40)
|
|
with self.assertRaises(ValueError):
|
|
self.d.deploy(io.BytesIO(data))
|
|
self.assertTrue(self.d.running)
|
|
self.assertFalse(self.d.flag.exists())
|
|
self.assertEqual(b'bookings before first upgrade', self.db.read_bytes())
|
|
|
|
|
|
if __name__ == '__main__':
|
|
unittest.main()
|
|
|
|
class ReadinessTest(unittest.TestCase):
|
|
def test_login_probe_checks_version_and_flow_then_logs_out(self):
|
|
from unittest.mock import Mock
|
|
def response(path, body, content_type='text/html'):
|
|
value = io.BytesIO(body.encode())
|
|
value.url = 'http://127.0.0.1:18090' + path
|
|
value.headers = {'Content-Type': content_type}
|
|
return value
|
|
opener = Mock()
|
|
opener.open.side_effect = [
|
|
response('/login', 'login'),
|
|
response('/login', '<input name="_csrf" value="first">'),
|
|
response('/', 'app'),
|
|
response('/api/v1/info', '{"application":"werkjournal","version":"test"}', 'application/json'),
|
|
response('/', '{"appConfig":{}}', 'application/json'),
|
|
response('/login', '<input name="_csrf" value="second">'),
|
|
response('/login', 'logged out'),
|
|
]
|
|
with patch.object(update.urllib.request, 'build_opener', return_value=opener):
|
|
self.assertTrue(update.Deployment(Path('/tmp')).ready_once({'version': 'test'}))
|
|
calls = opener.open.call_args_list
|
|
self.assertTrue(calls[2].args[0].endswith('/login/guest'))
|
|
self.assertEqual(b'_csrf=first', calls[2].kwargs['data'])
|
|
self.assertTrue(calls[-1].args[0].endswith('/logout'))
|
|
self.assertEqual(b'_csrf=second', calls[-1].kwargs['data'])
|
|
|
|
def test_guest_mode_wrong_version_is_not_ready(self):
|
|
from unittest.mock import Mock
|
|
value = io.BytesIO(b'{"application":"werkjournal","version":"wrong"}')
|
|
value.url = 'http://127.0.0.1:18090/api/v1/info'
|
|
opener = Mock()
|
|
opener.open.return_value = value
|
|
with patch.object(update.urllib.request, 'build_opener', return_value=opener):
|
|
self.assertFalse(update.Deployment(Path('/tmp')).ready_once({'version': 'test'}))
|
|
|
|
def test_public_info_still_requires_flow_login_and_legacy_info_can_be_denied(self):
|
|
from unittest.mock import Mock
|
|
for public_info in (True, False):
|
|
with self.subTest(public_info=public_info), tempfile.TemporaryDirectory() as directory:
|
|
home = Path(directory)
|
|
jar = home / 'opt/werkjournal/current/werkjournal.jar'
|
|
jar.parent.mkdir(parents=True)
|
|
with update.zipfile.ZipFile(jar, 'w') as archive:
|
|
archive.writestr('META-INF/build-info.properties', 'build.version=test\n')
|
|
def response(path, body, kind='text/html'):
|
|
value = io.BytesIO(body.encode())
|
|
value.url = 'http://127.0.0.1:18090' + path
|
|
value.headers = {'Content-Type': kind}
|
|
return value
|
|
info = '{"application":"werkjournal","version":"test"}'
|
|
replies = ([response('/api/v1/info', info, 'application/json'), response('/login', 'login')]
|
|
if public_info else [response('/login', 'login')])
|
|
replies += [response('/login', '<input name="_csrf" value="first">'), response('/', 'app'),
|
|
response('/api/v1/info', info, 'application/json') if public_info else
|
|
update.urllib.error.HTTPError('http://127.0.0.1:18090/api/v1/info', 403, 'denied', {}, None),
|
|
response('/', '{"appConfig":{}}', 'application/json'),
|
|
response('/login', '<input name="_csrf" value="second">'), response('/login', 'out')]
|
|
opener = Mock(); opener.open.side_effect = replies
|
|
with patch.object(update.urllib.request, 'build_opener', return_value=opener):
|
|
self.assertTrue(update.Deployment(home).ready_once({'version': 'test'}))
|
|
self.assertTrue(opener.open.call_args.args[0].endswith('/logout'))
|