497 lines
19 KiB
Python
497 lines
19 KiB
Python
from __future__ import annotations
|
|
|
|
import asyncio
|
|
import json
|
|
import stat
|
|
from pathlib import Path
|
|
|
|
import httpx
|
|
import pytest
|
|
from alembic import command
|
|
from backup_tool.cli import (
|
|
build_alembic_config,
|
|
recovery_export_payload,
|
|
rotate_repository_key,
|
|
)
|
|
from backup_tool.cli import main as cli_main
|
|
from backup_tool.config import Settings
|
|
from backup_tool.db.engine import create_engine
|
|
from backup_tool.db.models import Backup, Execution, Repository, RepositoryDataKeyEpoch
|
|
from backup_tool.repository import (
|
|
begin_key_rotation,
|
|
initialize,
|
|
inspect_repository,
|
|
replace_active_data_key,
|
|
)
|
|
from backup_tool.security.repository_crypto import create_data_key
|
|
from backup_tool.worker import Worker
|
|
from sqlalchemy import select
|
|
|
|
from .test_repository_safety import settings_for
|
|
|
|
PASSWORD = "correct-horse-battery-staple"
|
|
|
|
|
|
async def login(client: httpx.AsyncClient) -> dict[str, str]:
|
|
response = await client.post("/api/v2/setup", json={"username": "admin", "password": PASSWORD})
|
|
assert response.status_code == 201
|
|
return {"X-CSRF-Token": client.cookies["backup_tool_csrf"]}
|
|
|
|
|
|
def test_encrypted_initialization_creates_private_data_key(tmp_path: Path) -> None:
|
|
settings = settings_for(tmp_path)
|
|
initialized = initialize(settings, "encrypted", "none", "aes-256-gcm")
|
|
assert initialized.data_key_id is not None
|
|
assert initialized.data_key_path is not None
|
|
assert stat.S_IMODE(initialized.data_key_path.stat().st_mode) == 0o600
|
|
try:
|
|
payload = json.loads((initialized.root / "repository.json").read_text(encoding="utf-8"))
|
|
except (OSError, json.JSONDecodeError) as error:
|
|
raise AssertionError("encrypted repository metadata is unreadable") from error
|
|
assert payload["encryption"] == {
|
|
"mode": "aes-256-gcm",
|
|
"key_id": initialized.data_key_id,
|
|
}
|
|
inspected = inspect_repository(settings, initialized.root)
|
|
assert inspected.encryption == "aes-256-gcm"
|
|
assert inspected.data_key_id == initialized.data_key_id
|
|
|
|
|
|
def test_rotation_rejects_metadata_database_epoch_mismatch(tmp_path: Path) -> None:
|
|
settings = settings_for(tmp_path)
|
|
command.upgrade(build_alembic_config(settings), "head")
|
|
initialized = initialize(settings, "encrypted", "none", "aes-256-gcm")
|
|
assert initialized.data_key_id is not None
|
|
|
|
async def create_repository() -> str:
|
|
engine = create_engine(settings)
|
|
try:
|
|
from sqlalchemy.ext.asyncio import async_sessionmaker
|
|
|
|
sessions = async_sessionmaker(engine, expire_on_commit=False)
|
|
async with sessions() as db:
|
|
repository = Repository(
|
|
name="encrypted",
|
|
root=str(initialized.root),
|
|
format_version=initialized.format_version,
|
|
compression=initialized.compression,
|
|
encryption=initialized.encryption,
|
|
signing_key_id=initialized.signing_key_id,
|
|
signing_public_key=initialized.signing_public_key,
|
|
active_data_key_id=initialized.data_key_id,
|
|
)
|
|
db.add(repository)
|
|
await db.flush()
|
|
db.add(
|
|
RepositoryDataKeyEpoch(
|
|
repository_id=repository.id,
|
|
key_id=initialized.data_key_id,
|
|
state="active",
|
|
)
|
|
)
|
|
await db.commit()
|
|
return repository.id
|
|
finally:
|
|
await engine.dispose()
|
|
|
|
repository_id = asyncio.run(create_repository())
|
|
replacement_id, replacement_path = create_data_key(settings, initialized.repository_id)
|
|
try:
|
|
replace_active_data_key(initialized.root, initialized.data_key_id, replacement_id)
|
|
with pytest.raises(ValueError, match="repository encryption metadata is invalid"):
|
|
asyncio.run(rotate_repository_key(settings, repository_id))
|
|
finally:
|
|
replacement_path.unlink(missing_ok=True)
|
|
|
|
|
|
def test_rotation_reconciliation_clears_stale_rollback_journal_after_key_removal(
|
|
tmp_path: Path,
|
|
) -> None:
|
|
settings = settings_for(tmp_path)
|
|
command.upgrade(build_alembic_config(settings), "head")
|
|
initialized = initialize(settings, "encrypted", "none", "aes-256-gcm")
|
|
assert initialized.data_key_id is not None
|
|
|
|
async def create_repository() -> str:
|
|
engine = create_engine(settings)
|
|
try:
|
|
from sqlalchemy.ext.asyncio import async_sessionmaker
|
|
|
|
sessions = async_sessionmaker(engine, expire_on_commit=False)
|
|
async with sessions() as db:
|
|
repository = Repository(
|
|
name="encrypted",
|
|
root=str(initialized.root),
|
|
format_version=initialized.format_version,
|
|
compression=initialized.compression,
|
|
encryption=initialized.encryption,
|
|
signing_key_id=initialized.signing_key_id,
|
|
signing_public_key=initialized.signing_public_key,
|
|
active_data_key_id=initialized.data_key_id,
|
|
)
|
|
db.add(repository)
|
|
await db.flush()
|
|
db.add(
|
|
RepositoryDataKeyEpoch(
|
|
repository_id=repository.id,
|
|
key_id=initialized.data_key_id,
|
|
state="active",
|
|
)
|
|
)
|
|
await db.commit()
|
|
return repository.id
|
|
finally:
|
|
await engine.dispose()
|
|
|
|
repository_id = asyncio.run(create_repository())
|
|
new_key_id, new_key_path = create_data_key(settings, initialized.repository_id)
|
|
begin_key_rotation(
|
|
initialized.root,
|
|
repository_id,
|
|
initialized.repository_id,
|
|
initialized.data_key_id,
|
|
new_key_id,
|
|
)
|
|
new_key_path.unlink()
|
|
assert (initialized.root / ".key-rotation.json").is_file()
|
|
|
|
async def reconcile_and_assert() -> None:
|
|
worker = Worker(settings, owner="stale-rollback-journal-worker")
|
|
try:
|
|
assert await worker.startup() == 1
|
|
finally:
|
|
await worker.engine.dispose()
|
|
engine = create_engine(settings)
|
|
try:
|
|
from sqlalchemy.ext.asyncio import async_sessionmaker
|
|
|
|
sessions = async_sessionmaker(engine, expire_on_commit=False)
|
|
async with sessions() as db:
|
|
repository = await db.get(Repository, repository_id)
|
|
assert repository is not None
|
|
assert repository.active_data_key_id == initialized.data_key_id
|
|
inspected = inspect_repository(settings, initialized.root)
|
|
assert inspected.data_key_id == initialized.data_key_id
|
|
finally:
|
|
await engine.dispose()
|
|
|
|
asyncio.run(reconcile_and_assert())
|
|
assert not (initialized.root / ".key-rotation.json").exists()
|
|
|
|
|
|
def test_rotation_crash_after_db_commit_recovers_on_worker_startup(
|
|
tmp_path: Path,
|
|
) -> None:
|
|
settings = settings_for(tmp_path)
|
|
command.upgrade(build_alembic_config(settings), "head")
|
|
initialized = initialize(settings, "encrypted", "none", "aes-256-gcm")
|
|
assert initialized.data_key_id is not None
|
|
|
|
async def create_repository() -> str:
|
|
engine = create_engine(settings)
|
|
try:
|
|
from sqlalchemy.ext.asyncio import async_sessionmaker
|
|
|
|
sessions = async_sessionmaker(engine, expire_on_commit=False)
|
|
async with sessions() as db:
|
|
repository = Repository(
|
|
name="encrypted",
|
|
root=str(initialized.root),
|
|
format_version=initialized.format_version,
|
|
compression=initialized.compression,
|
|
encryption=initialized.encryption,
|
|
signing_key_id=initialized.signing_key_id,
|
|
signing_public_key=initialized.signing_public_key,
|
|
active_data_key_id=initialized.data_key_id,
|
|
)
|
|
db.add(repository)
|
|
await db.flush()
|
|
db.add(
|
|
RepositoryDataKeyEpoch(
|
|
repository_id=repository.id,
|
|
key_id=initialized.data_key_id,
|
|
state="active",
|
|
)
|
|
)
|
|
await db.commit()
|
|
return repository.id
|
|
finally:
|
|
await engine.dispose()
|
|
|
|
repository_id = asyncio.run(create_repository())
|
|
|
|
def interrupted_after_database_commit() -> None:
|
|
raise OSError("simulated process loss after database commit")
|
|
|
|
with pytest.raises(OSError, match="simulated process loss"):
|
|
asyncio.run(
|
|
rotate_repository_key(
|
|
settings,
|
|
repository_id,
|
|
after_database_commit=interrupted_after_database_commit,
|
|
)
|
|
)
|
|
assert inspect_repository(settings, initialized.root).data_key_id == initialized.data_key_id
|
|
assert (initialized.root / ".key-rotation.json").is_file()
|
|
|
|
async def reconcile_and_assert() -> None:
|
|
worker = Worker(settings, owner="rotation-recovery-worker")
|
|
try:
|
|
assert await worker.startup() == 1
|
|
finally:
|
|
await worker.engine.dispose()
|
|
engine = create_engine(settings)
|
|
try:
|
|
from sqlalchemy.ext.asyncio import async_sessionmaker
|
|
|
|
sessions = async_sessionmaker(engine, expire_on_commit=False)
|
|
async with sessions() as db:
|
|
repository = await db.get(Repository, repository_id)
|
|
assert repository is not None
|
|
assert repository.active_data_key_id is not None
|
|
inspected = inspect_repository(settings, initialized.root)
|
|
assert inspected.data_key_id == repository.active_data_key_id
|
|
active_data_key_id = repository.active_data_key_id
|
|
exported = await recovery_export_payload(settings)
|
|
exported_repository = exported["catalog"]["repositories"][0]
|
|
assert exported_repository["active_data_key_id"] == active_data_key_id
|
|
finally:
|
|
await engine.dispose()
|
|
|
|
asyncio.run(reconcile_and_assert())
|
|
assert not (initialized.root / ".key-rotation.json").exists()
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_encrypted_repository_worker_backup_and_restore(
|
|
app_client: tuple[httpx.AsyncClient, Settings],
|
|
) -> None:
|
|
client, settings = app_client
|
|
source_root = settings.local_source_roots[0] / "project"
|
|
source_root.mkdir()
|
|
plaintext = b"encrypted backup content\n"
|
|
(source_root / "hello.txt").write_bytes(plaintext)
|
|
initialized = initialize(settings, "encrypted", "none", "aes-256-gcm")
|
|
assert initialized.data_key_id is not None
|
|
|
|
engine = create_engine(settings)
|
|
try:
|
|
from sqlalchemy.ext.asyncio import async_sessionmaker
|
|
|
|
sessions = async_sessionmaker(engine, expire_on_commit=False)
|
|
async with sessions() as db:
|
|
repository = Repository(
|
|
name="encrypted",
|
|
root=str(initialized.root),
|
|
format_version=initialized.format_version,
|
|
compression=initialized.compression,
|
|
encryption=initialized.encryption,
|
|
signing_key_id=initialized.signing_key_id,
|
|
signing_public_key=initialized.signing_public_key,
|
|
active_data_key_id=initialized.data_key_id,
|
|
)
|
|
db.add(repository)
|
|
await db.flush()
|
|
db.add(
|
|
RepositoryDataKeyEpoch(
|
|
repository_id=repository.id,
|
|
key_id=initialized.data_key_id,
|
|
state="active",
|
|
)
|
|
)
|
|
await db.commit()
|
|
repository_id = repository.id
|
|
finally:
|
|
await engine.dispose()
|
|
|
|
headers = await login(client)
|
|
source_response = await client.post(
|
|
"/api/v2/sources",
|
|
json={
|
|
"name": "local",
|
|
"kind": "local",
|
|
"public_config": {"root": str(source_root)},
|
|
},
|
|
headers=headers,
|
|
)
|
|
assert source_response.status_code == 201
|
|
job_response = await client.post(
|
|
"/api/v2/jobs",
|
|
json={
|
|
"name": "encrypted-backup",
|
|
"source_id": source_response.json()["id"],
|
|
"repository_id": repository_id,
|
|
"requested_mode": "full",
|
|
"exclusions": [],
|
|
"retention": {},
|
|
"enabled": True,
|
|
"allow_empty": False,
|
|
},
|
|
headers=headers,
|
|
)
|
|
assert job_response.status_code == 201
|
|
execution_response = await client.post(
|
|
f"/api/v2/jobs/{job_response.json()['id']}/executions", headers=headers
|
|
)
|
|
assert execution_response.status_code == 202
|
|
execution_id = execution_response.json()["id"]
|
|
|
|
worker = Worker(settings, owner="encrypted-backup-worker")
|
|
try:
|
|
assert await worker.run_once()
|
|
finally:
|
|
await worker.engine.dispose()
|
|
|
|
engine = create_engine(settings)
|
|
try:
|
|
from sqlalchemy.ext.asyncio import async_sessionmaker
|
|
|
|
sessions = async_sessionmaker(engine, expire_on_commit=False)
|
|
async with sessions() as db:
|
|
execution = await db.get(Execution, execution_id)
|
|
backup = await db.scalar(select(Backup).where(Backup.execution_id == execution_id))
|
|
finally:
|
|
await engine.dispose()
|
|
assert execution is not None
|
|
assert execution.state == "committed"
|
|
assert backup is not None
|
|
|
|
blob = next((initialized.root / "blobs" / "sha256").iterdir())
|
|
assert blob.read_bytes().startswith(b"BTENC\x01")
|
|
assert plaintext not in blob.read_bytes()
|
|
manifest_path = initialized.root / "manifests" / f"{backup.manifest_id}.json"
|
|
stored_manifest = manifest_path.read_bytes()
|
|
assert stored_manifest.startswith(b"BTENC\x01")
|
|
assert b'"entries"' not in stored_manifest
|
|
assert b"hello.txt" not in stored_manifest
|
|
|
|
assert (
|
|
await asyncio.to_thread(
|
|
cli_main,
|
|
["admin", "repository-key", "rotate", "--repository-id", repository_id],
|
|
settings=settings,
|
|
)
|
|
== 0
|
|
)
|
|
engine = create_engine(settings)
|
|
try:
|
|
from sqlalchemy.ext.asyncio import async_sessionmaker
|
|
|
|
sessions = async_sessionmaker(engine, expire_on_commit=False)
|
|
async with sessions() as db:
|
|
repository = await db.get(Repository, repository_id)
|
|
epochs = list(
|
|
(
|
|
await db.scalars(
|
|
select(RepositoryDataKeyEpoch).where(
|
|
RepositoryDataKeyEpoch.repository_id == repository_id
|
|
)
|
|
)
|
|
).all()
|
|
)
|
|
finally:
|
|
await engine.dispose()
|
|
assert repository is not None
|
|
assert repository.active_data_key_id != initialized.data_key_id
|
|
epoch_states = {f"{epoch.key_id}:{epoch.state}" for epoch in epochs}
|
|
expected_epoch_states = {
|
|
f"{initialized.data_key_id}:retired",
|
|
f"{repository.active_data_key_id}:active",
|
|
}
|
|
assert epoch_states == expected_epoch_states
|
|
assert (
|
|
inspect_repository(settings, initialized.root).data_key_id == repository.active_data_key_id
|
|
)
|
|
|
|
plaintext_after_rotation = b"encrypted content after rotation\n"
|
|
(source_root / "hello.txt").write_bytes(plaintext_after_rotation)
|
|
second_execution_response = await client.post(
|
|
f"/api/v2/jobs/{job_response.json()['id']}/executions", headers=headers
|
|
)
|
|
assert second_execution_response.status_code == 202
|
|
second_execution_id = second_execution_response.json()["id"]
|
|
second_worker = Worker(settings, owner="encrypted-rotated-backup-worker")
|
|
try:
|
|
assert await second_worker.run_once()
|
|
finally:
|
|
await second_worker.engine.dispose()
|
|
engine = create_engine(settings)
|
|
try:
|
|
from sqlalchemy.ext.asyncio import async_sessionmaker
|
|
|
|
sessions = async_sessionmaker(engine, expire_on_commit=False)
|
|
async with sessions() as db:
|
|
second_backup = await db.scalar(
|
|
select(Backup).where(Backup.execution_id == second_execution_id)
|
|
)
|
|
finally:
|
|
await engine.dispose()
|
|
assert second_backup is not None
|
|
assert second_backup.data_key_id == repository.active_data_key_id
|
|
|
|
destination = settings.restore_roots[0] / "restored"
|
|
restore_response = await client.post(
|
|
f"/api/v2/backups/{backup.id}/restores",
|
|
json={
|
|
"destination": str(destination),
|
|
"selection": [],
|
|
"overwrite_policy": "fail",
|
|
},
|
|
headers=headers,
|
|
)
|
|
assert restore_response.status_code == 202
|
|
source_root.rename(settings.data_dir / "removed-source")
|
|
restore_worker = Worker(settings, owner="encrypted-restore-worker")
|
|
try:
|
|
assert await restore_worker.run_once()
|
|
finally:
|
|
await restore_worker.engine.dispose()
|
|
|
|
restored = await client.get(
|
|
f"/api/v2/restores/{restore_response.json()['id']}", headers=headers
|
|
)
|
|
assert restored.json()["state"] == "committed"
|
|
assert (destination / "hello.txt").read_bytes() == plaintext
|
|
|
|
rotated_destination = settings.restore_roots[0] / "rotated-restored"
|
|
rotated_restore_response = await client.post(
|
|
f"/api/v2/backups/{second_backup.id}/restores",
|
|
json={
|
|
"destination": str(rotated_destination),
|
|
"selection": [],
|
|
"overwrite_policy": "fail",
|
|
},
|
|
headers=headers,
|
|
)
|
|
assert rotated_restore_response.status_code == 202
|
|
rotated_restore_worker = Worker(settings, owner="encrypted-rotated-restore-worker")
|
|
try:
|
|
assert await rotated_restore_worker.run_once()
|
|
finally:
|
|
await rotated_restore_worker.engine.dispose()
|
|
assert (rotated_destination / "hello.txt").read_bytes() == plaintext_after_rotation
|
|
|
|
manifest_path.write_bytes(stored_manifest[:-1] + bytes([stored_manifest[-1] ^ 1]))
|
|
corrupt_destination = settings.restore_roots[0] / "corrupt-manifest"
|
|
corrupt_restore = await client.post(
|
|
f"/api/v2/backups/{backup.id}/restores",
|
|
json={
|
|
"destination": str(corrupt_destination),
|
|
"selection": [],
|
|
"overwrite_policy": "fail",
|
|
},
|
|
headers=headers,
|
|
)
|
|
assert corrupt_restore.status_code == 202
|
|
corrupt_worker = Worker(settings, owner="encrypted-corrupt-manifest-worker")
|
|
try:
|
|
assert await corrupt_worker.run_once()
|
|
finally:
|
|
await corrupt_worker.engine.dispose()
|
|
|
|
corrupt_status = await client.get(
|
|
f"/api/v2/restores/{corrupt_restore.json()['id']}", headers=headers
|
|
)
|
|
assert corrupt_status.json()["state"] == "failed"
|
|
assert not corrupt_destination.exists()
|