use std::collections::HashMap;
use std::path::{Path, PathBuf};
use std::sync::{Arc, RwLock};
use crate::cache::{BoundedCache, CacheWeight};
use crate::content::node::NodeState;
use crate::content::provider::SegmentProvider;
use crate::content::template::{Template, read_template};
use crate::content::value::read_string;
use crate::error::{Error, Result};
use crate::journal::{JournalEntry, read_journal};
use crate::progress::{DiscardedProgress, ProgressObserver, Step, WorkUnit};
use crate::segment::identifier::SegmentIdentifier;
use crate::segment::parsed_segment::ParsedSegment;
use crate::segment::record::RecordIdentifier;
use crate::segment::view::SegmentView;
use crate::tar_archive::archive::TarArchiveReader;
use crate::tar_archive::file_name::group_file_generations_newest_first;
mod archives;
mod manifest;
pub use archives::*;
pub(crate) use manifest::*;
pub(crate) const SEGMENT_CACHE_BUDGET_BYTES: usize = 192 * 1024 * 1024;
pub(crate) const STRING_CACHE_BUDGET_BYTES: usize = 48 * 1024 * 1024;
pub(crate) const TEMPLATE_CACHE_BUDGET_BYTES: usize = 48 * 1024 * 1024;
pub struct Repository {
pub(crate) directory: PathBuf,
pub(crate) archives: Vec<TarArchiveReader>,
pub(crate) segment_locations: HashMap<SegmentIdentifier, usize>,
pub(crate) journal_entries: Vec<JournalEntry>,
pub(crate) head_record_identifier: RecordIdentifier,
pub(crate) parsed_segment_cache: RwLock<BoundedCache<SegmentIdentifier, Arc<ParsedSegment>>>,
pub(crate) string_cache: RwLock<BoundedCache<RecordIdentifier, Arc<str>>>,
pub(crate) template_cache: RwLock<BoundedCache<RecordIdentifier, Arc<Template>>>,
}
impl Repository {
pub fn open(directory: &Path) -> Result<Self> {
Self::open_with_progress(directory, &mut DiscardedProgress)
}
pub fn open_with_progress(
directory: &Path,
observer: &mut dyn ProgressObserver,
) -> Result<Self> {
if !directory.is_dir() {
return Err(Error::InvalidFormat {
details: format!("{} is not a directory", directory.display()),
});
}
let archive_file_names = list_archive_file_names(directory)?;
check_manifest(directory, ArchivePresence::of(&archive_file_names))?;
let archives = open_archives_newest_valid_first(directory, &archive_file_names, observer)?;
let expected_segments: usize = archives.iter().map(TarArchiveReader::segment_count).sum();
let mut segment_locations = HashMap::with_capacity(expected_segments);
for (archive_position, archive) in archives.iter().enumerate() {
for segment_identifier in archive.segment_identifiers() {
segment_locations
.entry(segment_identifier)
.or_insert(archive_position);
}
}
let journal_entries = read_journal(&directory.join("journal.log"))?;
let head_record_identifier = journal_entries
.iter()
.filter_map(JournalEntry::record_identifier)
.find(|identifier| segment_locations.contains_key(&identifier.segment))
.ok_or_else(|| Error::InvalidFormat {
details: format!(
"no journal revision in {} references an existing segment; \
cannot open a read-only store from an empty journal",
directory.display()
),
})?;
Ok(Self {
directory: directory.to_owned(),
archives,
segment_locations,
journal_entries,
head_record_identifier,
parsed_segment_cache: RwLock::new(BoundedCache::new(SEGMENT_CACHE_BUDGET_BYTES)),
string_cache: RwLock::new(BoundedCache::new(STRING_CACHE_BUDGET_BYTES)),
template_cache: RwLock::new(BoundedCache::new(TEMPLATE_CACHE_BUDGET_BYTES)),
})
}
#[must_use]
pub fn directory(&self) -> &Path {
&self.directory
}
#[must_use]
pub fn archives(&self) -> &[TarArchiveReader] {
&self.archives
}
#[must_use]
pub fn journal_entries(&self) -> &[JournalEntry] {
&self.journal_entries
}
#[must_use]
pub fn head_record_identifier(&self) -> RecordIdentifier {
self.head_record_identifier
}
#[must_use]
pub fn segment_count(&self) -> usize {
self.segment_locations.len()
}
#[must_use]
pub fn contains_segment(&self, segment_identifier: SegmentIdentifier) -> bool {
self.segment_locations.contains_key(&segment_identifier)
}
pub fn segment_identifiers(&self) -> impl Iterator<Item = SegmentIdentifier> + '_ {
self.archives
.iter()
.flat_map(TarArchiveReader::segment_identifiers)
}
pub fn distinct_segment_identifiers(&self) -> impl Iterator<Item = SegmentIdentifier> + '_ {
self.archives
.iter()
.enumerate()
.flat_map(move |(position, archive)| {
archive.segment_identifiers().filter(move |identifier| {
self.segment_locations.get(identifier) == Some(&position)
})
})
}
pub fn segment_bytes(&self, segment_identifier: SegmentIdentifier) -> Result<&[u8]> {
let archive_position = *self
.segment_locations
.get(&segment_identifier)
.ok_or(Error::SegmentNotFound { segment_identifier })?;
self.archives[archive_position]
.segment_data(segment_identifier)
.ok_or(Error::SegmentNotFound { segment_identifier })
}
#[must_use]
pub fn node(&self, record_identifier: RecordIdentifier) -> NodeState<'_> {
NodeState::new(self, record_identifier)
}
#[must_use]
pub fn head(&self) -> NodeState<'_> {
self.node(self.head_record_identifier)
}
pub fn content_root(&self) -> Result<NodeState<'_>> {
self.head()
.child_node("root")?
.ok_or_else(|| Error::InvalidFormat {
details: "the super-root has no \"root\" child node".to_owned(),
})
}
pub fn node_at_path(&self, path: &str) -> Result<Option<NodeState<'_>>> {
let mut current = self.content_root()?;
for name in path.split('/').filter(|name| !name.is_empty()) {
match current.child_node(name)? {
Some(child) => current = child,
None => return Ok(None),
}
}
Ok(Some(current))
}
pub fn checkpoints(&self) -> Result<Vec<(String, NodeState<'_>)>> {
match self.head().child_node("checkpoints")? {
None => Ok(Vec::new()),
Some(checkpoints) => checkpoints.child_node_entries(),
}
}
}
impl SegmentProvider for Repository {
fn segment(&self, segment_identifier: SegmentIdentifier) -> Result<SegmentView<'_>> {
let bytes = self.segment_bytes(segment_identifier)?;
let structure =
load_through_cache(&self.parsed_segment_cache, &segment_identifier, || {
ParsedSegment::parse(segment_identifier, bytes).map(Arc::new)
})?;
Ok(SegmentView {
structure,
bytes: bytes.into(),
})
}
fn string(&self, record_identifier: RecordIdentifier) -> Result<Arc<str>> {
load_through_cache(&self.string_cache, &record_identifier, || {
read_string(self, record_identifier).map(Arc::from)
})
}
fn template(&self, record_identifier: RecordIdentifier) -> Result<Arc<Template>> {
load_through_cache(&self.template_cache, &record_identifier, || {
read_template(self, record_identifier).map(Arc::new)
})
}
}
impl std::fmt::Debug for Repository {
fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
write!(
formatter,
"Repository({}, {} archives, {} segments)",
self.directory.display(),
self.archives.len(),
self.segment_locations.len()
)
}
}
pub(crate) fn load_through_cache<Key, Value>(
cache: &RwLock<BoundedCache<Key, Value>>,
key: &Key,
load: impl FnOnce() -> Result<Value>,
) -> Result<Value>
where
Key: Clone + Eq + std::hash::Hash,
Value: Clone + CacheWeight,
{
if let Some(cached) = cache
.read()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.get(key)
{
return Ok(cached);
}
let value = load()?;
cache
.write()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.insert(key.clone(), value.clone());
Ok(value)
}