use std::collections::HashMap;
use std::sync::Mutex;
use uuid::Uuid;
use super::event_error::{EventError, EventResult};
use crate::infrastructure::persistence::seat_repository::{RegisterCommand, RegistrationRow};
pub const EVENT_TRUSTED_PROXY_ENV: &str = "EVENT_TRUSTED_PROXY";
#[derive(Debug, Clone, serde::Deserialize)]
#[serde(deny_unknown_fields)]
pub struct IntakePayload {
pub name: String,
pub email: String,
#[serde(default)]
pub phone: Option<String>,
#[serde(default)]
pub company_name: Option<String>,
#[serde(default)]
pub event_slot_id: Option<Uuid>,
#[serde(default)]
pub event_ticket_id: Option<Uuid>,
#[serde(default)]
pub answers: Vec<IntakeAnswer>,
}
#[derive(Debug, Clone, serde::Deserialize)]
#[serde(deny_unknown_fields)]
pub struct IntakeAnswer {
pub question_id: Uuid,
pub value_text: Option<String>,
pub value_answer_id: Option<Uuid>,
}
#[derive(Debug, Default)]
pub struct FixedWindows {
inner: Mutex<HashMap<String, (u64, u64)>>, }
impl FixedWindows {
pub fn new() -> Self {
Self::default()
}
pub fn allow(&self, key: &str, max: u64, window_secs: u64) -> bool {
let now = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map(|d| d.as_secs())
.unwrap_or(0);
let mut guard = match self.inner.lock() {
Ok(g) => g,
Err(poisoned) => poisoned.into_inner(),
};
let open = match guard.get_mut(key) {
Some((start, count)) => {
if now.saturating_sub(*start) < window_secs {
true
} else {
*start = now;
*count = 0;
true
}
}
None => {
guard.insert(key.to_string(), (now, 0));
true
}
};
if !open {
return false;
}
let Some((_, count)) = guard.get_mut(key) else {
return false;
};
if *count >= max {
return false;
}
*count += 1;
true
}
}
#[derive(Debug, Clone)]
pub struct IntakeThrottles {
pub per_identity_max: u64,
pub per_identity_window_secs: u64,
pub per_ip_max: u64,
pub per_ip_window_secs: u64,
}
impl Default for IntakeThrottles {
fn default() -> Self {
Self {
per_identity_max: 5,
per_identity_window_secs: 3600,
per_ip_max: 30,
per_ip_window_secs: 3600,
}
}
}
pub struct IntakeService {
registrations: super::registration_service::RegistrationCommandService,
windows: FixedWindows,
throttles: IntakeThrottles,
}
impl IntakeService {
pub fn new(
registrations: super::registration_service::RegistrationCommandService,
throttles: IntakeThrottles,
) -> Self {
Self {
registrations,
windows: FixedWindows::new(),
throttles,
}
}
pub async fn intake(
&self,
event_id: Uuid,
payload: IntakePayload,
client_ip: &str,
) -> EventResult<RegistrationRow> {
if !self.windows.allow(
&format!("identity:{}", payload.email.to_lowercase()),
self.throttles.per_identity_max,
self.throttles.per_identity_window_secs,
) || !self.windows.allow(
&format!("ip:{client_ip}"),
self.throttles.per_ip_max,
self.throttles.per_ip_window_secs,
) {
let retry = self
.throttles
.per_identity_window_secs
.min(self.throttles.per_ip_window_secs)
.max(1) as u32;
return Err(EventError::EventThrottled {
retry_after_secs: retry,
});
}
self.registrations
.register(RegisterCommand {
event_id,
event_slot_id: payload.event_slot_id,
event_ticket_id: payload.event_ticket_id,
name: payload.name,
email: payload.email,
phone: payload.phone,
company_name: payload.company_name,
partner_id: None, actor: None, lead_rule_skip: false, })
.await
}
}