182 lines
6.4 KiB
Python
182 lines
6.4 KiB
Python
from __future__ import annotations
|
|
|
|
from pathlib import Path
|
|
import sqlite3
|
|
import tempfile
|
|
import unittest
|
|
|
|
from sqlalchemy import create_engine
|
|
from sqlalchemy.orm import Session
|
|
|
|
from govoplan_core.core.recovery import (
|
|
RecoveryCheckpoint,
|
|
RecoveryOperation,
|
|
RecoveryStatus,
|
|
verify_recovery_evidence_chain,
|
|
)
|
|
from govoplan_core.core.runtime_coordination import (
|
|
DistributedLease,
|
|
RuntimeIdentity,
|
|
bind_process_runtime_identity,
|
|
)
|
|
from govoplan_records.backend.recovery import (
|
|
RecordRecoveryError,
|
|
begin_record_atomic_recovery,
|
|
)
|
|
|
|
|
|
class RecordsRecoveryTests(unittest.TestCase):
|
|
def setUp(self) -> None:
|
|
self.directory = tempfile.TemporaryDirectory(prefix="records-recovery-")
|
|
self.database_path = Path(self.directory.name) / "recovery.db"
|
|
self.engine = create_engine(f"sqlite:///{self.database_path}")
|
|
for table in (
|
|
DistributedLease.__table__,
|
|
RecoveryOperation.__table__,
|
|
RecoveryCheckpoint.__table__,
|
|
):
|
|
table.create(self.engine)
|
|
self.session = Session(self.engine)
|
|
bind_process_runtime_identity(
|
|
RuntimeIdentity(
|
|
installation_id="test-installation",
|
|
node_id="records-test-node",
|
|
incarnation="11111111-1111-4111-8111-111111111111",
|
|
role="api",
|
|
software_version="test",
|
|
composition_hash="a" * 64,
|
|
)
|
|
)
|
|
|
|
def tearDown(self) -> None:
|
|
bind_process_runtime_identity(None)
|
|
self.session.close()
|
|
self.engine.dispose()
|
|
self.directory.cleanup()
|
|
|
|
def test_atomic_evidence_chain_replays_only_the_same_request(self) -> None:
|
|
started = begin_record_atomic_recovery(
|
|
self.session,
|
|
tenant_id="tenant-1",
|
|
operation_type="record.close",
|
|
idempotency_key="close-1",
|
|
request={"record_id": "record-1", "expected_revision": 1},
|
|
resource_type="record",
|
|
resource_id="record-1",
|
|
)
|
|
started.commit_success(
|
|
self.session,
|
|
result={"record_id": "record-1", "revision": 2},
|
|
resource_id="record-1",
|
|
)
|
|
|
|
operation = self.session.get(RecoveryOperation, started.operation_id)
|
|
self.assertIsNotNone(operation)
|
|
self.assertEqual(RecoveryStatus.SUCCEEDED.value, operation.status)
|
|
self.assertTrue(
|
|
verify_recovery_evidence_chain(self.session, started.operation_id)
|
|
)
|
|
|
|
replay = begin_record_atomic_recovery(
|
|
self.session,
|
|
tenant_id="tenant-1",
|
|
operation_type="record.close",
|
|
idempotency_key="close-1",
|
|
request={"record_id": "record-1", "expected_revision": 1},
|
|
resource_type="record",
|
|
resource_id="record-1",
|
|
)
|
|
self.assertTrue(replay.replayed)
|
|
self.assertIsNone(replay.operation)
|
|
|
|
with self.assertRaisesRegex(RecordRecoveryError, "not started"):
|
|
begin_record_atomic_recovery(
|
|
self.session,
|
|
tenant_id="tenant-1",
|
|
operation_type="record.close",
|
|
idempotency_key="close-1",
|
|
request={"record_id": "record-1", "expected_revision": 9},
|
|
resource_type="record",
|
|
resource_id="record-1",
|
|
)
|
|
|
|
def test_unresolved_resource_blocks_another_runtime_effect(self) -> None:
|
|
started = begin_record_atomic_recovery(
|
|
self.session,
|
|
tenant_id="tenant-1",
|
|
operation_type="transfer.simulate",
|
|
idempotency_key="dispatch-1",
|
|
request={"record_id": "record-1", "package_id": "package-1"},
|
|
resource_type="record",
|
|
resource_id="record-1",
|
|
)
|
|
bind_process_runtime_identity(
|
|
RuntimeIdentity(
|
|
installation_id="test-installation",
|
|
node_id="other-records-test-node",
|
|
incarnation="22222222-2222-4222-8222-222222222222",
|
|
role="api",
|
|
software_version="test",
|
|
composition_hash="a" * 64,
|
|
)
|
|
)
|
|
with self.assertRaisesRegex(RecordRecoveryError, "Another runtime"):
|
|
begin_record_atomic_recovery(
|
|
self.session,
|
|
tenant_id="tenant-1",
|
|
operation_type="record.reopen",
|
|
idempotency_key="reopen-1",
|
|
request={"record_id": "record-1"},
|
|
resource_type="record",
|
|
resource_id="record-1",
|
|
)
|
|
bind_process_runtime_identity(
|
|
RuntimeIdentity(
|
|
installation_id="test-installation",
|
|
node_id="records-test-node",
|
|
incarnation="11111111-1111-4111-8111-111111111111",
|
|
role="api",
|
|
software_version="test",
|
|
composition_hash="a" * 64,
|
|
)
|
|
)
|
|
started.fail(summary="Test cleanup", error_type="TestInterruption")
|
|
|
|
def test_recovery_evidence_survives_database_restore(self) -> None:
|
|
started = begin_record_atomic_recovery(
|
|
self.session,
|
|
tenant_id="tenant-1",
|
|
operation_type="record.close",
|
|
idempotency_key="restore-close-1",
|
|
request={"record_id": "record-restore", "expected_revision": 1},
|
|
resource_type="record",
|
|
resource_id="record-restore",
|
|
)
|
|
started.commit_success(
|
|
self.session,
|
|
result={"record_id": "record-restore", "revision": 2},
|
|
resource_id="record-restore",
|
|
)
|
|
|
|
restored_path = Path(self.directory.name) / "restored.db"
|
|
with (
|
|
sqlite3.connect(self.database_path) as source,
|
|
sqlite3.connect(restored_path) as target,
|
|
):
|
|
source.backup(target)
|
|
restored_engine = create_engine(f"sqlite:///{restored_path}")
|
|
try:
|
|
with Session(restored_engine) as restored:
|
|
operation = restored.get(RecoveryOperation, started.operation_id)
|
|
self.assertIsNotNone(operation)
|
|
self.assertEqual(RecoveryStatus.SUCCEEDED.value, operation.status)
|
|
self.assertTrue(
|
|
verify_recovery_evidence_chain(restored, started.operation_id)
|
|
)
|
|
finally:
|
|
restored_engine.dispose()
|
|
|
|
|
|
if __name__ == "__main__":
|
|
unittest.main()
|