Compare commits

...

4 Commits

Author SHA1 Message Date
Blake 97af0f829b ops: Assert Host allowlist reject in validate-check
CI / vulncheck (pull_request) Successful in 16s
CI / check-and-test (pull_request) Successful in 53s
CI / vulncheck (push) Successful in 22s
CI / check-and-test (push) Successful in 1m0s
Release Tag / release (push) Successful in 15s
Empty-upstream Host allowlist is load-bearing: if that gate regresses,
the cache becomes an open LAN reverse proxy again. Unit tests already
cover hostAllowedForDirectFetch, but make validate-check did not probe
the live reject path.

Extend validate-check to GET a depot-like path with Host: evil.example
and a Steam User-Agent, requiring HTTP 400 Invalid URL. Document the
expected reject in README and the validate-config comment.

Fixes #37
2026-09-14 13:39:32 +00:00
pike 25622cbf25 ops: Fix goimports on fair-share bandwidth PR
CI / vulncheck (pull_request) Successful in 14s
CI / check-and-test (pull_request) Successful in 51s
CI / vulncheck (push) Successful in 14s
CI / check-and-test (push) Successful in 50s
Release Tag / release (push) Successful in 20s
2026-09-09 20:27:32 +00:00
pike b7710de0ca cache: Per-client fair-share bandwidth on table uplink
CI / vulncheck (pull_request) Successful in 20s
CI / check-and-test (pull_request) Failing after 21s
2026-09-09 20:22:23 +00:00
pike 50ca0a071c ops: Disk tier occupancy on /metrics
CI / vulncheck (pull_request) Successful in 19s
CI / check-and-test (pull_request) Successful in 47s
CI / vulncheck (push) Successful in 16s
CI / check-and-test (push) Successful in 46s
Release Tag / release (push) Successful in 14s
Operators could see attach-ready and capacity-pressure events but not how
full the configured disk (or memory) tier was without reading filesystems.
Expose size/capacity gauges and disk_cache_full_ratio next to disk_tier_ready.
Capacity is a config read, so it stays available while attach is pending.
2026-09-09 16:06:12 +00:00
16 changed files with 703 additions and 25 deletions
+22 -3
View File
@@ -62,7 +62,7 @@ validate run-validation: build clean-disk ## Start steamcache2 on :80 with small
fi; \
exec "$$BINARY" --config docs/examples/validate-config.yaml --log-level info
validate-check: ## Curl local /metrics (full dump + hit/miss + upstream/write/rate fields) and /lancache-heartbeat (default :80)
validate-check: ## Curl local /metrics (full dump + hit/miss + upstream/write/rate fields), /lancache-heartbeat, and non-Steam Host reject probe (empty upstream; default :80)
@echo "=== http://localhost/metrics ==="
@metrics=$$(curl -sf --max-time 5 http://localhost/metrics) || { \
echo "ERROR: could not fetch http://localhost/metrics"; \
@@ -84,7 +84,26 @@ validate-check: ## Curl local /metrics (full dump + hit/miss + upstream/write/ra
echo "$$hb" | grep -q '204' && echo "$$hb" | grep -qi 'X-LanCache-Processed-By' || { \
echo "ERROR: expected HTTP 204 and X-LanCache-Processed-By on /lancache-heartbeat"; \
exit 1; \
}
}; \
echo ""; \
echo "=== http://localhost/depot/allowlist-probe/chunk (Host: evil.example, Steam UA; expect 400 reject) ==="; \
allowlist_body=$$(mktemp); \
allowlist_code=$$(curl -s --max-time 5 -o "$$allowlist_body" -w '%{http_code}' -H 'Host: evil.example' -H 'User-Agent: Valve/Steam HTTP Client 1.0' http://localhost/depot/allowlist-probe/chunk) || { \
rm -f "$$allowlist_body"; \
echo "ERROR: could not probe http://localhost/depot/allowlist-probe/chunk"; \
echo "Is steamcache2 running on the default listen address :80?"; \
exit 1; \
}; \
printf 'HTTP %s\n' "$$allowlist_code"; \
printf '%s\n' "$$(cat "$$allowlist_body")"; \
if [ "$$allowlist_code" != "400" ] || ! grep -q 'Invalid URL' "$$allowlist_body"; then \
rm -f "$$allowlist_body"; \
echo "ERROR: expected HTTP 400 'Invalid URL' rejecting non-Steam Host with empty upstream (got $$allowlist_code)"; \
echo "Host allowlist gate regressed: steamcache2 may act as an open LAN reverse proxy."; \
exit 1; \
fi; \
rm -f "$$allowlist_body"; \
echo "Host allowlist reject OK (non-Steam Host -> 400)"
validate-kill: ## Kill leftover steamcache2 processes (safer, checks process name)
@echo "Looking for steamcache2 processes on common validation ports (80 is primary)..."
@@ -133,7 +152,7 @@ help: ## Show this help message
@echo " clean-disk Remove disk cache"
@echo " bench Run low-level VFS microbenchmarks"
@echo " validate / run-validation Start server on :80 (builds, auto-setcaps fresh binary, then runs as normal user, cleans disk cache first)"
@echo " validate-check Curl local /metrics (full dump + hit/miss + upstream/write/rate fields) and /lancache-heartbeat (default :80)"
@echo " validate-check Curl local /metrics (full dump + hit/miss + upstream/write/rate fields), /lancache-heartbeat, and non-Steam Host reject probe (empty upstream; default :80)"
@echo " setcap Explicitly set cap on current build (for port 80 use outside validate)"
@echo " validate-kill Kill leftover steamcache2 processes (safer)"
@echo " prefill Download latest SteamPrefill into bin/steam-prefill/SteamPrefill (gitignored)"
+15 -2
View File
@@ -76,7 +76,7 @@ curl -s http://localhost/metrics
curl -s -i http://localhost/lancache-heartbeat
```
`make validate-check` prints the full `/metrics` dump, highlights hit/miss plus `upstream_errors` / `cache_write_failures` / `rate_limited`, and curls `/lancache-heartbeat`. Read these fields:
`make validate-check` prints the full `/metrics` dump, highlights hit/miss plus `upstream_errors` / `cache_write_failures` / `rate_limited`, and curls `/lancache-heartbeat`. It also asserts the empty-upstream Host allowlist: a non-Steam `Host` sent with a Steam `User-Agent` must be rejected with HTTP 400. Read these fields:
| Field | Meaning |
| --- | --- |
@@ -88,6 +88,9 @@ curl -s -i http://localhost/lancache-heartbeat
| `total_requests` / `errors` | Volume and failures |
| `upstream_errors` / `cache_write_failures` / `rate_limited` | Upstream pipe, cache write, and rate-limit pressure (Quick check highlights these next to hit/miss) |
| `disk_tier_ready` | `0` while disk slow-tier attach pending; `1` when attached, or when no disk configured (N/A — not waiting) |
| `memory_cache_size` / `disk_cache_size` | Current cache occupancy per tier (bytes) |
| `memory_cache_capacity` / `disk_cache_capacity` | Configured capacity per tier (bytes); `disk_cache_capacity` is `0` when no disk is configured |
| `disk_cache_full_ratio` | `disk_cache_size / disk_cache_capacity` in [0,1]; 0 when no disk is configured or capacity is 0. Tells you "95% full" vs "barely filled" without reading the filesystem |
| `capacity_pressure_events` | Soft eviction under the memory or disk cap, and/or disk Create/Write/Mkdir hitting ENOSPC (volume full). Distinct from cold-cache misses and from the existing `evictions` counter. Logs `tier` (memory or disk) and `reason` (eviction or enospc). |
A first pass through new content is mostly misses (`hit_rate` near 0). Repeat the same content and `cache_hits` / `hit_rate` should rise.
@@ -181,7 +184,7 @@ curl -s http://localhost/metrics
curl -s -i http://localhost/lancache-heartbeat
```
`make validate-check` prints the full `/metrics` dump, highlights hit/miss fields, and curls `/lancache-heartbeat`. Look for:
`make validate-check` prints the full `/metrics` dump, highlights hit/miss fields, and curls `/lancache-heartbeat`. It also asserts a non-Steam Host is rejected with HTTP 400 when upstream is empty. Look for:
- High cache hit rate after the warmup pass (`cache_hits`, `hit_rate`, plus `memory_cache_hits` / `disk_cache_hits`)
- Non-zero `coalesced` and `disk` activity
- Zero unexpected `errors`, and quiet `upstream_errors` / `cache_write_failures` / `rate_limited`
@@ -233,6 +236,10 @@ While most configuration is done via the YAML file, some runtime options are sti
./steamcache2 --max-concurrent-requests 8
./steamcache2 --max-requests-per-client 4
# Table-tier uplink shaping (empty/0 = use config / disabled)
./steamcache2 --uplink-bandwidth 10MB
./steamcache2 --max-bytes-per-client-per-sec 2500000
# Show help
./steamcache2 --help
```
@@ -249,6 +256,11 @@ listen_address: :80
max_object_size: "0" # 0=unlimited; set e.g. "256MB" for response size DoS protection
trusted_proxies: [] # empty = safe (ignore XFF for rate limit); set CIDRs for trusted proxies
# Table-tier uplink bandwidth shaping (bytes/sec). Empty/0 = disabled (unlimited).
# Distinct from max_requests_per_client (concurrency). See "Table-tier uplink fair-share".
uplink_bandwidth: "" # e.g. "10MB" = 10e6 bytes/sec shared fairly across active clients
max_bytes_per_client_per_sec: 0 # optional absolute per-client cap; 0 = no absolute cap
# Cache configuration
cache:
# Memory cache settings
@@ -311,6 +323,7 @@ See `config.Validate()` and `steamcache.New` error paths. This ensures the LAN a
- Godoc on `disk.New` and `DiskFS.Size` expanded with the barrier/attach behavior.
- Startup logs: Info "Disk slow tier attach pending..." then later "Disk slow tier attached (...)" for disk-only and mixed modes.
- `/metrics` exposes `disk_tier_ready` 0/1 and stays responsive during attach (GetMetrics does not block on Size while pending).
- `/metrics` tier occupancy: `memory_cache_size` / `disk_cache_size` (bytes in use) next to `memory_cache_capacity` / `disk_cache_capacity` (configured capacity; `disk_cache_capacity` is 0 when no disk is configured), plus `disk_cache_full_ratio` (size/capacity in [0,1]). Capacity is a config read, so it is reported even while the disk attach is pending (size stays 0 until attach).
- `/lancache-heartbeat` header `X-SteamCache-Disk-Tier` mirrors that state.
- `/metrics` `capacity_pressure_events` counts times the cache dropped data under capacity pressure (soft eviction at the memory or disk cap, or disk Create/Write/Mkdir returning ENOSPC). Logs include `tier=memory|disk` and `reason=eviction|enospc` so operators can grep and tell this apart from a cold cache. The existing `evictions` counter is unchanged.
+12
View File
@@ -22,6 +22,8 @@ var (
maxConcurrentRequests int64
maxRequestsPerClient int64
uplinkBandwidth string
maxBytesPerClientPerSec int64
)
var rootCmd = &cobra.Command{
@@ -107,6 +109,12 @@ var rootCmd = &cobra.Command{
if maxRequestsPerClient > 0 {
finalMaxRequestsPerClient = maxRequestsPerClient
}
if uplinkBandwidth != "" {
cfg.UplinkBandwidth = uplinkBandwidth
}
if maxBytesPerClientPerSec > 0 {
cfg.MaxBytesPerClientPerSec = maxBytesPerClientPerSec
}
// Validate after loading and applying CLI overrides (fail fast, do not create default on validate error)
if err := cfg.Validate(); err != nil {
@@ -130,6 +138,8 @@ var rootCmd = &cobra.Command{
cfg.MaxObjectSize,
cfg.TrustedProxies,
cfg.Cache.NegativeTTL,
cfg.UplinkBandwidth,
cfg.MaxBytesPerClientPerSec,
)
if err != nil {
logger.Logger.Error().
@@ -170,4 +180,6 @@ func init() {
rootCmd.Flags().Int64Var(&maxConcurrentRequests, "max-concurrent-requests", 0, "Maximum concurrent requests (0 = use config file value)")
rootCmd.Flags().Int64Var(&maxRequestsPerClient, "max-requests-per-client", 0, "Maximum concurrent requests per client IP (0 = use config file value)")
rootCmd.Flags().StringVar(&uplinkBandwidth, "uplink-bandwidth", "", "Table uplink bandwidth bytes/sec human size e.g. 10MB (empty = use config; 0 disables)")
rootCmd.Flags().Int64Var(&maxBytesPerClientPerSec, "max-bytes-per-client-per-sec", 0, "Absolute per-client bytes/sec cap (0 = use config file value)")
}
+13
View File
@@ -19,6 +19,11 @@ type Config struct {
MaxConcurrentRequests int64 `yaml:"max_concurrent_requests" default:"200"`
MaxRequestsPerClient int64 `yaml:"max_requests_per_client" default:"5"`
// Table-tier uplink bandwidth shaping (bytes/sec). Distinct from MaxRequestsPerClient.
// Empty/"0" uplink and 0 max_bytes_per_client_per_sec = disabled (current unlimited behavior).
UplinkBandwidth string `yaml:"uplink_bandwidth"` // e.g. "10MB" via go-units = bytes/sec
MaxBytesPerClientPerSec int64 `yaml:"max_bytes_per_client_per_sec"` // absolute per-client cap; 0 = none
// Hardening limits (security/correctness)
MaxObjectSize string `yaml:"max_object_size" default:"0"` // 0=unlimited; e.g. "256MB" protects against OOM from huge/malicious upstream responses
TrustedProxies []string `yaml:"trusted_proxies"` // CIDR list; empty=never trust X-Forwarded-For (safe default). See README security notes.
@@ -186,6 +191,14 @@ func (c Config) Validate() error {
if c.MaxRequestsPerClient < 0 {
return fmt.Errorf("negative per-client limit not allowed")
}
if c.MaxBytesPerClientPerSec < 0 {
return fmt.Errorf("negative max_bytes_per_client_per_sec not allowed")
}
if c.UplinkBandwidth != "" && c.UplinkBandwidth != "0" {
if _, err := units.FromHumanSize(c.UplinkBandwidth); err != nil {
return fmt.Errorf("invalid uplink_bandwidth: %w", err)
}
}
if c.Cache.Memory.GCAlgorithm != "" {
switch c.Cache.Memory.GCAlgorithm {
+68
View File
@@ -211,3 +211,71 @@ func TestValidate(t *testing.T) {
})
}
}
func TestValidateUplinkBandwidth(t *testing.T) {
cases := []struct {
name string
mutate func(*Config)
wantErr bool
errSub string
}{
{
name: "empty uplink ok",
mutate: func(c *Config) {
c.UplinkBandwidth = ""
c.MaxBytesPerClientPerSec = 0
},
},
{
name: "zero uplink ok",
mutate: func(c *Config) {
c.UplinkBandwidth = "0"
},
},
{
name: "valid human size",
mutate: func(c *Config) {
c.UplinkBandwidth = "10MB"
},
},
{
name: "invalid uplink",
mutate: func(c *Config) {
c.UplinkBandwidth = "not-a-size"
},
wantErr: true,
errSub: "uplink_bandwidth",
},
{
name: "negative max bytes",
mutate: func(c *Config) {
c.MaxBytesPerClientPerSec = -1
},
wantErr: true,
errSub: "max_bytes_per_client_per_sec",
},
}
for _, tt := range cases {
t.Run(tt.name, func(t *testing.T) {
c := GetDefaultConfig()
tt.mutate(&c)
err := c.Validate()
if tt.wantErr {
if err == nil {
t.Fatalf("Validate() error = nil, wantErr")
}
if tt.errSub != "" && !contains(err.Error(), tt.errSub) {
t.Fatalf("Validate() error %q does not contain %q", err.Error(), tt.errSub)
}
return
}
if err != nil {
t.Fatalf("Validate() unexpected error: %v", err)
}
})
}
}
func contains(s, sub string) bool {
return strings.Contains(s, sub)
}
+2
View File
@@ -28,6 +28,7 @@
#
# After the benchmark run, inspect with:
# make validate-check # full /metrics + hit/miss fields + /lancache-heartbeat
# # also asserts a non-Steam Host is rejected (400) while upstream is empty
# # or, manually:
# curl -s http://localhost/metrics
# curl -s -i http://localhost/lancache-heartbeat # GET, not HEAD
@@ -38,6 +39,7 @@
listen_address: :80
max_concurrent_requests: 1000
# uplink_bandwidth / max_bytes_per_client_per_sec default off (unlimited)
max_requests_per_client: 10
max_object_size: "0" # unlimited for validation (real Steam files can be large)
+1
View File
@@ -9,6 +9,7 @@ require (
github.com/spf13/cobra v1.8.1
golang.org/x/sync v0.16.0
golang.org/x/sys v0.12.0
golang.org/x/time v0.16.0
gopkg.in/yaml.v3 v3.0.1
)
+2
View File
@@ -27,6 +27,8 @@ golang.org/x/sys v0.0.0-20220811171246-fbc7d0a398ab/go.mod h1:oPkhp1MJrh7nUepCBc
golang.org/x/sys v0.6.0/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg=
golang.org/x/sys v0.12.0 h1:CM0HF96J0hcLAwsHPJZjfdNzs0gftsLfgKt57wWHJ0o=
golang.org/x/sys v0.12.0/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg=
golang.org/x/time v0.16.0 h1:vMb6ptszcQMkcwiRTAuNNU50gom6++Q/6gY2hDM6VDE=
golang.org/x/time v0.16.0/go.mod h1:rVKOqvZeKvrDKTQiAHJ7wmwP0RzleSphoEA9RcdLA0s=
gopkg.in/check.v1 v0.0.0-20161208181325-20d25e280405 h1:yhCVgyC4o1eVCa2tZl7eS0r+SDo693bJlVdllGtEeKM=
gopkg.in/check.v1 v0.0.0-20161208181325-20d25e280405/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0=
gopkg.in/yaml.v3 v3.0.1 h1:fxVm/GzAzEWqLHuvctI91KS9hhNmmWOoWu0XTYJS7CA=
+158
View File
@@ -0,0 +1,158 @@
// steamcache/bandwidth.go
// Per-client fair-share / absolute bandwidth shaping for table-tier uplink.
// Distinct from max_requests_per_client concurrency (semaphores in ratelimit.go).
package steamcache
import (
"context"
"net/http"
"sync"
"golang.org/x/time/rate"
)
const bandwidthWriteChunk = 32 * 1024
// clientBandwidthLimiter fair-shares uplinkBytesPerSec among active clients and/or
// applies an absolute per-client bytes/sec cap. Both 0 disables shaping.
type clientBandwidthLimiter struct {
uplinkBytesPerSec int64
absoluteCap int64
mu sync.Mutex
active map[string]int // refcount of in-flight shaped responses per client IP
limiters map[string]*rate.Limiter
}
func newClientBandwidthLimiter(uplinkBytesPerSec, absoluteCap int64) *clientBandwidthLimiter {
if uplinkBytesPerSec < 0 {
uplinkBytesPerSec = 0
}
if absoluteCap < 0 {
absoluteCap = 0
}
return &clientBandwidthLimiter{
uplinkBytesPerSec: uplinkBytesPerSec,
absoluteCap: absoluteCap,
active: make(map[string]int),
limiters: make(map[string]*rate.Limiter),
}
}
func (b *clientBandwidthLimiter) enabled() bool {
return b != nil && (b.uplinkBytesPerSec > 0 || b.absoluteCap > 0)
}
// acquire registers clientIP as actively downloading and returns its limiter
// (nil if shaping disabled) plus a release func that must be deferred.
func (b *clientBandwidthLimiter) acquire(clientIP string) (*rate.Limiter, func()) {
if !b.enabled() {
return nil, func() {}
}
b.mu.Lock()
b.active[clientIP]++
lim := b.ensureLimiterLocked(clientIP)
b.recomputeRatesLocked()
b.mu.Unlock()
var once sync.Once
release := func() {
once.Do(func() {
b.mu.Lock()
defer b.mu.Unlock()
if n := b.active[clientIP]; n <= 1 {
delete(b.active, clientIP)
} else {
b.active[clientIP] = n - 1
}
b.recomputeRatesLocked()
})
}
return lim, release
}
func (b *clientBandwidthLimiter) ensureLimiterLocked(clientIP string) *rate.Limiter {
if lim, ok := b.limiters[clientIP]; ok {
return lim
}
// Start with a placeholder; recomputeRatesLocked sets the real rate.
lim := rate.NewLimiter(rate.Limit(1), 1)
b.limiters[clientIP] = lim
return lim
}
func (b *clientBandwidthLimiter) recomputeRatesLocked() {
n := len(b.active)
if n == 0 {
return
}
var fair int64
if b.uplinkBytesPerSec > 0 {
fair = b.uplinkBytesPerSec / int64(n)
if fair < 1 {
fair = 1
}
}
for ip := range b.active {
r := fair
if b.absoluteCap > 0 {
if r == 0 || b.absoluteCap < r {
r = b.absoluteCap
}
}
if r < 1 {
r = 1
}
lim := b.ensureLimiterLocked(ip)
burst := int(r)
if burst < bandwidthWriteChunk {
burst = bandwidthWriteChunk
}
// Cap burst to avoid huge memory spikes on huge uplinks.
if burst > 4*bandwidthWriteChunk {
burst = 4 * bandwidthWriteChunk
}
lim.SetLimit(rate.Limit(r))
lim.SetBurst(burst)
}
}
// limitedResponseWriter rate-limits response body Write calls. Headers/WriteHeader
// are unlimited. Implements http.ResponseWriter (+ optional Flusher/Hijacker passthrough
// is intentionally omitted — SteamCache body path only needs Write).
type limitedResponseWriter struct {
http.ResponseWriter
lim *rate.Limiter
ctx context.Context
}
func (w *limitedResponseWriter) Write(p []byte) (int, error) {
if w.lim == nil || len(p) == 0 {
return w.ResponseWriter.Write(p)
}
ctx := w.ctx
if ctx == nil {
ctx = context.Background()
}
total := 0
for total < len(p) {
chunk := p[total:]
if len(chunk) > bandwidthWriteChunk {
chunk = chunk[:bandwidthWriteChunk]
}
if err := w.lim.WaitN(ctx, len(chunk)); err != nil {
return total, err
}
n, err := w.ResponseWriter.Write(chunk)
total += n
if err != nil {
return total, err
}
}
return total, nil
}
// Unwrap exposes the underlying ResponseWriter for http.ResponseController etc.
func (w *limitedResponseWriter) Unwrap() http.ResponseWriter {
return w.ResponseWriter
}
+128
View File
@@ -0,0 +1,128 @@
package steamcache
import (
"bytes"
"context"
"net/http"
"net/http/httptest"
"testing"
"time"
)
func TestBandwidthFairShareRates(t *testing.T) {
b := newClientBandwidthLimiter(1000, 0)
lim1, rel1 := b.acquire("1.1.1.1")
defer rel1()
lim2, rel2 := b.acquire("2.2.2.2")
defer rel2()
if lim1 == nil || lim2 == nil {
t.Fatal("expected limiters")
}
// With 2 active clients, each should get ~500 bytes/sec.
got1 := float64(lim1.Limit())
got2 := float64(lim2.Limit())
if got1 < 400 || got1 > 600 || got2 < 400 || got2 > 600 {
t.Fatalf("fair-share rates = %v,%v want ~500", got1, got2)
}
rel2()
// After release, sole client should get full uplink.
lim1b, rel1b := b.acquire("1.1.1.1")
defer rel1b()
if float64(lim1b.Limit()) < 900 {
t.Fatalf("after release limit=%v want ~1000", lim1b.Limit())
}
}
func TestBandwidthAbsoluteCap(t *testing.T) {
b := newClientBandwidthLimiter(0, 250)
lim, rel := b.acquire("9.9.9.9")
defer rel()
if lim == nil {
t.Fatal("expected limiter")
}
if float64(lim.Limit()) != 250 {
t.Fatalf("limit=%v want 250", lim.Limit())
}
}
func TestBandwidthDisabled(t *testing.T) {
b := newClientBandwidthLimiter(0, 0)
lim, rel := b.acquire("9.9.9.9")
defer rel()
if lim != nil {
t.Fatal("expected nil limiter when disabled")
}
}
func TestLimitedResponseWriterShapes(t *testing.T) {
var buf bytes.Buffer
rec := httptest.NewRecorder()
// Use a custom writer sink via ResponseRecorder is fine; WaitN will delay.
lim := newClientBandwidthLimiter(0, 2000) // 2KB/s
l, rel := lim.acquire("127.0.0.1")
defer rel()
w := &limitedResponseWriter{ResponseWriter: rec, lim: l, ctx: context.Background()}
payload := bytes.Repeat([]byte("x"), 4000)
start := time.Now()
n, err := w.Write(payload)
elapsed := time.Since(start)
if err != nil {
t.Fatal(err)
}
if n != len(payload) {
t.Fatalf("wrote %d want %d", n, len(payload))
}
_ = buf
// 4000 bytes at 2000 B/s should take ~2s (allow slack for CI).
if elapsed < 1500*time.Millisecond {
t.Fatalf("elapsed %v too fast for 2KB/s shaping of 4KB", elapsed)
}
if elapsed > 8*time.Second {
t.Fatalf("elapsed %v unexpectedly slow", elapsed)
}
}
func TestServeHTTPBandwidthCap(t *testing.T) {
body := bytes.Repeat([]byte("a"), 3000)
upstream := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
w.Header().Set("Content-Type", "application/octet-stream")
w.Header().Set("Content-Length", "3000")
w.WriteHeader(200)
_, _ = w.Write(body)
}))
t.Cleanup(upstream.Close)
sc, err := NewWithOptions(Options{
Address: "127.0.0.1:0",
MemorySize: "1MB",
DiskSize: "0",
Upstream: upstream.URL,
MemoryGC: "lru",
DiskGC: "lru",
MaxConcurrentRequests: 20,
MaxRequestsPerClient: 10,
MaxObjectSize: "0",
MaxBytesPerClientPerSec: 1500, // 1.5KB/s
})
if err != nil {
t.Fatalf("NewWithOptions: %v", err)
}
t.Cleanup(func() { sc.Shutdown() })
req := httptest.NewRequest(http.MethodGet, "/depot/bw/chunk", nil)
req.Header.Set("User-Agent", "Valve/Steam HTTP Client 1.0")
rr := httptest.NewRecorder()
start := time.Now()
sc.ServeHTTP(rr, req)
elapsed := time.Since(start)
if rr.Code != 200 {
t.Fatalf("status=%d body=%s", rr.Code, rr.Body.String())
}
got := rr.Body.Bytes()
if len(got) != len(body) {
t.Fatalf("body len=%d want %d", len(got), len(body))
}
if elapsed < time.Second {
t.Fatalf("elapsed %v too fast for shaping", elapsed)
}
}
+14 -1
View File
@@ -38,10 +38,14 @@ type Options struct {
// NegativeTTL is a Go duration string for 404/410 negative cache entries.
// Empty defaults to 5m. "0" / "0s" disables storing negatives.
NegativeTTL string
// Table-tier uplink bandwidth shaping (bytes/sec). Empty/0 = disabled.
UplinkBandwidth string
MaxBytesPerClientPerSec int64
}
func NewWithOptions(o Options) (*SteamCache, error) {
return New(o.Address, o.MemorySize, o.DiskSize, o.DiskPath, o.Upstream, o.MemoryGC, o.DiskGC, o.MaxConcurrentRequests, o.MaxRequestsPerClient, o.MaxObjectSize, o.TrustedProxies, o.NegativeTTL)
return New(o.Address, o.MemorySize, o.DiskSize, o.DiskPath, o.Upstream, o.MemoryGC, o.DiskGC, o.MaxConcurrentRequests, o.MaxRequestsPerClient, o.MaxObjectSize, o.TrustedProxies, o.NegativeTTL, o.UplinkBandwidth, o.MaxBytesPerClientPerSec)
}
// handleSpecialEndpoints handles non-content paths (health, heartbeat, metrics) and
@@ -413,6 +417,15 @@ func (sc *SteamCache) ServeHTTP(w http.ResponseWriter, r *http.Request) {
return
}
// Per-client uplink bandwidth shaping (table-tier). Distinct from concurrency limits above.
if sc.bandwidth != nil && sc.bandwidth.enabled() {
lim, release := sc.bandwidth.acquire(clientIP)
defer release()
if lim != nil {
w = &limitedResponseWriter{ResponseWriter: w, lim: lim, ctx: r.Context()}
}
}
// Check if this is a request from a supported service
if service, isSupported := sc.detectService(r); isSupported {
// Cache key is the path only, never the Host: Steam rotates CDN hostnames
+42 -2
View File
@@ -32,6 +32,8 @@ type Metrics struct {
// Cache metrics
MemoryCacheSize int64
DiskCacheSize int64
MemoryCacheCapacity int64 // configured memory capacity (bytes)
DiskCacheCapacity int64 // configured disk capacity (bytes); 0 when no disk
MemoryCacheHits int64
DiskCacheHits int64
Promotions int64
@@ -137,6 +139,17 @@ func (m *Metrics) SetDiskCacheSize(size int64) {
atomic.StoreInt64(&m.DiskCacheSize, size)
}
// SetMemoryCacheCapacity sets the configured memory cache capacity in bytes.
func (m *Metrics) SetMemoryCacheCapacity(capacity int64) {
atomic.StoreInt64(&m.MemoryCacheCapacity, capacity)
}
// SetDiskCacheCapacity sets the configured disk cache capacity in bytes
// (0 when no disk is configured).
func (m *Metrics) SetDiskCacheCapacity(capacity int64) {
atomic.StoreInt64(&m.DiskCacheCapacity, capacity)
}
// SetDiskTierReady sets whether the disk slow tier is attached (1) or still pending (0).
// Memory-only (no disk) also uses 1 — meaning "not waiting on disk attach". Reset does not clear this.
func (m *Metrics) SetDiskTierReady(ready int64) {
@@ -248,6 +261,11 @@ func (m *Metrics) GetStats() *Stats {
serviceErrors[k] = v
}
memoryCacheSize := atomic.LoadInt64(&m.MemoryCacheSize)
diskCacheSize := atomic.LoadInt64(&m.DiskCacheSize)
memoryCacheCapacity := atomic.LoadInt64(&m.MemoryCacheCapacity)
diskCacheCapacity := atomic.LoadInt64(&m.DiskCacheCapacity)
return &Stats{
TotalRequests: totalRequests,
CacheHits: cacheHits,
@@ -262,8 +280,11 @@ func (m *Metrics) GetStats() *Stats {
AvgResponseTime: avgResponseTime,
TotalBytesServed: atomic.LoadInt64(&m.TotalBytesServed),
TotalBytesSaved: atomic.LoadInt64(&m.TotalBytesSaved),
MemoryCacheSize: atomic.LoadInt64(&m.MemoryCacheSize),
DiskCacheSize: atomic.LoadInt64(&m.DiskCacheSize),
MemoryCacheSize: memoryCacheSize,
DiskCacheSize: diskCacheSize,
MemoryCacheCapacity: memoryCacheCapacity,
DiskCacheCapacity: diskCacheCapacity,
DiskCacheFullRatio: diskFullRatio(diskCacheSize, diskCacheCapacity),
DiskTierReady: atomic.LoadInt64(&m.DiskTierReady),
MemoryCacheHits: atomic.LoadInt64(&m.MemoryCacheHits),
DiskCacheHits: atomic.LoadInt64(&m.DiskCacheHits),
@@ -279,6 +300,19 @@ func (m *Metrics) GetStats() *Stats {
}
}
// diskFullRatio is size / capacity clamped to [0,1].
// It is 0 when no disk is configured or the capacity is 0.
func diskFullRatio(size, capacity int64) float64 {
if size <= 0 || capacity <= 0 {
return 0
}
ratio := float64(size) / float64(capacity)
if ratio > 1 {
return 1
}
return ratio
}
// Reset resets all metrics to zero
func (m *Metrics) Reset() {
atomic.StoreInt64(&m.TotalRequests, 0)
@@ -330,6 +364,9 @@ type Stats struct {
MemoryCacheSize int64
DiskCacheSize int64
MemoryCacheCapacity int64 // configured memory capacity (bytes)
DiskCacheCapacity int64 // configured disk capacity (bytes); 0 when no disk
DiskCacheFullRatio float64 // disk_cache_size / disk_cache_capacity, clamped to [0,1]; 0 when no disk or capacity is 0
DiskTierReady int64
MemoryCacheHits int64
DiskCacheHits int64
@@ -382,7 +419,10 @@ func WriteText(w http.ResponseWriter, stats *Stats) {
writeInt(w, "total_bytes_served", "Total bytes sent to clients.", "counter", stats.TotalBytesServed)
writeInt(w, "total_bytes_saved", "Bytes served from cache instead of being re-downloaded from upstream.", "counter", stats.TotalBytesSaved)
writeInt(w, "memory_cache_size", "Current memory cache size in bytes.", "gauge", stats.MemoryCacheSize)
writeInt(w, "memory_cache_capacity", "Configured memory cache capacity in bytes.", "gauge", stats.MemoryCacheCapacity)
writeInt(w, "disk_cache_size", "Current disk cache size in bytes.", "gauge", stats.DiskCacheSize)
writeInt(w, "disk_cache_capacity", "Configured disk cache capacity in bytes; 0 when no disk is configured.", "gauge", stats.DiskCacheCapacity)
writeFloat(w, "disk_cache_full_ratio", "disk_cache_size / disk_cache_capacity in [0,1]; 0 when no disk or capacity is 0.", "gauge", "%.4f", stats.DiskCacheFullRatio)
writeInt(w, "disk_tier_ready", "1 if the disk tier is attached or no disk is configured; 0 while attach is pending.", "gauge", stats.DiskTierReady)
writeFloat(w, "uptime_seconds", "Process uptime in seconds.", "gauge", "%.2f", stats.Uptime.Seconds())
}
+61
View File
@@ -100,6 +100,9 @@ func TestWriteTextPrometheusExposition(t *testing.T) {
TotalBytesSaved: 50,
MemoryCacheSize: 8,
DiskCacheSize: 16,
MemoryCacheCapacity: 8,
DiskCacheCapacity: 32,
DiskCacheFullRatio: 0.5,
DiskTierReady: 1,
Uptime: 3 * time.Second,
}
@@ -120,6 +123,7 @@ func TestWriteTextPrometheusExposition(t *testing.T) {
}
gauges := []string{
"hit_rate", "avg_response_time_ms", "memory_cache_size", "disk_cache_size",
"memory_cache_capacity", "disk_cache_capacity", "disk_cache_full_ratio",
"disk_tier_ready", "uptime_seconds",
}
for _, name := range counters {
@@ -138,6 +142,15 @@ func TestWriteTextPrometheusExposition(t *testing.T) {
if !strings.Contains(body, "disk_tier_ready 1\n") {
t.Errorf("missing disk_tier_ready 1 sample: %q", body)
}
if !strings.Contains(body, "memory_cache_capacity 8\n") {
t.Errorf("missing memory_cache_capacity 8 sample: %q", body)
}
if !strings.Contains(body, "disk_cache_capacity 32\n") {
t.Errorf("missing disk_cache_capacity 32 sample: %q", body)
}
if !strings.Contains(body, "disk_cache_full_ratio 0.5000\n") {
t.Errorf("missing disk_cache_full_ratio 0.5000 sample: %q", body)
}
if !strings.Contains(body, "range_cache 3\n") {
t.Errorf("missing range_cache 3 sample: %q", body)
}
@@ -146,6 +159,54 @@ func TestWriteTextPrometheusExposition(t *testing.T) {
}
}
func TestDiskCacheFullRatioInGetStats(t *testing.T) {
t.Parallel()
cases := []struct {
name string
size int64
capacity int64
wantRatio float64
wantCapacity int64
}{
{"no disk (capacity 0)", 0, 0, 0, 0},
{"zero size with capacity", 0, 1024, 0, 1024},
{"half full", 512, 1024, 0.5, 1024},
{"exact full", 1024, 1024, 1, 1024},
{"size above capacity clamps to 1", 2048, 1024, 1, 1024},
}
for _, tc := range cases {
tc := tc
t.Run(tc.name, func(t *testing.T) {
t.Parallel()
m := NewMetrics()
m.SetDiskCacheSize(tc.size)
m.SetDiskCacheCapacity(tc.capacity)
st := m.GetStats()
if st.DiskCacheCapacity != tc.wantCapacity {
t.Fatalf("DiskCacheCapacity=%d, want %d", st.DiskCacheCapacity, tc.wantCapacity)
}
if st.DiskCacheFullRatio != tc.wantRatio {
t.Fatalf("DiskCacheFullRatio=%v, want %v", st.DiskCacheFullRatio, tc.wantRatio)
}
if st.DiskCacheFullRatio < 0 || st.DiskCacheFullRatio > 1 {
t.Fatalf("DiskCacheFullRatio=%v outside [0,1]", st.DiskCacheFullRatio)
}
})
}
// Memory capacity is a plain passthrough, and capacity survives Reset
// (re-derived by GetMetrics, like MemoryCacheSize/DiskTierReady).
m := NewMetrics()
m.SetMemoryCacheCapacity(4096)
if got := m.GetStats().MemoryCacheCapacity; got != 4096 {
t.Fatalf("MemoryCacheCapacity=%d, want 4096", got)
}
m.Reset()
if got := m.GetStats().MemoryCacheCapacity; got != 4096 {
t.Fatalf("MemoryCacheCapacity=%d after Reset, want 4096 (config snapshot, like size gauges)", got)
}
}
func assertHelpType(t *testing.T, body, name, typ string) {
t.Helper()
help := "# HELP " + name + " "
+1 -1
View File
@@ -170,7 +170,7 @@ func TestSerializeNegativeHeader(t *testing.T) {
}
func TestNewInvalidNegativeTTL(t *testing.T) {
sc, err := New("127.0.0.1:0", "1MB", "0", t.TempDir(), "", "lru", "lru", 10, 5, "0", nil, "not-a-duration")
sc, err := New("127.0.0.1:0", "1MB", "0", t.TempDir(), "", "lru", "lru", 10, 5, "0", nil, "not-a-duration", "", 0)
if err == nil {
if sc != nil {
sc.Shutdown()
+21 -1
View File
@@ -55,6 +55,9 @@ type SteamCache struct {
clientRateLimiter *clientRateLimiter
maxRequestsPerClient int64
// Per-client uplink bandwidth shaping (see bandwidth.go); nil/disabled = unlimited
bandwidth *clientBandwidthLimiter
// Hardening config fields (plumbed)
maxObjectSize int64
trustedProxies []string
@@ -85,7 +88,7 @@ const DefaultNegativeTTL = 5 * time.Minute
// negativeTTL is a Go duration string for 404/410 negative cache entries; empty means 5m.
// Callers must check the returned error.
// Prefer NewWithOptions (or config file) for forward compatibility. See README migration notes.
func New(address string, memorySize string, diskSize string, diskPath, upstream, memoryGC, diskGC string, maxConcurrentRequests int64, maxRequestsPerClient int64, maxObjectSize string, trustedProxies []string, negativeTTL string) (*SteamCache, error) {
func New(address string, memorySize string, diskSize string, diskPath, upstream, memoryGC, diskGC string, maxConcurrentRequests int64, maxRequestsPerClient int64, maxObjectSize string, trustedProxies []string, negativeTTL string, uplinkBandwidth string, maxBytesPerClientPerSec int64) (*SteamCache, error) {
memorysize, err := units.FromHumanSize(memorySize)
if err != nil {
return nil, fmt.Errorf("invalid memory size: %w", err)
@@ -114,6 +117,17 @@ func New(address string, memorySize string, diskSize string, diskPath, upstream,
return nil, err
}
var uplinkBytes int64
if uplinkBandwidth != "" && uplinkBandwidth != "0" {
uplinkBytes, err = units.FromHumanSize(uplinkBandwidth)
if err != nil {
return nil, fmt.Errorf("invalid uplink bandwidth: %w", err)
}
}
if maxBytesPerClientPerSec < 0 {
return nil, fmt.Errorf("negative max_bytes_per_client_per_sec not allowed")
}
c := cache.New()
var m *memory.MemoryFS
@@ -178,6 +192,7 @@ func New(address string, memorySize string, diskSize string, diskPath, upstream,
requestSemaphore: semaphore.NewWeighted(maxConcurrentRequests),
clientRateLimiter: newClientRateLimiter(maxRequestsPerClient),
maxRequestsPerClient: maxRequestsPerClient,
bandwidth: newClientBandwidthLimiter(uplinkBytes, maxBytesPerClientPerSec),
shutdownCh: make(chan struct{}),
// Hardening config plumbed
@@ -352,6 +367,11 @@ func (sc *SteamCache) Shutdown() {
func (sc *SteamCache) GetMetrics() *metrics.Stats {
if sc.memory != nil {
sc.metrics.SetMemoryCacheSize(sc.memory.Size())
sc.metrics.SetMemoryCacheCapacity(sc.memory.Capacity())
}
if sc.disk != nil {
// Capacity() is a plain config field — safe to read even while disk attach is pending.
sc.metrics.SetDiskCacheCapacity(sc.disk.Capacity())
}
// Skip disk.Size() while attach pending — Size() blocks on initDone and would hang /metrics.
if sc.disk != nil && sc.metrics.GetDiskTierReady() == 1 {
+141 -13
View File
@@ -28,7 +28,7 @@ import (
func TestCaching(t *testing.T) {
td := t.TempDir()
sc, err := New("localhost:8080", "1G", "1G", td, "", "lru", "lru", 200, 5, "0", nil, "")
sc, err := New("localhost:8080", "1G", "1G", td, "", "lru", "lru", 200, 5, "0", nil, "", "", 0)
if err != nil {
t.Fatalf("failed to create SteamCache: %v", err)
}
@@ -133,7 +133,7 @@ func TestCaching(t *testing.T) {
}
func TestCacheMissAndHit(t *testing.T) {
sc, err := New("localhost:8080", "1MB", "1G", t.TempDir(), "", "lru", "lru", 200, 5, "0", nil, "")
sc, err := New("localhost:8080", "1MB", "1G", t.TempDir(), "", "lru", "lru", 200, 5, "0", nil, "", "", 0)
if err != nil {
t.Fatalf("failed to create SteamCache: %v", err)
}
@@ -376,7 +376,7 @@ func TestServiceManagerExpandability(t *testing.T) {
// Removed hash calculation tests since we switched to lightweight validation
func TestSteamKeySharding(t *testing.T) {
sc, err := New("localhost:8080", "1MB", "1G", t.TempDir(), "", "lru", "lru", 200, 5, "0", nil, "")
sc, err := New("localhost:8080", "1MB", "1G", t.TempDir(), "", "lru", "lru", 200, 5, "0", nil, "", "", 0)
if err != nil {
t.Fatalf("failed to create SteamCache: %v", err)
}
@@ -483,7 +483,7 @@ func TestErrorTypes(t *testing.T) {
// TestMetrics tests the metrics functionality
func TestMetrics(t *testing.T) {
td := t.TempDir()
sc, err := New("localhost:8080", "1G", "1G", td, "", "lru", "lru", 200, 5, "0", nil, "")
sc, err := New("localhost:8080", "1G", "1G", td, "", "lru", "lru", 200, 5, "0", nil, "", "", 0)
if err != nil {
t.Fatalf("failed to create SteamCache: %v", err)
}
@@ -587,7 +587,7 @@ func newTestCacheWithFakeUpstream(t *testing.T, h http.HandlerFunc, mem, disk st
s := httptest.NewServer(h)
t.Cleanup(s.Close)
d := t.TempDir()
sc, err := New("127.0.0.1:0", mem, disk, d, s.URL, "lru", "lru", 200, 10, "0", nil, "")
sc, err := New("127.0.0.1:0", mem, disk, d, s.URL, "lru", "lru", 200, 10, "0", nil, "", "", 0)
if err != nil {
t.Fatalf("failed to create SteamCache: %v", err)
}
@@ -749,7 +749,7 @@ func TestErrorMetrics(t *testing.T) {
// Cover 503 capacity path + accounting skew: force Acquire err via canceled ctx.
// Asserts Errors+RateLimited inc, Total unchanged (per documented design in code comment).
tdCap := t.TempDir()
scCap, err := New("127.0.0.1:0", "1MB", "0", tdCap, "", "lru", "lru", 200, 5, "0", nil, "")
scCap, err := New("127.0.0.1:0", "1MB", "0", tdCap, "", "lru", "lru", 200, 5, "0", nil, "", "", 0)
if err != nil {
t.Fatalf("cap sc: %v", err)
}
@@ -813,7 +813,7 @@ func TestErrorMetrics(t *testing.T) {
func TestExpandedErrorMetrics(t *testing.T) {
t.Parallel()
td := t.TempDir()
sc, err := New("localhost:0", "1MB", "0", td, "", "lru", "lru", 10, 5, "0", nil, "")
sc, err := New("localhost:0", "1MB", "0", td, "", "lru", "lru", 10, 5, "0", nil, "", "", 0)
if err != nil {
t.Fatalf("create: %v", err)
}
@@ -903,7 +903,7 @@ func TestNewInvalidSizes(t *testing.T) {
}
for _, c := range cases {
t.Run(c.mem+"_"+c.disk, func(t *testing.T) {
sc, err := New("127.0.0.1:0", c.mem, c.disk, t.TempDir(), "", "lru", "lru", 10, 5, c.maxobj, nil, "")
sc, err := New("127.0.0.1:0", c.mem, c.disk, t.TempDir(), "", "lru", "lru", 10, 5, c.maxobj, nil, "", "", 0)
if err == nil {
t.Fatal("expected error for bad size, got nil")
}
@@ -924,7 +924,7 @@ func TestNewRunShutdownHygiene(t *testing.T) {
t.Skip("skips Run hygiene in -short per existing pattern")
}
d := t.TempDir()
sc, err := New("127.0.0.1:0", "1MB", "0", d, "", "lru", "lru", 10, 5, "0", nil, "")
sc, err := New("127.0.0.1:0", "1MB", "0", d, "", "lru", "lru", 10, 5, "0", nil, "", "", 0)
if err != nil {
t.Fatalf("new: %v", err)
}
@@ -1073,7 +1073,7 @@ func TestDiskOnlyDelayedAttach(t *testing.T) {
})
// mem=0, disk>0 -> pure disk delayed path (go func)
sc, err := New("localhost:0", "0", "10MB", diskPath, "", "lru", "lru", 10, 1, "0", nil, "")
sc, err := New("localhost:0", "0", "10MB", diskPath, "", "lru", "lru", 10, 1, "0", nil, "", "", 0)
if err != nil {
t.Fatalf("New disk-only: %v", err)
}
@@ -1145,7 +1145,7 @@ func TestDiskOnlyDelayedAttach(t *testing.T) {
// TestDiskTierSignalMemoryOnly covers memory-only mode: DiskTierReady=1 (N/A, not
// waiting on disk attach) and heartbeat header X-SteamCache-Disk-Tier: disabled.
func TestDiskTierSignalMemoryOnly(t *testing.T) {
sc, err := New("127.0.0.1:0", "1MB", "0", t.TempDir(), "", "lru", "lru", 10, 5, "0", nil, "")
sc, err := New("127.0.0.1:0", "1MB", "0", t.TempDir(), "", "lru", "lru", 10, 5, "0", nil, "", "", 0)
if err != nil {
t.Fatalf("New memory-only: %v", err)
}
@@ -1189,7 +1189,7 @@ func TestDiskTierSignalMixedPendingReady(t *testing.T) {
disk.ClearInitHold(diskPath)
})
sc, err := New("127.0.0.1:0", "1MB", "10MB", diskPath, "", "lru", "lru", 10, 1, "0", nil, "")
sc, err := New("127.0.0.1:0", "1MB", "10MB", diskPath, "", "lru", "lru", 10, 1, "0", nil, "", "", 0)
if err != nil {
t.Fatalf("New mixed: %v", err)
}
@@ -1341,7 +1341,7 @@ func TestHostAllowedForDirectFetch(t *testing.T) {
func TestDirectFetchRejectsNonSteamHost(t *testing.T) {
td := t.TempDir()
sc, err := New("127.0.0.1:0", "1MB", "0", td, "", "lru", "lru", 200, 5, "0", nil, "")
sc, err := New("127.0.0.1:0", "1MB", "0", td, "", "lru", "lru", 200, 5, "0", nil, "", "", 0)
if err != nil {
t.Fatalf("New: %v", err)
}
@@ -1497,3 +1497,131 @@ func TestCacheKeySharedAcrossCDNHostAliases(t *testing.T) {
t.Errorf("upstream fetched %d times across host aliases, want 1", got)
}
}
// TestGetMetricsCapacityGauges covers the tier-occupancy gauges: GetMetrics
// sets memory/disk capacity from the configured sizes, and WriteText emits
// memory_cache_capacity / disk_cache_capacity / disk_cache_full_ratio. During a
// pending disk attach, GetMetrics must return quickly (no disk.Size() call) and
// still report the configured disk capacity.
func TestGetMetricsCapacityGauges(t *testing.T) {
t.Run("memory-only", func(t *testing.T) {
sc, err := New("127.0.0.1:0", "1MB", "0", t.TempDir(), "", "lru", "lru", 10, 5, "0", nil, "", "", 0)
if err != nil {
t.Fatalf("New memory-only: %v", err)
}
t.Cleanup(func() { sc.Shutdown() })
st := sc.GetMetrics()
if st.MemoryCacheCapacity != 1000000 {
t.Errorf("MemoryCacheCapacity=%d, want 1000000 (configured 1MB)", st.MemoryCacheCapacity)
}
if st.DiskCacheCapacity != 0 {
t.Errorf("DiskCacheCapacity=%d, want 0 (no disk configured)", st.DiskCacheCapacity)
}
if st.DiskCacheFullRatio != 0 {
t.Errorf("DiskCacheFullRatio=%v, want 0 (no disk)", st.DiskCacheFullRatio)
}
rec := httptest.NewRecorder()
metrics.WriteText(rec, sc.GetMetrics())
body := rec.Body.String()
if !strings.Contains(body, "memory_cache_capacity 1000000\n") {
t.Errorf("WriteText missing memory_cache_capacity 1000000: %q", body)
}
if !strings.Contains(body, "disk_cache_capacity 0\n") {
t.Errorf("WriteText missing disk_cache_capacity 0: %q", body)
}
if !strings.Contains(body, "# TYPE disk_cache_full_ratio gauge") {
t.Errorf("WriteText missing # TYPE disk_cache_full_ratio gauge: %q", body)
}
})
t.Run("mixed pending attach reports capacity without Size", func(t *testing.T) {
td := t.TempDir()
diskPath := filepath.Join(td, "disk")
if err := os.MkdirAll(diskPath, 0755); err != nil {
t.Fatal(err)
}
hold := make(chan struct{})
var holdOnce sync.Once
closeHold := func() { holdOnce.Do(func() { close(hold) }) }
disk.RegisterInitHold(diskPath, hold)
t.Cleanup(func() {
closeHold()
disk.ClearInitHold(diskPath)
})
sc, err := New("127.0.0.1:0", "1MB", "10MB", diskPath, "", "lru", "lru", 10, 1, "0", nil, "", "", 0)
if err != nil {
t.Fatalf("New mixed: %v", err)
}
t.Cleanup(func() { sc.Shutdown() })
t.Cleanup(closeHold) // before Shutdown: attach is blocked in Size() until the hold closes
// Pending window is held open; if GetMetrics called disk.Size() it would
// block on the barrier, so a bounded wait proves non-blocking behavior.
done := make(chan *metrics.Stats, 1)
go func() { done <- sc.GetMetrics() }()
select {
case st := <-done:
if got := st.DiskTierReady; got != 0 {
t.Fatalf("immediate DiskTierReady=%d, want 0 (pending)", got)
}
if st.DiskCacheCapacity != 10000000 {
t.Errorf("pending DiskCacheCapacity=%d, want 10000000 (configured 10MB)", st.DiskCacheCapacity)
}
if st.MemoryCacheCapacity != 1000000 {
t.Errorf("pending MemoryCacheCapacity=%d, want 1000000 (configured 1MB)", st.MemoryCacheCapacity)
}
if st.DiskCacheFullRatio != 0 {
t.Errorf("pending DiskCacheFullRatio=%v, want 0 (size not reported while pending)", st.DiskCacheFullRatio)
}
case <-time.After(2 * time.Second):
t.Fatal("GetMetrics blocked during pending attach (must not call disk.Size())")
}
closeHold()
_ = sc.disk.Size()
deadline := time.Now().Add(2 * time.Second)
for time.Now().Before(deadline) {
if sc.GetMetrics().DiskTierReady == 1 {
break
}
time.Sleep(1 * time.Millisecond)
}
if got := sc.GetMetrics().DiskTierReady; got != 1 {
t.Fatalf("DiskTierReady=%d after barrier, want 1 (ready)", got)
}
// Post-attach writes prefer the slow (disk) tier, so a write produces a
// non-zero disk size and hence a non-zero occupancy ratio.
w, err := sc.vfs.Create("occupancy-key", 128)
if err != nil {
t.Fatalf("Create failed after attach: %v", err)
}
if _, err := w.Write(make([]byte, 128)); err != nil {
t.Fatalf("Write failed: %v", err)
}
if err := w.Close(); err != nil {
t.Fatalf("Close failed: %v", err)
}
st := sc.GetMetrics()
if st.DiskCacheCapacity != 10000000 {
t.Errorf("post-attach DiskCacheCapacity=%d, want 10000000", st.DiskCacheCapacity)
}
if st.DiskCacheSize <= 0 {
t.Errorf("post-attach DiskCacheSize=%d, want > 0 after a write", st.DiskCacheSize)
}
if st.DiskCacheFullRatio <= 0 || st.DiskCacheFullRatio > 1 {
t.Errorf("post-attach DiskCacheFullRatio=%v, want in (0,1]", st.DiskCacheFullRatio)
}
rec := httptest.NewRecorder()
metrics.WriteText(rec, st)
body := rec.Body.String()
if !strings.Contains(body, "disk_cache_capacity 10000000\n") {
t.Errorf("WriteText missing disk_cache_capacity 10000000: %q", body)
}
})
}