666 lines
23 KiB
Python
666 lines
23 KiB
Python
from __future__ import annotations
|
|
|
|
import logging
|
|
|
|
from fastapi import APIRouter, Depends, HTTPException, status
|
|
from sqlalchemy.orm import Session
|
|
|
|
from govoplan_campaign.backend.schemas import (
|
|
AppendSentRequest,
|
|
CampaignActionResponse,
|
|
CampaignRetryJobsRequest,
|
|
CampaignSendJobRequest,
|
|
CampaignSendUnattemptedRequest,
|
|
CampaignResolveOutcomeRequest,
|
|
CampaignDeliveryOptionsResponse,
|
|
MockCampaignSendRequest,
|
|
MockCampaignSendResponse,
|
|
QueueCampaignRequest,
|
|
QueueCampaignResponse,
|
|
SendCampaignNowRequest,
|
|
SendCampaignNowResponse,
|
|
)
|
|
from govoplan_core.auth import ApiPrincipal, require_any_scope, require_scope
|
|
from govoplan_core.audit.logging import audit_from_principal
|
|
from govoplan_campaign.backend.db.models import (
|
|
CampaignJob,
|
|
JobImapStatus,
|
|
JobQueueStatus,
|
|
JobSendStatus,
|
|
)
|
|
from govoplan_campaign.backend.integrations import (
|
|
postbox_integration,
|
|
)
|
|
from govoplan_core.db.session import get_session
|
|
from govoplan_campaign.backend.response_security import (
|
|
public_send_campaign_now_result,
|
|
send_campaign_now_audit_details,
|
|
)
|
|
from govoplan_campaign.backend.persistence.campaigns import (
|
|
CampaignPersistenceError,
|
|
)
|
|
from govoplan_campaign.backend.persistence.versions import (
|
|
is_user_locked_version,
|
|
)
|
|
|
|
from govoplan_campaign.backend.dev.mock_campaign import (
|
|
MockCampaignSendError,
|
|
run_mock_campaign_send,
|
|
)
|
|
from govoplan_campaign.backend.sending.execution import ExecutionSnapshotError
|
|
from govoplan_campaign.backend.sending.jobs import (
|
|
QueueingError,
|
|
SynchronousSendRejected,
|
|
cancel_campaign_jobs,
|
|
enqueue_pending_imap_appends,
|
|
pause_campaign_jobs,
|
|
queue_campaign_jobs,
|
|
queue_failed_jobs_for_retry,
|
|
queue_unattempted_jobs,
|
|
reconcile_job_outcome,
|
|
resume_campaign_jobs,
|
|
send_campaign_now,
|
|
send_single_campaign_job,
|
|
synchronous_send_options,
|
|
)
|
|
|
|
from govoplan_campaign.backend.route_support import (
|
|
_get_campaign_for_principal,
|
|
_get_campaign_for_tenant,
|
|
_get_version_for_tenant,
|
|
_require_campaign_profile_use_if_needed,
|
|
_require_campaign_versions_profile_use,
|
|
_require_mail_profile_use_if_needed,
|
|
_require_permission,
|
|
)
|
|
|
|
router = APIRouter(prefix="/campaigns", tags=["campaigns"])
|
|
logger = logging.getLogger(__name__)
|
|
|
|
|
|
@router.get(
|
|
"/{campaign_id}/delivery-options", response_model=CampaignDeliveryOptionsResponse
|
|
)
|
|
def campaign_delivery_options(
|
|
campaign_id: str,
|
|
version_id: str | None = None,
|
|
session: Session = Depends(get_session),
|
|
principal: ApiPrincipal = Depends(
|
|
require_any_scope("campaigns:campaign:send", "campaigns:campaign:queue")
|
|
),
|
|
):
|
|
_get_campaign_for_principal(session, campaign_id, principal, write=True)
|
|
_require_permission(principal, "campaigns:recipient:read")
|
|
try:
|
|
return CampaignDeliveryOptionsResponse(
|
|
**synchronous_send_options(
|
|
session,
|
|
tenant_id=principal.tenant_id,
|
|
campaign_id=campaign_id,
|
|
version_id=version_id,
|
|
),
|
|
postbox_available=postbox_integration().available,
|
|
)
|
|
except QueueingError as exc:
|
|
raise HTTPException(
|
|
status_code=status.HTTP_422_UNPROCESSABLE_CONTENT, detail=str(exc)
|
|
) from exc
|
|
|
|
|
|
@router.post("/{campaign_id}/queue", response_model=QueueCampaignResponse)
|
|
def queue_campaign(
|
|
campaign_id: str,
|
|
payload: QueueCampaignRequest | None = None,
|
|
session: Session = Depends(get_session),
|
|
principal: ApiPrincipal = Depends(require_scope("campaigns:campaign:queue")),
|
|
):
|
|
_get_campaign_for_principal(session, campaign_id, principal, write=True)
|
|
_require_permission(principal, "campaigns:recipient:read")
|
|
payload = payload or QueueCampaignRequest()
|
|
_require_campaign_profile_use_if_needed(
|
|
session, principal, campaign_id, payload.version_id
|
|
)
|
|
try:
|
|
result = queue_campaign_jobs(
|
|
session,
|
|
tenant_id=principal.tenant_id,
|
|
campaign_id=campaign_id,
|
|
version_id=payload.version_id,
|
|
include_warnings=payload.include_warnings,
|
|
enqueue_celery=payload.enqueue_celery,
|
|
dry_run=payload.dry_run,
|
|
)
|
|
audit_from_principal(
|
|
session,
|
|
principal,
|
|
action="campaign.queued"
|
|
if not payload.dry_run
|
|
else "campaign.queue_dry_run",
|
|
object_type="campaign",
|
|
object_id=campaign_id,
|
|
details=result.as_dict(),
|
|
commit=True,
|
|
)
|
|
return QueueCampaignResponse(**result.as_dict())
|
|
except QueueingError as exc:
|
|
raise HTTPException(
|
|
status_code=status.HTTP_422_UNPROCESSABLE_CONTENT, detail=str(exc)
|
|
) from exc
|
|
|
|
|
|
@router.post("/{campaign_id}/jobs/retry", response_model=CampaignActionResponse)
|
|
def retry_campaign_jobs(
|
|
campaign_id: str,
|
|
payload: CampaignRetryJobsRequest | None = None,
|
|
session: Session = Depends(get_session),
|
|
principal: ApiPrincipal = Depends(require_scope("campaigns:campaign:retry")),
|
|
):
|
|
_get_campaign_for_principal(session, campaign_id, principal, write=True)
|
|
_require_permission(principal, "campaigns:recipient:read")
|
|
payload = payload or CampaignRetryJobsRequest()
|
|
_require_campaign_profile_use_if_needed(
|
|
session, principal, campaign_id, payload.version_id
|
|
)
|
|
try:
|
|
result = queue_failed_jobs_for_retry(
|
|
session,
|
|
tenant_id=principal.tenant_id,
|
|
campaign_id=campaign_id,
|
|
version_id=payload.version_id,
|
|
job_ids=payload.job_ids or None,
|
|
include_permanent=payload.include_permanent,
|
|
force_max_attempts=payload.force_max_attempts,
|
|
enqueue_celery=payload.enqueue_celery,
|
|
dry_run=payload.dry_run,
|
|
)
|
|
audit_from_principal(
|
|
session,
|
|
principal,
|
|
action="campaign.jobs_retry_queued"
|
|
if not payload.dry_run
|
|
else "campaign.jobs_retry_dry_run",
|
|
object_type="campaign",
|
|
object_id=campaign_id,
|
|
details=result,
|
|
commit=True,
|
|
)
|
|
return CampaignActionResponse(result=result)
|
|
except (QueueingError, ExecutionSnapshotError) as exc:
|
|
raise HTTPException(
|
|
status_code=status.HTTP_422_UNPROCESSABLE_CONTENT, detail=str(exc)
|
|
) from exc
|
|
|
|
|
|
@router.post(
|
|
"/{campaign_id}/jobs/send-unattempted", response_model=CampaignActionResponse
|
|
)
|
|
def send_unattempted_campaign_jobs(
|
|
campaign_id: str,
|
|
payload: CampaignSendUnattemptedRequest | None = None,
|
|
session: Session = Depends(get_session),
|
|
principal: ApiPrincipal = Depends(require_scope("campaigns:campaign:queue")),
|
|
):
|
|
_get_campaign_for_principal(session, campaign_id, principal, write=True)
|
|
_require_permission(principal, "campaigns:recipient:read")
|
|
payload = payload or CampaignSendUnattemptedRequest()
|
|
_require_campaign_profile_use_if_needed(
|
|
session, principal, campaign_id, payload.version_id
|
|
)
|
|
try:
|
|
result = queue_unattempted_jobs(
|
|
session,
|
|
tenant_id=principal.tenant_id,
|
|
campaign_id=campaign_id,
|
|
version_id=payload.version_id,
|
|
job_ids=payload.job_ids or None,
|
|
enqueue_celery=payload.enqueue_celery,
|
|
dry_run=payload.dry_run,
|
|
)
|
|
audit_from_principal(
|
|
session,
|
|
principal,
|
|
action="campaign.unattempted_jobs_queued"
|
|
if not payload.dry_run
|
|
else "campaign.unattempted_jobs_dry_run",
|
|
object_type="campaign",
|
|
object_id=campaign_id,
|
|
details=result,
|
|
commit=True,
|
|
)
|
|
return CampaignActionResponse(result=result)
|
|
except (QueueingError, ExecutionSnapshotError) as exc:
|
|
raise HTTPException(
|
|
status_code=status.HTTP_422_UNPROCESSABLE_CONTENT, detail=str(exc)
|
|
) from exc
|
|
|
|
|
|
@router.post("/{campaign_id}/jobs/{job_id}/send", response_model=CampaignActionResponse)
|
|
def send_single_campaign_job_endpoint(
|
|
campaign_id: str,
|
|
job_id: str,
|
|
payload: CampaignSendJobRequest,
|
|
session: Session = Depends(get_session),
|
|
principal: ApiPrincipal = Depends(
|
|
require_any_scope(
|
|
"campaigns:campaign:send",
|
|
"campaigns:campaign:send_test",
|
|
)
|
|
),
|
|
):
|
|
_get_campaign_for_principal(session, campaign_id, principal, write=True)
|
|
_require_permission(principal, "campaigns:recipient:read")
|
|
_require_permission(
|
|
principal,
|
|
(
|
|
"campaigns:campaign:send_test"
|
|
if payload.kind == "test"
|
|
else "campaigns:campaign:send"
|
|
),
|
|
)
|
|
_require_campaign_profile_use_if_needed(session, principal, campaign_id, None)
|
|
try:
|
|
result = send_single_campaign_job(
|
|
session,
|
|
tenant_id=principal.tenant_id,
|
|
campaign_id=campaign_id,
|
|
job_id=job_id,
|
|
kind=payload.kind,
|
|
idempotency_key=payload.idempotency_key,
|
|
actor_user_id=principal.user.id,
|
|
actor_api_key_id=getattr(principal, "api_key_id", None),
|
|
reason=payload.reason,
|
|
action_context=payload.context,
|
|
include_warnings=payload.include_warnings,
|
|
use_rate_limit=payload.use_rate_limit,
|
|
enqueue_imap_task=payload.enqueue_imap_task,
|
|
)
|
|
audit_from_principal(
|
|
session,
|
|
principal,
|
|
action=f"campaign.message_{payload.kind}",
|
|
object_type="campaign_job",
|
|
object_id=job_id,
|
|
details=result,
|
|
commit=True,
|
|
)
|
|
return CampaignActionResponse(result=result)
|
|
except (QueueingError, ExecutionSnapshotError) as exc:
|
|
raise HTTPException(
|
|
status_code=status.HTTP_422_UNPROCESSABLE_CONTENT, detail=str(exc)
|
|
) from exc
|
|
except Exception as exc:
|
|
logger.exception(
|
|
"Unexpected single-message campaign action failure",
|
|
extra={
|
|
"campaign_id": campaign_id,
|
|
"job_id": job_id,
|
|
"action_kind": payload.kind,
|
|
},
|
|
)
|
|
raise HTTPException(
|
|
status_code=status.HTTP_500_INTERNAL_SERVER_ERROR,
|
|
detail="The message action failed because of an internal error.",
|
|
) from exc
|
|
|
|
|
|
@router.post(
|
|
"/{campaign_id}/jobs/{job_id}/resolve-outcome",
|
|
response_model=CampaignActionResponse,
|
|
)
|
|
def resolve_campaign_job_outcome(
|
|
campaign_id: str,
|
|
job_id: str,
|
|
payload: CampaignResolveOutcomeRequest,
|
|
session: Session = Depends(get_session),
|
|
principal: ApiPrincipal = Depends(require_scope("campaigns:campaign:reconcile")),
|
|
):
|
|
_get_campaign_for_principal(session, campaign_id, principal, write=True)
|
|
_require_permission(principal, "campaigns:recipient:read")
|
|
try:
|
|
result = reconcile_job_outcome(
|
|
session,
|
|
tenant_id=principal.tenant_id,
|
|
campaign_id=campaign_id,
|
|
job_id=job_id,
|
|
decision=payload.decision,
|
|
note=payload.note,
|
|
attempt_id=payload.attempt_id,
|
|
commit=False,
|
|
)
|
|
audit_from_principal(
|
|
session,
|
|
principal,
|
|
action="campaign.job_outcome_reconciled",
|
|
object_type="campaign_job",
|
|
object_id=job_id,
|
|
details=result,
|
|
commit=True,
|
|
)
|
|
return CampaignActionResponse(result=result)
|
|
except (QueueingError, ExecutionSnapshotError) as exc:
|
|
session.rollback()
|
|
raise HTTPException(
|
|
status_code=status.HTTP_422_UNPROCESSABLE_CONTENT, detail=str(exc)
|
|
) from exc
|
|
except Exception:
|
|
session.rollback()
|
|
raise
|
|
|
|
|
|
@router.post("/{campaign_id}/mock-send", response_model=MockCampaignSendResponse)
|
|
def mock_send_campaign(
|
|
campaign_id: str,
|
|
payload: MockCampaignSendRequest | None = None,
|
|
session: Session = Depends(get_session),
|
|
principal: ApiPrincipal = Depends(require_scope("campaigns:campaign:send_test")),
|
|
):
|
|
"""Run a fully visible mock delivery flow without mutating campaign state.
|
|
|
|
The route validates and builds the selected version, then optionally records
|
|
mock SMTP deliveries and mock IMAP appends. It never talks to the configured
|
|
real SMTP/IMAP servers and it does not mark the version sent/final.
|
|
"""
|
|
_get_campaign_for_principal(session, campaign_id, principal, write=True)
|
|
_require_permission(principal, "campaigns:recipient:read")
|
|
|
|
payload = payload or MockCampaignSendRequest()
|
|
_require_campaign_profile_use_if_needed(
|
|
session, principal, campaign_id, payload.version_id
|
|
)
|
|
try:
|
|
result = run_mock_campaign_send(
|
|
session,
|
|
tenant_id=principal.tenant_id,
|
|
campaign_id=campaign_id,
|
|
version_id=payload.version_id,
|
|
send=payload.send,
|
|
include_warnings=payload.include_warnings,
|
|
include_needs_review=payload.include_needs_review,
|
|
append_sent=payload.append_sent,
|
|
clear_mailbox=payload.clear_mailbox,
|
|
check_files=payload.check_files,
|
|
)
|
|
audit_from_principal(
|
|
session,
|
|
principal,
|
|
action="campaign.mock_send"
|
|
if payload.send
|
|
else "campaign.mock_send_review",
|
|
object_type="campaign",
|
|
object_id=campaign_id,
|
|
details={
|
|
"version_id": result.get("version_id"),
|
|
"send_requested": payload.send,
|
|
"sent_count": result.get("send", {}).get("sent_count"),
|
|
"failed_count": result.get("send", {}).get("failed_count"),
|
|
},
|
|
commit=True,
|
|
)
|
|
return MockCampaignSendResponse(result=result)
|
|
except MockCampaignSendError as exc:
|
|
raise HTTPException(
|
|
status_code=status.HTTP_422_UNPROCESSABLE_CONTENT, detail=str(exc)
|
|
) from exc
|
|
except Exception as exc:
|
|
raise HTTPException(
|
|
status_code=status.HTTP_422_UNPROCESSABLE_CONTENT, detail=str(exc)
|
|
) from exc
|
|
|
|
|
|
@router.post("/{campaign_id}/send-now", response_model=SendCampaignNowResponse)
|
|
def send_campaign_now_endpoint(
|
|
campaign_id: str,
|
|
payload: SendCampaignNowRequest | None = None,
|
|
session: Session = Depends(get_session),
|
|
principal: ApiPrincipal = Depends(require_scope("campaigns:campaign:send")),
|
|
):
|
|
"""Preflight and synchronously send a policy-bounded built execution."""
|
|
_get_campaign_for_principal(session, campaign_id, principal, write=True)
|
|
_require_permission(principal, "campaigns:recipient:read")
|
|
|
|
payload = payload or SendCampaignNowRequest()
|
|
try:
|
|
campaign = _get_campaign_for_tenant(session, campaign_id, principal.tenant_id)
|
|
version_id = payload.version_id or campaign.current_version_id
|
|
if not version_id:
|
|
raise HTTPException(
|
|
status_code=status.HTTP_422_UNPROCESSABLE_CONTENT,
|
|
detail="Campaign has no current version",
|
|
)
|
|
|
|
version = _get_version_for_tenant(session, version_id, principal.tenant_id)
|
|
_require_mail_profile_use_if_needed(
|
|
principal, version.raw_json if isinstance(version.raw_json, dict) else {}
|
|
)
|
|
validation_result: dict[str, object] | None = (
|
|
version.validation_summary
|
|
if isinstance(version.validation_summary, dict)
|
|
else None
|
|
)
|
|
build_result: dict[str, object] | None = (
|
|
version.build_summary if isinstance(version.build_summary, dict) else None
|
|
)
|
|
if is_user_locked_version(version):
|
|
raise HTTPException(
|
|
status_code=status.HTTP_409_CONFLICT,
|
|
detail="User-locked audit-safe versions cannot be dry-run or sent. Create an editable copy and validate it instead.",
|
|
)
|
|
if (
|
|
not version.locked_at
|
|
or not validation_result
|
|
or validation_result.get("ok") is not True
|
|
):
|
|
raise HTTPException(
|
|
status_code=status.HTTP_422_UNPROCESSABLE_CONTENT,
|
|
detail="Campaign version must be validated and locked before dry-run or sending.",
|
|
)
|
|
if not build_result:
|
|
raise HTTPException(
|
|
status_code=status.HTTP_422_UNPROCESSABLE_CONTENT,
|
|
detail="Campaign version must be built before dry-run or sending.",
|
|
)
|
|
|
|
delivery_result = send_campaign_now(
|
|
session,
|
|
tenant_id=principal.tenant_id,
|
|
campaign_id=campaign_id,
|
|
version_id=version_id,
|
|
include_warnings=payload.include_warnings,
|
|
dry_run=payload.dry_run,
|
|
use_rate_limit=payload.use_rate_limit,
|
|
enqueue_imap_task=payload.enqueue_imap_task,
|
|
).as_dict()
|
|
response_result = public_send_campaign_now_result(
|
|
delivery_result,
|
|
validation_summary=validation_result,
|
|
build_summary=build_result,
|
|
)
|
|
audit_from_principal(
|
|
session,
|
|
principal,
|
|
action="campaign.sent_now"
|
|
if not payload.dry_run
|
|
else "campaign.send_now_dry_run",
|
|
object_type="campaign",
|
|
object_id=campaign_id,
|
|
details=send_campaign_now_audit_details(delivery_result),
|
|
commit=True,
|
|
)
|
|
return SendCampaignNowResponse(result=response_result)
|
|
except SynchronousSendRejected as exc:
|
|
# A synchronous request stages queue state before the all-message
|
|
# preflight can run. Rejecting that preflight must not leave work
|
|
# eligible for a background worker when no provider effect occurred.
|
|
session.rollback()
|
|
audit_from_principal(
|
|
session,
|
|
principal,
|
|
action="campaign.send_now_rejected",
|
|
object_type="campaign",
|
|
object_id=campaign_id,
|
|
details={
|
|
**exc.audit_details(),
|
|
"version_id": payload.version_id,
|
|
},
|
|
commit=True,
|
|
)
|
|
raise HTTPException(
|
|
status_code=status.HTTP_422_UNPROCESSABLE_CONTENT, detail=str(exc)
|
|
) from exc
|
|
except HTTPException:
|
|
raise
|
|
except (CampaignPersistenceError, QueueingError) as exc:
|
|
raise HTTPException(
|
|
status_code=status.HTTP_422_UNPROCESSABLE_CONTENT, detail=str(exc)
|
|
) from exc
|
|
except Exception as exc:
|
|
raise HTTPException(
|
|
status_code=status.HTTP_422_UNPROCESSABLE_CONTENT, detail=str(exc)
|
|
) from exc
|
|
|
|
|
|
@router.post("/{campaign_id}/pause", response_model=CampaignActionResponse)
|
|
def pause_campaign(
|
|
campaign_id: str,
|
|
session: Session = Depends(get_session),
|
|
principal: ApiPrincipal = Depends(require_scope("campaigns:campaign:control")),
|
|
):
|
|
_get_campaign_for_principal(session, campaign_id, principal, write=True)
|
|
try:
|
|
result = pause_campaign_jobs(
|
|
session, tenant_id=principal.tenant_id, campaign_id=campaign_id
|
|
)
|
|
audit_from_principal(
|
|
session,
|
|
principal,
|
|
action="campaign.paused",
|
|
object_type="campaign",
|
|
object_id=campaign_id,
|
|
details=result,
|
|
commit=True,
|
|
)
|
|
return CampaignActionResponse(result=result)
|
|
except QueueingError as exc:
|
|
raise HTTPException(
|
|
status_code=status.HTTP_422_UNPROCESSABLE_CONTENT, detail=str(exc)
|
|
) from exc
|
|
|
|
|
|
@router.post("/{campaign_id}/resume", response_model=CampaignActionResponse)
|
|
def resume_campaign(
|
|
campaign_id: str,
|
|
session: Session = Depends(get_session),
|
|
principal: ApiPrincipal = Depends(require_scope("campaigns:campaign:control")),
|
|
):
|
|
_get_campaign_for_principal(session, campaign_id, principal, write=True)
|
|
version_ids = {
|
|
row[0]
|
|
for row in session.query(CampaignJob.campaign_version_id)
|
|
.filter(
|
|
CampaignJob.tenant_id == principal.tenant_id,
|
|
CampaignJob.campaign_id == campaign_id,
|
|
CampaignJob.queue_status == JobQueueStatus.PAUSED.value,
|
|
)
|
|
.distinct()
|
|
.all()
|
|
}
|
|
_require_campaign_versions_profile_use(session, principal, campaign_id, version_ids)
|
|
try:
|
|
result = resume_campaign_jobs(
|
|
session, tenant_id=principal.tenant_id, campaign_id=campaign_id
|
|
)
|
|
audit_from_principal(
|
|
session,
|
|
principal,
|
|
action="campaign.resumed",
|
|
object_type="campaign",
|
|
object_id=campaign_id,
|
|
details=result,
|
|
commit=True,
|
|
)
|
|
return CampaignActionResponse(result=result)
|
|
except QueueingError as exc:
|
|
raise HTTPException(
|
|
status_code=status.HTTP_422_UNPROCESSABLE_CONTENT, detail=str(exc)
|
|
) from exc
|
|
|
|
|
|
@router.post("/{campaign_id}/cancel", response_model=CampaignActionResponse)
|
|
def cancel_campaign(
|
|
campaign_id: str,
|
|
session: Session = Depends(get_session),
|
|
principal: ApiPrincipal = Depends(require_scope("campaigns:campaign:control")),
|
|
):
|
|
_get_campaign_for_principal(session, campaign_id, principal, write=True)
|
|
try:
|
|
result = cancel_campaign_jobs(
|
|
session, tenant_id=principal.tenant_id, campaign_id=campaign_id
|
|
)
|
|
audit_from_principal(
|
|
session,
|
|
principal,
|
|
action="campaign.cancelled",
|
|
object_type="campaign",
|
|
object_id=campaign_id,
|
|
details=result,
|
|
commit=True,
|
|
)
|
|
return CampaignActionResponse(result=result)
|
|
except QueueingError as exc:
|
|
raise HTTPException(
|
|
status_code=status.HTTP_422_UNPROCESSABLE_CONTENT, detail=str(exc)
|
|
) from exc
|
|
|
|
|
|
@router.post("/{campaign_id}/append-sent", response_model=CampaignActionResponse)
|
|
def append_sent(
|
|
campaign_id: str,
|
|
payload: AppendSentRequest | None = None,
|
|
session: Session = Depends(get_session),
|
|
principal: ApiPrincipal = Depends(require_scope("campaigns:campaign:send")),
|
|
):
|
|
_get_campaign_for_principal(session, campaign_id, principal, write=True)
|
|
payload = payload or AppendSentRequest()
|
|
version_ids = {
|
|
row[0]
|
|
for row in session.query(CampaignJob.campaign_version_id)
|
|
.filter(
|
|
CampaignJob.tenant_id == principal.tenant_id,
|
|
CampaignJob.campaign_id == campaign_id,
|
|
CampaignJob.send_status.in_(
|
|
[JobSendStatus.SMTP_ACCEPTED.value, JobSendStatus.SENT.value]
|
|
),
|
|
CampaignJob.imap_status.in_(
|
|
[JobImapStatus.PENDING.value, JobImapStatus.FAILED.value]
|
|
),
|
|
)
|
|
.distinct()
|
|
.all()
|
|
}
|
|
_require_campaign_versions_profile_use(session, principal, campaign_id, version_ids)
|
|
try:
|
|
result = enqueue_pending_imap_appends(
|
|
session,
|
|
tenant_id=principal.tenant_id,
|
|
campaign_id=campaign_id,
|
|
enqueue_celery=payload.enqueue_celery,
|
|
run_inline=payload.run_inline,
|
|
dry_run=payload.dry_run,
|
|
)
|
|
audit_from_principal(
|
|
session,
|
|
principal,
|
|
action="campaign.append_sent_enqueued"
|
|
if not payload.dry_run
|
|
else "campaign.append_sent_dry_run",
|
|
object_type="campaign",
|
|
object_id=campaign_id,
|
|
details=result,
|
|
commit=True,
|
|
)
|
|
return CampaignActionResponse(result=result)
|
|
except QueueingError as exc:
|
|
raise HTTPException(
|
|
status_code=status.HTTP_422_UNPROCESSABLE_CONTENT, detail=str(exc)
|
|
) from exc
|