use std::collections::BTreeMap;
use serde::{Deserialize, Serialize};
use crate::types::EventPayload;
pub mod email;
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize)]
#[non_exhaustive]
pub enum Transport {
Push,
Pull,
}
#[derive(Debug, Clone, PartialEq, Eq, Hash, Serialize, Deserialize)]
#[non_exhaustive]
pub enum IngestSource {
Webhook {
provider: String,
},
Email,
}
impl IngestSource {
#[must_use]
pub fn as_key(&self) -> String {
match self {
IngestSource::Webhook { provider } => format!("webhook:{provider}"),
IngestSource::Email => "email".to_string(),
}
}
#[must_use]
pub fn from_key(key: &str) -> Option<Self> {
if key == "email" {
return Some(IngestSource::Email);
}
key.strip_prefix("webhook:")
.filter(|provider| !provider.is_empty())
.map(|provider| IngestSource::Webhook {
provider: provider.to_string(),
})
}
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct StorageRef {
pub bucket: String,
pub key: String,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct Attachment {
pub storage: StorageRef,
pub content_type: String,
pub filename: String,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize)]
#[serde(rename_all = "kebab-case")]
#[non_exhaustive]
pub enum Classification {
Human,
OutOfOffice,
Bounce,
Challenge,
AutoGenerated,
}
impl Classification {
#[must_use]
pub const fn is_automated(self) -> bool {
!matches!(self, Classification::Human)
}
}
#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
pub struct InboundRouting {
pub entity_type: Option<String>,
pub entity_id: Option<String>,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct InboundMessage {
pub source: IngestSource,
pub idempotency_key: String,
pub thread_key: Option<String>,
pub from: Option<String>,
pub to: Vec<String>,
pub cc: Vec<String>,
pub subject: Option<String>,
pub body_text: Option<String>,
pub body_html: Option<String>,
pub payload: Option<serde_json::Value>,
pub attachments: Vec<Attachment>,
pub headers: BTreeMap<String, String>,
pub received_at: chrono::DateTime<chrono::Utc>,
pub raw_ref: Option<StorageRef>,
pub routing: InboundRouting,
#[serde(default)]
pub classification: Option<Classification>,
}
impl InboundMessage {
#[must_use]
pub fn new(
source: IngestSource,
idempotency_key: impl Into<String>,
received_at: chrono::DateTime<chrono::Utc>,
) -> Self {
Self {
source,
idempotency_key: idempotency_key.into(),
thread_key: None,
from: None,
to: Vec::new(),
cc: Vec::new(),
subject: None,
body_text: None,
body_html: None,
payload: None,
attachments: Vec::new(),
headers: BTreeMap::new(),
received_at,
raw_ref: None,
routing: InboundRouting::default(),
classification: None,
}
}
#[must_use]
pub fn trigger_type(&self) -> String {
format!("after:ingest:{}", self.source.as_key())
}
}
pub struct RawDelivery<'a> {
pub event_id: &'a str,
pub event_type: &'a str,
pub payload: &'a serde_json::Value,
pub headers: &'a BTreeMap<String, String>,
pub received_at: chrono::DateTime<chrono::Utc>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct IngestError {
pub message: String,
}
impl IngestError {
#[must_use]
pub fn new(message: impl Into<String>) -> Self {
Self {
message: message.into(),
}
}
}
impl std::fmt::Display for IngestError {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
write!(f, "{}", self.message)
}
}
impl std::error::Error for IngestError {}
pub trait Source: Send + Sync {
fn source(&self) -> IngestSource;
fn transport(&self) -> Transport;
}
pub trait PushSource: Source {
fn normalize(&self, delivery: &RawDelivery<'_>) -> Result<InboundMessage, IngestError>;
}
#[derive(Debug, Clone, Default)]
pub struct PullContext {
pub cursor: Option<serde_json::Value>,
}
#[derive(Debug, Clone)]
pub struct PullBatch {
pub messages: Vec<InboundMessage>,
pub next_cursor: serde_json::Value,
}
impl PullBatch {
#[must_use]
pub fn empty(current_cursor: Option<serde_json::Value>) -> Self {
Self {
messages: Vec::new(),
next_cursor: current_cursor.unwrap_or(serde_json::Value::Null),
}
}
}
#[allow(async_fn_in_trait)] pub trait PullSource: Source {
async fn poll(&self, ctx: &PullContext) -> Result<PullBatch, IngestError>;
}
#[derive(Debug, Clone)]
pub struct IngestTrigger {
pub function_name: String,
pub source: Option<IngestSource>,
}
impl IngestTrigger {
#[must_use]
pub fn matches(&self, message: &InboundMessage) -> bool {
self.source.as_ref().is_none_or(|source| *source == message.source)
}
#[must_use]
pub fn build_payload(&self, message: &InboundMessage) -> EventPayload {
EventPayload {
trigger_type: message.trigger_type(),
entity: message.source.as_key(),
event_kind: "ingest".to_string(),
data: serde_json::to_value(message).unwrap_or(serde_json::Value::Null),
timestamp: message.received_at,
}
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct Recipient {
pub base: String,
pub tag: Option<String>,
}
#[must_use]
pub fn parse_recipient(address: &str) -> Option<Recipient> {
let (local, domain) = address.rsplit_once('@')?;
if local.is_empty() || domain.is_empty() {
return None;
}
match local.split_once('+') {
Some((base_local, tag)) if !base_local.is_empty() && !tag.is_empty() => Some(Recipient {
base: format!("{base_local}@{domain}"),
tag: Some(tag.to_string()),
}),
_ => Some(Recipient {
base: format!("{local}@{domain}"),
tag: None,
}),
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct RoutingRule {
pub address: String,
pub entity_type: String,
}
#[must_use]
pub fn resolve_routing(message: &InboundMessage, rules: &[RoutingRule]) -> InboundRouting {
for address in message.to.iter().chain(message.cc.iter()) {
let Some(recipient) = parse_recipient(address) else {
continue;
};
if let Some(rule) =
rules.iter().find(|rule| rule.address.eq_ignore_ascii_case(&recipient.base))
{
return InboundRouting {
entity_type: Some(rule.entity_type.clone()),
entity_id: recipient.tag,
};
}
}
InboundRouting::default()
}
#[cfg(test)]
mod tests;