use std::collections::BTreeMap;
use std::sync::Arc;
use uuid::Uuid;
use super::event_error::{EventError, EventResult};
use super::lead_sink::{EventLeadSink, LeadGroup, LeadRegistrationView};
use crate::infrastructure::persistence::lead_command_repository::{
EligibleRegistration, LeadCommandRepository,
};
pub const DEFAULT_LEAD_BATCH: i64 = 200;
pub const DEFAULT_LEAD_CRON_LIMIT: i64 = 1000;
pub const LEASE_SECONDS: i64 = 900;
pub const GROUPING_PER_ORDER: &str = "per_order";
pub const GROUPING_PER_EVENT_DAY: &str = "per_event_day";
#[derive(Debug, Clone, Default, serde::Serialize)]
pub struct LeadRunSummary {
pub requests_claimed: i64,
pub rules_evaluated: i64,
pub registrations_matched: i64,
pub groups_created: i64,
pub groups_grown: i64,
pub leads_created: i64,
pub leads_updated: i64,
pub requests_completed: i64,
pub requests_parked: i64,
}
pub struct LeadGenerationService {
leads: LeadCommandRepository,
sink: Arc<dyn EventLeadSink>,
batch: i64,
cron_limit: i64,
}
impl LeadGenerationService {
pub fn new(leads: LeadCommandRepository, sink: Arc<dyn EventLeadSink>) -> Self {
Self {
leads,
sink,
batch: DEFAULT_LEAD_BATCH,
cron_limit: DEFAULT_LEAD_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_lead_requests(&self) -> EventResult<LeadRunSummary> {
let mut summary = LeadRunSummary::default();
let mut parked: Vec<Uuid> = Vec::new();
loop {
if summary.requests_claimed >= self.cron_limit {
break;
}
let Some(request) = self
.leads
.claim_next_due_request(LEASE_SECONDS, &parked)
.await?
else {
break;
};
summary.requests_claimed += 1;
if self.run_one(request.event_id, &mut summary).await? {
parked.push(request.event_id);
}
}
Ok(summary)
}
pub async fn run_request_of_event(&self, event_id: Uuid) -> EventResult<LeadRunSummary> {
let mut summary = LeadRunSummary::default();
let request = self
.leads
.try_lease_request_of_event(event_id, LEASE_SECONDS)
.await?;
summary.requests_claimed = 1;
let _parked = self.run_one(request.event_id, &mut summary).await?;
Ok(summary)
}
pub async fn request_of_event(
&self,
event_id: Uuid,
) -> EventResult<Option<crate::infrastructure::persistence::lead_command_repository::LeadRequestRow>>
{
self.leads.find_request_of_event(event_id).await
}
async fn run_one(&self, event_id: Uuid, summary: &mut LeadRunSummary) -> EventResult<bool> {
let rules = self.leads.active_rules_of_event(event_id).await?;
let mut park_reason: Option<String> = None;
for rule in &rules {
summary.rules_evaluated += 1;
let predicates = self.leads.predicates_of_rule(rule.id).await?;
loop {
let eligible = self
.leads
.eligible_registrations(event_id, rule, &predicates, self.batch)
.await?;
if eligible.is_empty() {
break;
}
summary.registrations_matched += eligible.len() as i64;
let refusal = self
.generate_groups(rule.id, event_id, &eligible, summary)
.await?;
if let Some(reason) = refusal {
park_reason = Some(reason);
break;
}
if (eligible.len() as i64) < self.batch {
break;
}
}
if park_reason.is_some() {
break;
}
}
let parked = park_reason.is_some();
match park_reason {
Some(reason) => {
self.leads.finish_request_by_event(event_id, Some(&reason)).await?;
summary.requests_parked += 1;
}
None => {
self.leads.finish_request_by_event(event_id, None).await?;
summary.requests_completed += 1;
}
}
Ok(parked)
}
async fn generate_groups(
&self,
rule_id: Uuid,
event_id: Uuid,
eligible: &[EligibleRegistration],
summary: &mut LeadRunSummary,
) -> Result<Option<String>, EventError> {
let mut groups: BTreeMap<(String, String), Vec<&EligibleRegistration>> = BTreeMap::new();
for r in eligible {
match r.sale_order_id {
Some(order_id) => {
groups
.entry((GROUPING_PER_ORDER.into(), order_id.to_string()))
.or_default()
.push(r);
}
None => {
let day = r
.created_at
.map(|ts| ts.date_naive().to_string())
.unwrap_or_else(|| "unknown-day".to_string());
groups
.entry((
GROUPING_PER_EVENT_DAY.into(),
format!("{event_id}:{day}"),
))
.or_default()
.push(r);
}
}
}
for ((grouping, group_key), members) in groups {
let (provenance_id, existing_lead) = self
.leads
.upsert_provenance(rule_id, event_id, &group_key, &grouping)
.await?;
let mut views: Vec<LeadRegistrationView> = self
.leads
.group_members(provenance_id)
.await?
.iter()
.map(view_of)
.collect();
let known: std::collections::BTreeSet<Uuid> =
views.iter().map(|v| v.id).collect();
let mut new_ids = Vec::new();
for m in &members {
if !known.contains(&m.id) {
new_ids.push(m.id);
views.push(view_of(m));
}
}
let group = LeadGroup {
rule_id,
event_id,
group_key: group_key.clone(),
grouping: grouping.clone(),
registrations: views,
};
match existing_lead {
Some(lead_id) => {
if !new_ids.is_empty() {
if let Err(reason) = self.sink.update_lead(lead_id, &group).await {
return Ok(Some(reason));
}
self.leads.junction_add(provenance_id, &new_ids).await?;
summary.groups_grown += 1;
summary.leads_updated += 1;
}
}
None => {
let lead_id = match self.sink.create_lead(&group).await {
Ok(id) => id,
Err(reason) => return Ok(Some(reason)),
};
self.leads.set_provenance_lead(provenance_id, lead_id).await?;
self.leads.junction_add(provenance_id, &new_ids).await?;
summary.groups_created += 1;
summary.leads_created += 1;
}
}
}
Ok(None)
}
}
fn view_of(m: &EligibleRegistration) -> LeadRegistrationView {
LeadRegistrationView {
id: m.id,
name: m.name.clone(),
email: m.email.clone(),
phone: m.phone.clone(),
company_name: m.company_name.clone(),
partner_id: m.partner_id,
sale_order_id: m.sale_order_id,
}
}