from __future__ import annotations from datetime import UTC, datetime, timedelta from zoneinfo import ZoneInfo, ZoneInfoNotFoundError from fastapi import APIRouter, Depends, HTTPException, status from sqlalchemy.orm import Session from govoplan_campaign.backend.campaign.scheduling import ( campaign_schedule_source_snapshot, canonical_configuration_hash, validate_autonomous_schedule_source, ) from govoplan_campaign.backend.db.models import ( CampaignSchedule, CampaignScheduleOccurrence, CampaignShare, CampaignVersion, ) from govoplan_campaign.backend.route_support import ( _get_campaign_for_principal, _require_permission, ) from govoplan_campaign.backend.schemas import ( CampaignScheduleCreateRequest, CampaignScheduleListResponse, CampaignScheduleOccurrenceResponse, CampaignScheduleResponse, CampaignScheduleStateRequest, ) from govoplan_core.auth import ApiPrincipal, require_scope from govoplan_core.audit.logging import audit_from_principal from govoplan_core.db.session import get_session router = APIRouter(prefix="/campaigns", tags=["campaign-schedules"]) @router.get( "/{campaign_id}/schedules", response_model=CampaignScheduleListResponse, ) def list_campaign_schedules( campaign_id: str, session: Session = Depends(get_session), principal: ApiPrincipal = Depends(require_scope("campaigns:campaign:read")), ): _get_campaign_for_principal(session, campaign_id, principal) schedules = ( session.query(CampaignSchedule) .filter( CampaignSchedule.tenant_id == principal.tenant_id, CampaignSchedule.campaign_id == campaign_id, ) .order_by(CampaignSchedule.created_at.desc(), CampaignSchedule.id.asc()) .all() ) occurrences = _occurrences_by_schedule(session, schedules) return CampaignScheduleListResponse( items=[ _schedule_response(item, occurrences.get(item.id, [])) for item in schedules ] ) @router.post( "/{campaign_id}/schedules", response_model=CampaignScheduleResponse, status_code=status.HTTP_201_CREATED, ) def create_campaign_schedule( campaign_id: str, payload: CampaignScheduleCreateRequest, session: Session = Depends(get_session), principal: ApiPrincipal = Depends(require_scope("campaigns:campaign:schedule")), ): campaign = _get_campaign_for_principal(session, campaign_id, principal) _require_permission(principal, "campaigns:campaign:copy") if payload.include_recipients: _require_permission(principal, "campaigns:recipient:read") if payload.include_shares: _require_permission(principal, "campaigns:campaign:share") if payload.delivery_mode == "autonomous": _require_permission(principal, "campaigns:campaign:queue") _require_permission(principal, "campaigns:campaign:send") _require_permission(principal, "mail:profile:use") source_version = ( session.query(CampaignVersion) .filter( CampaignVersion.id == payload.source_version_id, CampaignVersion.campaign_id == campaign.id, ) .one_or_none() ) if source_version is None: raise HTTPException(status_code=404, detail="Campaign version not found") starts_at = payload.starts_at.astimezone(UTC) if starts_at < datetime.now(UTC) - timedelta(minutes=5): raise HTTPException( status_code=status.HTTP_422_UNPROCESSABLE_CONTENT, detail="Campaign schedules cannot start in the past.", ) try: ZoneInfo(payload.timezone) except ZoneInfoNotFoundError as exc: raise HTTPException( status_code=status.HTTP_422_UNPROCESSABLE_CONTENT, detail="Unknown campaign schedule timezone.", ) from exc source_shares = ( session.query(CampaignShare) .filter( CampaignShare.tenant_id == principal.tenant_id, CampaignShare.campaign_id == campaign.id, CampaignShare.revoked_at.is_(None), ) .order_by(CampaignShare.id.asc()) .all() if payload.include_shares else [] ) snapshot = campaign_schedule_source_snapshot( configuration=source_version.raw_json, campaign_settings=campaign.settings or {}, mail_profile_policy=campaign.mail_profile_policy or {}, shares=[ { "target_type": item.target_type, "target_id": item.target_id, "permission": item.permission, } for item in source_shares ], ) autonomous_evidence: dict[str, object] | None = None if payload.delivery_mode == "autonomous": try: autonomous_evidence = validate_autonomous_schedule_source( session, campaign=campaign, version=source_version, ) except (RuntimeError, ValueError) as exc: raise HTTPException( status_code=status.HTTP_409_CONFLICT, detail=str(exc), ) from exc snapshot["autonomous_delivery"] = autonomous_evidence schedule = CampaignSchedule( tenant_id=principal.tenant_id, campaign_id=campaign.id, source_version_id=source_version.id, created_by_user_id=principal.user.id, name=payload.name.strip(), delivery_mode=payload.delivery_mode, recurrence_kind=payload.recurrence_kind, interval_count=payload.interval_count, timezone=payload.timezone, starts_at=starts_at, next_fire_at=starts_at, ends_at=payload.ends_at.astimezone(UTC) if payload.ends_at else None, max_occurrences=payload.max_occurrences, copy_options={ "include_recipients": payload.include_recipients, "include_files": payload.include_files, "include_shares": payload.include_shares, "include_policies": payload.include_policies, "include_mail_profile": payload.include_mail_profile, }, source_snapshot=snapshot, source_snapshot_hash=canonical_configuration_hash(snapshot), approved_execution_snapshot_hash=( str(autonomous_evidence["execution_snapshot_hash"]) if autonomous_evidence is not None else None ), source_base_path=source_version.source_base_path, ) session.add(schedule) session.flush() audit_from_principal( session, principal, action="campaign.schedule.created", object_type="campaign_schedule", object_id=schedule.id, details={ "campaign_id": campaign.id, "source_version_id": source_version.id, "recurrence_kind": schedule.recurrence_kind, "interval_count": schedule.interval_count, "starts_at": schedule.starts_at.isoformat(), "ends_at": schedule.ends_at.isoformat() if schedule.ends_at else None, "max_occurrences": schedule.max_occurrences, "delivery_mode": schedule.delivery_mode, "delivery_started": False, "autonomous_delivery_opted_in": ( schedule.delivery_mode == "autonomous" ), "approved_execution_snapshot_hash": ( schedule.approved_execution_snapshot_hash ), "approval_request_id": ( autonomous_evidence.get("approval_request_id") if autonomous_evidence is not None else None ), }, commit=True, ) session.refresh(schedule) return _schedule_response(schedule, []) @router.patch( "/{campaign_id}/schedules/{schedule_id}", response_model=CampaignScheduleResponse, ) def set_campaign_schedule_state( campaign_id: str, schedule_id: str, payload: CampaignScheduleStateRequest, session: Session = Depends(get_session), principal: ApiPrincipal = Depends(require_scope("campaigns:campaign:schedule")), ): _get_campaign_for_principal(session, campaign_id, principal) schedule = _schedule_for_campaign( session, tenant_id=principal.tenant_id, campaign_id=campaign_id, schedule_id=schedule_id, for_update=True, ) if schedule.resource_revision != payload.base_revision: raise HTTPException( status_code=status.HTTP_409_CONFLICT, detail="Campaign schedule changed. Reload it before changing its state.", ) if payload.active and schedule.next_fire_at is None: raise HTTPException( status_code=status.HTTP_409_CONFLICT, detail="A completed campaign schedule cannot be resumed.", ) if payload.active: unresolved = ( session.query(CampaignScheduleOccurrence.id) .filter( CampaignScheduleOccurrence.schedule_id == schedule.id, CampaignScheduleOccurrence.status == "uncertain", ) .first() ) if unresolved is not None: raise HTTPException( status_code=status.HTTP_409_CONFLICT, detail=( "Reconcile the autonomous delivery outcome in Mail before " "resuming this schedule." ), ) schedule.active = payload.active schedule.last_error = None if payload.active else schedule.last_error schedule.resource_revision += 1 session.add(schedule) audit_from_principal( session, principal, action=( "campaign.schedule.resumed" if payload.active else "campaign.schedule.paused" ), object_type="campaign_schedule", object_id=schedule.id, details={"campaign_id": campaign_id}, commit=True, ) session.refresh(schedule) occurrences = _occurrences_by_schedule(session, [schedule]).get(schedule.id, []) return _schedule_response(schedule, occurrences) def _schedule_for_campaign( session: Session, *, tenant_id: str, campaign_id: str, schedule_id: str, for_update: bool = False, ) -> CampaignSchedule: query = session.query(CampaignSchedule) if for_update: query = query.with_for_update() schedule = ( query .filter( CampaignSchedule.id == schedule_id, CampaignSchedule.tenant_id == tenant_id, CampaignSchedule.campaign_id == campaign_id, ) .one_or_none() ) if schedule is None: raise HTTPException(status_code=404, detail="Campaign schedule not found") return schedule def _occurrences_by_schedule( session: Session, schedules: list[CampaignSchedule], ) -> dict[str, list[CampaignScheduleOccurrence]]: ids = [item.id for item in schedules] if not ids: return {} rows = ( session.query(CampaignScheduleOccurrence) .filter(CampaignScheduleOccurrence.schedule_id.in_(ids)) .order_by( CampaignScheduleOccurrence.scheduled_for.desc(), CampaignScheduleOccurrence.id.asc(), ) .all() ) grouped: dict[str, list[CampaignScheduleOccurrence]] = {} for row in rows: grouped.setdefault(row.schedule_id, []).append(row) return grouped def _schedule_response( schedule: CampaignSchedule, occurrences: list[CampaignScheduleOccurrence], ) -> CampaignScheduleResponse: response = CampaignScheduleResponse.model_validate(schedule) return response.model_copy( update={ "occurrences": [ CampaignScheduleOccurrenceResponse.model_validate(item) for item in occurrences ] } )