#![cfg(not(target_arch = "wasm32"))]
use std::collections::BinaryHeap;
use std::fs::File;
use std::io::{BufRead, BufReader, BufWriter, Read, Seek, SeekFrom, Write};
use std::path::{Path, PathBuf};
use crate::dict::env_restart_interval;
use crate::index::{GroupSizer, INDEX_TILE_BUDGET};
use crate::ingest::{BuildStats, IngestError, RawQuad};
use crate::triples::TripleBlockBuilder;
use crate::varint::write_uvarint;
use crate::DictionaryBuilder;
type IdTriple = (u32, u32, u32);
type Synopsis = (u32, u32, u32, u32);
pub type MetadataFn = Box<dyn FnOnce(&BuildStats) -> Vec<u8>>;
type TermCarriers = (String, Vec<(usize, u32)>);
type MergeEntry = std::cmp::Reverse<(IdTriple, usize)>;
type TileDirEntry = (u32, u32, u64, Synopsis);
type EncodedTile = (u32, u32, Vec<u8>, Synopsis);
pub struct ExternalBuildOptions {
pub memory_budget: u64,
pub tmp_dir: Option<PathBuf>,
pub metadata: MetadataFn,
}
impl Default for ExternalBuildOptions {
fn default() -> Self {
Self {
memory_budget: 4 << 30,
tmp_dir: None,
metadata: Box::new(|_| Vec::new()),
}
}
}
#[derive(Debug, thiserror::Error)]
pub enum ExtBuildError {
#[error("ingest: {0}")]
Ingest(#[from] IngestError),
#[error("io: {0}")]
Io(#[from] std::io::Error),
#[error(
"external build supports the default graph only (named graph {0} found); \
use the standard build, or strip graph terms (.nq -> .nt) first"
)]
NamedGraph(String),
#[error("internal: {0}")]
Internal(&'static str),
}
const CHUNK_BUDGET_FRACTION: f64 = 0.5;
const PER_QUAD_OVERHEAD: u64 = 96;
const TILE_COMPRESS_BATCH: usize = 512;
pub fn build_external<S>(
mut stream: S,
output: &Path,
opts: ExternalBuildOptions,
) -> Result<BuildStats, ExtBuildError>
where
S: FnMut(&mut dyn FnMut(RawQuad) -> Result<(), ExtBuildError>) -> Result<(), ExtBuildError>,
{
let tmp_parent = opts
.tmp_dir
.clone()
.or_else(|| output.parent().map(|p| p.to_path_buf()))
.unwrap_or_else(|| PathBuf::from("."));
let tmp = TmpDir::create(&tmp_parent)?;
let budget = opts.memory_budget.max(64 << 20);
let chunk_budget = (budget as f64 * CHUNK_BUDGET_FRACTION) as u64;
eprintln!(
"extbuild: budget {} MiB -> chunk target {} MiB",
budget >> 20,
chunk_budget >> 20
);
let mut chunker = Chunker::new(&tmp, chunk_budget);
stream(&mut |q: RawQuad| chunker.push(q))?;
let chunks = chunker.finish()?;
let statements: u64 = chunks.iter().map(|c| c.triple_count).sum();
eprintln!(
"extbuild: {} chunk(s), {} statement(s) spilled",
chunks.len(),
statements
);
let mut merged = merge_dictionaries(&tmp, &chunks)?;
eprintln!(
"extbuild: merged dictionary — {} term(s)",
merged.term_count
);
let remaps = std::mem::take(&mut merged.remaps);
let global_tri = tmp.path("global.tri");
{
let mut out = BufWriter::new(File::create(&global_tri)?);
for (ci, _chunk) in chunks.iter().enumerate() {
let maps = &remaps[ci];
let mut rd = BufReader::new(File::open(tmp.path(&format!("c{ci}.tri")))?);
let mut buf = [0u8; 12];
loop {
match rd.read_exact(&mut buf) {
Ok(()) => {}
Err(e) if e.kind() == std::io::ErrorKind::UnexpectedEof => break,
Err(e) => return Err(e.into()),
}
let s = u32::from_le_bytes(buf[0..4].try_into().unwrap());
let p = u32::from_le_bytes(buf[4..8].try_into().unwrap());
let o = u32::from_le_bytes(buf[8..12].try_into().unwrap());
let gs = maps.subj[(s - 1) as usize];
let gp = maps.pred[(p - 1) as usize];
let go = maps.obj[(o - 1) as usize];
out.write_all(&gs.to_le_bytes())?;
out.write_all(&gp.to_le_bytes())?;
out.write_all(&go.to_le_bytes())?;
}
let _ = std::fs::remove_file(tmp.path(&format!("c{ci}.tri")));
}
out.flush()?;
}
drop(remaps);
let codec = crate::file::writer_codec();
let run_len = ((budget / 2) / 24).clamp(1 << 16, u32::MAX as u64) as usize;
let mut perm_sections: Vec<SectionFile> = Vec::with_capacity(6);
let mut deduped_count: Option<u64> = None;
for perm in crate::index::ALL_PERMS {
let (section, n) = build_permutation_section(&tmp, &global_tri, perm, run_len, codec)?;
if let Some(prev) = deduped_count {
if prev != n {
return Err(ExtBuildError::Internal("permutation dedup counts diverge"));
}
}
deduped_count = Some(n);
eprintln!(
"extbuild: permutation {} indexed ({} unique triple(s))",
perm.name(),
n
);
perm_sections.push(section);
}
let _ = std::fs::remove_file(&global_tri);
let quad_count = deduped_count.unwrap_or(0);
let mut stats = BuildStats {
statements: statements as usize,
default_triples: statements as usize,
named_graphs: 0,
terms: merged.term_count as usize,
pyramid_levels: 0,
};
let metadata = (opts.metadata)(&stats);
write_final_file(
output,
&metadata,
&merged,
&perm_sections,
quad_count,
codec,
)?;
stats.pyramid_levels = 0;
Ok(stats)
}
struct TmpDir {
dir: PathBuf,
}
impl TmpDir {
fn create(parent: &Path) -> Result<Self, std::io::Error> {
static SEQ: std::sync::atomic::AtomicU64 = std::sync::atomic::AtomicU64::new(0);
let seq = SEQ.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
let dir = parent.join(format!(".rete-extbuild-{}-{}", std::process::id(), seq));
std::fs::create_dir_all(&dir)?;
Ok(TmpDir { dir })
}
fn path(&self, name: &str) -> PathBuf {
self.dir.join(name)
}
}
impl Drop for TmpDir {
fn drop(&mut self) {
let _ = std::fs::remove_dir_all(&self.dir);
}
}
struct ChunkInfo {
triple_count: u64,
section_terms: [u32; 4],
}
struct Chunker<'a> {
tmp: &'a TmpDir,
chunk_budget: u64,
acc_bytes: u64,
quads: Vec<(String, String, String)>,
chunks: Vec<ChunkInfo>,
has_quoted: bool,
}
impl<'a> Chunker<'a> {
fn new(tmp: &'a TmpDir, chunk_budget: u64) -> Self {
Chunker {
tmp,
chunk_budget,
acc_bytes: 0,
quads: Vec::new(),
chunks: Vec::new(),
has_quoted: false,
}
}
fn push(&mut self, q: RawQuad) -> Result<(), ExtBuildError> {
let (s, p, o, g) = q;
if let Some(graph) = g {
return Err(ExtBuildError::NamedGraph(graph));
}
self.acc_bytes += (s.len() + p.len() + o.len()) as u64 + PER_QUAD_OVERHEAD;
self.quads.push((s, p, o));
if self.acc_bytes >= self.chunk_budget {
self.seal()?;
}
Ok(())
}
fn seal(&mut self) -> Result<(), ExtBuildError> {
if self.quads.is_empty() {
return Ok(());
}
let ci = self.chunks.len();
let quads = std::mem::take(&mut self.quads);
self.acc_bytes = 0;
let mut db = DictionaryBuilder::new();
for (s, p, o) in &quads {
db.observe(s, p, o);
}
let dict = db.build();
if dict.has_quoted_triples() {
self.has_quoted = true;
}
let mut tri = BufWriter::new(File::create(self.tmp.path(&format!("c{ci}.tri")))?);
for (s, p, o) in &quads {
let (si, pi, oi) = dict
.encode(s, p, o)
.ok_or(ExtBuildError::Internal("chunk term missing from own dict"))?;
tri.write_all(&si.to_le_bytes())?;
tri.write_all(&pi.to_le_bytes())?;
tri.write_all(&oi.to_le_bytes())?;
}
tri.flush()?;
let triple_count = quads.len() as u64;
drop(quads);
let shared = dict.shared_count();
let subj_only = dict.subject_only_count();
let obj_only = dict.object_only_count();
let preds = {
let mut n = 1u32;
while dict.predicate_term(n).is_some() {
n += 1;
}
n - 1
};
let write_terms = |name: &str,
count: u32,
term_of: &dyn Fn(u32) -> Option<String>|
-> Result<(), ExtBuildError> {
let mut w = BufWriter::new(File::create(self.tmp.path(name))?);
for i in 1..=count {
let t = term_of(i).ok_or(ExtBuildError::Internal("dict id out of range"))?;
w.write_all(t.as_bytes())?;
w.write_all(b"\n")?;
}
w.flush()?;
Ok(())
};
write_terms(&format!("c{ci}.shared"), shared, &|i| dict.subject_term(i))?;
write_terms(&format!("c{ci}.subj"), subj_only, &|i| {
dict.subject_term(shared + i)
})?;
write_terms(&format!("c{ci}.obj"), obj_only, &|i| {
dict.object_term(shared + i)
})?;
write_terms(&format!("c{ci}.pred"), preds, &|i| dict.predicate_term(i))?;
eprintln!(
"extbuild: chunk {ci} sealed — {triple_count} statement(s), {} term(s)",
shared + subj_only + obj_only + preds
);
self.chunks.push(ChunkInfo {
triple_count,
section_terms: [shared, subj_only, obj_only, preds],
});
Ok(())
}
fn finish(mut self) -> Result<Vec<ChunkInfo>, ExtBuildError> {
self.seal()?;
if self.chunks.is_empty() {
self.chunks.push(ChunkInfo {
triple_count: 0,
section_terms: [0, 0, 0, 0],
});
for name in ["c0.tri", "c0.shared", "c0.subj", "c0.obj", "c0.pred"] {
File::create(self.tmp.path(name))?;
}
}
Ok(self.chunks)
}
}
struct ChunkRemap {
subj: Vec<u32>,
obj: Vec<u32>,
pred: Vec<u32>,
}
struct MergedDict {
section_files: [SectionFile; 4],
term_count: u64,
has_quoted: bool,
remaps: Vec<ChunkRemap>,
}
struct SectionFile {
path: PathBuf,
len: u64,
}
struct SpaceStream {
a: TermFileReader, b: TermFileReader, a_base: u32, b_base: u32, a_next: Option<String>,
b_next: Option<String>,
a_rank: u32,
b_rank: u32,
}
impl SpaceStream {
fn new(a: TermFileReader, b: TermFileReader, shared: u32) -> Result<Self, std::io::Error> {
let mut s = SpaceStream {
a,
b,
a_base: 0,
b_base: shared,
a_next: None,
b_next: None,
a_rank: 0,
b_rank: 0,
};
s.a_next = s.a.next()?;
s.b_next = s.b.next()?;
Ok(s)
}
fn next(&mut self) -> Result<Option<(String, u32)>, std::io::Error> {
let take_a = match (&self.a_next, &self.b_next) {
(None, None) => return Ok(None),
(Some(_), None) => true,
(None, Some(_)) => false,
(Some(x), Some(y)) => x < y,
};
if take_a {
let t = self.a_next.take().unwrap();
self.a_rank += 1;
let id = self.a_base + self.a_rank;
self.a_next = self.a.next()?;
Ok(Some((t, id)))
} else {
let t = self.b_next.take().unwrap();
self.b_rank += 1;
let id = self.b_base + self.b_rank;
self.b_next = self.b.next()?;
Ok(Some((t, id)))
}
}
}
struct TermFileReader {
rd: BufReader<File>,
buf: String,
}
impl TermFileReader {
fn open(path: &Path) -> Result<Self, std::io::Error> {
Ok(TermFileReader {
rd: BufReader::with_capacity(1 << 20, File::open(path)?),
buf: String::new(),
})
}
fn next(&mut self) -> Result<Option<String>, std::io::Error> {
self.buf.clear();
let n = self.rd.read_line(&mut self.buf)?;
if n == 0 {
return Ok(None);
}
if self.buf.ends_with('\n') {
self.buf.pop();
}
Ok(Some(self.buf.clone()))
}
}
struct HeapEntry {
term: String,
chunk: usize,
local_id: u32,
}
impl PartialEq for HeapEntry {
fn eq(&self, other: &Self) -> bool {
self.term == other.term && self.chunk == other.chunk
}
}
impl Eq for HeapEntry {}
impl PartialOrd for HeapEntry {
fn partial_cmp(&self, other: &Self) -> Option<std::cmp::Ordering> {
Some(self.cmp(other))
}
}
impl Ord for HeapEntry {
fn cmp(&self, other: &Self) -> std::cmp::Ordering {
other
.term
.cmp(&self.term)
.then_with(|| other.chunk.cmp(&self.chunk))
}
}
struct KWayTerms {
heap: BinaryHeap<HeapEntry>,
streams: Vec<SpaceStream>,
}
impl KWayTerms {
fn new(mut streams: Vec<SpaceStream>) -> Result<Self, std::io::Error> {
let mut heap = BinaryHeap::new();
for (ci, s) in streams.iter_mut().enumerate() {
if let Some((term, id)) = s.next()? {
heap.push(HeapEntry {
term,
chunk: ci,
local_id: id,
});
}
}
Ok(KWayTerms { heap, streams })
}
fn next(&mut self) -> Result<Option<TermCarriers>, std::io::Error> {
let first = match self.heap.pop() {
Some(e) => e,
None => return Ok(None),
};
let term = first.term;
let mut carriers = vec![(first.chunk, first.local_id)];
if let Some((t, id)) = self.streams[first.chunk].next()? {
self.heap.push(HeapEntry {
term: t,
chunk: first.chunk,
local_id: id,
});
}
while let Some(top) = self.heap.peek() {
if top.term != term {
break;
}
let e = self.heap.pop().unwrap();
carriers.push((e.chunk, e.local_id));
if let Some((t, id)) = self.streams[e.chunk].next()? {
self.heap.push(HeapEntry {
term: t,
chunk: e.chunk,
local_id: id,
});
}
}
Ok(Some((term, carriers)))
}
}
const CLASS_SHARED: u32 = 0b00 << 30;
const CLASS_SUBJ_ONLY: u32 = 0b01 << 30;
const CLASS_OBJ_ONLY: u32 = 0b10 << 30;
const CLASS_MASK: u32 = 0b11 << 30;
fn merge_dictionaries(tmp: &TmpDir, chunks: &[ChunkInfo]) -> Result<MergedDict, ExtBuildError> {
let k = chunks.len();
let mut remaps: Vec<ChunkRemap> = chunks
.iter()
.map(|c| ChunkRemap {
subj: vec![0; (c.section_terms[0] + c.section_terms[1]) as usize],
obj: vec![0; (c.section_terms[0] + c.section_terms[2]) as usize],
pred: vec![0; c.section_terms[3] as usize],
})
.collect();
let mut subj_streams = Vec::with_capacity(k);
let mut obj_streams = Vec::with_capacity(k);
for (ci, c) in chunks.iter().enumerate() {
subj_streams.push(SpaceStream::new(
TermFileReader::open(&tmp.path(&format!("c{ci}.shared")))?,
TermFileReader::open(&tmp.path(&format!("c{ci}.subj")))?,
c.section_terms[0],
)?);
obj_streams.push(SpaceStream::new(
TermFileReader::open(&tmp.path(&format!("c{ci}.shared")))?,
TermFileReader::open(&tmp.path(&format!("c{ci}.obj")))?,
c.section_terms[0],
)?);
}
let mut subjects = KWayTerms::new(subj_streams)?;
let mut objects = KWayTerms::new(obj_streams)?;
let mut shared_sec = RawSectionWriter::create(tmp.path("g.shared.raw"))?;
let mut subj_sec = RawSectionWriter::create(tmp.path("g.subj.raw"))?;
let mut obj_sec = RawSectionWriter::create(tmp.path("g.obj.raw"))?;
let mut has_quoted = false;
let mut ranks = [0u32; 3]; let mut s_item = subjects.next()?;
let mut o_item = objects.next()?;
loop {
enum Class {
Shared,
SubjOnly,
ObjOnly,
}
let class = match (&s_item, &o_item) {
(None, None) => break,
(Some(_), None) => Class::SubjOnly,
(None, Some(_)) => Class::ObjOnly,
(Some((st, _)), Some((ot, _))) => match st.cmp(ot) {
std::cmp::Ordering::Less => Class::SubjOnly,
std::cmp::Ordering::Greater => Class::ObjOnly,
std::cmp::Ordering::Equal => Class::Shared,
},
};
match class {
Class::Shared => {
let (term, s_carriers) = s_item.take().unwrap();
let (_, o_carriers) = o_item.take().unwrap();
let enc = CLASS_SHARED | ranks[0];
for (ci, lid) in s_carriers {
remaps[ci].subj[(lid - 1) as usize] = enc;
}
for (ci, lid) in o_carriers {
remaps[ci].obj[(lid - 1) as usize] = enc;
}
has_quoted |= term.starts_with("<<");
shared_sec.push(&term)?;
ranks[0] += 1;
s_item = subjects.next()?;
o_item = objects.next()?;
}
Class::SubjOnly => {
let (term, carriers) = s_item.take().unwrap();
let enc = CLASS_SUBJ_ONLY | ranks[1];
for (ci, lid) in carriers {
remaps[ci].subj[(lid - 1) as usize] = enc;
}
has_quoted |= term.starts_with("<<");
subj_sec.push(&term)?;
ranks[1] += 1;
s_item = subjects.next()?;
}
Class::ObjOnly => {
let (term, carriers) = o_item.take().unwrap();
let enc = CLASS_OBJ_ONLY | ranks[2];
for (ci, lid) in carriers {
remaps[ci].obj[(lid - 1) as usize] = enc;
}
has_quoted |= term.starts_with("<<");
obj_sec.push(&term)?;
ranks[2] += 1;
o_item = objects.next()?;
}
}
}
let (n_shared, n_subj, n_obj) = (ranks[0], ranks[1], ranks[2]);
let mut pred_sec = RawSectionWriter::create(tmp.path("g.pred.raw"))?;
{
let mut streams = Vec::with_capacity(k);
for (ci, _) in chunks.iter().enumerate() {
streams.push(SpaceStream::new(
TermFileReader::open(&tmp.path(&format!("c{ci}.pred")))?,
TermFileReader::open(&tmp.path(&format!("c{ci}.pred.empty",))).or_else(
|_| -> Result<TermFileReader, std::io::Error> {
let p = tmp.path(&format!("c{ci}.pred.empty"));
File::create(&p)?;
TermFileReader::open(&p)
},
)?,
chunks[ci].section_terms[3],
)?);
}
let mut kway = KWayTerms::new(streams)?;
let mut rank = 0u32;
while let Some((term, carriers)) = kway.next()? {
for (ci, lid) in carriers {
remaps[ci].pred[(lid - 1) as usize] = rank + 1; }
pred_sec.push(&term)?;
rank += 1;
}
}
for rm in &mut remaps {
for v in rm.subj.iter_mut() {
let rank = *v & !CLASS_MASK;
*v = match *v & CLASS_MASK {
CLASS_SHARED => rank + 1,
CLASS_SUBJ_ONLY => n_shared + rank + 1,
_ => return Err(ExtBuildError::Internal("subject remap class corrupt")),
};
}
for v in rm.obj.iter_mut() {
let rank = *v & !CLASS_MASK;
*v = match *v & CLASS_MASK {
CLASS_SHARED => rank + 1,
CLASS_OBJ_ONLY => n_shared + rank + 1,
_ => return Err(ExtBuildError::Internal("object remap class corrupt")),
};
}
}
let codec = crate::file::writer_codec();
let section_files = [
shared_sec.finish_chunked(tmp, "g.shared.sec", codec)?,
subj_sec.finish_chunked(tmp, "g.subj.sec", codec)?,
obj_sec.finish_chunked(tmp, "g.obj.sec", codec)?,
pred_sec.finish_chunked(tmp, "g.pred.sec", codec)?,
];
for ci in 0..k {
for suffix in ["shared", "subj", "obj", "pred", "pred.empty"] {
let _ = std::fs::remove_file(tmp.path(&format!("c{ci}.{suffix}")));
}
}
let term_count =
n_shared as u64 + n_subj as u64 + n_obj as u64 + section_files[3].term_count as u64;
Ok(MergedDict {
section_files: section_files.map(|s| s.file),
term_count,
has_quoted,
remaps,
})
}
struct RawSectionWriter {
w: BufWriter<File>,
path: PathBuf,
prev: String,
n: u64,
body_len: u64,
restart_offsets: Vec<u64>,
}
struct FinishedSection {
file: SectionFile,
term_count: u32,
}
impl RawSectionWriter {
fn create(path: PathBuf) -> Result<Self, std::io::Error> {
Ok(RawSectionWriter {
w: BufWriter::with_capacity(1 << 20, File::create(&path)?),
path,
prev: String::new(),
n: 0,
body_len: 0,
restart_offsets: Vec::new(),
})
}
fn push(&mut self, term: &str) -> Result<(), std::io::Error> {
let r = env_restart_interval() as u64;
let mut entry = Vec::with_capacity(term.len() + 8);
if self.n.is_multiple_of(r) {
self.restart_offsets.push(self.body_len);
write_uvarint(&mut entry, 0);
write_uvarint(&mut entry, term.len() as u64);
entry.extend_from_slice(term.as_bytes());
} else {
let shared = common_prefix_len(&self.prev, term);
let suffix = &term.as_bytes()[shared..];
write_uvarint(&mut entry, shared as u64);
write_uvarint(&mut entry, suffix.len() as u64);
entry.extend_from_slice(suffix);
}
self.w.write_all(&entry)?;
self.body_len += entry.len() as u64;
self.prev.clear();
self.prev.push_str(term);
self.n += 1;
Ok(())
}
fn finish_chunked(
mut self,
tmp: &TmpDir,
out_name: &str,
codec: u8,
) -> Result<FinishedSection, ExtBuildError> {
self.w.flush()?;
drop(self.w);
let term_count = self.n as u32;
let mut header = Vec::new();
write_uvarint(&mut header, self.n);
write_uvarint(&mut header, env_restart_interval() as u64);
write_uvarint(&mut header, self.restart_offsets.len() as u64);
for off in &self.restart_offsets {
write_uvarint(&mut header, *off);
}
let budget = 64 * 1024u64; let offs = &self.restart_offsets;
let mut bounds: Vec<(usize, u64, u64)> = Vec::new(); let mut r = 0usize;
while r < offs.len() {
let start = offs[r];
let mut r2 = r + 1;
while r2 < offs.len() && offs[r2] - start < budget {
r2 += 1;
}
let end = if r2 < offs.len() {
offs[r2]
} else {
self.body_len
};
bounds.push((r, start, end));
r = r2;
}
let mut body = File::open(&self.path)?;
let comp_path = tmp.path(&format!("{out_name}.chunks"));
let mut comp_out = BufWriter::new(File::create(&comp_path)?);
let mut dir: Vec<(usize, Vec<u8>, u64)> = Vec::with_capacity(bounds.len());
for &(first_run, start, end) in &bounds {
body.seek(SeekFrom::Start(start))?;
let mut raw = vec![0u8; (end - start) as usize];
body.read_exact(&mut raw)?;
let first_term = read_restart_term(&raw)
.ok_or(ExtBuildError::Internal("restart entry unreadable"))?;
let comp = crate::file::compress(codec, &raw);
comp_out.write_all(&comp)?;
dir.push((first_run, first_term, comp.len() as u64));
}
comp_out.flush()?;
drop(body);
let _ = std::fs::remove_file(&self.path);
let out_path = tmp.path(out_name);
let mut out = BufWriter::new(File::create(&out_path)?);
let mut head = Vec::new();
write_uvarint(&mut head, header.len() as u64);
head.extend_from_slice(&header);
write_uvarint(&mut head, dir.len() as u64);
let mut prev_run = 0usize;
for (first_run, first_term, comp_len) in &dir {
write_uvarint(&mut head, (*first_run - prev_run) as u64);
write_uvarint(&mut head, first_term.len() as u64);
head.extend_from_slice(first_term);
write_uvarint(&mut head, *comp_len);
prev_run = *first_run;
}
out.write_all(&head)?;
let mut comp_in = File::open(&comp_path)?;
let copied = std::io::copy(&mut comp_in, &mut out)?;
out.flush()?;
let _ = std::fs::remove_file(&comp_path);
let len = head.len() as u64 + copied;
Ok(FinishedSection {
file: SectionFile {
path: out_path,
len,
},
term_count,
})
}
}
fn common_prefix_len(a: &str, b: &str) -> usize {
a.as_bytes()
.iter()
.zip(b.as_bytes())
.take_while(|(x, y)| x == y)
.count()
}
fn read_restart_term(raw: &[u8]) -> Option<Vec<u8>> {
let (shared, n1) = crate::varint::read_uvarint(raw)?;
if shared != 0 {
return None;
}
let (len, n2) = crate::varint::read_uvarint(&raw[n1..])?;
let start = n1 + n2;
raw.get(start..start + len as usize).map(|s| s.to_vec())
}
fn build_permutation_section(
tmp: &TmpDir,
global_tri: &Path,
perm: crate::index::IndexPermutation,
run_len: usize,
codec: u8,
) -> Result<(SectionFile, u64), ExtBuildError> {
let mut runs: Vec<PathBuf> = Vec::new();
{
let mut rd = BufReader::with_capacity(1 << 20, File::open(global_tri)?);
let mut buf = [0u8; 12];
let mut run: Vec<(u32, u32, u32)> = Vec::with_capacity(run_len.min(1 << 22));
loop {
let eof = match rd.read_exact(&mut buf) {
Ok(()) => false,
Err(e) if e.kind() == std::io::ErrorKind::UnexpectedEof => true,
Err(e) => return Err(e.into()),
};
if !eof {
let s = u32::from_le_bytes(buf[0..4].try_into().unwrap());
let p = u32::from_le_bytes(buf[4..8].try_into().unwrap());
let o = u32::from_le_bytes(buf[8..12].try_into().unwrap());
run.push(perm.forward((s, p, o)));
}
if run.len() >= run_len || (eof && !run.is_empty()) {
sort_triples(&mut run);
run.dedup();
let path = tmp.path(&format!("{}.run{}", perm.name(), runs.len()));
let mut w = BufWriter::with_capacity(1 << 20, File::create(&path)?);
for &(a, b, c) in &run {
w.write_all(&a.to_le_bytes())?;
w.write_all(&b.to_le_bytes())?;
w.write_all(&c.to_le_bytes())?;
}
w.flush()?;
runs.push(path);
run.clear();
}
if eof {
break;
}
}
}
let mut tiler = StreamingTiler::new(tmp, perm.name(), codec)?;
{
let mut readers: Vec<RunReader> = runs
.iter()
.map(|p| RunReader::open(p))
.collect::<Result<_, _>>()?;
let mut heap: BinaryHeap<MergeEntry> = BinaryHeap::new();
for (i, r) in readers.iter_mut().enumerate() {
if let Some(t) = r.next()? {
heap.push(std::cmp::Reverse((t, i)));
}
}
let mut last: Option<(u32, u32, u32)> = None;
while let Some(std::cmp::Reverse((t, i))) = heap.pop() {
if let Some(n) = readers[i].next()? {
heap.push(std::cmp::Reverse((n, i)));
}
if last != Some(t) {
tiler.push(t)?;
last = Some(t);
}
}
}
let (section, count) = tiler.finish(tmp)?;
for p in runs {
let _ = std::fs::remove_file(p);
}
Ok((section, count))
}
#[cfg(feature = "parallel")]
fn sort_triples(v: &mut [IdTriple]) {
use rayon::slice::ParallelSliceMut;
v.par_sort_unstable();
}
#[cfg(not(feature = "parallel"))]
fn sort_triples(v: &mut [IdTriple]) {
v.sort_unstable();
}
struct RunReader {
rd: BufReader<File>,
}
impl RunReader {
fn open(path: &Path) -> Result<Self, std::io::Error> {
Ok(RunReader {
rd: BufReader::with_capacity(1 << 20, File::open(path)?),
})
}
fn next(&mut self) -> Result<Option<(u32, u32, u32)>, std::io::Error> {
let mut buf = [0u8; 12];
match self.rd.read_exact(&mut buf) {
Ok(()) => Ok(Some((
u32::from_le_bytes(buf[0..4].try_into().unwrap()),
u32::from_le_bytes(buf[4..8].try_into().unwrap()),
u32::from_le_bytes(buf[8..12].try_into().unwrap()),
))),
Err(e) if e.kind() == std::io::ErrorKind::UnexpectedEof => Ok(None),
Err(e) => Err(e),
}
}
}
struct StreamingTiler {
tile_budget: usize,
codec: u8,
tile: Vec<(u32, u32, u32)>,
tile_size: usize,
group: Vec<(u32, u32, u32)>,
sizer: GroupSizer,
gtotal: usize,
prev_a: u32,
pending: Vec<Vec<(u32, u32, u32)>>,
comp_out: BufWriter<File>,
comp_path: PathBuf,
dir: Vec<TileDirEntry>,
count: u64,
}
impl StreamingTiler {
fn new(tmp: &TmpDir, name: &str, codec: u8) -> Result<Self, ExtBuildError> {
let comp_path = tmp.path(&format!("{name}.tiles"));
Ok(StreamingTiler {
tile_budget: INDEX_TILE_BUDGET,
codec,
tile: Vec::new(),
tile_size: 0,
group: Vec::new(),
sizer: GroupSizer::start(0, 0),
gtotal: 0,
prev_a: 0,
pending: Vec::new(),
comp_out: BufWriter::with_capacity(1 << 20, File::create(&comp_path)?),
comp_path,
dir: Vec::new(),
count: 0,
})
}
fn push(&mut self, t: (u32, u32, u32)) -> Result<(), ExtBuildError> {
self.count += 1;
if let Some(&(ga, _, _)) = self.group.first() {
if t.0 != ga {
self.close_group()?;
}
}
if self.group.is_empty() {
self.sizer = GroupSizer::start(t.0, self.prev_a);
}
self.group.push(t);
self.gtotal = self.sizer.push(t.1, t.2);
if self.gtotal > self.tile_budget {
self.tile.append(&mut self.group);
self.flush_tile()?;
self.tile_size = 0;
self.prev_a = t.0;
self.gtotal = 0;
}
Ok(())
}
fn close_group(&mut self) -> Result<(), ExtBuildError> {
if self.group.is_empty() {
return Ok(());
}
let a = self.group[0].0;
let gsize = self.sizer.total();
if !self.tile.is_empty() && self.tile_size + gsize > self.tile_budget {
self.flush_tile()?;
}
self.tile_size += gsize;
self.prev_a = a;
self.tile.append(&mut self.group);
Ok(())
}
fn flush_tile(&mut self) -> Result<(), ExtBuildError> {
if self.tile.is_empty() {
return Ok(());
}
self.pending.push(std::mem::take(&mut self.tile));
self.tile_size = 0;
if self.pending.len() >= TILE_COMPRESS_BATCH {
self.compress_pending()?;
}
Ok(())
}
fn compress_pending(&mut self) -> Result<(), ExtBuildError> {
if self.pending.is_empty() {
return Ok(());
}
let batch = std::mem::take(&mut self.pending);
let codec = self.codec;
let encode_one = |run: &Vec<(u32, u32, u32)>| -> (u32, u32, Vec<u8>, (u32, u32, u32, u32)) {
let mut b = TripleBlockBuilder::new();
for &t in run {
b.push(t);
}
let bytes = b.build();
let syn = match crate::triples::TripleBlock::parse(&bytes) {
Ok(blk) => {
let z = blk.zone();
(z.min_b, z.max_b, z.min_c, z.max_c)
}
Err(_) => (0, u32::MAX, 0, u32::MAX),
};
let comp = crate::file::compress(codec, &bytes);
(run[0].0, run[run.len() - 1].0, comp, syn)
};
#[cfg(feature = "parallel")]
let encoded: Vec<EncodedTile> = {
use rayon::prelude::*;
batch.par_iter().map(encode_one).collect()
};
#[cfg(not(feature = "parallel"))]
let encoded: Vec<EncodedTile> = batch.iter().map(encode_one).collect();
for (min_a, max_a, comp, syn) in encoded {
self.comp_out.write_all(&comp)?;
self.dir.push((min_a, max_a, comp.len() as u64, syn));
}
Ok(())
}
fn finish(mut self, tmp: &TmpDir) -> Result<(SectionFile, u64), ExtBuildError> {
self.close_group()?;
self.flush_tile()?;
self.compress_pending()?;
self.comp_out.flush()?;
drop(self.comp_out);
let name = self
.comp_path
.file_name()
.and_then(|n| n.to_str())
.unwrap_or("perm")
.to_string();
let out_path = tmp.path(&format!("{name}.sec"));
let mut out = BufWriter::with_capacity(1 << 20, File::create(&out_path)?);
let mut head = Vec::new();
write_uvarint(&mut head, self.dir.len() as u64);
let mut prev_min = 0u32;
for &(min_a, max_a, comp_len, _) in &self.dir {
write_uvarint(&mut head, (min_a - prev_min) as u64);
write_uvarint(&mut head, (max_a - min_a) as u64);
write_uvarint(&mut head, comp_len);
prev_min = min_a;
}
out.write_all(&head)?;
let mut tiles_in = File::open(&self.comp_path)?;
let copied = std::io::copy(&mut tiles_in, &mut out)?;
let mut trailer = Vec::new();
for &(_, _, _, (min_b, max_b, min_c, max_c)) in &self.dir {
write_uvarint(&mut trailer, min_b as u64);
write_uvarint(&mut trailer, (max_b - min_b) as u64);
write_uvarint(&mut trailer, min_c as u64);
write_uvarint(&mut trailer, (max_c - min_c) as u64);
}
out.write_all(&trailer)?;
out.flush()?;
let _ = std::fs::remove_file(&self.comp_path);
let len = head.len() as u64 + copied + trailer.len() as u64;
Ok((
SectionFile {
path: out_path,
len,
},
self.count,
))
}
}
fn write_final_file(
output: &Path,
metadata: &[u8],
dict: &MergedDict,
perm_sections: &[SectionFile],
quad_count: u64,
codec: u8,
) -> Result<(), ExtBuildError> {
use crate::header::{Header, FLAG_HAS_QUOTED_TRIPLES, FLAG_TILE_SYNOPSIS, HEADER_LEN};
let mut dict_frame = Vec::new();
write_uvarint(&mut dict_frame, 4);
let mut dict_len = dict_frame.len() as u64;
let mut dict_section_heads = Vec::with_capacity(4);
for s in &dict.section_files {
let mut h = Vec::new();
write_uvarint(&mut h, s.len);
dict_len += h.len() as u64 + s.len;
dict_section_heads.push(h);
}
let mut index_frame = Vec::new();
write_uvarint(&mut index_frame, perm_sections.len() as u64);
let mut index_len = index_frame.len() as u64;
let mut index_section_heads = Vec::with_capacity(perm_sections.len());
for s in perm_sections {
let mut h = Vec::new();
write_uvarint(&mut h, s.len);
index_len += h.len() as u64 + s.len;
index_section_heads.push(h);
}
let meta_len = metadata.len() as u64;
let dict_offset = HEADER_LEN as u64 + meta_len;
let index_offset = dict_offset + dict_len;
let mut out = BufWriter::with_capacity(1 << 20, File::create(output)?);
let mut hasher = blake3::Hasher::new();
out.write_all(&[0u8; HEADER_LEN])?;
if meta_len > 0 {
out.write_all(metadata)?;
hasher.update(metadata);
}
let write_hashed = |out: &mut BufWriter<File>,
hasher: &mut blake3::Hasher,
bytes: &[u8]|
-> Result<(), std::io::Error> {
out.write_all(bytes)?;
hasher.update(bytes);
Ok(())
};
write_hashed(&mut out, &mut hasher, &dict_frame)?;
for (s, head) in dict.section_files.iter().zip(&dict_section_heads) {
write_hashed(&mut out, &mut hasher, head)?;
copy_hashed(&s.path, &mut out, &mut hasher)?;
}
write_hashed(&mut out, &mut hasher, &index_frame)?;
for (s, head) in perm_sections.iter().zip(&index_section_heads) {
write_hashed(&mut out, &mut hasher, head)?;
copy_hashed(&s.path, &mut out, &mut hasher)?;
}
out.write_all(&crate::header::MAGIC)?; out.flush()?;
let mut hash = [0u8; 16];
hash.copy_from_slice(&hasher.finalize().as_bytes()[..16]);
let header = Header {
version: crate::header::CURRENT_FORMAT_VERSION,
flags: FLAG_TILE_SYNOPSIS
| if dict.has_quoted {
FLAG_HAS_QUOTED_TRIPLES
} else {
0
},
metadata_offset: HEADER_LEN as u64,
metadata_len: meta_len,
dictionary_offset: dict_offset,
dictionary_len: dict_len,
root_dir_offset: index_offset,
root_dir_len: index_len,
pyramid_meta_offset: 0,
pyramid_meta_len: 0,
dict_codec: codec,
block_codec: codec,
pyramid_levels: 0,
quad_count,
term_count: dict.term_count,
content_hash: hash,
named_graphs_offset: 0,
named_graphs_len: 0,
schema_meta_len: 0,
text_index_offset: 0,
text_index_len: 0,
extra_sections: Vec::new(),
};
let mut f = out.into_inner().map_err(|e| e.into_error())?;
f.seek(SeekFrom::Start(0))?;
f.write_all(&header.to_bytes())?;
f.flush()?;
Ok(())
}
fn copy_hashed(
path: &Path,
out: &mut BufWriter<File>,
hasher: &mut blake3::Hasher,
) -> Result<(), std::io::Error> {
let mut rd = File::open(path)?;
let mut buf = vec![0u8; 1 << 20];
loop {
let n = rd.read(&mut buf)?;
if n == 0 {
break;
}
out.write_all(&buf[..n])?;
hasher.update(&buf[..n]);
}
Ok(())
}
#[cfg(test)]
mod tests {
use super::*;
use crate::ingest;
fn test_quads(n: usize) -> Vec<RawQuad> {
let mut quads = Vec::new();
for i in 0..n {
let s = format!("<http://ex/s{}>", i % (n / 7 + 1));
let p = format!("<http://ex/p{}>", i % 13);
let o = match i % 5 {
0 => format!("<http://ex/s{}>", (i + 3) % (n / 7 + 1)), 1 => format!("\"lit {i}\""),
2 => format!("\"{}\"^^<http://www.w3.org/2001/XMLSchema#integer>", i),
3 => format!("\"v{}\"@en", i % 50),
_ => format!("_:b{}", i % 97),
};
quads.push((s, p, o, None));
}
for i in 0..(n / 20) {
let j = i * 17 % n;
quads.push(quads[j].clone());
}
quads
}
fn build_reference(quads: Vec<RawQuad>) -> Vec<u8> {
let qs = quads.clone();
let (bytes, _stats) = ingest::assemble_dataset_streaming_algo(
move |visit: &mut dyn FnMut(RawQuad)| {
for q in qs.iter().cloned() {
visit(q);
}
Ok(())
},
false,
false,
None,
crate::PyramidAlgo::Louvain,
|_, _, _| Vec::new(),
)
.unwrap();
bytes
}
fn build_ext(quads: Vec<RawQuad>, budget: u64) -> (Vec<u8>, BuildStats) {
static TEST_SEQ: std::sync::atomic::AtomicU64 = std::sync::atomic::AtomicU64::new(0);
let seq = TEST_SEQ.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
let dir = std::env::temp_dir().join(format!(
"rete-extbuild-test-{}-{budget}-{seq}",
std::process::id()
));
std::fs::create_dir_all(&dir).unwrap();
let out = dir.join("out.rete");
let stats = build_external(
|visit| {
for q in quads.iter().cloned() {
visit(q)?;
}
Ok(())
},
&out,
ExternalBuildOptions {
memory_budget: budget,
tmp_dir: Some(dir.clone()),
metadata: Box::new(|_| Vec::new()),
},
)
.unwrap();
let bytes = std::fs::read(&out).unwrap();
let _ = std::fs::remove_dir_all(&dir);
(bytes, stats)
}
#[test]
fn external_build_is_byte_identical_to_streaming() {
let quads = test_quads(3000);
let reference = build_reference(quads.clone());
let (bytes_floor, stats) = build_ext(quads.clone(), 0);
assert_eq!(stats.statements, quads.len());
assert_eq!(
bytes_floor, reference,
"single-chunk external build must be byte-identical"
);
}
#[test]
fn skewed_external_build_is_byte_identical() {
let mut quads: Vec<RawQuad> = Vec::new();
for i in 0..30_000usize {
quads.push((
format!("<http://ex/s{}>", i % 500),
"<http://ex/cites>".to_string(),
format!("<http://ex/o{i}>"),
None,
));
}
for i in 0..500usize {
quads.push((
format!("<http://ex/s{i}>"),
format!("<http://ex/p{}>", i % 7),
format!("\"lit {i}\""),
None,
));
}
let reference = build_reference(quads.clone());
let (bytes, stats) = build_ext(quads.clone(), 0);
assert_eq!(stats.statements, quads.len());
assert_eq!(
bytes, reference,
"mega-group external build must be byte-identical"
);
}
#[test]
fn multi_chunk_external_build_is_byte_identical() {
let quads = test_quads(3000);
let reference = build_reference(quads.clone());
let dir = std::env::temp_dir().join(format!("rete-extbuild-mc-{}", std::process::id()));
std::fs::create_dir_all(&dir).unwrap();
let out = dir.join("out.rete");
let tmp = TmpDir::create(&dir).unwrap();
let mut chunker = Chunker::new(&tmp, 4 * 1024); for q in quads.iter().cloned() {
chunker.push(q).unwrap();
}
let chunks = chunker.finish().unwrap();
assert!(
chunks.len() >= 4,
"expected many chunks, got {}",
chunks.len()
);
let statements: u64 = chunks.iter().map(|c| c.triple_count).sum();
let merged = merge_dictionaries(&tmp, &chunks).unwrap();
let global_tri = tmp.path("global.tri");
{
let mut w = BufWriter::new(File::create(&global_tri).unwrap());
for (ci, _c) in chunks.iter().enumerate() {
let maps = &merged.remaps[ci];
let mut rd = BufReader::new(File::open(tmp.path(&format!("c{ci}.tri"))).unwrap());
let mut buf = [0u8; 12];
while rd.read_exact(&mut buf).is_ok() {
let s = u32::from_le_bytes(buf[0..4].try_into().unwrap());
let p = u32::from_le_bytes(buf[4..8].try_into().unwrap());
let o = u32::from_le_bytes(buf[8..12].try_into().unwrap());
w.write_all(&maps.subj[(s - 1) as usize].to_le_bytes())
.unwrap();
w.write_all(&maps.pred[(p - 1) as usize].to_le_bytes())
.unwrap();
w.write_all(&maps.obj[(o - 1) as usize].to_le_bytes())
.unwrap();
}
}
w.flush().unwrap();
}
let codec = crate::file::writer_codec();
let mut sections = Vec::new();
let mut count = None;
for perm in crate::index::ALL_PERMS {
let (sec, n) = build_permutation_section(&tmp, &global_tri, perm, 256, codec).unwrap();
if let Some(prev) = count {
assert_eq!(prev, n);
}
count = Some(n);
sections.push(sec);
}
write_final_file(&out, &[], &merged, §ions, count.unwrap(), codec).unwrap();
let bytes = std::fs::read(&out).unwrap();
assert_eq!(statements as usize, quads.len());
assert_eq!(
bytes, reference,
"multi-chunk external build must be byte-identical"
);
drop(tmp);
let _ = std::fs::remove_dir_all(&dir);
}
#[test]
fn external_build_output_is_queryable() {
let quads = test_quads(1200);
let (bytes, _) = build_ext(quads.clone(), 0);
let rete = crate::Rete::open(&bytes).unwrap();
let out = crate::eval_query(
&rete,
"SELECT (COUNT(*) AS ?n) WHERE { ?s <http://ex/p1> ?o }",
)
.unwrap();
match out {
crate::QueryOutput::Select(_, rows) => assert_eq!(rows.len(), 1),
other => panic!("expected select, got {other:?}"),
}
}
#[test]
#[ignore = "operational tool, driven by RETE_RESUME_* env vars"]
fn resume_from_spill() {
let spill = PathBuf::from(std::env::var("RETE_RESUME_SPILL").expect("RETE_RESUME_SPILL"));
let out = PathBuf::from(std::env::var("RETE_RESUME_OUT").expect("RETE_RESUME_OUT"));
let term_count: u64 = std::env::var("RETE_RESUME_TERMS").unwrap().parse().unwrap();
let quad_count: u64 = std::env::var("RETE_RESUME_QUADS").unwrap().parse().unwrap();
let metadata: Vec<u8> = std::env::var("RETE_RESUME_CARD")
.ok()
.map(|p| std::fs::read(p).unwrap())
.unwrap_or_default();
let sec = |name: &str| -> SectionFile {
let path = spill.join(name);
let len = std::fs::metadata(&path).expect(name).len();
SectionFile { path, len }
};
let merged = MergedDict {
section_files: [
sec("g.shared.sec"),
sec("g.subj.sec"),
sec("g.obj.sec"),
sec("g.pred.sec"),
],
term_count,
has_quoted: false,
remaps: Vec::new(),
};
let tmp = TmpDir { dir: spill.clone() };
let global_tri = spill.join("global.tri");
let codec = crate::file::writer_codec();
let budget_mb: u64 = std::env::var("RETE_RESUME_BUDGET_MB")
.ok()
.and_then(|v| v.parse().ok())
.unwrap_or(16384);
let run_len = ((budget_mb << 20) / 2 / 24) as usize;
let mut sections = Vec::new();
for perm in crate::index::ALL_PERMS {
let done = spill.join(format!("{}.tiles.sec", perm.name()));
if done.exists() {
eprintln!("resume: reusing {}", perm.name());
let len = std::fs::metadata(&done).unwrap().len();
sections.push(SectionFile { path: done, len });
continue;
}
for entry in std::fs::read_dir(&spill).unwrap().flatten() {
let n = entry.file_name().to_string_lossy().into_owned();
if n.starts_with(&format!("{}.", perm.name())) {
let _ = std::fs::remove_file(entry.path());
}
}
eprintln!("resume: rebuilding {}", perm.name());
let (s, n) =
build_permutation_section(&tmp, &global_tri, perm, run_len, codec).unwrap();
assert_eq!(n, quad_count, "permutation count must match");
eprintln!("resume: {} done", perm.name());
sections.push(s);
}
write_final_file(&out, &metadata, &merged, §ions, quad_count, codec).unwrap();
eprintln!("resume: wrote {}", out.display());
std::mem::forget(tmp); }
#[test]
fn named_graphs_are_rejected() {
let dir = std::env::temp_dir().join(format!("rete-extbuild-ng-{}", std::process::id()));
std::fs::create_dir_all(&dir).unwrap();
let out = dir.join("out.rete");
let err = build_external(
|visit| {
visit((
"<http://ex/s>".into(),
"<http://ex/p>".into(),
"<http://ex/o>".into(),
Some("<http://ex/g>".into()),
))
},
&out,
ExternalBuildOptions {
memory_budget: 0,
tmp_dir: Some(dir.clone()),
metadata: Box::new(|_| Vec::new()),
},
)
.unwrap_err();
assert!(matches!(err, ExtBuildError::NamedGraph(_)));
let _ = std::fs::remove_dir_all(&dir);
}
}