use std::{collections::HashMap, future::Future, pin::Pin, sync::Arc, time::Duration};
use fraiseql_functions::host::live::storage::StorageBackend;
use tracing::{info, warn};
pub mod admin;
pub mod config;
pub mod correlation;
pub mod cursor;
pub mod imap;
pub mod poller;
pub mod probe;
pub mod sink;
pub mod smtp;
pub mod source;
pub mod tracking;
pub mod warming;
pub use admin::{SuppressionAdminState, suppression_admin_router};
pub use config::{
ImapConfig, MailboxConfig, MailboxSmtpConfig, ReturnPathConfig, RoutingRuleConfig,
SendSettings, SmtpTlsMode,
};
pub use correlation::correlate;
pub use cursor::Cursor;
pub use imap::{FetchBatch, FetchedMessage, ImapMailboxFetcher, MailboxFetcher};
pub use poller::MailboxPoller;
pub use probe::{ProbeOutcome, probe_recipient, run_return_path_probe};
pub use sink::EmailIngestSink;
pub use smtp::{SmtpMailboxTransport, build_email_transport};
pub use source::ImapSource;
pub use tracking::{
CorrelatedSend, PgSendTracker, RecordedSend, SendCorrelator, SendTracker, SentRecord,
SuppressionReason,
};
pub use warming::{SendCounter, WarmingState, warming_daily_limit};
use crate::subsystems::BeforeMutationHooks;
pub async fn init_cursor_store(pool: &sqlx::PgPool) -> fraiseql_error::Result<()> {
fraiseql_observers::PostgresSourceCursorStore::new(pool.clone())
.init()
.await
.map_err(|error| fraiseql_error::FraiseQLError::database(error.to_string()))
}
#[must_use]
#[allow(clippy::too_many_arguments)] pub fn build_pollers<S: std::hash::BuildHasher>(
mailboxes: &HashMap<String, MailboxConfig, S>,
pool: &sqlx::PgPool,
hooks: Option<&Arc<BeforeMutationHooks>>,
query_executor_factory: Option<&crate::routes::after_mutation::QueryExecutorFactory>,
attachment_sink: Option<&Arc<dyn StorageBackend>>,
correlator: Option<&Arc<dyn SendCorrelator>>,
address_hash_key: Option<&Arc<[u8]>>,
challenge_suppress_after: u32,
get_env: impl Fn(&str) -> Option<String>,
) -> Vec<(MailboxPoller, Duration)> {
let mut pollers = Vec::new();
for (name, mailbox) in mailboxes {
let Some(imap) = mailbox.imap.as_ref() else {
continue;
};
let Some(password) = get_env(&imap.password_env) else {
warn!(
mailbox = %name,
password_env = %imap.password_env,
"poll-IMAP mailbox not started: password env is unset"
);
continue;
};
let fetcher = match ImapMailboxFetcher::new(
&imap.host,
imap.port,
&imap.username,
password,
&imap.mailbox,
) {
Ok(fetcher) => Arc::new(fetcher),
Err(error) => {
warn!(mailbox = %name, %error, "poll-IMAP mailbox not started: TLS setup failed");
continue;
},
};
let routing_rules = imap.routing.iter().map(RoutingRuleConfig::to_rule).collect();
let source = ImapSource::new(
name.clone(),
fetcher,
routing_rules,
imap.batch_size,
imap.attachment_bucket.clone(),
attachment_sink.cloned(),
);
let ingest_sink = EmailIngestSink::new(
name.clone(),
pool.clone(),
hooks.cloned(),
query_executor_factory.cloned(),
correlator.cloned(),
address_hash_key.cloned(),
challenge_suppress_after,
);
let store = fraiseql_observers::PostgresSourceCursorStore::new(pool.clone());
let runner = fraiseql_observers::LeaseGuardedRunner::postgres(pool.clone(), name.clone());
let poller = MailboxPoller::new(source, ingest_sink, store, runner);
let interval = Duration::from_secs(imap.poll_interval_secs.max(1));
info!(
mailbox = %name,
host = %imap.host,
folder = %imap.mailbox,
"poll-IMAP mailbox configured"
);
pollers.push((poller, interval));
}
pollers
}
const PROBE_TIMEOUT: Duration = Duration::from_secs(120);
const PROBE_INTERVAL: Duration = Duration::from_secs(10);
pub async fn run_startup_probes<S: std::hash::BuildHasher + Sync>(
mailboxes: &HashMap<String, MailboxConfig, S>,
get_env: impl Fn(&str) -> Option<String> + Send,
) {
for (name, mailbox) in mailboxes {
let (Some(imap), Some(smtp)) = (mailbox.imap.as_ref(), mailbox.smtp.as_ref()) else {
continue;
};
let Some(transport) =
SmtpMailboxTransport::build(std::iter::once((name.as_str(), smtp)), &get_env)
else {
continue; };
let Some(password) = get_env(&imap.password_env) else {
continue;
};
let fetcher = match ImapMailboxFetcher::new(
&imap.host,
imap.port,
&imap.username,
password,
&imap.mailbox,
) {
Ok(fetcher) => fetcher,
Err(error) => {
warn!(mailbox = %name, %error, "Return-Path probe skipped: IMAP TLS setup failed");
continue;
},
};
let nonce = uuid::Uuid::new_v4().simple().to_string();
let sender = fraiseql_functions::SenderIdentity {
address: smtp.address.clone(),
display_name: None,
};
let probe_to = probe::probe_recipient(
smtp.return_path_local_part(),
smtp.return_path_domain(),
&nonce,
);
match probe::run_return_path_probe(
&transport,
&fetcher,
&sender,
&probe_to,
&nonce,
PROBE_TIMEOUT,
PROBE_INTERVAL,
)
.await
{
Ok(ProbeOutcome::Confirmed) => info!(
mailbox = %name,
"VERP Return-Path probe confirmed — delivery correlation is active"
),
Ok(ProbeOutcome::NotObserved) => warn!(
mailbox = %name,
"VERP Return-Path probe did NOT land within the window — plus-addressing may be \
stripped by the provider; delivery correlation may not work (sends still go out, \
but bounces/challenges/replies will not be tracked)"
),
Err(error) => {
warn!(mailbox = %name, %error, "VERP Return-Path probe failed to run");
},
}
}
}
pub struct LegacyStorageSink {
backend: Arc<dyn crate::storage::StorageBackend>,
}
impl LegacyStorageSink {
#[must_use]
pub fn new(backend: Arc<dyn crate::storage::StorageBackend>) -> Self {
Self { backend }
}
fn flat_key(bucket: &str, key: &str) -> String {
format!("{bucket}/{key}")
}
}
impl StorageBackend for LegacyStorageSink {
fn get(
&self,
bucket: &str,
key: &str,
) -> Pin<Box<dyn Future<Output = fraiseql_error::Result<Vec<u8>>> + Send + '_>> {
let full = Self::flat_key(bucket, key);
let backend = Arc::clone(&self.backend);
Box::pin(async move { backend.download(&full).await.map_err(Into::into) })
}
fn put(
&self,
bucket: &str,
key: &str,
body: &[u8],
content_type: &str,
) -> Pin<Box<dyn Future<Output = fraiseql_error::Result<()>> + Send + '_>> {
let full = Self::flat_key(bucket, key);
let body = body.to_vec();
let content_type = content_type.to_string();
let backend = Arc::clone(&self.backend);
Box::pin(async move {
backend
.upload(&full, &body, &content_type)
.await
.map(|_| ())
.map_err(Into::into)
})
}
}