Release govoplan-mail v0.1.27: stabilize credentials, folder encoding and transport progress
Module Package Release / publish-packages (push) Successful in 11s
Module Package Release / publish-packages (push) Successful in 11s
This commit is contained in:
@@ -5,12 +5,14 @@ from contextvars import ContextVar
|
||||
from dataclasses import dataclass
|
||||
from email.message import EmailMessage
|
||||
from email.utils import formatdate, make_msgid
|
||||
from threading import get_ident
|
||||
from typing import Any, Iterator
|
||||
|
||||
from sqlalchemy.orm import Session
|
||||
|
||||
from govoplan_core.core.mail import NotificationMailDeliveryRequest
|
||||
from govoplan_core.core.modules import ModuleContext
|
||||
from govoplan_mail.backend.config import ImapConfig
|
||||
from govoplan_mail.backend.mail_profiles import (
|
||||
MailProfileError,
|
||||
_assert_campaign_inherits_profile_credentials,
|
||||
@@ -38,6 +40,7 @@ from govoplan_mail.backend.recovery import (
|
||||
)
|
||||
from govoplan_mail.backend.sending.imap import (
|
||||
ImapAppendError,
|
||||
ImapBatchSession,
|
||||
ImapConfigurationError,
|
||||
append_message_to_sent,
|
||||
)
|
||||
@@ -72,6 +75,78 @@ class CampaignSmtpDeliveryResult:
|
||||
@dataclass(frozen=True, slots=True)
|
||||
class CampaignImapAppendResult:
|
||||
folder: str
|
||||
connection_sequence: int = 1
|
||||
session_reused: bool = False
|
||||
reconnect_count: int = 0
|
||||
|
||||
|
||||
class CampaignImapBatchState:
|
||||
"""Lazy transport reuse; authorization and recovery remain per message."""
|
||||
|
||||
def __init__(self, *, tenant_id: str, campaign_id: str):
|
||||
self.tenant_id = tenant_id
|
||||
self.campaign_id = campaign_id
|
||||
self._session: ImapBatchSession | None = None
|
||||
self._binding: tuple[Any, ...] | None = None
|
||||
self._previous_connections = 0
|
||||
self._previous_reconnects = 0
|
||||
self._closed = False
|
||||
self._owner_thread = get_ident()
|
||||
|
||||
@property
|
||||
def connection_count(self) -> int:
|
||||
return self._previous_connections + (self._session.connection_count if self._session else 0)
|
||||
|
||||
@property
|
||||
def reconnect_count(self) -> int:
|
||||
return self._previous_reconnects + (self._session.reconnect_count if self._session else 0)
|
||||
|
||||
def assert_scope(self, *, tenant_id: str, campaign_id: str) -> None:
|
||||
if (
|
||||
self._closed or get_ident() != self._owner_thread
|
||||
or tenant_id != self.tenant_id or campaign_id != self.campaign_id
|
||||
):
|
||||
raise ImapConfigurationError("The IMAP batch does not match this campaign scope")
|
||||
|
||||
def session_for(self, config: ImapConfig, *, binding: tuple[Any, ...]) -> ImapBatchSession:
|
||||
if self._closed:
|
||||
raise ImapConfigurationError("The IMAP batch is closed")
|
||||
if self._session is not None and (
|
||||
self._binding != binding or not self._session.matches_config(config)
|
||||
):
|
||||
self._release_session()
|
||||
if self._session is None:
|
||||
self._session = ImapBatchSession(config)
|
||||
self._binding = binding
|
||||
return self._session
|
||||
|
||||
def _release_session(self) -> None:
|
||||
if self._session is not None:
|
||||
self._previous_connections += self._session.connection_count
|
||||
self._previous_reconnects += self._session.reconnect_count
|
||||
self._session.close()
|
||||
self._session = None
|
||||
|
||||
def close(self) -> None:
|
||||
self._closed = True
|
||||
self._release_session()
|
||||
|
||||
|
||||
_ACTIVE_IMAP_BATCH: ContextVar[CampaignImapBatchState | None] = ContextVar(
|
||||
"govoplan_mail_active_imap_batch", default=None,
|
||||
)
|
||||
|
||||
|
||||
@contextmanager
|
||||
def campaign_imap_batch(*, tenant_id: str, campaign_id: str) -> Iterator[CampaignImapBatchState]:
|
||||
"""Open no connection until an individual append passes its current checks."""
|
||||
state = CampaignImapBatchState(tenant_id=tenant_id, campaign_id=campaign_id)
|
||||
token = _ACTIVE_IMAP_BATCH.set(state)
|
||||
try:
|
||||
yield state
|
||||
finally:
|
||||
_ACTIVE_IMAP_BATCH.reset(token)
|
||||
state.close()
|
||||
|
||||
|
||||
@dataclass(frozen=True, slots=True)
|
||||
@@ -155,6 +230,7 @@ def _authorized_campaign_profile(
|
||||
campaign_id: str,
|
||||
profile_id: str,
|
||||
selection: dict[str, str | None] | None = None,
|
||||
credential_protocol: str | None = None,
|
||||
):
|
||||
profile = ensure_mail_profile_allowed_for_campaign(
|
||||
session,
|
||||
@@ -164,7 +240,7 @@ def _authorized_campaign_profile(
|
||||
require_active=True,
|
||||
)
|
||||
policy = effective_mail_profile_policy(session, tenant_id=tenant_id, campaign_id=campaign_id)
|
||||
_assert_campaign_inherits_profile_credentials(profile, policy, selection)
|
||||
_assert_campaign_inherits_profile_credentials(profile, policy, selection, protocol=credential_protocol)
|
||||
return profile
|
||||
|
||||
|
||||
@@ -343,6 +419,7 @@ def campaign_smtp_batch(
|
||||
campaign_id=campaign_id,
|
||||
profile_id=profile_id,
|
||||
selection=selection,
|
||||
credential_protocol="smtp",
|
||||
)
|
||||
except MailProfileError:
|
||||
raise
|
||||
@@ -444,6 +521,7 @@ def send_campaign_email_bytes(
|
||||
campaign_id=campaign_id,
|
||||
profile_id=profile_id,
|
||||
selection=selection,
|
||||
credential_protocol="smtp",
|
||||
)
|
||||
except MailProfileError:
|
||||
raise
|
||||
@@ -593,6 +671,9 @@ def append_campaign_message_to_sent(
|
||||
recovery_resource_type: str | None = None,
|
||||
recovery_resource_id: str | None = None,
|
||||
) -> CampaignImapAppendResult:
|
||||
batch = _ACTIVE_IMAP_BATCH.get()
|
||||
if batch is not None:
|
||||
batch.assert_scope(tenant_id=tenant_id, campaign_id=campaign_id)
|
||||
selection = _selection_payload(
|
||||
profile_id=profile_id,
|
||||
smtp_server_id=smtp_server_id,
|
||||
@@ -607,6 +688,7 @@ def append_campaign_message_to_sent(
|
||||
campaign_id=campaign_id,
|
||||
profile_id=profile_id,
|
||||
selection=selection,
|
||||
credential_protocol="imap",
|
||||
)
|
||||
except MailProfileError:
|
||||
raise
|
||||
@@ -681,6 +763,15 @@ def append_campaign_message_to_sent(
|
||||
)
|
||||
except MailProfileError:
|
||||
raise MailProfileError("Appending to Sent is blocked by the effective Mail policy.") from None
|
||||
batch_session = None
|
||||
if batch is not None:
|
||||
batch_session = batch.session_for(
|
||||
imap,
|
||||
binding=(
|
||||
tenant_id, campaign_id, profile_id, smtp_server_id, smtp_credential_id,
|
||||
imap_server_id, imap_credential_id, smtp_revision, imap_revision, folder,
|
||||
),
|
||||
)
|
||||
try:
|
||||
recovery = begin_provider_effect_recovery(
|
||||
kind="imap-append",
|
||||
@@ -701,7 +792,12 @@ def append_campaign_message_to_sent(
|
||||
outcome_unknown=True,
|
||||
)
|
||||
try:
|
||||
result = append_message_to_sent(message_bytes, imap_config=imap, folder=folder)
|
||||
if batch_session is None:
|
||||
result = append_message_to_sent(message_bytes, imap_config=imap, folder=folder)
|
||||
else:
|
||||
result = append_message_to_sent(
|
||||
message_bytes, imap_config=imap, folder=folder, batch_session=batch_session,
|
||||
)
|
||||
except ImapAppendError as exc:
|
||||
sanitized = _sanitized_imap_error(exc)
|
||||
if recovery is not None:
|
||||
@@ -728,11 +824,18 @@ def append_campaign_message_to_sent(
|
||||
try:
|
||||
recovery.succeed_imap(folder=result.folder)
|
||||
except Exception:
|
||||
if batch is not None:
|
||||
batch.close()
|
||||
raise ImapAppendError(
|
||||
"IMAP APPEND returned success, but durable recovery evidence could not be finalized.",
|
||||
outcome_unknown=True,
|
||||
) from None
|
||||
return CampaignImapAppendResult(folder=result.folder)
|
||||
return CampaignImapAppendResult(
|
||||
folder=result.folder,
|
||||
connection_sequence=batch.connection_count if batch else getattr(result, "connection_sequence", 1),
|
||||
session_reused=getattr(result, "session_reused", False),
|
||||
reconnect_count=batch.reconnect_count if batch else getattr(result, "reconnect_count", 0),
|
||||
)
|
||||
|
||||
|
||||
class MailCampaignCapability:
|
||||
@@ -745,6 +848,7 @@ class MailCampaignCapability:
|
||||
mail_profile_id_from_campaign_json = staticmethod(mail_profile_id_from_campaign_json)
|
||||
campaign_profile_delivery_summary = staticmethod(campaign_profile_delivery_summary)
|
||||
campaign_smtp_batch = staticmethod(campaign_smtp_batch)
|
||||
campaign_imap_batch = staticmethod(campaign_imap_batch)
|
||||
send_campaign_email_bytes = staticmethod(send_campaign_email_bytes)
|
||||
append_campaign_message_to_sent = staticmethod(append_campaign_message_to_sent)
|
||||
wait_for_rate_limit = staticmethod(wait_for_rate_limit)
|
||||
|
||||
Reference in New Issue
Block a user