mirror of
https://github.com/infiniflow/ragflow.git
synced 2026-08-16 21:50:58 +08:00
## Problem During ingestion, `indexdoc.ProcessChunksForPipeline` stamped `ck["kb_id"]` on every chunk. This was both: - **a dead write** — `elasticsearch.InsertChunks` unconditionally overwrites the value with `datasetID` (`chunk.go:211`), so the producer's value never reached the index; - **the wrong shape** — it was emitted as `[]string`, while both engines actually need a single string. This is the `kb_id` slice of the ingestion -> engine schema leak tracked in #17371: ingestion was carrying index-physical schema knowledge it should not own. ## Fix Make the search engines the single owner of `kb_id` at the write boundary, and stop ingestion from emitting it: - **Elasticsearch** (`chunk.go:211`) already sets `docCopy["kb_id"] = datasetID` — unchanged. - **Infinity** (`chunk.go`) `InsertChunks` now stamps `insertChunks[i]["kb_id"] = datasetID` right after `transformChunkFields` (previously it only *read/normalized* the producer value, which forced ingestion to supply it). Both engines are now consistent. - `ProcessChunksForPipeline` no longer stamps `kb_id` and the now-leaky `kbID` parameter is removed. The same removal is propagated to `ProcessPipelineOutputForGolden` and the `compare_pipeline_golden` dev tool (its `-kb-id` flag is dropped). The stored `kb_id` value is byte-for-byte unchanged: `datasetID` passed to `InsertChunks` is `taskCtx.Doc.KbID`, i.e. the same id that was previously set on the producer chunk. ## Verification - `bash build.sh --test ./internal/ingestion/task/indexdoc/... ./internal/engine/infinity/...` — both green. - `internal/ingestion/task` has **two pre-existing** failures (`TestPipelineExecutor_Run_RealCanvasDSL_UsesGeneralPipeline`, `TestRunPipeline_RealPipelineOutput_ProducesIndexFields`) that assert `inserted chunk count = 1, want 2` — a parser/assertion mismatch (the Go parser merges the 2-paragraph fixture into 1 chunk). They are unrelated to this change, which never touches chunk counting. The `kb_id`-related test failure this change would otherwise introduce is fixed by updating the tests below. - Updated the pinning unit test: `TestProcessChunksForPipeline_SetsDocID` (formerly `...SetsDocIDAndKBID`) now asserts `kb_id` is **not** set by the producer. Removed the `kb_id` assertion and the now-dead `taskChunkFieldEqualsStr` helper from `pipeline_real_integration_test.go`. ## Scope This closes only the `kb_id` portion of #17371. The remaining index-physical fields (`docnm_kwd`, `create_timestamp_flt`, `page_num_int`/`top_int`/ `position_int`, etc.) are intentionally left for a follow-up (P2).
59 lines
1.8 KiB
Go
59 lines
1.8 KiB
Go
//
|
|
// 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 (
|
|
"time"
|
|
|
|
indexdoc "ragflow/internal/ingestion/task/indexdoc"
|
|
)
|
|
|
|
// GoldenCompareResult is the structured output used by the local golden tools.
|
|
type GoldenCompareResult struct {
|
|
NormalizedChunks []map[string]any `json:"normalized_chunks"`
|
|
ProcessedChunks []map[string]any `json:"processed_chunks"`
|
|
MergedMetadata map[string]any `json:"merged_metadata"`
|
|
}
|
|
|
|
// ProcessPipelineOutputForGolden replays the deterministic pipeline post-processing
|
|
// steps from a pipeline.run()-style output without embedding or external writes.
|
|
func ProcessPipelineOutputForGolden(
|
|
pipelineOutput map[string]any,
|
|
docID string,
|
|
docName string,
|
|
) (GoldenCompareResult, error) {
|
|
normalized := indexdoc.NormalizeChunks(pipelineOutput)
|
|
if normalized == nil {
|
|
normalized = []map[string]any{}
|
|
}
|
|
|
|
processed := indexdoc.DeepCopyChunks(normalized)
|
|
metadata, err := indexdoc.ProcessChunksForPipeline(processed, docID, docName, time.Now())
|
|
if err != nil {
|
|
return GoldenCompareResult{}, err
|
|
}
|
|
if metadata == nil {
|
|
metadata = map[string]any{}
|
|
}
|
|
|
|
return GoldenCompareResult{
|
|
NormalizedChunks: normalized,
|
|
ProcessedChunks: processed,
|
|
MergedMetadata: metadata,
|
|
}, nil
|
|
}
|