mod sha256;
use std::fmt;
use std::path::Path;
use std::sync::OnceLock;
pub const FORMAT_VERSION: u32 = 1;
pub const SUPPORTED_SCHEMA: u32 = 1;
const MAGIC: &str = "flux-connectors-catalog-pack";
static EMBEDDED: &[u8] = include_bytes!("../catalog.pack");
#[derive(Debug)]
#[non_exhaustive]
pub enum Error {
NotAPack,
UnsupportedFormat {
found: u32,
},
UnsupportedSchema {
found: u32,
},
DigestMismatch {
stated: String,
computed: String,
},
NotText,
Malformed(String),
Io(String),
}
impl fmt::Display for Error {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
match self {
Error::NotAPack => write!(f, "not a {MAGIC} file"),
Error::UnsupportedFormat { found } => write!(
f,
"the pack declares container format {found}, but this reader implements \
{FORMAT_VERSION}; a newer pack needs a newer reader"
),
Error::UnsupportedSchema { found } => write!(
f,
"the pack carries document schema {found}, but this reader serves \
{SUPPORTED_SCHEMA}; refusing rather than handing out records it cannot vouch for"
),
Error::DigestMismatch { stated, computed } => write!(
f,
"the pack's stated digest {stated} is not the content's digest {computed}; the \
file is truncated, corrupted or edited"
),
Error::NotText => write!(f, "the pack is not UTF-8 text"),
Error::Malformed(what) => write!(f, "malformed pack: {what}"),
Error::Io(what) => write!(f, "cannot read the pack: {what}"),
}
}
}
impl std::error::Error for Error {}
enum Bytes {
Embedded(&'static [u8]),
Owned(Vec<u8>),
}
impl Bytes {
fn as_slice(&self) -> &[u8] {
match self {
Bytes::Embedded(bytes) => bytes,
Bytes::Owned(bytes) => bytes,
}
}
}
struct ProviderRow {
id: String,
start: usize,
len: usize,
}
struct OperationRow {
id: String,
provider: String,
service: String,
start: usize,
len: usize,
}
pub struct Pack {
bytes: Bytes,
payload_start: usize,
schema_version: u32,
digest: String,
providers: Vec<ProviderRow>,
operations: Vec<OperationRow>,
}
impl fmt::Debug for Pack {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
f.debug_struct("Pack")
.field("schema_version", &self.schema_version)
.field("digest", &self.digest)
.field("providers", &self.providers.len())
.field("operations", &self.operations.len())
.finish_non_exhaustive()
}
}
#[derive(Clone, Copy)]
pub struct Provider<'a> {
pack: &'a Pack,
row: &'a ProviderRow,
}
impl fmt::Debug for Provider<'_> {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
f.debug_struct("Provider")
.field("id", &self.row.id)
.finish_non_exhaustive()
}
}
impl<'a> Provider<'a> {
pub fn id(&self) -> &'a str {
&self.row.id
}
pub fn document(&self) -> &'a str {
self.pack.payload_slice(self.row.start, self.row.len)
}
pub fn operations(self) -> impl Iterator<Item = Operation<'a>> + 'a {
let pack = self.pack;
let id = self.row.id.as_str();
pack.operations
.iter()
.filter(move |row| row.provider == id)
.map(move |row| Operation { pack, row })
}
}
#[derive(Clone, Copy)]
pub struct Operation<'a> {
pack: &'a Pack,
row: &'a OperationRow,
}
impl fmt::Debug for Operation<'_> {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
f.debug_struct("Operation")
.field("id", &self.row.id)
.field("provider", &self.row.provider)
.field("service", &self.row.service)
.finish_non_exhaustive()
}
}
impl<'a> Operation<'a> {
pub fn id(&self) -> &'a str {
&self.row.id
}
pub fn provider(&self) -> &'a str {
&self.row.provider
}
pub fn service(&self) -> &'a str {
&self.row.service
}
pub fn record(&self) -> &'a str {
self.pack.payload_slice(self.row.start, self.row.len)
}
pub fn document(&self) -> &'a str {
self.pack
.provider(&self.row.provider)
.expect("an operation's provider row exists; verified at construction")
.document()
}
}
impl Pack {
pub fn load(path: impl AsRef<Path>) -> Result<Pack, Error> {
let path = path.as_ref();
let bytes = std::fs::read(path)
.map_err(|error| Error::Io(format!("{}: {error}", path.display())))?;
Self::from_bytes(bytes)
}
pub fn from_bytes(bytes: Vec<u8>) -> Result<Pack, Error> {
Self::parse(Bytes::Owned(bytes))
}
pub fn schema_version(&self) -> u32 {
self.schema_version
}
pub fn digest(&self) -> &str {
&self.digest
}
pub fn providers(&self) -> impl ExactSizeIterator<Item = Provider<'_>> {
self.providers
.iter()
.map(|row| Provider { pack: self, row })
}
pub fn provider(&self, id: &str) -> Option<Provider<'_>> {
self.providers
.binary_search_by(|row| row.id.as_str().cmp(id))
.ok()
.map(|index| Provider {
pack: self,
row: &self.providers[index],
})
}
pub fn operations(&self) -> impl ExactSizeIterator<Item = Operation<'_>> {
self.operations
.iter()
.map(|row| Operation { pack: self, row })
}
pub fn operation(&self, id: &str) -> Option<Operation<'_>> {
self.operations
.binary_search_by(|row| row.id.as_str().cmp(id))
.ok()
.map(|index| Operation {
pack: self,
row: &self.operations[index],
})
}
pub fn operations_of<'s>(&'s self, provider: &str) -> impl Iterator<Item = Operation<'s>> + 's {
let provider = provider.to_owned();
self.operations
.iter()
.filter(move |row| row.provider == provider)
.map(move |row| Operation { pack: self, row })
}
fn payload_slice(&self, start: usize, len: usize) -> &str {
let bytes = &self.bytes.as_slice()[self.payload_start + start..][..len];
std::str::from_utf8(bytes).expect("every span was boundary-checked at construction")
}
fn parse(bytes: Bytes) -> Result<Pack, Error> {
let text = std::str::from_utf8(bytes.as_slice()).map_err(|_| Error::NotText)?;
let mut offset = 0usize;
let mut next_line = |what: &'static str| -> Result<(&str, usize), Error> {
let rest = &text[offset..];
let end = rest
.find('\n')
.ok_or_else(|| Error::Malformed(format!("the file ends before its {what} line")))?;
let line = &rest[..end];
offset += end + 1;
Ok((line, offset))
};
let (magic_line, _) = next_line("magic")?;
let mut words = magic_line.split(' ');
if words.next() != Some(MAGIC) {
return Err(Error::NotAPack);
}
let found: u32 = words
.next()
.and_then(|version| version.parse().ok())
.ok_or_else(|| Error::Malformed(format!("no format version in `{magic_line}`")))?;
if found != FORMAT_VERSION {
return Err(Error::UnsupportedFormat { found });
}
let (digest_line, digested_from) = next_line("digest")?;
let stated = digest_line
.strip_prefix("digest sha256 ")
.ok_or_else(|| Error::Malformed(format!("not a digest line: `{digest_line}`")))?
.to_owned();
let computed = sha256::hex_digest(&bytes.as_slice()[digested_from..]);
if stated != computed {
return Err(Error::DigestMismatch { stated, computed });
}
let mut schema_version: Option<u32> = None;
let mut declared_providers: Option<usize> = None;
let mut declared_operations: Option<usize> = None;
let mut providers: Vec<ProviderRow> = Vec::new();
let mut operations: Vec<OperationRow> = Vec::new();
let payload_start;
let declared_payload;
loop {
let (line, after) = next_line("payload")?;
let mut fields = line.split(' ');
match fields.next() {
Some("schema") => {
let found = parse_field(line, fields.next())?;
if found != SUPPORTED_SCHEMA {
return Err(Error::UnsupportedSchema { found });
}
schema_version = Some(found);
}
Some("providers") => declared_providers = Some(parse_field(line, fields.next())?),
Some("operations") => declared_operations = Some(parse_field(line, fields.next())?),
Some("p") => {
let id = required(line, fields.next())?.to_owned();
let start = parse_field(line, fields.next())?;
let len = parse_field(line, fields.next())?;
providers.push(ProviderRow { id, start, len });
}
Some("o") => {
let id = required(line, fields.next())?.to_owned();
let provider = required(line, fields.next())?.to_owned();
let service = required(line, fields.next())?.to_owned();
let start = parse_field(line, fields.next())?;
let len = parse_field(line, fields.next())?;
operations.push(OperationRow {
id,
provider,
service,
start,
len,
});
}
Some("payload") => {
declared_payload = parse_field(line, fields.next())?;
payload_start = after;
break;
}
_ => {}
}
}
if schema_version.is_none() {
return Err(Error::Malformed("no schema line".into()));
}
let payload = &text[payload_start..];
if payload.len() != declared_payload {
return Err(Error::Malformed(format!(
"the payload is {} bytes where the header declares {declared_payload}",
payload.len()
)));
}
if declared_providers != Some(providers.len()) {
return Err(Error::Malformed(format!(
"{} provider rows where the header declares {declared_providers:?}",
providers.len()
)));
}
if declared_operations != Some(operations.len()) {
return Err(Error::Malformed(format!(
"{} operation rows where the header declares {declared_operations:?}",
operations.len()
)));
}
for (start, len, what) in providers
.iter()
.map(|row| (row.start, row.len, row.id.as_str()))
.chain(
operations
.iter()
.map(|row| (row.start, row.len, row.id.as_str())),
)
{
let span = start
.checked_add(len)
.and_then(|end| payload.get(start..end));
if span.is_none() {
return Err(Error::Malformed(format!(
"`{what}`'s span {start}+{len} is not a slice of the {declared_payload}-byte \
payload"
)));
}
}
providers.sort_by(|a, b| a.id.cmp(&b.id));
operations.sort_by(|a, b| a.id.cmp(&b.id));
for row in &operations {
if providers
.binary_search_by(|held| held.id.cmp(&row.provider))
.is_err()
{
return Err(Error::Malformed(format!(
"operation `{}` names provider `{}`, which has no row",
row.id, row.provider
)));
}
}
Ok(Pack {
bytes,
payload_start,
schema_version: schema_version.expect("checked above"),
digest: stated,
providers,
operations,
})
}
}
fn required<'a>(line: &str, field: Option<&'a str>) -> Result<&'a str, Error> {
field
.filter(|value| !value.is_empty())
.ok_or_else(|| Error::Malformed(format!("a field is missing in `{line}`")))
}
fn parse_field<T: std::str::FromStr>(line: &str, field: Option<&str>) -> Result<T, Error> {
field
.and_then(|value| value.parse().ok())
.ok_or_else(|| Error::Malformed(format!("not a number where one is required: `{line}`")))
}
pub fn embedded() -> &'static Pack {
static PACK: OnceLock<Pack> = OnceLock::new();
PACK.get_or_init(|| {
Pack::parse(Bytes::Embedded(EMBEDDED))
.expect("the embedded catalog.pack verifies; it is committed beside this crate")
})
}
pub fn providers() -> impl ExactSizeIterator<Item = Provider<'static>> {
embedded().providers()
}
pub fn provider(id: &str) -> Option<Provider<'static>> {
embedded().provider(id)
}
pub fn operation(id: &str) -> Option<Operation<'static>> {
embedded().operation(id)
}
pub fn operations_of(provider: &str) -> impl Iterator<Item = Operation<'static>> {
embedded().operations_of(provider)
}