use super::*;
#[derive(Clone, Copy)]
pub(crate) struct EmbeddedMigration {
pub(crate) version: i64,
pub(crate) description: &'static str,
pub(crate) sql: &'static str,
}
include!(concat!(env!("OUT_DIR"), "/migration_inventory.rs"));
pub(crate) fn embedded_migrator(files: &[EmbeddedMigration]) -> Migrator {
Migrator::with_migrations(
files
.iter()
.map(|migration| {
Migration::new(
migration.version,
migration.description.into(),
MigrationType::Simple,
sqlx::SqlSafeStr::into_sql_str(migration.sql),
false,
)
})
.collect(),
)
}
#[expect(
clippy::items_after_test_module,
reason = "migration parity tests stay beside generated registration data"
)]
#[cfg(test)]
mod tests {
use super::*;
#[cfg(feature = "sqlite")]
#[test]
fn generated_sqlite_inventory_preserves_order_descriptions_and_bytes() {
let versions = SQLITE_MIGRATIONS
.iter()
.map(|migration| migration.version)
.collect::<Vec<_>>();
let descriptions = SQLITE_MIGRATIONS
.iter()
.map(|migration| migration.description)
.collect::<Vec<_>>();
let sql = SQLITE_MIGRATIONS
.iter()
.map(|migration| migration.sql)
.collect::<Vec<_>>();
assert_eq!(versions, vec![1, 2, 3, 4]);
assert_eq!(
descriptions,
vec![
"initial",
"command ledger",
"projection protocol",
"command ledger atomic state"
]
);
assert_eq!(
sql,
vec![
include_str!(concat!(
env!("CARGO_MANIFEST_DIR"),
"/migrations/sqlite/0001_initial.sql"
)),
include_str!(concat!(
env!("CARGO_MANIFEST_DIR"),
"/migrations/sqlite/0002_command_ledger.sql"
)),
include_str!(concat!(
env!("CARGO_MANIFEST_DIR"),
"/migrations/sqlite/0003_projection_protocol.sql"
)),
include_str!(concat!(
env!("CARGO_MANIFEST_DIR"),
"/migrations/sqlite/0004_command_ledger_atomic_state.sql"
)),
]
);
}
#[cfg(feature = "postgres")]
#[test]
fn generated_postgres_inventory_preserves_order_descriptions_and_bytes() {
let versions = POSTGRES_MIGRATIONS
.iter()
.map(|migration| migration.version)
.collect::<Vec<_>>();
let descriptions = POSTGRES_MIGRATIONS
.iter()
.map(|migration| migration.description)
.collect::<Vec<_>>();
let sql = POSTGRES_MIGRATIONS
.iter()
.map(|migration| migration.sql)
.collect::<Vec<_>>();
assert_eq!(versions, vec![1, 2, 3, 4]);
assert_eq!(
descriptions,
vec![
"initial",
"command ledger",
"projection protocol",
"command ledger atomic state"
]
);
assert_eq!(
sql,
vec![
include_str!(concat!(
env!("CARGO_MANIFEST_DIR"),
"/migrations/postgres/0001_initial.sql"
)),
include_str!(concat!(
env!("CARGO_MANIFEST_DIR"),
"/migrations/postgres/0002_command_ledger.sql"
)),
include_str!(concat!(
env!("CARGO_MANIFEST_DIR"),
"/migrations/postgres/0003_projection_protocol.sql"
)),
include_str!(concat!(
env!("CARGO_MANIFEST_DIR"),
"/migrations/postgres/0004_command_ledger_atomic_state.sql"
)),
]
);
}
}
pub(super) fn ids_by_type(identities: &[StreamIdentity]) -> BTreeMap<&str, Vec<&str>> {
let mut groups: BTreeMap<&str, Vec<&str>> = BTreeMap::new();
for identity in identities {
groups
.entry(identity.aggregate_type())
.or_default()
.push(identity.aggregate_id());
}
groups
}
pub(super) const EVENT_BIND_COLUMNS: usize = 10;
pub(super) const OUTBOX_BIND_COLUMNS: usize = 19;
pub trait SqlxRepoBackend: SqlxReadModelBackend {
fn migrator() -> &'static Migrator;
const MAX_BIND_PARAMS: usize;
const CONFLICT_REREAD_IN_TX: bool;
const NOW: &'static str;
const COMMAND_LEDGER_SELECT: &'static str;
const COMMAND_LEDGER_LOCK_SUFFIX: &'static str;
const COMMAND_LEDGER_COMPACTION_LOCK_SUFFIX: &'static str;
const EVENT_SELECT: &'static str;
const SNAPSHOT_SELECT: &'static str;
const OUTBOX_SELECT: &'static str;
const ORDER_BY_CREATED_AT: &'static str;
const OUTBOX_OLDEST_CREATED_AT_SELECT: &'static str;
const TABLE_DIALECT: TableSqlDialect;
type TimestampValue: Send + Sync + 'static;
fn default_pool_size(database_url: &str) -> u32 {
let _ = database_url;
5
}
fn is_unique_violation(err: &sqlx::Error) -> bool;
fn timestamp_value(timestamp: SystemTime) -> Result<Self::TimestampValue, RepositoryError>;
fn push_timestamp(sep: &mut Separated<'_, Self, &'static str>, value: &Self::TimestampValue);
fn push_optional_timestamp(
sep: &mut Separated<'_, Self, &'static str>,
value: Option<&Self::TimestampValue>,
);
fn push_timestamp_assign(builder: &mut QueryBuilder<Self>, value: &Self::TimestampValue);
fn push_timestamp_cmp(
builder: &mut QueryBuilder<Self>,
column: &'static str,
op: &'static str,
epoch_secs: f64,
);
fn push_command_ledger_now(builder: &mut QueryBuilder<Self>);
fn push_command_ledger_now_epoch(builder: &mut QueryBuilder<Self>);
fn push_command_ledger_deadline(builder: &mut QueryBuilder<Self>, duration: Duration);
fn push_command_ledger_deadline_is_live(
builder: &mut QueryBuilder<Self>,
deadline: &Self::TimestampValue,
);
fn push_command_ledger_json(builder: &mut QueryBuilder<Self>, json: &str);
fn decode_timestamp(
row: &Self::Row,
column: &'static str,
) -> Result<SystemTime, RepositoryError>;
fn decode_optional_timestamp(
row: &Self::Row,
column: &'static str,
) -> Result<Option<SystemTime>, RepositoryError>;
fn push_metadata(sep: &mut Separated<'_, Self, &'static str>, json: &str);
fn push_id_filter(builder: &mut QueryBuilder<Self>, ids: &[&str]);
fn inbox_purge_query(age: Duration) -> QueryBuilder<Self>;
fn claim_outbox<'a>(
pool: &'a Pool<Self>,
request: ClaimOutboxMessages,
) -> impl Future<Output = Result<Vec<OutboxMessage>, RepositoryError>> + Send + 'a;
}