Compare commits
6 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 5e22f1054b | |||
| 46aa59e877 | |||
| 8eb31143e5 | |||
| 337741e06e | |||
| 2c8ec276e7 | |||
| 144a68dcd7 |
@@ -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 fields) and /lancache-heartbeat (default :80)
|
||||
validate-check: ## Curl local /metrics (full dump + hit/miss + upstream/write/rate fields) and /lancache-heartbeat (default :80)
|
||||
@echo "=== http://localhost/metrics ==="
|
||||
@metrics=$$(curl -sf --max-time 5 http://localhost/metrics) || { \
|
||||
echo "ERROR: could not fetch http://localhost/metrics"; \
|
||||
@@ -71,8 +71,8 @@ validate-check: ## Curl local /metrics (full dump + hit/miss fields) and /lancac
|
||||
}; \
|
||||
printf '%s\n' "$$metrics"; \
|
||||
echo ""; \
|
||||
echo "=== hit/miss fields ==="; \
|
||||
printf '%s\n' "$$metrics" | grep -E '^(total_requests|cache_hits|cache_misses|hit_rate|memory_cache_hits|disk_cache_hits|errors) ' || true; \
|
||||
echo "=== hit/miss + upstream/write/rate fields ==="; \
|
||||
printf '%s\n' "$$metrics" | grep -E '^(total_requests|cache_hits|cache_misses|hit_rate|memory_cache_hits|disk_cache_hits|errors|upstream_errors|cache_write_failures|rate_limited) ' || true; \
|
||||
echo ""; \
|
||||
echo "=== http://localhost/lancache-heartbeat (GET; expect 204 + X-LanCache-Processed-By: SteamCache2) ==="; \
|
||||
hb=$$(curl -sD - -o /dev/null --max-time 5 http://localhost/lancache-heartbeat) || { \
|
||||
@@ -133,7 +133,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 fields) and /lancache-heartbeat (default :80)"
|
||||
@echo " validate-check Curl local /metrics (full dump + hit/miss + upstream/write/rate fields) and /lancache-heartbeat (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)"
|
||||
|
||||
@@ -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 fields, 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`. Read these fields:
|
||||
|
||||
| Field | Meaning |
|
||||
| --- | --- |
|
||||
@@ -84,11 +84,14 @@ curl -s -i http://localhost/lancache-heartbeat
|
||||
| `range_cache` / `range_upstream` | Range GETs served as 206 from a cached object vs after a full upstream fetch |
|
||||
| `memory_cache_hits` / `disk_cache_hits` | Which tier served the hits |
|
||||
| `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) |
|
||||
| `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.
|
||||
|
||||
Cache entries are keyed by depot object path (not the CDN `Host` header), so when Steam rotates CDN hostnames for the same depot path, hits still climb across the aliases.
|
||||
|
||||
Steam clients lean on Range requests. When an object is already cached, a Range GET is served locally as 206 from that full object (`range_cache`). On a Range miss the cache still fetches and stores the full upstream body, then returns the requested byte range as 206 (`range_upstream`).
|
||||
|
||||
To confirm the process is up (HTTP 204 and `X-LanCache-Processed-By: SteamCache2`):
|
||||
@@ -173,7 +176,7 @@ 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:
|
||||
- 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`
|
||||
- Zero unexpected `errors`, and quiet `upstream_errors` / `cache_write_failures` / `rate_limited`
|
||||
|
||||
Heartbeat should be HTTP 204 with `X-LanCache-Processed-By: SteamCache2`. Use GET (`curl -i`), not HEAD (`curl -I`).
|
||||
|
||||
|
||||
@@ -279,9 +279,12 @@ func (sc *SteamCache) ServeHTTP(w http.ResponseWriter, r *http.Request) {
|
||||
|
||||
// Check if this is a request from a supported service
|
||||
if service, isSupported := sc.detectService(r); isSupported {
|
||||
// trim the query parameters from the URL path
|
||||
// this is necessary because the cache key should not include query parameters
|
||||
urlPath := strings.SplitN(r.URL.String(), "?", 2)[0] // trim query for cache key (SplitN makes intent explicit vs Cut + ignored bool)
|
||||
// Cache key is the path only, never the Host: Steam rotates CDN hostnames
|
||||
// for the same depot object, so different Host headers (or absolute-form
|
||||
// request targets) for the same path must share one cache entry. r.URL.Path
|
||||
// is the decoded path (query is never part of it); validateURLPath checks
|
||||
// this decoded form and url.JoinPath re-escapes it for the upstream join.
|
||||
urlPath := r.URL.Path
|
||||
|
||||
// Validate URL path for security
|
||||
if err := validateURLPath(urlPath); err != nil {
|
||||
|
||||
@@ -2,21 +2,25 @@
|
||||
package steamcache
|
||||
|
||||
import (
|
||||
"bufio"
|
||||
"bytes"
|
||||
"context"
|
||||
"fmt"
|
||||
"io"
|
||||
"net"
|
||||
"net/http"
|
||||
"net/http/httptest"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"runtime"
|
||||
"s1d3sw1ped/steamcache2/steamcache/metrics"
|
||||
"s1d3sw1ped/steamcache2/vfs/disk"
|
||||
"s1d3sw1ped/steamcache2/vfs/eviction"
|
||||
"s1d3sw1ped/steamcache2/vfs/memory"
|
||||
"s1d3sw1ped/steamcache2/vfs/vfserror"
|
||||
"strings"
|
||||
"sync"
|
||||
"sync/atomic"
|
||||
"testing"
|
||||
"time"
|
||||
)
|
||||
@@ -1136,19 +1140,31 @@ func TestDiskTierSignalMemoryOnly(t *testing.T) {
|
||||
// TestDiskTierSignalMixedPendingReady covers mixed mode: DiskTierReady=0 (header
|
||||
// pending) while the disk attach is in the Size barrier, then DiskTierReady=1
|
||||
// (header ready) after the barrier opens and the attach goroutine sets SetSlow.
|
||||
// An init hold keeps the empty-dir attach from finishing between the pending
|
||||
// metric check and the heartbeat, which otherwise races under CI load.
|
||||
func TestDiskTierSignalMixedPendingReady(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)
|
||||
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
|
||||
|
||||
// Immediately in the pending window
|
||||
// Pending window is held open until closeHold; attach cannot finish.
|
||||
if got := sc.GetMetrics().DiskTierReady; got != 0 {
|
||||
t.Errorf("immediate DiskTierReady=%d, want 0 (pending)", got)
|
||||
}
|
||||
@@ -1159,7 +1175,7 @@ func TestDiskTierSignalMixedPendingReady(t *testing.T) {
|
||||
t.Errorf("heartbeat header=%q, want pending", got)
|
||||
}
|
||||
|
||||
// Wait the barrier, then retry until the attach goroutine flips the flag
|
||||
closeHold()
|
||||
_ = sc.disk.Size()
|
||||
deadline := time.Now().Add(2 * time.Second)
|
||||
for time.Now().Before(deadline) {
|
||||
@@ -1317,3 +1333,135 @@ func TestDirectFetchRejectsNonSteamHost(t *testing.T) {
|
||||
t.Errorf("non-CDN Host: expected 400, got %d", rec2.Code)
|
||||
}
|
||||
}
|
||||
|
||||
// TestCacheKeySharedAcrossCDNHostAliases verifies that one cache entry serves
|
||||
// the same depot object across different Steam CDN host aliases: the key uses
|
||||
// only the request path (never the Host header or an absolute-form target
|
||||
// host), so hits climb instead of re-stamping upstream per host rotation.
|
||||
func TestCacheKeySharedAcrossCDNHostAliases(t *testing.T) {
|
||||
body := []byte("depot chunk body for host-alias keying")
|
||||
var upstreamCalls atomic.Int64
|
||||
f := func(w http.ResponseWriter, r *http.Request) {
|
||||
upstreamCalls.Add(1)
|
||||
w.Header().Set("Content-Type", "application/octet-stream")
|
||||
_, _ = w.Write(body)
|
||||
}
|
||||
sc, _ := newTestCacheWithFakeUpstream(t, f, "1MB", "0")
|
||||
srv := newCacheServer(t, sc)
|
||||
const depotPath = "/depot/1684171/chunk/abc123"
|
||||
const ua = "Valve/Steam HTTP Client 1.0"
|
||||
c := &http.Client{Timeout: 5 * time.Second}
|
||||
|
||||
// 1) MISS under the first CDN alias (origin-form target, Host: cdn1).
|
||||
req1, err := http.NewRequest("GET", srv.URL+depotPath, nil)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
req1.Host = "cdn1.steamcontent.com"
|
||||
req1.Header.Set("User-Agent", ua)
|
||||
resp1, err := c.Do(req1)
|
||||
if err != nil {
|
||||
t.Fatalf("host-alias MISS request: %v", err)
|
||||
}
|
||||
data1, err := io.ReadAll(resp1.Body)
|
||||
resp1.Body.Close()
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if resp1.StatusCode != http.StatusOK {
|
||||
t.Fatalf("host-alias MISS: expected 200, got %d", resp1.StatusCode)
|
||||
}
|
||||
if got := resp1.Header.Get("X-LanCache-Status"); got != "MISS" {
|
||||
t.Fatalf("host-alias MISS: expected X-LanCache-Status MISS, got %q", got)
|
||||
}
|
||||
if !bytes.Equal(data1, body) {
|
||||
t.Fatalf("host-alias MISS: body mismatch: got %q", data1)
|
||||
}
|
||||
|
||||
// Bounded wait for the entry to be visible before hitting the next alias
|
||||
// (the MISS handler streams the body to the client before the VFS write).
|
||||
key, err := generateServiceCacheKey(depotPath, "steam")
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
deadline := time.Now().Add(2 * time.Second)
|
||||
for {
|
||||
if rc, e := sc.vfs.Open(key); e == nil {
|
||||
_ = rc.Close()
|
||||
break
|
||||
}
|
||||
if time.Now().After(deadline) {
|
||||
t.Fatalf("cache entry %q not visible after MISS", key)
|
||||
}
|
||||
time.Sleep(5 * time.Millisecond)
|
||||
}
|
||||
|
||||
// 2) Same depot path under the second CDN alias (Host header only) -> HIT.
|
||||
req2, err := http.NewRequest("GET", srv.URL+depotPath, nil)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
req2.Host = "cdn2.steamcontent.com"
|
||||
req2.Header.Set("User-Agent", ua)
|
||||
resp2, err := c.Do(req2)
|
||||
if err != nil {
|
||||
t.Fatalf("second host-alias request: %v", err)
|
||||
}
|
||||
data2, err := io.ReadAll(resp2.Body)
|
||||
resp2.Body.Close()
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if resp2.StatusCode != http.StatusOK {
|
||||
t.Fatalf("second host-alias: expected 200, got %d", resp2.StatusCode)
|
||||
}
|
||||
if got := resp2.Header.Get("X-LanCache-Status"); got != "HIT" {
|
||||
t.Fatalf("second host-alias: expected X-LanCache-Status HIT, got %q", got)
|
||||
}
|
||||
if !bytes.Equal(data2, body) {
|
||||
t.Fatalf("second host-alias: body mismatch: got %q", data2)
|
||||
}
|
||||
|
||||
// 3) Same depot path with an absolute-form target embedding a third CDN
|
||||
// hostname in the URL itself -> still a HIT on the same entry.
|
||||
// (Go's http client always sends origin-form targets, so use raw HTTP.)
|
||||
conn, err := net.Dial("tcp", srv.Listener.Addr().String())
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
defer conn.Close()
|
||||
rawRequest := "GET http://cdn3.steamcontent.com" + depotPath + " HTTP/1.1\r\n" +
|
||||
"Host: cdn3.steamcontent.com\r\n" +
|
||||
"User-Agent: " + ua + "\r\n" +
|
||||
"Connection: close\r\n\r\n"
|
||||
if _, err := conn.Write([]byte(rawRequest)); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
rawReq, err := http.NewRequest("GET", "http://cdn3.steamcontent.com"+depotPath, nil)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
resp3, err := http.ReadResponse(bufio.NewReader(conn), rawReq)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
defer resp3.Body.Close()
|
||||
data3, err := io.ReadAll(resp3.Body)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if resp3.StatusCode != http.StatusOK {
|
||||
t.Fatalf("absolute-form host-alias: expected 200, got %d", resp3.StatusCode)
|
||||
}
|
||||
if got := resp3.Header.Get("X-LanCache-Status"); got != "HIT" {
|
||||
t.Fatalf("absolute-form host-alias: expected X-LanCache-Status HIT, got %q", got)
|
||||
}
|
||||
if !bytes.Equal(data3, body) {
|
||||
t.Fatalf("absolute-form host-alias: body mismatch: got %q", data3)
|
||||
}
|
||||
|
||||
// All three aliases must have shared one upstream fetch.
|
||||
if got := upstreamCalls.Load(); got != 1 {
|
||||
t.Errorf("upstream fetched %d times across host aliases, want 1", got)
|
||||
}
|
||||
}
|
||||
|
||||
+39
-2
@@ -45,6 +45,25 @@ type DiskFS struct {
|
||||
initCloseOnce sync.Once
|
||||
startupEvict func(vfs.VFS, uint) uint // passed to New (via gc.GetGCAlgorithm); invoked as last step of bg init if over cap (no post-ctor race)
|
||||
metrics *metrics.Metrics
|
||||
// initHold, if non-nil, is received on before closing initDone (test pending-window hold).
|
||||
initHold <-chan struct{}
|
||||
}
|
||||
|
||||
// initHolds is a per-root registry of optional init holds (root path -> <-chan struct{}).
|
||||
// Tests call RegisterInitHold before New so that instance copies the channel and waits
|
||||
// before closing initDone; production never registers, so init is unchanged.
|
||||
var initHolds sync.Map
|
||||
|
||||
// RegisterInitHold registers a channel that DiskFS.New for this root copies onto that
|
||||
// instance. calculateSizeAndPopulateIndex receives on it before closing initDone, so
|
||||
// tests can observe the pending-attach window. Other DiskFS roots are unaffected.
|
||||
func RegisterInitHold(root string, ch <-chan struct{}) {
|
||||
initHolds.Store(root, ch)
|
||||
}
|
||||
|
||||
// ClearInitHold removes a previously registered hold for root.
|
||||
func ClearInitHold(root string) {
|
||||
initHolds.Delete(root)
|
||||
}
|
||||
|
||||
// shardPath converts a Steam cache key to a sharded directory path to reduce inode pressure
|
||||
@@ -129,6 +148,12 @@ func New(root string, capacity int64, evict func(vfs.VFS, uint) uint) (*DiskFS,
|
||||
startupEvict: evict,
|
||||
}
|
||||
|
||||
if v, ok := initHolds.Load(root); ok {
|
||||
if ch, ok := v.(<-chan struct{}); ok {
|
||||
d.initHold = ch
|
||||
}
|
||||
}
|
||||
|
||||
d.initDone = make(chan struct{})
|
||||
// Launch heavy population asynchronously so New returns fast (scans millions of files without blocking ctor or using O(N) temp RAM).
|
||||
// The initDone barrier ensures first Size() and subsequent ops (including late tier attach) see fully populated + post-eviction state.
|
||||
@@ -152,7 +177,7 @@ func (d *DiskFS) calculateSizeAndPopulateIndex() {
|
||||
if r := recover(); r != nil {
|
||||
logger.Logger.Error().Interface("recovered_panic", r).Msg("calculateSizeAndPopulateIndex panicked; ensuring initDone closed to unblock Size waiters and prevent hang")
|
||||
}
|
||||
d.initCloseOnce.Do(func() { close(d.initDone) })
|
||||
d.closeInitDone()
|
||||
}()
|
||||
|
||||
tstart := time.Now()
|
||||
@@ -243,7 +268,19 @@ func (d *DiskFS) calculateSizeAndPopulateIndex() {
|
||||
|
||||
// Signal readiness: Size() and callers (late tier attach + Evict*) now see correct populated + post-eviction state.
|
||||
// Use Once (recover path also uses it) to guarantee exactly one close even under panic.
|
||||
d.initCloseOnce.Do(func() { close(d.initDone) })
|
||||
d.closeInitDone()
|
||||
}
|
||||
|
||||
// closeInitDone receives on a copied test hold (if any) then closes initDone once.
|
||||
// Both the normal end of calculateSizeAndPopulateIndex and the panic-recovery defer
|
||||
// call this so Size() waiters unblock in either path.
|
||||
func (d *DiskFS) closeInitDone() {
|
||||
d.initCloseOnce.Do(func() {
|
||||
if d.initHold != nil {
|
||||
<-d.initHold
|
||||
}
|
||||
close(d.initDone)
|
||||
})
|
||||
}
|
||||
|
||||
// insertBatch populates info/LRU under lock for a bounded batch (follows maxEvictBatch pattern for short critical sections).
|
||||
|
||||
@@ -662,3 +662,57 @@ func TestDiskFS_NewMkdirError(t *testing.T) {
|
||||
t.Errorf("expected mkdir failure error for file-as-dir, got: %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
// TestDiskFS_InitHoldBlocksOnlyRegisteredRoot covers the per-root init hold:
|
||||
// Size() stays blocked while the hold is open, and a DiskFS on a different root
|
||||
// does not wait on that hold.
|
||||
func TestDiskFS_InitHoldBlocksOnlyRegisteredRoot(t *testing.T) {
|
||||
td := t.TempDir()
|
||||
hold := make(chan struct{})
|
||||
var holdOnce sync.Once
|
||||
closeHold := func() { holdOnce.Do(func() { close(hold) }) }
|
||||
RegisterInitHold(td, hold)
|
||||
t.Cleanup(func() {
|
||||
closeHold()
|
||||
ClearInitHold(td)
|
||||
})
|
||||
|
||||
d, err := New(td, 10*1024*1024, nil)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
|
||||
blocked := make(chan struct{})
|
||||
go func() {
|
||||
_ = d.Size()
|
||||
close(blocked)
|
||||
}()
|
||||
select {
|
||||
case <-blocked:
|
||||
t.Fatal("Size returned while init hold still open")
|
||||
case <-time.After(50 * time.Millisecond):
|
||||
}
|
||||
|
||||
td2 := t.TempDir()
|
||||
d2, err := New(td2, 10*1024*1024, nil)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
other := make(chan struct{})
|
||||
go func() {
|
||||
_ = d2.Size()
|
||||
close(other)
|
||||
}()
|
||||
select {
|
||||
case <-other:
|
||||
case <-time.After(2 * time.Second):
|
||||
t.Fatal("unrelated DiskFS Size hung; init hold leaked across roots")
|
||||
}
|
||||
|
||||
closeHold()
|
||||
select {
|
||||
case <-blocked:
|
||||
case <-time.After(2 * time.Second):
|
||||
t.Fatal("Size did not return after init hold released")
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user