InitRAGJobEngine 初始化RAG任务引擎
()
| 36 | |
| 37 | // InitRAGJobEngine 初始化RAG任务引擎 |
| 38 | func InitRAGJobEngine() { |
| 39 | initOnce.Do(func() { |
| 40 | if !common.IsRedisEnabled() { |
| 41 | logger.SysLog("Redis not enabled, skipping RAG job engine initialization") |
| 42 | return |
| 43 | } |
| 44 | |
| 45 | // 初始化 V2 引擎 |
| 46 | ragJobFactoryV2 = v2factory.NewJobFactory(model.DB, common.RDB) |
| 47 | ragJobEngineV2 = v2engines.NewRagJobEngineV2(common.RDB, model.DB, ragJobFactoryV2) |
| 48 | |
| 49 | // 注册 V2 Handler |
| 50 | ragJobEngineV2.RegisterHandler("document_parsing", v2steps.NewDocumentParsingHandler(model.DB)) |
| 51 | ragJobEngineV2.RegisterHandler("content_cleaning", func(ctx context.Context, job *model.RagJob, config json.RawMessage) error { |
| 52 | logger.Info(ctx, "V2 Handler: content_cleaning executed") |
| 53 | return nil |
| 54 | }) |
| 55 | ragJobEngineV2.RegisterHandler("summary_generation", v2steps.NewSummaryGenerationHandler(model.DB)) |
| 56 | ragJobEngineV2.RegisterHandler("document_chunking", v2steps.NewDocumentChunkingHandler(model.DB)) |
| 57 | ragJobEngineV2.RegisterHandler("vector_indexing", v2steps.NewVectorIndexingHandler(model.DB)) |
| 58 | ragJobEngineV2.RegisterHandler("graph_generation", v2steps.NewGraphGenerationHandler(model.DB)) |
| 59 | |
| 60 | // 注册 V2 Recovery Handler |
| 61 | ragJobEngineV2.RegisterRecoveryHandler("document_parsing", v2steps.RecoverDocumentParsing(model.DB)) |
| 62 | ragJobEngineV2.RegisterRecoveryHandler("content_cleaning", v2steps.RecoverContentCleaning()) |
| 63 | ragJobEngineV2.RegisterRecoveryHandler("summary_generation", v2steps.RecoverSummaryGeneration(model.DB)) |
| 64 | ragJobEngineV2.RegisterRecoveryHandler("document_chunking", v2steps.RecoverDocumentChunking(model.DB)) |
| 65 | ragJobEngineV2.RegisterRecoveryHandler("vector_indexing", v2steps.RecoverVectorIndexing(model.DB)) |
| 66 | ragJobEngineV2.RegisterRecoveryHandler("graph_generation", v2steps.RecoverGraphGeneration(model.DB)) |
| 67 | |
| 68 | // 在后台执行恢复+启动 Worker,不阻塞主流程启动 |
| 69 | go func() { |
| 70 | ragJobEngineV2.RecoverStuckJobs(context.Background()) |
| 71 | rag.StartSiteEmbeddingReindexCoordinator(context.Background(), model.DB, 30*time.Second) |
| 72 | ragJobEngineV2.StartWorkers() |
| 73 | // 每小时清理一次超过 24h 无心跳的死 job |
| 74 | ragJobEngineV2.StartStaleJobCleaner(context.Background(), 1*time.Hour) |
| 75 | logger.SysLog("RAG job engine recovery, workers, and stale-job cleaner started (background)") |
| 76 | }() |
| 77 | |
| 78 | logger.SysLog("RAG job engine initialized (recovery running in background)") |
| 79 | }) |
| 80 | } |
| 81 | |
| 82 | type RetryJobStepOptionsV2 struct { |
| 83 | Continue bool |
nothing calls this directly
no test coverage detected