Coverage for open_webui/retrieval/loaders/mineru.py: 6%
255 statements
« prev ^ index » next coverage.py v7.15.2, created at 2026-10-07 05:07 +0000
« prev ^ index » next coverage.py v7.15.2, created at 2026-10-07 05:07 +0000
1import logging
2import os
3import tempfile
4import time
5import zipfile
6from typing import List, Optional
8import requests
9from fastapi import HTTPException, status
10from langchain_core.documents import Document
12log = logging.getLogger(__name__)
15class MinerULoader:
16 """
17 MinerU document parser loader supporting both Cloud API and Local API modes.
19 Cloud API: Uses MinerU managed service with async task-based processing
20 Local API: Uses self-hosted MinerU API with synchronous processing
21 """
23 def __init__(
24 self,
25 file_path: str,
26 api_mode: str = 'local',
27 api_url: str = 'http://localhost:8000',
28 api_key: str = '',
29 params: dict = None,
30 timeout: Optional[int] = 300,
31 max_markdown_bytes: Optional[int] = None,
32 ):
33 self.file_path = file_path
34 self.api_mode = api_mode.lower()
35 self.api_url = api_url.rstrip('/')
36 self.api_key = api_key
37 self.timeout = timeout
38 self.max_markdown_bytes = max_markdown_bytes
40 # Parse params dict with defaults
41 self.params = params or {}
42 self.enable_ocr = self.params.get('enable_ocr', False)
43 self.enable_formula = self.params.get('enable_formula', True)
44 self.enable_table = self.params.get('enable_table', True)
45 self.language = self.params.get('language', 'en')
46 self.model_version = self.params.get('model_version', 'pipeline')
48 self.page_ranges = self.params.pop('page_ranges', '')
50 # Validate API mode
51 if self.api_mode not in ['local', 'cloud']:
52 raise ValueError(f"Invalid API mode: {self.api_mode}. Must be 'local' or 'cloud'")
54 # Validate Cloud API requirements
55 if self.api_mode == 'cloud' and not self.api_key:
56 raise ValueError('API key is required for Cloud API mode')
58 def load(self) -> List[Document]:
59 """
60 Main entry point for loading and parsing the document.
61 Routes to Cloud or Local API based on api_mode.
62 """
63 try:
64 if self.api_mode == 'cloud':
65 return self._load_cloud_api()
66 else:
67 return self._load_local_api()
68 except Exception as e:
69 log.error(f'Error loading document with MinerU: {e}')
70 raise
72 def _load_local_api(self) -> List[Document]:
73 """
74 Load document using Local API (synchronous).
75 Posts file to /file_parse endpoint and gets immediate response.
76 """
77 log.info('Using MinerU Local API at %s', self.api_url)
79 filename = os.path.basename(self.file_path)
81 # Build form data for Local API
82 form_data = {
83 **self.params,
84 'return_md': 'true',
85 }
87 # Page ranges (Local API uses start_page_id and end_page_id)
88 if self.page_ranges:
89 # For simplicity, if page_ranges is specified, log a warning
90 # Full page range parsing would require parsing the string
91 log.warning(
92 f"Page ranges '{self.page_ranges}' specified but Local API uses different format. "
93 'Consider using start_page_id/end_page_id parameters if needed.'
94 )
96 try:
97 with open(self.file_path, 'rb') as f:
98 files = {'files': (filename, f, 'application/octet-stream')}
100 log.info('Sending file to MinerU Local API: %s', filename)
101 log.debug('Local API parameters: %s', form_data)
103 response = requests.post(
104 f'{self.api_url}/file_parse',
105 data=form_data,
106 files=files,
107 timeout=self.timeout,
108 )
109 response.raise_for_status()
111 except FileNotFoundError:
112 raise HTTPException(status.HTTP_404_NOT_FOUND, detail=f'File not found: {self.file_path}')
113 except requests.Timeout:
114 raise HTTPException(
115 status.HTTP_504_GATEWAY_TIMEOUT,
116 detail='MinerU Local API request timed out',
117 )
118 except requests.HTTPError as e:
119 error_detail = f'MinerU Local API request failed: {e}'
120 if e.response is not None:
121 try:
122 error_data = e.response.json()
123 error_detail += f' - {error_data}'
124 except Exception:
125 error_detail += f' - {e.response.text}'
126 raise HTTPException(status.HTTP_400_BAD_REQUEST, detail=error_detail)
127 except Exception as e:
128 raise HTTPException(
129 status.HTTP_500_INTERNAL_SERVER_ERROR,
130 detail=f'Error calling MinerU Local API: {str(e)}',
131 )
133 # Parse response
134 try:
135 result = response.json()
136 except ValueError as e:
137 raise HTTPException(
138 status.HTTP_502_BAD_GATEWAY,
139 detail=f'Invalid JSON response from MinerU Local API: {e}',
140 )
142 # Extract markdown content from response
143 if 'results' not in result:
144 raise HTTPException(
145 status.HTTP_502_BAD_GATEWAY,
146 detail="MinerU Local API response missing 'results' field",
147 )
149 results = result['results']
150 if not results:
151 raise HTTPException(
152 status.HTTP_400_BAD_REQUEST,
153 detail='MinerU returned empty results',
154 )
156 # Get the first (and typically only) result
157 file_result = list(results.values())[0]
158 markdown_content = file_result.get('md_content', '')
160 if not markdown_content:
161 raise HTTPException(
162 status.HTTP_400_BAD_REQUEST,
163 detail='MinerU returned empty markdown content',
164 )
166 log.info('Successfully parsed document with MinerU Local API: %s', filename)
168 # Create metadata
169 metadata = {
170 'source': filename,
171 'api_mode': 'local',
172 'backend': result.get('backend', 'unknown'),
173 'version': result.get('version', 'unknown'),
174 }
176 return [Document(page_content=markdown_content, metadata=metadata)]
178 def _load_cloud_api(self) -> List[Document]:
179 """
180 Load document using Cloud API (asynchronous).
181 Uses batch upload endpoint to avoid need for public file URLs.
182 """
183 log.info('Using MinerU Cloud API at %s', self.api_url)
185 filename = os.path.basename(self.file_path)
187 # Step 1: Request presigned upload URL
188 batch_id, upload_url = self._request_upload_url(filename)
190 # Step 2: Upload file to presigned URL
191 self._upload_to_presigned_url(upload_url)
193 # Step 3: Poll for results
194 result = self._poll_batch_status(batch_id, filename)
196 # Step 4: Download and extract markdown from ZIP
197 markdown_content = self._download_and_extract_zip(result['full_zip_url'], filename)
199 log.info('Successfully parsed document with MinerU Cloud API: %s', filename)
201 # Create metadata
202 metadata = {
203 'source': filename,
204 'api_mode': 'cloud',
205 'batch_id': batch_id,
206 }
208 return [Document(page_content=markdown_content, metadata=metadata)]
210 def _request_upload_url(self, filename: str) -> tuple:
211 """
212 Request presigned upload URL from Cloud API.
213 Returns (batch_id, upload_url).
214 """
215 headers = {
216 'Authorization': f'Bearer {self.api_key}',
217 'Content-Type': 'application/json',
218 }
220 # Build request body
221 request_body = {
222 **self.params,
223 'files': [
224 {
225 'name': filename,
226 'is_ocr': self.enable_ocr,
227 }
228 ],
229 }
231 # Add page ranges if specified
232 if self.page_ranges:
233 request_body['files'][0]['page_ranges'] = self.page_ranges
235 log.info('Requesting upload URL for: %s', filename)
236 log.debug('Cloud API request body: %s', request_body)
238 try:
239 response = requests.post(
240 f'{self.api_url}/file-urls/batch',
241 headers=headers,
242 json=request_body,
243 timeout=30,
244 )
245 response.raise_for_status()
246 except requests.HTTPError as e:
247 error_detail = f'Failed to request upload URL: {e}'
248 if e.response is not None:
249 try:
250 error_data = e.response.json()
251 error_detail += f' - {error_data.get("msg", error_data)}'
252 except Exception:
253 error_detail += f' - {e.response.text}'
254 raise HTTPException(status.HTTP_400_BAD_REQUEST, detail=error_detail)
255 except Exception as e:
256 raise HTTPException(
257 status.HTTP_500_INTERNAL_SERVER_ERROR,
258 detail=f'Error requesting upload URL: {str(e)}',
259 )
261 try:
262 result = response.json()
263 except ValueError as e:
264 raise HTTPException(
265 status.HTTP_502_BAD_GATEWAY,
266 detail=f'Invalid JSON response: {e}',
267 )
269 # Check for API error response
270 if result.get('code') != 0:
271 raise HTTPException(
272 status.HTTP_400_BAD_REQUEST,
273 detail=f'MinerU Cloud API error: {result.get("msg", "Unknown error")}',
274 )
276 data = result.get('data', {})
277 batch_id = data.get('batch_id')
278 file_urls = data.get('file_urls', [])
280 if not batch_id or not file_urls:
281 raise HTTPException(
282 status.HTTP_502_BAD_GATEWAY,
283 detail='MinerU Cloud API response missing batch_id or file_urls',
284 )
286 upload_url = file_urls[0]
287 log.info('Received upload URL for batch: %s', batch_id)
289 return batch_id, upload_url
291 def _upload_to_presigned_url(self, upload_url: str) -> None:
292 """
293 Upload file to presigned URL (no authentication needed).
294 """
295 log.info(f'Uploading file to presigned URL')
297 try:
298 with open(self.file_path, 'rb') as f:
299 response = requests.put(
300 upload_url,
301 data=f,
302 timeout=self.timeout,
303 )
304 response.raise_for_status()
305 except FileNotFoundError:
306 raise HTTPException(status.HTTP_404_NOT_FOUND, detail=f'File not found: {self.file_path}')
307 except requests.Timeout:
308 raise HTTPException(
309 status.HTTP_504_GATEWAY_TIMEOUT,
310 detail='File upload to presigned URL timed out',
311 )
312 except requests.HTTPError as e:
313 raise HTTPException(
314 status.HTTP_400_BAD_REQUEST,
315 detail=f'Failed to upload file to presigned URL: {e}',
316 )
317 except Exception as e:
318 raise HTTPException(
319 status.HTTP_500_INTERNAL_SERVER_ERROR,
320 detail=f'Error uploading file: {str(e)}',
321 )
323 log.info('File uploaded successfully')
325 def _poll_batch_status(self, batch_id: str, filename: str) -> dict:
326 """
327 Poll batch status until completion.
328 Returns the result dict for the file.
329 """
330 headers = {
331 'Authorization': f'Bearer {self.api_key}',
332 }
334 max_iterations = 300 # 10 minutes max (2 seconds per iteration)
335 poll_interval = 2 # seconds
337 log.info('Polling batch status: %s', batch_id)
339 for iteration in range(max_iterations):
340 try:
341 response = requests.get(
342 f'{self.api_url}/extract-results/batch/{batch_id}',
343 headers=headers,
344 timeout=30,
345 )
346 response.raise_for_status()
347 except requests.HTTPError as e:
348 error_detail = f'Failed to poll batch status: {e}'
349 if e.response is not None:
350 try:
351 error_data = e.response.json()
352 error_detail += f' - {error_data.get("msg", error_data)}'
353 except Exception:
354 error_detail += f' - {e.response.text}'
355 raise HTTPException(status.HTTP_400_BAD_REQUEST, detail=error_detail)
356 except Exception as e:
357 raise HTTPException(
358 status.HTTP_500_INTERNAL_SERVER_ERROR,
359 detail=f'Error polling batch status: {str(e)}',
360 )
362 try:
363 result = response.json()
364 except ValueError as e:
365 raise HTTPException(
366 status.HTTP_502_BAD_GATEWAY,
367 detail=f'Invalid JSON response while polling: {e}',
368 )
370 # Check for API error response
371 if result.get('code') != 0:
372 raise HTTPException(
373 status.HTTP_400_BAD_REQUEST,
374 detail=f'MinerU Cloud API error: {result.get("msg", "Unknown error")}',
375 )
377 data = result.get('data', {})
378 extract_result = data.get('extract_result', [])
380 # Find our file in the batch results
381 file_result = None
382 for item in extract_result:
383 if item.get('file_name') == filename:
384 file_result = item
385 break
387 if not file_result:
388 raise HTTPException(
389 status.HTTP_502_BAD_GATEWAY,
390 detail=f'File {filename} not found in batch results',
391 )
393 state = file_result.get('state')
395 if state == 'done':
396 log.info('Processing complete for %s', filename)
397 return file_result
398 elif state == 'failed':
399 error_msg = file_result.get('err_msg', 'Unknown error')
400 raise HTTPException(
401 status.HTTP_400_BAD_REQUEST,
402 detail=f'MinerU processing failed: {error_msg}',
403 )
404 elif state in ['waiting-file', 'pending', 'running', 'converting']:
405 # Still processing
406 if iteration % 10 == 0: # Log every 20 seconds
407 log.info('Processing status: %s (iteration %s/%s)', state, iteration + 1, max_iterations)
408 time.sleep(poll_interval)
409 else:
410 log.warning(f'Unknown state: {state}')
411 time.sleep(poll_interval)
413 # Timeout
414 raise HTTPException(
415 status.HTTP_504_GATEWAY_TIMEOUT,
416 detail='MinerU processing timed out after 10 minutes',
417 )
419 def _download_and_extract_zip(self, zip_url: str, filename: str) -> str:
420 """
421 Download ZIP file from CDN and extract markdown content.
422 Returns the markdown content as a string.
423 """
424 log.info('Downloading results from: %s', zip_url)
426 try:
427 response = requests.get(zip_url, timeout=60)
428 response.raise_for_status()
429 except requests.HTTPError as e:
430 raise HTTPException(
431 status.HTTP_400_BAD_REQUEST,
432 detail=f'Failed to download results ZIP: {e}',
433 )
434 except Exception as e:
435 raise HTTPException(
436 status.HTTP_500_INTERNAL_SERVER_ERROR,
437 detail=f'Error downloading results: {str(e)}',
438 )
440 # Save ZIP to temporary file before reading.
441 tmp_zip_path = None
442 markdown_content = None
443 try:
444 with tempfile.NamedTemporaryFile(delete=False, suffix='.zip') as tmp_zip:
445 tmp_zip.write(response.content)
446 tmp_zip_path = tmp_zip.name
448 with zipfile.ZipFile(tmp_zip_path, 'r') as zip_ref:
449 members = zip_ref.infolist()
450 all_files = [member.filename for member in members]
451 md_members = [member for member in members if member.filename.endswith('.md')]
452 read_errors = []
454 for member in md_members:
455 log.info('Found markdown file in ZIP: %s', member.filename)
456 try:
457 with zip_ref.open(member, 'r') as f:
458 if self.max_markdown_bytes is None:
459 content = f.read()
460 else:
461 content = f.read(self.max_markdown_bytes + 1)
462 if len(content) > self.max_markdown_bytes:
463 raise HTTPException(
464 status.HTTP_502_BAD_GATEWAY,
465 detail=f'Markdown file in results ZIP is too large: {member.filename}',
466 )
467 markdown_content = content.decode('utf-8')
468 except UnicodeDecodeError as e:
469 read_errors.append(f'{member.filename}: {e}')
470 log.warning(f'Failed to decode {member.filename}: {e}')
471 continue
472 except HTTPException:
473 raise
474 except Exception as e:
475 read_errors.append(f'{member.filename}: {e}')
476 log.warning(f'Failed to read {member.filename}: {e}')
477 continue
478 if markdown_content:
479 break
481 if markdown_content is None:
482 log.error(f'Available files in ZIP: {all_files}')
483 if read_errors:
484 error_msg = f"Found .md files but couldn't read them: {read_errors}"
485 else:
486 error_msg = f'No .md files found in ZIP. Available files: {all_files}'
487 raise HTTPException(
488 status.HTTP_502_BAD_GATEWAY,
489 detail=error_msg,
490 )
491 except zipfile.BadZipFile as e:
492 raise HTTPException(
493 status.HTTP_502_BAD_GATEWAY,
494 detail=f'Invalid ZIP file received: {e}',
495 )
496 except HTTPException:
497 raise
498 except Exception as e:
499 raise HTTPException(
500 status.HTTP_500_INTERNAL_SERVER_ERROR,
501 detail=f'Error extracting ZIP: {str(e)}',
502 )
503 finally:
504 if tmp_zip_path:
505 try:
506 os.unlink(tmp_zip_path)
507 except FileNotFoundError:
508 pass
509 except Exception as e:
510 log.warning(f'Failed to remove temporary ZIP file {tmp_zip_path}: {e}')
512 if not markdown_content:
513 raise HTTPException(
514 status.HTTP_400_BAD_REQUEST,
515 detail='Extracted markdown content is empty',
516 )
518 log.info('Successfully extracted markdown content (%s characters)', len(markdown_content))
519 return markdown_content