Files
ragflow/internal/parser/parser/pdf_parser_mineru.go

117 lines
3.1 KiB
Go
Raw Normal View History

package parser
import (
"fmt"
"ragflow/internal/common"
"strings"
"time"
models "ragflow/internal/entity/models"
)
const minerUPollTimeout = 30 * time.Second
const minerUPollInterval = 200 * time.Millisecond
func parsePDFWithMinerU(filename string, data []byte, parser *PDFParser) ParseResult {
if len(data) == 0 {
return emptyPDFResult(filename)
}
apiServer := strings.TrimSpace(parser.MinerUAPIServer)
if apiServer == "" {
apiServer = strings.TrimSpace(common.GetEnv(common.EnvMineruApiServer))
}
if apiServer == "" {
return ParseResult{Err: fmt.Errorf("parser: MinerU requires mineru_apiserver or MINERU_APISERVER")}
}
apiKey := parser.MinerUAPIKey
if strings.TrimSpace(apiKey) == "" {
apiKey = strings.TrimSpace(common.GetEnv(common.EnvMineruApiKey))
}
backend := strings.TrimSpace(parser.MinerUBackend)
if backend == "" {
backend = strings.TrimSpace(common.GetEnv(common.EnvMineruBackend))
}
if backend == "" {
backend = "pipeline"
}
timeout := parser.MinerUPollTimeout
if timeout <= 0 {
timeout = minerUPollTimeout
}
driver := models.NewMinerLocalUModel(
map[string]string{"default": apiServer},
models.URLSuffix{DocumentParse: "file_parse", Task: "tasks"},
)
apiConfig := &models.APIConfig{
BaseURL: &apiServer,
}
if apiKey != "" {
apiConfig.ApiKey = &apiKey
}
task, err := driver.ParseFile(&backend, data, nil, apiConfig, &models.ParseFileConfig{})
if err != nil {
return ParseResult{Err: fmt.Errorf("parser: MinerU submit: %w", err)}
}
content, err := pollMinerUTask(driver, task.TaskID, apiConfig, timeout)
if err != nil {
return ParseResult{Err: fmt.Errorf("parser: MinerU result: %w", err)}
}
pageCount := 0
if strings.TrimSpace(content) != "" {
pageCount = 1
}
return parseMinerUMarkdownResult(filename, content, parser.OutputFormat, pageCount)
}
func pollMinerUTask(driver *models.MinerULocalModel, taskID string, apiConfig *models.APIConfig, timeout time.Duration) (string, error) {
if timeout <= 0 {
timeout = minerUPollTimeout
}
deadline := time.Now().Add(timeout)
var lastErr error
for {
task, err := driver.ShowTask(taskID, apiConfig)
if err == nil {
for _, segment := range task.Segments {
if strings.TrimSpace(segment.Content) != "" {
return segment.Content, nil
}
}
lastErr = fmt.Errorf("empty MinerU task content")
} else {
lastErr = err
}
if time.Now().After(deadline) {
if lastErr == nil {
lastErr = fmt.Errorf("timed out waiting for MinerU task %s", taskID)
}
return "", lastErr
}
time.Sleep(minerUPollInterval)
}
}
func parseMinerUMarkdownResult(filename, markdown, outputFormat string, pageCount int) ParseResult {
fileMeta := pdfFileMeta(filename, pageCount)
switch strings.ToLower(strings.TrimSpace(outputFormat)) {
case "", "json":
feat(agent): Go ingestion pipeline progress mirroring and DeepDOC parser hardening (#16795) feat(ingestion): mirror Go pipeline progress into the document table; harden resume guards - pipeline: bind the owning document via WithDocumentID; after each TrackProgress event aggregate ingestion_task_log progress and mirror progress/run/progress_msg back into the document table, so GET /api/v1/datasets/{dataset_id}/documents reflects live Go pipeline progress without a bespoke endpoint. - canvas: extend the S3 resume guard to reject legacy no-op nodes (e.g. ExitLoop) so component_total equals the count of progress-reporting components and the aggregate percent can reach 100%. - runtime/canvas: route progress through TrackProgress; add interrupt test coverage (r3_interrupt_test.go). - dao/entity: add IngestionTask.DocumentID column and AggregateProgress support used by the mirror; IngestionTaskLog keeps a Checkpoint column alongside the progress fields. feat(deepdoc): cache DocAnalyzer inference results in Redis (1h TTL) - Redis-backed DocAnalyzerCache decorator over inference.Client; cache key = "ddoc:cache:<method>:" + sha256 of the JPEG-encoded image bytes (deterministic). - TTL = 1h; hits skip the inner HTTP call and return cached JSON; inner errors are not cached. refactor(deepdoc): align figure cropping with Python cropout + bounded page caches - CropSectionByDLA mirrors Python cropout: best-overlap DLA figure/equation region, fallback to section bbox per page, vertical concat on gray background. - sliding-window page-image cache bounds peak memory to the recent window instead of the whole PDF. - rename DLADebug -> DLARegions across parser/chunker/tests. refactor(parser): drop lib_type selector; align NewXxxParser with NewPDFParser - remove config["lib_type"] lookup and the libType param/field/switch from all nine constructors; surface the CGO-required error at ParseWithResult time instead of construction time; drop resolveLibType, its test, and the four lib_type constants. feat(utility): add a reusable workerpool for bounded concurrent execution - internal/utility/workerpool.go (+ tests). refactor: translate Chinese prose comments to English in non-harness Go files. chore: upgrade github.com/cloudwego/eino from v0.9.9 to v0.9.12.
2026-07-10 10:36:10 +08:00
mp := NewMarkdownParser()
res := mp.ParseWithResult(filename, []byte(markdown))
if res.Err != nil {
return res
}
res.File = fileMeta
return res
case "markdown":
return ParseResult{
OutputFormat: "markdown",
File: fileMeta,
Markdown: markdown,
}
default:
return ParseResult{Err: fmt.Errorf("parser: unsupported PDF output_format %q", outputFormat)}
}
}