(eid, fileID, libraryID int64)
| 34 | } |
| 35 | |
| 36 | func _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. |
no test coverage detected