use std::env::VarError;
use std::time::Duration;
use reliar_core::SettingsError;
#[derive(Clone, Debug)]
#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
#[cfg_attr(feature = "serde", serde(default, deny_unknown_fields))]
#[non_exhaustive]
pub struct PostgresOutboxSettings {
pub schema: String,
pub enqueue_sets_search_path: bool,
#[cfg_attr(
feature = "serde",
serde(rename = "statement_timeout_ms", with = "crate::duration_serde")
)]
pub statement_timeout: Duration,
}
impl Default for PostgresOutboxSettings {
fn default() -> Self {
Self {
schema: "reliar".to_owned(),
enqueue_sets_search_path: false,
statement_timeout: Duration::ZERO,
}
}
}
impl PostgresOutboxSettings {
#[must_use]
pub fn schema(mut self, schema: impl Into<String>) -> Self {
self.schema = schema.into();
self
}
#[must_use]
pub const fn enqueue_sets_search_path(mut self, enabled: bool) -> Self {
self.enqueue_sets_search_path = enabled;
self
}
#[must_use]
pub const fn statement_timeout(mut self, timeout: Duration) -> Self {
self.statement_timeout = timeout;
self
}
pub fn from_env(prefix: &str) -> Result<Self, SettingsError> {
let mut settings = Self::default();
if let Some(v) = env_raw(prefix, "SCHEMA")? {
settings.schema = v;
}
if let Some(v) = env_bool(prefix, "ENQUEUE_SETS_SEARCH_PATH")? {
settings.enqueue_sets_search_path = v;
}
if let Some(v) = env_duration_ms(prefix, "STATEMENT_TIMEOUT_MS")? {
settings.statement_timeout = v;
}
Ok(settings)
}
}
#[derive(Clone, Debug)]
#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
#[cfg_attr(feature = "serde", serde(default, deny_unknown_fields))]
#[non_exhaustive]
pub struct PostgresInboxSettings {
pub schema: String,
pub claim_sets_search_path: bool,
#[cfg_attr(
feature = "serde",
serde(rename = "statement_timeout_ms", with = "crate::duration_serde")
)]
pub statement_timeout: Duration,
pub max_attempts: u32,
}
const DEFAULT_MAX_ATTEMPTS: u32 = 10;
impl Default for PostgresInboxSettings {
fn default() -> Self {
Self {
schema: "reliar".to_owned(),
claim_sets_search_path: false,
statement_timeout: Duration::ZERO,
max_attempts: DEFAULT_MAX_ATTEMPTS,
}
}
}
impl PostgresInboxSettings {
#[must_use]
pub fn schema(mut self, schema: impl Into<String>) -> Self {
self.schema = schema.into();
self
}
#[must_use]
pub const fn claim_sets_search_path(mut self, enabled: bool) -> Self {
self.claim_sets_search_path = enabled;
self
}
#[must_use]
pub const fn statement_timeout(mut self, timeout: Duration) -> Self {
self.statement_timeout = timeout;
self
}
#[must_use]
pub const fn max_attempts(mut self, max_attempts: u32) -> Self {
self.max_attempts = max_attempts;
self
}
pub fn from_env(prefix: &str) -> Result<Self, SettingsError> {
let mut settings = Self::default();
if let Some(v) = env_raw(prefix, "SCHEMA")? {
settings.schema = v;
}
if let Some(v) = env_bool(prefix, "CLAIM_SETS_SEARCH_PATH")? {
settings.claim_sets_search_path = v;
}
if let Some(v) = env_duration_ms(prefix, "STATEMENT_TIMEOUT_MS")? {
settings.statement_timeout = v;
}
if let Some(v) = env_u32(prefix, "MAX_ATTEMPTS")? {
settings.max_attempts = v;
}
Ok(settings)
}
pub(crate) fn validate(&self) -> Result<(), crate::PostgresInboxError> {
if self.max_attempts == 0 {
return Err(crate::PostgresInboxError::InvalidSettings {
message: "max_attempts must not be 0 (reads as \"no retries\"; use u32::MAX for \
unbounded)"
.to_owned(),
});
}
Ok(())
}
}
fn env_raw(prefix: &str, suffix: &str) -> Result<Option<String>, SettingsError> {
let key = format!("{prefix}{suffix}");
match std::env::var(&key) {
Ok(value) => Ok(Some(value)),
Err(VarError::NotPresent) => Ok(None),
Err(VarError::NotUnicode(_)) => Err(SettingsError::parse(key, "a UTF-8 string")),
}
}
fn env_bool(prefix: &str, suffix: &str) -> Result<Option<bool>, SettingsError> {
let Some(raw) = env_raw(prefix, suffix)? else {
return Ok(None);
};
match raw.trim().to_ascii_lowercase().as_str() {
"true" | "1" => Ok(Some(true)),
"false" | "0" => Ok(Some(false)),
_ => Err(SettingsError::parse(
format!("{prefix}{suffix}"),
"bool (\"true\"/\"false\"/\"1\"/\"0\")",
)),
}
}
fn env_duration_ms(prefix: &str, suffix: &str) -> Result<Option<Duration>, SettingsError> {
let Some(raw) = env_raw(prefix, suffix)? else {
return Ok(None);
};
let ms = raw
.trim()
.parse::<u64>()
.map_err(|_| SettingsError::parse(format!("{prefix}{suffix}"), "milliseconds"))?;
Ok(Some(Duration::from_millis(ms)))
}
fn env_u32(prefix: &str, suffix: &str) -> Result<Option<u32>, SettingsError> {
let Some(raw) = env_raw(prefix, suffix)? else {
return Ok(None);
};
raw.trim()
.parse::<u32>()
.map(Some)
.map_err(|_| SettingsError::parse(format!("{prefix}{suffix}"), "u32"))
}