Files
gochat/internal/service/captain_document_worker.go
T

169 lines
6.3 KiB
Go

package service
import (
"context"
"encoding/json"
"fmt"
"sync"
"time"
"github.com/gochat/gochat/internal/model"
"github.com/gochat/gochat/internal/worker"
)
const (
TaskTypeCaptainDocumentSync = "captain:document_sync"
TaskTypeCaptainDocumentCrawl = "captain:document_crawl"
TaskTypeCaptainDocumentPageCrawlParse = "captain:document_page_crawl_parse"
TaskTypeCaptainDocumentScheduleSyncs = "captain:documents_schedule_syncs"
TaskTypeCaptainDocumentResponseBuilder = "captain:document_response_builder"
TaskTypeCaptainLLMUpdateEmbedding = "captain:llm_update_embedding"
)
const captainDocumentScheduleInterval = 24 * time.Hour
type captainDocumentSyncJob struct {
AccountID uint `json:"account_id"`
DocumentID uint `json:"document_id"`
}
type captainDocumentCrawlJob struct {
AccountID uint `json:"account_id"`
DocumentID uint `json:"document_id"`
}
type captainDocumentPageCrawlParseJob struct {
AccountID uint `json:"account_id"`
AssistantID uint `json:"assistant_id"`
PageLink string `json:"page_link"`
}
type captainDocumentScheduleSyncsJob struct {
PlanName string `json:"plan_name,omitempty"`
}
type captainDocumentResponseBuilderJob struct {
AccountID uint `json:"account_id"`
DocumentID uint `json:"document_id"`
}
type captainLLMUpdateEmbeddingJob struct {
AccountID uint `json:"account_id"`
ResponseID uint `json:"response_id"`
Content string `json:"content,omitempty"`
}
var captainDocumentRegistrations sync.Map
// RegisterCaptainDocumentJobs wires Captain::Documents::PerformSyncJob into
// the durable worker. The sync backend remains fakeable for tests and disabled
// deployments.
func RegisterCaptainDocumentJobs(wp *worker.WorkerPool, svc *CaptainDocumentService) {
if wp == nil || svc == nil {
return
}
if _, loaded := captainDocumentRegistrations.LoadOrStore(wp, struct{}{}); loaded {
return
}
wp.Register(TaskTypeCaptainDocumentSync, svc.performDocumentSyncJob)
wp.Register(TaskTypeCaptainDocumentCrawl, svc.performDocumentCrawlJob)
wp.Register(TaskTypeCaptainDocumentPageCrawlParse, svc.performDocumentPageCrawlParseJob)
wp.Register(TaskTypeCaptainDocumentScheduleSyncs, svc.performDocumentScheduleSyncsJob)
wp.Register(TaskTypeCaptainDocumentResponseBuilder, svc.performDocumentResponseBuilderJob)
wp.Register(TaskTypeCaptainLLMUpdateEmbedding, svc.performCaptainLLMUpdateEmbeddingJob)
}
func EnqueueCaptainDocumentScheduleSyncs(ctx context.Context, wp *worker.WorkerPool, scheduledAt time.Time) (*model.BackgroundJob, error) {
if wp == nil {
return nil, nil
}
if scheduledAt.IsZero() {
scheduledAt = time.Now()
}
return wp.Enqueue(ctx, TaskTypeCaptainDocumentScheduleSyncs, captainDocumentScheduleSyncsJob{},
worker.WithQueue("scheduled_jobs"),
worker.WithScheduledAt(scheduledAt),
worker.WithMaxAttempts(3),
worker.WithIdempotencyKey(captainDocumentScheduleIdempotencyKey(scheduledAt)),
)
}
func captainDocumentScheduleIdempotencyKey(scheduledAt time.Time) string {
bucket := scheduledAt.UTC().Truncate(captainDocumentScheduleInterval).Unix()
return fmt.Sprintf("captain:documents_schedule_syncs:%d", bucket)
}
func (s *CaptainDocumentService) performDocumentSyncJob(ctx context.Context, job *model.BackgroundJob) error {
var payload captainDocumentSyncJob
if err := json.Unmarshal(job.Payload, &payload); err != nil {
return fmt.Errorf("unmarshal captain document sync job: %w", err)
}
if payload.AccountID == 0 || payload.DocumentID == 0 {
return fmt.Errorf("invalid captain document sync job payload: %#v", payload)
}
_, err := s.SyncDocumentByAccount(ctx, payload.AccountID, payload.DocumentID)
return err
}
func (s *CaptainDocumentService) performDocumentCrawlJob(ctx context.Context, job *model.BackgroundJob) error {
var payload captainDocumentCrawlJob
if err := json.Unmarshal(job.Payload, &payload); err != nil {
return fmt.Errorf("unmarshal captain document crawl job: %w", err)
}
if payload.AccountID == 0 || payload.DocumentID == 0 {
return fmt.Errorf("invalid captain document crawl job payload: %#v", payload)
}
_, err := s.CrawlDocumentByAccount(ctx, payload.AccountID, payload.DocumentID)
return err
}
func (s *CaptainDocumentService) performDocumentPageCrawlParseJob(ctx context.Context, job *model.BackgroundJob) error {
var payload captainDocumentPageCrawlParseJob
if err := json.Unmarshal(job.Payload, &payload); err != nil {
return fmt.Errorf("unmarshal captain document page crawl parse job: %w", err)
}
if payload.AccountID == 0 || payload.AssistantID == 0 || payload.PageLink == "" {
return fmt.Errorf("invalid captain document page crawl parse job payload: %#v", payload)
}
_, err := s.ParseCrawledPage(ctx, payload.AccountID, payload.AssistantID, payload.PageLink)
return err
}
func (s *CaptainDocumentService) performDocumentScheduleSyncsJob(ctx context.Context, job *model.BackgroundJob) error {
if len(job.Payload) > 0 {
var payload captainDocumentScheduleSyncsJob
if err := json.Unmarshal(job.Payload, &payload); err != nil {
return fmt.Errorf("unmarshal captain document schedule syncs job: %w", err)
}
}
if _, err := s.ScheduleDueDocumentSyncs(ctx, time.Now()); err != nil {
return err
}
_, err := EnqueueCaptainDocumentScheduleSyncs(ctx, s.worker, time.Now().Add(captainDocumentScheduleInterval))
return err
}
func (s *CaptainDocumentService) performDocumentResponseBuilderJob(ctx context.Context, job *model.BackgroundJob) error {
var payload captainDocumentResponseBuilderJob
if err := json.Unmarshal(job.Payload, &payload); err != nil {
return fmt.Errorf("unmarshal captain document response builder job: %w", err)
}
if payload.AccountID == 0 || payload.DocumentID == 0 {
return fmt.Errorf("invalid captain document response builder job payload: %#v", payload)
}
_, err := s.BuildResponsesForDocumentByAccount(ctx, payload.AccountID, payload.DocumentID)
return err
}
func (s *CaptainDocumentService) performCaptainLLMUpdateEmbeddingJob(ctx context.Context, job *model.BackgroundJob) error {
var payload captainLLMUpdateEmbeddingJob
if err := json.Unmarshal(job.Payload, &payload); err != nil {
return fmt.Errorf("unmarshal captain llm update embedding job: %w", err)
}
if payload.AccountID == 0 || payload.ResponseID == 0 {
return fmt.Errorf("invalid captain llm update embedding job payload: %#v", payload)
}
_, err := s.UpdateAssistantResponseEmbeddingByAccount(ctx, payload.AccountID, payload.ResponseID, payload.Content)
return err
}