use anyhow::{Context, Result, anyhow};
use clap::Args;
use smol::{
fs::{File, OpenOptions},
io::{AsyncBufReadExt, BufReader},
};
use std::{collections::BTreeSet, ffi::OsStr, path::PathBuf};
use ufotofu::{
Consumer, ConsumerExt, IntoConsumer, IntoProducer,
consumer::compat::writer::writer_to_bulk_consumer, queues::new_unbounded_elastic,
};
use willow25::{
drop_format::{DropEncoder, EncodeDropError, ExportDropError, export_drop},
entry::NamespaceId,
groupings::{Area, Keylike},
prelude::AuthorisedEntry,
storage::{PersistentStore, Store},
};
use crate::util::{SNEAKERWEB_NAMESPACE_ID_BYTES, get_domain, sneakerweb_dir};
#[derive(Args)]
pub struct ExportArgs {
pub dest: PathBuf,
#[arg(short, long)]
pub collection: Option<PathBuf>,
}
pub async fn export_sneak(args: &ExportArgs) -> Result<PathBuf> {
let mut dest = args.dest.clone();
if dest.is_dir() {
dest = dest.join("exported");
}
if dest.extension() != Some(OsStr::new("snk")) {
dest.add_extension("snk");
}
let namespace = NamespaceId::from_bytes(&SNEAKERWEB_NAMESPACE_ID_BYTES);
let sneakerweb_fs_path = sneakerweb_dir().await?;
let mut store = PersistentStore::new(sneakerweb_fs_path).await?;
let destination = OpenOptions::new()
.create(true)
.write(true)
.truncate(true)
.open(&dest)
.await?;
let drop_consumer = writer_to_bulk_consumer(destination, new_unbounded_elastic());
let mut encoder = DropEncoder::new(drop_consumer);
let included = match &args.collection {
Some(collection) => {
let mut included = BTreeSet::new();
let mut collection_reader = BufReader::new(
File::open(collection)
.await
.context(format!("could not open {}", collection.to_string_lossy()))?,
);
let mut line = String::new();
while let Ok(n) = collection_reader.read_line(&mut line).await
&& n > 0
{
let (domain_id, _) = get_domain(Some(&line), "N/A").context(format!(
"invalid domain '{}' found in collection '{}'",
&line.trim(),
collection.to_string_lossy()
))?;
included.insert(domain_id);
line.clear();
}
Some(included)
}
None => None,
};
let mut entries: Vec<AuthorisedEntry> = vec![];
store
.get_area(
&namespace,
&Area::full(),
&mut (&mut entries)
.into_consumer()
.to_filter(async |entry| match &included {
Some(collection) => collection.contains(entry.subspace_id()),
None => true,
}),
)
.await?;
export_drop(&mut store, &mut entries.into_producer(), &mut encoder)
.await
.map_err(|err| match err {
ExportDropError::StoreError(err) => anyhow!(err),
ExportDropError::ConsumerError(err) => match err {
EncodeDropError::ConsumerError(err) => anyhow!(err),
EncodeDropError::ConsumedBytesMismatch => {
anyhow!("the encoder received a different number of bytes than expected")
}
EncodeDropError::ArchitectureTooSmall => {
anyhow!("encountered an entry too large for this device")
}
},
ExportDropError::EntryDeleted => {
anyhow!("an entry was deleted concurrently with writing the export")
}
})?;
encoder.flush().await.map_err(|err| match err {
EncodeDropError::ConsumerError(err) => anyhow!(err),
EncodeDropError::ConsumedBytesMismatch => {
anyhow!("the encoder receieved a different number of bytes than expected")
}
EncodeDropError::ArchitectureTooSmall => {
anyhow!("encountered an entry too large for this device")
}
})?;
Ok(dest)
}