Repository navigation
Expand file tree
/
Copy pathlifecycle_worker.go
More file actions
101 lines (90 loc) · 2.73 KB
/
Copy pathlifecycle_worker.go
File metadata and controls
101 lines (90 loc) · 2.73 KB
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
95
96
97
98
99
100
101
package ledger
import (
"context"
"runtime/debug"
"time"
log "github.com/xraph/go-utils/log"
)
// lifecycleWorker runs Advance every lifecycleInterval until Stop. The first
// run is one interval after Start. Each run has a deadline of one interval
// and is detached from the Start context, which a caller may cancel once
// Start returns: Stop is what ends the worker, and it also cancels a run that
// is in flight, so Stop never waits out a slow Advance.
func (l *Ledger) lifecycleWorker(ctx context.Context) {
defer l.wg.Done()
runCtx, cancel := context.WithCancel(context.WithoutCancel(ctx))
defer cancel()
// Stop closes stopChan; this turns that into a cancelled run context. The
// watcher is counted, so Stop waits for it too. Adding to the group from
// here is safe: this goroutine still holds its own count.
l.wg.Add(1)
go func() {
defer l.wg.Done()
select {
case <-l.stopChan:
cancel()
case <-runCtx.Done():
}
}()
ticker := time.NewTicker(l.lifecycleInterval)
defer ticker.Stop()
for {
if l.stopping() {
return
}
select {
case <-l.stopChan:
return
case <-ticker.C:
// A tick and a stop can be ready together, and select picks
// either. Stop wins: no run starts once Stop has been called.
if l.stopping() {
return
}
l.runLifecycle(runCtx)
}
}
}
// stopping reports whether Stop has been called.
func (l *Ledger) stopping() bool {
select {
case <-l.stopChan:
return true
default:
return false
}
}
// runLifecycle is one tick. A panic in Advance, which a store driver or a
// store wrapper can raise, is logged and ends only this run: the next tick
// starts clean. The recover does not reach plugin hooks. They run on their own
// goroutines (plugin.Registry callWithTimeout), so a panic in a hook is not
// caught here, and was not caught before the clock existed either.
func (l *Ledger) runLifecycle(ctx context.Context) {
defer func() {
if r := recover(); r != nil {
l.logger.Error("ledger: lifecycle run panicked",
log.Any("panic", r),
log.String("stack", string(debug.Stack())),
)
}
}()
ctx, cancel := context.WithTimeout(ctx, l.lifecycleInterval)
defer cancel()
report, err := l.Advance(ctx, l.now())
if err != nil {
// A run that Stop cancelled is not work left undone.
if l.stopping() {
l.logger.Debug("ledger: lifecycle run cut short by stop", log.Error(err))
} else {
l.logger.Warn("ledger: lifecycle run left work undone", log.Error(err))
}
}
if !report.Empty() {
l.logger.Info("ledger: lifecycle run",
log.Int("periods_advanced", len(report.PeriodsAdvanced)),
log.Int("cancels_enacted", len(report.CancelsEnacted)),
log.Int("trials_ended", len(report.TrialsEnded)),
log.Int("invoices_past_due", len(report.InvoicesPastDue)),
)
}
}