Skip to main content

cloudiful_docling_convert/api/
docling.rs

1use std::time::{Duration, Instant};
2
3use reqwest::multipart;
4use serde_json::Value;
5
6use crate::document::{ChunkerKind, ChunkingOptions, InputDocument, OutputFormat, PipelineKind};
7use crate::error::{PdfConvertError, Result};
8use crate::models::{ConversionStatus, TaskPostResponse, TaskStatusResponse};
9
10use super::chunk::{build_convert_file_form, build_file_form};
11use super::result::{DoclingResult, DoclingTaskResult, parse_response};
12use super::source::{chunk_source_request, source_request};
13use super::transport::{
14    Transport, default_request_timeout, default_task_timeout, handle_response, retry_with_backoff,
15};
16
17#[derive(Debug, Clone)]
18pub struct DoclingConfig {
19    pub base_url: String,
20    pub openai_base_url: String,
21    pub vlm_pipeline_model: String,
22    pub picture_description_model: String,
23    pub code_formula_model: String,
24    pub api_key: Option<String>,
25    pub openai_api_key: Option<String>,
26    pub tenant_id: Option<String>,
27    pub request_timeout: Option<Duration>,
28    pub task_timeout: Option<Duration>,
29}
30
31impl DoclingConfig {
32    pub fn without_vlm(base_url: impl Into<String>) -> Self {
33        Self {
34            base_url: base_url.into(),
35            openai_base_url: String::new(),
36            vlm_pipeline_model: String::new(),
37            picture_description_model: String::new(),
38            code_formula_model: String::new(),
39            api_key: None,
40            openai_api_key: None,
41            tenant_id: None,
42            request_timeout: None,
43            task_timeout: None,
44        }
45    }
46}
47
48#[derive(Debug, Clone)]
49pub struct DoclingConvertRequest {
50    pub output_formats: Vec<OutputFormat>,
51    pub page_range: Option<(u32, u32)>,
52    pub chunker: ChunkerKind,
53    pub chunking: ChunkingOptions,
54    pub pipeline: Option<PipelineKind>,
55    /// Preset ID for picture description, forwarded to Docling Serve as
56    /// `picture_description_preset`. When set, the legacy
57    /// `picture_description_custom_config` form field is suppressed so the
58    /// preset wins. When `None`, the legacy custom VLM configuration is
59    /// preserved unchanged.
60    pub picture_description_preset: Option<String>,
61}
62
63impl DoclingConvertRequest {
64    pub fn for_outputs(output_formats: Vec<OutputFormat>) -> Self {
65        Self {
66            output_formats,
67            page_range: None,
68            chunker: ChunkerKind::None,
69            chunking: ChunkingOptions::hybrid_defaults(),
70            pipeline: None,
71            picture_description_preset: None,
72        }
73    }
74
75    pub fn with_chunker(mut self, chunker: ChunkerKind, options: ChunkingOptions) -> Self {
76        self.chunker = chunker;
77        self.chunking = options;
78        self
79    }
80
81    pub fn with_pipeline(mut self, pipeline: Option<PipelineKind>) -> Self {
82        self.pipeline = pipeline;
83        self
84    }
85
86    pub fn with_picture_description_preset(
87        mut self,
88        picture_description_preset: Option<String>,
89    ) -> Self {
90        self.picture_description_preset = picture_description_preset;
91        self
92    }
93}
94
95#[derive(Clone)]
96pub struct DoclingClient {
97    transport: Transport,
98    result_body_limit: Option<usize>,
99}
100
101impl std::fmt::Debug for DoclingClient {
102    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
103        f.debug_struct("DoclingClient")
104            .field("base_url", &self.transport.config().base_url)
105            .field("tenant_id", &self.transport.config().tenant_id)
106            .finish_non_exhaustive()
107    }
108}
109
110impl DoclingClient {
111    pub fn new(config: DoclingConfig) -> Result<Self> {
112        Ok(Self {
113            transport: Transport::new(config)?,
114            result_body_limit: None,
115        })
116    }
117
118    pub fn new_with_result_body_limit(
119        config: DoclingConfig,
120        result_body_limit: usize,
121    ) -> Result<Self> {
122        if result_body_limit == 0 {
123            return Err(PdfConvertError::validation_error(
124                "result_body_limit",
125                "value must be greater than 0",
126            ));
127        }
128        Ok(Self {
129            transport: Transport::new(config)?,
130            result_body_limit: Some(result_body_limit),
131        })
132    }
133
134    pub fn config(&self) -> &DoclingConfig {
135        self.transport.config()
136    }
137
138    pub fn request_timeout(&self) -> Duration {
139        self.config()
140            .request_timeout
141            .unwrap_or_else(default_request_timeout)
142    }
143
144    pub fn task_timeout(&self) -> Duration {
145        self.config()
146            .task_timeout
147            .unwrap_or_else(default_task_timeout)
148    }
149
150    pub async fn convert_file(
151        &self,
152        input: &InputDocument,
153        request: &DoclingConvertRequest,
154    ) -> Result<DoclingResult> {
155        let operation = || async {
156            let form = self.build_form(input, request)?;
157            let path = match request.chunker {
158                ChunkerKind::None => "convert/file",
159                ChunkerKind::Hybrid => "chunk/hybrid/file",
160                ChunkerKind::Hierarchical => "chunk/hierarchical/file",
161            };
162            let response = self
163                .transport
164                .request(reqwest::Method::POST, path)
165                .multipart(form)
166                .send()
167                .await
168                .map_err(PdfConvertError::from)?;
169            parse_response(response, "Docling file conversion", self.result_body_limit).await
170        };
171
172        retry_with_backoff(operation, "docling_convert_file").await
173    }
174
175    pub async fn submit_file_async(
176        &self,
177        input: &InputDocument,
178        request: &DoclingConvertRequest,
179    ) -> Result<String> {
180        let operation = || async {
181            let form = self.build_form(input, request)?;
182            let path = match request.chunker {
183                ChunkerKind::None => "convert/file/async",
184                ChunkerKind::Hybrid => "chunk/hybrid/file/async",
185                ChunkerKind::Hierarchical => "chunk/hierarchical/file/async",
186            };
187            let response = self
188                .transport
189                .request(reqwest::Method::POST, path)
190                .multipart(form)
191                .send()
192                .await
193                .map_err(PdfConvertError::from)?;
194            let response = handle_response(response, "Docling async submission").await?;
195            let task = response.json::<TaskPostResponse>().await.map_err(|error| {
196                PdfConvertError::parse_error("Docling async submission response", error.to_string())
197            })?;
198            Ok(task.task_id)
199        };
200
201        retry_with_backoff(operation, "docling_submit_file_async").await
202    }
203
204    pub async fn convert_source(
205        &self,
206        url: &str,
207        input_kind: crate::document::InputKind,
208        request: &DoclingConvertRequest,
209    ) -> Result<DoclingResult> {
210        let operation = || async {
211            let path = match request.chunker {
212                ChunkerKind::None => "convert/source",
213                ChunkerKind::Hybrid => "chunk/hybrid/source",
214                ChunkerKind::Hierarchical => "chunk/hierarchical/source",
215            };
216            let body = if request.chunker == ChunkerKind::None {
217                serde_json::to_value(source_request(url, input_kind, request))?
218            } else {
219                serde_json::to_value(chunk_source_request(url, input_kind, request)?)?
220            };
221            let response = self
222                .transport
223                .request(reqwest::Method::POST, path)
224                .json(&body)
225                .send()
226                .await
227                .map_err(PdfConvertError::from)?;
228            parse_response(
229                response,
230                "Docling source conversion",
231                self.result_body_limit,
232            )
233            .await
234        };
235
236        retry_with_backoff(operation, "docling_convert_source").await
237    }
238
239    pub async fn submit_source_async(
240        &self,
241        url: &str,
242        input_kind: crate::document::InputKind,
243        request: &DoclingConvertRequest,
244    ) -> Result<String> {
245        let operation = || async {
246            let path = match request.chunker {
247                ChunkerKind::None => "convert/source/async",
248                ChunkerKind::Hybrid => "chunk/hybrid/source/async",
249                ChunkerKind::Hierarchical => "chunk/hierarchical/source/async",
250            };
251            let body = if request.chunker == ChunkerKind::None {
252                serde_json::to_value(source_request(url, input_kind, request))?
253            } else {
254                serde_json::to_value(chunk_source_request(url, input_kind, request)?)?
255            };
256            let response = self
257                .transport
258                .request(reqwest::Method::POST, path)
259                .json(&body)
260                .send()
261                .await
262                .map_err(PdfConvertError::from)?;
263            let response = handle_response(response, "Docling source async submission").await?;
264            let task = response.json::<TaskPostResponse>().await.map_err(|error| {
265                PdfConvertError::parse_error(
266                    "Docling source async submission response",
267                    error.to_string(),
268                )
269            })?;
270            Ok(task.task_id)
271        };
272
273        retry_with_backoff(operation, "docling_submit_source_async").await
274    }
275
276    pub async fn wait_for_result(&self, task_id: &str) -> Result<DoclingTaskResult> {
277        self.wait_for_result_with_progress(task_id, |_| async {})
278            .await
279    }
280
281    pub async fn wait_for_result_with_progress<F, Fut>(
282        &self,
283        task_id: &str,
284        mut on_status: F,
285    ) -> Result<DoclingTaskResult>
286    where
287        F: FnMut(TaskStatusResponse) -> Fut + Send,
288        Fut: std::future::Future<Output = ()> + Send,
289    {
290        let deadline = Instant::now() + self.task_timeout();
291        loop {
292            let status = self.poll_task_status(task_id).await?;
293            on_status(status.clone()).await;
294            if status.task_status.is_terminal() {
295                if matches!(
296                    status.task_status,
297                    ConversionStatus::Failure | ConversionStatus::Skipped
298                ) {
299                    return Err(task_failure_error(&status));
300                }
301                let result = self.get_task_result(task_id).await?;
302                return Ok(DoclingTaskResult {
303                    status: status.task_status,
304                    result,
305                    errors: task_status_errors(&status),
306                });
307            }
308
309            if Instant::now() >= deadline {
310                return Err(PdfConvertError::operation_error(
311                    "waiting for Docling task",
312                    format!(
313                        "task {task_id} did not reach a terminal state within {:?}",
314                        self.task_timeout()
315                    ),
316                ));
317            }
318        }
319    }
320
321    /// Build a [`DoclingTaskResult`] for a task whose terminal status was already
322    /// observed via [`Self::poll_task_status`]. This is the resumable counterpart of
323    /// [`Self::wait_for_result`]: the caller controls polling cadence and deadline.
324    pub async fn fetch_task_result(
325        &self,
326        task_id: &str,
327        status: &TaskStatusResponse,
328    ) -> Result<DoclingTaskResult> {
329        if matches!(
330            status.task_status,
331            ConversionStatus::Failure | ConversionStatus::Skipped
332        ) {
333            return Err(task_failure_error(status));
334        }
335        if !status.task_status.is_terminal() {
336            return Err(PdfConvertError::operation_error(
337                "fetching Docling task result",
338                format!(
339                    "task {task_id} is not terminal yet: {:?}",
340                    status.task_status
341                ),
342            ));
343        }
344        let result = self.get_task_result(task_id).await?;
345        Ok(DoclingTaskResult {
346            status: status.task_status,
347            result,
348            errors: task_status_errors(status),
349        })
350    }
351
352    pub async fn poll_task_status(&self, task_id: &str) -> Result<TaskStatusResponse> {
353        let operation = || async {
354            let path = format!("status/poll/{task_id}");
355            let response = self
356                .transport
357                .request(reqwest::Method::GET, &format!("{path}?wait=30"))
358                .send()
359                .await
360                .map_err(PdfConvertError::from)?;
361            let response = handle_response(response, "Polling task status").await?;
362            let text = response.text().await.map_err(PdfConvertError::from)?;
363            serde_json::from_str::<TaskStatusResponse>(&text).map_err(|error| {
364                PdfConvertError::parse_error(
365                    "task status response",
366                    format!("task {task_id} returned invalid response: {error}; body: {text}"),
367                )
368            })
369        };
370
371        retry_with_backoff(operation, &format!("check_task_status({task_id})")).await
372    }
373
374    pub async fn check_task_status(&self, task_id: &str) -> Result<bool> {
375        let status = self.poll_task_status(task_id).await?;
376        if matches!(
377            status.task_status,
378            ConversionStatus::Failure | ConversionStatus::Skipped
379        ) {
380            return Err(PdfConvertError::api_task_failed(
381                format!("{:?}", status.task_status),
382                task_status_error(&status),
383            ));
384        }
385        Ok(status.task_status.is_successful())
386    }
387
388    pub async fn get_task_result(&self, task_id: &str) -> Result<DoclingResult> {
389        self.get_task_result_with_connection(task_id, false).await
390    }
391
392    pub async fn get_task_result_with_connection(
393        &self,
394        task_id: &str,
395        close_connection: bool,
396    ) -> Result<DoclingResult> {
397        let operation = || async {
398            let path = format!("result/{task_id}");
399            let mut request = self.transport.request(reqwest::Method::GET, &path);
400            if close_connection {
401                request = request.header(reqwest::header::CONNECTION, "close");
402            }
403            let response = request.send().await.map_err(PdfConvertError::from)?;
404            parse_response(response, "Fetching task result", self.result_body_limit).await
405        };
406
407        retry_with_backoff(operation, &format!("get_task_result({task_id})")).await
408    }
409
410    pub async fn get_task_result_value(&self, task_id: &str) -> Result<Value> {
411        self.get_task_result(task_id)
412            .await?
413            .into_json()
414            .ok_or_else(|| {
415                PdfConvertError::operation_error(
416                    "reading task result",
417                    "task result is a ZIP response and has no JSON value",
418                )
419            })
420    }
421
422    pub(crate) fn build_form(
423        &self,
424        input: &InputDocument,
425        request: &DoclingConvertRequest,
426    ) -> Result<multipart::Form> {
427        match request.chunker {
428            ChunkerKind::None => build_convert_file_form(self, input, request),
429            ChunkerKind::Hybrid | ChunkerKind::Hierarchical => build_file_form(input, request),
430        }
431    }
432}
433
434fn task_status_error(status: &TaskStatusResponse) -> String {
435    status
436        .error_message
437        .clone()
438        .or_else(|| {
439            status
440                .failure
441                .as_ref()
442                .map(|failure| failure.message.clone())
443        })
444        .unwrap_or_else(|| format!("Docling task status is {:?}", status.task_status))
445}
446
447fn task_status_errors(status: &TaskStatusResponse) -> Vec<String> {
448    let mut errors = Vec::new();
449    if let Some(message) = status.error_message.as_deref() {
450        errors.push(message.to_string());
451    }
452    if let Some(failure) = status.failure.as_ref() {
453        errors.push(failure.message.clone());
454    }
455    errors.sort();
456    errors.dedup();
457    errors
458}
459
460fn task_failure_error(status: &TaskStatusResponse) -> PdfConvertError {
461    PdfConvertError::api_task_failed(
462        format!("{:?}", status.task_status),
463        task_status_error(status),
464    )
465}