fix(campaign): query reports before pagination

This commit is contained in:
2026-07-22 08:31:43 +02:00
parent 60efd1cb5d
commit aa4ec66b7b
9 changed files with 603 additions and 131 deletions

View File

@@ -2,9 +2,10 @@ from __future__ import annotations
import copy
import dataclasses
import json
import logging
from collections.abc import Callable
from typing import Any
from typing import Any, Literal
from fastapi import APIRouter, Depends, HTTPException, Query, Response, status
from sqlalchemy import and_, exists, func, or_
@@ -99,6 +100,7 @@ from govoplan_campaign.backend.db.models import (
JobImapStatus,
JobQueueStatus,
JobSendStatus,
JobValidationStatus,
RecipientImportMappingProfile,
SendAttempt,
)
@@ -2252,6 +2254,125 @@ def _status_counts(session: Session, filters: list[object]) -> dict[str, dict[st
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,
"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}),
"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,
@@ -2262,6 +2383,7 @@ def _campaign_jobs_query_context(
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")
@@ -2290,6 +2412,7 @@ def _campaign_jobs_query_context(
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
@@ -2308,6 +2431,9 @@ def _campaign_jobs_page_response(
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:
@@ -2322,11 +2448,20 @@ def _campaign_jobs_page_response(
validation_status=validation_status,
imap_status=imap_status,
query_text=query_text,
grid_filters=grid_filters,
sort_by=sort_by,
sort_direction=sort_direction,
)
ordered_query = session.query(CampaignJob).filter(*filtered).order_by(CampaignJob.entry_index.asc(), CampaignJob.id.asc())
ordering = _campaign_jobs_ordering(sort_by, sort_direction)
ordered_query = session.query(CampaignJob).filter(*filtered).order_by(*ordering)
start_cursor: str | None = None
effective_offset = 0
if cursor:
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:
@@ -2338,14 +2473,14 @@ def _campaign_jobs_page_response(
else:
effective_offset = (page - 1) * page_size
page_query = session.query(CampaignJob).filter(*filtered)
if effective_offset > 0:
if effective_offset > 0 and sort_by == "number" and sort_direction == "asc":
previous_row = ordered_query.offset(effective_offset - 1).limit(1).first()
if previous_row is not None:
start_cursor = _campaign_jobs_cursor_for_row(previous_row, fingerprint=fingerprint)
rows_plus_one = (
page_query
.order_by(CampaignJob.entry_index.asc(), CampaignJob.id.asc())
.order_by(*ordering)
.offset(effective_offset)
.limit(page_size + 1)
.all()
@@ -2353,7 +2488,12 @@ def _campaign_jobs_page_response(
jobs = rows_plus_one[:page_size]
next_cursor = (
_campaign_jobs_cursor_for_row(jobs[-1], fingerprint=fingerprint)
if changed_job_ids is None and len(rows_plus_one) > page_size and jobs else None
if changed_job_ids is None
and sort_by == "number"
and sort_direction == "asc"
and len(rows_plus_one) > page_size
and jobs
else None
)
if changed_job_ids is not None:
jobs = [job for job in jobs if job.id in changed_job_ids]
@@ -2381,6 +2521,9 @@ def _campaign_jobs_cursor_fingerprint(
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,
@@ -2392,7 +2535,8 @@ def _campaign_jobs_cursor_fingerprint(
"validation_status": sorted(validation_status or []),
"imap_status": sorted(imap_status or []),
"query": (query_text or "").strip(),
"order": ["entry_index:asc", "id:asc"],
"grid_filters": sorted((grid_filters or {}).items()),
"order": [f"{sort_by}:{sort_direction}", "id:asc"],
},
)
@@ -2439,6 +2583,16 @@ class CampaignJobsQuery:
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", "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_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
@@ -2448,6 +2602,22 @@ class CampaignJobsQuery:
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,
"imap": filter_imap,
"attempts": filter_attempts,
"evidence": filter_evidence,
}.items()
if value is not None and value.strip()
}
@router.get("/{campaign_id}/jobs", response_model=CampaignJobsResponse)
@@ -2472,6 +2642,7 @@ def list_jobs(
validation_status=filters.validation_status,
imap_status=filters.imap_status,
query_text=filters.query_text,
grid_filters=filters.grid_filters,
)
return _campaign_jobs_page_response(
session,
@@ -2487,6 +2658,9 @@ def list_jobs(
validation_status=filters.validation_status,
imap_status=filters.imap_status,
query_text=filters.query_text,
grid_filters=filters.grid_filters,
sort_by=filters.sort_by,
sort_direction=filters.sort_direction,
cursor=filters.cursor,
)
@@ -2503,6 +2677,9 @@ def _campaign_jobs_full_delta_response(
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,
cursor: str | None = None,
) -> CampaignJobsDeltaResponse:
_campaign, base_filters, filtered, review_metadata, reviewed_keys = _campaign_jobs_query_context(
@@ -2514,6 +2691,7 @@ def _campaign_jobs_full_delta_response(
validation_status=validation_status,
imap_status=imap_status,
query_text=query_text,
grid_filters=grid_filters,
)
payload = _campaign_jobs_page_response(
session,
@@ -2529,6 +2707,9 @@ def _campaign_jobs_full_delta_response(
validation_status=validation_status,
imap_status=imap_status,
query_text=query_text,
grid_filters=grid_filters,
sort_by=sort_by,
sort_direction=sort_direction,
cursor=cursor,
)
return CampaignJobsDeltaResponse(
@@ -2546,8 +2727,19 @@ def _job_filter_membership_can_shift(
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,
) -> bool:
return bool(send_status or validation_status or imap_status or (query_text and query_text.strip()))
return bool(
send_status
or validation_status
or imap_status
or (query_text and query_text.strip())
or grid_filters
or sort_by != "number"
or sort_direction != "asc"
)
@router.get("/{campaign_id}/jobs/delta", response_model=CampaignJobsDeltaResponse)
@@ -2571,6 +2763,9 @@ def list_jobs_delta(
validation_status=filters.validation_status,
imap_status=filters.imap_status,
query_text=filters.query_text,
grid_filters=filters.grid_filters,
sort_by=filters.sort_by,
sort_direction=filters.sort_direction,
cursor=filters.cursor,
)
@@ -2583,6 +2778,7 @@ def list_jobs_delta(
validation_status=filters.validation_status,
imap_status=filters.imap_status,
query_text=filters.query_text,
grid_filters=filters.grid_filters,
)
try:
since_sequence = decode_sequence_watermark(since)
@@ -2606,6 +2802,9 @@ def list_jobs_delta(
validation_status=filters.validation_status,
imap_status=filters.imap_status,
query_text=filters.query_text,
grid_filters=filters.grid_filters,
sort_by=filters.sort_by,
sort_direction=filters.sort_direction,
cursor=filters.cursor,
)
@@ -2632,6 +2831,9 @@ def list_jobs_delta(
validation_status=filters.validation_status,
imap_status=filters.imap_status,
query_text=filters.query_text,
grid_filters=filters.grid_filters,
sort_by=filters.sort_by,
sort_direction=filters.sort_direction,
)
or any(entry.operation in {"created", "deleted"} for entry in relevant_entries)
):
@@ -2646,6 +2848,9 @@ def list_jobs_delta(
validation_status=filters.validation_status,
imap_status=filters.imap_status,
query_text=filters.query_text,
grid_filters=filters.grid_filters,
sort_by=filters.sort_by,
sort_direction=filters.sort_direction,
cursor=filters.cursor,
)
@@ -2668,6 +2873,9 @@ def list_jobs_delta(
validation_status=filters.validation_status,
imap_status=filters.imap_status,
query_text=filters.query_text,
grid_filters=filters.grid_filters,
sort_by=filters.sort_by,
sort_direction=filters.sort_direction,
cursor=filters.cursor,
changed_job_ids=changed_job_ids,
)