use crate::{
build_pyramid_meta_algo, Dictionary, DictionaryBuilder, GraphIndexBuilder, PyramidAlgo,
DEFAULT_TILE_BUDGET,
};
pub type RawTriple = (String, String, String);
pub type RawQuad = (String, String, String, Option<String>);
#[derive(Debug, thiserror::Error)]
#[non_exhaustive]
pub enum IngestError {
#[error("line {0}: {1}")]
Line(usize, &'static str),
#[error("turtle: {0}")]
Turtle(String),
#[error("rdf/xml: {0}")]
RdfXml(String),
#[error("unknown input format: {0} (expected nt, nq, ttl, or rdf/xml)")]
UnknownFormat(String),
#[error("io: {0}")]
Io(String),
}
fn estimate_statements(input: &str) -> usize {
bytecount_newlines(input).max(1)
}
fn bytecount_newlines(input: &str) -> usize {
input.as_bytes().iter().filter(|&&b| b == b'\n').count()
}
fn parse_nq_line(raw: &str, lineno: usize) -> Result<Option<RawQuad>, IngestError> {
let line = raw.trim();
if line.is_empty() || line.starts_with('#') {
return Ok(None);
}
let stripped = line
.strip_suffix('.')
.ok_or(IngestError::Line(lineno, "missing trailing '.'"))?
.trim_end();
let (s, rest) = take_term(stripped).ok_or(IngestError::Line(lineno, "bad subject"))?;
let (p, rest) =
take_term(rest.trim_start()).ok_or(IngestError::Line(lineno, "bad predicate"))?;
let (o, rest) = take_term(rest.trim_start()).ok_or(IngestError::Line(lineno, "bad object"))?;
let rest = rest.trim();
let graph = if rest.is_empty() {
None
} else {
let (g, tail) = take_term(rest).ok_or(IngestError::Line(lineno, "bad graph"))?;
if !tail.trim().is_empty() {
return Err(IngestError::Line(lineno, "trailing content after graph"));
}
Some(g)
};
Ok(Some((s, p, o, graph)))
}
fn parse_nt_line(raw: &str, lineno: usize) -> Result<Option<RawTriple>, IngestError> {
let line = raw.trim();
if line.is_empty() || line.starts_with('#') {
return Ok(None);
}
let stripped = line
.strip_suffix('.')
.ok_or(IngestError::Line(lineno, "missing trailing '.'"))?
.trim_end();
let (s, rest) = take_term(stripped).ok_or(IngestError::Line(lineno, "bad subject"))?;
let (p, rest) =
take_term(rest.trim_start()).ok_or(IngestError::Line(lineno, "bad predicate"))?;
let (o, rest) = take_term(rest.trim_start()).ok_or(IngestError::Line(lineno, "bad object"))?;
if !rest.trim().is_empty() {
return Err(IngestError::Line(lineno, "trailing content after object"));
}
Ok(Some((s, p, o)))
}
pub fn parse_quads(input: &str) -> Result<Vec<RawQuad>, IngestError> {
let mut out = Vec::with_capacity(estimate_statements(input));
for (i, raw) in input.lines().enumerate() {
if let Some(q) = parse_nq_line(raw, i + 1)? {
out.push(q);
}
}
Ok(out)
}
pub fn parse(input: &str) -> Result<Vec<RawTriple>, IngestError> {
let mut out = Vec::with_capacity(estimate_statements(input));
for (i, raw) in input.lines().enumerate() {
if let Some(t) = parse_nt_line(raw, i + 1)? {
out.push(t);
}
}
Ok(out)
}
pub fn parse_reader<R: std::io::BufRead>(
reader: R,
format: &str,
cap: usize,
) -> Result<Vec<RawQuad>, IngestError> {
let mut out = Vec::with_capacity(cap);
stream_reader(reader, format, &mut |q| out.push(q))?;
Ok(out)
}
pub fn stream_reader<R: std::io::BufRead>(
reader: R,
format: &str,
f: &mut dyn FnMut(RawQuad),
) -> Result<(), IngestError> {
if format != "nt" && format != "nq" {
return Err(IngestError::UnknownFormat(format.to_string()));
}
for (i, line) in reader.lines().enumerate() {
let line = line.map_err(|e| IngestError::Io(e.to_string()))?;
if format == "nq" {
if let Some(q) = parse_nq_line(&line, i + 1)? {
f(q);
}
} else if let Some((s, p, o)) = parse_nt_line(&line, i + 1)? {
f((s, p, o, None));
}
}
Ok(())
}
pub fn parse_turtle(text: &str) -> Result<Vec<RawTriple>, IngestError> {
let mut out = Vec::new();
for r in oxttl::TurtleParser::new()
.with_quoted_triples()
.for_reader(text.as_bytes())
{
let t = r.map_err(|e| IngestError::Turtle(e.to_string()))?;
out.push((
t.subject.to_string(),
t.predicate.to_string(),
t.object.to_string(),
));
}
Ok(out)
}
pub fn parse_rdfxml(text: &str) -> Result<Vec<RawTriple>, IngestError> {
let mut out = Vec::new();
for r in oxrdfxml::RdfXmlParser::new().for_reader(text.as_bytes()) {
let t = r.map_err(|e| IngestError::RdfXml(e.to_string()))?;
out.push((
t.subject.to_string(),
t.predicate.to_string(),
t.object.to_string(),
));
}
Ok(out)
}
pub fn parse_statements(text: &str, format: &str) -> Result<Vec<RawQuad>, IngestError> {
match format {
"nq" => parse_quads(text),
"ttl" => Ok(parse_turtle(text)?
.into_iter()
.map(|(s, p, o)| (s, p, o, None))
.collect()),
"rdfxml" => Ok(parse_rdfxml(text)?
.into_iter()
.map(|(s, p, o)| (s, p, o, None))
.collect()),
"nt" => Ok(parse(text)?
.into_iter()
.map(|(s, p, o)| (s, p, o, None))
.collect()),
other => Err(IngestError::UnknownFormat(other.to_string())),
}
}
pub(crate) fn take_term(s: &str) -> Option<(String, &str)> {
let bytes = s.as_bytes();
let first = *bytes.first()?;
match first {
b'<' if bytes.get(1) == Some(&b'<') => {
let (inner, rdf12) = match s[2..].strip_prefix('(') {
Some(after) => (after.trim_start(), true),
None => (s[2..].trim_start(), false),
};
let (subj, r) = take_term(inner)?;
let (pred, r) = take_term(r.trim_start())?;
let (obj, r) = take_term(r.trim_start())?;
let r = r.trim_start();
let rest = if rdf12 {
r.strip_prefix(")>>")?
} else {
r.strip_prefix(">>")?
};
Some((format!("<<{subj} {pred} {obj}>>"), rest))
}
b'<' => {
let end = s.find('>')?;
Some((s[..=end].to_string(), &s[end + 1..]))
}
b'_' => {
let end = s.find(char::is_whitespace).unwrap_or(s.len());
Some((s[..end].to_string(), &s[end..]))
}
b'"' => {
let mut i = 1;
let b = s.as_bytes();
while i < b.len() {
match b[i] {
b'\\' => i += 2, b'"' => break,
_ => i += 1,
}
}
if i >= b.len() {
return None; }
let mut end = i + 1; if s[end..].starts_with("^^<") {
let close = s[end..].find('>')? + end;
end = close + 1;
} else if s[end..].starts_with('@') {
let mut j = end + 1; while j < b.len() && (b[j].is_ascii_alphanumeric() || b[j] == b'-') {
j += 1;
}
end = j;
}
Some((s[..end].to_string(), &s[end..]))
}
_ => None,
}
}
pub(crate) fn quoted_triple_parts(t: &str) -> Option<(String, String, String)> {
let inner = t.strip_prefix("<<")?.strip_suffix(">>")?.trim();
let (s, r) = take_term(inner)?;
let (p, r) = take_term(r.trim_start())?;
let (o, r) = take_term(r.trim_start())?;
if !r.trim().is_empty() {
return None;
}
Some((s, p, o))
}
#[derive(Debug, Clone, Copy)]
pub struct BuildStats {
pub statements: usize,
pub default_triples: usize,
pub named_graphs: usize,
pub terms: usize,
pub pyramid_levels: u16,
}
pub fn assemble_dataset(quads: Vec<RawQuad>, metadata: &[u8]) -> (Vec<u8>, BuildStats) {
let blob = metadata.to_vec();
assemble_dataset_with(quads, move |_, _| blob)
}
pub fn assemble_dataset_with(
quads: Vec<RawQuad>,
metadata: impl FnOnce(&BuildStats, &[RawQuad]) -> Vec<u8>,
) -> (Vec<u8>, BuildStats) {
assemble_dataset_with_opts(quads, true, false, None, metadata)
}
pub fn assemble_dataset_with_opts(
quads: Vec<RawQuad>,
with_pyramid: bool,
with_text_index: bool,
type_override: Option<&str>,
metadata: impl FnOnce(&BuildStats, &[RawQuad]) -> Vec<u8>,
) -> (Vec<u8>, BuildStats) {
assemble_dataset_with_opts_algo(
quads,
with_pyramid,
with_text_index,
type_override,
PyramidAlgo::Louvain,
metadata,
)
}
#[allow(clippy::too_many_arguments)]
pub fn assemble_dataset_with_opts_algo(
quads: Vec<RawQuad>,
with_pyramid: bool,
with_text_index: bool,
type_override: Option<&str>,
algo: PyramidAlgo,
metadata: impl FnOnce(&BuildStats, &[RawQuad]) -> Vec<u8>,
) -> (Vec<u8>, BuildStats) {
use std::collections::BTreeMap;
let mut db = DictionaryBuilder::new();
for (s, p, o, _) in &quads {
db.observe(s, p, o);
}
let dict = db.build();
let mut default_triples = Vec::new();
let mut named: BTreeMap<String, Vec<(u32, u32, u32)>> = BTreeMap::new();
for (s, p, o, g) in &quads {
let t = dict.encode(s, p, o).expect("observed term");
match g {
None => default_triples.push(t),
Some(graph) => named.entry(graph.clone()).or_default().push(t),
}
}
let stats = BuildStats {
statements: quads.len(),
default_triples: default_triples.len(),
named_graphs: named.len(),
terms: dict.term_count() as usize,
pyramid_levels: 0,
};
let blob = metadata(&stats, &quads);
drop(quads);
finish_assembly(
dict,
default_triples,
named,
with_pyramid,
with_text_index,
type_override,
algo,
blob,
stats,
)
}
pub fn assemble_dataset_streaming<S>(
stream: S,
with_pyramid: bool,
with_text_index: bool,
type_override: Option<&str>,
metadata: impl FnOnce(&BuildStats, &Dictionary, &[(u32, u32, u32)]) -> Vec<u8>,
) -> Result<(Vec<u8>, BuildStats), IngestError>
where
S: FnMut(&mut dyn FnMut(RawQuad)) -> Result<(), IngestError>,
{
assemble_dataset_streaming_algo(
stream,
with_pyramid,
with_text_index,
type_override,
PyramidAlgo::Louvain,
metadata,
)
}
#[allow(clippy::too_many_arguments)]
pub fn assemble_dataset_streaming_algo<S>(
mut stream: S,
with_pyramid: bool,
with_text_index: bool,
type_override: Option<&str>,
algo: PyramidAlgo,
metadata: impl FnOnce(&BuildStats, &Dictionary, &[(u32, u32, u32)]) -> Vec<u8>,
) -> Result<(Vec<u8>, BuildStats), IngestError>
where
S: FnMut(&mut dyn FnMut(RawQuad)) -> Result<(), IngestError>,
{
use std::collections::BTreeMap;
let mut db = DictionaryBuilder::new();
stream(&mut |(s, p, o, _g)| db.observe(&s, &p, &o))?;
let dict = db.build();
let mut default_triples: Vec<(u32, u32, u32)> = Vec::new();
let mut named: BTreeMap<String, Vec<(u32, u32, u32)>> = BTreeMap::new();
stream(&mut |(s, p, o, g)| {
let t = dict.encode(&s, &p, &o).expect("observed term");
match g {
None => default_triples.push(t),
Some(graph) => named.entry(graph).or_default().push(t),
}
})?;
let statements = default_triples.len() + named.values().map(Vec::len).sum::<usize>();
let stats = BuildStats {
statements,
default_triples: default_triples.len(),
named_graphs: named.len(),
terms: dict.term_count() as usize,
pyramid_levels: 0,
};
let blob = metadata(&stats, &dict, &default_triples);
Ok(finish_assembly(
dict,
default_triples,
named,
with_pyramid,
with_text_index,
type_override,
algo,
blob,
stats,
))
}
#[allow(clippy::too_many_arguments)]
fn finish_assembly(
dict: Dictionary,
default_triples: Vec<(u32, u32, u32)>,
named: std::collections::BTreeMap<String, Vec<(u32, u32, u32)>>,
with_pyramid: bool,
with_text_index: bool,
type_override: Option<&str>,
algo: PyramidAlgo,
blob: Vec<u8>,
mut stats: BuildStats,
) -> (Vec<u8>, BuildStats) {
let has_named = !named.is_empty();
let (meta, levels) = if with_pyramid {
build_pyramid_meta_algo(
&dict,
&default_triples,
DEFAULT_TILE_BUDGET,
type_override,
algo,
)
} else {
(Vec::new(), 0)
};
stats.pyramid_levels = levels;
let text_index = if with_text_index {
crate::file::compute_text_index(&dict, &default_triples)
} else {
Vec::new()
};
let codec = crate::file::writer_codec();
let dict_container = crate::file::encode_dict_container(&dict, codec);
let term_count = dict.term_count() as u64;
let has_quoted_triples = dict.has_quoted_triples();
drop(dict);
let build_index = |triples: Vec<(u32, u32, u32)>| -> crate::GraphIndex {
let n = triples.len();
let b = GraphIndexBuilder::from_triples(triples);
if n > LOWMEM_TRIPLE_THRESHOLD {
b.build_seq()
} else {
b.build()
}
};
let def = build_index(default_triples);
let named_indexes: Vec<(String, crate::GraphIndex)> = named
.into_iter()
.map(|(g, ts)| (g, build_index(ts)))
.collect();
let bytes = crate::file::write_dataset_from_parts(
&dict_container,
term_count,
&def,
&named_indexes,
has_named,
has_quoted_triples,
&meta,
levels,
&blob,
&text_index,
codec,
);
(bytes, stats)
}
const LOWMEM_TRIPLE_THRESHOLD: usize = 30_000_000;
#[cfg(test)]
mod tests {
use super::*;
use crate::Rete;
#[test]
fn parses_iris_bnodes_literals() {
let input = r#"
# a comment
<http://ex/Alice> <http://ex/knows> <http://ex/Bob> .
<http://ex/Alice> <http://ex/age> "30"^^<http://www.w3.org/2001/XMLSchema#integer> .
<http://ex/Bob> <http://ex/label> "Bob"@en .
_:b0 <http://ex/p> "plain" .
"#;
let t = parse(input).unwrap();
assert_eq!(t.len(), 4);
assert_eq!(t[0].0, "<http://ex/Alice>");
assert_eq!(t[0].2, "<http://ex/Bob>");
assert_eq!(t[1].2, "\"30\"^^<http://www.w3.org/2001/XMLSchema#integer>");
assert_eq!(t[2].2, "\"Bob\"@en");
assert_eq!(t[3].0, "_:b0");
assert_eq!(t[3].2, "\"plain\"");
}
#[test]
fn literal_with_spaces_and_escaped_quote() {
let input = r#"<http://ex/s> <http://ex/p> "a \"quoted\" phrase here" ."#;
let t = parse(input).unwrap();
assert_eq!(t.len(), 1);
assert_eq!(t[0].2, r#""a \"quoted\" phrase here""#);
}
#[test]
fn quoted_triple_langtagged_object() {
let s = r#"<<<http://ex/sp> <http://ex/name> "Hirondelle rustique"@fr>>"#;
let (tok, rest) = take_term(s).unwrap();
assert_eq!(
tok,
r#"<<<http://ex/sp> <http://ex/name> "Hirondelle rustique"@fr>>"#
);
assert_eq!(rest, "");
let line = r#"<<<http://ex/sp> <http://ex/name> "Oreneta vulgar"@ca>> <http://purl.org/dc/terms/source> "Catalogue of Life" ."#;
let t = parse(line).unwrap();
assert_eq!(t.len(), 1);
assert_eq!(
t[0].0,
r#"<<<http://ex/sp> <http://ex/name> "Oreneta vulgar"@ca>>"#
);
assert_eq!(t[0].2, r#""Catalogue of Life""#);
assert_eq!(take_term(r#""x"@pt-BR ."#).unwrap().0, r#""x"@pt-BR"#);
}
#[test]
fn rdf12_triple_term_same_token_as_rdf_star() {
let rdf12 = r#"<http://ex/r> <http://ex/p> <<( <http://ex/s> <http://ex/q> "o" )>> ."#;
let star = r#"<http://ex/r> <http://ex/p> << <http://ex/s> <http://ex/q> "o" >> ."#;
let a = parse(rdf12).unwrap();
let b = parse(star).unwrap();
assert_eq!(a.len(), 1);
assert_eq!(a[0].2, r#"<<<http://ex/s> <http://ex/q> "o">>"#);
assert_eq!(
a[0].2, b[0].2,
"RDF 1.2 and RDF-star must dedupe to one token"
);
let (tok, rest) =
take_term(r#"<<( <http://ex/s> <http://ex/q> <http://ex/o> )>>"#).unwrap();
assert_eq!(tok, "<<<http://ex/s> <http://ex/q> <http://ex/o>>>");
assert_eq!(rest, "");
}
#[test]
fn rejects_missing_dot() {
assert!(parse("<a> <b> <c>").is_err());
}
#[test]
fn parses_quads_with_and_without_graph() {
let input = "<http://ex/a> <http://ex/p> <http://ex/b> .\n\
<http://ex/a> <http://ex/p> <http://ex/c> <http://ex/g> .";
let q = parse_quads(input).unwrap();
assert_eq!(q.len(), 2);
assert_eq!(q[0].3, None); assert_eq!(q[1].3, Some("<http://ex/g>".to_string()));
}
#[test]
fn parse_reader_matches_text_parse() {
let nt = "<http://ex/a> <http://ex/p> <http://ex/b> .\r\n\
# comment\n\
\n\
_:b0 <http://ex/q> \"lit\"@en .\n";
let via_text = parse_statements(nt, "nt").unwrap();
let via_reader = parse_reader(std::io::Cursor::new(nt), "nt", 0).unwrap();
assert_eq!(via_text, via_reader);
let nq = "<http://ex/a> <http://ex/p> <http://ex/b> .\n\
<http://ex/a> <http://ex/p> <http://ex/c> <http://ex/g> .\n";
assert_eq!(
parse_quads(nq).unwrap(),
parse_reader(std::io::Cursor::new(nq), "nq", 0).unwrap()
);
assert!(parse_reader(std::io::Cursor::new(nt), "ttl", 0).is_err());
}
#[test]
fn parse_statements_dispatches_by_format() {
let nt = "<http://ex/a> <http://ex/p> <http://ex/b> .";
assert_eq!(parse_statements(nt, "nt").unwrap().len(), 1);
let ttl = "@prefix ex: <http://ex/> .\nex:A ex:knows ex:B , ex:C .";
assert_eq!(parse_statements(ttl, "ttl").unwrap().len(), 2);
assert!(parse_statements(nt, "trig").is_err());
}
#[test]
fn parses_rdfxml_owl() {
let xml = r#"<?xml version="1.0"?>
<rdf:RDF xmlns:rdf="http://www.w3.org/1999/02/22-rdf-syntax-ns#"
xmlns:rdfs="http://www.w3.org/2000/01/rdf-schema#"
xmlns:owl="http://www.w3.org/2002/07/owl#">
<owl:Class rdf:about="http://ex/Dog">
<rdfs:subClassOf rdf:resource="http://ex/Animal"/>
<rdfs:label>Dog</rdfs:label>
</owl:Class>
</rdf:RDF>"#;
let triples = parse_rdfxml(xml).unwrap();
assert_eq!(triples.len(), 3);
assert!(triples.iter().any(|(s, p, o)| s == "<http://ex/Dog>"
&& p == "<http://www.w3.org/2000/01/rdf-schema#subClassOf>"
&& o == "<http://ex/Animal>"));
assert_eq!(parse_statements(xml, "rdfxml").unwrap().len(), 3);
assert!(matches!(
parse_statements("<rdf:RDF><not closed", "rdfxml"),
Err(IngestError::RdfXml(_))
));
}
#[test]
fn rdfxml_namespace_fanout_is_bounded() {
let declarations = (0..2_000)
.map(|i| format!(r#" xmlns:p{i}="http://example.test/{i}/""#))
.collect::<String>();
let xml = format!(
r#"<rdf:RDF xmlns:rdf="http://www.w3.org/1999/02/22-rdf-syntax-ns#"{declarations}/>"#
);
let result = parse_statements(&xml, "rdfxml");
assert!(result.is_ok() || result.unwrap_err().to_string().contains("namespace"));
}
#[test]
fn rdfxml_attribute_fanout_completes() {
let attributes = (0..2_000)
.map(|i| format!(r#" p:a{i}="{i}""#))
.collect::<String>();
let xml = format!(
r#"<rdf:RDF xmlns:rdf="http://www.w3.org/1999/02/22-rdf-syntax-ns#" xmlns:p="http://example.test/"><rdf:Description rdf:about="http://example.test/s"{attributes}/></rdf:RDF>"#
);
assert!(parse_statements(&xml, "rdfxml").is_ok());
}
#[test]
fn assemble_dataset_round_trips() {
let text = "<http://ex/a> <http://ex/knows> <http://ex/b> .\n\
<http://ex/b> <http://ex/knows> <http://ex/c> .\n\
<http://ex/a> <http://ex/age> \"30\"^^<http://www.w3.org/2001/XMLSchema#integer> .\n\
<http://ex/a> <http://ex/p> <http://ex/d> <http://ex/g1> .";
let quads = parse_statements(text, "nq").unwrap();
let (bytes, stats) = assemble_dataset(quads, &[]);
assert_eq!(stats.statements, 4);
assert_eq!(stats.default_triples, 3);
assert_eq!(stats.named_graphs, 1);
assert!(stats.terms >= 7);
let rete = Rete::open(&bytes).unwrap();
assert_eq!(
rete.query(None, Some("<http://ex/knows>"), None).len(),
2,
"default-graph pattern query"
);
assert_eq!(rete.graph_names(), vec!["<http://ex/g1>"]);
let out = crate::eval_query(
&rete,
"SELECT ?x WHERE { ?x <http://ex/knows> ?y . ?y <http://ex/knows> ?z }",
)
.unwrap();
match out {
crate::QueryOutput::Select(_, rows) => assert_eq!(rows.len(), 1),
other => panic!("expected select result, got {other:?}"),
}
}
#[test]
fn assemble_minimal_typeless_graph() {
let text = "<http://ex/A> <http://ex/knows> <http://ex/B> .\n\
<http://ex/B> <http://ex/knows> <http://ex/C> .\n";
let quads = parse_statements(text, "nt").unwrap();
let (bytes, stats) = assemble_dataset(quads, &[]);
assert_eq!(stats.default_triples, 2);
let rete = Rete::open(&bytes).unwrap();
assert_eq!(rete.query(None, Some("<http://ex/knows>"), None).len(), 2);
}
}