Files
Zhichang Yu 0784bef5b0 Port dataset-level structure merge for timeline/graph/mindmap (#18201)
Ports Python dataset-level structure aggregation (timeline, graph, mindmap) to Go. Mindmap emits entity/relation rows and merges like graph. Adds dataset_merge guard, engine gate, resolveDatasetStructureKind, kind-required structure graph GET/DELETE API, per-index task-id fields.
2026-08-13 18:37:47 +08:00

241 lines
7.4 KiB
Go

// Package mindmap implements the "mindmap" variant of KnowledgeCompiler,
// mirroring Python's MindMapExtractor: the source chunks are packed into
// token-budget batches; each batch gets one LLM call (system = the rendered
// MIND_MAP_EXTRACTION_PROMPT, user = "Output:") whose markdown reply is
// parsed (dictify semantics), list-to-kv converted, merged across batches,
// and shaped into the {"id","children"} mind-map tree. The tree emits as one
// product per node with parent links.
//
// Per PORT_PLAN.md the Markdown source is the LLM's reply, NOT the source
// document Markdown, so the Parser component is not reused.
package mindmap
import (
"context"
"fmt"
"strings"
"ragflow/internal/ingestion/component/knowledge_compiler/common"
"ragflow/internal/utility"
)
// batchSubmitter fans out the batch extraction jobs on the process-wide
// knowledge-compilation pool. It is injected by the knowledge_compiler wiring
// (component.go) so every stage shares one vCPU-sized concurrency bound; when
// nil the batches run sequentially (the historic default).
var batchSubmitter func(ctx context.Context, jobs []func() error) error
// SetBatchSubmitter installs the shared-pool fan-out used by Run's extraction
// stage. Pass nil to revert to serial execution.
func SetBatchSubmitter(submit func(ctx context.Context, jobs []func() error) error) {
batchSubmitter = submit
}
// runBatches mirrors structure.runBatches: concurrent under the wired global
// compiler pool, or serial when no submitter is set. The first error is
// returned after all jobs settle; the global pool is never StopWait'd.
func runBatches(ctx context.Context, jobs []func() error) error {
if len(jobs) == 0 {
return nil
}
if batchSubmitter != nil {
return batchSubmitter(ctx, jobs)
}
for _, j := range jobs {
if err := j(); err != nil {
return err
}
}
return nil
}
// Run executes the mindmap variant.
func Run(ctx context.Context, deps common.Deps, param common.Param, inputs common.Inputs) (common.Outputs, error) {
docID := firstNonEmpty(inputs.DocID, deps.DatasetID)
if docID == "" {
docID = "unknown"
}
llmID := firstNonEmpty(param.LLMID, inputs.LLMID)
tenantID := deps.TenantID
if deps.Chat == nil {
return common.Outputs{}, fmt.Errorf("mindmap: chat model required")
}
sections := chunkTexts(inputs.Chunks)
// One LLM task per token-budget batch (mirrors __call__'s task fan-out).
batches := packSections(sections, deps.Tokenizer)
results := make([]utility.OMap, len(batches))
jobs := make([]func() error, 0, len(batches))
for i, text := range batches {
i, text := i, text
jobs = append(jobs, func() error {
resp, err := deps.Chat.Chat(ctx, common.ChatRequest{
LLMID: llmID,
SystemPrompt: renderPrompt(text),
UserPrompt: userMessage,
})
if err != nil {
return err
}
// Distinct slice index per batch → no cross-goroutine contention.
results[i] = utility.Todict(utility.Dictify(utility.StripFences(resp.Content)))
return nil
})
}
// The extraction batches are LLM-bounded, not CPU-bounded: run them on the
// shared global compiler pool (vCPU-sized) when a submitter is wired in,
// otherwise fall back to serial execution (historic default).
if err := runBatches(ctx, jobs); err != nil {
return common.Outputs{}, err
}
// Merge batch dicts in batch order (mirrors reduce(self._merge, res)) and
// shape the final tree. Python returns a bare root when nothing parsed.
var merged utility.OMap
if len(results) > 0 {
merged = results[0]
for _, r := range results[1:] {
merged = utility.MergeDicts(merged, r)
}
}
root := utility.ShapeTree(merged)
products := treeToProducts(tenantID, docID, root)
// Batched embedding of each node's content for downstream vector search.
if len(products) > 0 && deps.Embed != nil {
texts := make([]string, len(products))
for i, p := range products {
texts[i] = p.Content
}
vectors, err := deps.Embed.Encode(ctx, texts)
if err != nil {
return common.Outputs{}, err
}
for i := range products {
if i < len(vectors) {
products[i].Vector = vectors[i]
}
}
}
// Buffer every tree node in one slice; the component merges them into the
// upstream chunk stream (matching Python, which appends compiled units onto
// the chunk list).
out := common.Outputs{
Products: products,
}
return out, nil
}
// treeToProducts flattens the shaped mind-map tree into entity/relation Products
// so mindmap participates in dataset-level merge exactly like graph (plan §1.2,
// aligning with Python's dataset_structure_merger which merges
// knowledge_graph_kwd IN {entity,relation} rows for structure_mindmap too).
//
// Mapping:
// - each node (including the root) → an entity product (kind="entity",
// name = node id, type = "mindmap").
// - each parent→child edge → a relation product (kind="relation",
// from = parent id, to = child id, type = "related" — Python's default).
//
// The entity/relation discriminator is carried in Meta["kind"] so the consumer's
// mergeStructureDataset buckets entities by (name,type) and relations by
// (from,type,to), matching graph/timeline.
func treeToProducts(tenantID, docID string, root *utility.Node) []common.Product {
var out []common.Product
if root == nil || root.ID == "" {
return out
}
seen := map[string]bool{}
// Entity: root node.
out = append(out, common.Product{
ID: common.StableRowID(tenantID, docID, string(common.VariantMindmap), "entity", root.ID),
DocID: docID,
TenantID: tenantID,
Variant: common.VariantMindmap,
Content: root.ID,
Meta: map[string]any{
"kind": "entity",
"name": root.ID,
"entity_type": "mindmap",
"compile_kwd": "mindmap",
},
})
seen[root.ID] = true
type pending struct {
node *utility.Node
parent string
}
queue := []pending{{root, root.ID}}
for len(queue) > 0 {
p := queue[0]
queue = queue[1:]
for _, child := range p.node.Children {
if child.ID == "" {
continue
}
// Entity: child node (dedup by id so a DAG-shaped tree does not emit
// the same node twice).
if !seen[child.ID] {
seen[child.ID] = true
out = append(out, common.Product{
ID: common.StableRowID(tenantID, docID, string(common.VariantMindmap), "entity", child.ID),
DocID: docID,
TenantID: tenantID,
Variant: common.VariantMindmap,
Content: child.ID,
Meta: map[string]any{
"kind": "entity",
"name": child.ID,
"entity_type": "mindmap",
"compile_kwd": "mindmap",
},
})
}
// Relation: parent → child edge (type = "related", Python default).
out = append(out, common.Product{
ID: common.StableRowID(tenantID, docID, string(common.VariantMindmap), "relation", p.parent, child.ID),
DocID: docID,
TenantID: tenantID,
Variant: common.VariantMindmap,
Content: p.parent + " related " + child.ID,
Meta: map[string]any{
"kind": "relation",
"from": p.parent,
"to": child.ID,
"relation_type": "related",
"compile_kwd": "mindmap",
},
})
queue = append(queue, pending{child, child.ID})
}
}
return out
}
func chunkTexts(chunks []common.Chunk) []string {
var out []string
for _, c := range chunks {
t := firstNonEmpty(c.Text, c.Content)
if strings.TrimSpace(t) == "" {
continue
}
out = append(out, t)
}
return out
}
func firstNonEmpty(vals ...string) string {
for _, v := range vals {
if v != "" {
return v
}
}
return ""
}