use rusqlite::TransactionBehavior;
use super::*;
#[derive(Default, Debug)]
#[repr(C)]
pub struct Limits {
pub max_value_length_sum: Option<u64>,
pub disable_hole_punching: bool,
}
type DeletedValuesSender = std::sync::mpsc::SyncSender<Vec<NonzeroValueLocation>>;
#[derive(Debug)]
pub struct Handle {
pub(crate) conn: Mutex<Connection>,
pub(crate) exclusive_files: Mutex<HashMap<FileId, ExclusiveFile>>,
pub(crate) dir: Dir,
pub(crate) clones: Mutex<FileCloneCache>,
pub(crate) instance_limits: Limits,
deleted_values: Option<DeletedValuesSender>,
_value_puncher: Option<std::thread::JoinHandle<()>>,
}
type ManifestUserVersion = u32;
impl Handle {
pub fn dir_supports_file_cloning(&self) -> bool {
self.dir.supports_file_cloning()
}
pub fn set_instance_limits(&mut self, limits: Limits) -> Result<()> {
self.instance_limits = limits;
self.start_deferred_transaction()?.apply_limits()
}
pub fn dir(&self) -> &Path {
self.dir.as_ref()
}
pub(crate) fn get_exclusive_file(&self) -> Result<ExclusiveFile> {
{
let mut files = self.exclusive_files.lock().unwrap();
if let Some(id) = files.keys().next().cloned() {
let file = files.remove(&id).unwrap();
debug_assert_eq!(id, file.id);
debug!("using exclusive file {} from handle", &file.id);
return Ok(file);
}
}
trace!("about to open existing files");
if let Some(file) = self.open_existing_exclusive_file()? {
debug!("opened existing values file {}", file.id);
return Ok(file);
}
trace!("here");
let ret = ExclusiveFile::new(&self.dir);
if let Ok(file) = &ret {
debug!("created new exclusive file {}", file.id);
}
ret
}
fn open_existing_exclusive_file(&self) -> Result<Option<ExclusiveFile>> {
for res in read_dir(&self.dir)? {
let entry = res?;
if !entry.file_type()?.is_file() {
continue;
}
if !valid_file_name(entry.file_name().to_str().unwrap()) {
continue;
}
let path = entry.path();
debug!(?path, "opening existing file");
match ExclusiveFile::open(path.clone()) {
Ok(ef) => return Ok(ef),
Err(err) => {
debug!(?path, ?err, "open");
}
}
}
Ok(None)
}
const USER_VERSION: u32 = 2;
pub fn new(dir: PathBuf) -> Result<Self> {
let sqlite_version = rusqlite::version_number();
if sqlite_version < 3042000 {
bail!(
"sqlite version {} below minimum {}",
rusqlite::version(),
"3.42"
);
}
let dir = Dir::new(dir)?;
let mut conn = Connection::open(dir.path().join(MANIFEST_DB_FILE_NAME))?;
Self::init_sqlite_conn(&mut conn)?;
let (deleted_values, receiver) = std::sync::mpsc::sync_channel(10);
let handle = Self {
conn: Mutex::new(conn),
exclusive_files: Default::default(),
dir: dir.clone(),
clones: Default::default(),
instance_limits: Default::default(),
deleted_values: Some(deleted_values),
_value_puncher: Some(std::thread::spawn(|| -> () {
if let Err(err) = Self::value_puncher(dir, receiver) {
error!("value puncher thread failed with {err:?}");
}
})),
};
Ok(handle)
}
fn init_sqlite_conn(conn: &mut Connection) -> rusqlite::Result<()> {
conn.pragma_update(None, "synchronous", "off")?;
let get_user_version = |conn: &Connection| -> Result<ManifestUserVersion, _> {
conn.pragma_query_value(None, "user_version", |row| row.get(0))
};
let user_version: ManifestUserVersion = get_user_version(conn)?;
if user_version == Self::USER_VERSION {
return Ok(());
}
conn.pragma_update(None, "journal_mode", "wal")?;
let tx = conn.transaction_with_behavior(TransactionBehavior::Immediate)?;
if get_user_version(&tx)? < Self::USER_VERSION {
init_manifest_schema(&tx)?;
tx.pragma_update(None, "user_version", Self::USER_VERSION)?;
}
if false {
tx.pragma_update(None, "locking_mode", "exclusive")?;
}
tx.commit()
}
pub fn cleanup_snapshots(&self) -> PubResult<()> {
delete_unused_snapshots(self.dir.path()).map_err(Into::into)
}
pub fn block_size(&self) -> u64 {
self.dir.block_size()
}
pub fn new_writer(&self) -> Result<BatchWriter> {
Ok(BatchWriter {
handle: self,
exclusive_files: Default::default(),
pending_writes: Default::default(),
value_renames: Default::default(),
})
}
fn start_transaction<'h, T, O>(
&'h self,
make_tx: impl FnOnce(&'h mut Connection, &'h Handle) -> rusqlite::Result<T>,
) -> rusqlite::Result<O>
where
O: From<OwnedTxInner<'h, T>>,
{
let guard = self.conn.lock().unwrap();
Ok(owned_cell::OwnedCell::try_make(guard, |conn| make_tx(conn, self))?.into())
}
pub(crate) fn start_immediate_transaction(&self) -> rusqlite::Result<OwnedTx> {
self.start_writable_transaction_with_behaviour(TransactionBehavior::Immediate)
}
pub(crate) fn start_writable_transaction_with_behaviour(
&self,
behaviour: TransactionBehavior,
) -> rusqlite::Result<OwnedTx> {
self.start_transaction(|conn, handle| {
let rtx = conn.transaction_with_behavior(behaviour)?;
Ok(Transaction::new(rtx, handle))
})
}
pub fn start_deferred_transaction_for_read(&self) -> rusqlite::Result<OwnedReadTx> {
self.start_transaction(|conn, _handle| {
let rtx = conn.transaction_with_behavior(TransactionBehavior::Deferred)?;
Ok(ReadTransaction {
tx: ReadOnlyRusqliteTransaction { conn: rtx },
})
})
}
pub(crate) fn start_deferred_transaction(&self) -> rusqlite::Result<OwnedTx> {
self.start_writable_transaction_with_behaviour(TransactionBehavior::Deferred)
}
pub fn read(&self) -> rusqlite::Result<Reader> {
let reader = Reader {
owned_tx: self.start_deferred_transaction()?,
handle: self,
reads: Default::default(),
};
Ok(reader)
}
pub fn read_single(&self, key: &[u8]) -> Result<Option<SnapshotValue<Value>>> {
let mut reader = self.read()?;
let Some(value) = reader.add(key)? else {
return Ok(None);
};
let snapshot = reader.begin()?;
Ok(Some(snapshot.value(value)))
}
pub fn single_write_from(
&self,
key: Vec<u8>,
r: impl Read,
) -> Result<(u64, WriteCommitResult)> {
let mut writer = self.new_writer()?;
let mut value = writer.new_value().begin()?;
trace!("got value writer");
let n = value.copy_from(r)?;
writer.stage_write(key, value)?;
let commit = writer.commit()?;
Ok((n, commit))
}
pub fn single_delete(&self, key: &[u8]) -> PubResult<Option<c_api::PossumStat>> {
let mut tx = self.start_deferred_transaction()?;
let deleted = tx.delete_key(key)?;
if deleted.is_some() {
tx.commit(())?.complete()?;
}
Ok(deleted)
}
pub fn clone_from_file(&mut self, key: Vec<u8>, file: &mut File) -> Result<u64> {
let mut writer = self.new_writer()?;
let mut value = writer.new_value().clone_file(file)?;
let n = value.value_length()?;
writer.stage_write(key, value)?;
writer.commit()?;
Ok(n)
}
pub fn rename_item(&mut self, from: &[u8], to: &[u8]) -> PubResult<Timestamp> {
let mut tx = self.start_immediate_transaction()?;
let last_used = tx.rename_item(from, to)?;
Ok(tx.commit(last_used)?.complete()?)
}
pub fn walk_dir(&self) -> Result<Vec<walk::Entry>> {
crate::walk::walk_dir(&self.dir)
}
pub fn list_items(&self, prefix: &[u8]) -> PubResult<Vec<Item>> {
self.start_deferred_transaction_for_read()?
.list_items(prefix)
}
fn value_puncher(
dir: Dir,
values_receiver: std::sync::mpsc::Receiver<Vec<NonzeroValueLocation>>,
) -> Result<()> {
let manifest_path = dir.path().join(MANIFEST_DB_FILE_NAME);
use rusqlite::OpenFlags;
let mut conn = Connection::open_with_flags(
manifest_path,
OpenFlags::SQLITE_OPEN_READ_ONLY
| OpenFlags::SQLITE_OPEN_NO_MUTEX
| OpenFlags::SQLITE_OPEN_URI,
)?;
while let Ok(mut values) = values_receiver.recv() {
while let Ok(mut more_values) = values_receiver.try_recv() {
values.append(&mut more_values);
}
let tx = conn.transaction_with_behavior(TransactionBehavior::Deferred)?;
let tx = ReadTransaction {
tx: ReadOnlyRusqliteTransaction { conn: tx },
};
Self::punch_values(&dir, &values, &tx)?;
}
Ok(())
}
pub(crate) fn punch_values(
dir: &Dir,
values: &[NonzeroValueLocation],
transaction: &ReadTransactionOwned,
) -> PubResult<()> {
for v in values {
let NonzeroValueLocation {
file_id,
file_offset,
length,
..
} = v;
let value_length = length;
let msg = format!(
"deleting value at {:?} {} {}",
file_id, file_offset, value_length
);
debug!("{}", msg);
punch_value(PunchValueOptions {
dir: dir.path(),
file_id,
offset: *file_offset,
length: *value_length,
tx: transaction,
block_size: dir.block_size(),
constraints: Default::default(),
})
.context(msg)?;
}
Ok(())
}
pub(crate) fn send_values_for_delete(&self, values: Vec<NonzeroValueLocation>) {
use std::sync::mpsc::TrySendError::*;
let sender = self.deleted_values.as_ref().unwrap();
match sender.try_send(values) {
Ok(()) => (),
Err(Disconnected(values)) => {
error!("sending {values:?}: channel disconnected");
}
Err(Full(values)) => {
warn!("channel full while sending values. blocking.");
sender.send(values).unwrap()
}
}
}
}
use item::Item;
use crate::dir::Dir;
use crate::ownedtx::{OwnedReadTx, OwnedTxInner};
use crate::tx::{ReadOnlyRusqliteTransaction, ReadTransaction};
impl Drop for Handle {
fn drop(&mut self) {
}
}