use std::sync::mpsc::{sync_channel, Receiver};
use std::thread::JoinHandle;
use docling_core::{ImageMode, MarkdownStreamer};
use crate::converter::DocumentConverter;
use crate::error::ConversionError;
use crate::format::InputFormat;
use crate::source::SourceDocument;
const CHANNEL_DEPTH: usize = 8;
pub struct MarkdownStream {
rx: Option<Receiver<Result<String, ConversionError>>>,
handle: Option<JoinHandle<()>>,
}
impl Iterator for MarkdownStream {
type Item = Result<String, ConversionError>;
fn next(&mut self) -> Option<Self::Item> {
match self.rx.as_ref()?.recv() {
Ok(item) => Some(item),
Err(_) => {
if let Some(h) = self.handle.take() {
let _ = h.join();
}
None
}
}
}
}
impl Drop for MarkdownStream {
fn drop(&mut self) {
self.rx = None;
if let Some(h) = self.handle.take() {
let _ = h.join();
}
}
}
pub(crate) fn spawn(
converter: DocumentConverter,
source: SourceDocument,
image_mode: ImageMode,
strict: bool,
no_table_former: bool,
no_ocr: bool,
enrich: docling_pdf::EnrichmentOptions,
) -> MarkdownStream {
let (tx, rx) = sync_channel::<Result<String, ConversionError>>(CHANNEL_DEPTH);
let handle = std::thread::spawn(move || {
run(
converter,
source,
image_mode,
strict,
no_table_former,
no_ocr,
enrich,
&tx,
)
});
MarkdownStream {
rx: Some(rx),
handle: Some(handle),
}
}
#[allow(clippy::too_many_arguments)]
fn run(
converter: DocumentConverter,
source: SourceDocument,
image_mode: ImageMode,
strict: bool,
no_table_former: bool,
no_ocr: bool,
enrich: docling_pdf::EnrichmentOptions,
tx: &std::sync::mpsc::SyncSender<Result<String, ConversionError>>,
) {
match source.format {
InputFormat::Pdf => run_pdf(
&source,
image_mode,
strict,
no_table_former,
no_ocr,
enrich,
tx,
),
_ => run_buffered(converter, source, image_mode, strict, tx),
}
}
fn run_pdf(
source: &SourceDocument,
image_mode: ImageMode,
strict: bool,
no_table_former: bool,
no_ocr: bool,
enrich: docling_pdf::EnrichmentOptions,
tx: &std::sync::mpsc::SyncSender<Result<String, ConversionError>>,
) {
let mut streamer = MarkdownStreamer::new(strict, image_mode, false);
let mut pipeline = match docling_pdf::Pipeline::new().map(|p| {
p.no_table_former(no_table_former)
.no_ocr(no_ocr)
.enrichments(enrich)
}) {
Ok(p) => p,
Err(e) => {
let _ = tx.send(Err(ConversionError::Parse(e.to_string())));
return;
}
};
let result = pipeline.convert_streaming(&source.bytes, None, &source.name, |nodes, links| {
let chunk = streamer.push(&nodes, &links);
if !chunk.is_empty() && tx.send(Ok(chunk)).is_err() {
return Err(docling_pdf::PdfError::Pdfium(
"markdown stream consumer dropped".into(),
));
}
Ok(())
});
match result {
Ok(()) => {
let tail = streamer.finish();
if !tail.is_empty() {
let _ = tx.send(Ok(tail));
}
}
Err(e) => {
let _ = tx.send(Err(ConversionError::Parse(e.to_string())));
}
}
}
fn run_buffered(
converter: DocumentConverter,
source: SourceDocument,
image_mode: ImageMode,
strict: bool,
tx: &std::sync::mpsc::SyncSender<Result<String, ConversionError>>,
) {
let doc = match converter.convert(source) {
Ok(result) => result.document,
Err(e) => {
let _ = tx.send(Err(e));
return;
}
};
let mut streamer = MarkdownStreamer::new(strict, image_mode, doc.compact_tables);
let chunk = streamer.push(&doc.nodes, &doc.links);
if !chunk.is_empty() && tx.send(Ok(chunk)).is_err() {
return;
}
let tail = streamer.finish();
if !tail.is_empty() {
let _ = tx.send(Ok(tail));
}
}