use std::collections::HashMap;
use std::fs::File;
use std::os::unix::fs::FileExt;
use std::path::{Path, PathBuf};
use std::sync::atomic::{AtomicU64, Ordering};
use std::sync::mpsc::{Receiver, Sender, TryRecvError, channel};
use std::sync::{Arc, Condvar, Mutex};
use anyhow::{Result, anyhow};
use sha1::{Digest, Sha1};
use znippy_common::arrow::array::{ArrayRef, StringArray, UInt32Array, UInt64Array};
use znippy_common::arrow::datatypes::{DataType, Field, Schema, SchemaRef};
use znippy_common::arrow::record_batch::RecordBatch;
use znippy_zoomies::background::Job;
use znippy_zoomies::gatling_forkjoin::gatling_for_each;
use crate::archive_write::{ArchiveWrite, Extent, read_journal};
const MAX_DRAIN_BATCH: usize = 4096;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct IndexJob {
pub pack_id: u64,
pub offset: u64,
pub len: u64,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct IndexRow {
pub pack_id: u64,
pub offset: u64,
pub len: u64,
pub version: u32,
pub object_count: u32,
pub sha1: String,
}
pub fn index_schema() -> SchemaRef {
Arc::new(Schema::new(vec![
Field::new("pack_id", DataType::UInt64, false),
Field::new("blob_offset", DataType::UInt64, false),
Field::new("blob_size", DataType::UInt64, false),
Field::new("pack_version", DataType::UInt32, false),
Field::new("object_count", DataType::UInt32, false),
Field::new("pack_sha1", DataType::Utf8, false),
]))
}
fn rows_to_batch(rows: &[IndexRow]) -> Result<RecordBatch> {
let ids: ArrayRef = Arc::new(UInt64Array::from_iter_values(rows.iter().map(|r| r.pack_id)));
let offs: ArrayRef = Arc::new(UInt64Array::from_iter_values(rows.iter().map(|r| r.offset)));
let lens: ArrayRef = Arc::new(UInt64Array::from_iter_values(rows.iter().map(|r| r.len)));
let vers: ArrayRef = Arc::new(UInt32Array::from_iter_values(rows.iter().map(|r| r.version)));
let cnts: ArrayRef = Arc::new(UInt32Array::from_iter_values(
rows.iter().map(|r| r.object_count),
));
let sha: ArrayRef = Arc::new(StringArray::from_iter_values(
rows.iter().map(|r| r.sha1.as_str()),
));
RecordBatch::try_new(index_schema(), vec![ids, offs, lens, vers, cnts, sha])
.map_err(|e| anyhow!("index batch: {e}"))
}
fn build_row(archive: &File, job: IndexJob) -> IndexRow {
let mut buf = vec![0u8; job.len as usize];
if archive.read_exact_at(&mut buf, job.offset).is_err() {
return IndexRow {
pack_id: job.pack_id,
offset: job.offset,
len: job.len,
version: 0,
object_count: 0,
sha1: String::new(),
};
}
let (version, object_count) = if buf.len() >= 12 && &buf[0..4] == b"PACK" {
(
u32::from_be_bytes([buf[4], buf[5], buf[6], buf[7]]),
u32::from_be_bytes([buf[8], buf[9], buf[10], buf[11]]),
)
} else {
(0, 0)
};
let mut h = Sha1::new();
h.update(&buf);
IndexRow {
pack_id: job.pack_id,
offset: job.offset,
len: job.len,
version,
object_count,
sha1: hex::encode(h.finalize()),
}
}
pub trait ObjectAbsorb: Send + Sync {
fn absorb(&self, job: IndexJob) -> Result<()>;
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum Lookup {
Indexed(Box<IndexRow>),
ScanJournal(Vec<Extent>),
Unknowable,
}
#[derive(Default)]
struct Built {
batches: Vec<RecordBatch>,
published: HashMap<u64, IndexRow>,
absorb_failures: u64,
last_absorb_error: Option<String>,
}
struct Progress {
outstanding: Mutex<u64>,
caught_up: Condvar,
}
pub struct AccountIndexer {
account: String,
tx: Option<Sender<IndexJob>>,
built: Arc<Mutex<Built>>,
progress: Arc<Progress>,
journal: Option<PathBuf>,
worker: Option<Job<()>>,
}
impl AccountIndexer {
pub fn start(account: &str, archive: Arc<File>, journal: Option<PathBuf>) -> Self {
Self::start_with_absorber(account, archive, journal, None)
}
pub fn start_with_absorber(
account: &str,
archive: Arc<File>,
journal: Option<PathBuf>,
absorber: Option<Arc<dyn ObjectAbsorb>>,
) -> Self {
let (tx, rx): (Sender<IndexJob>, Receiver<IndexJob>) = channel();
let built = Arc::new(Mutex::new(Built::default()));
let progress = Arc::new(Progress {
outstanding: Mutex::new(0),
caught_up: Condvar::new(),
});
let w_built = built.clone();
let w_progress = progress.clone();
let worker = Job::spawn(move || index_worker(archive, rx, w_built, w_progress, absorber));
Self {
account: account.to_string(),
tx: Some(tx),
built,
progress,
journal,
worker: Some(worker),
}
}
pub fn account(&self) -> &str {
&self.account
}
pub fn submit(&self, job: IndexJob) -> Result<()> {
{
let mut o = self
.progress
.outstanding
.lock()
.map_err(|_| anyhow!("indexer progress poisoned"))?;
*o += 1;
}
self.tx
.as_ref()
.ok_or_else(|| anyhow!("indexer already closed"))?
.send(job)
.map_err(|_| anyhow!("indexer worker is gone"))
}
pub fn wait_caught_up(&self) {
let mut o = self.progress.outstanding.lock().unwrap();
while *o > 0 {
o = self.progress.caught_up.wait(o).unwrap();
}
}
pub fn is_indexed(&self, pack_id: u64) -> bool {
self.built.lock().unwrap().published.contains_key(&pack_id)
}
pub fn rows(&self) -> usize {
self.built.lock().unwrap().published.len()
}
pub fn absorb_failures(&self) -> u64 {
self.built.lock().unwrap().absorb_failures
}
pub fn last_absorb_error(&self) -> Option<String> {
self.built.lock().unwrap().last_absorb_error.clone()
}
pub fn tables(&self) -> Vec<RecordBatch> {
self.built.lock().unwrap().batches.clone()
}
pub fn lookup(&self, pack_id: u64) -> Lookup {
{
let b = self.built.lock().unwrap();
if let Some(r) = b.published.get(&pack_id) {
return Lookup::Indexed(Box::new(r.clone()));
}
}
match self.journal.as_deref() {
Some(p) => match read_journal(p) {
Ok(extents) => Lookup::ScanJournal(extents),
Err(_) => Lookup::Unknowable,
},
None => Lookup::Unknowable,
}
}
pub fn finish(mut self) -> Vec<RecordBatch> {
self.tx = None;
if let Some(w) = self.worker.take() {
let _ = w.join();
}
self.built.lock().unwrap().batches.clone()
}
}
impl Drop for AccountIndexer {
fn drop(&mut self) {
self.tx = None;
if let Some(w) = self.worker.take() {
let _ = w.join();
}
}
}
fn index_worker(
archive: Arc<File>,
rx: Receiver<IndexJob>,
built: Arc<Mutex<Built>>,
progress: Arc<Progress>,
absorber: Option<Arc<dyn ObjectAbsorb>>,
) {
loop {
let Ok(first) = rx.recv() else { return };
let mut batch = Vec::with_capacity(64);
batch.push(first);
loop {
if batch.len() >= MAX_DRAIN_BATCH {
break;
}
match rx.try_recv() {
Ok(j) => batch.push(j),
Err(TryRecvError::Empty) | Err(TryRecvError::Disconnected) => break,
}
}
let n = batch.len();
let arch = archive.as_ref();
let jobs = &batch;
let rows = gatling_for_each(n, 0, |i| build_row(arch, jobs[i]));
let mut ok: Vec<IndexRow> = Vec::with_capacity(rows.len());
let mut failed: Vec<String> = Vec::new();
for (job, row) in batch.iter().zip(rows) {
match absorber.as_deref() {
Some(a) => match a.absorb(*job) {
Ok(()) => ok.push(row),
Err(e) => failed.push(format!("pack {}: {e:#}", job.pack_id)),
},
None => ok.push(row),
}
}
let encoded = if ok.is_empty() {
None
} else {
rows_to_batch(&ok).ok()
};
{
let mut b = built.lock().unwrap();
if let Some(rb) = encoded {
b.batches.push(rb);
for r in ok {
b.published.insert(r.pack_id, r);
}
}
b.absorb_failures += failed.len() as u64;
if let Some(last) = failed.pop() {
b.last_absorb_error = Some(last);
}
}
let mut o = progress.outstanding.lock().unwrap();
*o = o.saturating_sub(n as u64);
if *o == 0 {
progress.caught_up.notify_all();
}
}
}
pub struct IndexerPool {
archive: Arc<File>,
journal: Option<PathBuf>,
absorber: Option<Arc<dyn ObjectAbsorb>>,
accounts: Mutex<HashMap<String, Arc<AccountIndexer>>>,
}
impl IndexerPool {
pub fn new(archive: &Path, journal: Option<PathBuf>) -> Result<Self> {
Self::build(archive, journal, None)
}
pub fn with_absorber(
archive: &Path,
journal: Option<PathBuf>,
absorber: Arc<dyn ObjectAbsorb>,
) -> Result<Self> {
Self::build(archive, journal, Some(absorber))
}
fn build(
archive: &Path,
journal: Option<PathBuf>,
absorber: Option<Arc<dyn ObjectAbsorb>>,
) -> Result<Self> {
let f = File::open(archive)
.map_err(|e| anyhow!("indexer: open {}: {e}", archive.display()))?;
Ok(Self {
archive: Arc::new(f),
journal,
absorber,
accounts: Mutex::new(HashMap::new()),
})
}
pub fn indexer(&self, account: &str) -> Arc<AccountIndexer> {
let mut m = self.accounts.lock().unwrap();
m.entry(account.to_string())
.or_insert_with(|| {
Arc::new(AccountIndexer::start_with_absorber(
account,
self.archive.clone(),
self.journal.clone(),
self.absorber.clone(),
))
})
.clone()
}
pub fn accounts(&self) -> usize {
self.accounts.lock().unwrap().len()
}
pub fn wait_caught_up(&self) {
let all: Vec<Arc<AccountIndexer>> =
self.accounts.lock().unwrap().values().cloned().collect();
for a in all {
a.wait_caught_up();
}
}
}
fn packs_already_acked(journal: Option<&Path>) -> Result<u64> {
match journal {
Some(p) if p.exists() => {
Ok(crate::archive_write::acked_packs(&read_journal(p)?).len() as u64)
}
_ => Ok(0),
}
}
pub struct PushPath {
writer: Box<dyn ArchiveWrite>,
pool: IndexerPool,
next_pack_id: AtomicU64,
}
impl PushPath {
pub fn new(writer: Box<dyn ArchiveWrite>, archive: &Path, journal: Option<PathBuf>) -> Result<Self> {
Ok(Self {
next_pack_id: AtomicU64::new(packs_already_acked(journal.as_deref())?),
writer,
pool: IndexerPool::new(archive, journal)?,
})
}
pub fn with_absorber(
writer: Box<dyn ArchiveWrite>,
archive: &Path,
journal: Option<PathBuf>,
absorber: Arc<dyn ObjectAbsorb>,
) -> Result<Self> {
Ok(Self {
next_pack_id: AtomicU64::new(packs_already_acked(journal.as_deref())?),
writer,
pool: IndexerPool::with_absorber(archive, journal, absorber)?,
})
}
pub fn name(&self) -> &'static str {
self.writer.name()
}
pub fn durability(&self) -> &'static str {
self.writer.durability()
}
pub fn push_pack(&self, account: &str, bytes: &[u8]) -> Result<(u64, Extent)> {
let (offset, len) = self.writer.append(bytes)?;
let pack_id = self.next_pack_id.fetch_add(1, Ordering::SeqCst);
self.pool.indexer(account).submit(IndexJob {
pack_id,
offset,
len,
})?;
Ok((pack_id, (offset, len)))
}
pub fn indexer(&self, account: &str) -> Arc<AccountIndexer> {
self.pool.indexer(account)
}
pub fn pool(&self) -> &IndexerPool {
&self.pool
}
}