mirror of
https://github.com/infiniflow/ragflow.git
synced 2026-07-29 20:19:24 +08:00
## Summary
Continuation of the Python→Go ingestion pipeline migration (File →
Parser → Chunker → Extractor → Tokenizer). Fixes cover Parser, Chunker,
and Tokenizer gaps identified. Fix page number (0-indexed and 1-index
mixed before fix; use 1-indexed after fix) and chunk order issues.
### Parser
- **Slides TCADP (1.7):** `pptx_tcadp.go` + TCADP branch in
`pptx_parser.go`/`ppt_parser.go` — PowerPoint files now support
`parse_method="tcadp"` via the TCADP cloud service, matching the
spreadsheet-family TCADP pattern. PPT containers pass `"PPT"` as
fileType (not hardcoded `"PPTX"`).
- **Audio default output_format (2.11):** `defaultSetups()` audio
default changed from `"text"` to `"json"`, aligning with Python
`parser.py:232` and `AllowedOutputFormat["audio"]={"json"}`.
- **PDF VLM enhancement (1.1):** `maybeDispatchPDFVisionEnhancement` in
`pdf_vision_dispatch.go` enriches image/table items with IMAGE2TEXT
model descriptions after PDF parsing, mirroring Python
`enhance_media_sections_with_vision`. Semaphore fix: acquire before
goroutine start to prevent unbounded goroutine creation.
- **json family (2.3):** reclassified as Keep Go — `json_parser.go` is a
functional enhancement, not a parity gap.
- **page number:** changed from "mixed use of 1-indexed & 0-indexed" to
"1-indexed"
### Chunker
- **BULLET_PATTERN fallback (1.7):** 4th-level fallback in
`resolveTitleLevels` (`title.go`) detects bullet/numbered-list patterns
(Chinese legal, numbering, English) when outline + regex levels produce
only bodyLevel. Guarded by `allBodyLevel` to never override existing
structure.
- **Tag/One chunker fields (1.8):** `tag.go` sets `TopInt` from source
row index; `one.go` preserves `Positions`/`PDFPositions` from source
items. TSV multi-line RowNum fix: tracks `contentStart` for correct row
attribution.
- **Overlapped_percent normalization (2.6):**
`NormalizeOverlappedPercent` in `schema/chunker.go` mirrors Python
`common/float_utils.py:50-58` — accepts `[0,1)` fraction or `[0,90]`
percent, normalizes to canonical `[0,90]`.
- **Paragraph splitting (2.7):** aligned to Python flow `naive_merge` —
`CRLF` normalization, `splitKeepingDelimiter` preserves sentence
delimiters, single-section merge with token-budget-governed chunking.
- **chunk order:** sort by reading order
### Tokenizer
- **Phantom chunk filtering (Omission 2):** `isPhantomChunk` + filter
loop in `chunksFromTokenizerUpstream` skips zero-value ChunkDocs.
- **Batch size env var (Omission 3):** `embeddingBatchSize()` reads
`TOKENIZER_EMBEDDING_BATCH_SIZE`, defaults to 16.
- **Summary empty check (Diff 5):** `TrimSpace(s) != ""` → `s != ""`,
matching Python truthy check.
- **chunk_order_int all paths (Diff 8):** set unconditionally before
full_text/embedding branching.
- **Timeout default (Diff 10):** `600s` → `60s`, matching Python
`@timeout(60)`.
- **Small maxTokens truncation (Diff 14):** `truncateForEmbedding`
returns `""` when `maxTokens <= 10`, matching Python.
### Code review fixes
- Semaphore acquire moved before goroutine in `pdf_vision_dispatch.go`
(concurrency control)
- Context propagation in `pptx_tcadp.go` (cancellation support)
- Test resolver leak fix in `media_dispatch_test.go` (defer restore)
- Migration history comments removed per AGENTS.md
## Test plan
```
bash build.sh --test ./internal/parser/parser/... ./internal/ingestion/component/...
```
## Notes
- Migration diff tracking: `docs/migration_python_go_diff.md`
- Remaining gaps: Extractor component only (21 items)
178 lines
5.9 KiB
Go
178 lines
5.9 KiB
Go
//go:build integration
|
|
|
|
//
|
|
// Copyright 2026 The InfiniFlow Authors. All Rights Reserved.
|
|
//
|
|
// Licensed under the Apache License, Version 2.0 (the "License");
|
|
// you may not use this file except in compliance with the License.
|
|
// You may obtain a copy of the License at
|
|
//
|
|
// http://www.apache.org/licenses/LICENSE-2.0
|
|
//
|
|
// Unless required by applicable law or agreed to in writing, software
|
|
// distributed under the License is distributed on an "AS IS" BASIS,
|
|
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
|
// See the License for the specific language governing permissions and
|
|
// limitations under the License.
|
|
//
|
|
|
|
package task
|
|
|
|
import (
|
|
"context"
|
|
"encoding/json"
|
|
"testing"
|
|
|
|
"ragflow/internal/common"
|
|
"ragflow/internal/dao"
|
|
"ragflow/internal/entity"
|
|
"ragflow/internal/ingestion/testutil"
|
|
)
|
|
|
|
// TestRealProducerConsumer exercises the project's real producer and consumer code paths:
|
|
//
|
|
// Producer: document.go pattern — Create(IngestionTask) → PublishTask(NATS)
|
|
// Consumer: Ingestor.Start() core logic — calls each actual function in sequence
|
|
func TestRealProducerConsumer(t *testing.T) {
|
|
// ── 1. NATS (embedded in-process server) ──
|
|
natsEngine := testutil.SetupNatsEngine(t)
|
|
if err := natsEngine.InitConsumer("tasks.>"); err != nil {
|
|
t.Fatalf("InitConsumer: %v", err)
|
|
}
|
|
|
|
// Purge stale messages
|
|
for {
|
|
h, _ := natsEngine.GetMessages(1)
|
|
if len(h) == 0 {
|
|
break
|
|
}
|
|
h[0].Ack()
|
|
}
|
|
|
|
// ── 2. SQLite DB ──
|
|
db := testutil.SetupTestDB(t)
|
|
cleanup := testutil.ReplaceDBForTest(t, db)
|
|
defer cleanup()
|
|
|
|
db.Create(&entity.Tenant{ID: "t1", LLMID: "gpt-4", Status: testutil.StrPtr("1")})
|
|
db.Create(&entity.Knowledgebase{ID: "kb1", TenantID: "t1", EmbdID: "e1", Status: testutil.StrPtr("1"), ParserConfig: entity.JSONMap{}})
|
|
docName := "doc-real.pdf"
|
|
db.Create(&entity.Document{ID: "doc-real", KbID: "kb1", ParserID: "naive", ParserConfig: entity.JSONMap{}, Name: &docName})
|
|
|
|
// ── 3. Producer: Mirrors document.go:1062-1085 exactly ──
|
|
ingestionTask := &entity.IngestionTask{
|
|
ID: "ingest-task-1",
|
|
UserID: "u1",
|
|
DocumentID: "doc-real",
|
|
DatasetID: "kb1",
|
|
Status: common.CREATED,
|
|
}
|
|
created, err := dao.NewIngestionTaskDAO().Create(context.Background(), db, ingestionTask)
|
|
if err != nil {
|
|
t.Fatalf("Create: %v", err)
|
|
}
|
|
t.Logf("Producer: IngestionTask created id=%s status=%s", created.ID, created.Status)
|
|
|
|
taskMessage := common.TaskMessage{
|
|
TaskID: created.ID,
|
|
TaskType: common.TaskTypeIngestionTask,
|
|
}
|
|
payload, _ := json.Marshal(taskMessage)
|
|
if err := natsEngine.PublishTask("tasks.RAGFLOW", payload); err != nil {
|
|
t.Fatalf("PublishTask: %v", err)
|
|
}
|
|
t.Logf("Producer: Published %s", payload)
|
|
|
|
// ── 4. Consumer: Mirrors Ingestor.Start():131-189 exactly ──
|
|
handles, err := natsEngine.GetMessages(1)
|
|
if err != nil {
|
|
t.Fatalf("GetMessages: %v", err)
|
|
}
|
|
if len(handles) != 1 {
|
|
t.Fatalf("expected 1 message, got %d", len(handles))
|
|
}
|
|
taskHandle := handles[0]
|
|
taskMsg := taskHandle.GetMessage()
|
|
t.Logf("Consumer: Received TaskID=%s TaskType=%s", taskMsg.TaskID, taskMsg.TaskType)
|
|
|
|
// Mirrors Start():133 — type filter
|
|
if taskMsg.TaskType != common.TaskTypeIngestionTask {
|
|
taskHandle.Ack()
|
|
t.Fatalf("unexpected task type: %s", taskMsg.TaskType)
|
|
}
|
|
|
|
// Mirrors Start():142-143 — UpdateStatusIfCurrent
|
|
ingestionTaskDAO := dao.NewIngestionTaskDAO()
|
|
_, err = ingestionTaskDAO.UpdateStatusIfCurrent(context.Background(), db, taskMsg.TaskID, common.CREATED, common.RUNNING)
|
|
if err != nil {
|
|
t.Fatalf("UpdateStatusIfCurrent: %v", err)
|
|
}
|
|
task, err := ingestionTaskDAO.GetByID(context.Background(), db, taskMsg.TaskID)
|
|
if err != nil {
|
|
t.Fatalf("GetByID: %v", err)
|
|
}
|
|
if task == nil {
|
|
t.Logf("Consumer: task %s not found in ingestion_task table — skipped", taskMsg.TaskID)
|
|
taskHandle.Ack()
|
|
return
|
|
}
|
|
t.Logf("Consumer: UpdateStatusIfCurrent status=%s", task.Status)
|
|
|
|
// Mirrors Start():167-180 — status check
|
|
switch task.Status {
|
|
case common.COMPLETED, common.STOPPED, common.FAILED:
|
|
taskHandle.Ack()
|
|
t.Fatalf("task already terminal: %s", task.Status)
|
|
case common.STOPPING, common.CREATED:
|
|
t.Fatalf("unexpected status: %s", task.Status)
|
|
case common.RUNNING:
|
|
t.Logf("Consumer: task is RUNNING — dispatching to executeTask")
|
|
}
|
|
|
|
// ── 5. executeTask (our modified version) ──
|
|
// Set a pipeline ID so the handler can resolve the canvas.
|
|
if err := db.Model(&entity.Document{}).Where("id = ?", "doc-real").Update("pipeline_id", "pipeline-real").Error; err != nil {
|
|
t.Fatalf("set pipeline_id: %v", err)
|
|
}
|
|
|
|
tc, err := LoadFromIngestionTask(context.Background(), task)
|
|
if err != nil {
|
|
t.Fatalf("LoadFromIngestionTask: %v", err)
|
|
}
|
|
t.Logf("Consumer: Loaded Doc=%s Parser=%s KB=%s Tenant=%s",
|
|
tc.Doc.ID, tc.Doc.ParserID, tc.KB.ID, tc.Tenant.ID)
|
|
|
|
svc, err := NewPipelineExecutor(tc, tc.PipelineID, 0)
|
|
if err != nil {
|
|
t.Fatalf("NewPipelineExecutor: %v", err)
|
|
}
|
|
svc.WithLoadDSLFunc(func(ctx context.Context, canvasID string) (string, string, error) {
|
|
return `{"nodes":[{"id":"test","type":"parser"}],"edges":[]}`, canvasID, nil
|
|
})
|
|
svc.WithRunPipelineFunc(func(ctx context.Context, dsl string) (map[string]any, string, error) {
|
|
return nil, "", nil
|
|
})
|
|
svc.WithInsertFunc(func(ctx context.Context, chunks []map[string]any, baseName, datasetID string) ([]string, error) {
|
|
return nil, nil
|
|
})
|
|
if _, err := svc.Execute(tc.Ctx); err != nil {
|
|
t.Fatalf("Execute: %v", err)
|
|
}
|
|
t.Log("Consumer: PipelineExecutor.Execute() - OK")
|
|
|
|
// Mirrors executeTask — mark as completed
|
|
if _, err := ingestionTaskDAO.UpdateStatusIfCurrent(context.Background(), db, task.ID, common.RUNNING, common.COMPLETED); err != nil {
|
|
t.Fatalf("UpdateStatus: %v", err)
|
|
}
|
|
|
|
// Mirrors Start():135 — Ack
|
|
taskHandle.Ack()
|
|
|
|
// ── 6. Verify ──
|
|
final, _ := ingestionTaskDAO.GetByID(context.Background(), db, task.ID)
|
|
if final.Status != common.COMPLETED {
|
|
t.Errorf("final status = %s, want %s", final.Status, common.COMPLETED)
|
|
}
|
|
t.Logf("Final: IngestionTask status=%s ✅", final.Status)
|
|
}
|