use std::sync::Arc;
use uuid::Uuid;
use super::event_error::EventResult;
use super::sms_port::{EventSmsQueue, RenderedSms};
use super::template_port::{EventMailQueue, EventTemplateRenderer, RenderContext};
use crate::infrastructure::persistence::event_command_repository::EventCommandRepository;
use crate::infrastructure::persistence::scheduler_repository::{SchedulerRepository, SchedulerRow};
pub const DEFAULT_MAIL_BATCH: i64 = 50;
pub const DEFAULT_CRON_LIMIT: i64 = 1000;
pub const OUTCOME_QUEUED: &str = "queued";
pub const OUTCOME_DROPPED_WINDOW_CLOSED: &str = "dropped_window_closed";
#[derive(Debug, Clone, Default, serde::Serialize)]
pub struct SchedulerRunSummary {
pub schedulers_claimed: i64,
pub receipts_materialized: i64,
pub receipts_queued: i64,
pub sms_queued: i64,
pub receipts_dropped_window_closed: i64,
pub typed_failures_recorded: i64,
pub cancellation_receipts_deleted: i64,
pub schedulers_rearmed: i64,
pub schedulers_completed: i64,
}
pub struct SchedulerService {
schedulers: SchedulerRepository,
events: EventCommandRepository,
renderer: Arc<dyn EventTemplateRenderer>,
queue: Arc<dyn EventMailQueue>,
sms_queue: Arc<dyn EventSmsQueue>,
batch: i64,
cron_limit: i64,
}
impl SchedulerService {
pub fn new(
schedulers: SchedulerRepository,
events: EventCommandRepository,
renderer: Arc<dyn EventTemplateRenderer>,
queue: Arc<dyn EventMailQueue>,
sms_queue: Arc<dyn EventSmsQueue>,
) -> Self {
Self {
schedulers,
events,
renderer,
queue,
sms_queue,
batch: DEFAULT_MAIL_BATCH,
cron_limit: DEFAULT_CRON_LIMIT,
}
}
pub fn with_caps(mut self, batch: i64, cron_limit: i64) -> Self {
self.batch = batch.max(1);
self.cron_limit = cron_limit.max(1);
self
}
pub async fn run_due_schedulers(&self) -> EventResult<SchedulerRunSummary> {
let mut summary = SchedulerRunSummary::default();
summary.cancellation_receipts_deleted = self
.schedulers
.propagate_cancellations(self.cron_limit)
.await?;
let claimed = self.schedulers.claim_due(self.cron_limit).await?;
summary.schedulers_claimed = claimed.len() as i64;
for scheduler in &claimed {
self.run_one(scheduler, &mut summary).await?;
}
Ok(summary)
}
async fn run_one(&self, scheduler: &SchedulerRow, summary: &mut SchedulerRunSummary) -> EventResult<()> {
if scheduler.template_ref.is_none() {
self.schedulers
.record_failure(scheduler.id, "template_unresolved")
.await?;
summary.typed_failures_recorded += 1;
return Ok(());
}
summary.receipts_materialized += self
.schedulers
.materialize_receipts(scheduler, self.batch)
.await?
.len() as i64;
let event = self.events.find(scheduler.event_id).await?;
loop {
let due = self.schedulers.due_receipts(scheduler.id, self.batch).await?;
if due.is_empty() {
break;
}
for receipt in &due {
if event.date_end <= chrono::Utc::now() {
self.schedulers
.mark_receipt_dropped_window_closed(receipt.receipt_id)
.await?;
summary.receipts_dropped_window_closed += 1;
continue;
}
let ctx = RenderContext {
event_id: event.id,
event_name: event.name.clone(),
event_date_begin: event.date_begin,
event_date_end: event.date_end,
event_date_tz: event.date_tz.clone(),
registration_id: receipt.registration_id,
attendee_name: receipt.attendee_name.clone(),
attendee_email: receipt.attendee_email.clone(),
attendee_phone: receipt.attendee_phone.clone(),
registration_barcode: receipt.barcode.clone(),
};
let rendered = match self
.renderer
.render(scheduler.template_ref, scheduler.template_kind.as_deref(), &ctx)
.await
{
Ok(mail) => mail,
Err(failure) => {
let kind = match failure {
super::template_port::RenderFailure::TemplateUnresolved => {
"template_unresolved"
}
super::template_port::RenderFailure::RendererNotComposed => {
"template_renderer_not_composed"
}
super::template_port::RenderFailure::RenderFailed(_) => "render_failed",
};
self.schedulers.record_failure(scheduler.id, kind).await?;
summary.typed_failures_recorded += 1;
continue;
}
};
if scheduler.notification_channel == "sms" {
let phone = receipt
.attendee_phone
.as_deref()
.map(str::trim)
.unwrap_or("");
if phone.is_empty() {
self.schedulers
.record_failure(scheduler.id, "recipient_invalid")
.await?;
summary.typed_failures_recorded += 1;
continue;
}
let sms = RenderedSms {
body_text: rendered.body_text,
};
match self.sms_queue.enqueue(phone, &sms).await {
Ok(()) => {
self.schedulers
.mark_receipt_queued(receipt.receipt_id)
.await?;
summary.sms_queued += 1;
}
Err(_refusal) => {
self.schedulers
.record_failure(scheduler.id, "sms_enqueue_refused")
.await?;
summary.typed_failures_recorded += 1;
}
}
continue;
}
if receipt.attendee_email.is_empty() || !receipt.attendee_email.contains('@') {
self.schedulers
.record_failure(scheduler.id, "recipient_invalid")
.await?;
summary.typed_failures_recorded += 1;
continue;
}
match self.queue.enqueue(&receipt.attendee_email, &rendered).await {
Ok(()) => {
self.schedulers
.mark_receipt_queued(receipt.receipt_id)
.await?;
summary.receipts_queued += 1;
}
Err(refusal) => {
let _ = refusal;
self.schedulers
.record_failure(scheduler.id, "enqueue_refused")
.await?;
summary.typed_failures_recorded += 1;
}
}
}
if (due.len() as i64) < self.batch {
break;
}
}
let done = self.schedulers.recompute_mail_done(scheduler.id).await?;
if done {
summary.schedulers_completed += 1;
self.schedulers.clear_failure(scheduler.id).await?;
}
let unmaterialized = self.schedulers.unmaterialized_count(scheduler.id).await?;
let still_due = self
.schedulers
.due_receipts(scheduler.id, 1)
.await?
.len() as i64;
if !done && (unmaterialized > 0 || still_due > 0) {
self.schedulers.rearm(scheduler.id).await?;
summary.schedulers_rearmed += 1;
}
Ok(())
}
pub async fn run_scheduler_by_id(
&self,
scheduler_id: Uuid,
) -> EventResult<SchedulerRunSummary> {
let mut summary = SchedulerRunSummary::default();
let row = self.schedulers.find(scheduler_id).await?;
summary.schedulers_claimed = 1;
self.run_one(&row, &mut summary).await?;
Ok(summary)
}
pub async fn list_for_event(
&self,
event_id: Uuid,
) -> EventResult<Vec<crate::infrastructure::persistence::scheduler_repository::SchedulerRow>>
{
self.schedulers.list_for_event(event_id).await
}
pub async fn sweep_mark_done(&self, limit: i64) -> EventResult<Vec<Uuid>> {
self.schedulers.sweep_mark_done(limit).await
}
pub async fn on_template_deleted(
&self,
template_kind: &str,
template_ref: Uuid,
actor: Option<Uuid>,
) -> EventResult<(i64, i64)> {
if template_kind.trim().is_empty() {
return Err(super::event_error::EventError::Validation(
"template cascade requires a template_kind".into(),
));
}
self.schedulers
.cascade_delete_template_dependents(template_kind, template_ref, actor)
.await
}
}