use std::path::{Path, PathBuf};
use std::sync::Arc;
use arrow::array::{RecordBatch, StringArray, UInt32Array, UInt64Array};
use arrow::datatypes::{DataType, Field, Schema};
use parquet::arrow::ArrowWriter;
use parquet::arrow::arrow_reader::ParquetRecordBatchReaderBuilder;
use parquet::basic::{Compression, ZstdLevel};
use parquet::file::metadata::SortingColumn;
use parquet::file::properties::WriterProperties;
use polyc_eventlog_host::encode_partition;
use polyc_state::revision::PartitionIncarnation;
use super::{IndexedMessage, PostingsRecord};
const META_INDEXED_THROUGH: &str = "polychrome.search.indexed_through";
const META_INCARNATION: &str = "polychrome.search.incarnation";
const META_AVAILABLE: &str = "polychrome.search.available";
const META_EXCISION_SCANNED_THROUGH: &str = "polychrome.search.excision_scanned_through";
const META_DESTROYED: &str = "polychrome.search.destroyed";
const META_KEY_ID: &str = "polychrome.search.key_id";
const META_FORMAT: &str = "polychrome.search.format";
const FORMAT_VERSION: &str = "2";
const TERM_HASH_COLUMN: &str = "term_hash";
const NO_TERMS_SENTINEL: u32 = u32::MAX;
#[derive(Debug, thiserror::Error)]
pub(crate) enum StoreError {
#[error("search projection io: {0}")]
Io(String),
#[error("search projection parquet: {0}")]
Parquet(String),
#[error("search projection metadata is unreadable and must be rebuilt from the journal")]
Corrupt,
#[error("search projection exceeds the {MAX_POSTINGS_ROWS}-row read cap")]
TooLarge,
#[error("search projection for a destroyed conversation cannot be published to")]
Destroyed,
}
impl StoreError {
pub(crate) const fn is_unreadable(&self) -> bool {
matches!(self, Self::Corrupt | Self::Parquet(_))
}
}
impl From<std::io::Error> for StoreError {
fn from(err: std::io::Error) -> Self {
Self::Io(err.to_string())
}
}
impl From<parquet::errors::ParquetError> for StoreError {
fn from(err: parquet::errors::ParquetError) -> Self {
Self::Parquet(err.to_string())
}
}
impl From<arrow::error::ArrowError> for StoreError {
fn from(err: arrow::error::ArrowError) -> Self {
Self::Parquet(err.to_string())
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub(crate) struct Coverage {
pub(crate) indexed_through: u64,
pub(crate) source_incarnation: PartitionIncarnation,
pub(crate) available: bool,
pub(crate) excision_scanned_through: u64,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub(crate) enum CoverageState {
NeverIndexed,
Indexed(Coverage),
Stale,
Destroyed,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(crate) enum ExcisionScan {
Clear,
Pending,
Unknown,
}
impl ExcisionScan {
const fn is_clear(self) -> bool {
matches!(self, Self::Clear)
}
}
pub(crate) trait JournalState: Sync {
fn source_incarnation(
&self,
partition: &str,
) -> impl std::future::Future<Output = Option<PartitionIncarnation>> + Send;
fn excision_since(
&self,
partition: &str,
scanned_through: u64,
) -> impl std::future::Future<Output = ExcisionScan> + Send;
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum Marker {
Live,
Destroyed,
}
fn projection_schema() -> Arc<Schema> {
Arc::new(Schema::new(vec![
Field::new("turn_id", DataType::Utf8, false),
Field::new("position", DataType::UInt64, false),
Field::new("term_hash", DataType::UInt32, false),
]))
}
const MAX_POSTINGS_ROWS: u64 = 10_000_000;
#[derive(Debug, Clone)]
struct Segment {
path: PathBuf,
state: CoverageState,
rows: u64,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(crate) struct SegmentStats {
pub(crate) segments: u32,
pub(crate) rows: u64,
pub(crate) unmerged_rows: u64,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(crate) struct Appended {
pub(crate) stats: SegmentStats,
pub(crate) clamped_unavailable: bool,
}
fn stats_from_rows(rows: &[u64]) -> SegmentStats {
let total = rows.iter().copied().fold(0u64, u64::saturating_add);
let base = rows.first().copied().unwrap_or(0);
SegmentStats {
segments: u32::try_from(rows.len()).unwrap_or(u32::MAX),
rows: total,
unmerged_rows: total.saturating_sub(base),
}
}
#[derive(Clone)]
pub(crate) struct SearchProjection {
root: Arc<PathBuf>,
#[cfg(test)]
fail_next_remove: Arc<std::sync::atomic::AtomicBool>,
}
impl SearchProjection {
pub(crate) fn open(root: PathBuf) -> Result<Self, StoreError> {
std::fs::create_dir_all(&root)?;
Ok(Self {
root: Arc::new(root),
#[cfg(test)]
fail_next_remove: Arc::new(std::sync::atomic::AtomicBool::new(false)),
})
}
#[cfg(test)]
pub(crate) fn fail_next_remove(&self) {
self.fail_next_remove
.store(true, std::sync::atomic::Ordering::SeqCst);
}
fn dir_for(root: &Path, partition: &str) -> PathBuf {
let encoded = encode_partition(partition)
.expect("the journal accepted this partition name, so it encodes");
root.join(format!("conversation_id={encoded}"))
}
fn segment_paths(root: &Path, partition: &str) -> Result<Vec<PathBuf>, StoreError> {
let dir = Self::dir_for(root, partition);
if !dir.exists() {
return Ok(Vec::new());
}
let mut paths = Vec::new();
for entry in std::fs::read_dir(&dir)? {
let path = entry?.path();
if path.extension().and_then(std::ffi::OsStr::to_str) == Some("parquet") {
paths.push(path);
}
}
paths.sort();
Ok(paths)
}
fn segments(root: &Path, partition: &str, key_id: &str) -> Result<Vec<Segment>, StoreError> {
let dir = Self::dir_for(root, partition);
if !dir.exists() {
return Ok(Vec::new());
}
let mut segments = Vec::new();
for entry in std::fs::read_dir(&dir)? {
let path = entry?.path();
if path.extension().and_then(std::ffi::OsStr::to_str) != Some("parquet") {
continue;
}
let file = std::fs::File::open(&path)?;
let builder = ParquetRecordBatchReaderBuilder::try_new(file)?;
let state = state_from(
builder.metadata().file_metadata().key_value_metadata(),
key_id,
)?;
let rows = u64::try_from(builder.metadata().file_metadata().num_rows())
.map_err(|_| StoreError::Corrupt)?;
segments.push(Segment { path, state, rows });
}
segments.sort_by(|a, b| a.path.cmp(&b.path));
Ok(segments)
}
pub(crate) async fn append(
&self,
partition: &str,
messages: &[IndexedMessage],
coverage: &Coverage,
key_id: &str,
) -> Result<Appended, StoreError> {
let (root, partition, messages, coverage, key_id) = (
Arc::clone(&self.root),
partition.to_owned(),
messages.to_vec(),
coverage.clone(),
key_id.to_owned(),
);
blocking(move || {
let state = Self::newest_state(&root, &partition, &key_id);
if matches!(state, Ok(CoverageState::Destroyed)) {
return Err(StoreError::Destroyed);
}
let newest = match state {
Ok(CoverageState::Indexed(newest)) => Some(newest),
_ => None,
};
let clamped_unavailable =
coverage.available && newest.as_ref().is_some_and(|newest| !newest.available);
let coverage = Coverage {
indexed_through: coverage.indexed_through,
source_incarnation: coverage.source_incarnation,
available: coverage.available && !clamped_unavailable,
excision_scanned_through: coverage
.excision_scanned_through
.max(newest.map_or(0, |newest| newest.excision_scanned_through)),
};
Self::write_segment(
&root,
&partition,
&messages,
&coverage,
&key_id,
Marker::Live,
)?;
Ok(Appended {
stats: Self::stats_of(&Self::dir_for(&root, &partition))?,
clamped_unavailable,
})
})
.await
}
pub(crate) async fn rebuild(
&self,
partition: &str,
messages: &[IndexedMessage],
coverage: &Coverage,
key_id: &str,
) -> Result<SegmentStats, StoreError> {
let (root, partition, messages, coverage, key_id) = (
Arc::clone(&self.root),
partition.to_owned(),
messages.to_vec(),
coverage.clone(),
key_id.to_owned(),
);
blocking(move || {
if matches!(
Self::newest_state(&root, &partition, &key_id),
Ok(CoverageState::Destroyed)
) {
return Err(StoreError::Destroyed);
}
let existing = Self::segment_paths(&root, &partition)?;
Self::write_segment(
&root,
&partition,
&messages,
&coverage,
&key_id,
Marker::Live,
)?;
for path in existing {
std::fs::remove_file(&path)?;
}
Self::stats_of(&Self::dir_for(&root, &partition))
})
.await
}
pub(crate) async fn compact(&self, partition: &str, key_id: &str) -> Result<(), StoreError> {
let (root, partition, key_id) = (
Arc::clone(&self.root),
partition.to_owned(),
key_id.to_owned(),
);
blocking(move || {
let segments = Self::segments(&root, &partition, &key_id)?;
if segments.len() < 2 {
return Ok(());
}
let Some(CoverageState::Indexed(coverage)) = segments.last().map(|s| s.state.clone())
else {
return Ok(());
};
let merged = Self::read_segments(&segments)?;
Self::write_segment(&root, &partition, &merged, &coverage, &key_id, Marker::Live)?;
for segment in segments {
std::fs::remove_file(&segment.path)?;
}
Ok(())
})
.await
}
pub(crate) async fn coverage(
&self,
partition: &str,
key_id: &str,
) -> Result<CoverageState, StoreError> {
let (root, partition, key_id) = (
Arc::clone(&self.root),
partition.to_owned(),
key_id.to_owned(),
);
blocking(move || Self::newest_state(&root, &partition, &key_id)).await
}
pub(crate) async fn verified_coverage(
&self,
partition: &str,
key_id: &str,
journal: &impl JournalState,
) -> Result<CoverageState, StoreError> {
let coverage = match self.coverage(partition, key_id).await? {
CoverageState::Indexed(coverage) => coverage,
other => return Ok(other),
};
let live = journal.source_incarnation(partition).await;
if live != Some(coverage.source_incarnation) {
return Ok(CoverageState::Stale);
}
let scan = journal
.excision_since(partition, coverage.excision_scanned_through)
.await;
let still_live = journal.source_incarnation(partition).await;
if scan.is_clear() && still_live == Some(coverage.source_incarnation) {
Ok(CoverageState::Indexed(coverage))
} else {
tracing::warn!(
partition,
?scan,
scanned_through = coverage.excision_scanned_through,
"search index holds a conversation whose journal has an excision it never \
applied; refusing until a rebuild strips it"
);
Ok(CoverageState::Stale)
}
}
fn newest_state(
root: &Path,
partition: &str,
key_id: &str,
) -> Result<CoverageState, StoreError> {
let Some(newest) = Self::segment_paths(root, partition)?.pop() else {
return Ok(CoverageState::NeverIndexed);
};
let file = std::fs::File::open(&newest)?;
let builder = ParquetRecordBatchReaderBuilder::try_new(file)?;
state_from(
builder.metadata().file_metadata().key_value_metadata(),
key_id,
)
}
pub(crate) async fn postings(
&self,
partition: &str,
key_id: &str,
) -> Result<Option<PostingsRecord>, StoreError> {
let (root, partition, key_id) = (
Arc::clone(&self.root),
partition.to_owned(),
key_id.to_owned(),
);
blocking(move || {
let segments = Self::segments(&root, &partition, &key_id)?;
if segments.is_empty()
|| matches!(
segments.last().map(|s| &s.state),
Some(&CoverageState::Destroyed)
)
{
return Ok(None);
}
Ok(Some(PostingsRecord {
messages: Self::read_segments(&segments)?,
}))
})
.await
}
pub(crate) async fn segment_stats(
&self,
partition: &str,
key_id: &str,
) -> Result<SegmentStats, StoreError> {
let (root, partition, key_id) = (
Arc::clone(&self.root),
partition.to_owned(),
key_id.to_owned(),
);
blocking(move || {
let segments = Self::segments(&root, &partition, &key_id)?;
let rows: Vec<u64> = segments.iter().map(|segment| segment.rows).collect();
let stats = stats_from_rows(&rows);
if stats.rows > MAX_POSTINGS_ROWS {
return Err(StoreError::TooLarge);
}
Ok(stats)
})
.await
}
pub(crate) async fn mark_unavailable(
&self,
partition: &str,
key_id: &str,
) -> Result<(), StoreError> {
let CoverageState::Indexed(coverage) = self.coverage(partition, key_id).await? else {
return Ok(());
};
if !coverage.available {
return Ok(());
}
self.append(
partition,
&[],
&Coverage {
available: false,
..coverage
},
key_id,
)
.await
.map(|_| ())
}
pub(crate) async fn destroy(&self, partition: &str, key_id: &str) -> Result<(), StoreError> {
let (root, partition, key_id) = (
Arc::clone(&self.root),
partition.to_owned(),
key_id.to_owned(),
);
blocking(move || {
let dir = Self::dir_for(&root, &partition);
if dir.exists() {
std::fs::remove_dir_all(&dir)?;
}
Self::write_segment(
&root,
&partition,
&[],
&Coverage {
indexed_through: 0,
source_incarnation: PartitionIncarnation::from_bytes([0; 32]),
available: false,
excision_scanned_through: 0,
},
&key_id,
Marker::Destroyed,
)
})
.await
}
fn stats_of(dir: &Path) -> Result<SegmentStats, StoreError> {
let mut paths = Vec::new();
if dir.exists() {
for entry in std::fs::read_dir(dir)? {
let path = entry?.path();
if path.extension().and_then(std::ffi::OsStr::to_str) == Some("parquet") {
paths.push(path);
}
}
}
paths.sort();
let mut rows = Vec::with_capacity(paths.len());
for path in &paths {
let file = std::fs::File::open(path)?;
let builder = ParquetRecordBatchReaderBuilder::try_new(file)?;
rows.push(
u64::try_from(builder.metadata().file_metadata().num_rows())
.map_err(|_| StoreError::Corrupt)?,
);
}
Ok(stats_from_rows(&rows))
}
pub(crate) async fn remove(&self, partition: &str) -> Result<(), StoreError> {
#[cfg(test)]
if self
.fail_next_remove
.swap(false, std::sync::atomic::Ordering::SeqCst)
{
return Err(std::io::Error::other("injected remove failure").into());
}
let (root, partition) = (Arc::clone(&self.root), partition.to_owned());
blocking(move || {
let dir = Self::dir_for(&root, &partition);
if dir.exists() {
std::fs::remove_dir_all(&dir)?;
}
Ok(())
})
.await
}
pub(crate) async fn authorized_paths(
&self,
authorized: Vec<String>,
) -> Result<Vec<PathBuf>, StoreError> {
let root = Arc::clone(&self.root);
blocking(move || {
let mut paths = Vec::new();
for partition in &authorized {
let dir = Self::dir_for(&root, partition);
if !dir.exists() {
continue;
}
for entry in std::fs::read_dir(&dir)? {
let path = entry?.path();
if path.extension().and_then(std::ffi::OsStr::to_str) == Some("parquet") {
paths.push(path);
}
}
}
paths.sort();
Ok(paths)
})
.await
}
pub(crate) fn root(&self) -> &Path {
&self.root
}
#[cfg(test)]
pub(crate) async fn write_format_for_test(
&self,
partition: &str,
messages: &[IndexedMessage],
coverage: &Coverage,
key_id: &str,
format: &str,
) -> Result<(), StoreError> {
let (root, partition, messages, coverage, key_id, format) = (
Arc::clone(&self.root),
partition.to_owned(),
messages.to_vec(),
coverage.clone(),
key_id.to_owned(),
format.to_owned(),
);
blocking(move || {
Self::write_segment_with_format(
&root,
&partition,
&messages,
&coverage,
&key_id,
Marker::Live,
&format,
)
})
.await
}
fn write_segment(
root: &Path,
partition: &str,
messages: &[IndexedMessage],
coverage: &Coverage,
key_id: &str,
marker: Marker,
) -> Result<(), StoreError> {
Self::write_segment_with_format(
root,
partition,
messages,
coverage,
key_id,
marker,
FORMAT_VERSION,
)
}
fn write_segment_with_format(
root: &Path,
partition: &str,
messages: &[IndexedMessage],
coverage: &Coverage,
key_id: &str,
marker: Marker,
format: &str,
) -> Result<(), StoreError> {
let dir = Self::dir_for(root, partition);
std::fs::create_dir_all(&dir)?;
let temp = dir.join(format!(
"seg.{}.{}.tmp",
std::process::id(),
TEMP_COUNTER.fetch_add(1, std::sync::atomic::Ordering::Relaxed)
));
let sequence = next_sequence(&dir)?;
let final_path = dir.join(format!(
"part-{:020}-{:010}.parquet",
coverage.indexed_through, sequence
));
if final_path.exists() {
return Err(StoreError::Io(format!(
"search projection segment {} already exists",
final_path.display()
)));
}
let props = WriterProperties::builder()
.set_compression(Compression::ZSTD(ZstdLevel::default()))
.set_column_bloom_filter_enabled(TERM_HASH_COLUMN.into(), true)
.set_sorting_columns(Some(vec![SortingColumn {
column_idx: 2,
descending: false,
nulls_first: false,
}]))
.set_key_value_metadata(Some(vec![
kv(META_INDEXED_THROUGH, &coverage.indexed_through.to_string()),
kv(
META_INCARNATION,
&hex(coverage.source_incarnation.as_bytes()),
),
kv(META_AVAILABLE, &coverage.available.to_string()),
kv(
META_EXCISION_SCANNED_THROUGH,
&coverage.excision_scanned_through.to_string(),
),
kv(META_KEY_ID, key_id),
kv(META_FORMAT, format),
kv(
META_DESTROYED,
&matches!(marker, Marker::Destroyed).to_string(),
),
]))
.build();
let written = (|| -> Result<(), StoreError> {
let file = open_owner_only(&temp)?;
let schema = projection_schema();
let mut writer =
ArrowWriter::try_new(file.try_clone()?, Arc::clone(&schema), Some(props))?;
if let Some(batch) = batch_for(&schema, messages)? {
writer.write(&batch)?;
}
writer.close()?;
file.sync_all()?;
Ok(())
})();
if let Err(err) = written {
let _ = std::fs::remove_file(&temp);
return Err(err);
}
std::fs::rename(&temp, &final_path)?;
Ok(())
}
fn read_segments(segments: &[Segment]) -> Result<Vec<IndexedMessage>, StoreError> {
let rows: u64 = segments.iter().map(|segment| segment.rows).sum();
if rows > MAX_POSTINGS_ROWS {
return Err(StoreError::TooLarge);
}
let mut by_position: std::collections::BTreeMap<u64, (String, Vec<u32>)> =
std::collections::BTreeMap::new();
for segment in segments {
let mut this_segment: std::collections::BTreeMap<u64, (String, Vec<u32>)> =
std::collections::BTreeMap::new();
let file = std::fs::File::open(&segment.path)?;
let reader = ParquetRecordBatchReaderBuilder::try_new(file)?.build()?;
for batch in reader {
let batch = batch?;
let turns = column::<StringArray>(&batch, 0)?;
let positions = column::<UInt64Array>(&batch, 1)?;
let hashes = column::<UInt32Array>(&batch, 2)?;
for row in 0..batch.num_rows() {
let entry = this_segment
.entry(positions.value(row))
.or_insert_with(|| (turns.value(row).to_owned(), Vec::new()));
let hash = hashes.value(row);
if hash != NO_TERMS_SENTINEL {
entry.1.push(hash);
}
}
}
by_position.extend(this_segment);
}
Ok(by_position
.into_iter()
.map(|(position, (turn_id, mut term_hashes))| {
term_hashes.sort_unstable();
term_hashes.dedup();
IndexedMessage {
position,
turn_id,
term_hashes,
}
})
.collect())
}
}
async fn blocking<T, F>(work: F) -> Result<T, StoreError>
where
F: FnOnce() -> Result<T, StoreError> + Send + 'static,
T: Send + 'static,
{
tokio::task::spawn_blocking(work)
.await
.map_err(|err| StoreError::Io(format!("search projection task: {err}")))?
}
fn batch_for(
schema: &Arc<Schema>,
messages: &[IndexedMessage],
) -> Result<Option<RecordBatch>, StoreError> {
let mut turns: Vec<&str> = Vec::new();
let mut positions: Vec<u64> = Vec::new();
let mut hashes: Vec<u32> = Vec::new();
for message in messages {
if message.term_hashes.is_empty() {
turns.push(&message.turn_id);
positions.push(message.position);
hashes.push(NO_TERMS_SENTINEL);
continue;
}
for &hash in &message.term_hashes {
turns.push(&message.turn_id);
positions.push(message.position);
hashes.push(hash);
}
}
if turns.is_empty() {
return Ok(None);
}
let mut rows: Vec<(usize, &&str)> = turns.iter().enumerate().collect();
rows.sort_by_key(|&(i, _)| (hashes[i], positions[i], turns[i]));
let order: Vec<usize> = rows.into_iter().map(|(i, _)| i).collect();
let turns: Vec<&str> = order.iter().map(|&i| turns[i]).collect();
let positions: Vec<u64> = order.iter().map(|&i| positions[i]).collect();
let hashes: Vec<u32> = order.iter().map(|&i| hashes[i]).collect();
RecordBatch::try_new(
Arc::clone(schema),
vec![
Arc::new(StringArray::from(turns)),
Arc::new(UInt64Array::from(positions)),
Arc::new(UInt32Array::from(hashes)),
],
)
.map(Some)
.map_err(Into::into)
}
fn column<T: 'static>(batch: &RecordBatch, index: usize) -> Result<&T, StoreError> {
batch
.column(index)
.as_any()
.downcast_ref::<T>()
.ok_or(StoreError::Corrupt)
}
fn state_from(
metadata: Option<&Vec<parquet::file::metadata::KeyValue>>,
expected_key_id: &str,
) -> Result<CoverageState, StoreError> {
let entries = metadata.ok_or(StoreError::Corrupt)?;
let get = |key: &str| {
entries
.iter()
.find(|kv| kv.key == key)
.and_then(|kv| kv.value.as_deref())
};
if get(META_DESTROYED) == Some("true") {
return Ok(CoverageState::Destroyed);
}
let indexed_through = get(META_INDEXED_THROUGH)
.and_then(|v| v.parse::<u64>().ok())
.ok_or(StoreError::Corrupt)?;
let source_incarnation =
PartitionIncarnation::parse_hex(get(META_INCARNATION).ok_or(StoreError::Corrupt)?)
.map_err(|_| StoreError::Corrupt)?;
let available = get(META_AVAILABLE)
.and_then(|v| v.parse::<bool>().ok())
.ok_or(StoreError::Corrupt)?;
let excision_scanned_through = get(META_EXCISION_SCANNED_THROUGH)
.and_then(|v| v.parse::<u64>().ok())
.ok_or(StoreError::Corrupt)?;
if get(META_FORMAT) != Some(FORMAT_VERSION) || get(META_KEY_ID) != Some(expected_key_id) {
return Err(StoreError::Corrupt);
}
Ok(CoverageState::Indexed(Coverage {
indexed_through,
source_incarnation,
available,
excision_scanned_through,
}))
}
fn next_sequence(dir: &Path) -> Result<u64, StoreError> {
let mut highest: Option<u64> = None;
for entry in std::fs::read_dir(dir)? {
let path = entry?.path();
if path.extension().and_then(std::ffi::OsStr::to_str) != Some("parquet") {
continue;
}
let Some(sequence) = path
.file_stem()
.and_then(std::ffi::OsStr::to_str)
.and_then(|stem| stem.rsplit_once('-'))
.and_then(|(_, seq)| seq.parse::<u64>().ok())
else {
continue;
};
highest = Some(highest.map_or(sequence, |seen: u64| seen.max(sequence)));
}
Ok(highest.map_or(0, |seen| seen.saturating_add(1)))
}
static TEMP_COUNTER: std::sync::atomic::AtomicU64 = std::sync::atomic::AtomicU64::new(0);
fn open_owner_only(path: &Path) -> Result<std::fs::File, StoreError> {
let mut options = std::fs::OpenOptions::new();
options.write(true).create_new(true);
#[cfg(unix)]
{
use std::os::unix::fs::OpenOptionsExt as _;
options.mode(0o600);
}
Ok(options.open(path)?)
}
#[cfg(test)]
pub(crate) fn unhex_for_test(incarnation: &str) -> Result<Vec<u8>, StoreError> {
unhex(incarnation)
}
fn kv(key: &str, value: &str) -> parquet::file::metadata::KeyValue {
parquet::file::metadata::KeyValue::new(key.to_owned(), value.to_owned())
}
fn hex(bytes: &[u8]) -> String {
bytes.iter().fold(String::new(), |mut out, byte| {
use std::fmt::Write as _;
let _ = write!(out, "{byte:02x}");
out
})
}
fn unhex(text: &str) -> Result<Vec<u8>, StoreError> {
let bytes = text.as_bytes();
if !bytes.len().is_multiple_of(2) {
return Err(StoreError::Corrupt);
}
bytes
.chunks_exact(2)
.map(|pair| {
let digits = std::str::from_utf8(pair).map_err(|_| StoreError::Corrupt)?;
u8::from_str_radix(digits, 16).map_err(|_| StoreError::Corrupt)
})
.collect()
}