use std::{io, io::Write, path::Path, sync::Arc};
use block_on_place::HandleExt;
use derive_more::Debug;
use eyre::Result;
use sqlx::PgPool;
use tantivy::{
Directory,
directory::{
DirectoryLock, FileHandle, Lock, TerminatingWrite, WatchCallback, WatchHandle, WritePtr,
error::{DeleteError, LockError, OpenReadError, OpenWriteError},
},
};
use tokio::runtime::Handle;
use uuid::Uuid;
use crate::{
bundle::{self, Bundler},
context::Context,
directory::is_metadata,
empty::Empty,
file::{File, Slice},
metadata::{FileRecord, MetadataStore},
utils::{PathExt, WrapIoErrorExt},
writer::{OnDone, OpenSink, Outcome, Sink, Writer},
};
#[derive(Clone, Debug)]
#[debug("LightDirectory {{ index: {}, inner: {inner:?} }}", metadata.context.index)]
pub struct LightDirectory<D> {
rt: Handle,
metadata: MetadataStore,
bundle: Bundler,
inner: D,
}
impl<D> LightDirectory<D> {
pub async fn open(
inner: D,
index: Uuid,
operator: opendal::Operator,
pool: PgPool,
) -> Result<Self> {
let context = Context::new(index);
let metadata = MetadataStore::open(&context, pool, operator).await?;
Ok(Self {
rt: Handle::current(),
metadata,
bundle: Bundler::default(),
inner,
})
}
pub fn with_threshold(mut self, threshold: usize) -> Self {
self.metadata.context.threshold = threshold;
self
}
pub fn with_read_chunks(mut self, chunks: usize) -> Self {
self.metadata.context.read_chunks = Some(chunks);
self
}
pub fn with_write_chunks(mut self, chunks: usize) -> Self {
self.metadata.context.write_chunks = Some(chunks);
self
}
pub fn with_read_concurrency(mut self, concurrency: usize) -> Self {
self.metadata.context.read_concurrency = Some(concurrency);
self
}
pub fn with_write_concurrency(mut self, concurrency: usize) -> Self {
self.metadata.context.write_concurrency = Some(concurrency);
self
}
pub fn with_bundling(mut self) -> Self {
self.metadata.context.bundle = true;
self
}
pub fn with_bundle_max_file_bytes(mut self, bytes: usize) -> Self {
self.metadata.context.bundle_max_file_bytes = bytes;
self
}
}
impl<D: Directory + Clone> Directory for LightDirectory<D> {
fn get_file_handle(&self, filepath: &Path) -> Result<Arc<dyn FileHandle>, OpenReadError> {
match self.inner.get_file_handle(filepath) {
Err(OpenReadError::FileDoesNotExist(_)) => (),
result => return result,
}
let path = filepath.try_to_str::<OpenReadError>()?;
if let Some(bytes) = self.rt.block_on_place(self.bundle.get(filepath)) {
return Ok(File::memory_owned(path, bytes));
}
let record = self
.rt
.block_on_place(self.metadata.file_lookup(path))
.map_err(OpenReadError::wrapper(filepath))?;
match record {
Some(record) if record.is_empty => {
let empty = Empty::for_path(filepath)
.ok_or_else(|| OpenReadError::FileDoesNotExist(filepath.into()))?;
Ok(File::memory(path, empty.bytes()))
}
Some(FileRecord {
byte_offset,
byte_length: Some(byte_length),
..
}) => {
let object = bundle::object(filepath)
.ok_or_else(|| OpenReadError::FileDoesNotExist(filepath.into()))?;
let handle = self.inner.get_file_handle(&object)?;
Ok(Slice::new(
handle,
byte_offset as usize,
byte_length as usize,
))
}
_ => Err(OpenReadError::FileDoesNotExist(filepath.into())),
}
}
fn delete(&self, filepath: &Path) -> Result<(), DeleteError> {
match self.inner.delete(filepath) {
Err(DeleteError::FileDoesNotExist(_)) => (),
result => return result,
}
let path = filepath.try_to_str::<DeleteError>()?;
let deleted = self
.rt
.block_on_place(self.metadata.delete_file(path))
.map_err(DeleteError::wrapper(filepath))?;
if deleted {
Ok(())
} else {
Err(DeleteError::FileDoesNotExist(filepath.into()))
}
}
fn exists(&self, filepath: &Path) -> Result<bool, OpenReadError> {
if is_metadata(filepath) {
let path = filepath.try_to_str::<OpenReadError>()?;
return self
.rt
.block_on_place(self.metadata.metadata_exists(path))
.map_err(OpenReadError::wrapper(filepath));
}
if self.inner.exists(filepath)? {
return Ok(true);
}
let path = filepath.try_to_str::<OpenReadError>()?;
let found = self
.rt
.block_on_place(self.metadata.file_lookup(path))
.map_err(OpenReadError::wrapper(filepath))?;
Ok(found.is_some())
}
fn open_write(&self, filepath: &Path) -> Result<WritePtr, OpenWriteError> {
let inner = self.inner.clone();
let target = filepath.to_path_buf();
let open: OpenSink = Box::new(move || {
let writer = inner.open_write(&target).map_err(io::Error::other)?;
Ok(Box::new(writer) as Sink)
});
let eligible = self.metadata.context.bundle && bundle::is_bundleable(filepath);
let cap = if eligible {
self.metadata.context.bundle_max_file_bytes
} else {
Empty::max_len()
};
let metadata = self.metadata.clone();
let bundler = self.bundle.clone();
let rt = self.rt.clone();
let path = filepath.to_path_buf();
let on_done: OnDone = Box::new(move |outcome| {
match outcome {
Outcome::Empty(_) => {
let path = path.try_to_str::<io::Error>()?;
rt.block_on_place(metadata.create_file(path, true, None))
.map_err(io::Error::other)?;
}
Outcome::Standalone => {}
Outcome::Bundled(bytes) => rt.block_on_place(bundler.buffer(path, bytes)),
}
Ok(())
});
let writer = Writer::new(filepath.to_path_buf(), cap, eligible, open, on_done);
Ok(WritePtr::new(Box::new(writer)))
}
fn atomic_read(&self, filepath: &Path) -> Result<Vec<u8>, OpenReadError> {
self.rt
.block_on_place(self.metadata.read_metadata(filepath))
.map_err(OpenReadError::wrapper(filepath))?
.ok_or_else(|| OpenReadError::FileDoesNotExist(filepath.into()))
}
fn atomic_write(&self, filepath: &Path, data: &[u8]) -> io::Result<()> {
self.rt
.block_on_place(self.metadata.write_metadata(filepath, data))
}
fn sync_directory(&self) -> io::Result<()> {
for bundle in self.rt.block_on_place(self.bundle.drain()) {
let mut writer = self
.inner
.open_write(&bundle.path)
.map_err(io::Error::other)?;
writer.write_all(&bundle.bytes)?;
writer.terminate()?;
for entry in bundle.entries {
let path = entry.path.try_to_str::<io::Error>()?;
self.rt
.block_on_place(self.metadata.create_file(
path,
false,
Some((entry.offset, entry.length)),
))
.map_err(io::Error::wrapper(&entry.path))?;
}
}
self.inner.sync_directory()
}
#[inline]
fn watch(&self, callback: WatchCallback) -> tantivy::Result<WatchHandle> {
self.inner.watch(callback)
}
#[inline]
fn acquire_lock(&self, lock: &Lock) -> Result<DirectoryLock, LockError> {
self.inner.acquire_lock(lock)
}
}