Skip to content

Commit 0d2214c

Browse files
committed
perf(cache): batch reader lease releases
1 parent 4cc988f commit 0d2214c

10 files changed

Lines changed: 625 additions & 9 deletions

File tree

cmd/server/main.go

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -162,11 +162,15 @@ func main() {
162162
}()
163163

164164
shutdownErr := <-shutdownErrCh
165+
readerLeaseReleaseErr := cacheService.ShutdownReaderLeaseReleaser(shutdownCtx)
165166
mergeErr := <-mergeErrCh
166167
backgroundErr := <-backgroundErrCh
167168
if shutdownErr != nil {
168169
logger.Fatal().Err(shutdownErr).Msg("graceful shutdown failed")
169170
}
171+
if readerLeaseReleaseErr != nil {
172+
logger.Error().Err(readerLeaseReleaseErr).Msg("waiting for reader lease releases failed")
173+
}
170174
if mergeErr != nil {
171175
logger.Fatal().Err(mergeErr).Msg("waiting for in-flight merges failed")
172176
}

internal/cache/leased_reader.go

Lines changed: 5 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -14,6 +14,7 @@ type leasedReadCloser struct {
1414
stream io.ReadCloser
1515
lifecycle *storagelifecycle.Service
1616
lease *storagelifecycle.ReaderLease
17+
release func(string)
1718

1819
renewCancel context.CancelFunc
1920
renewDone chan struct{}
@@ -23,12 +24,13 @@ type leasedReadCloser struct {
2324
leaseErr error
2425
}
2526

26-
func newLeasedReadCloser(stream io.ReadCloser, lifecycle *storagelifecycle.Service, lease *storagelifecycle.ReaderLease) *leasedReadCloser {
27+
func newLeasedReadCloser(stream io.ReadCloser, lifecycle *storagelifecycle.Service, lease *storagelifecycle.ReaderLease, release func(string)) *leasedReadCloser {
2728
renewCtx, renewCancel := context.WithCancel(context.Background())
2829
reader := &leasedReadCloser{
2930
stream: stream,
3031
lifecycle: lifecycle,
3132
lease: lease,
33+
release: release,
3234
renewCancel: renewCancel,
3335
renewDone: make(chan struct{}),
3436
}
@@ -55,10 +57,8 @@ func (r *leasedReadCloser) Close() error {
5557
r.renewCancel()
5658
<-r.renewDone
5759
streamErr := r.stream.Close()
58-
cleanupCtx, cancel := context.WithTimeout(context.Background(), mergeCleanupTimeout)
59-
releaseErr := r.lifecycle.ReleaseReader(cleanupCtx, r.lease.ID)
60-
cancel()
61-
r.closeErr = errors.Join(streamErr, releaseErr, r.currentLeaseError())
60+
r.release(r.lease.ID)
61+
r.closeErr = errors.Join(streamErr, r.currentLeaseError())
6262
})
6363
return r.closeErr
6464
}
Lines changed: 207 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,207 @@
1+
package cache
2+
3+
import (
4+
"context"
5+
"sync"
6+
"time"
7+
8+
"github.com/MxOrbit/GitHubActionCacheServer/internal/storagelifecycle"
9+
"github.com/rs/zerolog"
10+
)
11+
12+
const (
13+
readerLeaseReleaseBatchSize = 64
14+
readerLeaseReleaseDelay = 5 * time.Millisecond
15+
readerLeaseReleaseMaxPending = 4096
16+
)
17+
18+
type readerLeaseReleaseFunc func(context.Context, []string) (int, error)
19+
20+
type readerLeaseReleaser struct {
21+
release readerLeaseReleaseFunc
22+
logger zerolog.Logger
23+
batchSize int
24+
flushDelay time.Duration
25+
26+
mu sync.Mutex
27+
accepting bool
28+
workerStarted bool
29+
queue chan string
30+
workerCancel context.CancelFunc
31+
workerDone chan struct{}
32+
}
33+
34+
func newReaderLeaseReleaser(lifecycle *storagelifecycle.Service, logger zerolog.Logger) *readerLeaseReleaser {
35+
return newReaderLeaseReleaserWithOptions(
36+
lifecycle.ReleaseReaders,
37+
logger,
38+
readerLeaseReleaseBatchSize,
39+
readerLeaseReleaseDelay,
40+
readerLeaseReleaseMaxPending,
41+
)
42+
}
43+
44+
func newReaderLeaseReleaserWithOptions(
45+
release readerLeaseReleaseFunc,
46+
logger zerolog.Logger,
47+
batchSize int,
48+
flushDelay time.Duration,
49+
maxPending int,
50+
) *readerLeaseReleaser {
51+
if batchSize < 1 {
52+
batchSize = 1
53+
}
54+
if flushDelay <= 0 {
55+
flushDelay = readerLeaseReleaseDelay
56+
}
57+
if maxPending < batchSize {
58+
maxPending = batchSize
59+
}
60+
return &readerLeaseReleaser{
61+
release: release,
62+
logger: logger,
63+
batchSize: batchSize,
64+
flushDelay: flushDelay,
65+
accepting: true,
66+
queue: make(chan string, maxPending),
67+
}
68+
}
69+
70+
func (r *readerLeaseReleaser) Enqueue(leaseID string) bool {
71+
if leaseID == "" {
72+
return true
73+
}
74+
75+
r.mu.Lock()
76+
defer r.mu.Unlock()
77+
if !r.accepting {
78+
return false
79+
}
80+
if !r.workerStarted {
81+
r.startWorkerLocked()
82+
}
83+
select {
84+
case r.queue <- leaseID:
85+
return true
86+
default:
87+
return false
88+
}
89+
}
90+
91+
func (r *readerLeaseReleaser) Shutdown(ctx context.Context) error {
92+
r.mu.Lock()
93+
if r.accepting {
94+
r.accepting = false
95+
if r.workerStarted {
96+
close(r.queue)
97+
}
98+
}
99+
if !r.workerStarted {
100+
r.mu.Unlock()
101+
return nil
102+
}
103+
done := r.workerDone
104+
cancel := r.workerCancel
105+
r.mu.Unlock()
106+
107+
select {
108+
case <-done:
109+
return nil
110+
case <-ctx.Done():
111+
cancel()
112+
<-done
113+
return ctx.Err()
114+
}
115+
}
116+
117+
func (r *readerLeaseReleaser) startWorkerLocked() {
118+
workerCtx, workerCancel := context.WithCancel(context.Background())
119+
r.workerStarted = true
120+
r.workerCancel = workerCancel
121+
r.workerDone = make(chan struct{})
122+
go r.runWorker(workerCtx, r.workerDone)
123+
}
124+
125+
func (r *readerLeaseReleaser) runWorker(ctx context.Context, done chan struct{}) {
126+
defer close(done)
127+
batch := make([]string, 0, r.batchSize)
128+
timer := time.NewTimer(r.flushDelay)
129+
stopTimer(timer)
130+
var timerC <-chan time.Time
131+
defer timer.Stop()
132+
133+
for {
134+
select {
135+
case <-ctx.Done():
136+
r.logAbandoned(len(batch)+len(r.queue), ctx.Err())
137+
return
138+
case leaseID, ok := <-r.queue:
139+
if !ok {
140+
stopTimer(timer)
141+
if len(batch) > 0 {
142+
r.releaseBatch(ctx, batch)
143+
}
144+
return
145+
}
146+
batch = append(batch, leaseID)
147+
if len(batch) == 1 {
148+
timer.Reset(r.flushDelay)
149+
timerC = timer.C
150+
}
151+
if len(batch) == r.batchSize {
152+
stopTimer(timer)
153+
timerC = nil
154+
if !r.releaseBatch(ctx, batch) {
155+
return
156+
}
157+
batch = batch[:0]
158+
}
159+
case <-timerC:
160+
timerC = nil
161+
if !r.releaseBatch(ctx, batch) {
162+
return
163+
}
164+
batch = batch[:0]
165+
}
166+
}
167+
}
168+
169+
func (r *readerLeaseReleaser) releaseBatch(ctx context.Context, batch []string) bool {
170+
cleanupCtx, cancel := context.WithTimeout(ctx, mergeCleanupTimeout)
171+
_, err := r.release(cleanupCtx, batch)
172+
cancel()
173+
if err != nil {
174+
if ctx.Err() != nil {
175+
r.logAbandoned(len(batch)+len(r.queue), ctx.Err())
176+
return false
177+
}
178+
r.logger.Error().
179+
Err(err).
180+
Int("reader_lease_count", len(batch)).
181+
Msg("cache reader lease batch release failed")
182+
}
183+
if ctx.Err() != nil {
184+
r.logAbandoned(len(r.queue), ctx.Err())
185+
return false
186+
}
187+
return true
188+
}
189+
190+
func (r *readerLeaseReleaser) logAbandoned(count int, err error) {
191+
if count == 0 {
192+
return
193+
}
194+
r.logger.Warn().
195+
Err(err).
196+
Int("reader_lease_count", count).
197+
Msg("cache reader lease releases abandoned during shutdown")
198+
}
199+
200+
func stopTimer(timer *time.Timer) {
201+
if !timer.Stop() {
202+
select {
203+
case <-timer.C:
204+
default:
205+
}
206+
}
207+
}

0 commit comments

Comments
 (0)