use super::*;
#[must_use]
pub(crate) struct PostCommitWork<'h, T> {
handle: &'h Handle,
deleted_values: Vec<NonzeroValueLocation>,
altered_files: HashSet<FileId>,
reward: T,
}
#[repr(transparent)]
pub(crate) struct ReadOnlyRusqliteTransaction<T> {
pub(crate) conn: T,
}
impl<'t, T> ReadOnlyRusqliteTransaction<T>
where
T: Borrow<rusqlite::Transaction<'t>>,
{
pub fn prepare_cached<'a>(&'a self, sql: &str) -> rusqlite::Result<CachedStatement<'a>>
where
't: 'a,
{
let stmt = self.conn.borrow().prepare_cached(sql)?;
assert!(stmt.readonly());
Ok(stmt)
}
}
impl ReadTransactionOwned<'_> {
pub fn as_ref(&self) -> ReadTransactionRef {
ReadTransactionRef {
tx: ReadOnlyRusqliteTransaction {
conn: &self.tx.conn,
},
}
}
}
pub type ReadTransactionRef<'a> = ReadTransaction<&'a rusqlite::Transaction<'a>>;
pub type ReadTransactionOwned<'a> = ReadTransaction<rusqlite::Transaction<'a>>;
#[repr(transparent)]
pub struct ReadTransaction<T> {
pub(crate) tx: ReadOnlyRusqliteTransaction<T>,
}
impl<'a, T> ReadTransaction<T>
where
T: Borrow<rusqlite::Transaction<'a>>,
{
pub fn file_values(
&'a self,
file_id: &'a FileIdFancy,
) -> rusqlite::Result<FileValues<'a, CachedStatement<'a>>> {
let stmt = self.tx.prepare_cached(&format!(
"select {} from keys where file_id=? order by file_offset",
value_columns_sql()
))?;
let iter = FileValues { stmt, file_id };
Ok(iter)
}
pub fn sum_value_length(&self) -> rusqlite::Result<u64> {
self.tx
.prepare_cached("select value from sums where key='value_length'")?
.query_row([], |row| row.get(0))
.map_err(Into::into)
}
pub fn query_last_end_offset(&self, file_id: &FileId, offset: u64) -> rusqlite::Result<u64> {
self.tx
.prepare_cached(
"select max(file_offset+value_length) as last_offset \
from keys \
where file_id=? and file_offset+value_length <= ?",
)?
.query_row(params![file_id.deref(), offset], |row| {
let res: rusqlite::Result<Option<_>> = row.get(0);
res.map(|v| v.unwrap_or_default())
})
}
pub fn next_value_offset(
&self,
file_id: &FileId,
min_offset: u64,
) -> rusqlite::Result<Option<u64>> {
self.tx
.prepare_cached(
"select min(file_offset) \
from keys \
where file_id=? and file_offset >= ?",
)?
.query_row(params![file_id.deref(), min_offset], |row| row.get(0))
}
pub fn list_items(&self, prefix: &[u8]) -> PubResult<Vec<Item>> {
let range_end = {
let mut prefix = prefix.to_owned();
if inc_big_endian_array(&mut prefix) {
Some(prefix)
} else {
None
}
};
match range_end {
None => self.list_items_inner(
&format!(
"select {}, key from keys where key >= ?",
value_columns_sql()
),
[prefix],
),
Some(range_end) => self.list_items_inner(
&format!(
"select {}, key from keys where key >= ? and key < ?",
value_columns_sql()
),
rusqlite::params![prefix, range_end],
),
}
}
fn list_items_inner(&self, sql: &str, params: impl rusqlite::Params) -> PubResult<Vec<Item>> {
self.tx
.prepare_cached(sql)
.unwrap()
.query_map(params, |row| {
Ok(Item {
value: Value::from_row(row)?,
key: row.get(VALUE_COLUMN_NAMES.len())?,
})
})?
.collect::<rusqlite::Result<Vec<_>>>()
.map_err(Into::into)
}
}
impl<'h, T> PostCommitWork<'h, T> {
pub fn complete(self) -> Result<T> {
if !self.handle.instance_limits.disable_hole_punching {
self.handle.send_values_for_delete(self.deleted_values);
}
for file_id in self.altered_files {
self.handle.clones.lock().unwrap().remove(&file_id);
}
Ok(self.reward)
}
}
pub(crate) struct Transaction<'h> {
tx: rusqlite::Transaction<'h>,
handle: &'h Handle,
deleted_values: Vec<NonzeroValueLocation>,
altered_files: HashSet<FileId>,
}
impl<'h> Transaction<'h> {
pub fn read(&self) -> ReadTransaction<&rusqlite::Transaction> {
ReadTransaction {
tx: ReadOnlyRusqliteTransaction { conn: &self.tx },
}
}
pub fn new(tx: rusqlite::Transaction<'h>, handle: &'h Handle) -> Self {
Self {
tx,
handle,
deleted_values: vec![],
altered_files: Default::default(),
}
}
pub(crate) fn commit<T>(mut self, reward: T) -> Result<PostCommitWork<'h, T>> {
self.apply_limits()?;
self.tx.commit()?;
Ok(PostCommitWork {
handle: self.handle,
deleted_values: self.deleted_values,
altered_files: self.altered_files,
reward,
})
}
pub fn touch_for_read(&mut self, key: &[u8]) -> rusqlite::Result<Value> {
self.tx
.prepare_cached(&format!(
"update keys \
set last_used=cast(unixepoch('subsec')*1e3 as integer) \
where key=? \
returning {}",
value_columns_sql()
))?
.query_row([key], Value::from_row)
}
pub fn rename_value(&mut self, value: &Value, new_key: Vec<u8>) -> PubResult<bool> {
match self
.tx
.prepare_cached(&format!(
"delete from keys where key=? returning {}",
value_columns_sql()
))?
.query_row(params![&new_key], Value::from_row)
{
Err(QueryReturnedNoRows) => {}
Err(err) => return Err(err.into()),
Ok(existing_value) => {
match existing_value.location {
Nonzero(a) => {
let b = value;
if Some(a.file_offset) == b.file_offset()
&& Some(&*a.file_id) == b.file_id()
{
assert_eq!(a.length, b.length());
return Ok(true);
}
self.deleted_values.push(a);
}
ZeroLength => {}
}
}
};
let res: rusqlite::Result<ValueLength> = self
.tx
.prepare_cached(
"update keys set key=? where file_id=? and file_offset=?\
returning value_length",
)?
.query_row(
params![new_key, value.file_id(), value.file_offset()],
|row| row.get(0),
);
match res {
Err(QueryReturnedNoRows) => Ok(false),
Err(err) => Err(err).context("updating value key").map_err(Into::into),
Ok(value_length) => {
assert_eq!(value_length, value.length());
Ok(true)
}
}
}
pub fn rename_item(&mut self, from: &[u8], to: &[u8]) -> PubResult<Timestamp> {
let row_result = self.tx.query_row(
"update keys set key=? where key=? returning last_used",
[to, from],
|row| {
let ts: Timestamp = row.get(0)?;
Ok(ts)
},
);
let last_used = match row_result {
Err(QueryReturnedNoRows) => Err(Error::NoSuchKey),
Ok(ok) => Ok(ok),
Err(err) => Err(err.into()),
}?;
assert_eq!(self.tx.changes(), 1);
Ok(last_used)
}
pub(crate) fn insert_key(&mut self, pw: PendingWrite) -> rusqlite::Result<()> {
let mut file_id = Some(pw.value_file_id.deref());
let mut file_offset = Some(pw.value_file_offset);
if pw.value_length == 0 {
file_id = None;
file_offset = None;
}
let inserted = self
.tx
.prepare_cached(
"insert into keys (key, file_id, file_offset, value_length)\
values (?, ?, ?, ?)",
)?
.execute(rusqlite::params!(
pw.key,
file_id,
file_offset,
pw.value_length
))?;
assert_eq!(inserted, 1);
if pw.value_length != 0 {
self.altered_files.insert(pw.value_file_id);
}
Ok(())
}
fn push_value_for_deletion(&mut self, value: Value) {
match value.location {
Nonzero(location) => self.deleted_values.push(location),
ZeroLength => {}
}
}
pub fn delete_key(&mut self, key: &[u8]) -> rusqlite::Result<Option<c_api::PossumStat>> {
let res = self
.tx
.prepare_cached(&format!(
"delete from keys where key=? returning {}",
value_columns_sql()
))?
.query_row([key], Value::from_row);
match res {
Err(QueryReturnedNoRows) => Ok(None),
Ok(value) => {
let stat = value.as_ref().into();
self.push_value_for_deletion(value);
Ok(Some(stat))
}
Err(err) => Err(err),
}
}
pub fn apply_limits(&mut self) -> Result<()> {
if self.tx.transaction_state(None)? != rusqlite::TransactionState::Write {
return Ok(());
}
if let Some(max) = self.handle.instance_limits.max_value_length_sum {
loop {
let actual = self
.read()
.sum_value_length()
.context("reading value_length sum")?;
if actual <= max {
break;
}
self.evict_values(actual - max)?;
}
}
Ok(())
}
pub fn evict_values(&mut self, target_bytes: u64) -> Result<()> {
let mut stmt = self.tx.prepare_cached(&format!(
"delete from keys where key_id in (\
select key_id from keys order by last_used limit 1\
)\
returning {}",
value_columns_sql()
))?;
let mut value_bytes_deleted = 0;
let mut values_deleted = vec![];
while value_bytes_deleted < target_bytes {
let value = stmt.query_row([], Value::from_row)?;
value_bytes_deleted += value.length();
info!("evicting {:?}", &value);
values_deleted.push(value);
}
drop(stmt);
for value in values_deleted {
self.push_value_for_deletion(value);
}
Ok(())
}
}