105 lines
3.4 KiB
Go
105 lines
3.4 KiB
Go
package multitable
|
|
|
|
import (
|
|
"context"
|
|
"errors"
|
|
"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))
|
|
// A missing source row cannot be repaired by retrying. This commonly occurs
|
|
// for audit events created before a source record was deleted.
|
|
if errors.Is(cause, gorm.ErrRecordNotFound) || 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)
|
|
}
|