use std::num::NonZeroUsize;
use std::sync::atomic::{AtomicBool, Ordering};
use foundationdb::{
RangeOption, Transaction,
future::{FdbSlice, FdbValues},
options::{MutationType, TransactionOption},
};
use thiserror::Error;
const RESERVED_METADATA_START: &[u8] = b"__meta";
const RESERVED_METADATA_END: &[u8] = b"__metb";
const MAX_TENANT_KEY_BYTES: usize = 10_000;
const MAX_TENANT_VALUE_BYTES: usize = 100_000;
const MAX_RANGE_RESULTS: usize = 64;
const MAX_RANGE_TARGET_BYTES: usize = 1_000_000;
const MAX_TRANSACTION_SIZE_LIMIT: i32 = 10_000_000;
const VERSIONSTAMP_BYTES: usize = 10;
const VERSIONSTAMP_OFFSET_BYTES: usize = 4;
const MIN_STABLE_KEY_PREFIX_BYTES: usize = 6;
#[derive(Debug, Clone, Copy, PartialEq, Eq, Error)]
pub enum TenantDataAccessErrorReason {
#[error("tenant data key is empty")]
EmptyKey,
#[error("tenant data key exceeds the size limit")]
KeyTooLong,
#[error("tenant data key is reserved")]
ReservedKey,
#[error("tenant data value exceeds the size limit")]
ValueTooLong,
#[error("tenant data range bounds are invalid")]
InvalidRange,
#[error("tenant data range overlaps reserved metadata")]
ReservedRange,
#[error("tenant data limit is invalid")]
InvalidRangeLimit,
#[error("tenant data mutation is prohibited")]
KeyChangingMutation,
}
#[derive(Clone, Copy)]
pub struct TenantTransactionSizeLimit(i32);
impl TenantTransactionSizeLimit {
pub fn new(bytes: i32) -> Result<Self, TenantDataAccessErrorReason> {
if !(1..=MAX_TRANSACTION_SIZE_LIMIT).contains(&bytes) {
return Err(TenantDataAccessErrorReason::InvalidRangeLimit);
}
Ok(Self(bytes))
}
}
#[derive(Clone, Copy)]
pub struct TenantDataRangeTargetBytes(NonZeroUsize);
impl TenantDataRangeTargetBytes {
pub fn new(bytes: usize) -> Result<Self, TenantDataAccessErrorReason> {
let bytes = NonZeroUsize::new(bytes)
.filter(|bytes| bytes.get() <= MAX_RANGE_TARGET_BYTES)
.ok_or(TenantDataAccessErrorReason::InvalidRangeLimit)?;
Ok(Self(bytes))
}
}
#[derive(Clone, Copy)]
pub struct TenantDataKey<'a> {
bytes: &'a [u8],
}
impl<'a> TenantDataKey<'a> {
pub fn new(bytes: &'a [u8]) -> Result<Self, TenantDataAccessErrorReason> {
if bytes.is_empty() {
return Err(TenantDataAccessErrorReason::EmptyKey);
}
if bytes.len() > MAX_TENANT_KEY_BYTES {
return Err(TenantDataAccessErrorReason::KeyTooLong);
}
if (RESERVED_METADATA_START..RESERVED_METADATA_END).contains(&bytes) || bytes[0] == 0xff {
return Err(TenantDataAccessErrorReason::ReservedKey);
}
Ok(Self { bytes })
}
pub fn as_bytes(self) -> &'a [u8] {
self.bytes
}
}
pub struct TenantDataRange<'a> {
begin: TenantDataKey<'a>,
end: TenantDataKey<'a>,
}
impl<'a> TenantDataRange<'a> {
pub fn new(begin: &'a [u8], end: &'a [u8]) -> Result<Self, TenantDataAccessErrorReason> {
let begin = TenantDataKey::new(begin)?;
let end = TenantDataKey::new(end)?;
if begin.as_bytes() >= end.as_bytes() {
return Err(TenantDataAccessErrorReason::InvalidRange);
}
if begin.as_bytes() < RESERVED_METADATA_END && end.as_bytes() > RESERVED_METADATA_START {
return Err(TenantDataAccessErrorReason::ReservedRange);
}
Ok(Self { begin, end })
}
pub fn begin(&self) -> &[u8] {
self.begin.as_bytes()
}
pub fn end(&self) -> &[u8] {
self.end.as_bytes()
}
}
#[derive(Clone, Copy)]
pub struct TenantDataRangeLimit(NonZeroUsize);
impl TenantDataRangeLimit {
pub const MAX: usize = MAX_RANGE_RESULTS;
pub fn new(value: usize) -> Result<Self, TenantDataAccessErrorReason> {
let value = NonZeroUsize::new(value)
.filter(|value| value.get() <= MAX_RANGE_RESULTS)
.ok_or(TenantDataAccessErrorReason::InvalidRangeLimit)?;
Ok(Self(value))
}
pub fn get(self) -> usize {
self.0.get()
}
}
pub struct TenantDataTransaction<'a> {
inner: &'a Transaction,
writable: bool,
mutation_rejected: AtomicBool,
}
impl<'a> TenantDataTransaction<'a> {
pub(super) fn new(inner: &'a Transaction, writable: bool) -> Self {
Self {
inner,
writable,
mutation_rejected: AtomicBool::new(false),
}
}
pub(super) fn mutation_rejected(&self) -> bool {
self.mutation_rejected.load(Ordering::Relaxed)
}
fn require_write(&self) -> Result<(), TenantDataAccessErrorReason> {
if self.writable {
Ok(())
} else {
self.mutation_rejected.store(true, Ordering::Relaxed);
Err(TenantDataAccessErrorReason::KeyChangingMutation)
}
}
pub fn set_size_limit(&self, limit: TenantTransactionSizeLimit) -> foundationdb::FdbResult<()> {
self.inner.set_option(TransactionOption::SizeLimit(limit.0))
}
pub async fn get(
&self,
key: TenantDataKey<'_>,
snapshot: bool,
) -> foundationdb::FdbResult<Option<FdbSlice>> {
self.inner.get(key.as_bytes(), snapshot).await
}
pub fn set(
&self,
key: TenantDataKey<'_>,
value: &[u8],
) -> Result<(), TenantDataAccessErrorReason> {
self.require_write()?;
if value.len() > MAX_TENANT_VALUE_BYTES {
return Err(TenantDataAccessErrorReason::ValueTooLong);
}
self.inner.set(key.as_bytes(), value);
Ok(())
}
pub fn clear(&self, key: TenantDataKey<'_>) {
let _ = self.try_clear(key);
}
pub fn try_clear(&self, key: TenantDataKey<'_>) -> Result<(), TenantDataAccessErrorReason> {
self.require_write()?;
self.inner.clear(key.as_bytes());
Ok(())
}
pub fn clear_range(&self, range: &TenantDataRange<'_>) {
let _ = self.try_clear_range(range);
}
pub fn try_clear_range(
&self,
range: &TenantDataRange<'_>,
) -> Result<(), TenantDataAccessErrorReason> {
self.require_write()?;
self.inner.clear_range(range.begin(), range.end());
Ok(())
}
pub fn atomic_op(
&self,
key: TenantDataKey<'_>,
parameter: &[u8],
mutation: MutationType,
) -> Result<(), TenantDataAccessErrorReason> {
self.require_write()?;
validate_atomic_mutation(mutation)?;
if parameter.len() > MAX_TENANT_VALUE_BYTES {
return Err(TenantDataAccessErrorReason::ValueTooLong);
}
self.inner.atomic_op(key.as_bytes(), parameter, mutation);
Ok(())
}
pub fn set_versionstamped_key(
&self,
key_template: &[u8],
value: &[u8],
) -> Result<(), TenantDataAccessErrorReason> {
self.require_write()?;
validate_versionstamped_key_template(key_template)?;
if value.len() > MAX_TENANT_VALUE_BYTES {
return Err(TenantDataAccessErrorReason::ValueTooLong);
}
self.inner
.atomic_op(key_template, value, MutationType::SetVersionstampedKey);
Ok(())
}
pub async fn get_range(
&self,
range: &TenantDataRange<'_>,
limit: TenantDataRangeLimit,
snapshot: bool,
) -> foundationdb::FdbResult<FdbValues> {
self.get_range_with_target_bytes(range, limit, MAX_RANGE_TARGET_BYTES, snapshot)
.await
}
pub async fn get_range_with_target(
&self,
range: &TenantDataRange<'_>,
limit: TenantDataRangeLimit,
target_bytes: TenantDataRangeTargetBytes,
snapshot: bool,
) -> foundationdb::FdbResult<FdbValues> {
self.get_range_with_target_bytes(range, limit, target_bytes.0.get(), snapshot)
.await
}
async fn get_range_with_target_bytes(
&self,
range: &TenantDataRange<'_>,
limit: TenantDataRangeLimit,
target_bytes: usize,
snapshot: bool,
) -> foundationdb::FdbResult<FdbValues> {
let mut options = RangeOption::from((range.begin(), range.end()));
options.limit = Some(limit.get());
options.target_bytes = target_bytes;
self.inner.get_range(&options, 1, snapshot).await
}
}
fn validate_versionstamped_key_template(
key_template: &[u8],
) -> Result<(), TenantDataAccessErrorReason> {
let key_bytes = key_template
.len()
.checked_sub(VERSIONSTAMP_OFFSET_BYTES)
.ok_or(TenantDataAccessErrorReason::KeyChangingMutation)?;
TenantDataKey::new(&key_template[..key_bytes])?;
let offset_bytes: [u8; VERSIONSTAMP_OFFSET_BYTES] = key_template[key_bytes..]
.try_into()
.map_err(|_| TenantDataAccessErrorReason::KeyChangingMutation)?;
let offset = usize::try_from(u32::from_le_bytes(offset_bytes))
.map_err(|_| TenantDataAccessErrorReason::KeyChangingMutation)?;
let end = offset
.checked_add(VERSIONSTAMP_BYTES)
.ok_or(TenantDataAccessErrorReason::KeyChangingMutation)?;
if offset < MIN_STABLE_KEY_PREFIX_BYTES
|| end > key_bytes
|| key_template[offset..end].iter().any(|byte| *byte != 0xff)
{
return Err(TenantDataAccessErrorReason::KeyChangingMutation);
}
Ok(())
}
fn validate_atomic_mutation(mutation: MutationType) -> Result<(), TenantDataAccessErrorReason> {
match mutation {
MutationType::Add
| MutationType::And
| MutationType::BitAnd
| MutationType::Or
| MutationType::BitOr
| MutationType::Xor
| MutationType::BitXor
| MutationType::AppendIfFits
| MutationType::Max
| MutationType::Min
| MutationType::SetVersionstampedValue
| MutationType::ByteMin
| MutationType::ByteMax
| MutationType::CompareAndClear => Ok(()),
MutationType::SetVersionstampedKey => Err(TenantDataAccessErrorReason::KeyChangingMutation),
_ => Err(TenantDataAccessErrorReason::KeyChangingMutation),
}
}
#[cfg(test)]
#[path = "data_tests.rs"]
mod tests;