package service import ( "context" "fmt" "path/filepath" "time" "rag-local/common" "rag-local/kb/consts" "rag-local/kb/dao" "rag-local/kb/model/entity" "github.com/gogf/gf/v2/errors/gerror" "github.com/gogf/gf/v2/frame/g" ) var ParseTaskService = &parseTaskService{} type parseTaskService struct{} // StartParsePoller 启动任务轮询:单 goroutine 串行消费待处理任务(与 video-factory StartVideoPoller 同模式) func (s *parseTaskService) StartParsePoller(ctx context.Context) { go func() { g.Log().Info(ctx, "parse task poller started") for { select { case <-ctx.Done(): return case <-time.After(consts.ParsePollIntervalSeconds * time.Second): s.processOne(ctx) } } }() } // processOne 处理一个待处理任务:解析 → 分块 → 落库(向量化在 M3 接入) func (s *parseTaskService) processOne(ctx context.Context) { task, err := dao.ParseTask.NextPending(ctx) if err != nil { g.Log().Errorf(ctx, "next parse task failed: %v", err) return } if task == nil { return } if err := dao.ParseTask.UpdateStatus(ctx, task.Id, consts.TaskStatusRunning, ""); err != nil { g.Log().Errorf(ctx, "mark task running failed: %v", err) return } doc, err := dao.Document.GetOne(ctx, task.DocumentId) if err != nil { s.fail(ctx, task, "读取文档失败: "+err.Error()) return } if doc == nil { s.fail(ctx, task, "文档不存在") return } if task.TaskType == consts.TaskTypeReembed { s.processReembed(ctx, task) return } if err := dao.Document.UpdateFields(ctx, doc.Id, g.Map{"status": consts.DocumentStatusParsing}); err != nil { s.fail(ctx, task, "更新文档状态失败: "+err.Error()) return } text, err := common.ParseFile(filepath.Join("workspace", doc.FilePath)) if err != nil { s.fail(ctx, task, "解析文件失败: "+err.Error()) return } // 数据集分块配置(大小/重叠),未设置用默认值 chunkSize, chunkOverlap := consts.DefaultChunkSize, consts.DefaultChunkOverlap if ds, err := dao.Dataset.GetOne(ctx, task.DatasetId); err == nil && ds != nil { chunkSize, chunkOverlap = ds.ChunkSize, ds.ChunkOverlap } // 数据集必须绑定向量模型,构建失败直接失败任务(不降级全文索引) cfgId, err := dao.Dataset.GetEmbeddingCfgId(ctx, task.DatasetId) if err != nil { s.fail(ctx, task, "读取数据集向量模型配置失败: "+err.Error()) return } if cfgId <= 0 { s.fail(ctx, task, "数据集未绑定向量模型") return } embedder, err := BuildEmbedder(ctx, cfgId) if err != nil { s.fail(ctx, task, "构建向量模型失败: "+err.Error()) return } if dim := embedder.Dim(); dim != g.Cfg().MustGet(ctx, "vector.dim", consts.DefaultEmbeddingDim).Int() { g.Log().Warningf(ctx, "embedding 维度 %d 与 vec0 表维度不一致,请确认 vector.dim 配置", dim) } // 重新解析时先清空旧分块,避免新旧分块并存 if err := ChunkService.DeleteByDocument(ctx, doc.Id); err != nil { s.fail(ctx, task, "清理旧分块失败: "+err.Error()) return } // 原始全文落库(分块抽屉展示用),失败不阻断流水线 if err := dao.Document.UpdateFields(ctx, doc.Id, g.Map{"content": text}); err != nil { g.Log().Warningf(ctx, "save document content failed: %v", err) } chunks, unitPattern, ctxPattern, err := ChunkService.SplitAuto(ctx, text, chunkSize, chunkOverlap, embedder) if err != nil { s.fail(ctx, task, "分块失败: "+err.Error()) return } // 结构识别出的模式落库(列表展示用);未识别出结构时保留旧值,避免清空已展示的模式 if unitPattern != "" { if err := dao.Dataset.UpdateFields(ctx, task.DatasetId, g.Map{ "unit_pattern": unitPattern, "context_pattern": ctxPattern, }); err != nil { g.Log().Warningf(ctx, "save detected structure patterns failed: %v", err) } } if err := dao.Document.UpdateFields(ctx, doc.Id, g.Map{"status": consts.DocumentStatusEmbedding}); err != nil { s.fail(ctx, task, "更新文档状态失败: "+err.Error()) return } if err := ChunkService.InsertAll(ctx, task.DatasetId, doc.Id, chunks, embedder); err != nil { s.fail(ctx, task, "写入分块失败: "+err.Error()) return } // 第 3.5 步:知识图谱抽取(失败不阻断流水线,但文档标记图谱未构建,避免误报已完成) if err := dao.Document.UpdateFields(ctx, doc.Id, g.Map{"status": consts.DocumentStatusKgBuilding}); err != nil { g.Log().Warningf(ctx, "mark kg building failed: %v", err) } kgFailed, kgErr := KgEntityService.ExtractDocument(ctx, task.DatasetId, doc.Id) kgMsg := "" switch { case kgErr != nil: kgMsg = "知识图谱未构建: " + kgErr.Error() case kgFailed > 0: kgMsg = fmt.Sprintf("知识图谱未构建:%d 个分块抽取失败", kgFailed) } if kgMsg != "" { g.Log().Warningf(ctx, "kg extract incomplete for doc %d: %s", doc.Id, kgMsg) } if err := dao.Document.UpdateFields(ctx, doc.Id, g.Map{"status": consts.DocumentStatusDone, "error_msg": kgMsg}); err != nil { g.Log().Warningf(ctx, "mark document done failed: %v", err) } if err := dao.ParseTask.UpdateStatus(ctx, task.Id, consts.TaskStatusDone, kgMsg); err != nil { g.Log().Errorf(ctx, "mark task done failed: %v", err) } } // processReembed 重新向量化任务:分块文本不变,用数据集当前绑定模型重算全部向量;失败不动文档状态 func (s *parseTaskService) processReembed(ctx context.Context, task *entity.ParseTask) { if _, err := DocumentService.Reembed(ctx, task.DocumentId); err != nil { _ = dao.ParseTask.UpdateStatus(ctx, task.Id, consts.TaskStatusFailed, "重新向量化失败: "+err.Error()) g.Log().Errorf(ctx, "reembed task %d failed: %s", task.Id, err.Error()) return } if err := dao.ParseTask.UpdateStatus(ctx, task.Id, consts.TaskStatusDone, ""); err != nil { g.Log().Errorf(ctx, "mark task done failed: %v", err) } } func (s *parseTaskService) fail(ctx context.Context, task *entity.ParseTask, msg string) { _ = dao.ParseTask.UpdateStatus(ctx, task.Id, consts.TaskStatusFailed, msg) _ = dao.Document.UpdateFields(ctx, task.DocumentId, g.Map{ "status": consts.DocumentStatusFailed, "error_msg": msg, }) g.Log().Errorf(ctx, "parse task %d failed: %s", task.Id, msg) } // List 任务列表 func (s *parseTaskService) List(ctx context.Context, page, pageSize int) ([]*entity.ParseTask, int, error) { return dao.ParseTask.List(ctx, page, pageSize) } // Retry 失败任务重置为待处理(文档状态同步重置) func (s *parseTaskService) Retry(ctx context.Context, id int64) error { task, err := dao.ParseTask.GetOne(ctx, id) if err != nil { return err } if task == nil { return gerror.New("任务不存在") } if task.Status != consts.TaskStatusFailed { return gerror.New("仅失败任务可重试") } if err := dao.ParseTask.UpdateStatus(ctx, id, consts.TaskStatusPending, ""); err != nil { return err } return dao.Document.UpdateFields(ctx, task.DocumentId, g.Map{ "status": consts.DocumentStatusPending, "error_msg": "", }) }