omnix-sopiga/worker/broadcast_worker.go
2026-08-07 12:36:57 +07:00

86 lines
1.8 KiB
Go

package worker
import (
"context"
"time"
fwlogger "repository.promas.id/prana/go-dw-framework/logger"
"repository.promas.id/prana/omnix-sopiga/domain"
"repository.promas.id/prana/omnix-sopiga/service"
)
const maxConcurrency = 5
type BroadcastWorker struct {
service *service.Service
logger *fwlogger.Logger
checkInterval time.Duration
batchSize int
}
func New(svc *service.Service, log *fwlogger.Logger, checkInterval time.Duration, batchSize int) *BroadcastWorker {
return &BroadcastWorker{
service: svc,
logger: log,
checkInterval: checkInterval,
batchSize: batchSize,
}
}
func (w *BroadcastWorker) Start(ctx context.Context) {
w.logger.Info("broadcast worker started")
ticker := time.NewTicker(w.checkInterval)
defer ticker.Stop()
for {
select {
case <-ctx.Done():
w.logger.Info("broadcast worker stopped")
return
case <-ticker.C:
w.processPendingRecords(ctx)
}
}
}
func (w *BroadcastWorker) processPendingRecords(ctx context.Context) {
broadcasts, err := w.service.FetchPending(ctx, w.batchSize)
if err != nil {
w.logger.Error("fetch pending failed", "error", err)
return
}
if len(broadcasts) == 0 {
return
}
w.logger.Info("processing pending broadcasts", "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.dispatchOne(ctx, b)
}(b)
}
for i := 0; i < len(broadcasts); i++ {
<-done
}
}
func (w *BroadcastWorker) dispatchOne(ctx context.Context, b *domain.Broadcast) {
if err := w.service.Dispatch(ctx, b); err != nil {
w.logger.Error("dispatch failed", "broadcast_id", b.ID, "error", err)
return
}
w.logger.Info("broadcast dispatched", "broadcast_id", b.ID)
}