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 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 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}