mirror of
https://github.com/infiniflow/ragflow.git
synced 2026-07-26 02:13:29 +08:00
Implement Create/Drop Index/Metadata index in GO (#13791)
### What problem does this PR solve? Implement Create/Drop Index/Metadata index in GO New API handling in GO: POST/kb/index DELETE /kb/index POST /tenant/doc_meta_index DELETE /tenant/doc_meta_index CREATE INDEX FOR DATASET 'dataset_name' VECTOR_SIZE 1024; DROP INDEX FOR DATASET 'dataset_name'; CREATE INDEX DOC_META; DROP INDEX DOC_META; ### Type of change - [x] Refactoring
This commit is contained in:
@@ -19,11 +19,14 @@ package infinity
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
"ragflow/internal/server"
|
||||
"reflect"
|
||||
"strconv"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
infinity "github.com/infiniflow/infinity-go-sdk"
|
||||
"ragflow/internal/server"
|
||||
"ragflow/internal/logger"
|
||||
)
|
||||
|
||||
// infinityClient Infinity SDK client wrapper
|
||||
@@ -48,21 +51,80 @@ func NewInfinityClient(cfg *server.InfinityConfig) (*infinityClient, error) {
|
||||
}
|
||||
}
|
||||
|
||||
conn, err := infinity.Connect(infinity.NetworkAddress{IP: host, Port: port})
|
||||
// Retry connecting for up to 120 seconds (24 attempts * 5 seconds)
|
||||
logger.Info("Connecting to Infinity")
|
||||
var conn *infinity.InfinityConnection
|
||||
var err error
|
||||
for i := 0; i < 24; i++ {
|
||||
conn, err = infinity.Connect(infinity.NetworkAddress{IP: host, Port: port})
|
||||
if err == nil {
|
||||
break
|
||||
}
|
||||
if i < 23 {
|
||||
time.Sleep(5 * time.Second)
|
||||
}
|
||||
}
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("failed to connect to Infinity: %w", err)
|
||||
return nil, fmt.Errorf("Failed to connect to Infinity after 120s: %w", err)
|
||||
}
|
||||
|
||||
return &infinityClient{
|
||||
client := &infinityClient{
|
||||
conn: conn,
|
||||
dbName: cfg.DBName,
|
||||
}, nil
|
||||
}
|
||||
|
||||
return client, nil
|
||||
}
|
||||
|
||||
// WaitForHealthy blocks until Infinity is healthy or timeout
|
||||
func (c *infinityClient) WaitForHealthy(ctx context.Context, timeout time.Duration) error {
|
||||
logger.Info("Waiting for Infinity to be healthy")
|
||||
deadline := time.Now().Add(timeout)
|
||||
for time.Now().Before(deadline) {
|
||||
select {
|
||||
case <-ctx.Done():
|
||||
return ctx.Err()
|
||||
default:
|
||||
}
|
||||
|
||||
res, err := c.conn.ShowCurrentNode()
|
||||
if err != nil {
|
||||
time.Sleep(5 * time.Second)
|
||||
continue
|
||||
}
|
||||
// Use reflection to access ErrorCode and ServerStatus fields
|
||||
// since ShowCurrentNodeResponse is in an internal package
|
||||
v := reflect.ValueOf(res)
|
||||
if v.Kind() != reflect.Ptr {
|
||||
time.Sleep(5 * time.Second)
|
||||
continue
|
||||
}
|
||||
v = v.Elem()
|
||||
errorCode := v.FieldByName("ErrorCode")
|
||||
serverStatus := v.FieldByName("ServerStatus")
|
||||
if !errorCode.IsValid() || !serverStatus.IsValid() {
|
||||
time.Sleep(5 * time.Second)
|
||||
continue
|
||||
}
|
||||
// ErrorCode 0 means OK, ServerStatus "started" or "alive" means healthy
|
||||
if errorCode.Int() == 0 {
|
||||
status := serverStatus.String()
|
||||
if status == "started" || status == "alive" {
|
||||
logger.Info("Infinity is healthy")
|
||||
return nil
|
||||
}
|
||||
}
|
||||
time.Sleep(5 * time.Second)
|
||||
}
|
||||
return fmt.Errorf("Infinity not healthy after %v", timeout)
|
||||
}
|
||||
|
||||
// Engine Infinity engine implementation using Go SDK
|
||||
type infinityEngine struct {
|
||||
config *server.InfinityConfig
|
||||
client *infinityClient
|
||||
config *server.InfinityConfig
|
||||
client *infinityClient
|
||||
mappingFileName string
|
||||
docMetaMappingFileName string
|
||||
}
|
||||
|
||||
// NewEngine creates an Infinity engine
|
||||
@@ -77,9 +139,30 @@ func NewEngine(cfg interface{}) (*infinityEngine, error) {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
mappingFileName := infConfig.MappingFileName
|
||||
if mappingFileName == "" {
|
||||
mappingFileName = "infinity_mapping.json"
|
||||
}
|
||||
docMetaMappingFileName := infConfig.DocMetaMappingFileName
|
||||
if docMetaMappingFileName == "" {
|
||||
docMetaMappingFileName = "doc_meta_infinity_mapping.json"
|
||||
}
|
||||
|
||||
engine := &infinityEngine{
|
||||
config: infConfig,
|
||||
client: client,
|
||||
config: infConfig,
|
||||
client: client,
|
||||
mappingFileName: mappingFileName,
|
||||
docMetaMappingFileName: docMetaMappingFileName,
|
||||
}
|
||||
|
||||
// Wait for Infinity to be healthy
|
||||
if err := client.WaitForHealthy(context.Background(), 120*time.Second); err != nil {
|
||||
return nil, fmt.Errorf("Infinity not healthy: %w", err)
|
||||
}
|
||||
|
||||
// MigrateDB creates the database if it doesn't exist
|
||||
if err := engine.MigrateDB(context.Background()); err != nil {
|
||||
return nil, fmt.Errorf("failed to migrate database: %w", err)
|
||||
}
|
||||
|
||||
return engine, nil
|
||||
@@ -109,3 +192,12 @@ func (e *infinityEngine) Close() error {
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// MigrateDB creates the database if it doesn't exist
|
||||
func (e *infinityEngine) MigrateDB(ctx context.Context) error {
|
||||
_, err := e.client.conn.CreateDatabase(e.client.dbName, infinity.ConflictTypeIgnore, "")
|
||||
if err != nil {
|
||||
return fmt.Errorf("failed to create database: %w", err)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
@@ -17,21 +17,330 @@
|
||||
package infinity
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"context"
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"regexp"
|
||||
"strings"
|
||||
|
||||
infinity "github.com/infiniflow/infinity-go-sdk"
|
||||
"ragflow/internal/logger"
|
||||
"ragflow/internal/utility"
|
||||
|
||||
"go.uber.org/zap"
|
||||
)
|
||||
|
||||
// CreateIndex creates a table/index
|
||||
func (e *infinityEngine) CreateIndex(ctx context.Context, indexName string, mapping interface{}) error {
|
||||
return fmt.Errorf("infinity create table not implemented: waiting for official Go SDK")
|
||||
// fieldInfo represents a field in the infinity mapping schema
|
||||
type fieldInfo struct {
|
||||
Type string `json:"type"`
|
||||
Default interface{} `json:"default"`
|
||||
Analyzer interface{} `json:"analyzer"` // string or []string
|
||||
IndexType interface{} `json:"index_type"` // string or map
|
||||
Comment string `json:"comment"`
|
||||
}
|
||||
|
||||
// orderedFields preserves the order of fields as defined in JSON
|
||||
type orderedFields struct {
|
||||
Keys []string
|
||||
Fields map[string]fieldInfo
|
||||
}
|
||||
|
||||
func (o *orderedFields) UnmarshalJSON(data []byte) error {
|
||||
// Parse JSON manually to preserve key order
|
||||
// Look for key names by scanning the JSON string
|
||||
// This is a simple approach: find {"key": value, "key2": value2...}
|
||||
o.Fields = make(map[string]fieldInfo)
|
||||
o.Keys = make([]string, 0)
|
||||
|
||||
// Use a streaming JSON parser approach
|
||||
dec := json.NewDecoder(bytes.NewReader(data))
|
||||
tok, err := dec.Token()
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if delim, ok := tok.(json.Delim); ok && delim == '{' {
|
||||
for dec.More() {
|
||||
// Read key
|
||||
tok, err := dec.Token()
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
key, ok := tok.(string)
|
||||
if !ok {
|
||||
continue
|
||||
}
|
||||
o.Keys = append(o.Keys, key)
|
||||
|
||||
// Read value into fieldInfo
|
||||
var field fieldInfo
|
||||
if err := dec.Decode(&field); err != nil {
|
||||
return err
|
||||
}
|
||||
o.Fields[key] = field
|
||||
}
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// CreateIndex creates a table/index in Infinity
|
||||
// indexName is the table name prefix (e.g., "ragflow_<tenant_id>")
|
||||
// The full table name is built as "{indexName}_{datasetID}"
|
||||
func (e *infinityEngine) CreateIndex(ctx context.Context, indexName, datasetID string, vectorSize int, parserID string) error {
|
||||
vecSize := vectorSize
|
||||
|
||||
// Build full table name: {indexName}_{datasetID}
|
||||
tableName := fmt.Sprintf("%s_%s", indexName, datasetID)
|
||||
|
||||
// Use configured schema
|
||||
fpMapping := filepath.Join(utility.GetProjectRoot(), "conf", e.mappingFileName)
|
||||
|
||||
schemaData, err := os.ReadFile(fpMapping)
|
||||
if err != nil {
|
||||
return fmt.Errorf("Failed to read mapping file: %w", err)
|
||||
}
|
||||
|
||||
var schema orderedFields
|
||||
if err := json.Unmarshal(schemaData, &schema); err != nil {
|
||||
return fmt.Errorf("Failed to parse mapping file: %w", err)
|
||||
}
|
||||
|
||||
// Get database
|
||||
db, err := e.client.conn.GetDatabase(e.client.dbName)
|
||||
if err != nil {
|
||||
return fmt.Errorf("Failed to get database: %w", err)
|
||||
}
|
||||
|
||||
// Build column definitions (preserving JSON order)
|
||||
var columns infinity.TableSchema
|
||||
for _, fieldName := range schema.Keys {
|
||||
fieldInfo := schema.Fields[fieldName]
|
||||
col := infinity.ColumnDefinition{
|
||||
Name: fieldName,
|
||||
DataType: fieldInfo.Type,
|
||||
Default: fieldInfo.Default,
|
||||
// Comment: fieldInfo.Comment,
|
||||
}
|
||||
columns = append(columns, &col)
|
||||
}
|
||||
|
||||
// Add vector column
|
||||
vectorColName := fmt.Sprintf("q_%d_vec", vecSize)
|
||||
columns = append(columns, &infinity.ColumnDefinition{
|
||||
Name: vectorColName,
|
||||
DataType: fmt.Sprintf("vector,%d,float", vecSize),
|
||||
})
|
||||
|
||||
// Add chunk_data column for table parser
|
||||
if parserID == "table" {
|
||||
columns = append(columns, &infinity.ColumnDefinition{
|
||||
Name: "chunk_data",
|
||||
DataType: "json",
|
||||
Default: "{}",
|
||||
})
|
||||
}
|
||||
|
||||
// Create table
|
||||
table, err := db.CreateTable(tableName, columns, infinity.ConflictTypeIgnore)
|
||||
if err != nil {
|
||||
return fmt.Errorf("Failed to create table: %w", err)
|
||||
}
|
||||
logger.Debug("Infinity created table", zap.String("tableName", tableName))
|
||||
|
||||
// Create HNSW index on vector column
|
||||
_, err = table.CreateIndex(
|
||||
"q_vec_idx",
|
||||
infinity.NewIndexInfo(vectorColName, infinity.IndexTypeHnsw, map[string]string{
|
||||
"M": "16",
|
||||
"ef_construction": "50",
|
||||
"metric": "cosine",
|
||||
"encode": "lvq",
|
||||
}),
|
||||
infinity.ConflictTypeIgnore,
|
||||
"",
|
||||
)
|
||||
if err != nil {
|
||||
return fmt.Errorf("Failed to create HNSW index: %w", err)
|
||||
}
|
||||
|
||||
// Create full-text indexes for varchar fields with analyzers
|
||||
for _, fieldName := range schema.Keys {
|
||||
fieldInfo := schema.Fields[fieldName]
|
||||
if fieldInfo.Type != "varchar" || fieldInfo.Analyzer == nil {
|
||||
continue
|
||||
}
|
||||
|
||||
analyzers := []string{}
|
||||
switch a := fieldInfo.Analyzer.(type) {
|
||||
case string:
|
||||
analyzers = []string{a}
|
||||
case []interface{}:
|
||||
for _, v := range a {
|
||||
if s, ok := v.(string); ok {
|
||||
analyzers = append(analyzers, s)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
for _, analyzer := range analyzers {
|
||||
indexNameFt := fmt.Sprintf("ft_%s_%s",
|
||||
regexp.MustCompile(`[^a-zA-Z0-9]`).ReplaceAllString(fieldName, "_"),
|
||||
regexp.MustCompile(`[^a-zA-Z0-9]`).ReplaceAllString(analyzer, "_"),
|
||||
)
|
||||
_, err = table.CreateIndex(
|
||||
indexNameFt,
|
||||
infinity.NewIndexInfo(fieldName, infinity.IndexTypeFullText, map[string]string{"ANALYZER": analyzer}),
|
||||
infinity.ConflictTypeIgnore,
|
||||
"",
|
||||
)
|
||||
if err != nil {
|
||||
return fmt.Errorf("Failed to create fulltext index %s: %w", indexNameFt, err)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// Create secondary indexes for fields with index_type
|
||||
for _, fieldName := range schema.Keys {
|
||||
fieldInfo := schema.Fields[fieldName]
|
||||
if fieldInfo.IndexType == nil {
|
||||
continue
|
||||
}
|
||||
|
||||
indexTypeStr := ""
|
||||
params := map[string]string{}
|
||||
|
||||
switch it := fieldInfo.IndexType.(type) {
|
||||
case string:
|
||||
indexTypeStr = it
|
||||
case map[string]interface{}:
|
||||
if t, ok := it["type"].(string); ok {
|
||||
indexTypeStr = t
|
||||
}
|
||||
if card, ok := it["cardinality"].(string); ok {
|
||||
params["cardinality"] = card
|
||||
}
|
||||
}
|
||||
|
||||
if indexTypeStr == "secondary" {
|
||||
indexNameSec := fmt.Sprintf("sec_%s", fieldName)
|
||||
_, err = table.CreateIndex(
|
||||
indexNameSec,
|
||||
infinity.NewIndexInfo(fieldName, infinity.IndexTypeSecondary, params),
|
||||
infinity.ConflictTypeIgnore,
|
||||
"",
|
||||
)
|
||||
if err != nil {
|
||||
return fmt.Errorf("Failed to create secondary index %s: %w", indexNameSec, err)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
_ = table // suppress unused variable warning
|
||||
return nil
|
||||
}
|
||||
|
||||
// DeleteIndex deletes a table/index
|
||||
func (e *infinityEngine) DeleteIndex(ctx context.Context, indexName string) error {
|
||||
return fmt.Errorf("infinity drop table not implemented: waiting for official Go SDK")
|
||||
db, err := e.client.conn.GetDatabase(e.client.dbName)
|
||||
if err != nil {
|
||||
return fmt.Errorf("Failed to get database: %w", err)
|
||||
}
|
||||
|
||||
_, err = db.DropTable(indexName, infinity.ConflictTypeIgnore)
|
||||
if err != nil {
|
||||
return fmt.Errorf("Failed to drop table: %w", err)
|
||||
}
|
||||
logger.Debug("Infinity deleted table", zap.String("tableName", indexName))
|
||||
return nil
|
||||
}
|
||||
|
||||
// IndexExists checks if table/index exists
|
||||
func (e *infinityEngine) IndexExists(ctx context.Context, indexName string) (bool, error) {
|
||||
return false, fmt.Errorf("infinity check table existence not implemented: waiting for official Go SDK")
|
||||
db, err := e.client.conn.GetDatabase(e.client.dbName)
|
||||
if err != nil {
|
||||
return false, fmt.Errorf("Failed to get database: %w", err)
|
||||
}
|
||||
|
||||
_, err = db.GetTable(indexName)
|
||||
if err != nil {
|
||||
// Check if error is "table not found"
|
||||
if strings.Contains(err.Error(), "not found") || strings.Contains(err.Error(), "NotFound") {
|
||||
return false, nil
|
||||
}
|
||||
return false, err
|
||||
}
|
||||
return true, nil
|
||||
}
|
||||
|
||||
// CreateDocMetaIndex creates the document metadata table/index
|
||||
func (e *infinityEngine) CreateDocMetaIndex(ctx context.Context, indexName string) error {
|
||||
// Get database
|
||||
db, err := e.client.conn.GetDatabase(e.client.dbName)
|
||||
if err != nil {
|
||||
return fmt.Errorf("Failed to get database: %w", err)
|
||||
}
|
||||
|
||||
// Use configured doc_meta mapping file
|
||||
fpMapping := filepath.Join(utility.GetProjectRoot(), "conf", e.docMetaMappingFileName)
|
||||
|
||||
schemaData, err := os.ReadFile(fpMapping)
|
||||
if err != nil {
|
||||
return fmt.Errorf("Failed to read mapping file: %w", err)
|
||||
}
|
||||
|
||||
var schema map[string]fieldInfo
|
||||
if err := json.Unmarshal(schemaData, &schema); err != nil {
|
||||
return fmt.Errorf("Failed to parse mapping file: %w", err)
|
||||
}
|
||||
|
||||
// Build column definitions
|
||||
var columns infinity.TableSchema
|
||||
for fieldName, fieldInfo := range schema {
|
||||
col := infinity.ColumnDefinition{
|
||||
Name: fieldName,
|
||||
DataType: fieldInfo.Type,
|
||||
Default: fieldInfo.Default,
|
||||
// Comment: fieldInfo.Comment,
|
||||
}
|
||||
columns = append(columns, &col)
|
||||
}
|
||||
|
||||
// Create table
|
||||
_, err = db.CreateTable(indexName, columns, infinity.ConflictTypeIgnore)
|
||||
if err != nil {
|
||||
return fmt.Errorf("Failed to create doc meta table: %w", err)
|
||||
}
|
||||
logger.Debug("Infinity created doc meta table", zap.String("tableName", indexName))
|
||||
|
||||
// Get table for creating indexes
|
||||
table, err := db.GetTable(indexName)
|
||||
if err != nil {
|
||||
return fmt.Errorf("Failed to get table: %w", err)
|
||||
}
|
||||
|
||||
// Create secondary index on id
|
||||
_, err = table.CreateIndex(
|
||||
fmt.Sprintf("idx_%s_id", indexName),
|
||||
infinity.NewIndexInfo("id", infinity.IndexTypeSecondary, nil),
|
||||
infinity.ConflictTypeIgnore,
|
||||
"",
|
||||
)
|
||||
if err != nil {
|
||||
return fmt.Errorf("Failed to create secondary index on id: %w", err)
|
||||
}
|
||||
|
||||
// Create secondary index on kb_id
|
||||
_, err = table.CreateIndex(
|
||||
fmt.Sprintf("idx_%s_kb_id", indexName),
|
||||
infinity.NewIndexInfo("kb_id", infinity.IndexTypeSecondary, nil),
|
||||
infinity.ConflictTypeIgnore,
|
||||
"",
|
||||
)
|
||||
if err != nil {
|
||||
return fmt.Errorf("Failed to create secondary index on kb_id: %w", err)
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user