mirror of
https://github.com/infiniflow/ragflow.git
synced 2026-08-02 05:47:31 +08:00
## Summary
Adds a side-effect-free DataFlow canvas **debug (dry-run) mode** plus a
**debug run log with a "View result" panel**, so a canvas can be
executed synchronously and inspected end-to-end (per-component progress
and parsed chunks) without persisting anything.
### Dry-run execution (inline parsed chunks)
- `task/debug.go`: `NewDebugTaskContext` builds an in-memory
`TaskContext` with **`KB.ID == ""` — the single debug signal used across
the ingestion pipeline**. A canvas debug run has no knowledgebase, and
production ingestion always supplies one, so `kb_id == ""` occurs ONLY
in debug mode. Components gate their own side effects on this signal
without any dedicated debug vocabulary (the former `CANVAS_DEBUG_DOC_ID`
marker constant is removed).
- `task/pipeline_executor.go`: `validateTaskContext` no longer requires
a KB when `KB.ID == ""` (debug); debug runs return `collectDebugOutput`
(chunks) instead of a no-op; uploaded bytes are delivered as
`inputs['binary']` for doc-less runs; `injectDebugPageCap` caps the
parser to the first pages for a fast preview via the production
`override_params` channel (Parser cpnID + family).
- `component/tokenizer.go`: `shouldHaveEmbedding` skips embedding when
`kb_id == ""` — the embedder is configured on the knowledgebase, so a
debug run has nothing to resolve against and stays side-effect free.
- `chunker/register.go`: chunk images are uploaded to MinIO only when a
KB is present (persist run). **In debug mode the raw image bytes are
intentionally dropped (`delete(ck, "image")`)** — the debug preview does
not render chunk images, and dropping the bytes keeps them out of memory
and out of the Redis-stored debug log. This is a deliberate trade-off,
not an oversight.
- `component/file.go`: pass through in-memory binary bytes, skipping
`doc_id` -> storage resolution.
- `handler/agent.go` + `agent_webhook.go`: detect `dataflow_canvas` and
run a sync debug returning chunks inline on the existing
chat/completions endpoint; reject DataFlow canvases from webhooks (fixes
the previously dead `== "DataFlow"` check; mirrors Python
`agent_api.py`).
- `parser_dispatch.go`: export `ParserFileFamily` for the executor's
page-cap injection.
### Debug run log + "View result"
Mirrors Python's debug-log contract so the front-end can replay each
component's progress and parsed output:
- `task/debug_log_sink.go`: a `DebugLogSink` records every component's
lifecycle into a `[{component_id, trace}]` array (each trace entry
carries `message`, `progress`, `timestamp`, `elapsed_time`). `Flush`
appends a terminal `END` marker whose first trace message is non-empty
so the front-end detects completion. On failure the END marker is
prefixed `[ERROR]` yet still carries the run, so the failure timeline
renders instead of being stuck empty. Timestamps and `elapsed_time` are
in seconds (matching the rest of the app).
- `task/debug_result_dsl.go`: `BuildDebugResultDSL` builds the `dsl` the
END marker carries — the Go analogue of Python's `Graph.__str__` +
END-marker `dsl` in `rag/flow/pipeline.py`. It combines the static DSL
structure (component_name / downstream / params / graph.nodes) with the
run output map (`output["state"][<id>]`) to emit, per component,
`obj.params.outputs[<format>].value` (chunks / text / json / html /
markdown) — the exact keys the front-end `dataflow-result` page reads to
render each step's parsed chunks. Raw embedding vectors (including the
dimension-scoped `q_<dim>_vec` keys) are stripped so the stored log
stays Python-scale.
- `task/pipeline_executor.go`: after the run, attach the built `dsl` to
the END marker via the `ResultSink` capability.
- `handler/agent.go`: `runCanvasPipelineDebug` generates a stable
`message_id` up-front and always flushes the log (success or failure);
`respondWithDebugResult` returns `message_id` in **both** the success
and the error envelope so the front-end can poll the log. The debug-log
endpoint `GET /agents/:id/logs/:message_id` serves the array.
- `web/src/pages/agent/hooks/use-run-dataflow.ts`: on a run failure,
also surface `message_id` via `setMessageId` so the log sheet renders
the failure timeline (the `[ERROR]` END marker is already written).
Guarded by `if (msgId)`, so it is a safe no-op when the back-end does
not return an id.
## Behavioral notes
- Debug parses only the first pages (`debugPageCapPages`) for a fast
preview; an explicit `pages` cap already present in the ParserConfig is
respected.
- Debug mode does not keep chunk images (see above) and does not compute
embeddings — it exercises parse + chunk only.
## Test plan
- Go: `debug_test.go`, `debug_log_sink_test.go` (trace pairing, END
marker, `[ERROR]` prefix, fractional-second timestamp/elapsed_time, size
caps, and a real-pipeline test asserting the END-marker `dsl` carries
non-empty per-component `params.outputs` with chunks),
`debug_result_dsl_test.go` (flat and real nested `output["state"]`
shapes, vector stripping, format priority),
`debug_pages_integration_test.go`, `pipeline_executor_persist_test.go`,
`handler/agent_pipeline_debug_test.go`, `handler/agent_logs_test.go`
(incl. `TestRunCanvasPipelineDebug_ErrorStillExposesMessageID` /
`TestRespondWithDebugResult_ErrorCarriesMessageID` locking `message_id`
on failure), plus updates to `agent_test.go` / `agent_webhook_test.go` /
`chunker/image_upload_test.go` / `tokenizer*.go`.
- `go build ./...` and `./build.sh --test` for affected packages.
🤖 Generated with [CodeBuddy Code](https://cnb.cool/codebuddy)
---------
Co-authored-by: CodeBuddy Code <noreply@tencent.com>
590 lines
21 KiB
Go
590 lines
21 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 (
|
|
"context"
|
|
"encoding/json"
|
|
"fmt"
|
|
"ragflow/internal/utility"
|
|
"strings"
|
|
"time"
|
|
|
|
"ragflow/internal/common"
|
|
"ragflow/internal/dao"
|
|
"ragflow/internal/engine"
|
|
"ragflow/internal/entity"
|
|
"ragflow/internal/ingestion/component"
|
|
"ragflow/internal/ingestion/knowledge_compile"
|
|
pipelinepkg "ragflow/internal/ingestion/pipeline"
|
|
|
|
"gorm.io/gorm"
|
|
)
|
|
|
|
// PipelineResult is the outcome of a pipeline run: chunks have been
|
|
// indexed, and these bookkeeping inputs remain for the caller to apply to
|
|
// document state (metadata merge + chunk/token counter bumps).
|
|
type PipelineResult struct {
|
|
DocID string
|
|
KbID string
|
|
Metadata map[string]any
|
|
Chunks []map[string]any // populated only in debug (dry-run) mode
|
|
ChunkCount int
|
|
TokenConsumption int
|
|
Duration float64 // pipeline wall-clock seconds
|
|
// MessageID is the polling key for the debug-run log. The front-end reads
|
|
// it from the run response and polls GET /agents/:id/logs/:message_id to
|
|
// render progress; it is empty for non-debug (persist) runs.
|
|
MessageID string
|
|
}
|
|
|
|
type PipelineExecutor struct {
|
|
taskCtx *TaskContext
|
|
canvasID string
|
|
docBulkSize int
|
|
|
|
indexWriter *chunkIndexWriter
|
|
logCreateFunc func(ctx context.Context, db *gorm.DB, log *entity.PipelineOperationLog) error
|
|
loadDSLFunc func(ctx context.Context, canvasID string) (string, string, error)
|
|
runPipelineFunc func(ctx context.Context, dsl string) (map[string]any, string, error)
|
|
progressSink pipelinepkg.ProgressSink
|
|
requireResume bool // when true, the pipeline run passes WithRequireResume
|
|
}
|
|
|
|
func validateTaskContext(taskCtx *TaskContext) error {
|
|
if taskCtx == nil {
|
|
return fmt.Errorf("pipeline executor: nil task context")
|
|
}
|
|
if taskCtx.Doc.ID == "" {
|
|
return fmt.Errorf("pipeline executor: empty document id")
|
|
}
|
|
// A debug (dry-run) context carries no knowledgebase (see
|
|
// TaskContext.IsDebug), so it must not be required to supply one.
|
|
if !taskCtx.IsDebug() && taskCtx.Doc.KbID == "" {
|
|
return fmt.Errorf("pipeline executor: empty document knowledgebase id")
|
|
}
|
|
if taskCtx.Doc.Name == nil || *taskCtx.Doc.Name == "" {
|
|
return fmt.Errorf("pipeline executor: empty document name")
|
|
}
|
|
if taskCtx.Tenant.ID == "" {
|
|
return fmt.Errorf("pipeline executor: empty tenant id")
|
|
}
|
|
return nil
|
|
}
|
|
|
|
func NewPipelineExecutor(
|
|
taskCtx *TaskContext,
|
|
canvasID string,
|
|
docBulkSize int,
|
|
) (*PipelineExecutor, error) {
|
|
if err := validateTaskContext(taskCtx); err != nil {
|
|
return nil, err
|
|
}
|
|
if strings.TrimSpace(canvasID) == "" {
|
|
return nil, fmt.Errorf("pipeline executor: empty canvas id")
|
|
}
|
|
svc := &PipelineExecutor{
|
|
taskCtx: taskCtx,
|
|
canvasID: canvasID,
|
|
docBulkSize: docBulkSize,
|
|
indexWriter: newChunkIndexWriter(
|
|
func(ctx context.Context, chunks []map[string]any, baseName string, datasetID string) ([]string, error) {
|
|
return engine.Get().InsertChunks(ctx, chunks, baseName, datasetID)
|
|
},
|
|
fmt.Sprintf("ragflow_%s", taskCtx.Tenant.ID),
|
|
taskCtx.Doc.KbID,
|
|
docBulkSize,
|
|
),
|
|
logCreateFunc: dao.NewPipelineOperationLogDAO().Create,
|
|
}
|
|
svc.loadDSLFunc = svc.loadDSLFromCanvas
|
|
svc.runPipelineFunc = svc.runPipelineWithDSL
|
|
return svc, nil
|
|
}
|
|
|
|
func (s *PipelineExecutor) WithInsertFunc(f InsertFunc) *PipelineExecutor {
|
|
s.indexWriter.insertFunc = f
|
|
return s
|
|
}
|
|
|
|
func (s *PipelineExecutor) WithLogCreateFunc(f func(ctx context.Context, db *gorm.DB, log *entity.PipelineOperationLog) error) *PipelineExecutor {
|
|
s.logCreateFunc = f
|
|
return s
|
|
}
|
|
|
|
func (s *PipelineExecutor) WithLoadDSLFunc(f func(ctx context.Context, canvasID string) (string, string, error)) *PipelineExecutor {
|
|
s.loadDSLFunc = f
|
|
return s
|
|
}
|
|
|
|
func (s *PipelineExecutor) WithRunPipelineFunc(f func(ctx context.Context, dsl string) (map[string]any, string, error)) *PipelineExecutor {
|
|
s.runPipelineFunc = f
|
|
return s
|
|
}
|
|
|
|
// WithProgressSink injects a sink that receives pipeline component progress
|
|
// events. The sink owns all document/ingestion_task_log persistence; when
|
|
// unset, the pipeline runs DB-independent (progress events are dropped).
|
|
func (s *PipelineExecutor) WithProgressSink(sink pipelinepkg.ProgressSink) *PipelineExecutor {
|
|
s.progressSink = sink
|
|
return s
|
|
}
|
|
|
|
// WithRequireResume makes the pipeline refuse to start when no checkpoint
|
|
// store is resolvable (Redis down or not configured). Production ingestion
|
|
// sets this; tests skip it so they can exercise runPlain without Redis.
|
|
func (s *PipelineExecutor) WithRequireResume() *PipelineExecutor {
|
|
s.requireResume = true
|
|
return s
|
|
}
|
|
|
|
func (s *PipelineExecutor) KB() *entity.Knowledgebase { return &s.taskCtx.KB }
|
|
func (s *PipelineExecutor) Doc() *entity.Document { return &s.taskCtx.Doc }
|
|
func (s *PipelineExecutor) Tenant() *entity.Tenant { return &s.taskCtx.Tenant }
|
|
|
|
func (s *PipelineExecutor) Execute(ctx context.Context) (*PipelineResult, error) {
|
|
start := time.Now()
|
|
if err := ctx.Err(); err != nil {
|
|
return nil, err
|
|
}
|
|
|
|
dsl, correctedID, err := s.loadDSLFunc(ctx, s.canvasID)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
if correctedID != "" {
|
|
s.canvasID = correctedID
|
|
}
|
|
|
|
pipelineOutput, pipelineDSL, err := s.runPipelineFunc(ctx, dsl)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
|
|
// A debug (dry-run) run produces no persistent side effect (no MinIO
|
|
// image upload, no index insert, no pipeline log); see TaskContext.IsDebug.
|
|
if s.taskCtx.IsDebug() {
|
|
return s.collectDebugOutput(ctx, pipelineOutput, start)
|
|
}
|
|
|
|
result, err := s.processOutput(ctx, pipelineOutput, start)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
|
|
if pipelineDSL != "" {
|
|
s.recordPipelineLog(ctx, dao.DB, s.taskCtx.Doc.ID, pipelineDSL, "done")
|
|
}
|
|
|
|
return result, nil
|
|
}
|
|
|
|
// collectDebugOutput builds a PipelineResult for a debug (dry-run) run.
|
|
// It surfaces the pipeline's chunks so a debug endpoint can render them, but
|
|
// 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)
|
|
return &PipelineResult{
|
|
DocID: s.taskCtx.Doc.ID,
|
|
KbID: s.taskCtx.Doc.KbID,
|
|
Chunks: chunks,
|
|
ChunkCount: countDistinctChunkIDs(chunks),
|
|
TokenConsumption: GetEmbeddingTokenConsumption(pipelineOutput),
|
|
Duration: time.Since(start).Seconds(),
|
|
}, nil
|
|
}
|
|
|
|
func (s *PipelineExecutor) processOutput(ctx context.Context, pipelineOutput map[string]any, start time.Time) (*PipelineResult, error) {
|
|
if pipelineOutput == nil {
|
|
return nil, nil
|
|
}
|
|
if err := ctx.Err(); err != nil {
|
|
return nil, err
|
|
}
|
|
|
|
chunks := NormalizeChunks(pipelineOutput)
|
|
if len(chunks) == 0 {
|
|
return nil, nil
|
|
}
|
|
|
|
embeddingTokenConsumption := GetEmbeddingTokenConsumption(pipelineOutput)
|
|
metadata, err := ProcessChunksForPipeline(
|
|
chunks,
|
|
s.taskCtx.Doc.ID,
|
|
s.taskCtx.Doc.KbID,
|
|
*s.taskCtx.Doc.Name,
|
|
time.Now(),
|
|
)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
|
|
tableMeta := AggregateTableDocMetadata(chunks, map[string]interface{}(s.taskCtx.Doc.ParserConfig))
|
|
if tableMeta != nil {
|
|
if metadata == nil {
|
|
metadata = make(map[string]any)
|
|
}
|
|
for k, v := range tableMeta {
|
|
if _, exists := metadata[k]; !exists {
|
|
metadata[k] = v
|
|
}
|
|
}
|
|
}
|
|
|
|
// Per-document compiled knowledge products (those emitted by the
|
|
// KnowledgeCompiler component and stamped with `compile_kwd`) are persisted
|
|
// as available_int=0: invisible to the normal retriever until the dataset-level
|
|
// post-processing consumer (§11) merges them into available_int=1 products.
|
|
// Ordinary source chunks stay available_int=1 (the index default).
|
|
markCompiledProductsHidden(chunks)
|
|
|
|
if err := s.indexWriter.Write(ctx, chunks); err != nil {
|
|
return nil, err
|
|
}
|
|
|
|
// All chunks are now persisted. Notify the dataset-level post-processing consumer
|
|
// (§11) that this document is complete: its compiled products were written
|
|
// available_int=0 and the consumer later merges them into dataset-level products
|
|
// (available_int=1). The notification is sent only after a successful persist
|
|
// and is best-effort / non-fatal — a delivery failure is logged but does not
|
|
// fail the pipeline task.
|
|
if err := knowledge_compile.PublishCompleted(ctx, s.taskCtx.Tenant.ID, s.taskCtx.Doc.KbID, s.taskCtx.Doc.ID, 0); err != nil {
|
|
common.Logger.Warn(fmt.Sprintf("knowledge_compile: publish doc_completed for %s failed: %v", s.taskCtx.Doc.ID, err))
|
|
}
|
|
|
|
chunkCount := countDistinctChunkIDs(chunks)
|
|
|
|
return &PipelineResult{
|
|
DocID: s.taskCtx.Doc.ID,
|
|
KbID: s.taskCtx.Doc.KbID,
|
|
Metadata: metadata,
|
|
ChunkCount: chunkCount,
|
|
TokenConsumption: embeddingTokenConsumption,
|
|
Duration: time.Since(start).Seconds(),
|
|
}, nil
|
|
}
|
|
|
|
// countDistinctChunkIDs returns the number of distinct chunk IDs in the slice.
|
|
// After ProcessChunksForPipeline, every chunk carries an "id" field computed
|
|
// from xxhash(text+docID). Chunks with identical text share the same id and
|
|
// the search engine treats them as upserts, so the effective stored chunk count
|
|
// is the number of unique ids — not len(chunks). Mirrors the index-side
|
|
// deduplication that happens at write time.
|
|
func countDistinctChunkIDs(chunks []map[string]any) int {
|
|
seen := make(map[string]struct{}, len(chunks))
|
|
for _, ck := range chunks {
|
|
id, _ := ck["id"].(string)
|
|
if id == "" {
|
|
continue
|
|
}
|
|
seen[id] = struct{}{}
|
|
}
|
|
return len(seen)
|
|
}
|
|
|
|
// markCompiledProductsHidden sets available_int=0 on the per-document compiled
|
|
// knowledge products so they are hidden from the retriever until the dataset-level
|
|
// post-processing consumer merges them into available_int=1 products (§11). A
|
|
// chunk is a compiled product iff it carries the compile_kwd discriminator the
|
|
// KnowledgeCompiler component stamps; ordinary source chunks (no compile_kwd)
|
|
// keep the index default available_int=1 and remain immediately searchable.
|
|
// Merged dataset-level products are written by the consumer, never here, so they are
|
|
// never double-marked.
|
|
func markCompiledProductsHidden(chunks []map[string]any) {
|
|
for _, ck := range chunks {
|
|
if _, ok := ck["compile_kwd"]; !ok {
|
|
continue
|
|
}
|
|
ck["available_int"] = 0
|
|
}
|
|
}
|
|
|
|
func (s *PipelineExecutor) recordPipelineLog(ctx context.Context, db *gorm.DB, docID, dsl, status string) {
|
|
var dslMap entity.JSONMap
|
|
if err := json.Unmarshal([]byte(dsl), &dslMap); err != nil {
|
|
dslMap = entity.JSONMap{"raw": dsl}
|
|
}
|
|
log := &entity.PipelineOperationLog{
|
|
ID: utility.GenerateUUID(),
|
|
TenantID: s.Tenant().ID,
|
|
KbID: s.KB().ID,
|
|
DocumentID: docID,
|
|
PipelineID: &s.canvasID,
|
|
TaskType: string(entity.PipelineTaskTypeParse),
|
|
DSL: dslMap,
|
|
ParserID: s.taskCtx.Doc.ParserID,
|
|
DocumentName: *s.Doc().Name,
|
|
DocumentSuffix: s.taskCtx.Doc.Suffix,
|
|
DocumentType: s.taskCtx.Doc.Type,
|
|
SourceFrom: s.taskCtx.Doc.SourceType,
|
|
OperationStatus: status,
|
|
}
|
|
if err := s.logCreateFunc(ctx, db, log); err != nil {
|
|
common.Warn(fmt.Sprintf("failed to record pipeline log: %v", err))
|
|
}
|
|
}
|
|
|
|
func (s *PipelineExecutor) loadDSLFromCanvas(ctx context.Context, canvasID string) (string, string, error) {
|
|
if s == nil || s.taskCtx == nil {
|
|
return "", "", fmt.Errorf("pipeline executor: nil task context")
|
|
}
|
|
if canvasID == "" {
|
|
return "", "", fmt.Errorf("pipeline executor: empty canvas id")
|
|
}
|
|
canvas, err := dao.NewUserCanvasDAO().GetByID(ctx, dao.DB, canvasID)
|
|
if err != nil {
|
|
return "", "", fmt.Errorf("load canvas %s: %w", canvasID, err)
|
|
}
|
|
|
|
canvasTitle := ""
|
|
if canvas.Title != nil {
|
|
canvasTitle = *canvas.Title
|
|
}
|
|
common.Info(fmt.Sprintf("load canvas %s, name %s", canvasID, canvasTitle))
|
|
|
|
raw, err := json.Marshal(canvas.DSL)
|
|
if err != nil {
|
|
return "", "", fmt.Errorf("marshal canvas dsl %s: %w", canvasID, err)
|
|
}
|
|
return string(raw), canvasID, nil
|
|
}
|
|
|
|
// warnUnknownComponentParams logs a warning for any component id in the
|
|
// parserConfig whose id is absent from the pipeline DSL. The runtime merge
|
|
// (component params -> override_params) silently drops such entries, so we
|
|
// surface them here for operability. API-side validation
|
|
// already rejects unknown ids on write; this is purely a defensive guard
|
|
// for legacy/stale rows.
|
|
func warnUnknownComponentParams(dsl string, parserConfig map[string]any) {
|
|
if len(parserConfig) == 0 {
|
|
return
|
|
}
|
|
schemas, err := pipelinepkg.ExtractAllComponentParams([]byte(dsl))
|
|
if err != nil {
|
|
common.Warn(fmt.Sprintf("warnUnknownComponentParams: cannot parse DSL to validate component params: %v", err))
|
|
return
|
|
}
|
|
dslCPNs := make(map[string]struct{}, len(schemas))
|
|
for _, s := range schemas {
|
|
dslCPNs[s.CpnID] = struct{}{}
|
|
}
|
|
for cpnID := range parserConfig {
|
|
if _, ok := dslCPNs[cpnID]; !ok {
|
|
common.Warn(fmt.Sprintf(
|
|
"parser_config references cpnID %q not present in the pipeline DSL; it will be ignored at runtime", cpnID))
|
|
}
|
|
}
|
|
}
|
|
|
|
func (s *PipelineExecutor) runPipelineWithDSL(ctx context.Context, dsl string) (map[string]any, string, error) {
|
|
if s == nil || s.taskCtx == nil {
|
|
return nil, dsl, fmt.Errorf("pipeline executor: nil task context")
|
|
}
|
|
|
|
parserConfig := map[string]interface{}(s.taskCtx.Doc.ParserConfig)
|
|
if parserConfig == nil {
|
|
// Debug (dataflow dry-run) contexts intentionally carry no
|
|
// ParserConfig; start from an empty map so the debug page cap can be
|
|
// injected in place below without a nil-map assignment panic.
|
|
parserConfig = map[string]interface{}{}
|
|
}
|
|
common.InjectExtractorLLMID(parserConfig, s.taskCtx.Tenant.LLMID)
|
|
// When the dataset enables auto-metadata, ensure the Extractor node(s)
|
|
// carry the enable_metadata mode + field schema so the LLM extraction fires
|
|
// (mirrors Python task_executor.py:519 enabling gen_metadata_task). The
|
|
// dataset flag is authoritative: a node that already has enable_metadata
|
|
// turned on keeps its own config, but a shipped DSL defaulting it to 0 is
|
|
// still overridden so auto-metadata can activate.
|
|
common.InjectExtractorEnableMetadata(parserConfig)
|
|
|
|
// Surface component params whose cpnID is absent from the DSL. The
|
|
// runtime merge (override_params) silently drops such entries;
|
|
// API-side validation already rejects unknown ids on write, so this is a
|
|
// defensive guard for legacy/stale rows.
|
|
warnUnknownComponentParams(dsl, parserConfig)
|
|
|
|
pipelineID := "pipeline_" + s.taskCtx.Doc.ID
|
|
if s.taskCtx.IngestionTask != nil && s.taskCtx.IngestionTask.ID != "" {
|
|
pipelineID = s.taskCtx.IngestionTask.ID
|
|
}
|
|
pipe, err := pipelinepkg.NewPipelineFromDSL([]byte(dsl), pipelineID,
|
|
pipelinepkg.WithProgressSink(s.progressSink),
|
|
pipelinepkg.WithDocumentID(s.taskCtx.Doc.ID))
|
|
if err != nil {
|
|
return nil, dsl, fmt.Errorf("compile pipeline dsl: %w", err)
|
|
}
|
|
inputs := map[string]any{}
|
|
if s.taskCtx.Doc.ID != "" {
|
|
inputs["doc_id"] = s.taskCtx.Doc.ID
|
|
}
|
|
// Run-level metadata shared by both persist and debug (dataflow
|
|
// dry-run) runs. In debug the KB is absent (NewDebugTaskContext forces
|
|
// KB.ID == ""); in persist it is the document's own KB. Either way the
|
|
// Tokenizer reads kb_id from CanvasState.Globals.
|
|
inputs["tenant_id"] = s.taskCtx.Tenant.ID
|
|
inputs["kb_id"] = s.taskCtx.KB.ID
|
|
if s.taskCtx.KB.Language != nil {
|
|
inputs["lang"] = *s.taskCtx.KB.Language
|
|
}
|
|
|
|
// File delivery and doc metadata differ between the two run modes.
|
|
debug := s.taskCtx.IsDebug()
|
|
if debug {
|
|
// A debug (dry-run) run has no DB document row, so the parser
|
|
// cannot resolve its bytes via doc_id → storage. Deliver the
|
|
// uploaded bytes directly as `binary` (what the parser actually
|
|
// reads) and surface the doc name/type so family detection works.
|
|
if s.taskCtx.File != nil {
|
|
inputs["file"] = s.taskCtx.File
|
|
inputs["binary"] = s.taskCtx.File
|
|
}
|
|
if s.taskCtx.Doc.Name != nil && *s.taskCtx.Doc.Name != "" {
|
|
inputs["name"] = *s.taskCtx.Doc.Name
|
|
}
|
|
if s.taskCtx.Doc.Type != "" {
|
|
inputs["file_type"] = s.taskCtx.Doc.Type
|
|
}
|
|
} else {
|
|
if s.taskCtx.File != nil {
|
|
inputs["file"] = s.taskCtx.File
|
|
}
|
|
}
|
|
|
|
// A canvas-debug (dataflow dry-run) must return a fast preview, so it
|
|
// caps the parser to the first few pages. The cap is delivered through
|
|
// override_params (Run's 3rd argument) — the SAME channel the
|
|
// production ParserConfig uses — keyed by the Parser component's cpnID
|
|
// and the document's filetype family. It is NOT passed through pipeline
|
|
// inputs: the parser selects pages from ParserConfig[cpnID][family]
|
|
// ["pages"] (a list of 1-indexed inclusive ranges), exactly mirroring
|
|
// NormalizeParserConfigPages / pdf_pages_test.go. See injectDebugPageCap.
|
|
if debug {
|
|
injectDebugPageCap(dsl, parserConfig, s.taskCtx.Doc.Type)
|
|
}
|
|
|
|
// Component params from Doc.ParserConfig — including the tenant LLM id
|
|
// injected into Extractor components above — are passed to Run as
|
|
// override_params, keyed by cpnID with override-wins. The DSL itself is
|
|
// compiled unchanged.
|
|
output, err := pipe.Run(ctx, inputs, parserConfig)
|
|
if err != nil {
|
|
return nil, dsl, err
|
|
}
|
|
|
|
// Surface the debug-run result DSL to any sink that implements ResultSink
|
|
// (the DebugLogSink used by canvas-debug runs). This mirrors Python's
|
|
// END-marker `dsl` attachment (rag/flow/pipeline.py:98) so the front-end
|
|
// "View result" page can render parsed chunks. The probe is an optional
|
|
// capability: non-debug (DB-backed) sinks ignore it and the ProgressSink
|
|
// contract is unchanged, keeping the coupling one-directional.
|
|
if rs, ok := s.progressSink.(ResultSink); ok {
|
|
if resultDSL, e := BuildDebugResultDSL(dsl, output); e == nil {
|
|
rs.SetResult(resultDSL)
|
|
}
|
|
}
|
|
|
|
payload, err := pipelinepkg.ExtractPayload(dsl, output)
|
|
if err != nil {
|
|
return nil, dsl, err
|
|
}
|
|
return payload, dsl, nil
|
|
}
|
|
|
|
// debugPageCapPages is the number of leading pages a canvas-debug
|
|
// (dataflow dry-run) parses. The debug preview must return fast, so we cap
|
|
// the parser to the first few pages. The cap is expressed as the 1-indexed
|
|
// inclusive range [1, debugPageCapPages], matching the production
|
|
// ParserConfig[cpnID][filetype]["pages"] shape (see NormalizeParserConfigPages).
|
|
const debugPageCapPages = 2
|
|
|
|
// injectDebugPageCap wires the canvas-debug page cap into parserConfig using
|
|
// the SAME channel the production ParserConfig travels through: Run's
|
|
// override_params (the 3rd argument), keyed by the Parser component's cpnID
|
|
// and the document's filetype family. It must NOT be passed as a pipeline
|
|
// input — the parser selects pages from ParserConfig[cpnID][family]["pages"],
|
|
// which the deepdoc/pdf parser consumes as a list of page ranges.
|
|
//
|
|
// dsl is the raw pipeline DSL (used only to discover the Parser component's
|
|
// cpnID); docType is the uploaded file's extension (e.g. "pdf"). When no
|
|
// Parser component can be found, or the docType yields no known family, the
|
|
// call is a no-op (the run parses everything, which is safe).
|
|
//
|
|
// dsl arrives as the canvas ENVELOPE ({ "dsl": { "components": ... } }) in
|
|
// production (loadDSLFromCanvas marshals canvas.DSL), so it is unwrapped the
|
|
// same way NewPipelineFromDSL does before ExtractAllComponentParams runs.
|
|
// An explicit page cap already present in parserConfig (keyed by cpnID +
|
|
// family) is respected and left untouched — the debug default is only a
|
|
// fallback, so a debug run can honour a narrower or wider caller-supplied cap.
|
|
func injectDebugPageCap(dsl string, parserConfig map[string]any, docType string) {
|
|
// Unwrap the canvas envelope to the inner components map.
|
|
var raw map[string]any
|
|
if err := json.Unmarshal([]byte(dsl), &raw); err != nil {
|
|
return
|
|
}
|
|
if env, ok := raw["dsl"].(map[string]any); ok && len(env) > 0 {
|
|
raw = env
|
|
}
|
|
inner, err := json.Marshal(raw)
|
|
if err != nil {
|
|
return
|
|
}
|
|
schemas, err := pipelinepkg.ExtractAllComponentParams(inner)
|
|
if err != nil {
|
|
return
|
|
}
|
|
var parserCpnID string
|
|
for _, s := range schemas {
|
|
if s.ComponentName == component.ComponentNameParser {
|
|
parserCpnID = s.CpnID
|
|
break
|
|
}
|
|
}
|
|
if parserCpnID == "" {
|
|
return
|
|
}
|
|
family := component.ParserFileFamily(docType)
|
|
if family == "" {
|
|
return
|
|
}
|
|
// Respect an explicit page cap already present under cpnID + family.
|
|
if cpnEntry, ok := parserConfig[parserCpnID].(map[string]any); ok {
|
|
if famEntry, ok := cpnEntry[family].(map[string]any); ok {
|
|
if _, has := famEntry["pages"]; has {
|
|
return
|
|
}
|
|
}
|
|
}
|
|
cpnEntry, ok := parserConfig[parserCpnID].(map[string]any)
|
|
if !ok {
|
|
cpnEntry = map[string]any{}
|
|
parserConfig[parserCpnID] = cpnEntry
|
|
}
|
|
famEntry, ok := cpnEntry[family].(map[string]any)
|
|
if !ok {
|
|
famEntry = map[string]any{}
|
|
cpnEntry[family] = famEntry
|
|
}
|
|
// pages is delivered as a JSON-decoded map — the very shape a ParserConfig
|
|
// arrives in from the API/storage JSON round-trip: a []any of [from,to]
|
|
// pairs. The deepdoc/pdf parser's NormalizePDFPages requires this
|
|
// []any-of-[]any form (not a Go [][]int), so it is built explicitly here.
|
|
// The shallow override_params merge in applyOverrideParams preserves this
|
|
// shape all the way to ConfigureFromSetup, so the cap is honoured.
|
|
famEntry["pages"] = []any{[]any{1, debugPageCapPages}}
|
|
}
|