from __future__ import annotations from fastapi import APIRouter, Depends, HTTPException, Query, status from sqlalchemy.orm import Session from govoplan_campaign.backend.schemas import ( CampaignJobsResponse, CampaignJobsDeltaResponse, CampaignJobDetailResponse, CampaignJobDiagnosticsResponse, ) from govoplan_core.auth import ApiPrincipal, require_scope from govoplan_core.core.change_sequence import ( decode_sequence_watermark, encode_sequence_watermark, sequence_entries_since, sequence_watermark_is_expired, ) from govoplan_campaign.backend.change_tracking import ( CAMPAIGNS_MODULE_ID, CAMPAIGN_JOBS_COLLECTION, ) from govoplan_campaign.backend.db.models import ( CampaignJob, CampaignMessageAction, CampaignMessageActionAttempt, ImapAppendAttempt, PostboxDeliveryAttempt, PrintOutputAttempt, SendAttempt, ) from govoplan_campaign.backend.integrations import postbox_integration from govoplan_core.db.session import get_session from govoplan_campaign.backend.route_support import ( _get_campaign_for_principal, _get_campaign_for_tenant, _require_permission, job_attempt_rows as _job_attempt_rows, ) from govoplan_campaign.backend.services.job_queries import ( CampaignJobsQuery, _campaign_jobs_delta_watermark, _campaign_jobs_page_response, _campaign_jobs_query_context, _job_attempts_payload, _calendar_invitations_for_jobs, _job_detail_payload, _job_diagnostics_payload, ) router = APIRouter(prefix="/campaigns", tags=["campaigns"]) def _postbox_receipts_for_attempts( session: Session, *, tenant_id: str, attempts: list[PostboxDeliveryAttempt], ): integration = postbox_integration() if not integration.receipt_evidence_available: return None delivery_ids = [ attempt.provider_delivery_id for attempt in attempts if attempt.provider_delivery_id ] return integration.delivery_receipt_summaries( session, tenant_id=tenant_id, delivery_ids=delivery_ids, ) @router.get("/{campaign_id}/jobs", response_model=CampaignJobsResponse) def list_jobs( campaign_id: str, filters: CampaignJobsQuery = Depends(CampaignJobsQuery), session: Session = Depends(get_session), principal: ApiPrincipal = Depends(require_scope("campaigns:campaign:read")), ): """Return a lightweight, paginated job list with server-side filters. Complete recipients, attachment metadata, issues and attempt history are available from the separate job-detail endpoint. """ _campaign, base_filters, filtered, review_metadata, reviewed_keys = ( _campaign_jobs_query_context( session, principal, campaign_id=campaign_id, version_id=filters.version_id, send_status=filters.send_status, 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, campaign_id=campaign_id, version_id=filters.version_id, base_filters=base_filters, filtered=filtered, reviewed_keys=reviewed_keys, review_metadata=review_metadata, page=filters.page, page_size=filters.page_size, send_status=filters.send_status, 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, ) def _campaign_jobs_full_delta_response( session: Session, *, principal: ApiPrincipal, campaign_id: str, version_id: str | None, page: int, 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, cursor: str | None = None, ) -> CampaignJobsDeltaResponse: _campaign, base_filters, filtered, review_metadata, reviewed_keys = ( _campaign_jobs_query_context( session, principal, campaign_id=campaign_id, version_id=version_id, send_status=send_status, validation_status=validation_status, imap_status=imap_status, query_text=query_text, grid_filters=grid_filters, ) ) payload = _campaign_jobs_page_response( session, campaign_id=campaign_id, version_id=version_id, base_filters=base_filters, filtered=filtered, reviewed_keys=reviewed_keys, review_metadata=review_metadata, page=page, 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, cursor=cursor, ) return CampaignJobsDeltaResponse( **payload.model_dump(), deleted=[], watermark=_campaign_jobs_delta_watermark(session, principal.tenant_id), has_more=False, full=True, ) def _job_filter_membership_can_shift( *, 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, ) -> bool: 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) def list_jobs_delta( campaign_id: str, filters: CampaignJobsQuery = Depends(CampaignJobsQuery), since: str | None = None, limit: int = Query(default=500, ge=1, le=1000), session: Session = Depends(get_session), principal: ApiPrincipal = Depends(require_scope("campaigns:campaign:read")), ): if since is None: return _campaign_jobs_full_delta_response( session, principal=principal, campaign_id=campaign_id, version_id=filters.version_id, page=filters.page, page_size=filters.page_size, send_status=filters.send_status, 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, ) campaign, base_filters, filtered, review_metadata, reviewed_keys = ( _campaign_jobs_query_context( session, principal, campaign_id=campaign_id, version_id=filters.version_id, send_status=filters.send_status, 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) except ValueError as exc: raise HTTPException( status_code=status.HTTP_400_BAD_REQUEST, detail=str(exc) ) from exc if sequence_watermark_is_expired( session, since=since_sequence, tenant_id=principal.tenant_id, module_id=CAMPAIGNS_MODULE_ID, collections=(CAMPAIGN_JOBS_COLLECTION,), ): return _campaign_jobs_full_delta_response( session, principal=principal, campaign_id=campaign_id, version_id=filters.version_id, page=filters.page, page_size=filters.page_size, send_status=filters.send_status, 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, ) entries_plus_one = sequence_entries_since( session, since=since_sequence, tenant_id=principal.tenant_id, module_id=CAMPAIGNS_MODULE_ID, collections=(CAMPAIGN_JOBS_COLLECTION,), limit=limit + 1, ) has_more = len(entries_plus_one) > limit entries = entries_plus_one[:limit] relevant_entries = [ entry for entry in entries if (entry.payload or {}).get("campaign_id") == campaign.id and ( not filters.version_id or (entry.payload or {}).get("version_id") == filters.version_id ) ] if relevant_entries and ( _job_filter_membership_can_shift( send_status=filters.send_status, 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) ): return _campaign_jobs_full_delta_response( session, principal=principal, campaign_id=campaign_id, version_id=filters.version_id, page=filters.page, page_size=filters.page_size, send_status=filters.send_status, 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 = { entry.resource_id for entry in relevant_entries if entry.resource_type == "campaign_job" and entry.operation != "deleted" } payload = _campaign_jobs_page_response( session, campaign_id=campaign_id, version_id=filters.version_id, base_filters=base_filters, filtered=filtered, reviewed_keys=reviewed_keys, review_metadata=review_metadata, page=filters.page, page_size=filters.page_size, send_status=filters.send_status, 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, ) watermark = ( encode_sequence_watermark(entries[-1].id) if has_more and entries else _campaign_jobs_delta_watermark(session, principal.tenant_id) ) return CampaignJobsDeltaResponse( **payload.model_dump(), deleted=[], watermark=watermark, has_more=has_more, full=False, ) @router.get("/{campaign_id}/jobs/{job_id}", response_model=CampaignJobDetailResponse) def get_job_detail( campaign_id: str, job_id: str, session: Session = Depends(get_session), principal: ApiPrincipal = Depends(require_scope("campaigns:campaign:read")), ): _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) job = session.get(CampaignJob, job_id) if ( not job or job.campaign_id != campaign.id or job.tenant_id != principal.tenant_id ): raise HTTPException( status_code=status.HTTP_404_NOT_FOUND, detail="Campaign job not found" ) send_attempts = _job_attempt_rows( session.query(SendAttempt) .filter(SendAttempt.job_id == job.id) .order_by(SendAttempt.attempt_number.asc()), label="SMTP attempts for this campaign job", ) imap_attempts = _job_attempt_rows( session.query(ImapAppendAttempt) .filter(ImapAppendAttempt.job_id == job.id) .order_by(ImapAppendAttempt.attempt_number.asc()), label="IMAP attempts for this campaign job", ) postbox_attempts = _job_attempt_rows( session.query(PostboxDeliveryAttempt) .filter(PostboxDeliveryAttempt.job_id == job.id) .order_by( PostboxDeliveryAttempt.target_index.asc(), PostboxDeliveryAttempt.attempt_number.asc(), ), label="Postbox attempts for this campaign job", ) print_attempts = _job_attempt_rows( session.query(PrintOutputAttempt) .filter(PrintOutputAttempt.job_id == job.id) .order_by(PrintOutputAttempt.attempt_number.asc()), label="Printable output attempts for this campaign job", ) message_actions = _job_attempt_rows( session.query(CampaignMessageAction) .filter(CampaignMessageAction.job_id == job.id) .order_by(CampaignMessageAction.created_at.asc()), label="Single-message actions for this campaign job", ) action_ids = [action.id for action in message_actions] message_action_attempts = ( _job_attempt_rows( session.query(CampaignMessageActionAttempt) .filter(CampaignMessageActionAttempt.action_id.in_(action_ids)) .order_by(CampaignMessageActionAttempt.started_at.asc()), label="Single-message action attempts for this campaign job", ) if action_ids else [] ) return CampaignJobDetailResponse( job=_job_detail_payload( job, calendar_invitation=_calendar_invitations_for_jobs( session, [job], ).get(job.id), ), attempts=_job_attempts_payload( send_attempts, imap_attempts, postbox_attempts=postbox_attempts, print_attempts=print_attempts, message_actions=message_actions, message_action_attempts=message_action_attempts, postbox_receipts=_postbox_receipts_for_attempts( session, tenant_id=principal.tenant_id, attempts=postbox_attempts, ), ), ) @router.get( "/{campaign_id}/jobs/{job_id}/diagnostics", response_model=CampaignJobDiagnosticsResponse, ) def get_job_diagnostics( campaign_id: str, job_id: str, session: Session = Depends(get_session), principal: ApiPrincipal = Depends(require_scope("campaigns:diagnostic:read")), ): """Return infrastructure details only to campaign operators/admins.""" _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) job = session.get(CampaignJob, job_id) if ( not job or job.campaign_id != campaign.id or job.tenant_id != principal.tenant_id ): raise HTTPException( status_code=status.HTTP_404_NOT_FOUND, detail="Campaign job not found" ) send_attempts = _job_attempt_rows( session.query(SendAttempt) .filter(SendAttempt.job_id == job.id) .order_by(SendAttempt.attempt_number.asc()), label="SMTP diagnostics for this campaign job", ) imap_attempts = _job_attempt_rows( session.query(ImapAppendAttempt) .filter(ImapAppendAttempt.job_id == job.id) .order_by(ImapAppendAttempt.attempt_number.asc()), label="IMAP diagnostics for this campaign job", ) postbox_attempts = _job_attempt_rows( session.query(PostboxDeliveryAttempt) .filter(PostboxDeliveryAttempt.job_id == job.id) .order_by( PostboxDeliveryAttempt.target_index.asc(), PostboxDeliveryAttempt.attempt_number.asc(), ), label="Postbox diagnostics for this campaign job", ) print_attempts = _job_attempt_rows( session.query(PrintOutputAttempt) .filter(PrintOutputAttempt.job_id == job.id) .order_by(PrintOutputAttempt.attempt_number.asc()), label="Printable output diagnostics for this campaign job", ) message_actions = _job_attempt_rows( session.query(CampaignMessageAction) .filter(CampaignMessageAction.job_id == job.id) .order_by(CampaignMessageAction.created_at.asc()), label="Single-message action diagnostics for this campaign job", ) action_ids = [action.id for action in message_actions] message_action_attempts = ( _job_attempt_rows( session.query(CampaignMessageActionAttempt) .filter(CampaignMessageActionAttempt.action_id.in_(action_ids)) .order_by(CampaignMessageActionAttempt.started_at.asc()), label="Single-message action-attempt diagnostics for this campaign job", ) if action_ids else [] ) return _job_diagnostics_payload( job, send_attempts, imap_attempts, postbox_attempts=postbox_attempts, print_attempts=print_attempts, message_actions=message_actions, message_action_attempts=message_action_attempts, postbox_receipts=_postbox_receipts_for_attempts( session, tenant_id=principal.tenant_id, attempts=postbox_attempts, ), )