use std::collections::HashMap;
#[cfg(any(feature = "postgres", feature = "sqlite"))]
use crate::read_model::ReadModelError;
use crate::repository::RepositoryError;
pub(crate) fn serialize_event_metadata(
metadata: &HashMap<String, String>,
) -> Result<String, RepositoryError> {
serde_json::to_string(metadata)
.map_err(|err| RepositoryError::Model(format!("serialize event metadata: {err}")))
}
pub(crate) fn deserialize_event_metadata(
metadata_json: &str,
) -> Result<HashMap<String, String>, RepositoryError> {
serde_json::from_str(metadata_json)
.map_err(|err| RepositoryError::Model(format!("deserialize event metadata: {err}")))
}
pub(crate) fn repository_i64_from_u64(
backend: &str,
value: u64,
field: &str,
storage: &str,
) -> Result<i64, RepositoryError> {
i64::try_from(value).map_err(|_| {
RepositoryError::Model(format!("{backend} {field} value {value} exceeds {storage}"))
})
}
#[cfg(feature = "postgres")]
pub(crate) fn repository_i32_from_u64(
backend: &str,
value: u64,
field: &str,
storage: &str,
) -> Result<i32, RepositoryError> {
i32::try_from(value).map_err(|_| {
RepositoryError::Model(format!("{backend} {field} value {value} exceeds {storage}"))
})
}
pub(crate) fn repository_u64_from_i64(
backend: &str,
value: i64,
field: &str,
) -> Result<u64, RepositoryError> {
u64::try_from(value)
.map_err(|_| RepositoryError::Model(format!("{backend} {field} value {value} is negative")))
}
#[cfg(feature = "postgres")]
pub(crate) fn repository_u64_from_i32(
backend: &str,
value: i32,
field: &str,
) -> Result<u64, RepositoryError> {
u64::try_from(value)
.map_err(|_| RepositoryError::Model(format!("{backend} {field} value {value} is negative")))
}
#[cfg(feature = "sqlite")]
pub(crate) fn repository_u16_from_i64(
backend: &str,
value: i64,
field: &str,
) -> Result<u16, RepositoryError> {
u16::try_from(value)
.map_err(|_| RepositoryError::Model(format!("{backend} {field} value {value} is invalid")))
}
#[cfg(feature = "postgres")]
pub(crate) fn repository_u16_from_i32(
backend: &str,
value: i32,
field: &str,
) -> Result<u16, RepositoryError> {
u16::try_from(value)
.map_err(|_| RepositoryError::Model(format!("{backend} {field} value {value} is invalid")))
}
#[cfg(any(feature = "postgres", feature = "sqlite"))]
pub(crate) fn read_model_i64_from_u64(
backend: &str,
value: u64,
field: &str,
storage: &str,
) -> Result<i64, ReadModelError> {
i64::try_from(value).map_err(|_| {
ReadModelError::Storage(format!("{backend} {field} value {value} exceeds {storage}"))
})
}
#[cfg(any(feature = "postgres", feature = "sqlite"))]
pub(crate) fn read_model_u64_from_i64(
backend: &str,
value: i64,
field: &str,
) -> Result<u64, ReadModelError> {
u64::try_from(value).map_err(|_| {
ReadModelError::Storage(format!("{backend} {field} value {value} is negative"))
})
}
#[cfg(any(feature = "postgres", feature = "sqlite"))]
pub(crate) fn audited_table_schema_sql(statement: String) -> sqlx::AssertSqlSafe<String> {
sqlx::AssertSqlSafe(statement)
}
#[cfg(feature = "sqlite")]
pub(crate) fn is_sqlite_unique_constraint(err: &sqlx::Error) -> bool {
match err {
sqlx::Error::Database(db_err) => {
let message = db_err.message();
let code = db_err.code().map(|code| code.into_owned());
message.contains("UNIQUE constraint failed")
|| message.contains("PRIMARY KEY")
|| matches!(code.as_deref(), Some("1555" | "2067"))
}
_ => false,
}
}
#[cfg(feature = "sqlite")]
pub(crate) fn is_sqlite_busy(err: &sqlx::Error) -> bool {
match err {
sqlx::Error::Database(db_err) => {
let message = db_err.message().to_ascii_lowercase();
let code = db_err.code().map(|code| code.into_owned());
message.contains("database is locked")
|| message.contains("database table is locked")
|| matches!(
code.as_deref(),
Some("5" | "6" | "261" | "262" | "517" | "518")
)
}
sqlx::Error::PoolTimedOut => true,
_ => false,
}
}
#[cfg(feature = "postgres")]
pub(crate) fn is_postgres_unique_violation(err: &sqlx::Error) -> bool {
match err {
sqlx::Error::Database(db_err) => db_err.code().as_deref() == Some("23505"),
_ => false,
}
}
pub(crate) fn is_sqlx_transient(err: &sqlx::Error) -> bool {
if matches!(
err,
sqlx::Error::PoolTimedOut | sqlx::Error::PoolClosed | sqlx::Error::Io(_)
) {
return true;
}
#[cfg(feature = "sqlite")]
if is_sqlite_busy(err) {
return true;
}
if let sqlx::Error::Database(db_err) = err {
if matches!(db_err.code().as_deref(), Some("40001" | "40P01")) {
return true;
}
}
false
}
pub(crate) fn repository_storage_error(
backend: &str,
operation: &str,
err: sqlx::Error,
) -> RepositoryError {
let retryable = is_sqlx_transient(&err);
RepositoryError::Storage {
operation: format!("{backend} {operation}"),
retryable,
source: Some(Box::new(err)),
}
}
#[cfg(any(feature = "postgres", feature = "sqlite"))]
pub(crate) fn read_model_storage_error(
backend: &str,
operation: &str,
err: sqlx::Error,
) -> ReadModelError {
ReadModelError::Storage(format!("{backend} {operation} failed: {err}"))
}