Skip to content

Commit 3eb1928

Browse files
leodidoona-agent
andcommitted
perf(cache): add environment variables for S3 cache tuning
Adds configurable environment variables to tune S3 cache performance: - LEEWAY_S3_WORKER_COUNT: Number of concurrent workers for uploads and existence checks (default: 10) - LEEWAY_S3_DOWNLOAD_WORKERS: Number of concurrent download workers (default: 30) - LEEWAY_S3_RATE_LIMIT: S3 API request rate limit in requests/second (default: 100) - LEEWAY_S3_BURST_LIMIT: S3 API burst limit for rate limiting (default: 200) These allow tuning cache performance without code changes: export LEEWAY_S3_WORKER_COUNT=20 export LEEWAY_S3_DOWNLOAD_WORKERS=50 export LEEWAY_S3_RATE_LIMIT=300 export LEEWAY_S3_BURST_LIMIT=500 Co-authored-by: Ona <no-reply@ona.com>
1 parent 11e76f6 commit 3eb1928

2 files changed

Lines changed: 184 additions & 5 deletions

File tree

pkg/leeway/cache/remote/s3.go

Lines changed: 42 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -8,6 +8,7 @@ import (
88
"os"
99
"path/filepath"
1010
"runtime"
11+
"strconv"
1112
"strings"
1213
"sync"
1314
"time"
@@ -29,16 +30,26 @@ const (
2930
// defaultS3PartSize is the default part size for S3 multipart operations
3031
defaultS3PartSize = 5 * 1024 * 1024
3132
// defaultWorkerCount is the default number of concurrent workers for general operations
33+
// (existence checks, uploads). Can be overridden via LEEWAY_S3_WORKER_COUNT.
3234
defaultWorkerCount = 10
3335
// defaultDownloadWorkerCount is the number of concurrent workers for download operations
34-
// Higher than default to maximize download throughput
36+
// Higher than default to maximize download throughput.
37+
// Can be overridden via LEEWAY_S3_DOWNLOAD_WORKERS environment variable.
3538
defaultDownloadWorkerCount = 30
3639
// defaultRateLimit is the default rate limit for S3 API calls (requests per second)
40+
// Can be overridden via LEEWAY_S3_RATE_LIMIT environment variable.
3741
defaultRateLimit = 100
3842
// defaultBurstLimit is the default burst limit for S3 API calls
43+
// Can be overridden via LEEWAY_S3_BURST_LIMIT environment variable.
3944
defaultBurstLimit = 200
4045
// maxConcurrentOperations is the maximum number of concurrent goroutines for parallel operations
4146
maxConcurrentOperations = 50
47+
48+
// Environment variable names for S3 cache tuning
49+
envvarS3WorkerCount = "LEEWAY_S3_WORKER_COUNT"
50+
envvarS3DownloadWorkers = "LEEWAY_S3_DOWNLOAD_WORKERS"
51+
envvarS3RateLimit = "LEEWAY_S3_RATE_LIMIT"
52+
envvarS3BurstLimit = "LEEWAY_S3_BURST_LIMIT"
4253
)
4354

4455
// downloadResult represents the result of a download operation with proper error attribution
@@ -77,6 +88,16 @@ type S3Cache struct {
7788
semaphore chan struct{} // Semaphore for limiting concurrent operations
7889
}
7990

91+
// getEnvInt reads an integer from an environment variable, returning the default if not set or invalid
92+
func getEnvInt(envvar string, defaultVal int) int {
93+
if v := os.Getenv(envvar); v != "" {
94+
if parsed, err := strconv.Atoi(v); err == nil && parsed > 0 {
95+
return parsed
96+
}
97+
}
98+
return defaultVal
99+
}
100+
80101
// NewS3Cache creates a new S3 cache implementation
81102
func NewS3Cache(cfg *cache.RemoteConfig) (*S3Cache, error) {
82103
if cfg.BucketName == "" {
@@ -104,17 +125,33 @@ func NewS3Cache(cfg *cache.RemoteConfig) (*S3Cache, error) {
104125
}).Debug("SLSA verification enabled for cache")
105126
}
106127

107-
// Initialize rate limiter with default limits
108-
rateLimiter := rate.NewLimiter(rate.Limit(defaultRateLimit), defaultBurstLimit)
128+
// Read tuning parameters from environment variables (with defaults)
129+
workerCount := getEnvInt(envvarS3WorkerCount, defaultWorkerCount)
130+
downloadWorkers := getEnvInt(envvarS3DownloadWorkers, defaultDownloadWorkerCount)
131+
rateLimit := getEnvInt(envvarS3RateLimit, defaultRateLimit)
132+
burstLimit := getEnvInt(envvarS3BurstLimit, defaultBurstLimit)
133+
134+
// Log if non-default values are used
135+
if workerCount != defaultWorkerCount || downloadWorkers != defaultDownloadWorkerCount || rateLimit != defaultRateLimit || burstLimit != defaultBurstLimit {
136+
log.WithFields(log.Fields{
137+
"workerCount": workerCount,
138+
"downloadWorkers": downloadWorkers,
139+
"rateLimit": rateLimit,
140+
"burstLimit": burstLimit,
141+
}).Debug("S3 cache using custom tuning parameters")
142+
}
143+
144+
// Initialize rate limiter
145+
rateLimiter := rate.NewLimiter(rate.Limit(rateLimit), burstLimit)
109146

110147
// Initialize semaphore for goroutine limiting
111148
semaphore := make(chan struct{}, maxConcurrentOperations)
112149

113150
return &S3Cache{
114151
storage: storage,
115152
cfg: cfg,
116-
workerCount: defaultWorkerCount,
117-
downloadWorkerCount: defaultDownloadWorkerCount,
153+
workerCount: workerCount,
154+
downloadWorkerCount: downloadWorkers,
118155
slsaVerifier: slsaVerifier,
119156
rateLimiter: rateLimiter,
120157
semaphore: semaphore,
Lines changed: 142 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,142 @@
1+
package remote
2+
3+
import (
4+
"os"
5+
"testing"
6+
7+
"github.com/stretchr/testify/assert"
8+
)
9+
10+
func TestGetEnvInt(t *testing.T) {
11+
tests := []struct {
12+
name string
13+
envVar string
14+
envValue string
15+
defaultVal int
16+
expected int
17+
}{
18+
{
19+
name: "returns default when env var not set",
20+
envVar: "TEST_UNSET_VAR",
21+
envValue: "",
22+
defaultVal: 30,
23+
expected: 30,
24+
},
25+
{
26+
name: "returns parsed value when env var is valid integer",
27+
envVar: "TEST_VALID_INT",
28+
envValue: "50",
29+
defaultVal: 30,
30+
expected: 50,
31+
},
32+
{
33+
name: "returns default when env var is invalid integer",
34+
envVar: "TEST_INVALID_INT",
35+
envValue: "not-a-number",
36+
defaultVal: 30,
37+
expected: 30,
38+
},
39+
{
40+
name: "returns default when env var is zero",
41+
envVar: "TEST_ZERO_INT",
42+
envValue: "0",
43+
defaultVal: 30,
44+
expected: 30,
45+
},
46+
{
47+
name: "returns default when env var is negative",
48+
envVar: "TEST_NEGATIVE_INT",
49+
envValue: "-10",
50+
defaultVal: 30,
51+
expected: 30,
52+
},
53+
{
54+
name: "returns parsed value for large numbers",
55+
envVar: "TEST_LARGE_INT",
56+
envValue: "1000",
57+
defaultVal: 100,
58+
expected: 1000,
59+
},
60+
}
61+
62+
for _, tt := range tests {
63+
t.Run(tt.name, func(t *testing.T) {
64+
// Clean up env var before and after test
65+
os.Unsetenv(tt.envVar)
66+
defer os.Unsetenv(tt.envVar)
67+
68+
if tt.envValue != "" {
69+
os.Setenv(tt.envVar, tt.envValue)
70+
}
71+
72+
result := getEnvInt(tt.envVar, tt.defaultVal)
73+
assert.Equal(t, tt.expected, result)
74+
})
75+
}
76+
}
77+
78+
func TestS3CacheEnvVarTuning(t *testing.T) {
79+
// Save original env vars
80+
origWorkerCount := os.Getenv(envvarS3WorkerCount)
81+
origDownloadWorkers := os.Getenv(envvarS3DownloadWorkers)
82+
origRateLimit := os.Getenv(envvarS3RateLimit)
83+
origBurstLimit := os.Getenv(envvarS3BurstLimit)
84+
85+
// Restore after test
86+
defer func() {
87+
if origWorkerCount != "" {
88+
os.Setenv(envvarS3WorkerCount, origWorkerCount)
89+
} else {
90+
os.Unsetenv(envvarS3WorkerCount)
91+
}
92+
if origDownloadWorkers != "" {
93+
os.Setenv(envvarS3DownloadWorkers, origDownloadWorkers)
94+
} else {
95+
os.Unsetenv(envvarS3DownloadWorkers)
96+
}
97+
if origRateLimit != "" {
98+
os.Setenv(envvarS3RateLimit, origRateLimit)
99+
} else {
100+
os.Unsetenv(envvarS3RateLimit)
101+
}
102+
if origBurstLimit != "" {
103+
os.Setenv(envvarS3BurstLimit, origBurstLimit)
104+
} else {
105+
os.Unsetenv(envvarS3BurstLimit)
106+
}
107+
}()
108+
109+
t.Run("uses defaults when env vars not set", func(t *testing.T) {
110+
os.Unsetenv(envvarS3WorkerCount)
111+
os.Unsetenv(envvarS3DownloadWorkers)
112+
os.Unsetenv(envvarS3RateLimit)
113+
os.Unsetenv(envvarS3BurstLimit)
114+
115+
workerCount := getEnvInt(envvarS3WorkerCount, defaultWorkerCount)
116+
downloadWorkers := getEnvInt(envvarS3DownloadWorkers, defaultDownloadWorkerCount)
117+
rateLimit := getEnvInt(envvarS3RateLimit, defaultRateLimit)
118+
burstLimit := getEnvInt(envvarS3BurstLimit, defaultBurstLimit)
119+
120+
assert.Equal(t, defaultWorkerCount, workerCount)
121+
assert.Equal(t, defaultDownloadWorkerCount, downloadWorkers)
122+
assert.Equal(t, defaultRateLimit, rateLimit)
123+
assert.Equal(t, defaultBurstLimit, burstLimit)
124+
})
125+
126+
t.Run("uses custom values when env vars set", func(t *testing.T) {
127+
os.Setenv(envvarS3WorkerCount, "20")
128+
os.Setenv(envvarS3DownloadWorkers, "50")
129+
os.Setenv(envvarS3RateLimit, "300")
130+
os.Setenv(envvarS3BurstLimit, "500")
131+
132+
workerCount := getEnvInt(envvarS3WorkerCount, defaultWorkerCount)
133+
downloadWorkers := getEnvInt(envvarS3DownloadWorkers, defaultDownloadWorkerCount)
134+
rateLimit := getEnvInt(envvarS3RateLimit, defaultRateLimit)
135+
burstLimit := getEnvInt(envvarS3BurstLimit, defaultBurstLimit)
136+
137+
assert.Equal(t, 20, workerCount)
138+
assert.Equal(t, 50, downloadWorkers)
139+
assert.Equal(t, 300, rateLimit)
140+
assert.Equal(t, 500, burstLimit)
141+
})
142+
}

0 commit comments

Comments
 (0)