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