use std::ops::Bound;
use reifydb_core::common::CommitVersion;
use reifydb_sqlite::batch::values_placeholders;
#[inline]
pub(super) fn version_to_bytes(version: CommitVersion) -> [u8; 8] {
version.0.to_be_bytes()
}
#[inline]
pub(super) fn version_from_bytes(bytes: &[u8]) -> CommitVersion {
CommitVersion(u64::from_be_bytes(bytes.try_into().expect("version must be 8 bytes")))
}
pub(super) fn build_create_current_sql(table_name: &str) -> String {
format!(
"CREATE TABLE IF NOT EXISTS \"{0}\" (\
key BLOB PRIMARY KEY,\
version BLOB NOT NULL,\
value BLOB,\
updated_at INTEGER\
) WITHOUT ROWID;\
CREATE INDEX IF NOT EXISTS \"{0}__version\" ON \"{0}\" (version);\
CREATE INDEX IF NOT EXISTS \"{0}__tombstone\" ON \"{0}\" (version) WHERE value IS NULL;\
CREATE INDEX IF NOT EXISTS \"{0}__expiry\" ON \"{0}\" (updated_at) \
WHERE value IS NOT NULL AND updated_at IS NOT NULL;",
table_name
)
}
pub(super) fn build_expired_keys_sql(table_name: &str, has_cursor: bool, limit: usize) -> String {
let mut sql = format!(
"SELECT key, updated_at FROM \"{0}\" \
WHERE value IS NOT NULL AND updated_at IS NOT NULL AND updated_at <= ?1",
table_name
);
if has_cursor {
sql.push_str(" AND (updated_at > ?2 OR (updated_at = ?2 AND key > ?3))");
}
sql.push_str(&format!(" ORDER BY updated_at, key LIMIT {}", limit));
sql
}
pub(super) fn build_reap_tombstones_sql(table_name: &str, limit: usize) -> String {
format!(
"DELETE FROM \"{0}\" WHERE key IN (\
SELECT key FROM \"{0}\" WHERE value IS NULL AND version <= ?1 LIMIT {1}\
) AND value IS NULL",
table_name, limit
)
}
pub(super) fn build_get_current_sql(table_name: &str) -> String {
format!("SELECT version, value FROM \"{}\" WHERE key = ?1", table_name)
}
pub(super) fn build_get_many_current_sql(table_name: &str, key_count: usize) -> String {
let placeholders = build_placeholders(key_count);
format!("SELECT key, version, value FROM \"{}\" WHERE key IN ({})", table_name, placeholders)
}
fn build_placeholders(key_count: usize) -> String {
let mut placeholders = String::with_capacity(key_count.saturating_mul(2));
for i in 0..key_count {
if i > 0 {
placeholders.push(',');
}
placeholders.push('?');
}
placeholders
}
pub(super) fn build_upsert_current_sql(table_name: &str) -> String {
format!(
"INSERT INTO \"{0}\" (key, version, value, updated_at) VALUES (?1, ?2, ?3, ?4) \
ON CONFLICT(key) DO UPDATE SET \
version = excluded.version, \
value = excluded.value, \
updated_at = excluded.updated_at \
WHERE excluded.version >= \"{0}\".version",
table_name
)
}
pub(super) fn build_chunked_upsert_sql(table_name: &str, chunk: usize) -> String {
format!(
"INSERT INTO \"{0}\" (key, version, value, updated_at) VALUES {1} \
ON CONFLICT(key) DO UPDATE SET \
version = excluded.version, \
value = excluded.value, \
updated_at = excluded.updated_at \
WHERE excluded.version >= \"{0}\".version \
RETURNING key",
table_name,
values_placeholders(chunk, 4)
)
}
pub(super) fn build_delete_below_version_sql(
table_name: &str,
has_prefix: bool,
has_cursor: bool,
limit: usize,
) -> String {
let mut inner = format!("SELECT key FROM \"{0}\" WHERE version <= ?1", table_name);
if has_prefix {
inner.push_str(" AND key >= ?2 AND key < ?3");
}
if has_cursor {
let param = if has_prefix {
4
} else {
2
};
inner.push_str(&format!(" AND key > ?{}", param));
}
inner.push_str(&format!(" ORDER BY key LIMIT {}", limit));
format!("DELETE FROM \"{0}\" WHERE key IN ({1}) RETURNING key", table_name, inner)
}
pub(super) fn prefix_upper_bound(prefix: &[u8]) -> Vec<u8> {
let mut upper = prefix.to_vec();
while let Some(last) = upper.last_mut() {
if *last < 0xFF {
*last += 1;
return upper;
}
upper.pop();
}
upper
}
pub(super) fn build_delete_keys_sql(table_name: &str, key_count: usize) -> String {
let placeholders = build_placeholders(key_count);
format!("DELETE FROM \"{}\" WHERE key IN ({})", table_name, placeholders)
}
pub(super) fn build_range_consistent_sql(table_name: &str, start: Bound<()>, end: Bound<()>) -> String {
let mut sql = format!("SELECT key, version, value FROM \"{}\" WHERE 1=1", table_name);
match start {
Bound::Included(()) => sql.push_str(" AND key >= ?"),
Bound::Excluded(()) => sql.push_str(" AND key > ?"),
Bound::Unbounded => {}
}
match end {
Bound::Included(()) => sql.push_str(" AND key <= ?"),
Bound::Excluded(()) => sql.push_str(" AND key < ?"),
Bound::Unbounded => {}
}
sql.push_str(" AND version <= ? ORDER BY key ASC");
sql
}
pub(super) fn build_range_current_sql(
table_name: &str,
start: Bound<()>,
end: Bound<()>,
has_last_key: bool,
descending: bool,
) -> String {
let mut sql = format!("SELECT key, version, value FROM \"{}\" WHERE 1=1", table_name);
match start {
Bound::Included(()) => sql.push_str(" AND key >= ?"),
Bound::Excluded(()) => sql.push_str(" AND key > ?"),
Bound::Unbounded => {}
}
match end {
Bound::Included(()) => sql.push_str(" AND key <= ?"),
Bound::Excluded(()) => sql.push_str(" AND key < ?"),
Bound::Unbounded => {}
}
if has_last_key {
sql.push_str(if descending {
" AND key < ?"
} else {
" AND key > ?"
});
}
sql.push_str(" AND value IS NOT NULL AND version <= ?");
if descending {
sql.push_str(" ORDER BY key DESC LIMIT ?");
} else {
sql.push_str(" ORDER BY key ASC LIMIT ?");
}
sql
}