MCPcopy Create free account
hub / github.com/53AI/53AIHub / _enqueueRetrievalChunksByFile

Function _enqueueRetrievalChunksByFile

api/service/rag/embedding_enqueue_helper.go:36–91  ·  view source on GitHub ↗
(eid, fileID, libraryID int64)

Source from the content-addressed store, hash-verified

34}
35
36func _enqueueRetrievalChunksByFile(eid, fileID, libraryID int64) error {
37 // 未配置向量化渠道则不入队
38 cfgSvc := NewChunkConfigService(model.DB)
39 cfg, cfgErr := cfgSvc.GetConfig(eid, &libraryID, model.ChunkTypeDefault)
40 if cfgErr != nil {
41 CheckEmbeddingStepStatusSave(eid, fileID, fmt.Sprintf("time:%d; 获取向量化配置失败 Err:%v", time.Now().Unix(), cfgErr))
42 logger.Warn(context.TODO(), fmt.Sprintf("[embEnqueueConfigCheckError][eid=%d][fileID=%d] %v", eid, fileID, cfgErr))
43 }
44 if cfg == nil || cfg.EmbeddingChannelID == nil {
45 CheckEmbeddingStepStatusSave(eid, fileID, fmt.Sprintf("time:%d; 未配置向量化渠道", time.Now().Unix()))
46 logger.Warn(context.TODO(), fmt.Sprintf("[embEnqueueSkipNoEmbeddingChannel][eid=%d][fileID=%d][libraryID=%d]", eid, fileID, libraryID))
47 return nil
48 }
49
50 q := GetDefaultEmbeddingQueue()
51 if q == nil {
52 CheckEmbeddingStepStatusSave(eid, fileID, fmt.Sprintf("time:%d; 未配置向量化渠道", time.Now().Unix()))
53 logger.Warn(context.TODO(), fmt.Sprintf("[embEnqueueSkipNoQueue][eid=%d][fileID=%d]", eid, fileID))
54 return nil
55 }
56
57 // 检查 API
58 retrievalService := NewRetrievalChunkService(model.DB)
59 err := retrievalService.CheckGenerateEmbeddingForChunk(eid, &libraryID, &fileID, "TestAPI")
60 if err != nil {
61 CheckEmbeddingStepStatusSave(eid, fileID, fmt.Sprintf("time:%d; embedding API 不可用 Err:%v", time.Now().Unix(), err))
62 return fmt.Errorf("embedding API 不可用: %v", err)
63 }
64
65 var rids []int64
66 if err := model.DB.Model(&model.RetrievalChunk{}).
67 Joins("JOIN files ON retrieval_chunks.file_id = files.id AND retrieval_chunks.eid = files.eid").
68 Where("retrieval_chunks.eid = ? AND retrieval_chunks.file_id = ? AND files.parsing_status != ?", eid, fileID, model.FileParsingStatusDisabled).
69 Pluck("retrieval_chunks.id", &rids).Error; err != nil {
70 CheckEmbeddingStepStatusSave(eid, fileID, fmt.Sprintf("time:%d; 获取检索块失败 Err:%v", time.Now().Unix(), err))
71 return err
72 }
73 CheckEmbeddingStepStatusSave(eid, fileID, fmt.Sprintf("time:%d; 获取检索块成功 count:%d", time.Now().Unix(), len(rids)))
74 for _, rid := range rids {
75 _, err := q.EnqueueIfNotExists(context.TODO(), EmbeddingTask{
76 Eid: eid,
77 RetrievalChunkID: rid,
78 FileID: fileID,
79 LibraryID: libraryID,
80 TraceID: "",
81 Retries: 0,
82 })
83 CheckEmbeddingStepStatusSave(eid, fileID, fmt.Sprintf("add_rid:%d", rid))
84 if err != nil {
85 logger.Warn(context.TODO(), fmt.Sprintf("[embEnqueueOneFail][eid=%d][fileID=%d][rid=%d]%+v", eid, fileID, rid, err))
86 }
87 }
88
89 logger.Info(context.TODO(), fmt.Sprintf("[embEnqueueByFileDone][eid=%d][fileID=%d][count=%d]", eid, fileID, len(rids)))
90 return nil
91}
92
93// EnqueueRetrievalChunk enqueues a single retrieval chunk id.

Callers 1

Calls 10

GetConfigMethod · 0.95
NewChunkConfigServiceFunction · 0.85
GetDefaultEmbeddingQueueFunction · 0.85
NewRetrievalChunkServiceFunction · 0.85
WarnMethod · 0.80
ErrorfMethod · 0.80
InfoMethod · 0.80
EnqueueIfNotExistsMethod · 0.65

Tested by

no test coverage detected