use std::collections::VecDeque;
use crate::chunk::Chunk;
use crate::error::{ChunkError, Result};
use crate::formats::md::block_stream::{resume_point, BlockStream};
use crate::formats::md::common::MIN_CHUNK_CHARS;
use crate::formats::md::structural::{ProseMerger, StructuralBuilder};
use crate::formats::pipeline;
use super::markdown::PAGE_SEPARATOR;
use super::{author_block, parse};
#[cfg(not(target_arch = "wasm32"))]
const CHANNEL_DEPTH: usize = 64;
pub struct PdfChunkStream {
backend: Backend,
}
enum Backend {
#[cfg(not(target_arch = "wasm32"))]
Threaded(std::sync::mpsc::Receiver<Result<Chunk>>),
#[allow(dead_code)]
Deferred(Box<Deferred>),
#[cfg(target_arch = "wasm32")]
Incremental(Box<Incremental>),
Draining(std::vec::IntoIter<Chunk>),
Failed(Option<ChunkError>),
Done,
}
struct Deferred {
bytes: Vec<u8>,
mode: String,
window_size: usize,
overlap: usize,
sentences_per_chunk: usize,
paragraphs_per_page: usize,
}
struct Incremental {
reader: parse::Reader,
metadata: serde_json::Value,
total_pages: usize,
pending: String,
blocks: BlockStream,
builder: StructuralBuilder,
merger: ProseMerger,
queue: VecDeque<Chunk>,
wrote_a_page: bool,
drained: bool,
}
impl Incremental {
fn open(bytes: &[u8]) -> Result<Incremental> {
let reader = parse::Reader::open(bytes, parse::Headings::PerPage).map_err(ChunkError::Parse)?;
let total_pages = reader.total_pages();
Ok(Incremental {
reader,
metadata: serde_json::json!({ "source_type": "pdf", "total_pages": total_pages }),
total_pages,
pending: String::new(),
blocks: BlockStream::new(),
builder: StructuralBuilder::new(),
merger: ProseMerger::new(),
queue: VecDeque::new(),
wrote_a_page: false,
drained: false,
})
}
fn drain(&mut self, flush: bool) {
let cut = if flush { self.pending.len() } else { resume_point(&self.pending) };
if cut == 0 {
return;
}
let rest = self.pending.split_off(cut);
let ready = std::mem::replace(&mut self.pending, rest);
let normalized = author_block::normalize(&ready);
for block in self.blocks.push(&normalized) {
self.builder.advance(block);
}
if flush {
for block in self.blocks.finish() {
self.builder.advance(block);
}
}
let finished = if flush { self.builder.finish() } else { self.builder.take() };
for record in finished {
if let Some(out) = self.merger.push(record, MIN_CHUNK_CHARS) {
self.emit(out);
}
}
if flush {
if let Some(out) = self.merger.finish() {
self.emit(out);
}
}
}
fn emit(&mut self, record: crate::formats::md::common::SpannedRecord) {
let rec = pipeline::stamp(record, &self.metadata, None);
self.queue.push_back(Chunk::new(rec.content, rec.content_type.as_str(), rec.metadata));
}
fn pump(&mut self) -> Option<Result<Chunk>> {
loop {
if let Some(chunk) = self.queue.pop_front() {
return Some(Ok(chunk));
}
if self.drained {
return None;
}
match self.reader.next_page() {
Some(page) => {
if !page.trim().is_empty() {
if self.wrote_a_page {
self.pending.push_str(PAGE_SEPARATOR);
}
self.wrote_a_page = true;
self.pending.push_str(&page);
}
if self.reader.has_text() {
self.drain(false);
}
}
None => {
self.drained = true;
if !self.reader.has_text() {
return Some(Err(ChunkError::Parse(format!(
"PDF contains no extractable text ({} page(s) scanned or image-only). OCR is not enabled; pass list_images to get one rendered image per page.",
self.total_pages
))));
}
self.drain(true);
}
}
}
}
}
impl Deferred {
fn run(&self) -> Result<Vec<Chunk>> {
let loaded = super::load(&self.bytes, false, super::headings_for(&self.mode))?;
pipeline::chunk(
&loaded,
&self.mode,
self.window_size,
self.overlap,
self.sentences_per_chunk,
self.paragraphs_per_page,
)
}
}
impl Iterator for PdfChunkStream {
type Item = Result<Chunk>;
fn next(&mut self) -> Option<Self::Item> {
loop {
match &mut self.backend {
#[cfg(not(target_arch = "wasm32"))]
Backend::Threaded(rx) => return rx.recv().ok(),
Backend::Draining(chunks) => return chunks.next().map(Ok),
#[cfg(target_arch = "wasm32")]
Backend::Incremental(work) => return work.pump(),
Backend::Failed(error) => {
let error = error.take();
self.backend = Backend::Done;
return error.map(Err);
}
Backend::Done => return None,
Backend::Deferred(work) => {
self.backend = match work.run() {
Ok(chunks) => Backend::Draining(chunks.into_iter()),
Err(error) => Backend::Failed(Some(error)),
};
}
}
}
}
}
#[cfg(not(target_arch = "wasm32"))]
fn incremental(bytes: Vec<u8>) -> PdfChunkStream {
let (tx, rx) = std::sync::mpsc::sync_channel(CHANNEL_DEPTH);
std::thread::spawn(move || {
let mut work = match Incremental::open(&bytes) {
Ok(work) => work,
Err(error) => {
let _ = tx.send(Err(error));
return;
}
};
while let Some(item) = work.pump() {
let failed = item.is_err();
if tx.send(item).is_err() || failed {
return;
}
}
});
PdfChunkStream { backend: Backend::Threaded(rx) }
}
#[cfg(target_arch = "wasm32")]
fn incremental(bytes: Vec<u8>) -> PdfChunkStream {
match Incremental::open(&bytes) {
Ok(work) => PdfChunkStream { backend: Backend::Incremental(Box::new(work)) },
Err(error) => PdfChunkStream { backend: Backend::Failed(Some(error)) },
}
}
pub fn stream_from_bytes(
bytes: Vec<u8>,
mode: &str,
window_size: usize,
overlap: usize,
sentences_per_chunk: usize,
paragraphs_per_page: usize,
) -> PdfChunkStream {
if mode == "default" {
return incremental(bytes);
}
let work = Deferred {
bytes,
mode: mode.to_string(),
window_size,
overlap,
sentences_per_chunk,
paragraphs_per_page,
};
#[cfg(not(target_arch = "wasm32"))]
let backend = {
let (tx, rx) = std::sync::mpsc::sync_channel(CHANNEL_DEPTH);
std::thread::spawn(move || match work.run() {
Ok(chunks) => {
for chunk in chunks {
if tx.send(Ok(chunk)).is_err() {
return;
}
}
}
Err(error) => {
let _ = tx.send(Err(error));
}
});
Backend::Threaded(rx)
};
#[cfg(target_arch = "wasm32")]
let backend = Backend::Deferred(Box::new(work));
PdfChunkStream { backend }
}
#[cfg(test)]
mod tests {
use super::*;
fn fixture(name: &str) -> std::path::PathBuf {
[env!("CARGO_MANIFEST_DIR"), "..", "test_files", "pdf", name].iter().collect()
}
#[test]
fn the_first_chunk_costs_a_few_pages_not_the_whole_document() {
let path = fixture("sample-5000-page.pdf");
if !path.exists() {
return;
}
let bytes = std::fs::read(&path).expect("read");
let mut work = Incremental::open(&bytes).expect("open");
assert_eq!(work.reader.pages_rendered(), 0, "opening a stream must render nothing");
assert_eq!(work.total_pages, 5_000, "fixture changed");
let first = work.pump().expect("a chunk").expect("chunked");
assert!(!first.content.is_empty());
let pages = work.reader.pages_rendered();
assert!(pages < 10, "the first chunk cost {pages} pages of {}", work.total_pages);
}
#[test]
fn a_scanned_document_raises_instead_of_chunking_its_image_references() {
let path = fixture("large-doc.pdf");
if !path.exists() {
return;
}
let bytes = std::fs::read(&path).expect("read");
let mut work = Incremental::open(&bytes).expect("open");
match work.pump() {
Some(Err(error)) => assert!(error.to_string().contains("no extractable text")),
other => panic!("expected an error, got {:?}", other.map(|r| r.is_ok())),
}
}
}