feat(captain): queue document syncs
This commit is contained in:
@@ -574,6 +574,7 @@ func Bootstrap(env string) (*App, error) {
|
||||
// Captain services (P10 M10 — Captain AI + Copilot)
|
||||
captainAssistantService := service.NewCaptainAssistantService(captainAssistantRepo, captainInboxRepo, captainDocumentRepo, captainAssistantResponseRepo, llmProvider)
|
||||
captainDocumentService := service.NewCaptainDocumentService(captainDocumentRepo, llmProvider, captainAssistantRepo)
|
||||
captainDocumentService.SetWorkerPool(workerPool)
|
||||
captainScenarioService := service.NewCaptainScenarioService(captainScenarioRepo, captainAssistantRepo)
|
||||
captainCustomToolService := service.NewCaptainCustomToolService(captainCustomToolRepo)
|
||||
copilotService := service.NewCopilotService(copilotThreadRepo, copilotMessageRepo, copilotSuggestionRepo, llmProvider, captainAssistantRepo)
|
||||
|
||||
@@ -194,7 +194,7 @@ func (h *CaptainDocumentHandler) SyncDocument(c *gin.Context) {
|
||||
return
|
||||
}
|
||||
|
||||
if _, err := h.svc.MarkSyncing(c.Request.Context(), accountID, id); err != nil {
|
||||
if _, err := h.svc.RequestSyncDocumentByAccount(c.Request.Context(), accountID, id); err != nil {
|
||||
applogger.L().Errorf("SyncDocument: %v", err)
|
||||
response.AbortWithStatusError(c, captainAssistantErrorStatus(err), response.ErrInternal, "failed to sync document")
|
||||
return
|
||||
|
||||
@@ -12,6 +12,7 @@ import (
|
||||
"github.com/gochat/gochat/internal/llm"
|
||||
"github.com/gochat/gochat/internal/model"
|
||||
"github.com/gochat/gochat/internal/repository"
|
||||
"github.com/gochat/gochat/internal/worker"
|
||||
applogger "github.com/gochat/gochat/pkg/logger"
|
||||
)
|
||||
|
||||
@@ -22,6 +23,7 @@ type CaptainDocumentService struct {
|
||||
assistantRepo *repository.CaptainAssistantRepo
|
||||
llmProvider llm.Provider
|
||||
syncBackend CaptainDocumentSyncBackend
|
||||
worker *worker.WorkerPool
|
||||
}
|
||||
|
||||
type CaptainDocumentSyncBackend interface {
|
||||
@@ -54,6 +56,11 @@ func (s *CaptainDocumentService) SetSyncBackend(syncBackend CaptainDocumentSyncB
|
||||
s.syncBackend = syncBackend
|
||||
}
|
||||
|
||||
func (s *CaptainDocumentService) SetWorkerPool(wp *worker.WorkerPool) {
|
||||
s.worker = wp
|
||||
RegisterCaptainDocumentJobs(wp, s)
|
||||
}
|
||||
|
||||
// --- Request DTOs ---
|
||||
|
||||
// CreateDocumentRequest is the DTO for creating a document.
|
||||
@@ -216,6 +223,20 @@ func (s *CaptainDocumentService) MarkSyncing(ctx context.Context, accountID, id
|
||||
return s.documentRepo.GetByAccountAndID(ctx, accountID, id)
|
||||
}
|
||||
|
||||
func (s *CaptainDocumentService) RequestSyncDocumentByAccount(ctx context.Context, accountID, id uint) (*model.CaptainDocument, error) {
|
||||
doc, err := s.MarkSyncing(ctx, accountID, id)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if s.worker == nil {
|
||||
return doc, nil
|
||||
}
|
||||
if _, err := s.worker.Enqueue(ctx, TaskTypeCaptainDocumentSync, captainDocumentSyncJob{AccountID: accountID, DocumentID: id}, worker.WithQueue("low"), worker.WithMaxAttempts(3)); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return doc, nil
|
||||
}
|
||||
|
||||
func (s *CaptainDocumentService) SyncDocumentByAccount(ctx context.Context, accountID, id uint) (*model.CaptainDocument, error) {
|
||||
doc, err := s.documentRepo.GetByAccountAndID(ctx, accountID, id)
|
||||
if err != nil {
|
||||
|
||||
@@ -5,9 +5,11 @@ import (
|
||||
"errors"
|
||||
"fmt"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/gochat/gochat/internal/model"
|
||||
"github.com/gochat/gochat/internal/repository"
|
||||
"github.com/gochat/gochat/internal/worker"
|
||||
"github.com/stretchr/testify/assert"
|
||||
"github.com/stretchr/testify/require"
|
||||
"gorm.io/driver/sqlite"
|
||||
@@ -19,7 +21,7 @@ func setupCaptainDocumentServiceTest(t *testing.T) (*gorm.DB, *CaptainDocumentSe
|
||||
dbName := fmt.Sprintf("file:%s?mode=memory&cache=private", t.Name())
|
||||
db, err := gorm.Open(sqlite.Open(dbName), &gorm.Config{})
|
||||
require.NoError(t, err)
|
||||
require.NoError(t, db.AutoMigrate(&model.Account{}, &model.CaptainAssistant{}, &model.CaptainDocument{}))
|
||||
require.NoError(t, db.AutoMigrate(&model.Account{}, &model.CaptainAssistant{}, &model.CaptainDocument{}, &model.BackgroundJob{}))
|
||||
t.Cleanup(func() {
|
||||
sqlDB, _ := db.DB()
|
||||
sqlDB.Close()
|
||||
@@ -104,6 +106,55 @@ func TestCaptainDocumentService_SyncDocumentByAccountScopesDocument(t *testing.T
|
||||
require.Error(t, err)
|
||||
}
|
||||
|
||||
func TestCaptainDocumentService_RequestSyncQueuesDurableJob(t *testing.T) {
|
||||
db, svc := setupCaptainDocumentServiceTest(t)
|
||||
account, _, doc := seedCaptainDocumentSyncFixture(t, db)
|
||||
backend := &captainDocumentFakeSyncBackend{result: &CaptainDocumentSyncResult{Title: "Durable Help", Content: "fresh durable content"}}
|
||||
svc.SetSyncBackend(backend)
|
||||
now := time.Date(2026, 6, 5, 23, 0, 0, 0, time.UTC)
|
||||
wp := worker.NewWorkerPoolWithOptions(db, worker.WithNow(func() time.Time { return now }))
|
||||
svc.SetWorkerPool(wp)
|
||||
|
||||
queued, err := svc.RequestSyncDocumentByAccount(context.Background(), account.ID, doc.ID)
|
||||
require.NoError(t, err)
|
||||
assert.Equal(t, model.DocumentSyncStatusPending, queued.SyncStatus)
|
||||
assert.Equal(t, uint(0), backend.documentID)
|
||||
|
||||
var jobCount int64
|
||||
require.NoError(t, db.Model(&model.BackgroundJob{}).Where("job_type = ? AND queue = ? AND status = ?", TaskTypeCaptainDocumentSync, "low", model.BackgroundJobStatusQueued).Count(&jobCount).Error)
|
||||
assert.Equal(t, int64(1), jobCount)
|
||||
|
||||
processed, err := wp.ProcessOne(context.Background())
|
||||
require.NoError(t, err)
|
||||
assert.True(t, processed)
|
||||
assert.Equal(t, doc.ID, backend.documentID)
|
||||
|
||||
var synced model.CaptainDocument
|
||||
require.NoError(t, db.First(&synced, doc.ID).Error)
|
||||
assert.Equal(t, "Durable Help", synced.Name)
|
||||
assert.Equal(t, "fresh durable content", synced.Content)
|
||||
assert.Equal(t, model.DocumentSyncStatusSynced, synced.SyncStatus)
|
||||
assert.Empty(t, synced.LastSyncErrorCode)
|
||||
}
|
||||
|
||||
func TestCaptainDocumentService_DocumentSyncJobRetriesMissingDocument(t *testing.T) {
|
||||
db, svc := setupCaptainDocumentServiceTest(t)
|
||||
now := time.Date(2026, 6, 5, 23, 15, 0, 0, time.UTC)
|
||||
wp := worker.NewWorkerPoolWithOptions(db, worker.WithNow(func() time.Time { return now }), worker.WithBackoff(func(attempt int) time.Duration { return time.Minute }))
|
||||
svc.SetWorkerPool(wp)
|
||||
|
||||
_, err := wp.Enqueue(context.Background(), TaskTypeCaptainDocumentSync, captainDocumentSyncJob{AccountID: 999, DocumentID: 9999}, worker.WithQueue("low"), worker.WithMaxAttempts(3))
|
||||
require.NoError(t, err)
|
||||
processed, err := wp.ProcessOne(context.Background())
|
||||
require.Error(t, err)
|
||||
assert.True(t, processed)
|
||||
|
||||
var job model.BackgroundJob
|
||||
require.NoError(t, db.Where("job_type = ?", TaskTypeCaptainDocumentSync).First(&job).Error)
|
||||
assert.Equal(t, model.BackgroundJobStatusRetrying, job.Status)
|
||||
assert.NotEmpty(t, job.LastError)
|
||||
}
|
||||
|
||||
type captainDocumentFakeSyncBackend struct {
|
||||
result *CaptainDocumentSyncResult
|
||||
err error
|
||||
|
||||
@@ -0,0 +1,45 @@
|
||||
package service
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"sync"
|
||||
|
||||
"github.com/gochat/gochat/internal/model"
|
||||
"github.com/gochat/gochat/internal/worker"
|
||||
)
|
||||
|
||||
const TaskTypeCaptainDocumentSync = "captain:document_sync"
|
||||
|
||||
type captainDocumentSyncJob struct {
|
||||
AccountID uint `json:"account_id"`
|
||||
DocumentID uint `json:"document_id"`
|
||||
}
|
||||
|
||||
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)
|
||||
}
|
||||
|
||||
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
|
||||
}
|
||||
Reference in New Issue
Block a user