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

1import logging 

2import os 

3import tempfile 

4import time 

5import zipfile 

6from typing import List, Optional 

7 

8import requests 

9from fastapi import HTTPException, status 

10from langchain_core.documents import Document 

11 

12log = logging.getLogger(__name__) 

13 

14 

15class MinerULoader: 

16 """ 

17 MinerU document parser loader supporting both Cloud API and Local API modes. 

18 

19 Cloud API: Uses MinerU managed service with async task-based processing 

20 Local API: Uses self-hosted MinerU API with synchronous processing 

21 """ 

22 

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 

39 

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') 

47 

48 self.page_ranges = self.params.pop('page_ranges', '') 

49 

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'") 

53 

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') 

57 

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 

71 

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) 

78 

79 filename = os.path.basename(self.file_path) 

80 

81 # Build form data for Local API 

82 form_data = { 

83 **self.params, 

84 'return_md': 'true', 

85 } 

86 

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 ) 

95 

96 try: 

97 with open(self.file_path, 'rb') as f: 

98 files = {'files': (filename, f, 'application/octet-stream')} 

99 

100 log.info('Sending file to MinerU Local API: %s', filename) 

101 log.debug('Local API parameters: %s', form_data) 

102 

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() 

110 

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 ) 

132 

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 ) 

141 

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 ) 

148 

149 results = result['results'] 

150 if not results: 

151 raise HTTPException( 

152 status.HTTP_400_BAD_REQUEST, 

153 detail='MinerU returned empty results', 

154 ) 

155 

156 # Get the first (and typically only) result 

157 file_result = list(results.values())[0] 

158 markdown_content = file_result.get('md_content', '') 

159 

160 if not markdown_content: 

161 raise HTTPException( 

162 status.HTTP_400_BAD_REQUEST, 

163 detail='MinerU returned empty markdown content', 

164 ) 

165 

166 log.info('Successfully parsed document with MinerU Local API: %s', filename) 

167 

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 } 

175 

176 return [Document(page_content=markdown_content, metadata=metadata)] 

177 

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) 

184 

185 filename = os.path.basename(self.file_path) 

186 

187 # Step 1: Request presigned upload URL 

188 batch_id, upload_url = self._request_upload_url(filename) 

189 

190 # Step 2: Upload file to presigned URL 

191 self._upload_to_presigned_url(upload_url) 

192 

193 # Step 3: Poll for results 

194 result = self._poll_batch_status(batch_id, filename) 

195 

196 # Step 4: Download and extract markdown from ZIP 

197 markdown_content = self._download_and_extract_zip(result['full_zip_url'], filename) 

198 

199 log.info('Successfully parsed document with MinerU Cloud API: %s', filename) 

200 

201 # Create metadata 

202 metadata = { 

203 'source': filename, 

204 'api_mode': 'cloud', 

205 'batch_id': batch_id, 

206 } 

207 

208 return [Document(page_content=markdown_content, metadata=metadata)] 

209 

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 } 

219 

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 } 

230 

231 # Add page ranges if specified 

232 if self.page_ranges: 

233 request_body['files'][0]['page_ranges'] = self.page_ranges 

234 

235 log.info('Requesting upload URL for: %s', filename) 

236 log.debug('Cloud API request body: %s', request_body) 

237 

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 ) 

260 

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 ) 

268 

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 ) 

275 

276 data = result.get('data', {}) 

277 batch_id = data.get('batch_id') 

278 file_urls = data.get('file_urls', []) 

279 

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 ) 

285 

286 upload_url = file_urls[0] 

287 log.info('Received upload URL for batch: %s', batch_id) 

288 

289 return batch_id, upload_url 

290 

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') 

296 

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 ) 

322 

323 log.info('File uploaded successfully') 

324 

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 } 

333 

334 max_iterations = 300 # 10 minutes max (2 seconds per iteration) 

335 poll_interval = 2 # seconds 

336 

337 log.info('Polling batch status: %s', batch_id) 

338 

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 ) 

361 

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 ) 

369 

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 ) 

376 

377 data = result.get('data', {}) 

378 extract_result = data.get('extract_result', []) 

379 

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 

386 

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 ) 

392 

393 state = file_result.get('state') 

394 

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) 

412 

413 # Timeout 

414 raise HTTPException( 

415 status.HTTP_504_GATEWAY_TIMEOUT, 

416 detail='MinerU processing timed out after 10 minutes', 

417 ) 

418 

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) 

425 

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 ) 

439 

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 

447 

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 = [] 

453 

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 

480 

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

511 

512 if not markdown_content: 

513 raise HTTPException( 

514 status.HTTP_400_BAD_REQUEST, 

515 detail='Extracted markdown content is empty', 

516 ) 

517 

518 log.info('Successfully extracted markdown content (%s characters)', len(markdown_content)) 

519 return markdown_content