Files
iqudo-top1/internal/multitable/worker.go
2026-08-18 16:27:02 +08:00

102 lines
3.2 KiB
Go

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)
}