reifydb-store-multi 0.9.0

Multi-version storage for OLTP operations with MVCC support
Documentation
// SPDX-License-Identifier: Apache-2.0
// Copyright (c) 2026 ReifyDB

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
}