diff --git a/internal/ingestion/task/constants.go b/internal/ingestion/task/constants.go index e0450e56f2..d7eebaa8b5 100644 --- a/internal/ingestion/task/constants.go +++ b/internal/ingestion/task/constants.go @@ -20,7 +20,4 @@ package task const ( // GRAPH_RAPTOR_FAKE_DOC_ID is the fake doc_id used for RAPTOR-generated chunks. GRAPH_RAPTOR_FAKE_DOC_ID = "graph_raptor_fake_doc" - - // EmbeddingTokenConsumptionKey is the key in pipeline output for embedding token count. - EmbeddingTokenConsumptionKey = "embedding_token_consumption" ) diff --git a/internal/ingestion/task/golden_compare.go b/internal/ingestion/task/golden_compare.go index dc2bdd4350..90252122cb 100644 --- a/internal/ingestion/task/golden_compare.go +++ b/internal/ingestion/task/golden_compare.go @@ -18,6 +18,8 @@ package task import ( "time" + + indexdoc "ragflow/internal/ingestion/task/indexdoc" ) // GoldenCompareResult is the structured output used by the local golden tools. @@ -35,13 +37,13 @@ func ProcessPipelineOutputForGolden( kbID string, docName string, ) (GoldenCompareResult, error) { - normalized := NormalizeChunks(pipelineOutput) + normalized := indexdoc.NormalizeChunks(pipelineOutput) if normalized == nil { normalized = []map[string]any{} } - processed := deepCopyChunks(normalized) - metadata, err := ProcessChunksForPipeline(processed, docID, kbID, docName, time.Now()) + processed := indexdoc.DeepCopyChunks(normalized) + metadata, err := indexdoc.ProcessChunksForPipeline(processed, docID, kbID, docName, time.Now()) if err != nil { return GoldenCompareResult{}, err } diff --git a/internal/ingestion/task/golden_compare_test.go b/internal/ingestion/task/golden_compare_test.go index eb2761bdc0..2ebc5a9004 100644 --- a/internal/ingestion/task/golden_compare_test.go +++ b/internal/ingestion/task/golden_compare_test.go @@ -3,6 +3,8 @@ package task import ( "testing" "time" + + indexdoc "ragflow/internal/ingestion/task/indexdoc" ) func TestProcessPipelineOutputForGolden_Markdown(t *testing.T) { @@ -44,7 +46,7 @@ func TestProcessChunksForPipeline_StableFields(t *testing.T) { }, } - meta, err := ProcessChunksForPipeline(chunks, "doc-1", "kb-1", "sample.md", now) + meta, err := indexdoc.ProcessChunksForPipeline(chunks, "doc-1", "kb-1", "sample.md", now) if err != nil { t.Fatalf("ProcessChunksForPipeline: %v", err) } diff --git a/internal/ingestion/task/indexdoc/constants.go b/internal/ingestion/task/indexdoc/constants.go new file mode 100644 index 0000000000..3111f6f6f6 --- /dev/null +++ b/internal/ingestion/task/indexdoc/constants.go @@ -0,0 +1,20 @@ +// +// 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 indexdoc + +// EmbeddingTokenConsumptionKey is the key in pipeline output for embedding token count. +const EmbeddingTokenConsumptionKey = "embedding_token_consumption" diff --git a/internal/ingestion/task/chunk_utils.go b/internal/ingestion/task/indexdoc/normalize.go similarity index 80% rename from internal/ingestion/task/chunk_utils.go rename to internal/ingestion/task/indexdoc/normalize.go index b79cbecc2e..cc678cf679 100644 --- a/internal/ingestion/task/chunk_utils.go +++ b/internal/ingestion/task/indexdoc/normalize.go @@ -14,7 +14,7 @@ // limitations under the License. // -package task +package indexdoc import ( "fmt" @@ -30,16 +30,16 @@ func NormalizeChunks(output map[string]any) []map[string]any { } if chunks, ok := output["chunks"].([]map[string]any); ok { - return deepCopyChunks(chunks) + return DeepCopyChunks(chunks) } if chunks, ok := toChunkMaps(output["chunks"]); ok { - return deepCopyChunks(chunks) + return DeepCopyChunks(chunks) } if json, ok := output["json"].([]map[string]any); ok { - return deepCopyChunks(json) + return DeepCopyChunks(json) } if json, ok := toChunkMaps(output["json"]); ok { - return deepCopyChunks(json) + return DeepCopyChunks(json) } if md, ok := output["markdown"].(string); ok && md != "" { return []map[string]any{{"text": md}} @@ -69,10 +69,16 @@ func toChunkMaps(v any) ([]map[string]any, bool) { return out, true } -// deepCopyChunks returns a deep copy of the chunk slice and each chunk map. -// Slice values (e.g. []float64 vectors) are fully copied, not shared. -// Mirrors Python: copy.deepcopy() -func deepCopyChunks(chunks []map[string]any) []map[string]any { +// DeepCopyChunks returns a copy of the chunk slice and each chunk map. +// The chunk maps themselves are freshly allocated, and the value types that +// the pipeline actually emits — []float64 (vectors), []int, and []string — are +// element-wise copied so callers cannot mutate the originals through them. +// Other value types (nested maps, [][]float64 positions, etc.) are shared by +// reference, not recursively deep-copied; positions are later flattened and +// copied independently by processChunkPositions. It is therefore NOT a full +// recursive deep copy (unlike Python's copy.deepcopy), only the copy needed +// for the post-processing pass over pipeline output. +func DeepCopyChunks(chunks []map[string]any) []map[string]any { if chunks == nil { return nil } diff --git a/internal/ingestion/task/chunk_utils_test.go b/internal/ingestion/task/indexdoc/normalize_test.go similarity index 99% rename from internal/ingestion/task/chunk_utils_test.go rename to internal/ingestion/task/indexdoc/normalize_test.go index 0698f09171..e27e9ab0cd 100644 --- a/internal/ingestion/task/chunk_utils_test.go +++ b/internal/ingestion/task/indexdoc/normalize_test.go @@ -1,4 +1,4 @@ -package task +package indexdoc import ( "testing" diff --git a/internal/ingestion/task/position.go b/internal/ingestion/task/indexdoc/position.go similarity index 99% rename from internal/ingestion/task/position.go rename to internal/ingestion/task/indexdoc/position.go index e9f53bf9ad..26ef0a65df 100644 --- a/internal/ingestion/task/position.go +++ b/internal/ingestion/task/indexdoc/position.go @@ -14,7 +14,7 @@ // limitations under the License. // -package task +package indexdoc // AddPositions adds position fields to a chunk map. // Input positions is a flat []float64 grouped as [pn, left, right, top, bottom] diff --git a/internal/ingestion/task/position_test.go b/internal/ingestion/task/indexdoc/position_test.go similarity index 99% rename from internal/ingestion/task/position_test.go rename to internal/ingestion/task/indexdoc/position_test.go index 2cb79159dc..50464b0d55 100644 --- a/internal/ingestion/task/position_test.go +++ b/internal/ingestion/task/indexdoc/position_test.go @@ -1,4 +1,4 @@ -package task +package indexdoc import ( "testing" diff --git a/internal/ingestion/task/chunk_process.go b/internal/ingestion/task/indexdoc/process.go similarity index 99% rename from internal/ingestion/task/chunk_process.go rename to internal/ingestion/task/indexdoc/process.go index d6c352a166..d5516e8d94 100644 --- a/internal/ingestion/task/chunk_process.go +++ b/internal/ingestion/task/indexdoc/process.go @@ -14,7 +14,7 @@ // limitations under the License. // -package task +package indexdoc import ( "fmt" diff --git a/internal/ingestion/task/chunk_process_test.go b/internal/ingestion/task/indexdoc/process_test.go similarity index 99% rename from internal/ingestion/task/chunk_process_test.go rename to internal/ingestion/task/indexdoc/process_test.go index 1f1e409777..61f1c0b2da 100644 --- a/internal/ingestion/task/chunk_process_test.go +++ b/internal/ingestion/task/indexdoc/process_test.go @@ -1,4 +1,4 @@ -package task +package indexdoc import ( "testing" diff --git a/internal/ingestion/task/pipeline_e2e_test.go b/internal/ingestion/task/pipeline_e2e_test.go index ac9cc690d8..b68970c335 100644 --- a/internal/ingestion/task/pipeline_e2e_test.go +++ b/internal/ingestion/task/pipeline_e2e_test.go @@ -30,6 +30,7 @@ import ( "ragflow/internal/engine" "ragflow/internal/engine/elasticsearch" "ragflow/internal/engine/infinity" + indexdoc "ragflow/internal/ingestion/task/indexdoc" "ragflow/internal/ingestion/testutil" "ragflow/internal/server" "ragflow/internal/service" @@ -247,7 +248,7 @@ func TestPipelineE2E_PipelineExecutor(t *testing.T) { "q_2_vec": []float64{0.3, 0.4}, // Pre-vectorized to skip embedding }, }, - EmbeddingTokenConsumptionKey: 100, + indexdoc.EmbeddingTokenConsumptionKey: 100, }, dsl, nil }) diff --git a/internal/ingestion/task/pipeline_executor.go b/internal/ingestion/task/pipeline_executor.go index f8a3fbd5b5..2d9a3168be 100644 --- a/internal/ingestion/task/pipeline_executor.go +++ b/internal/ingestion/task/pipeline_executor.go @@ -31,6 +31,7 @@ import ( "ragflow/internal/ingestion/component" "ragflow/internal/ingestion/knowledge_compile" pipelinepkg "ragflow/internal/ingestion/pipeline" + indexdoc "ragflow/internal/ingestion/task/indexdoc" "gorm.io/gorm" ) @@ -198,13 +199,13 @@ func (s *PipelineExecutor) Execute(ctx context.Context) (*PipelineResult, error) // performs no DB/index writes — the embedding vectors already computed by the // pipeline run are left on the chunks. This keeps debug runs side-effect free. func (s *PipelineExecutor) collectDebugOutput(ctx context.Context, pipelineOutput map[string]any, start time.Time) (*PipelineResult, error) { - chunks := NormalizeChunks(pipelineOutput) + chunks := indexdoc.NormalizeChunks(pipelineOutput) return &PipelineResult{ DocID: s.taskCtx.Doc.ID, KbID: s.taskCtx.Doc.KbID, Chunks: chunks, ChunkCount: countDistinctChunkIDs(chunks), - TokenConsumption: GetEmbeddingTokenConsumption(pipelineOutput), + TokenConsumption: indexdoc.GetEmbeddingTokenConsumption(pipelineOutput), Duration: time.Since(start).Seconds(), }, nil } @@ -217,13 +218,13 @@ func (s *PipelineExecutor) processOutput(ctx context.Context, pipelineOutput map return nil, err } - chunks := NormalizeChunks(pipelineOutput) + chunks := indexdoc.NormalizeChunks(pipelineOutput) if len(chunks) == 0 { return nil, nil } - embeddingTokenConsumption := GetEmbeddingTokenConsumption(pipelineOutput) - metadata, err := ProcessChunksForPipeline( + embeddingTokenConsumption := indexdoc.GetEmbeddingTokenConsumption(pipelineOutput) + metadata, err := indexdoc.ProcessChunksForPipeline( chunks, s.taskCtx.Doc.ID, s.taskCtx.Doc.KbID, @@ -234,7 +235,7 @@ func (s *PipelineExecutor) processOutput(ctx context.Context, pipelineOutput map return nil, err } - tableMeta := AggregateTableDocMetadata(chunks, map[string]interface{}(s.taskCtx.Doc.ParserConfig)) + tableMeta := indexdoc.AggregateTableDocMetadata(chunks, map[string]interface{}(s.taskCtx.Doc.ParserConfig)) if tableMeta != nil { if metadata == nil { metadata = make(map[string]any) diff --git a/internal/ingestion/task/pipeline_executor_test.go b/internal/ingestion/task/pipeline_executor_test.go index 1fce515841..76609ba8d4 100644 --- a/internal/ingestion/task/pipeline_executor_test.go +++ b/internal/ingestion/task/pipeline_executor_test.go @@ -14,6 +14,7 @@ import ( "ragflow/internal/dao" "ragflow/internal/entity" pipelinepkg "ragflow/internal/ingestion/pipeline" + indexdoc "ragflow/internal/ingestion/task/indexdoc" ) // ============================================================================= @@ -201,7 +202,7 @@ func TestKB_Doc_Tenant_Accessors(t *testing.T) { func TestPipelineExecutor_ProcessChunks_WrapsProcessChunksForPipeline(t *testing.T) { svc := mustNewPipelineExecutor(t, makeTaskCtx(), "flow-1", 0) chunks := []map[string]any{{"text": "hello world"}} - meta, err := ProcessChunksForPipeline(chunks, svc.taskCtx.Doc.ID, svc.taskCtx.Doc.KbID, *svc.taskCtx.Doc.Name, time.Now()) + meta, err := indexdoc.ProcessChunksForPipeline(chunks, svc.taskCtx.Doc.ID, svc.taskCtx.Doc.KbID, *svc.taskCtx.Doc.Name, time.Now()) if err != nil { t.Fatalf("ProcessChunksForPipeline: %v", err) }