mirror of
https://github.com/infiniflow/ragflow.git
synced 2026-08-15 05:04:27 +08:00
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.
241 lines
7.4 KiB
Go
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 ""
|
|
}
|