use std::collections::{BTreeMap, HashSet};
use std::fs::File;
use std::io::BufRead;
use std::sync::Arc;
use std::thread;
use std::time::Instant;
use crossbeam_channel::{Receiver, bounded};
use crate::config::ExtractorConfig;
use crate::dump::reader::{MultistreamSource, PageSource};
use crate::dump::xml::DumpParser;
use crate::dump::{Page, SiteInfo, multistream};
use crate::error::{Error, Result};
use crate::expand::TemplateDb;
use crate::output::format::render_page;
use crate::output::writer::DocSink;
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
pub struct Stats {
pub pages: u64,
pub docs: u64,
}
pub fn run(source: PageSource, config: &ExtractorConfig, sink: &mut dyn DocSink) -> Result<Stats> {
run_with_templates(source, &Arc::new(TemplateDb::default()), config, sink)
}
pub fn run_with_templates(
source: PageSource,
templates: &Arc<TemplateDb>,
config: &ExtractorConfig,
sink: &mut dyn DocSink,
) -> Result<Stats> {
match source {
PageSource::Sequential(reader) => run_sequential(reader, templates, config, sink),
PageSource::Multistream(source) => run_multistream(source, templates, config, sink),
}
}
struct PageFilter {
accepted: HashSet<String>,
}
impl PageFilter {
fn new(config: &ExtractorConfig) -> Self {
Self {
accepted: config.namespaces.iter().cloned().collect(),
}
}
fn keep(&self, page: &Page) -> bool {
if page.redirect.is_some() {
return false;
}
match page.title.find(':') {
None => true,
Some(colon) => self.accepted.contains(&page.title[..colon]),
}
}
}
fn run_sequential(
input: Box<dyn BufRead + Send>,
templates: &Arc<TemplateDb>,
config: &ExtractorConfig,
sink: &mut dyn DocSink,
) -> Result<Stats> {
let mut parser = DumpParser::new(input);
let site = Arc::new(site_info_or_warn(&mut parser)?);
let filter = PageFilter::new(config);
let workers = config.workers.max(1);
let capacity = workers * 8;
let (page_tx, page_rx) = bounded::<(u64, Page)>(capacity);
let (doc_tx, doc_rx) = bounded::<(u64, String)>(capacity);
thread::scope(|scope| -> Result<Stats> {
let reader = scope.spawn(move || -> Result<u64> {
let mut pages = 0;
let mut seq = 0;
while let Some(page) = parser.next_page()? {
pages += 1;
if filter.keep(&page) {
if page_tx.send((seq, page)).is_err() {
break; }
seq += 1;
}
}
Ok(pages)
});
for _ in 0..workers {
let page_rx = page_rx.clone();
let doc_tx = doc_tx.clone();
let site = Arc::clone(&site);
let templates = Arc::clone(templates);
scope.spawn(move || {
for (seq, page) in page_rx {
let doc = render_page(&page, &site, config, templates.as_ref());
if doc_tx.send((seq, doc)).is_err() {
break;
}
}
});
}
drop(page_rx);
drop(doc_tx);
let docs = write_ordered(doc_rx, sink)?;
let pages = reader.join().expect("reader thread panicked")?;
Ok(Stats { pages, docs })
})
}
fn write_ordered(rx: Receiver<(u64, String)>, sink: &mut dyn DocSink) -> Result<u64> {
let started = Instant::now();
let mut pending = BTreeMap::new();
let mut next = 0;
for (seq, doc) in rx {
pending.insert(seq, doc);
while let Some(doc) = pending.remove(&next) {
sink.write_doc(&doc)?;
next += 1;
log_progress(next, &started);
}
}
debug_assert!(pending.is_empty(), "gap in document sequence numbers");
Ok(next)
}
fn log_progress(docs: u64, started: &Instant) {
if docs.is_multiple_of(10_000) {
log::info!(
"{} documents extracted ({:.0} docs/s)",
docs,
docs as f64 / started.elapsed().as_secs_f64().max(1e-9)
);
}
}
struct StreamOutput {
docs: Vec<String>,
pages: u64,
}
fn run_multistream(
source: MultistreamSource,
templates: &Arc<TemplateDb>,
config: &ExtractorConfig,
sink: &mut dyn DocSink,
) -> Result<Stats> {
let mut file = File::open(&source.path)?;
let header = multistream::read_stream(&mut file, 0)?;
let mut header_parser = DumpParser::new(header.as_slice());
let site = Arc::new(site_info_or_warn(&mut header_parser)?);
let filter = Arc::new(PageFilter::new(config));
let workers = config.workers.max(1);
let stream_count = source.offsets.len() as u64;
let (job_tx, job_rx) = bounded::<(u64, u64)>(workers * 4);
let (result_tx, result_rx) = bounded::<(u64, Result<StreamOutput>)>(workers * 2);
thread::scope(|scope| -> Result<Stats> {
let offsets = source.offsets;
scope.spawn(move || {
for (seq, offset) in offsets.into_iter().enumerate() {
if job_tx.send((seq as u64, offset)).is_err() {
break;
}
}
});
for _ in 0..workers {
let job_rx = job_rx.clone();
let result_tx = result_tx.clone();
let site = Arc::clone(&site);
let filter = Arc::clone(&filter);
let templates = Arc::clone(templates);
let path = source.path.clone();
scope.spawn(move || -> Result<()> {
let mut file = File::open(&path)?;
for (seq, offset) in job_rx {
let output =
process_stream(&mut file, offset, &site, &filter, config, &templates);
let failed = output.is_err();
if result_tx.send((seq, output)).is_err() || failed {
break;
}
}
Ok(())
});
}
drop(job_rx);
drop(result_tx);
let started = Instant::now();
let mut pending: BTreeMap<u64, StreamOutput> = BTreeMap::new();
let mut next = 0;
let mut stats = Stats::default();
for (seq, output) in result_rx {
pending.insert(seq, output?);
while let Some(output) = pending.remove(&next) {
stats.pages += output.pages;
for doc in &output.docs {
sink.write_doc(doc)?;
stats.docs += 1;
log_progress(stats.docs, &started);
}
next += 1;
}
}
if next != stream_count {
return Err(Error::InvalidDump(format!(
"processed {next} of {stream_count} multistream blocks"
)));
}
Ok(stats)
})
}
fn process_stream(
file: &mut File,
offset: u64,
site: &SiteInfo,
filter: &PageFilter,
config: &ExtractorConfig,
templates: &TemplateDb,
) -> Result<StreamOutput> {
let bytes = multistream::read_stream(file, offset)?;
let mut parser = DumpParser::new(bytes.as_slice());
let mut output = StreamOutput {
docs: Vec::new(),
pages: 0,
};
while let Some(page) = parser.next_page()? {
output.pages += 1;
if filter.keep(&page) {
output
.docs
.push(render_page(&page, site, config, templates));
}
}
Ok(output)
}
fn site_info_or_warn<R: BufRead>(parser: &mut DumpParser<R>) -> Result<SiteInfo> {
let site = parser.site_info()?.cloned().unwrap_or_default();
if site.base.is_empty() {
log::warn!("dump has no <siteinfo>; document URLs will be incomplete");
}
Ok(site)
}