/
worker.go
94 lines (79 loc) · 2.26 KB
/
worker.go
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
package service
import (
"context"
"os"
"os/signal"
"syscall"
"time"
"github.com/allisson/postmand"
"go.uber.org/zap"
)
// Worker implements postmand.WorkerService interface.
type Worker struct {
deliveryRepository postmand.DeliveryRepository
logger *zap.Logger
pollingInterval time.Duration
isStop bool
}
func (w *Worker) run(ctx context.Context) {
for {
// Break forloop if isStop is true.
if w.isStop {
break
}
// Dispatch webhook.
deliveryAttempt, err := w.deliveryRepository.Dispatch(ctx)
if err != nil {
w.logger.Error("worker-dispatch-error", zap.Error(err))
time.Sleep(w.pollingInterval)
continue
}
if deliveryAttempt == nil {
time.Sleep(w.pollingInterval)
continue
}
// Log delivery attempt.
w.logger.Info(
"worker-delivery-attempt-created",
zap.String("id", deliveryAttempt.ID.String()),
zap.String("webhook_id", deliveryAttempt.WebhookID.String()),
zap.String("delivery_id", deliveryAttempt.DeliveryID.String()),
zap.Int("response_status_code", deliveryAttempt.ResponseStatusCode),
zap.Int("execution_duration", deliveryAttempt.ExecutionDuration),
zap.Bool("success", deliveryAttempt.Success),
)
}
w.logger.Info("worker-shutdown-completed")
}
// Run sending of webhooks until the Shutdown method is called.
func (w *Worker) Run(ctx context.Context) {
idleConnsClosed := make(chan struct{})
go func() {
sigint := make(chan os.Signal, 1)
// interrupt signal sent from terminal
signal.Notify(sigint, os.Interrupt)
// sigterm signal sent from kubernetes
signal.Notify(sigint, syscall.SIGTERM)
<-sigint
// We received an interrupt signal, shut down.
w.Shutdown(ctx)
close(idleConnsClosed)
}()
w.logger.Info("worker-started")
w.run(ctx)
<-idleConnsClosed
}
// Shutdown stops the forloop in Run method.
func (w *Worker) Shutdown(ctx context.Context) {
w.isStop = true
w.logger.Info("worker-shutdown-started")
}
// NewWorker will create an implementation of postmand.WorkerService.
func NewWorker(deliveryRepository postmand.DeliveryRepository, logger *zap.Logger, pollingInterval time.Duration) *Worker {
return &Worker{
deliveryRepository: deliveryRepository,
logger: logger,
pollingInterval: pollingInterval,
isStop: false,
}
}