use std::sync::Arc;
use pigeonhole_engine::{FamilyId, Predicate, TableId, Txn, ValueRef};
use pigeonhole_format::Durability;
use pigeonhole_format::key::MAX_KEY_PART;
use crate::db::Db;
use crate::table::TableCore;
use crate::{CellRef, Condition, Error, ErrorCode, Result, Table};
type EngineBatch = pigeonhole_engine::WriteBatch;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct CommitInfo {
pub seqno: u64,
pub durability: Durability,
}
impl From<pigeonhole_engine::CommitInfo> for CommitInfo {
fn from(c: pigeonhole_engine::CommitInfo) -> Self {
Self {
seqno: c.seqno,
durability: c.durability,
}
}
}
#[derive(Debug, Clone, Copy)]
struct Sizes {
row: usize,
qualifier: usize,
value: usize,
}
impl Sizes {
fn new(row: &[u8], qualifier: &[u8], value: usize) -> Self {
Self {
row: row.len(),
qualifier: qualifier.len(),
value,
}
}
}
fn size_error(e: pigeonhole_engine::Error, sizes: Sizes, max_value: usize) -> Error {
use pigeonhole_engine::Error as E;
match e {
E::KeyTooLarge => {
let (what, len) = if sizes.row > MAX_KEY_PART {
("row key", sizes.row)
} else {
("qualifier", sizes.qualifier)
};
Error::new(
ErrorCode::KeyTooLarge,
format!("{what} of {len} bytes exceeds the limit of {MAX_KEY_PART} bytes"),
)
}
E::ValueTooLarge => value_too_large(sizes.value, max_value),
e => e.into(),
}
}
fn value_too_large(len: usize, max_value: usize) -> Error {
Error::new(
ErrorCode::ValueTooLarge,
format!(
"value of {len} bytes exceeds the limit: {MAX_PUT_VALUE} bytes for a put, \
{max_value} bytes for a merge operand (the smallest of the WAL segment payload, \
64 MiB and half a shard's memtable arena; decision D16)"
),
)
}
const MAX_PUT_VALUE: usize = u32::MAX as usize - 1;
#[derive(Debug, Default)]
struct Builder {
batch: EngineBatch,
error: Option<Error>,
largest_value: usize,
}
impl Builder {
fn with(
&mut self,
table: &TableCore,
family: &str,
sizes: Sizes,
f: impl FnOnce(&mut EngineBatch, FamilyId) -> pigeonhole_engine::Result<()>,
) {
if self.error.is_some() {
return;
}
self.largest_value = self.largest_value.max(sizes.value);
let max_value = table.db.max_value;
let r = table
.family_id(family)
.and_then(|id| f(&mut self.batch, id).map_err(|e| size_error(e, sizes, max_value)));
if let Err(e) = r {
self.error = Some(e);
}
}
fn row_delete(&mut self, table: &TableCore, row: &[u8]) {
if self.error.is_none()
&& let Err(e) = self.batch.delete_row(table.info.id, row, None)
{
self.error = Some(size_error(e, Sizes::new(row, &[], 0), 0));
}
}
fn check_table(&mut self, db: &Arc<Db>, table: &Table) -> bool {
if self.error.is_none()
&& let Err(e) = check_table(db, table)
{
self.error = Some(e);
}
self.error.is_none()
}
fn commit(self, db: &Db, durability: Option<Durability>) -> Result<CommitInfo> {
if let Some(e) = self.error {
return Err(e);
}
let largest = self.largest_value;
db.engine
.commit(self.batch, durability)
.map(CommitInfo::from)
.map_err(|e| commit_error(e, largest, db.max_value))
}
}
fn commit_error(e: pigeonhole_engine::Error, largest_value: usize, max_value: usize) -> Error {
match e {
pigeonhole_engine::Error::ValueTooLarge => value_too_large(largest_value, max_value),
e => e.into(),
}
}
fn check_table(db: &Arc<Db>, table: &Table) -> Result<()> {
if Arc::ptr_eq(db, &table.core.db) {
Ok(())
} else {
Err(Error::new(
ErrorCode::InvalidArgument,
format!("table {:?} belongs to another database", table.name()),
))
}
}
#[derive(Debug)]
#[must_use = "a mutation does nothing until .commit()"]
pub struct RowMutation<'t> {
table: &'t TableCore,
row: Vec<u8>,
builder: Builder,
durability: Option<Durability>,
}
impl<'t> RowMutation<'t> {
pub(crate) fn new(table: &'t Table, row: &[u8]) -> Self {
Self {
table: &table.core,
row: row.to_vec(),
builder: Builder::default(),
durability: None,
}
}
fn op(
mut self,
family: &str,
qualifier: &[u8],
value: usize,
f: impl FnOnce(&mut EngineBatch, TableId, FamilyId, &[u8]) -> pigeonhole_engine::Result<()>,
) -> Self {
let (table, row) = (self.table.info.id, &self.row);
let sizes = Sizes::new(row, qualifier, value);
self.builder
.with(self.table, family, sizes, |b, fam| f(b, table, fam, row));
self
}
}
impl RowMutation<'_> {
pub fn put(self, family: &str, qualifier: &[u8], value: &[u8]) -> Self {
self.op(family, qualifier, value.len(), |b, t, f, row| {
b.put(t, f, row, qualifier, None, ValueRef::Bytes(value))
})
}
pub fn put_at(self, family: &str, qualifier: &[u8], ts: u64, value: &[u8]) -> Self {
self.op(family, qualifier, value.len(), |b, t, f, row| {
b.put(t, f, row, qualifier, Some(ts), ValueRef::Bytes(value))
})
}
pub fn put_i64(self, family: &str, qualifier: &[u8], value: i64) -> Self {
self.op(family, qualifier, 8, |b, t, f, row| {
b.put(t, f, row, qualifier, None, ValueRef::I64(value))
})
}
pub fn put_i64_at(self, family: &str, qualifier: &[u8], ts: u64, value: i64) -> Self {
self.op(family, qualifier, 8, |b, t, f, row| {
b.put(t, f, row, qualifier, Some(ts), ValueRef::I64(value))
})
}
pub fn put_f64(self, family: &str, qualifier: &[u8], value: f64) -> Self {
self.op(family, qualifier, 8, |b, t, f, row| {
b.put(t, f, row, qualifier, None, ValueRef::F64(value))
})
}
pub fn incr(self, family: &str, qualifier: &[u8], delta: i64) -> Self {
self.op(family, qualifier, 8, |b, t, f, row| {
b.merge(t, f, row, qualifier, ValueRef::I64(delta))
})
}
pub fn incr_at(self, family: &str, qualifier: &[u8], ts: u64, delta: i64) -> Self {
self.op(family, qualifier, 8, |b, t, f, row| {
b.merge_at(t, f, row, qualifier, ts, ValueRef::I64(delta))
})
}
pub fn merge(self, family: &str, qualifier: &[u8], operand: &[u8]) -> Self {
self.op(family, qualifier, operand.len(), |b, t, f, row| {
b.merge(t, f, row, qualifier, ValueRef::Bytes(operand))
})
}
pub fn delete_cell(self, family: &str, qualifier: &[u8], ts: u64) -> Self {
self.op(family, qualifier, 0, |b, t, f, row| {
b.delete_cell(t, f, row, qualifier, ts)
})
}
pub fn delete_column(self, family: &str, qualifier: &[u8]) -> Self {
self.op(family, qualifier, 0, |b, t, f, row| {
b.delete_column(t, f, row, qualifier, None)
})
}
pub fn delete_family(self, family: &str) -> Self {
self.op(family, &[], 0, |b, t, f, row| {
b.delete_family(t, f, row, None)
})
}
pub fn delete_row(mut self) -> Self {
self.builder.row_delete(self.table, &self.row);
self
}
pub fn durability(mut self, durability: Durability) -> Self {
self.durability = Some(durability);
self
}
pub fn commit(self) -> Result<CommitInfo> {
self.builder.commit(&self.table.db, self.durability)
}
pub fn commit_if(self, condition: &Condition) -> Result<Option<CommitInfo>> {
if let Some(e) = self.builder.error {
return Err(e);
}
let predicate = match condition {
Condition::Exists { family, qualifier } => Predicate::Exists {
family: self.table.family_id(family)?,
qualifier: qualifier.clone(),
},
Condition::Absent { family, qualifier } => Predicate::Absent {
family: self.table.family_id(family)?,
qualifier: qualifier.clone(),
},
Condition::Value {
family,
qualifier,
filter,
} => Predicate::Value {
family: self.table.family_id(family)?,
qualifier: qualifier.clone(),
predicate: filter.to_engine(),
},
};
let db = &self.table.db;
let largest = self.builder.largest_value;
let (applied, info) = db
.engine
.check_and_mutate(
self.table.info.id,
&self.row,
&predicate,
self.builder.batch,
self.durability,
)
.map_err(|e| commit_error(e, largest, db.max_value))?;
Ok(if applied {
info.map(CommitInfo::from)
} else {
None
})
}
}
#[derive(Debug)]
pub struct WriteBatch {
db: Arc<Db>,
builder: Builder,
}
impl WriteBatch {
pub(crate) fn new(db: Arc<Db>) -> Self {
Self {
db,
builder: Builder::default(),
}
}
fn op(
&mut self,
table: &Table,
row: &[u8],
family: &str,
qualifier: &[u8],
value: usize,
f: impl FnOnce(&mut EngineBatch, TableId, FamilyId) -> pigeonhole_engine::Result<()>,
) -> &mut Self {
if self.builder.check_table(&self.db, table) {
let id = table.core.info.id;
let sizes = Sizes::new(row, qualifier, value);
self.builder
.with(&table.core, family, sizes, |b, fam| f(b, id, fam));
}
self
}
pub fn put(
&mut self,
table: &Table,
row: &[u8],
family: &str,
qualifier: &[u8],
value: &[u8],
) -> &mut Self {
self.op(table, row, family, qualifier, value.len(), |b, t, f| {
b.put(t, f, row, qualifier, None, ValueRef::Bytes(value))
})
}
pub fn put_at(
&mut self,
table: &Table,
row: &[u8],
family: &str,
qualifier: &[u8],
ts: u64,
value: &[u8],
) -> &mut Self {
self.op(table, row, family, qualifier, value.len(), |b, t, f| {
b.put(t, f, row, qualifier, Some(ts), ValueRef::Bytes(value))
})
}
pub fn put_i64(
&mut self,
table: &Table,
row: &[u8],
family: &str,
qualifier: &[u8],
value: i64,
) -> &mut Self {
self.op(table, row, family, qualifier, 8, |b, t, f| {
b.put(t, f, row, qualifier, None, ValueRef::I64(value))
})
}
pub fn put_i64_at(
&mut self,
table: &Table,
row: &[u8],
family: &str,
qualifier: &[u8],
ts: u64,
value: i64,
) -> &mut Self {
self.op(table, row, family, qualifier, 8, |b, t, f| {
b.put(t, f, row, qualifier, Some(ts), ValueRef::I64(value))
})
}
pub fn put_f64(
&mut self,
table: &Table,
row: &[u8],
family: &str,
qualifier: &[u8],
value: f64,
) -> &mut Self {
self.op(table, row, family, qualifier, 8, |b, t, f| {
b.put(t, f, row, qualifier, None, ValueRef::F64(value))
})
}
pub fn incr(
&mut self,
table: &Table,
row: &[u8],
family: &str,
qualifier: &[u8],
delta: i64,
) -> &mut Self {
self.op(table, row, family, qualifier, 8, |b, t, f| {
b.merge(t, f, row, qualifier, ValueRef::I64(delta))
})
}
pub fn incr_at(
&mut self,
table: &Table,
row: &[u8],
family: &str,
qualifier: &[u8],
ts: u64,
delta: i64,
) -> &mut Self {
self.op(table, row, family, qualifier, 8, |b, t, f| {
b.merge_at(t, f, row, qualifier, ts, ValueRef::I64(delta))
})
}
pub fn merge(
&mut self,
table: &Table,
row: &[u8],
family: &str,
qualifier: &[u8],
operand: &[u8],
) -> &mut Self {
self.op(table, row, family, qualifier, operand.len(), |b, t, f| {
b.merge(t, f, row, qualifier, ValueRef::Bytes(operand))
})
}
pub fn delete_cell(
&mut self,
table: &Table,
row: &[u8],
family: &str,
qualifier: &[u8],
ts: u64,
) -> &mut Self {
self.op(table, row, family, qualifier, 0, |b, t, f| {
b.delete_cell(t, f, row, qualifier, ts)
})
}
pub fn delete_column(
&mut self,
table: &Table,
row: &[u8],
family: &str,
qualifier: &[u8],
) -> &mut Self {
self.op(table, row, family, qualifier, 0, |b, t, f| {
b.delete_column(t, f, row, qualifier, None)
})
}
pub fn delete_family(&mut self, table: &Table, row: &[u8], family: &str) -> &mut Self {
self.op(table, row, family, &[], 0, |b, t, f| {
b.delete_family(t, f, row, None)
})
}
pub fn delete_row(&mut self, table: &Table, row: &[u8]) -> &mut Self {
if self.builder.check_table(&self.db, table) {
self.builder.row_delete(&table.core, row);
}
self
}
pub fn len(&self) -> usize {
self.builder.batch.len()
}
pub fn is_empty(&self) -> bool {
self.builder.batch.is_empty()
}
pub fn commit(self) -> Result<CommitInfo> {
self.builder.commit(&self.db, None)
}
pub fn commit_with(self, durability: Durability) -> Result<CommitInfo> {
self.builder.commit(&self.db, Some(durability))
}
}
#[derive(Debug)]
pub struct Transaction {
db: Arc<Db>,
txn: Txn,
error: Option<Error>,
largest_value: usize,
}
impl Transaction {
pub(crate) fn new(db: Arc<Db>, txn: Txn) -> Self {
Self {
db,
txn,
error: None,
largest_value: 0,
}
}
fn op(
&mut self,
table: &Table,
row: &[u8],
family: &str,
qualifier: &[u8],
value: usize,
f: impl FnOnce(&mut EngineBatch, TableId, FamilyId) -> pigeonhole_engine::Result<()>,
) -> &mut Self {
if self.error.is_some() {
return self;
}
self.largest_value = self.largest_value.max(value);
let sizes = Sizes::new(row, qualifier, value);
let max_value = self.db.max_value;
let r = check_table(&self.db, table).and_then(|()| {
let fam = table.core.family_id(family)?;
f(self.txn.batch(), table.core.info.id, fam)
.map_err(|e| size_error(e, sizes, max_value))
});
if let Err(e) = r {
self.error = Some(e);
}
self
}
fn finish(self, durability: Option<Durability>) -> Result<CommitInfo> {
if let Some(e) = self.error {
return Err(e);
}
self.db.check_open()?;
let (largest, max) = (self.largest_value, self.db.max_value);
self.txn
.commit(durability)
.map(CommitInfo::from)
.map_err(|e| commit_error(e, largest, max))
}
pub fn get(
&mut self,
table: &Table,
row: &[u8],
family: &str,
qualifier: &[u8],
) -> Result<Option<CellRef<'_>>> {
check_table(&self.db, table)?;
self.db.check_open()?;
let family = table.core.family_id(family)?;
Ok(self
.txn
.get(table.core.info.id, family, row, qualifier)?
.map(CellRef::owned))
}
pub fn put(
&mut self,
table: &Table,
row: &[u8],
family: &str,
qualifier: &[u8],
value: &[u8],
) -> &mut Self {
self.op(table, row, family, qualifier, value.len(), |b, t, f| {
b.put(t, f, row, qualifier, None, ValueRef::Bytes(value))
})
}
pub fn delete_column(
&mut self,
table: &Table,
row: &[u8],
family: &str,
qualifier: &[u8],
) -> &mut Self {
self.op(table, row, family, qualifier, 0, |b, t, f| {
b.delete_column(t, f, row, qualifier, None)
})
}
pub fn commit(self) -> Result<CommitInfo> {
self.finish(None)
}
pub fn commit_with(self, durability: Durability) -> Result<CommitInfo> {
self.finish(Some(durability))
}
}