ruwex 0.1.0

Fast Rust rewrite of wikiextractor: extract and clean text from Wikimedia XML dumps
Documentation
//! Orchestration: source → worker pool → order-preserving writer.
//!
//! Two shapes, both preserving dump order in the output:
//!
//! - **Sequential** input: one reader thread parses pages and feeds them to
//!   render workers through a bounded channel; the calling thread reorders
//!   results by sequence number and writes them.
//! - **Multistream** input: each worker owns a file handle and processes
//!   whole bz2 streams (seek → decompress → parse → render), so decompression
//!   itself is parallel; the calling thread reorders by stream number.

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 {
    /// Pages read from the dump.
    pub pages: u64,
    /// Documents written (pages that passed filtering).
    pub docs: u64,
}

/// Extracts all pages from `source` into `sink` according to `config`,
/// without template definitions (like `--no-templates`).
pub fn run(source: PageSource, config: &ExtractorConfig, sink: &mut dyn DocSink) -> Result<Stats> {
    run_with_templates(source, &Arc::new(TemplateDb::default()), config, sink)
}

/// Extracts all pages from `source` into `sink`, expanding templates
/// against `templates` (see [`crate::expand::templates::load`]).
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),
    }
}

/// Decides which pages become output documents. Like wikiextractor this is
/// title-prefix based: pages without a namespace prefix are kept, prefixed
/// ones only if the prefix is in the accepted list (case-sensitive).
/// Divergence (documented): redirects are always skipped, where the original
/// emits an empty document for main-namespace redirects due to an operator
/// precedence slip.
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; // downstream gone (writer error); stop reading
                    }
                    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 })
    })
}

/// Drains `(sequence, document)` pairs, writing them in sequence order.
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)
        );
    }
}

/// Documents rendered from one multistream block, in dump order.
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)
    })
}

/// Decompresses one multistream block and renders its accepted pages.
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)
}