use core::future::Future;
use bytes::Bytes;
use super::error::StoreError;
use super::keys;
use super::partition::Partition;
use crate::rt::{MaybeSend, MaybeSync};
pub const MAX_KEY_BYTES: usize = 1024;
pub const MAX_VALUE_BYTES: usize = 512 * 1024;
pub const MAX_BATCH_OPS: usize = 100;
pub const MAX_BATCH_BYTES: usize = 1024 * 1024;
pub(crate) const MAX_TIMER_VALUE_BYTES: usize = (MAX_BATCH_BYTES - 4 * MAX_KEY_BYTES) / 2;
macro_rules! bytes_newtype {
($(#[$doc:meta])* $name:ident) => {
$(#[$doc])*
#[derive(Debug, Clone, PartialEq, Eq, Hash, PartialOrd, Ord, Default)]
pub struct $name(Bytes);
impl $name {
#[must_use]
pub fn new(bytes: impl Into<Bytes>) -> Self {
Self(bytes.into())
}
#[must_use]
pub fn as_bytes(&self) -> &[u8] {
&self.0
}
#[must_use]
pub fn into_bytes(self) -> Bytes {
self.0
}
}
};
}
bytes_newtype!(
Key
);
bytes_newtype!(
Value
);
bytes_newtype!(
Cursor
);
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum Precondition {
Absent(Key),
Present(Key),
Equals(Key, Value),
NotAfter(u64),
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum Write {
Put(Key, Value),
Delete(Key),
}
#[derive(Debug, Clone, Default, PartialEq, Eq)]
pub struct Batch {
pub preconditions: Vec<Precondition>,
pub writes: Vec<Write>,
}
impl Batch {
#[must_use]
pub fn new() -> Self {
Self::default()
}
#[must_use]
pub fn require(mut self, precondition: Precondition) -> Self {
self.preconditions.push(precondition);
self
}
#[must_use]
pub fn put(mut self, key: Key, value: Value) -> Self {
self.writes.push(Write::Put(key, value));
self
}
#[must_use]
pub fn delete(mut self, key: Key) -> Self {
self.writes.push(Write::Delete(key));
self
}
#[must_use]
pub fn has_put(&self) -> bool {
self.writes.iter().any(|w| matches!(w, Write::Put(..)))
}
pub fn validate(&self, caps: &StoreCapabilities) -> Result<(), StoreError> {
if self.preconditions.len() + self.writes.len()
> MAX_BATCH_OPS.saturating_sub(caps.reserved_batch_ops)
{
return Err(StoreError::Invalid("batch exceeds MAX_BATCH_OPS".into()));
}
let mut keys = Vec::new();
for pre in &self.preconditions {
match pre {
Precondition::Absent(k) | Precondition::Present(k) => keys.push((k, None)),
Precondition::Equals(k, v) => keys.push((k, Some(v))),
Precondition::NotAfter(_) => {}
}
}
let key_preconditions = keys.len();
for write in &self.writes {
match write {
Write::Put(k, v) => keys.push((k, Some(v))),
Write::Delete(k) => keys.push((k, None)),
}
}
let total: usize = keys
.iter()
.map(|(k, v)| k.as_bytes().len() + v.map_or(0, |v| v.as_bytes().len()))
.sum();
if total > MAX_BATCH_BYTES {
return Err(StoreError::Invalid("batch exceeds MAX_BATCH_BYTES".into()));
}
for (key, value) in &keys {
if key.as_bytes().len() > MAX_KEY_BYTES {
return Err(StoreError::Invalid("key exceeds MAX_KEY_BYTES".into()));
}
if value.is_some_and(|v| v.as_bytes().len() > MAX_VALUE_BYTES) {
return Err(StoreError::Invalid("value exceeds MAX_VALUE_BYTES".into()));
}
if value.is_some_and(|v| v.as_bytes().len() > MAX_TIMER_VALUE_BYTES)
&& key.as_bytes().starts_with(b"w\0")
&& matches!(keys::parse(key), Some(keys::ParsedKey::Timer { .. }))
{
return Err(StoreError::Invalid(
"timer value exceeds retry batch allowance".into(),
));
}
if caps.key_classes == KeyClasses::RefsOnly && !keys::is_ref_key(key) {
return Err(StoreError::Unsupported(
"this store holds only ref keys".into(),
));
}
}
if !caps.atomic_multi_key {
let (pre, writes) = keys.split_at(key_preconditions);
let one_key = match (pre, writes) {
([], [] | [_]) | ([_], []) => true,
([(p, _)], [(w, _)]) => p == w,
_ => false,
};
if !one_key {
return Err(StoreError::Unsupported(
"this store commits at most one key per batch".into(),
));
}
}
Ok(())
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum BatchOutcome {
Committed,
PreconditionFailed {
index: usize,
observed: Option<Value>,
},
DeadlinePassed {
backend_now: u64,
},
}
#[derive(Debug, Clone, PartialEq, Eq, Default)]
pub struct ScanPage {
pub entries: Vec<(Key, Value)>,
pub next: Option<Cursor>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
#[non_exhaustive]
pub struct RangeScan {
pub start: Key,
pub end: Key,
pub after: Option<Cursor>,
pub limit: u32,
}
impl RangeScan {
#[must_use]
pub fn new(start: Key, end: Key, after: Option<Cursor>, limit: u32) -> Self {
Self {
start,
end,
after,
limit,
}
}
}
pub const MAX_SCAN_RANGES: usize = 256;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
#[non_exhaustive]
pub enum KeyClasses {
All,
RefsOnly,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum MembershipMode {
StorePresence,
Explicit,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
#[non_exhaustive]
pub struct StoreCapabilities {
pub atomic_multi_key: bool,
pub reserved_batch_ops: usize,
pub key_classes: KeyClasses,
pub membership: MembershipMode,
pub implicit_layout_version: Option<u32>,
}
impl StoreCapabilities {
#[must_use]
pub const fn full() -> Self {
Self {
atomic_multi_key: true,
reserved_batch_ops: 0,
key_classes: KeyClasses::All,
membership: MembershipMode::StorePresence,
implicit_layout_version: None,
}
}
#[must_use]
pub const fn refs_only() -> Self {
Self {
atomic_multi_key: false,
reserved_batch_ops: 0,
key_classes: KeyClasses::RefsOnly,
membership: MembershipMode::StorePresence,
implicit_layout_version: Some(keys::LAYOUT_VERSION),
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct PartitionStats {
pub bytes: u64,
pub keys: Option<u64>,
}
pub trait NamespaceStore: MaybeSend + MaybeSync {
fn capabilities(&self) -> StoreCapabilities;
fn get(
&self,
p: &Partition,
key: &Key,
) -> impl Future<Output = Result<Option<Value>, StoreError>> + MaybeSend;
fn has(
&self,
p: &Partition,
key: &Key,
) -> impl Future<Output = Result<bool, StoreError>> + MaybeSend {
async move { Ok(self.get(p, key).await?.is_some()) }
}
fn get_many(
&self,
p: &Partition,
keys: &[Key],
) -> impl Future<Output = Result<Vec<Option<Value>>, StoreError>> + MaybeSend {
async move {
let mut values = Vec::with_capacity(keys.len());
for key in keys {
values.push(self.get(p, key).await?);
}
Ok(values)
}
}
fn scan(
&self,
p: &Partition,
start: &Key,
end: &Key,
after: Option<&Cursor>,
limit: u32,
) -> impl Future<Output = Result<ScanPage, StoreError>> + MaybeSend;
fn scan_many(
&self,
p: &Partition,
ranges: &[RangeScan],
) -> impl Future<Output = Result<Vec<ScanPage>, StoreError>> + MaybeSend {
async move {
if ranges.len() > MAX_SCAN_RANGES {
return Err(StoreError::Invalid("too many scan ranges".into()));
}
let mut pages = Vec::with_capacity(ranges.len());
for range in ranges {
pages.push(
self.scan(
p,
&range.start,
&range.end,
range.after.as_ref(),
range.limit,
)
.await?,
);
}
Ok(pages)
}
}
fn apply(
&self,
p: &Partition,
batch: Batch,
) -> impl Future<Output = Result<BatchOutcome, StoreError>> + MaybeSend;
fn stats(
&self,
p: &Partition,
) -> impl Future<Output = Result<PartitionStats, StoreError>> + MaybeSend;
fn probe(&self) -> impl Future<Output = Result<(), StoreError>> + MaybeSend;
}
impl<S: NamespaceStore + ?Sized> NamespaceStore for std::sync::Arc<S> {
fn capabilities(&self) -> StoreCapabilities {
(**self).capabilities()
}
async fn get(&self, p: &Partition, key: &Key) -> Result<Option<Value>, StoreError> {
(**self).get(p, key).await
}
async fn has(&self, p: &Partition, key: &Key) -> Result<bool, StoreError> {
(**self).has(p, key).await
}
async fn get_many(
&self,
p: &Partition,
keys: &[Key],
) -> Result<Vec<Option<Value>>, StoreError> {
(**self).get_many(p, keys).await
}
async fn scan(
&self,
p: &Partition,
start: &Key,
end: &Key,
after: Option<&Cursor>,
limit: u32,
) -> Result<ScanPage, StoreError> {
(**self).scan(p, start, end, after, limit).await
}
async fn scan_many(
&self,
p: &Partition,
ranges: &[RangeScan],
) -> Result<Vec<ScanPage>, StoreError> {
(**self).scan_many(p, ranges).await
}
async fn apply(&self, p: &Partition, batch: Batch) -> Result<BatchOutcome, StoreError> {
(**self).apply(p, batch).await
}
async fn stats(&self, p: &Partition) -> Result<PartitionStats, StoreError> {
(**self).stats(p).await
}
async fn probe(&self) -> Result<(), StoreError> {
(**self).probe().await
}
}