Files
configcenter/internal/outbox/worker.go

68 lines
1.8 KiB
Go

package outbox
import (
"context"
"log/slog"
"time"
metricspkg "github.com/longpeng/configcenter/internal/metrics"
runtimepkg "github.com/longpeng/configcenter/internal/runtime"
"github.com/longpeng/configcenter/internal/store"
)
type Worker struct {
store store.Store
runtime runtimepkg.Store
interval time.Duration
batch int
maxRetry int
logger *slog.Logger
metrics *metricspkg.Collector
}
func New(repository store.Store, runtimeStore runtimepkg.Store, interval time.Duration, batch, maxRetry int, metrics *metricspkg.Collector, logger *slog.Logger) *Worker {
return &Worker{store: repository, runtime: runtimeStore, interval: interval, batch: batch, maxRetry: maxRetry, metrics: metrics, logger: logger}
}
func (w *Worker) Run(ctx context.Context) {
ticker := time.NewTicker(w.interval)
defer ticker.Stop()
w.process(ctx)
for {
select {
case <-ctx.Done():
return
case <-ticker.C:
w.process(ctx)
}
}
}
func (w *Worker) process(ctx context.Context) {
entries, err := w.store.ClaimOutbox(ctx, w.batch)
if err != nil {
w.logger.Error("claim outbox", "error", err)
return
}
for _, entry := range entries {
release, err := w.store.GetRelease(ctx, entry.ReleaseID)
if err == nil {
var revision int64
revision, err = w.runtime.Put(ctx, entry.EtcdKey, entry.Payload, release)
if err == nil {
err = w.store.MarkOutboxDone(ctx, entry.ID, entry.ReleaseID, revision)
if err == nil {
w.metrics.OutboxApplied()
}
}
}
if err != nil {
w.metrics.OutboxFailed()
w.logger.Warn("apply release outbox", "outbox_id", entry.ID, "release_id", entry.ReleaseID, "retry", entry.RetryCount, "error", err)
if markErr := w.store.MarkOutboxFailed(ctx, entry.ID, entry.ReleaseID, err.Error(), w.maxRetry); markErr != nil {
w.logger.Error("mark outbox failed", "outbox_id", entry.ID, "error", markErr)
}
}
}
}