package multitable import ( "context" "fmt" "time" "git.iwork-ai.com/xdc/iqudo-top1/internal/config" "git.iwork-ai.com/xdc/iqudo-top1/internal/model" "go.uber.org/zap" "gorm.io/gorm" "gorm.io/gorm/clause" ) type Worker struct { db *gorm.DB projector *Projector config config.MultiTableConfig log *zap.Logger } func NewWorker(db *gorm.DB, cfg config.MultiTableConfig, log *zap.Logger) *Worker { return &Worker{db: db, projector: NewProjector(db, NewClient(cfg.BaseURL, cfg.APIKey, cfg.RequestTimeout), cfg.Tables), config: cfg, log: log.Named("multitable")} } func (w *Worker) Run(ctx context.Context) { // A prior process may have stopped while an event was claimed. w.db.Model(&model.MultiTableOutbox{}).Where("status = ?", "processing").Update("status", "pending") w.process(ctx) ticker := time.NewTicker(w.config.SyncInterval) defer ticker.Stop() for { select { case <-ctx.Done(): return case <-ticker.C: w.process(ctx) } } } func (w *Worker) process(ctx context.Context) { for _, event := range w.claim() { if err := w.projector.Project(ctx, event); err != nil { w.retry(event, err) continue } if err := w.db.Model(&model.MultiTableOutbox{}).Where("id = ?", event.ID).Updates(map[string]interface{}{"status": "succeeded", "last_error": ""}).Error; err != nil { w.log.Error("mark outbox succeeded", zap.Uint64("event_id", event.ID), zap.Error(err)) } } } func (w *Worker) claim() []model.MultiTableOutbox { items := make([]model.MultiTableOutbox, 0) err := w.db.Transaction(func(tx *gorm.DB) error { if err := tx.Clauses(clause.Locking{Strength: "UPDATE"}).Where("status = ? AND available_at <= ?", "pending", time.Now()).Order("id").Limit(w.config.BatchSize).Find(&items).Error; err != nil { return err } for i := range items { if err := tx.Model(&model.MultiTableOutbox{}).Where("id = ? AND status = ?", items[i].ID, "pending").Updates(map[string]interface{}{"status": "processing", "attempts": gorm.Expr("attempts + 1")}).Error; err != nil { return err } items[i].Attempts++ } return nil }) if err != nil { w.log.Error("claim multitable outbox", zap.Error(err)) return nil } return items } func (w *Worker) retry(event model.MultiTableOutbox, cause error) { status := "pending" availableAt := time.Now().Add(backoff(event.Attempts)) if event.Attempts >= w.config.MaxAttempts { status = "failed" } message := truncate(cause.Error(), 1024) if err := w.db.Model(&model.MultiTableOutbox{}).Where("id = ?", event.ID).Updates(map[string]interface{}{"status": status, "available_at": availableAt, "last_error": message}).Error; err != nil { w.log.Error("reschedule multitable event", zap.Uint64("event_id", event.ID), zap.Error(err)) return } w.log.Warn("multitable projection failed", zap.Uint64("event_id", event.ID), zap.String("resource", event.Resource), zap.String("action", event.Action), zap.String("status", status), zap.Error(cause)) } func backoff(attempt int) time.Duration { if attempt < 1 { attempt = 1 } if attempt > 8 { attempt = 8 } return time.Second * time.Duration(1<<(attempt-1)) } func (w *Worker) String() string { return fmt.Sprintf("multitable worker (every %s)", w.config.SyncInterval) }