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 }