refactor(api): split campaign workflow routers
This commit is contained in:
@@ -0,0 +1,884 @@
|
||||
from __future__ import annotations
|
||||
|
||||
import json
|
||||
from collections.abc import Sequence
|
||||
from typing import Literal
|
||||
|
||||
from fastapi import HTTPException, Query, status
|
||||
from sqlalchemy import and_, func, or_
|
||||
from sqlalchemy.orm import Session
|
||||
|
||||
from govoplan_campaign.backend.schemas import (
|
||||
CampaignJobsResponse,
|
||||
CampaignJobDiagnosticsResponse,
|
||||
)
|
||||
from govoplan_core.auth import ApiPrincipal
|
||||
from govoplan_core.core.change_sequence import (
|
||||
encode_sequence_watermark,
|
||||
max_sequence_id,
|
||||
)
|
||||
from govoplan_core.core.pagination import (
|
||||
KeysetCursorError,
|
||||
decode_keyset_cursor,
|
||||
encode_keyset_cursor,
|
||||
keyset_query_fingerprint,
|
||||
)
|
||||
from govoplan_campaign.backend.change_tracking import (
|
||||
CAMPAIGNS_MODULE_ID,
|
||||
CAMPAIGN_JOBS_COLLECTION,
|
||||
)
|
||||
from govoplan_campaign.backend.db.models import (
|
||||
Campaign,
|
||||
CampaignJob,
|
||||
CampaignVersion,
|
||||
ImapAppendAttempt,
|
||||
JobImapStatus,
|
||||
JobPostboxStatus,
|
||||
JobQueueStatus,
|
||||
JobSendStatus,
|
||||
JobValidationStatus,
|
||||
PostboxDeliveryAttempt,
|
||||
SendAttempt,
|
||||
)
|
||||
from govoplan_campaign.backend.response_security import (
|
||||
public_campaign_payload,
|
||||
public_delivery_result_message,
|
||||
)
|
||||
|
||||
|
||||
from govoplan_campaign.backend.route_support import (
|
||||
_get_campaign_for_principal,
|
||||
_get_campaign_for_tenant,
|
||||
_get_version_for_tenant,
|
||||
_require_permission,
|
||||
)
|
||||
|
||||
|
||||
CAMPAIGN_JOBS_CURSOR_SCOPE = "campaign.jobs"
|
||||
|
||||
|
||||
def _job_review_key(job: CampaignJob) -> str:
|
||||
return str(job.entry_id or job.entry_index)
|
||||
|
||||
|
||||
def _job_summary_payload(
|
||||
job: CampaignJob,
|
||||
*,
|
||||
reviewed_keys: set[str] | None = None,
|
||||
) -> dict[str, object]:
|
||||
review_key = _job_review_key(job)
|
||||
return {
|
||||
"id": job.id,
|
||||
"campaign_version_id": job.campaign_version_id,
|
||||
"entry_index": job.entry_index,
|
||||
"entry_id": job.entry_id,
|
||||
"recipient_email": job.recipient_email,
|
||||
"subject": job.subject,
|
||||
"message_id_header": job.message_id_header,
|
||||
"build_status": job.build_status,
|
||||
"validation_status": job.validation_status,
|
||||
"queue_status": job.queue_status,
|
||||
"send_status": job.send_status,
|
||||
"delivery_channel_policy": getattr(job, "delivery_channel_policy", "mail"),
|
||||
"postbox_status": getattr(job, "postbox_status", "not_requested"),
|
||||
"imap_status": job.imap_status,
|
||||
"eml_size_bytes": job.eml_size_bytes,
|
||||
"eml_sha256": job.eml_sha256,
|
||||
"attempt_count": job.attempt_count,
|
||||
"postbox_attempt_count": getattr(job, "postbox_attempt_count", 0),
|
||||
"postbox_target_count": len(
|
||||
getattr(job, "resolved_postbox_targets", None) or []
|
||||
),
|
||||
"last_error": public_delivery_result_message(
|
||||
last_error=job.last_error,
|
||||
send_status=job.send_status,
|
||||
imap_status=job.imap_status,
|
||||
postbox_status=getattr(job, "postbox_status", "not_requested"),
|
||||
),
|
||||
"queued_at": job.queued_at,
|
||||
"outcome_unknown_at": job.outcome_unknown_at,
|
||||
"sent_at": job.sent_at,
|
||||
"created_at": job.created_at,
|
||||
"updated_at": job.updated_at,
|
||||
"issues_count": len(job.issues_snapshot or []),
|
||||
"attachment_count": len(job.resolved_attachments or []),
|
||||
"review_key": review_key,
|
||||
"reviewed": review_key in reviewed_keys if reviewed_keys is not None else False,
|
||||
"matched_file_count": sum(
|
||||
len(item.get("matches") or [])
|
||||
for item in (job.resolved_attachments or [])
|
||||
if isinstance(item, dict)
|
||||
),
|
||||
}
|
||||
|
||||
|
||||
def _job_detail_payload(job: CampaignJob) -> dict[str, object]:
|
||||
return {
|
||||
**_job_summary_payload(job),
|
||||
"message_id_header": job.message_id_header,
|
||||
"issues": job.issues_snapshot or [],
|
||||
"attachments": public_campaign_payload(job.resolved_attachments or []),
|
||||
"resolved_recipients": job.resolved_recipients or {},
|
||||
"resolved_postbox_targets": getattr(job, "resolved_postbox_targets", None)
|
||||
or [],
|
||||
}
|
||||
|
||||
|
||||
def _job_attempts_payload(
|
||||
send_attempts: list[SendAttempt],
|
||||
imap_attempts: list[ImapAppendAttempt],
|
||||
postbox_attempts: Sequence[PostboxDeliveryAttempt] = (),
|
||||
*,
|
||||
include_diagnostics: bool = False,
|
||||
) -> dict[str, list[dict[str, object]]]:
|
||||
smtp_payloads: list[dict[str, object]] = []
|
||||
for attempt in send_attempts:
|
||||
payload: dict[str, object] = {
|
||||
"id": attempt.id,
|
||||
"attempt_number": attempt.attempt_number,
|
||||
"status": attempt.status,
|
||||
"smtp_status_code": attempt.smtp_status_code,
|
||||
"started_at": attempt.started_at,
|
||||
"finished_at": attempt.finished_at,
|
||||
}
|
||||
if include_diagnostics:
|
||||
payload["claim_token"] = attempt.claim_token
|
||||
payload["smtp_response"] = attempt.smtp_response
|
||||
payload["error_type"] = attempt.error_type
|
||||
payload["error_message"] = attempt.error_message
|
||||
smtp_payloads.append(payload)
|
||||
|
||||
imap_payloads: list[dict[str, object]] = []
|
||||
for attempt in imap_attempts:
|
||||
payload = {
|
||||
"id": attempt.id,
|
||||
"attempt_number": attempt.attempt_number,
|
||||
"status": attempt.status,
|
||||
"folder": attempt.folder,
|
||||
"created_at": attempt.created_at,
|
||||
"updated_at": attempt.updated_at,
|
||||
}
|
||||
if include_diagnostics:
|
||||
payload["claim_token"] = attempt.claim_token
|
||||
payload["error_message"] = attempt.error_message
|
||||
imap_payloads.append(payload)
|
||||
postbox_payloads: list[dict[str, object]] = []
|
||||
for attempt in postbox_attempts:
|
||||
payload = {
|
||||
"id": attempt.id,
|
||||
"target_index": attempt.target_index,
|
||||
"attempt_number": attempt.attempt_number,
|
||||
"status": attempt.status,
|
||||
"postbox_id": attempt.postbox_id,
|
||||
"address": attempt.address,
|
||||
"holder_count": attempt.holder_count,
|
||||
"vacant": attempt.vacant,
|
||||
"duplicate": attempt.duplicate,
|
||||
"target": attempt.target_snapshot or {},
|
||||
"started_at": attempt.started_at,
|
||||
"finished_at": attempt.finished_at,
|
||||
"error_code": attempt.error_code,
|
||||
}
|
||||
if include_diagnostics:
|
||||
payload["idempotency_key"] = attempt.idempotency_key
|
||||
payload["provider_delivery_id"] = attempt.provider_delivery_id
|
||||
payload["provider_message_id"] = attempt.provider_message_id
|
||||
payload["evidence"] = attempt.evidence or {}
|
||||
payload["error_type"] = attempt.error_type
|
||||
payload["error_message"] = attempt.error_message
|
||||
postbox_payloads.append(payload)
|
||||
return {
|
||||
"smtp": smtp_payloads,
|
||||
"imap": imap_payloads,
|
||||
"postbox": postbox_payloads,
|
||||
}
|
||||
|
||||
|
||||
def _job_diagnostics_payload(
|
||||
job: CampaignJob,
|
||||
send_attempts: list[SendAttempt],
|
||||
imap_attempts: list[ImapAppendAttempt],
|
||||
postbox_attempts: Sequence[PostboxDeliveryAttempt] = (),
|
||||
) -> CampaignJobDiagnosticsResponse:
|
||||
return CampaignJobDiagnosticsResponse(
|
||||
job_id=job.id,
|
||||
campaign_id=job.campaign_id,
|
||||
campaign_version_id=job.campaign_version_id,
|
||||
storage={
|
||||
"eml_local_path": job.eml_local_path,
|
||||
"eml_storage_key": job.eml_storage_key,
|
||||
"eml_size_bytes": job.eml_size_bytes,
|
||||
"eml_sha256": job.eml_sha256,
|
||||
},
|
||||
worker_claim={
|
||||
"claim_token": job.claim_token,
|
||||
"claimed_at": job.claimed_at,
|
||||
"smtp_started_at": job.smtp_started_at,
|
||||
"outcome_unknown_at": job.outcome_unknown_at,
|
||||
"last_error": job.last_error,
|
||||
},
|
||||
attempts=_job_attempts_payload(
|
||||
send_attempts,
|
||||
imap_attempts,
|
||||
postbox_attempts,
|
||||
include_diagnostics=True,
|
||||
),
|
||||
)
|
||||
|
||||
|
||||
def _review_metadata(
|
||||
session: Session,
|
||||
version: CampaignVersion | None,
|
||||
base_filters: list[object],
|
||||
) -> tuple[dict[str, object], set[str]]:
|
||||
if version is None:
|
||||
return _empty_review_metadata(), set()
|
||||
|
||||
review_state, state_is_current = _current_review_state(version)
|
||||
reviewed_keys = _current_reviewed_keys(
|
||||
review_state, state_is_current=state_is_current
|
||||
)
|
||||
counts = _review_metadata_counts(_review_rows(session, base_filters), reviewed_keys)
|
||||
return {
|
||||
"inspection_complete": bool(
|
||||
state_is_current and review_state.get("inspection_complete") is True
|
||||
),
|
||||
**counts,
|
||||
}, reviewed_keys
|
||||
|
||||
|
||||
def _empty_review_metadata() -> dict[str, object]:
|
||||
return {
|
||||
"inspection_complete": False,
|
||||
"blocking_count": 0,
|
||||
"required_count": 0,
|
||||
"reviewed_required_count": 0,
|
||||
"bulk_acceptable_count": 0,
|
||||
}
|
||||
|
||||
|
||||
def _current_review_state(version: CampaignVersion) -> tuple[dict[str, object], bool]:
|
||||
build_summary = (
|
||||
version.build_summary if isinstance(version.build_summary, dict) else {}
|
||||
)
|
||||
build_token = str(
|
||||
build_summary.get("build_token") or build_summary.get("built_at") or ""
|
||||
)
|
||||
editor_state = (
|
||||
version.editor_state if isinstance(version.editor_state, dict) else {}
|
||||
)
|
||||
review_state = (
|
||||
editor_state.get("review_send")
|
||||
if isinstance(editor_state.get("review_send"), dict)
|
||||
else {}
|
||||
)
|
||||
state_token = str(review_state.get("build_token") or "")
|
||||
return review_state, bool(build_token and state_token == build_token)
|
||||
|
||||
|
||||
def _current_reviewed_keys(
|
||||
review_state: dict[str, object], *, state_is_current: bool
|
||||
) -> set[str]:
|
||||
return {
|
||||
str(value)
|
||||
for value in (review_state.get("reviewed_message_keys") or [])
|
||||
if state_is_current and str(value).strip()
|
||||
}
|
||||
|
||||
|
||||
def _review_rows(
|
||||
session: Session, base_filters: list[object]
|
||||
) -> list[tuple[object, int, str, str]]:
|
||||
return (
|
||||
session.query(
|
||||
CampaignJob.entry_id,
|
||||
CampaignJob.entry_index,
|
||||
CampaignJob.build_status,
|
||||
CampaignJob.validation_status,
|
||||
)
|
||||
.filter(*base_filters)
|
||||
.all()
|
||||
)
|
||||
|
||||
|
||||
def _review_metadata_counts(
|
||||
review_rows: list[tuple[object, int, str, str]], reviewed_keys: set[str]
|
||||
) -> dict[str, int]:
|
||||
blocking_count = 0
|
||||
required_count = 0
|
||||
reviewed_required_count = 0
|
||||
bulk_acceptable_count = 0
|
||||
for entry_id, entry_index, build_status, validation_status in review_rows:
|
||||
key = str(entry_id or entry_index)
|
||||
if build_status != "built" or validation_status == "blocked":
|
||||
blocking_count += 1
|
||||
if validation_status == "needs_review":
|
||||
required_count += 1
|
||||
if key in reviewed_keys:
|
||||
reviewed_required_count += 1
|
||||
elif validation_status in {"warning", "excluded"}:
|
||||
bulk_acceptable_count += 1
|
||||
|
||||
return {
|
||||
"blocking_count": blocking_count,
|
||||
"required_count": required_count,
|
||||
"reviewed_required_count": reviewed_required_count,
|
||||
"bulk_acceptable_count": bulk_acceptable_count,
|
||||
}
|
||||
|
||||
|
||||
def _status_counts(
|
||||
session: Session, filters: list[object]
|
||||
) -> dict[str, dict[str, int]]:
|
||||
result: dict[str, dict[str, int]] = {}
|
||||
for field_name in (
|
||||
"build_status",
|
||||
"validation_status",
|
||||
"queue_status",
|
||||
"send_status",
|
||||
"postbox_status",
|
||||
"imap_status",
|
||||
):
|
||||
column = getattr(CampaignJob, field_name)
|
||||
rows = (
|
||||
session.query(column, func.count(CampaignJob.id))
|
||||
.filter(*filters)
|
||||
.group_by(column)
|
||||
.all()
|
||||
)
|
||||
result[field_name.removesuffix("_status")] = {
|
||||
str(value or "unknown"): int(count) for value, count in rows
|
||||
}
|
||||
return result
|
||||
|
||||
|
||||
CAMPAIGN_JOB_GRID_SORT_COLUMNS = {
|
||||
"number": CampaignJob.entry_index,
|
||||
"recipient": func.lower(func.coalesce(CampaignJob.recipient_email, "")),
|
||||
"subject": func.lower(func.coalesce(CampaignJob.subject, "")),
|
||||
"validation": CampaignJob.validation_status,
|
||||
"queue": CampaignJob.queue_status,
|
||||
"send": CampaignJob.send_status,
|
||||
"postbox": CampaignJob.postbox_status,
|
||||
"imap": CampaignJob.imap_status,
|
||||
"attempts": CampaignJob.attempt_count,
|
||||
"updated": CampaignJob.updated_at,
|
||||
}
|
||||
CAMPAIGN_JOB_GRID_LIST_FILTERS = {
|
||||
"validation": (
|
||||
CampaignJob.validation_status,
|
||||
{item.value for item in JobValidationStatus},
|
||||
),
|
||||
"queue": (CampaignJob.queue_status, {item.value for item in JobQueueStatus}),
|
||||
"send": (CampaignJob.send_status, {item.value for item in JobSendStatus}),
|
||||
"postbox": (
|
||||
CampaignJob.postbox_status,
|
||||
{item.value for item in JobPostboxStatus},
|
||||
),
|
||||
"imap": (CampaignJob.imap_status, {item.value for item in JobImapStatus}),
|
||||
}
|
||||
|
||||
|
||||
def _campaign_jobs_grid_filter_expressions(
|
||||
grid_filters: dict[str, str] | None,
|
||||
) -> list[object]:
|
||||
values = grid_filters or {}
|
||||
expressions: list[object] = []
|
||||
recipient = values.get("recipient", "").strip()
|
||||
if recipient:
|
||||
pattern = _contains_pattern(recipient)
|
||||
expressions.append(
|
||||
or_(
|
||||
CampaignJob.recipient_email.ilike(pattern, escape="\\"),
|
||||
CampaignJob.entry_id.ilike(pattern, escape="\\"),
|
||||
)
|
||||
)
|
||||
subject = values.get("subject", "").strip()
|
||||
if subject:
|
||||
expressions.append(
|
||||
CampaignJob.subject.ilike(_contains_pattern(subject), escape="\\")
|
||||
)
|
||||
evidence = values.get("evidence", "").strip()
|
||||
if evidence:
|
||||
pattern = _contains_pattern(evidence)
|
||||
expressions.append(
|
||||
or_(
|
||||
CampaignJob.message_id_header.ilike(pattern, escape="\\"),
|
||||
CampaignJob.eml_sha256.ilike(pattern, escape="\\"),
|
||||
)
|
||||
)
|
||||
attempts = values.get("attempts", "").strip()
|
||||
if attempts:
|
||||
expressions.append(
|
||||
_campaign_jobs_integer_filter(
|
||||
CampaignJob.attempt_count, attempts, column_id="attempts"
|
||||
)
|
||||
)
|
||||
for column_id, (column, allowed_values) in CAMPAIGN_JOB_GRID_LIST_FILTERS.items():
|
||||
raw_value = values.get(column_id, "").strip()
|
||||
if not raw_value:
|
||||
continue
|
||||
selected = _campaign_jobs_list_filter(
|
||||
raw_value, column_id=column_id, allowed_values=allowed_values
|
||||
)
|
||||
expressions.append(column.in_(selected))
|
||||
return expressions
|
||||
|
||||
|
||||
def _campaign_jobs_list_filter(
|
||||
raw_value: str, *, column_id: str, allowed_values: set[str]
|
||||
) -> list[str]:
|
||||
if raw_value.startswith("list:"):
|
||||
try:
|
||||
parsed = json.loads(raw_value[5:])
|
||||
except json.JSONDecodeError as exc:
|
||||
raise HTTPException(
|
||||
status_code=status.HTTP_422_UNPROCESSABLE_CONTENT,
|
||||
detail=f"Invalid {column_id} list filter",
|
||||
) from exc
|
||||
if not isinstance(parsed, list) or any(
|
||||
not isinstance(value, str) for value in parsed
|
||||
):
|
||||
raise HTTPException(
|
||||
status_code=status.HTTP_422_UNPROCESSABLE_CONTENT,
|
||||
detail=f"Invalid {column_id} list filter",
|
||||
)
|
||||
selected = list(
|
||||
dict.fromkeys(value.strip() for value in parsed if value.strip())
|
||||
)
|
||||
else:
|
||||
selected = list(
|
||||
dict.fromkeys(
|
||||
value.strip() for value in raw_value.split(",") if value.strip()
|
||||
)
|
||||
)
|
||||
if len(selected) > 50 or any(value not in allowed_values for value in selected):
|
||||
raise HTTPException(
|
||||
status_code=status.HTTP_422_UNPROCESSABLE_CONTENT,
|
||||
detail=f"Invalid {column_id} list filter",
|
||||
)
|
||||
return selected
|
||||
|
||||
|
||||
def _campaign_jobs_integer_filter(column: object, raw_value: str, *, column_id: str):
|
||||
operator, separator, value = raw_value.partition(":")
|
||||
if not separator:
|
||||
operator, value = "eq", operator
|
||||
if operator not in {"eq", "gt", "gte", "lt", "lte"}:
|
||||
raise HTTPException(
|
||||
status_code=status.HTTP_422_UNPROCESSABLE_CONTENT,
|
||||
detail=f"Invalid {column_id} filter operator",
|
||||
)
|
||||
try:
|
||||
expected = int(value)
|
||||
except ValueError as exc:
|
||||
raise HTTPException(
|
||||
status_code=status.HTTP_422_UNPROCESSABLE_CONTENT,
|
||||
detail=f"Invalid {column_id} filter value",
|
||||
) from exc
|
||||
if operator == "gt":
|
||||
return column > expected
|
||||
if operator == "gte":
|
||||
return column >= expected
|
||||
if operator == "lt":
|
||||
return column < expected
|
||||
if operator == "lte":
|
||||
return column <= expected
|
||||
return column == expected
|
||||
|
||||
|
||||
def _contains_pattern(value: str) -> str:
|
||||
escaped = value.replace("\\", "\\\\").replace("%", "\\%").replace("_", "\\_")
|
||||
return f"%{escaped}%"
|
||||
|
||||
|
||||
def _campaign_jobs_ordering(sort_by: str, sort_direction: str) -> list[object]:
|
||||
column = CAMPAIGN_JOB_GRID_SORT_COLUMNS.get(sort_by)
|
||||
if column is None:
|
||||
raise HTTPException(
|
||||
status_code=status.HTTP_422_UNPROCESSABLE_CONTENT,
|
||||
detail="Unsupported Campaign job sort column",
|
||||
)
|
||||
primary = column.desc() if sort_direction == "desc" else column.asc()
|
||||
return [primary, CampaignJob.id.asc()]
|
||||
|
||||
|
||||
def _campaign_jobs_query_context(
|
||||
session: Session,
|
||||
principal: ApiPrincipal,
|
||||
*,
|
||||
campaign_id: str,
|
||||
version_id: str | None,
|
||||
send_status: list[str] | None,
|
||||
validation_status: list[str] | None,
|
||||
imap_status: list[str] | None,
|
||||
query_text: str | None,
|
||||
grid_filters: dict[str, str] | None = None,
|
||||
) -> tuple[Campaign, list[object], list[object], dict[str, object], set[str]]:
|
||||
_get_campaign_for_principal(session, campaign_id, principal)
|
||||
_require_permission(principal, "campaigns:recipient:read")
|
||||
campaign = _get_campaign_for_tenant(session, campaign_id, principal.tenant_id)
|
||||
base_filters: list[object] = [
|
||||
CampaignJob.campaign_id == campaign.id,
|
||||
CampaignJob.tenant_id == principal.tenant_id,
|
||||
]
|
||||
selected_version: CampaignVersion | None = None
|
||||
if version_id:
|
||||
version = _get_version_for_tenant(session, version_id, principal.tenant_id)
|
||||
if version.campaign_id != campaign.id:
|
||||
raise HTTPException(
|
||||
status_code=status.HTTP_404_NOT_FOUND,
|
||||
detail="Campaign version not found",
|
||||
)
|
||||
selected_version = version
|
||||
base_filters.append(CampaignJob.campaign_version_id == version.id)
|
||||
|
||||
review_metadata, reviewed_keys = _review_metadata(
|
||||
session, selected_version, base_filters
|
||||
)
|
||||
filtered = list(base_filters)
|
||||
if send_status:
|
||||
filtered.append(CampaignJob.send_status.in_(send_status))
|
||||
if validation_status:
|
||||
filtered.append(CampaignJob.validation_status.in_(validation_status))
|
||||
if imap_status:
|
||||
filtered.append(CampaignJob.imap_status.in_(imap_status))
|
||||
if query_text and query_text.strip():
|
||||
pattern = f"%{query_text.strip()}%"
|
||||
filtered.append(
|
||||
or_(
|
||||
CampaignJob.recipient_email.ilike(pattern),
|
||||
CampaignJob.subject.ilike(pattern),
|
||||
CampaignJob.entry_id.ilike(pattern),
|
||||
)
|
||||
)
|
||||
filtered.extend(_campaign_jobs_grid_filter_expressions(grid_filters))
|
||||
return campaign, base_filters, filtered, review_metadata, reviewed_keys
|
||||
|
||||
|
||||
def _campaign_jobs_page_response(
|
||||
session: Session,
|
||||
*,
|
||||
campaign_id: str,
|
||||
version_id: str | None,
|
||||
base_filters: list[object],
|
||||
filtered: list[object],
|
||||
reviewed_keys: set[str],
|
||||
review_metadata: dict[str, object],
|
||||
page: int,
|
||||
page_size: int,
|
||||
send_status: list[str] | None = None,
|
||||
validation_status: list[str] | None = None,
|
||||
imap_status: list[str] | None = None,
|
||||
query_text: str | None = None,
|
||||
grid_filters: dict[str, str] | None = None,
|
||||
sort_by: str = "number",
|
||||
sort_direction: str = "asc",
|
||||
cursor: str | None = None,
|
||||
changed_job_ids: set[str] | None = None,
|
||||
) -> CampaignJobsResponse:
|
||||
total_unfiltered, total = _campaign_jobs_page_counts(
|
||||
session, base_filters=base_filters, filtered=filtered
|
||||
)
|
||||
pages = (total + page_size - 1) // page_size if total else 0
|
||||
fingerprint = _campaign_jobs_cursor_fingerprint(
|
||||
campaign_id=campaign_id,
|
||||
version_id=version_id,
|
||||
page_size=page_size,
|
||||
send_status=send_status,
|
||||
validation_status=validation_status,
|
||||
imap_status=imap_status,
|
||||
query_text=query_text,
|
||||
grid_filters=grid_filters,
|
||||
sort_by=sort_by,
|
||||
sort_direction=sort_direction,
|
||||
)
|
||||
ordering = _campaign_jobs_ordering(sort_by, sort_direction)
|
||||
rows_plus_one, start_cursor = _campaign_jobs_page_rows(
|
||||
session,
|
||||
filtered=filtered,
|
||||
ordering=ordering,
|
||||
page=page,
|
||||
page_size=page_size,
|
||||
sort_by=sort_by,
|
||||
sort_direction=sort_direction,
|
||||
cursor=cursor,
|
||||
fingerprint=fingerprint,
|
||||
)
|
||||
jobs = rows_plus_one[:page_size]
|
||||
next_cursor = _campaign_jobs_next_cursor(
|
||||
jobs,
|
||||
rows_plus_one=rows_plus_one,
|
||||
page_size=page_size,
|
||||
sort_by=sort_by,
|
||||
sort_direction=sort_direction,
|
||||
changed_job_ids=changed_job_ids,
|
||||
fingerprint=fingerprint,
|
||||
)
|
||||
if changed_job_ids is not None:
|
||||
jobs = [job for job in jobs if job.id in changed_job_ids]
|
||||
return CampaignJobsResponse(
|
||||
jobs=[_job_summary_payload(job, reviewed_keys=reviewed_keys) for job in jobs],
|
||||
page=page,
|
||||
page_size=page_size,
|
||||
total=total,
|
||||
total_unfiltered=total_unfiltered,
|
||||
pages=pages,
|
||||
cursor=start_cursor,
|
||||
next_cursor=next_cursor,
|
||||
counts=_status_counts(session, base_filters),
|
||||
filtered_counts=_status_counts(session, filtered),
|
||||
review=review_metadata,
|
||||
)
|
||||
|
||||
|
||||
def _campaign_jobs_page_counts(
|
||||
session: Session,
|
||||
*,
|
||||
base_filters: list[object],
|
||||
filtered: list[object],
|
||||
) -> tuple[int, int]:
|
||||
total_unfiltered = int(
|
||||
session.query(func.count(CampaignJob.id)).filter(*base_filters).scalar() or 0
|
||||
)
|
||||
total = int(
|
||||
session.query(func.count(CampaignJob.id)).filter(*filtered).scalar() or 0
|
||||
)
|
||||
return total_unfiltered, total
|
||||
|
||||
|
||||
def _campaign_jobs_page_rows(
|
||||
session: Session,
|
||||
*,
|
||||
filtered: list[object],
|
||||
ordering: list[object],
|
||||
page: int,
|
||||
page_size: int,
|
||||
sort_by: str,
|
||||
sort_direction: str,
|
||||
cursor: str | None,
|
||||
fingerprint: str,
|
||||
) -> tuple[list[CampaignJob], str | None]:
|
||||
if cursor:
|
||||
page_query = _campaign_jobs_query_after_cursor(
|
||||
session,
|
||||
filtered=filtered,
|
||||
cursor=cursor,
|
||||
fingerprint=fingerprint,
|
||||
sort_by=sort_by,
|
||||
sort_direction=sort_direction,
|
||||
)
|
||||
return page_query.order_by(*ordering).limit(page_size + 1).all(), cursor
|
||||
|
||||
effective_offset = (page - 1) * page_size
|
||||
page_query = session.query(CampaignJob).filter(*filtered)
|
||||
start_cursor = _campaign_jobs_offset_cursor(
|
||||
page_query,
|
||||
ordering=ordering,
|
||||
effective_offset=effective_offset,
|
||||
sort_by=sort_by,
|
||||
sort_direction=sort_direction,
|
||||
fingerprint=fingerprint,
|
||||
)
|
||||
rows = (
|
||||
page_query.order_by(*ordering)
|
||||
.offset(effective_offset)
|
||||
.limit(page_size + 1)
|
||||
.all()
|
||||
)
|
||||
return rows, start_cursor
|
||||
|
||||
|
||||
def _campaign_jobs_query_after_cursor(
|
||||
session: Session,
|
||||
*,
|
||||
filtered: list[object],
|
||||
cursor: str,
|
||||
fingerprint: str,
|
||||
sort_by: str,
|
||||
sort_direction: str,
|
||||
):
|
||||
if sort_by != "number" or sort_direction != "asc":
|
||||
raise HTTPException(
|
||||
status_code=status.HTTP_400_BAD_REQUEST,
|
||||
detail="Campaign job cursors require number ascending order",
|
||||
)
|
||||
try:
|
||||
cursor_values = decode_keyset_cursor(
|
||||
CAMPAIGN_JOBS_CURSOR_SCOPE, cursor, fingerprint=fingerprint
|
||||
)
|
||||
if cursor_values is None:
|
||||
raise KeysetCursorError("Invalid pagination cursor")
|
||||
except KeysetCursorError as exc:
|
||||
raise HTTPException(
|
||||
status_code=status.HTTP_400_BAD_REQUEST, detail=str(exc)
|
||||
) from exc
|
||||
return (
|
||||
session.query(CampaignJob)
|
||||
.filter(*filtered)
|
||||
.filter(_campaign_jobs_cursor_condition(cursor_values))
|
||||
)
|
||||
|
||||
|
||||
def _campaign_jobs_offset_cursor(
|
||||
query,
|
||||
*,
|
||||
ordering: list[object],
|
||||
effective_offset: int,
|
||||
sort_by: str,
|
||||
sort_direction: str,
|
||||
fingerprint: str,
|
||||
) -> str | None:
|
||||
if effective_offset <= 0 or sort_by != "number" or sort_direction != "asc":
|
||||
return None
|
||||
previous_row = (
|
||||
query.order_by(*ordering).offset(effective_offset - 1).limit(1).first()
|
||||
)
|
||||
if previous_row is None:
|
||||
return None
|
||||
return _campaign_jobs_cursor_for_row(previous_row, fingerprint=fingerprint)
|
||||
|
||||
|
||||
def _campaign_jobs_next_cursor(
|
||||
jobs: list[CampaignJob],
|
||||
*,
|
||||
rows_plus_one: list[CampaignJob],
|
||||
page_size: int,
|
||||
sort_by: str,
|
||||
sort_direction: str,
|
||||
changed_job_ids: set[str] | None,
|
||||
fingerprint: str,
|
||||
) -> str | None:
|
||||
cursor_supported = sort_by == "number" and sort_direction == "asc"
|
||||
if (
|
||||
changed_job_ids is not None
|
||||
or not cursor_supported
|
||||
or len(rows_plus_one) <= page_size
|
||||
or not jobs
|
||||
):
|
||||
return None
|
||||
return _campaign_jobs_cursor_for_row(jobs[-1], fingerprint=fingerprint)
|
||||
|
||||
|
||||
def _campaign_jobs_cursor_fingerprint(
|
||||
*,
|
||||
campaign_id: str,
|
||||
version_id: str | None,
|
||||
page_size: int,
|
||||
send_status: list[str] | None,
|
||||
validation_status: list[str] | None,
|
||||
imap_status: list[str] | None,
|
||||
query_text: str | None,
|
||||
grid_filters: dict[str, str] | None,
|
||||
sort_by: str,
|
||||
sort_direction: str,
|
||||
) -> str:
|
||||
return keyset_query_fingerprint(
|
||||
CAMPAIGN_JOBS_CURSOR_SCOPE,
|
||||
{
|
||||
"campaign_id": campaign_id,
|
||||
"version_id": version_id or "",
|
||||
"page_size": page_size,
|
||||
"send_status": sorted(send_status or []),
|
||||
"validation_status": sorted(validation_status or []),
|
||||
"imap_status": sorted(imap_status or []),
|
||||
"query": (query_text or "").strip(),
|
||||
"grid_filters": sorted((grid_filters or {}).items()),
|
||||
"order": [f"{sort_by}:{sort_direction}", "id:asc"],
|
||||
},
|
||||
)
|
||||
|
||||
|
||||
def _campaign_jobs_cursor_for_row(row: CampaignJob, *, fingerprint: str) -> str:
|
||||
return encode_keyset_cursor(
|
||||
CAMPAIGN_JOBS_CURSOR_SCOPE,
|
||||
fingerprint=fingerprint,
|
||||
values={"id": row.id, "entry_index": row.entry_index},
|
||||
)
|
||||
|
||||
|
||||
def _campaign_jobs_cursor_condition(cursor_values: dict[str, object]):
|
||||
cursor_id = cursor_values.get("id")
|
||||
cursor_index = cursor_values.get("entry_index")
|
||||
if not isinstance(cursor_id, str) or not cursor_id:
|
||||
raise KeysetCursorError("Invalid pagination cursor")
|
||||
try:
|
||||
entry_index = int(cursor_index)
|
||||
except (TypeError, ValueError) as exc:
|
||||
raise KeysetCursorError("Invalid pagination cursor") from exc
|
||||
return or_(
|
||||
CampaignJob.entry_index > entry_index,
|
||||
and_(CampaignJob.entry_index == entry_index, CampaignJob.id > cursor_id),
|
||||
)
|
||||
|
||||
|
||||
def _campaign_jobs_delta_watermark(session: Session, tenant_id: str) -> str:
|
||||
return encode_sequence_watermark(
|
||||
max_sequence_id(
|
||||
session,
|
||||
tenant_id=tenant_id,
|
||||
module_id=CAMPAIGNS_MODULE_ID,
|
||||
collections=(CAMPAIGN_JOBS_COLLECTION,),
|
||||
)
|
||||
)
|
||||
|
||||
|
||||
class CampaignJobsQuery:
|
||||
def __init__(
|
||||
self,
|
||||
version_id: str | None = None,
|
||||
page: int = Query(default=1, ge=1),
|
||||
page_size: int = Query(default=50, ge=1, le=200),
|
||||
cursor: str | None = Query(default=None),
|
||||
send_status: list[str] | None = Query(default=None),
|
||||
validation_status: list[str] | None = Query(default=None),
|
||||
imap_status: list[str] | None = Query(default=None),
|
||||
query_text: str | None = Query(default=None, alias="q", max_length=200),
|
||||
sort_by: Literal[
|
||||
"number",
|
||||
"recipient",
|
||||
"subject",
|
||||
"validation",
|
||||
"queue",
|
||||
"send",
|
||||
"postbox",
|
||||
"imap",
|
||||
"attempts",
|
||||
"updated",
|
||||
] = Query(default="number"),
|
||||
sort_direction: Literal["asc", "desc"] = Query(default="asc"),
|
||||
filter_recipient: str | None = Query(default=None, max_length=500),
|
||||
filter_subject: str | None = Query(default=None, max_length=1000),
|
||||
filter_validation: str | None = Query(default=None, max_length=1000),
|
||||
filter_queue: str | None = Query(default=None, max_length=1000),
|
||||
filter_send: str | None = Query(default=None, max_length=1000),
|
||||
filter_postbox: str | None = Query(default=None, max_length=1000),
|
||||
filter_imap: str | None = Query(default=None, max_length=1000),
|
||||
filter_attempts: str | None = Query(default=None, max_length=100),
|
||||
filter_evidence: str | None = Query(default=None, max_length=500),
|
||||
) -> None:
|
||||
self.version_id = version_id
|
||||
self.page = page
|
||||
self.page_size = page_size
|
||||
self.cursor = cursor
|
||||
self.send_status = send_status
|
||||
self.validation_status = validation_status
|
||||
self.imap_status = imap_status
|
||||
self.query_text = query_text
|
||||
self.sort_by = sort_by
|
||||
self.sort_direction = sort_direction
|
||||
self.grid_filters = {
|
||||
column_id: value
|
||||
for column_id, value in {
|
||||
"recipient": filter_recipient,
|
||||
"subject": filter_subject,
|
||||
"validation": filter_validation,
|
||||
"queue": filter_queue,
|
||||
"send": filter_send,
|
||||
"postbox": filter_postbox,
|
||||
"imap": filter_imap,
|
||||
"attempts": filter_attempts,
|
||||
"evidence": filter_evidence,
|
||||
}.items()
|
||||
if value is not None and value.strip()
|
||||
}
|
||||
Reference in New Issue
Block a user