use std::path::Path;
use std::sync::{Arc, Mutex};
use dynomite::embed::hooks::{BoxFuture, Datastore, DatastoreByteStream, DatastoreError, Protocol};
use dynomite::msg::{Msg, MsgType};
use noxu::{
Cursor, CursorConfig, Database, DatabaseConfig, DatabaseEntry, Environment, EnvironmentConfig,
Get, NoxuError, OperationStatus, Transaction,
};
use crate::txn::{TransactionalStore, TxnBatch, TxnOp, TxnOutcome, TxnStoreError};
const PRIMARY_TAG: &[u8] = b"K\0";
const FWD_TAG: &[u8] = b"I\0";
const REV_TAG: &[u8] = b"R\0";
const SEP: u8 = 0;
pub const INT_SUFFIX: &[u8] = b"_int";
pub const BIN_SUFFIX: &[u8] = b"_bin";
const INT_ENCODED_WIDTH: usize = 8;
#[derive(Debug, thiserror::Error)]
#[non_exhaustive]
pub enum NoxuDatastoreError {
#[error("noxu: {0}")]
Noxu(#[from] NoxuError),
#[error("noxu: invalid name: {what} contains a NUL byte")]
InvalidName {
what: &'static str,
},
#[error("noxu: bad _int index value: expected 8 bytes, got {got}")]
BadIntValue {
got: usize,
},
#[error("noxu: corrupt reverse index entry for bucket={bucket:?}, key={key:?}")]
CorruptReverse {
bucket: Vec<u8>,
key: Vec<u8>,
},
#[error("noxu: datastore is not transactional; open with open_transactional")]
NotTransactional,
#[error("noxu xa: {0}")]
Xa(String),
}
#[derive(Clone)]
pub struct NoxuDatastore {
inner: Arc<Inner>,
}
struct Inner {
env: Mutex<Environment>,
db: Mutex<Database>,
transactional: bool,
}
#[derive(Clone, Debug)]
pub struct EncodedIndexValue {
pub bytes: Vec<u8>,
}
impl NoxuDatastore {
pub fn open(path: &Path) -> Result<Self, NoxuDatastoreError> {
Self::open_with_db_name(path, "riak.objects")
}
pub fn open_transactional(path: &Path) -> Result<Self, NoxuDatastoreError> {
Self::open_with_options(path, "riak.objects", true)
}
pub fn open_with_db_name(path: &Path, db_name: &str) -> Result<Self, NoxuDatastoreError> {
Self::open_with_options(path, db_name, false)
}
fn open_with_options(
path: &Path,
db_name: &str,
transactional: bool,
) -> Result<Self, NoxuDatastoreError> {
let env_config = EnvironmentConfig::new(path.to_path_buf())
.with_allow_create(true)
.with_transactional(transactional);
let env = Environment::open(env_config)?;
let db_config = DatabaseConfig::new()
.with_allow_create(true)
.with_transactional(transactional);
let db = env.open_database(None, db_name, &db_config)?;
Ok(Self {
inner: Arc::new(Inner {
env: Mutex::new(env),
db: Mutex::new(db),
transactional,
}),
})
}
pub fn open_in(dir: &Path) -> Result<Self, NoxuDatastoreError> {
Self::open(dir)
}
pub fn get(&self, key: &[u8]) -> Result<Option<Vec<u8>>, NoxuDatastoreError> {
let mut value = DatabaseEntry::new();
let db = self.lock_db();
if db_get(&db, None, key, &mut value)? {
Ok(Some(value.data().to_vec()))
} else {
Ok(None)
}
}
pub fn put(&self, key: &[u8], value: &[u8]) -> Result<(), NoxuDatastoreError> {
let db = self.lock_db();
db_put(&db, None, key, value)?;
Ok(())
}
pub fn delete(&self, key: &[u8]) -> Result<bool, NoxuDatastoreError> {
let db = self.lock_db();
db_delete(&db, None, key)
}
pub fn get_object(
&self,
bucket: &[u8],
key: &[u8],
) -> Result<Option<Vec<u8>>, NoxuDatastoreError> {
let db = self.lock_db();
get_object_in(&db, None, bucket, key)
}
pub fn put_object(
&self,
bucket: &[u8],
key: &[u8],
value: &[u8],
indexes: &[(Vec<u8>, Vec<u8>)],
) -> Result<(), NoxuDatastoreError> {
let db = self.lock_db();
put_object_in(&db, None, bucket, key, value, indexes)
}
pub fn delete_object(&self, bucket: &[u8], key: &[u8]) -> Result<bool, NoxuDatastoreError> {
let db = self.lock_db();
delete_object_in(&db, None, bucket, key)
}
pub fn index_eq(
&self,
bucket: &[u8],
index_name: &[u8],
value: &[u8],
) -> Result<Vec<Vec<u8>>, NoxuDatastoreError> {
validate_no_separator(bucket, "bucket")?;
validate_no_separator(index_name, "index_name")?;
let encoded = encode_index_value(index_name, value)?;
let prefix = forward_prefix_with_value(bucket, index_name, &encoded);
let db = self.lock_db();
let mut results = Vec::new();
scan_prefix(&db, None, &prefix, |full_key, _val| {
let suffix = &full_key[prefix.len()..];
results.push(suffix.to_vec());
Ok(())
})?;
Ok(results)
}
pub fn index_range(
&self,
bucket: &[u8],
index_name: &[u8],
range_min: &[u8],
range_max: &[u8],
) -> Result<Vec<Vec<u8>>, NoxuDatastoreError> {
validate_no_separator(bucket, "bucket")?;
validate_no_separator(index_name, "index_name")?;
let min_enc = encode_index_value(index_name, range_min)?;
let max_enc = encode_index_value(index_name, range_max)?;
let prefix = forward_prefix(bucket, index_name);
let db = self.lock_db();
let mut results = Vec::new();
scan_prefix(&db, None, &prefix, |full_key, _val| {
let suffix = &full_key[prefix.len()..];
if suffix.len() < 4 {
return Ok(());
}
let mut len_bytes = [0u8; 4];
len_bytes.copy_from_slice(&suffix[..4]);
let value_len = u32::from_be_bytes(len_bytes) as usize;
if suffix.len() < 4 + value_len {
return Ok(());
}
let value_bytes = &suffix[4..4 + value_len];
if value_bytes < min_enc.as_slice() || value_bytes > max_enc.as_slice() {
return Ok(());
}
let key_bytes = &suffix[4 + value_len..];
results.push(key_bytes.to_vec());
Ok(())
})?;
Ok(results)
}
fn lock_db(&self) -> std::sync::MutexGuard<'_, Database> {
self.inner
.db
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
}
pub fn fold_primary<F>(&self, mut f: F) -> Result<(), NoxuDatastoreError>
where
F: FnMut(&[u8], &[u8], &[u8]) -> Result<(), NoxuDatastoreError>,
{
let db = self.lock_db();
scan_prefix(&db, None, PRIMARY_TAG, |full_key, value| {
let suffix = &full_key[PRIMARY_TAG.len()..];
let Some(sep_idx) = suffix.iter().position(|b| *b == SEP) else {
return Ok(());
};
let (bucket, rest) = suffix.split_at(sep_idx);
let key = &rest[1..];
f(bucket, key, value)
})
}
fn lock_env(&self) -> std::sync::MutexGuard<'_, Environment> {
self.inner
.env
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
}
pub fn transaction<F, R, E>(&self, f: F) -> Result<R, E>
where
F: FnOnce(&NoxuTxn<'_>) -> Result<R, E>,
E: From<NoxuDatastoreError>,
{
if !self.inner.transactional {
return Err(E::from(NoxuDatastoreError::NotTransactional));
}
let txn = {
let env = self.lock_env();
env.begin_transaction(None)
.map_err(NoxuDatastoreError::from)?
};
let db = self.lock_db();
let handle = NoxuTxn { db: &db, txn: &txn };
match f(&handle) {
Ok(value) => {
txn.commit().map_err(NoxuDatastoreError::from)?;
Ok(value)
}
Err(err) => {
let _ = txn.abort();
Err(err)
}
}
}
}
pub struct NoxuTxn<'a> {
db: &'a Database,
txn: &'a Transaction,
}
impl NoxuTxn<'_> {
pub fn put_object(
&self,
bucket: &[u8],
key: &[u8],
value: &[u8],
indexes: &[(Vec<u8>, Vec<u8>)],
) -> Result<(), NoxuDatastoreError> {
put_object_in(self.db, Some(self.txn), bucket, key, value, indexes)
}
pub fn delete_object(&self, bucket: &[u8], key: &[u8]) -> Result<bool, NoxuDatastoreError> {
delete_object_in(self.db, Some(self.txn), bucket, key)
}
pub fn get_object(
&self,
bucket: &[u8],
key: &[u8],
) -> Result<Option<Vec<u8>>, NoxuDatastoreError> {
get_object_in(self.db, Some(self.txn), bucket, key)
}
pub fn apply(&self, op: &TxnOp) -> Result<(), NoxuDatastoreError> {
match op {
TxnOp::Put {
bucket,
key,
value,
indexes,
} => self.put_object(bucket, key, value, indexes),
TxnOp::Delete { bucket, key } => self.delete_object(bucket, key).map(|_| ()),
}
}
}
impl Datastore for NoxuDatastore {
fn protocol(&self) -> Protocol {
Protocol::Custom
}
fn as_any(&self) -> Option<&(dyn std::any::Any + 'static)> {
Some(self)
}
fn dispatch(&self, req: Msg) -> BoxFuture<'_, Result<Msg, DatastoreError>> {
Box::pin(async move {
let mut rsp = Msg::new(req.id(), MsgType::Unknown, false);
rsp.set_parent_id(req.id());
Ok(rsp)
})
}
fn riak_get<'a>(
&'a self,
bucket: &'a [u8],
key: &'a [u8],
) -> BoxFuture<'a, Result<Option<Vec<u8>>, DatastoreError>> {
Box::pin(async move {
self.get_object(bucket, key)
.map_err(|e| DatastoreError::Backend(e.to_string()))
})
}
fn riak_put<'a>(
&'a self,
bucket: &'a [u8],
key: &'a [u8],
value: &'a [u8],
indexes: &'a [(Vec<u8>, Vec<u8>)],
) -> BoxFuture<'a, Result<(), DatastoreError>> {
Box::pin(async move {
self.put_object(bucket, key, value, indexes)
.map_err(|e| DatastoreError::Backend(e.to_string()))
})
}
fn riak_delete<'a>(
&'a self,
bucket: &'a [u8],
key: &'a [u8],
) -> BoxFuture<'a, Result<bool, DatastoreError>> {
Box::pin(async move {
self.delete_object(bucket, key)
.map_err(|e| DatastoreError::Backend(e.to_string()))
})
}
fn riak_index_eq<'a>(
&'a self,
bucket: &'a [u8],
index_name: &'a [u8],
value: &'a [u8],
) -> BoxFuture<'a, Result<Vec<Vec<u8>>, DatastoreError>> {
Box::pin(async move {
self.index_eq(bucket, index_name, value)
.map_err(|e| DatastoreError::Backend(e.to_string()))
})
}
fn riak_index_range<'a>(
&'a self,
bucket: &'a [u8],
index_name: &'a [u8],
min: &'a [u8],
max: &'a [u8],
) -> BoxFuture<'a, Result<Vec<Vec<u8>>, DatastoreError>> {
Box::pin(async move {
self.index_range(bucket, index_name, min, max)
.map_err(|e| DatastoreError::Backend(e.to_string()))
})
}
fn list_keys_stream(&self, bucket: &[u8]) -> DatastoreByteStream {
let want = bucket.to_vec();
let mut keys: Vec<bytes::Bytes> = Vec::new();
let result = self.fold_primary(|b, k, _v| {
if b == want.as_slice() {
keys.push(bytes::Bytes::copy_from_slice(k));
}
Ok(())
});
match result {
Ok(()) => Box::pin(futures_util::stream::iter(keys.into_iter().map(Ok))),
Err(e) => Box::pin(futures_util::stream::once(async move {
Err(DatastoreError::Backend(e.to_string()))
})),
}
}
}
pub(crate) fn map_txn_store_error(err: &NoxuDatastoreError) -> TxnStoreError {
if let NoxuDatastoreError::Noxu(inner) = err {
let text = inner.to_string();
let lowered = text.to_ascii_lowercase();
if lowered.contains("deadlock")
|| lowered.contains("lock timeout")
|| lowered.contains("locktimeout")
|| lowered.contains("conflict")
|| lowered.contains("would block")
{
return TxnStoreError::Conflict(text);
}
}
TxnStoreError::Backend(err.to_string())
}
enum BatchErr {
Abort,
Store(NoxuDatastoreError),
}
impl From<NoxuDatastoreError> for BatchErr {
fn from(e: NoxuDatastoreError) -> Self {
Self::Store(e)
}
}
impl TransactionalStore for NoxuDatastore {
fn execute_batch(&self, batch: &TxnBatch) -> Result<TxnOutcome, TxnStoreError> {
if batch.ops.is_empty() {
return Err(TxnStoreError::EmptyBatch);
}
let force_abort = batch.force_abort;
let result: Result<(), BatchErr> = self.transaction(|tx| {
for op in &batch.ops {
tx.apply(op).map_err(BatchErr::Store)?;
}
if force_abort {
return Err(BatchErr::Abort);
}
Ok(())
});
match result {
Ok(()) => Ok(TxnOutcome::Committed {
operations: batch.ops.len(),
}),
Err(BatchErr::Abort) => Ok(TxnOutcome::Aborted {
reason: "client requested abort".to_string(),
}),
Err(BatchErr::Store(e)) => Err(map_txn_store_error(&e)),
}
}
}
fn db_get(
db: &Database,
txn: Option<&Transaction>,
key: &[u8],
out: &mut DatabaseEntry,
) -> Result<bool, NoxuDatastoreError> {
Ok(db.get_into(txn, key, out)?)
}
fn db_put(
db: &Database,
txn: Option<&Transaction>,
key: &[u8],
value: &[u8],
) -> Result<(), NoxuDatastoreError> {
match txn {
Some(t) => db.put_in(t, key, value)?,
None => db.put(key, value)?,
}
Ok(())
}
fn db_delete(
db: &Database,
txn: Option<&Transaction>,
key: &[u8],
) -> Result<bool, NoxuDatastoreError> {
Ok(match txn {
Some(t) => db.delete_in(t, key)?,
None => db.delete(key)?,
})
}
pub(crate) fn get_object_in(
db: &Database,
txn: Option<&Transaction>,
bucket: &[u8],
key: &[u8],
) -> Result<Option<Vec<u8>>, NoxuDatastoreError> {
validate_no_separator(bucket, "bucket")?;
let storage_key = primary_key(bucket, key);
let mut value = DatabaseEntry::new();
if db_get(db, txn, &storage_key, &mut value)? {
Ok(Some(value.data().to_vec()))
} else {
Ok(None)
}
}
pub(crate) fn put_object_in(
db: &Database,
txn: Option<&Transaction>,
bucket: &[u8],
key: &[u8],
value: &[u8],
indexes: &[(Vec<u8>, Vec<u8>)],
) -> Result<(), NoxuDatastoreError> {
validate_no_separator(bucket, "bucket")?;
for (name, _) in indexes {
validate_no_separator(name, "index_name")?;
}
let encoded: Vec<(Vec<u8>, Vec<u8>)> = indexes
.iter()
.map(|(name, value)| {
let enc = encode_index_value(name, value)?;
Ok::<_, NoxuDatastoreError>((name.clone(), enc))
})
.collect::<Result<_, _>>()?;
clear_forward_for(db, txn, bucket, key)?;
let primary = primary_key(bucket, key);
db_put(db, txn, &primary, value)?;
for (name, enc) in &encoded {
let fk = forward_key(bucket, name, enc, key);
db_put(db, txn, &fk, b"")?;
}
let rk = reverse_key(bucket, key);
if encoded.is_empty() {
db_delete(db, txn, &rk)?;
} else {
let rv = encode_reverse_value(&encoded);
db_put(db, txn, &rk, &rv)?;
}
Ok(())
}
pub(crate) fn delete_object_in(
db: &Database,
txn: Option<&Transaction>,
bucket: &[u8],
key: &[u8],
) -> Result<bool, NoxuDatastoreError> {
validate_no_separator(bucket, "bucket")?;
let removed_2i = clear_forward_for(db, txn, bucket, key)?;
let primary = primary_key(bucket, key);
let removed_pri = db_delete(db, txn, &primary)?;
Ok(removed_pri || removed_2i)
}
fn clear_forward_for(
db: &Database,
txn: Option<&Transaction>,
bucket: &[u8],
key: &[u8],
) -> Result<bool, NoxuDatastoreError> {
let rk = reverse_key(bucket, key);
let mut value = DatabaseEntry::new();
if db_get(db, txn, &rk, &mut value)? {
let pairs = decode_reverse_value(value.data()).ok_or_else(|| {
NoxuDatastoreError::CorruptReverse {
bucket: bucket.to_vec(),
key: key.to_vec(),
}
})?;
for (name, encoded) in &pairs {
let fk = forward_key(bucket, name, encoded, key);
db_delete(db, txn, &fk)?;
}
db_delete(db, txn, &rk)?;
Ok(true)
} else {
Ok(false)
}
}
fn primary_key(bucket: &[u8], key: &[u8]) -> Vec<u8> {
let mut out = Vec::with_capacity(PRIMARY_TAG.len() + bucket.len() + 1 + key.len());
out.extend_from_slice(PRIMARY_TAG);
out.extend_from_slice(bucket);
out.push(SEP);
out.extend_from_slice(key);
out
}
fn forward_prefix(bucket: &[u8], index_name: &[u8]) -> Vec<u8> {
let mut out = Vec::with_capacity(FWD_TAG.len() + bucket.len() + 1 + index_name.len() + 1);
out.extend_from_slice(FWD_TAG);
out.extend_from_slice(bucket);
out.push(SEP);
out.extend_from_slice(index_name);
out.push(SEP);
out
}
fn forward_prefix_with_value(bucket: &[u8], index_name: &[u8], encoded_value: &[u8]) -> Vec<u8> {
let mut out = forward_prefix(bucket, index_name);
out.extend_from_slice(
&u32::try_from(encoded_value.len())
.unwrap_or(u32::MAX)
.to_be_bytes(),
);
out.extend_from_slice(encoded_value);
out
}
fn forward_key(bucket: &[u8], index_name: &[u8], encoded_value: &[u8], key: &[u8]) -> Vec<u8> {
let mut out = forward_prefix_with_value(bucket, index_name, encoded_value);
out.extend_from_slice(key);
out
}
fn reverse_key(bucket: &[u8], key: &[u8]) -> Vec<u8> {
let mut out = Vec::with_capacity(REV_TAG.len() + bucket.len() + 1 + key.len());
out.extend_from_slice(REV_TAG);
out.extend_from_slice(bucket);
out.push(SEP);
out.extend_from_slice(key);
out
}
fn encode_reverse_value(pairs: &[(Vec<u8>, Vec<u8>)]) -> Vec<u8> {
let mut out = Vec::new();
out.extend_from_slice(&u32::try_from(pairs.len()).unwrap_or(u32::MAX).to_be_bytes());
for (name, value) in pairs {
out.extend_from_slice(&u32::try_from(name.len()).unwrap_or(u32::MAX).to_be_bytes());
out.extend_from_slice(name);
out.extend_from_slice(&u32::try_from(value.len()).unwrap_or(u32::MAX).to_be_bytes());
out.extend_from_slice(value);
}
out
}
fn decode_reverse_value(buf: &[u8]) -> Option<Vec<(Vec<u8>, Vec<u8>)>> {
let mut p = 0usize;
let count = read_u32_be(buf, &mut p)?;
let mut out = Vec::with_capacity(count as usize);
for _ in 0..count {
let name_len = read_u32_be(buf, &mut p)? as usize;
if p + name_len > buf.len() {
return None;
}
let name = buf[p..p + name_len].to_vec();
p += name_len;
let value_len = read_u32_be(buf, &mut p)? as usize;
if p + value_len > buf.len() {
return None;
}
let value = buf[p..p + value_len].to_vec();
p += value_len;
out.push((name, value));
}
Some(out)
}
fn read_u32_be(buf: &[u8], p: &mut usize) -> Option<u32> {
if *p + 4 > buf.len() {
return None;
}
let mut bytes = [0u8; 4];
bytes.copy_from_slice(&buf[*p..*p + 4]);
*p += 4;
Some(u32::from_be_bytes(bytes))
}
fn encode_index_value(index_name: &[u8], value: &[u8]) -> Result<Vec<u8>, NoxuDatastoreError> {
if name_is_int(index_name) {
if value.len() == INT_ENCODED_WIDTH {
return Ok(value.to_vec());
}
let s = std::str::from_utf8(value)
.map_err(|_| NoxuDatastoreError::BadIntValue { got: value.len() })?;
let parsed: u64 = s
.parse()
.map_err(|_| NoxuDatastoreError::BadIntValue { got: value.len() })?;
Ok(parsed.to_be_bytes().to_vec())
} else {
Ok(value.to_vec())
}
}
fn name_is_int(index_name: &[u8]) -> bool {
index_name.ends_with(INT_SUFFIX)
}
fn validate_no_separator(value: &[u8], what: &'static str) -> Result<(), NoxuDatastoreError> {
if value.contains(&SEP) {
return Err(NoxuDatastoreError::InvalidName { what });
}
Ok(())
}
fn scan_prefix<F>(
db: &Database,
txn: Option<&Transaction>,
prefix: &[u8],
mut f: F,
) -> Result<(), NoxuDatastoreError>
where
F: FnMut(&[u8], &[u8]) -> Result<(), NoxuDatastoreError>,
{
let mut scan = |cursor: &mut Cursor<'_>| -> Result<(), NoxuDatastoreError> {
let mut key = DatabaseEntry::from_bytes(prefix);
let mut value = DatabaseEntry::new();
let mut status = cursor.get(&mut key, &mut value, Get::SearchGte, None)?;
while matches!(status, OperationStatus::Success) {
let k = key.data();
if !k.starts_with(prefix) {
break;
}
f(k, value.data())?;
status = cursor.get(&mut key, &mut value, Get::Next, None)?;
}
let _ = cursor.close();
Ok(())
};
if let Some(t) = txn {
let mut cursor = db.open_cursor_in(t, Some(&CursorConfig::new()))?;
scan(&mut cursor)
} else {
let mut cursor = db.open_cursor(Some(&CursorConfig::new()))?;
scan(&mut cursor)
}
}
#[doc(hidden)]
pub fn encode_index_value_for_test(
index_name: &[u8],
value: &[u8],
) -> Result<Vec<u8>, NoxuDatastoreError> {
encode_index_value(index_name, value)
}
#[cfg(test)]
mod tests {
use super::*;
use tempfile::TempDir;
#[test]
fn put_get_delete_round_trips() {
let dir = TempDir::new().expect("tempdir");
let ds = NoxuDatastore::open_in(dir.path()).expect("open");
ds.put(b"alpha", b"one").expect("put alpha");
ds.put(b"beta", b"two").expect("put beta");
assert_eq!(ds.get(b"alpha").unwrap().as_deref(), Some(&b"one"[..]));
assert_eq!(ds.get(b"beta").unwrap().as_deref(), Some(&b"two"[..]));
assert_eq!(ds.get(b"missing").unwrap(), None);
assert!(ds.delete(b"alpha").unwrap());
assert!(!ds.delete(b"alpha").unwrap());
assert_eq!(ds.get(b"alpha").unwrap(), None);
assert_eq!(ds.get(b"beta").unwrap().as_deref(), Some(&b"two"[..]));
}
#[test]
fn put_overwrites_existing_value() {
let dir = TempDir::new().expect("tempdir");
let ds = NoxuDatastore::open_in(dir.path()).expect("open");
ds.put(b"k", b"v1").expect("put v1");
ds.put(b"k", b"v2").expect("put v2");
assert_eq!(ds.get(b"k").unwrap().as_deref(), Some(&b"v2"[..]));
}
#[test]
fn open_transactional_round_trips() {
let dir = TempDir::new().expect("tempdir");
let ds = NoxuDatastore::open_transactional(dir.path()).expect("open transactional");
ds.put_object(
b"users",
b"alice",
b"hello",
&[(b"age_int".to_vec(), b"42".to_vec())],
)
.expect("put");
assert_eq!(
ds.get_object(b"users", b"alice").unwrap().as_deref(),
Some(&b"hello"[..])
);
assert_eq!(
ds.index_eq(b"users", b"age_int", b"42").unwrap(),
vec![b"alice".to_vec()]
);
}
#[tokio::test]
async fn datastore_dispatch_is_a_trampoline() {
let dir = TempDir::new().expect("tempdir");
let ds = NoxuDatastore::open_in(dir.path()).expect("open");
let req = Msg::new(7, MsgType::Unknown, true);
let rsp = <NoxuDatastore as Datastore>::dispatch(&ds, req)
.await
.expect("dispatch");
assert_eq!(rsp.parent_id(), 7);
}
#[test]
fn put_object_round_trips() {
let dir = TempDir::new().expect("tempdir");
let ds = NoxuDatastore::open_in(dir.path()).expect("open");
ds.put_object(b"users", b"alice", b"hello", &[])
.expect("put");
assert_eq!(
ds.get_object(b"users", b"alice").unwrap().as_deref(),
Some(&b"hello"[..])
);
assert!(ds.get_object(b"users", b"bob").unwrap().is_none());
assert!(ds.delete_object(b"users", b"alice").unwrap());
assert!(ds.get_object(b"users", b"alice").unwrap().is_none());
}
#[test]
fn index_eq_returns_matching_keys() {
let dir = TempDir::new().expect("tempdir");
let ds = NoxuDatastore::open_in(dir.path()).expect("open");
ds.put_object(
b"users",
b"alice",
b"v1",
&[(b"age_int".to_vec(), b"42".to_vec())],
)
.expect("put alice");
ds.put_object(
b"users",
b"bob",
b"v2",
&[(b"age_int".to_vec(), b"42".to_vec())],
)
.expect("put bob");
ds.put_object(
b"users",
b"carol",
b"v3",
&[(b"age_int".to_vec(), b"99".to_vec())],
)
.expect("put carol");
let mut hits = ds.index_eq(b"users", b"age_int", b"42").unwrap();
hits.sort();
assert_eq!(hits, vec![b"alice".to_vec(), b"bob".to_vec()]);
let hits99 = ds.index_eq(b"users", b"age_int", b"99").unwrap();
assert_eq!(hits99, vec![b"carol".to_vec()]);
let none = ds.index_eq(b"users", b"age_int", b"7").unwrap();
assert!(none.is_empty());
}
#[test]
fn index_eq_handles_bin_indexes() {
let dir = TempDir::new().expect("tempdir");
let ds = NoxuDatastore::open_in(dir.path()).expect("open");
ds.put_object(
b"users",
b"alice",
b"v1",
&[(b"city_bin".to_vec(), b"seattle".to_vec())],
)
.unwrap();
ds.put_object(
b"users",
b"bob",
b"v2",
&[(b"city_bin".to_vec(), b"portland".to_vec())],
)
.unwrap();
let hits = ds.index_eq(b"users", b"city_bin", b"seattle").unwrap();
assert_eq!(hits, vec![b"alice".to_vec()]);
}
#[test]
fn index_range_returns_keys_in_range() {
let dir = TempDir::new().expect("tempdir");
let ds = NoxuDatastore::open_in(dir.path()).expect("open");
for (key, age) in [
(&b"alice"[..], "10"),
(&b"bob"[..], "15"),
(&b"carol"[..], "20"),
(&b"dave"[..], "25"),
(&b"erin"[..], "30"),
] {
ds.put_object(
b"users",
key,
b"v",
&[(b"age_int".to_vec(), age.as_bytes().to_vec())],
)
.unwrap();
}
let mut hits = ds.index_range(b"users", b"age_int", b"15", b"25").unwrap();
hits.sort();
assert_eq!(
hits,
vec![b"bob".to_vec(), b"carol".to_vec(), b"dave".to_vec()]
);
}
#[test]
fn delete_object_clears_2i_entries() {
let dir = TempDir::new().expect("tempdir");
let ds = NoxuDatastore::open_in(dir.path()).expect("open");
ds.put_object(
b"users",
b"alice",
b"v",
&[
(b"age_int".to_vec(), b"42".to_vec()),
(b"city_bin".to_vec(), b"seattle".to_vec()),
],
)
.unwrap();
assert_eq!(
ds.index_eq(b"users", b"age_int", b"42").unwrap(),
vec![b"alice".to_vec()]
);
assert!(ds.delete_object(b"users", b"alice").unwrap());
assert!(ds.index_eq(b"users", b"age_int", b"42").unwrap().is_empty());
assert!(ds
.index_eq(b"users", b"city_bin", b"seattle")
.unwrap()
.is_empty());
assert!(ds.get_object(b"users", b"alice").unwrap().is_none());
}
#[test]
fn update_replaces_index_entries() {
let dir = TempDir::new().expect("tempdir");
let ds = NoxuDatastore::open_in(dir.path()).expect("open");
ds.put_object(
b"users",
b"alice",
b"v1",
&[(b"age_int".to_vec(), b"30".to_vec())],
)
.unwrap();
ds.put_object(
b"users",
b"alice",
b"v2",
&[(b"age_int".to_vec(), b"31".to_vec())],
)
.unwrap();
assert!(ds.index_eq(b"users", b"age_int", b"30").unwrap().is_empty());
assert_eq!(
ds.index_eq(b"users", b"age_int", b"31").unwrap(),
vec![b"alice".to_vec()]
);
}
#[test]
fn invalid_bucket_with_separator_rejected() {
let dir = TempDir::new().expect("tempdir");
let ds = NoxuDatastore::open_in(dir.path()).expect("open");
let bad = vec![b'b', 0u8, b'a', b'd'];
let err = ds.put_object(&bad, b"k", b"v", &[]);
assert!(matches!(
err,
Err(NoxuDatastoreError::InvalidName { what: "bucket" })
));
}
#[test]
fn bad_int_value_rejected() {
let dir = TempDir::new().expect("tempdir");
let ds = NoxuDatastore::open_in(dir.path()).expect("open");
let err = ds.put_object(
b"users",
b"alice",
b"v",
&[(b"age_int".to_vec(), b"not-a-number".to_vec())],
);
assert!(matches!(err, Err(NoxuDatastoreError::BadIntValue { .. })));
}
#[test]
fn reverse_codec_round_trip() {
let pairs = vec![
(b"a_int".to_vec(), 42u64.to_be_bytes().to_vec()),
(b"b_bin".to_vec(), b"hello".to_vec()),
];
let buf = encode_reverse_value(&pairs);
assert_eq!(decode_reverse_value(&buf), Some(pairs));
assert_eq!(decode_reverse_value(&[]), None);
}
#[test]
fn transaction_commits_multiple_keys_atomically() {
let dir = TempDir::new().expect("tempdir");
let ds = NoxuDatastore::open_transactional(dir.path()).expect("open");
ds.transaction(|tx| {
tx.put_object(b"users", b"alice", b"a", &[])?;
tx.put_object(b"users", b"bob", b"b", &[])?;
tx.put_object(b"users", b"carol", b"c", &[])?;
Ok::<_, NoxuDatastoreError>(())
})
.expect("commit");
assert_eq!(
ds.get_object(b"users", b"alice").unwrap().as_deref(),
Some(&b"a"[..])
);
assert_eq!(
ds.get_object(b"users", b"bob").unwrap().as_deref(),
Some(&b"b"[..])
);
assert_eq!(
ds.get_object(b"users", b"carol").unwrap().as_deref(),
Some(&b"c"[..])
);
}
#[test]
fn transaction_rolls_back_every_write_on_error() {
let dir = TempDir::new().expect("tempdir");
let ds = NoxuDatastore::open_transactional(dir.path()).expect("open");
ds.put_object(b"users", b"alice", b"original", &[])
.expect("seed");
let result: Result<(), NoxuDatastoreError> = ds.transaction(|tx| {
tx.put_object(b"users", b"alice", b"changed", &[])?;
tx.put_object(b"users", b"bob", b"b", &[])?;
Err(NoxuDatastoreError::NotTransactional)
});
assert!(result.is_err());
assert_eq!(
ds.get_object(b"users", b"alice").unwrap().as_deref(),
Some(&b"original"[..])
);
assert!(ds.get_object(b"users", b"bob").unwrap().is_none());
}
#[test]
fn transaction_requires_transactional_environment() {
let dir = TempDir::new().expect("tempdir");
let ds = NoxuDatastore::open_in(dir.path()).expect("open");
let result: Result<(), NoxuDatastoreError> =
ds.transaction(|tx| tx.put_object(b"users", b"alice", b"v", &[]));
assert!(matches!(result, Err(NoxuDatastoreError::NotTransactional)));
}
#[test]
fn transaction_2i_records_commit_and_roll_back_atomically() {
let dir = TempDir::new().expect("tempdir");
let ds = NoxuDatastore::open_transactional(dir.path()).expect("open");
ds.transaction(|tx| {
tx.put_object(
b"users",
b"alice",
b"v",
&[(b"age_int".to_vec(), b"42".to_vec())],
)
})
.expect("commit");
assert_eq!(
ds.get_object(b"users", b"alice").unwrap().as_deref(),
Some(&b"v"[..])
);
assert_eq!(
ds.index_eq(b"users", b"age_int", b"42").unwrap(),
vec![b"alice".to_vec()]
);
let result: Result<(), NoxuDatastoreError> = ds.transaction(|tx| {
tx.put_object(
b"users",
b"bob",
b"v",
&[(b"age_int".to_vec(), b"99".to_vec())],
)?;
Err(NoxuDatastoreError::NotTransactional)
});
assert!(result.is_err());
assert!(ds.get_object(b"users", b"bob").unwrap().is_none());
assert!(ds.index_eq(b"users", b"age_int", b"99").unwrap().is_empty());
}
#[test]
fn execute_batch_commits_and_reports_count() {
let dir = TempDir::new().expect("tempdir");
let ds = NoxuDatastore::open_transactional(dir.path()).expect("open");
let batch = TxnBatch {
ops: vec![
TxnOp::Put {
bucket: b"users".to_vec(),
key: b"alice".to_vec(),
value: b"a".to_vec(),
indexes: vec![(b"age_int".to_vec(), b"42".to_vec())],
},
TxnOp::Put {
bucket: b"users".to_vec(),
key: b"bob".to_vec(),
value: b"b".to_vec(),
indexes: vec![],
},
],
force_abort: false,
};
let outcome = ds.execute_batch(&batch).expect("commit");
assert_eq!(outcome, TxnOutcome::Committed { operations: 2 });
assert_eq!(
ds.get_object(b"users", b"alice").unwrap().as_deref(),
Some(&b"a"[..])
);
assert_eq!(
ds.index_eq(b"users", b"age_int", b"42").unwrap(),
vec![b"alice".to_vec()]
);
}
#[test]
fn execute_batch_force_abort_leaves_no_writes() {
let dir = TempDir::new().expect("tempdir");
let ds = NoxuDatastore::open_transactional(dir.path()).expect("open");
let batch = TxnBatch {
ops: vec![
TxnOp::Put {
bucket: b"users".to_vec(),
key: b"alice".to_vec(),
value: b"a".to_vec(),
indexes: vec![(b"age_int".to_vec(), b"42".to_vec())],
},
TxnOp::Delete {
bucket: b"users".to_vec(),
key: b"ghost".to_vec(),
},
],
force_abort: true,
};
let outcome = ds.execute_batch(&batch).expect("abort is not an error");
assert!(matches!(outcome, TxnOutcome::Aborted { .. }));
assert!(ds.get_object(b"users", b"alice").unwrap().is_none());
assert!(ds.index_eq(b"users", b"age_int", b"42").unwrap().is_empty());
}
#[test]
fn execute_batch_rejects_empty() {
let dir = TempDir::new().expect("tempdir");
let ds = NoxuDatastore::open_transactional(dir.path()).expect("open");
let batch = TxnBatch::default();
assert!(matches!(
ds.execute_batch(&batch),
Err(TxnStoreError::EmptyBatch)
));
}
#[test]
fn execute_batch_as_any_downcasts_to_transactional_store() {
let dir = TempDir::new().expect("tempdir");
let ds = NoxuDatastore::open_transactional(dir.path()).expect("open");
let erased: Arc<dyn Datastore> = Arc::new(ds);
let any = <dyn Datastore>::as_any(erased.as_ref()).expect("as_any");
let store = any
.downcast_ref::<NoxuDatastore>()
.expect("downcast to NoxuDatastore");
let batch = TxnBatch {
ops: vec![TxnOp::Put {
bucket: b"users".to_vec(),
key: b"alice".to_vec(),
value: b"a".to_vec(),
indexes: vec![],
}],
force_abort: false,
};
let outcome = store.execute_batch(&batch).expect("commit");
assert_eq!(outcome, TxnOutcome::Committed { operations: 1 });
}
}