313 lines
11 KiB
Python
313 lines
11 KiB
Python
from __future__ import annotations
|
|
|
|
import json
|
|
from pathlib import Path
|
|
import tempfile
|
|
import unittest
|
|
|
|
from sqlalchemy import Column, String, Table, create_engine, select
|
|
from sqlalchemy.orm import Session, sessionmaker
|
|
|
|
from govoplan_core.core.recovery import (
|
|
RecoveryCheckpoint,
|
|
RecoveryGuaranteeError,
|
|
RecoveryOperation,
|
|
RecoveryStatus,
|
|
)
|
|
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_mail.backend.db.models import (
|
|
MailMailboxFolderIndex,
|
|
MailMailboxMessageIndex,
|
|
MailServerProfile,
|
|
)
|
|
from govoplan_mail.backend.mailbox_index import cache_mailbox_folders
|
|
from govoplan_mail.backend.recovery import (
|
|
MailRecoveryError,
|
|
MailboxRefreshBusy,
|
|
begin_mailbox_refresh_recovery,
|
|
begin_provider_effect_recovery,
|
|
reconcile_outbox_provider_effect,
|
|
)
|
|
from govoplan_mail.backend.sending.imap import (
|
|
ImapFolderListResult,
|
|
ImapMailboxInfo,
|
|
)
|
|
|
|
|
|
class MailRecoveryTests(unittest.TestCase):
|
|
def setUp(self) -> None:
|
|
self.tempdir = tempfile.TemporaryDirectory()
|
|
self.addCleanup(self.tempdir.cleanup)
|
|
database_path = Path(self.tempdir.name) / "mail-recovery.sqlite3"
|
|
self.engine = create_engine(f"sqlite:///{database_path}")
|
|
access_users = Base.metadata.tables.get("access_users")
|
|
if access_users is None:
|
|
access_users = Table(
|
|
"access_users",
|
|
Base.metadata,
|
|
Column("id", String(36), primary_key=True),
|
|
)
|
|
Base.metadata.create_all(
|
|
self.engine,
|
|
tables=[
|
|
access_users,
|
|
DistributedLease.__table__,
|
|
RecoveryOperation.__table__,
|
|
RecoveryCheckpoint.__table__,
|
|
MailServerProfile.__table__,
|
|
MailMailboxFolderIndex.__table__,
|
|
MailMailboxMessageIndex.__table__,
|
|
],
|
|
)
|
|
configure_database(
|
|
f"sqlite:///{database_path}",
|
|
engine=self.engine,
|
|
dispose_previous=True,
|
|
)
|
|
bind_process_runtime_identity(
|
|
RuntimeIdentity(
|
|
installation_id="mail-recovery-tests",
|
|
node_id="mail-test-node",
|
|
incarnation="mail-test-incarnation",
|
|
role="worker",
|
|
software_version="test",
|
|
composition_hash="a" * 64,
|
|
)
|
|
)
|
|
self.SessionLocal = sessionmaker(
|
|
bind=self.engine,
|
|
class_=Session,
|
|
expire_on_commit=False,
|
|
)
|
|
with self.SessionLocal() as session:
|
|
session.add(
|
|
MailServerProfile(
|
|
id="profile-1",
|
|
tenant_id="tenant-1",
|
|
scope_type="tenant",
|
|
scope_id="tenant-1",
|
|
name="Recovery profile",
|
|
slug="recovery-profile",
|
|
smtp_config={},
|
|
)
|
|
)
|
|
session.commit()
|
|
self.addCleanup(self._cleanup_runtime)
|
|
|
|
def _cleanup_runtime(self) -> None:
|
|
bind_process_runtime_identity(None)
|
|
reset_database()
|
|
self.engine.dispose()
|
|
|
|
def test_smtp_effect_is_durable_before_provider_and_redacted_on_success(self) -> None:
|
|
recovery = begin_provider_effect_recovery(
|
|
kind="smtp-delivery",
|
|
effect_id="outbox:command-1:smtp-attempt:1",
|
|
tenant_id="tenant-1",
|
|
profile_id="profile-1",
|
|
message_bytes=b"Subject: Recovery\r\n\r\nBody",
|
|
expected_transport_revision="revision-1",
|
|
recipient_count=1,
|
|
resource_type="mail_delivery_command",
|
|
resource_id="command-1",
|
|
)
|
|
assert recovery is not None and recovery.operation is not None
|
|
with self.SessionLocal() as session:
|
|
operation = session.get(RecoveryOperation, recovery.operation_id)
|
|
assert operation is not None
|
|
self.assertEqual(RecoveryStatus.RUNNING.value, operation.status)
|
|
|
|
recovery.succeed_smtp(
|
|
accepted_count=1,
|
|
refused_recipients={},
|
|
)
|
|
|
|
with self.SessionLocal() as session:
|
|
operation = session.get(RecoveryOperation, recovery.operation_id)
|
|
assert operation is not None
|
|
self.assertEqual(RecoveryStatus.SUCCEEDED.value, operation.status)
|
|
evidence = json.dumps(
|
|
[
|
|
item.evidence
|
|
for item in session.scalars(
|
|
select(RecoveryCheckpoint).where(
|
|
RecoveryCheckpoint.operation_id == operation.id
|
|
)
|
|
)
|
|
]
|
|
)
|
|
self.assertNotIn("recipient@example.test", evidence)
|
|
self.assertNotIn("Subject: Recovery", evidence)
|
|
|
|
def test_unknown_smtp_outcome_blocks_replay_until_reconciled(self) -> None:
|
|
kwargs = {
|
|
"kind": "smtp-delivery",
|
|
"effect_id": "outbox:command-2:smtp-attempt:1",
|
|
"tenant_id": "tenant-1",
|
|
"profile_id": "profile-1",
|
|
"message_bytes": b"Subject: Unknown\r\n\r\nBody",
|
|
"expected_transport_revision": "revision-1",
|
|
"recipient_count": 1,
|
|
"resource_type": "mail_delivery_command",
|
|
"resource_id": "command-2",
|
|
}
|
|
recovery = begin_provider_effect_recovery(**kwargs)
|
|
assert recovery is not None
|
|
recovery.unknown(
|
|
code="socket_closed_after_data",
|
|
summary="SMTP outcome is unknown",
|
|
)
|
|
|
|
with self.assertRaises(MailRecoveryError):
|
|
begin_provider_effect_recovery(**kwargs)
|
|
|
|
self.assertTrue(
|
|
reconcile_outbox_provider_effect(
|
|
command_id="command-2",
|
|
effect_occurred=False,
|
|
evidence_reference="provider-case-42",
|
|
user_id="operator-1",
|
|
)
|
|
)
|
|
with self.SessionLocal() as session:
|
|
operation = session.get(RecoveryOperation, recovery.operation_id)
|
|
assert operation is not None
|
|
self.assertEqual(RecoveryStatus.RECOVERED.value, operation.status)
|
|
|
|
def test_mailbox_refresh_verifies_the_committed_index(self) -> None:
|
|
result = ImapFolderListResult(
|
|
host="imap.example.test",
|
|
port=993,
|
|
security="tls",
|
|
folders=[
|
|
ImapMailboxInfo(
|
|
name="INBOX",
|
|
flags=["\\HasNoChildren"],
|
|
message_count=4,
|
|
unseen_count=1,
|
|
)
|
|
],
|
|
)
|
|
recovery = begin_mailbox_refresh_recovery(
|
|
tenant_id="tenant-1",
|
|
profile_id="profile-1",
|
|
folder="*",
|
|
purpose="folders",
|
|
)
|
|
with self.SessionLocal() as session:
|
|
session.add(
|
|
MailMailboxFolderIndex(
|
|
tenant_id="tenant-1",
|
|
profile_id="profile-1",
|
|
folder="Removed",
|
|
flags=[],
|
|
)
|
|
)
|
|
session.commit()
|
|
cache_mailbox_folders(
|
|
session,
|
|
tenant_id="tenant-1",
|
|
profile_id="profile-1",
|
|
result=result,
|
|
)
|
|
session.commit()
|
|
recovery.complete_folders(result)
|
|
|
|
with self.SessionLocal() as session:
|
|
self.assertEqual(
|
|
["INBOX"],
|
|
list(
|
|
session.scalars(
|
|
select(MailMailboxFolderIndex.folder).order_by(
|
|
MailMailboxFolderIndex.folder
|
|
)
|
|
)
|
|
),
|
|
)
|
|
operation = session.get(
|
|
RecoveryOperation,
|
|
recovery.operation.operation_id,
|
|
)
|
|
assert operation is not None
|
|
self.assertEqual(RecoveryStatus.SUCCEEDED.value, operation.status)
|
|
|
|
def test_missing_runtime_identity_fails_before_a_provider_effect_can_start(self) -> None:
|
|
bind_process_runtime_identity(None)
|
|
with self.assertRaises(MailRecoveryError):
|
|
begin_provider_effect_recovery(
|
|
kind="imap-append",
|
|
effect_id="campaign-job:job-1:imap-attempt:1",
|
|
tenant_id="tenant-1",
|
|
profile_id="profile-1",
|
|
message_bytes=b"message",
|
|
expected_transport_revision="revision-1",
|
|
folder="Sent",
|
|
)
|
|
|
|
def test_tampered_provider_evidence_cannot_be_marked_successful(self) -> None:
|
|
recovery = begin_provider_effect_recovery(
|
|
kind="smtp-delivery",
|
|
effect_id="outbox:command-3:smtp-attempt:1",
|
|
tenant_id="tenant-1",
|
|
profile_id="profile-1",
|
|
message_bytes=b"message",
|
|
expected_transport_revision="revision-1",
|
|
recipient_count=1,
|
|
resource_type="mail_delivery_command",
|
|
resource_id="command-3",
|
|
)
|
|
assert recovery is not None
|
|
with self.SessionLocal() as session:
|
|
checkpoint = session.scalar(
|
|
select(RecoveryCheckpoint)
|
|
.where(RecoveryCheckpoint.operation_id == recovery.operation_id)
|
|
.order_by(RecoveryCheckpoint.sequence)
|
|
.limit(1)
|
|
)
|
|
assert checkpoint is not None
|
|
checkpoint.summary = "tampered"
|
|
session.commit()
|
|
|
|
with self.assertRaises(RecoveryGuaranteeError):
|
|
recovery.succeed_smtp(accepted_count=1, refused_recipients={})
|
|
with self.SessionLocal() as session:
|
|
operation = session.get(RecoveryOperation, recovery.operation_id)
|
|
assert operation is not None
|
|
self.assertNotEqual(RecoveryStatus.SUCCEEDED.value, operation.status)
|
|
|
|
def test_mailbox_refresh_has_a_cross_runtime_fence(self) -> None:
|
|
recovery = begin_mailbox_refresh_recovery(
|
|
tenant_id="tenant-1",
|
|
profile_id="profile-1",
|
|
folder="INBOX",
|
|
purpose="messages",
|
|
)
|
|
bind_process_runtime_identity(
|
|
RuntimeIdentity(
|
|
installation_id="mail-recovery-tests",
|
|
node_id="mail-test-node-2",
|
|
incarnation="mail-test-incarnation-2",
|
|
role="worker",
|
|
software_version="test",
|
|
composition_hash="a" * 64,
|
|
)
|
|
)
|
|
with self.assertRaises(MailboxRefreshBusy):
|
|
begin_mailbox_refresh_recovery(
|
|
tenant_id="tenant-1",
|
|
profile_id="profile-1",
|
|
folder="INBOX",
|
|
purpose="messages",
|
|
)
|
|
recovery.reject(summary="Test refresh stopped", code="test")
|
|
|
|
|
|
if __name__ == "__main__":
|
|
unittest.main()
|