Files
govoplan-files/tests/test_storage_recovery.py

647 lines
24 KiB
Python

from __future__ import annotations
import hashlib
from datetime import UTC, datetime
from pathlib import Path
import tempfile
import unittest
from unittest.mock import patch
from sqlalchemy import create_engine
from sqlalchemy.orm import sessionmaker
from govoplan_access.backend.db.models import Account, Group, User
from govoplan_core.core.recovery import (
RecoveryCheckpoint,
RecoveryOperation,
RecoveryStatus,
)
from govoplan_core.core.change_sequence import ChangeSequenceEntry
from govoplan_core.core.runtime_coordination import (
DistributedLease,
RuntimeIdentity,
bind_process_runtime_identity,
)
from govoplan_core.db.base import Base
from govoplan_core.db.session import configure_database, reset_database
from govoplan_files.backend.db.models import (
FileAsset,
FileBlob,
FileFormEvidenceGrant,
FileIntegrityFinding,
FileIntegrityScan,
FileShare,
FileVersion,
)
from govoplan_files.backend.storage.backends import (
LocalFilesystemStorageBackend,
)
from govoplan_files.backend.storage.common import FileStorageError
from govoplan_files.backend.storage.files import _get_or_create_blob
from govoplan_files.backend.storage.integrity import cleanup_orphan_finding
from govoplan_files.backend.storage.recovery import begin_blob_write_recovery
from govoplan_files.backend.storage.connector_profiles import ConnectorProfile
from govoplan_files.backend.storage.connector_writes import write_connector_file
from govoplan_files.backend.storage.lifecycle import (
execute_asset_purge,
garbage_collect_unreferenced_blobs,
preview_asset_purge,
)
TENANT_ID = "tenant-1"
USER_ID = "user-1"
class _MissingS3Object(Exception):
response = {
"Error": {"Code": "NoSuchKey"},
"ResponseMetadata": {"HTTPStatusCode": 404},
}
class _WriteS3Client:
def __init__(
self, *, tamper_metadata: bool = False, fail_post_write_probe: bool = False
) -> None:
self.object: dict[str, object] | None = None
self.tamper_metadata = tamper_metadata
self.fail_post_write_probe = fail_post_write_probe
self.put_request: dict[str, object] | None = None
self.closed = False
self.write_completed = False
def head_object(self, **_kwargs: object) -> dict[str, object]:
if self.write_completed and self.fail_post_write_probe:
raise RuntimeError("provider unavailable")
if self.object is None:
raise _MissingS3Object()
return dict(self.object)
def put_object(self, **kwargs: object) -> dict[str, object]:
self.put_request = dict(kwargs)
metadata = dict(kwargs.get("Metadata") or {})
if self.tamper_metadata:
metadata["govoplan-sha256"] = "0" * 64
self.object = {
"ETag": '"etag-written"',
"VersionId": "version-written",
"Metadata": metadata,
}
self.write_completed = True
return {"ETag": '"etag-written"', "VersionId": "version-written"}
def close(self) -> None:
self.closed = True
class StorageRecoveryTests(unittest.TestCase):
def setUp(self) -> None:
self.temporary_directory = tempfile.TemporaryDirectory()
self.addCleanup(self.temporary_directory.cleanup)
root = Path(self.temporary_directory.name)
self.backend = LocalFilesystemStorageBackend(root / "objects")
database_path = root / "recovery.sqlite3"
self.engine = create_engine(f"sqlite:///{database_path}", future=True)
Base.metadata.create_all(
bind=self.engine,
tables=[
Account.__table__,
User.__table__,
Group.__table__,
ChangeSequenceEntry.__table__,
DistributedLease.__table__,
RecoveryOperation.__table__,
RecoveryCheckpoint.__table__,
FileBlob.__table__,
FileAsset.__table__,
FileVersion.__table__,
FileFormEvidenceGrant.__table__,
FileShare.__table__,
FileIntegrityScan.__table__,
FileIntegrityFinding.__table__,
],
)
configure_database(
f"sqlite:///{database_path}",
engine=self.engine,
dispose_previous=True,
)
bind_process_runtime_identity(
RuntimeIdentity(
installation_id="files-recovery-test",
node_id="node-1",
incarnation="incarnation-1",
role="api",
software_version="test",
composition_hash="a" * 64,
)
)
self.Session = sessionmaker(
bind=self.engine,
expire_on_commit=False,
future=True,
)
self.enterContext(
patch(
"govoplan_files.backend.storage.files._storage_backend_name",
return_value=self.backend.name,
)
)
self.enterContext(
patch(
"govoplan_files.backend.storage.files._storage_bucket_name",
return_value="",
)
)
self.session = self.Session()
self.addCleanup(self._cleanup_runtime)
def _cleanup_runtime(self) -> None:
self.session.close()
bind_process_runtime_identity(None)
reset_database()
self.engine.dispose()
def test_committed_blob_write_is_independently_verified(self) -> None:
observed_running_operation: list[bool] = []
put_bytes = self.backend.put_bytes
class ObservingBackend:
name = self.backend.name
def __getattr__(backend_self, name):
return getattr(self.backend, name)
def put_bytes(backend_self, key, data, *, content_type=None):
with self.Session() as evidence_session:
operation = evidence_session.query(RecoveryOperation).one_or_none()
observed_running_operation.append(
operation is not None
and operation.status == RecoveryStatus.RUNNING.value
and len(operation.request_sha256) == 64
)
put_bytes(key, data, content_type=content_type)
observing_backend = ObservingBackend()
with patch(
"govoplan_files.backend.storage.files.get_storage_backend",
return_value=observing_backend,
):
blob = _get_or_create_blob(
self.session,
tenant_id=TENANT_ID,
data=b"durable",
filename="private-name.txt",
content_type="text/plain",
actor_id=USER_ID,
)
self.session.commit()
operation = self._only_operation()
self.assertEqual(RecoveryStatus.SUCCEEDED.value, operation.status)
# SQLite uses an explicit caller-transaction mode because it permits
# only one writer. The intent becomes visible with the business commit;
# PostgreSQL retains the independent pre-effect commit guarantee.
self.assertEqual([False], observed_running_operation)
self.assertEqual(
"sqlite_caller_transaction",
operation.metadata_["durability_mode"],
)
self.assertTrue(self.backend.exists(blob.storage_key))
self.assertNotIn("private-name", blob.storage_key)
self.assertEqual(".blob", Path(blob.storage_key).suffix)
def test_rolled_back_blob_write_is_compensated(self) -> None:
with patch(
"govoplan_files.backend.storage.files.get_storage_backend",
return_value=self.backend,
):
blob = _get_or_create_blob(
self.session,
tenant_id=TENANT_ID,
data=b"rollback",
filename="rollback.txt",
content_type="text/plain",
actor_id=USER_ID,
)
storage_key = blob.storage_key
self.assertTrue(self.backend.exists(storage_key))
self.session.rollback()
operation = self._only_operation()
self.assertEqual(RecoveryStatus.RECOVERED.value, operation.status)
self.assertEqual(
"sqlite_post_rollback_reconstruction",
operation.metadata_["durability_mode"],
)
self.assertFalse(self.backend.exists(storage_key))
with self.Session() as evidence_session:
self.assertIsNone(evidence_session.get(FileBlob, blob.id))
def test_multiple_blob_writes_share_one_sqlite_business_transaction(self) -> None:
with patch(
"govoplan_files.backend.storage.files.get_storage_backend",
return_value=self.backend,
):
first = _get_or_create_blob(
self.session,
tenant_id=TENANT_ID,
data=b"first archive member",
filename="first.txt",
content_type="text/plain",
actor_id=USER_ID,
)
second = _get_or_create_blob(
self.session,
tenant_id=TENANT_ID,
data=b"second archive member",
filename="second.txt",
content_type="text/plain",
actor_id=USER_ID,
)
self.session.commit()
self.assertTrue(self.backend.exists(first.storage_key))
self.assertTrue(self.backend.exists(second.storage_key))
with self.Session() as evidence_session:
operations = evidence_session.query(RecoveryOperation).all()
self.assertEqual(2, len(operations))
self.assertEqual(
{RecoveryStatus.SUCCEEDED.value},
{operation.status for operation in operations},
)
self.assertEqual(
{"sqlite_caller_transaction"},
{
operation.metadata_["durability_mode"]
for operation in operations
},
)
def test_post_write_tamper_is_quarantined_and_recovery_required(self) -> None:
with patch(
"govoplan_files.backend.storage.files.get_storage_backend",
return_value=self.backend,
):
blob = _get_or_create_blob(
self.session,
tenant_id=TENANT_ID,
data=b"expected",
filename="evidence.bin",
content_type="application/octet-stream",
actor_id=USER_ID,
)
self.backend.put_bytes(blob.storage_key, b"tampered")
self.session.commit()
operation = self._only_operation()
self.assertEqual(
RecoveryStatus.RECOVERY_REQUIRED.value,
operation.status,
)
with self.Session() as evidence_session:
persisted = evidence_session.get(FileBlob, blob.id)
self.assertIsNotNone(persisted)
self.assertEqual("checksum_mismatch", persisted.integrity_status)
self.assertIsNotNone(persisted.quarantined_at)
def test_missing_optional_encryption_fails_before_object_effect(self) -> None:
with patch(
"govoplan_files.backend.storage.files.get_storage_backend",
return_value=self.backend,
), patch(
"govoplan_files.backend.storage.content_protection.encryption_content_cipher",
return_value=None,
), self.assertRaisesRegex(FileStorageError, "Encryption module"):
_get_or_create_blob(
self.session,
tenant_id=TENANT_ID,
data=b"protected",
filename="protected.bin",
content_type="application/octet-stream",
actor_id=USER_ID,
encryption_vault_id="vault-1",
)
self.session.rollback()
operation = self._only_operation()
self.assertEqual(RecoveryStatus.REJECTED.value, operation.status)
objects = self.backend.list_objects(
prefix=f"tenants/{TENANT_ID}/files/",
limit=10,
)
self.assertEqual((), objects.objects)
def test_orphan_cleanup_forward_completes_after_business_rollback(self) -> None:
key = f"tenants/{TENANT_ID}/files/orphan.bin"
self.backend.put_bytes(key, b"orphan")
scan = FileIntegrityScan(
id="scan-1",
tenant_id=TENANT_ID,
storage_backend=self.backend.name,
storage_prefix=f"tenants/{TENANT_ID}/files/",
status="completed",
)
finding = FileIntegrityFinding(
id="finding-1",
scan_id=scan.id,
tenant_id=TENANT_ID,
kind="orphan_object",
state="open",
storage_key=key,
observed_size_bytes=6,
observed_checksum_sha256=hashlib.sha256(b"orphan").hexdigest(),
)
self.session.add_all([scan, finding])
self.session.commit()
cleanup_orphan_finding(
self.session,
finding,
user_id=USER_ID,
dry_run=False,
backend=self.backend,
)
self.assertFalse(self.backend.exists(key))
self.session.rollback()
operation = self._only_operation()
self.assertEqual(RecoveryStatus.SUCCEEDED.value, operation.status)
with self.Session() as evidence_session:
persisted = evidence_session.get(FileIntegrityFinding, finding.id)
self.assertIsNotNone(persisted)
self.assertEqual("deleted", persisted.state)
def test_missing_runtime_identity_blocks_before_object_write(self) -> None:
bind_process_runtime_identity(None)
with patch(
"govoplan_files.backend.storage.files.get_storage_backend",
return_value=self.backend,
), self.assertRaisesRegex(FileStorageError, "recovery ledger"):
_get_or_create_blob(
self.session,
tenant_id=TENANT_ID,
data=b"blocked",
filename="blocked.bin",
content_type="application/octet-stream",
actor_id=USER_ID,
)
self.session.rollback()
objects = self.backend.list_objects(
prefix=f"tenants/{TENANT_ID}/files/",
limit=10,
)
self.assertEqual((), objects.objects)
def test_blob_fence_blocks_a_competing_runtime_before_effect(self) -> None:
checksum = hashlib.sha256(b"fenced").hexdigest()
begin_blob_write_recovery(
self.session,
backend=self.backend,
tenant_id=TENANT_ID,
blob_id="blob-fenced",
storage_key=f"tenants/{TENANT_ID}/files/fenced.blob",
semantic_checksum_sha256=checksum,
semantic_size_bytes=6,
protection_discriminator="plaintext",
created_new=True,
)
with self.Session() as competing_session, self.assertRaisesRegex(
FileStorageError,
"already owned",
):
begin_blob_write_recovery(
competing_session,
backend=self.backend,
tenant_id=TENANT_ID,
blob_id="blob-fenced",
storage_key=f"tenants/{TENANT_ID}/files/fenced.blob",
semantic_checksum_sha256=checksum,
semantic_size_bytes=6,
protection_discriminator="plaintext",
created_new=True,
)
self.session.rollback()
self.assertEqual(RecoveryStatus.REJECTED.value, self._only_operation().status)
def test_purge_is_idempotent_and_gc_deletes_only_after_reference_check(self) -> None:
key = f"tenants/{TENANT_ID}/files/purge-me.blob"
self.backend.put_bytes(key, b"purge-me")
blob = FileBlob(
id="blob-purge",
tenant_id=TENANT_ID,
storage_backend=self.backend.name,
storage_key=key,
checksum_sha256=hashlib.sha256(b"purge-me").hexdigest(),
size_bytes=8,
ref_count=1,
)
asset = FileAsset(
id="asset-purge",
tenant_id=TENANT_ID,
owner_type="user",
owner_user_id=USER_ID,
current_version_id="version-purge",
display_path="purge-me.txt",
filename="purge-me.txt",
deleted_at=datetime.now(UTC),
)
version = FileVersion(
id="version-purge",
tenant_id=TENANT_ID,
file_asset_id=asset.id,
blob_id=blob.id,
version_number=1,
filename_at_upload=asset.filename,
display_path_at_upload=asset.display_path,
size_bytes=blob.size_bytes,
checksum_sha256=blob.checksum_sha256,
)
self.session.add_all([blob, asset, version])
self.session.commit()
preview = preview_asset_purge(
self.session, tenant_id=TENANT_ID, file_ids=[asset.id]
)
result = execute_asset_purge(
self.session,
tenant_id=TENANT_ID,
file_ids=[asset.id],
preview_sha256=preview.preview_sha256,
idempotency_key="purge-request-1",
approval_reference="approval-1",
)
self.assertEqual(1, result.purged_files)
self.assertEqual(1, result.released_blobs)
self.assertIsNone(self.session.get(FileAsset, asset.id))
retained_blob = self.session.get(FileBlob, blob.id)
self.assertIsNotNone(retained_blob)
self.assertEqual(0, retained_blob.ref_count)
self.assertTrue(self.backend.exists(key))
replay = execute_asset_purge(
self.session,
tenant_id=TENANT_ID,
file_ids=[asset.id],
preview_sha256=preview.preview_sha256,
idempotency_key="purge-request-1",
approval_reference="approval-1",
)
self.assertTrue(replay.replayed)
with patch(
"govoplan_files.backend.storage.lifecycle.get_storage_backend",
return_value=self.backend,
):
gc_result = garbage_collect_unreferenced_blobs(
self.session,
tenant_id=TENANT_ID,
limit=10,
approval_reference="approval-gc-1",
)
self.assertEqual(1, gc_result.deleted_blobs)
self.assertFalse(self.backend.exists(key))
self.assertIsNone(self.session.get(FileBlob, blob.id))
with self.Session() as evidence_session:
operations = evidence_session.query(RecoveryOperation).all()
self.assertEqual(2, len(operations))
self.assertEqual(
{RecoveryStatus.SUCCEEDED.value},
{operation.status for operation in operations},
)
def test_s3_write_is_conditional_verified_and_idempotent(self) -> None:
client = _WriteS3Client()
profile = ConnectorProfile(
id="s3-write",
label="S3 write",
provider="s3",
base_path="root",
capabilities=("browse", "write"),
metadata={"bucket": "files"},
)
with patch(
"govoplan_files.backend.storage.connector_writes._s3_client",
return_value=client,
):
result = write_connector_file(
profile,
tenant_id=TENANT_ID,
library_id=None,
remote_path="out/report.txt",
data=b"report",
content_type="text/plain",
expected_revision=None,
idempotency_key="connector-write-1",
)
replay = write_connector_file(
profile,
tenant_id=TENANT_ID,
library_id=None,
remote_path="out/report.txt",
data=b"report",
content_type="text/plain",
expected_revision=None,
idempotency_key="connector-write-1",
)
self.assertEqual(RecoveryStatus.SUCCEEDED.value, result.status)
self.assertTrue(replay.replayed)
self.assertEqual("*", client.put_request["IfNoneMatch"])
self.assertEqual("root/out/report.txt", client.put_request["Key"])
self.assertNotIn(b"report", repr(self._only_operation().metadata_).encode())
self.assertTrue(client.closed)
def test_s3_write_tamper_is_visible_and_blocks_another_writer(self) -> None:
client = _WriteS3Client(tamper_metadata=True)
profile = ConnectorProfile(
id="s3-write",
label="S3 write",
provider="s3",
capabilities=("write",),
metadata={"bucket": "files"},
)
with patch(
"govoplan_files.backend.storage.connector_writes._s3_client",
return_value=client,
):
result = write_connector_file(
profile,
tenant_id=TENANT_ID,
library_id=None,
remote_path="report.txt",
data=b"report",
content_type="text/plain",
expected_revision=None,
idempotency_key="connector-write-tampered",
)
with self.assertRaisesRegex(FileStorageError, "unresolved recovery"):
write_connector_file(
profile,
tenant_id=TENANT_ID,
library_id=None,
remote_path="report.txt",
data=b"replacement",
content_type="text/plain",
expected_revision="version-written",
idempotency_key="connector-write-after-tamper",
)
self.assertEqual(RecoveryStatus.RECOVERY_REQUIRED.value, result.status)
self.assertEqual(
RecoveryStatus.RECOVERY_REQUIRED.value, self._only_operation().status
)
def test_s3_write_probe_failure_is_outcome_unknown_not_confirmed_absent(self) -> None:
client = _WriteS3Client(fail_post_write_probe=True)
profile = ConnectorProfile(
id="s3-write",
label="S3 write",
provider="s3",
capabilities=("write",),
metadata={"bucket": "files"},
)
with patch(
"govoplan_files.backend.storage.connector_writes._s3_client",
return_value=client,
):
result = write_connector_file(
profile,
tenant_id=TENANT_ID,
library_id=None,
remote_path="report.txt",
data=b"report",
content_type="text/plain",
expected_revision=None,
idempotency_key="connector-write-probe-failed",
)
self.assertEqual(RecoveryStatus.OUTCOME_UNKNOWN.value, result.status)
operation = self._only_operation()
self.assertEqual(RecoveryStatus.OUTCOME_UNKNOWN.value, operation.status)
with self.Session() as session:
checkpoint = (
session.query(RecoveryCheckpoint)
.filter(RecoveryCheckpoint.operation_id == operation.id)
.order_by(RecoveryCheckpoint.sequence.desc())
.first()
)
self.assertIsNotNone(checkpoint)
self.assertFalse(checkpoint.evidence["observed"]["probe_verified"])
self.assertFalse(checkpoint.evidence["observed"]["present"])
def _only_operation(self) -> RecoveryOperation:
with self.Session() as session:
operations = session.query(RecoveryOperation).all()
self.assertEqual(1, len(operations))
session.expunge(operations[0])
return operations[0]
if __name__ == "__main__":
unittest.main()