omnix-sopiga/services/omnix-broadcast/worker/delivery_sync_worker.go
2026-08-07 12:21:42 +07:00

84 lines
1.9 KiB
Go

package worker
import (
"context"
"time"
fwlogger "github.com/yourorg/go-dw-platform/framework/logger"
"github.com/yourorg/go-dw-platform/services/omnix-broadcast/domain"
"github.com/yourorg/go-dw-platform/services/omnix-broadcast/service"
)
type DeliverySyncWorker struct {
service *service.Service
logger *fwlogger.Logger
checkInterval time.Duration
batchSize int
}
func NewDeliverySync(svc *service.Service, log *fwlogger.Logger, checkInterval time.Duration, batchSize int) *DeliverySyncWorker {
return &DeliverySyncWorker{
service: svc,
logger: log,
checkInterval: checkInterval,
batchSize: batchSize,
}
}
func (w *DeliverySyncWorker) Start(ctx context.Context) {
w.logger.Info("delivery sync worker started")
ticker := time.NewTicker(w.checkInterval)
defer ticker.Stop()
for {
select {
case <-ctx.Done():
w.logger.Info("delivery sync worker stopped")
return
case <-ticker.C:
w.syncDispatchedRecords(ctx)
}
}
}
func (w *DeliverySyncWorker) syncDispatchedRecords(ctx context.Context) {
broadcasts, err := w.service.FetchDispatched(ctx, w.batchSize)
if err != nil {
w.logger.Error("fetch dispatched failed", "error", err)
return
}
if len(broadcasts) == 0 {
return
}
w.logger.Info("syncing delivery status", "count", len(broadcasts))
sem := make(chan struct{}, maxConcurrency)
done := make(chan struct{}, len(broadcasts))
for _, b := range broadcasts {
sem <- struct{}{}
go func(b *domain.Broadcast) {
defer func() {
<-sem
done <- struct{}{}
}()
w.syncOne(ctx, b)
}(b)
}
for i := 0; i < len(broadcasts); i++ {
<-done
}
}
func (w *DeliverySyncWorker) syncOne(ctx context.Context, b *domain.Broadcast) {
if err := w.service.SyncDelivery(ctx, b); err != nil {
w.logger.Error("sync delivery failed", "broadcast_id", b.ID, "error", err)
return
}
w.logger.Info("delivery status synced", "broadcast_id", b.ID)
}