86 lines
1.9 KiB
Go
86 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"
|
||
|
|
)
|
||
|
|
|
||
|
|
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)
|
||
|
|
}
|