Compare commits

...

6 Commits

Author SHA1 Message Date
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
linus a3ea4806a7 Merge pull request 'cache: Coalesce in-flight identical upstream fetches' (#52) from cache/coalesce-inflight-fetches into develop
CI / vulncheck (push) Successful in 16s
CI / check-and-test (push) Successful in 48s
Release Tag / release (push) Successful in 15s
cache: Coalesce in-flight identical upstream fetches

Fixes #35
2026-09-08 15:19:26 -05:00
ash ea195993de cache: Coalesce in-flight identical upstream fetches
CI / vulncheck (pull_request) Successful in 14s
CI / check-and-test (pull_request) Successful in 48s
In-flight coalescing already dedups identical misses, but acceptance
still lacked a gated test (without a hold, waiters can become sequential
HITs after the first fill) and the Quick check metrics table omitted
cache_coalesced.

Add steamcache/coalesce_test.go: N concurrent identical GETs share one
upstream fill (leader MISS, waiters HIT-COALESCED) and a 5xx sibling.
Document cache_coalesced next to the other hit/miss fields.

Fixes #35
2026-09-08 20:00:53 +00:00
eva 7c34ff4538 Merge pull request 'metrics: Prometheus text exposition for /metrics' (#51) from metrics/prometheus-exposition into main
CI / vulncheck (push) Successful in 15s
CI / check-and-test (push) Successful in 41s
Release Tag / release (push) Successful in 16s
metrics: Prometheus text exposition for /metrics (#51)

Closes #49
2026-09-08 14:57:22 -05:00
pike dd72668c2d metrics: Prometheus text exposition for /metrics
CI / vulncheck (pull_request) Successful in 14s
CI / check-and-test (pull_request) Successful in 42s
/metrics was Prometheus-ish name/value text with a custom banner and
text/plain Content-Type, so strict scrapers could fail.

Emit Prometheus text 0.0.4 (# HELP, # TYPE, counter|gauge) with stable
metric names, and set Content-Type to text/plain; version=0.0.4; charset=utf-8.
2026-09-08 19:53:58 +00:00
linus 0db9943436 ops: Signal disk-full and eviction capacity pressure (#46)
CI / vulncheck (push) Successful in 14s
CI / check-and-test (push) Successful in 43s
Release Tag / release (push) Successful in 14s
capacity_pressure_events + ENOSPC/eviction logs. Closes #36.
2026-09-08 13:20:22 -05:00
9 changed files with 642 additions and 31 deletions
+9
View File
@@ -81,18 +81,24 @@ curl -s -i http://localhost/lancache-heartbeat
| Field | Meaning |
| --- | --- |
| `cache_hits` / `cache_misses` / `hit_rate` | Whether later requests were served from cache |
| `cache_coalesced` | Waiters on an in-flight identical miss share one upstream fill (`X-LanCache-Status: HIT-COALESCED`) |
| `negative_cache_hits` | 404/410 served from a still-valid negative cache entry (also counted in `cache_hits`) |
| `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) |
| `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.
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.
Concurrent identical misses for the same key share one upstream GET: the leader is a `MISS` and waiters are `HIT-COALESCED` (`cache_coalesced`).
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`).
Definitive upstream 404/410 (gone depot objects) are stored as a short-TTL negative entry in the **same** cache, under the same depot-path key as a positive object. Repeating the request within `cache.negative_ttl` (default `5m`) is served as 404/410 without re-hitting upstream (`negative_cache_hits`). 5xx is not cached as negative. When the TTL expires the entry is deleted and the next request fetches again.
@@ -109,6 +115,8 @@ Heartbeat also returns `X-SteamCache-Disk-Tier: pending|ready|disabled` (`disabl
These are the cache process's own `/metrics` and `/lancache-heartbeat` endpoints. There is no separate metrics daemon.
`/metrics` is Prometheus text exposition format 0.0.4 (`Content-Type: text/plain; version=0.0.4; charset=utf-8`) so Prometheus and compatible scrapers can pull it. Metric names in the table above are unchanged; each series is preceded by `# HELP` and `# TYPE`.
If you changed `listen_address`, point curl at that host:port instead. For a full SteamPrefill validation workflow (small caches, coalescing, GC), see [Validating Full Functionality](#validating-full-functionality-with-external-tools).
### Development Workflow
@@ -306,6 +314,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.
+199
View File
@@ -0,0 +1,199 @@
package steamcache
import (
"bytes"
"net/http"
"net/http/httptest"
"sync"
"sync/atomic"
"testing"
"time"
)
func coalescerWaiterCount(sc *SteamCache, cacheKey string) int32 {
sc.coalescer.mu.Lock()
defer sc.coalescer.mu.Unlock()
cr := sc.coalescer.requests[cacheKey]
if cr == nil {
return 0
}
return cr.waitingCount.Load()
}
func steamCoalesceRequest(path string) *http.Request {
req := httptest.NewRequest(http.MethodGet, path, nil)
req.Header.Set("User-Agent", "Valve/Steam HTTP Client 1.0")
return req
}
func waitForCoalescerJoin(t *testing.T, sc *SteamCache, cacheKey string, n int, release func(), wg *sync.WaitGroup, upstreamCalls *atomic.Int64) {
t.Helper()
deadline := time.Now().Add(2 * time.Second)
var waiters int32
for {
waiters = coalescerWaiterCount(sc, cacheKey)
if waiters >= int32(n) {
return
}
if time.Now().After(deadline) {
release()
wg.Wait()
t.Fatalf("coalescer waiters=%d want %d (upstreamCalls=%d)", waiters, n, upstreamCalls.Load())
}
time.Sleep(1 * time.Millisecond)
}
}
// TestCoalesceIdenticalMissesOneUpstreamGET holds the leader's upstream GET
// open until every concurrent client has joined the in-flight coalescer.
// Without that hold, later requests can become sequential HITs after the first
// miss fills, which would not prove coalescing.
func TestCoalesceIdenticalMissesOneUpstreamGET(t *testing.T) {
const nClients = 8
body := []byte("coalesced depot chunk body")
var upstreamCalls atomic.Int64
release := make(chan struct{})
var releaseOnce sync.Once
releaseUpstream := func() { releaseOnce.Do(func() { close(release) }) }
t.Cleanup(releaseUpstream)
f := func(w http.ResponseWriter, r *http.Request) {
upstreamCalls.Add(1)
select {
case <-release:
case <-r.Context().Done():
return
}
w.Header().Set("Content-Type", "application/octet-stream")
_, _ = w.Write(body)
}
sc, _ := newTestCacheWithFakeUpstream(t, f, "1MB", "0")
sc.ResetMetrics()
const depotPath = "/depot/1684171/chunk/coalesce-inflight"
cacheKey, err := generateServiceCacheKey(depotPath, "steam")
if err != nil {
t.Fatal(err)
}
type clientResult struct {
status int
hdr string
body []byte
}
results := make([]clientResult, nClients)
var wg sync.WaitGroup
start := make(chan struct{})
wg.Add(nClients)
for i := 0; i < nClients; i++ {
go func(i int) {
defer wg.Done()
<-start
rec := httptest.NewRecorder()
sc.ServeHTTP(rec, steamCoalesceRequest(depotPath))
results[i] = clientResult{
status: rec.Code,
hdr: rec.Header().Get("X-LanCache-Status"),
body: rec.Body.Bytes(),
}
}(i)
}
close(start)
waitForCoalescerJoin(t, sc, cacheKey, nClients, releaseUpstream, &wg, &upstreamCalls)
releaseUpstream()
wg.Wait()
if got := upstreamCalls.Load(); got != 1 {
t.Fatalf("expected exactly 1 upstream GET, got %d", got)
}
var miss, coalesced int
for i, r := range results {
if r.status != http.StatusOK {
t.Errorf("client %d: expected 200, got %d", i, r.status)
}
if !bytes.Equal(r.body, body) {
t.Errorf("client %d: body mismatch: got %q", i, r.body)
}
switch r.hdr {
case "MISS":
miss++
case "HIT-COALESCED":
coalesced++
default:
t.Errorf("client %d: unexpected X-LanCache-Status %q", i, r.hdr)
}
}
if miss != 1 {
t.Errorf("expected 1 MISS leader, got %d", miss)
}
if coalesced != nClients-1 {
t.Errorf("expected %d HIT-COALESCED waiters, got %d", nClients-1, coalesced)
}
if got := sc.GetMetrics().CacheCoalesced; got < int64(nClients-1) {
t.Errorf("CacheCoalesced=%d, want >= %d", got, nClients-1)
}
}
// TestCoalesceIdenticalMissesSharedUpstreamError is the 5xx sibling: waiters
// share the leader's failure instead of each hitting origin. Upstream 500 is
// retried, so the call count is the leader's retry budget (not N).
func TestCoalesceIdenticalMissesSharedUpstreamError(t *testing.T) {
const nClients = 8
var upstreamCalls atomic.Int64
release := make(chan struct{})
var releaseOnce sync.Once
releaseUpstream := func() { releaseOnce.Do(func() { close(release) }) }
t.Cleanup(releaseUpstream)
f := func(w http.ResponseWriter, r *http.Request) {
upstreamCalls.Add(1)
select {
case <-release:
case <-r.Context().Done():
return
}
w.WriteHeader(http.StatusInternalServerError)
}
sc, _ := newTestCacheWithFakeUpstream(t, f, "1MB", "0")
sc.ResetMetrics()
const depotPath = "/depot/1684171/chunk/coalesce-inflight-err"
cacheKey, err := generateServiceCacheKey(depotPath, "steam")
if err != nil {
t.Fatal(err)
}
codes := make([]int, nClients)
var wg sync.WaitGroup
start := make(chan struct{})
wg.Add(nClients)
for i := 0; i < nClients; i++ {
go func(i int) {
defer wg.Done()
<-start
rec := httptest.NewRecorder()
sc.ServeHTTP(rec, steamCoalesceRequest(depotPath))
codes[i] = rec.Code
}(i)
}
close(start)
waitForCoalescerJoin(t, sc, cacheKey, nClients, releaseUpstream, &wg, &upstreamCalls)
releaseUpstream()
wg.Wait()
if got := upstreamCalls.Load(); got < 1 || got >= int64(nClients) {
t.Fatalf("expected coalesced origin GETs (leader + retries, < %d waiters), got %d", nClients, got)
}
for i, code := range codes {
if code != http.StatusInternalServerError {
t.Errorf("client %d: expected 500, got %d", i, code)
}
}
if got := sc.GetMetrics().Errors; got < int64(nClients) {
t.Errorf("Errors=%d, want >= %d (once per client)", got, nClients)
}
}
+2 -2
View File
@@ -76,9 +76,9 @@ func (sc *SteamCache) handleSpecialEndpoints(w http.ResponseWriter, r *http.Requ
}
if r.URL.String() == "/metrics" {
// Return metrics in a simple text format
// Prometheus text exposition format 0.0.4
stats := sc.GetMetrics()
w.Header().Set("Content-Type", "text/plain")
w.Header().Set("Content-Type", "text/plain; version=0.0.4; charset=utf-8")
w.WriteHeader(http.StatusOK)
metrics.WriteText(w, stats)
return true
+86 -29
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
@@ -344,42 +381,62 @@ type Stats struct {
LastResetTime time.Time
}
// WriteText emits the Prometheus-style text metrics to the ResponseWriter.
// Promoted from internal handler per Phase 3 for better package ownership.
// WriteText emits Prometheus text exposition format 0.0.4 to the ResponseWriter.
// Each metric family is # HELP, then # TYPE, then one or more sample lines.
// Metric names are stable; labeled series keep service=%q (Prometheus-valid quotes).
// All fmt.Fprintf errors are intentionally discarded via _ = : this is a best-effort
// read-only debug endpoint; client disconnects or write errors during metrics dump
// are not actionable (do not affect cache correctness or require retries).
func WriteText(w http.ResponseWriter, stats *Stats) {
_, _ = fmt.Fprintf(w, "# SteamCache2 Metrics\n")
_, _ = fmt.Fprintf(w, "total_requests %d\n", stats.TotalRequests)
_, _ = fmt.Fprintf(w, "cache_hits %d\n", stats.CacheHits)
_, _ = fmt.Fprintf(w, "cache_misses %d\n", stats.CacheMisses)
_, _ = fmt.Fprintf(w, "negative_cache_hits %d\n", stats.NegativeCacheHits)
_, _ = fmt.Fprintf(w, "cache_coalesced %d\n", stats.CacheCoalesced)
_, _ = fmt.Fprintf(w, "range_cache %d\n", stats.RangeCache)
_, _ = fmt.Fprintf(w, "range_upstream %d\n", stats.RangeUpstream)
_, _ = fmt.Fprintf(w, "errors %d\n", stats.Errors)
_, _ = fmt.Fprintf(w, "rate_limited %d\n", stats.RateLimited)
_, _ = fmt.Fprintf(w, "upstream_errors %d\n", stats.UpstreamErrors)
_, _ = fmt.Fprintf(w, "cache_write_failures %d\n", stats.CacheWriteFailures)
_, _ = fmt.Fprintf(w, "memory_cache_hits %d\n", stats.MemoryCacheHits)
_, _ = fmt.Fprintf(w, "disk_cache_hits %d\n", stats.DiskCacheHits)
_, _ = fmt.Fprintf(w, "promotions %d\n", stats.Promotions)
_, _ = fmt.Fprintf(w, "evictions %d\n", stats.Evictions)
_, _ = fmt.Fprintf(w, "capacity_pressure_events %d\n", stats.CapacityPressureEvents)
writeInt(w, "total_requests", "Total HTTP requests handled.", "counter", stats.TotalRequests)
writeInt(w, "cache_hits", "Requests served from cache.", "counter", stats.CacheHits)
writeInt(w, "cache_misses", "Requests not found in cache.", "counter", stats.CacheMisses)
writeInt(w, "negative_cache_hits", "404/410 served from a still-valid negative cache entry.", "counter", stats.NegativeCacheHits)
writeInt(w, "cache_coalesced", "Requests coalesced onto an in-flight upstream fetch.", "counter", stats.CacheCoalesced)
writeInt(w, "range_cache", "Range requests served as 206 from an already-cached object.", "counter", stats.RangeCache)
writeInt(w, "range_upstream", "Range requests that required an upstream fetch, served as 206.", "counter", stats.RangeUpstream)
writeInt(w, "errors", "Request errors.", "counter", stats.Errors)
writeInt(w, "rate_limited", "Requests rejected by rate limiting.", "counter", stats.RateLimited)
writeInt(w, "upstream_errors", "Errors talking to upstream.", "counter", stats.UpstreamErrors)
writeInt(w, "cache_write_failures", "Failures writing objects into cache.", "counter", stats.CacheWriteFailures)
writeInt(w, "memory_cache_hits", "Hits served from the memory tier.", "counter", stats.MemoryCacheHits)
writeInt(w, "disk_cache_hits", "Hits served from the disk tier.", "counter", stats.DiskCacheHits)
writeInt(w, "promotions", "Objects promoted from disk to memory.", "counter", stats.Promotions)
writeInt(w, "evictions", "Objects evicted from cache.", "counter", stats.Evictions)
writeInt(w, "capacity_pressure_events", "Soft eviction under the memory or disk cap, and/or disk ENOSPC.", "counter", stats.CapacityPressureEvents)
writeHelpType(w, "service_errors", "Errors attributed to a named service.", "counter")
for svc, cnt := range stats.ServiceErrors {
_, _ = fmt.Fprintf(w, "service_errors{service=%q} %d\n", svc, cnt)
}
writeHelpType(w, "service_requests", "Requests attributed to a named service.", "counter")
for svc, cnt := range stats.ServiceRequests {
_, _ = fmt.Fprintf(w, "service_requests{service=%q} %d\n", svc, cnt)
}
_, _ = fmt.Fprintf(w, "hit_rate %.4f\n", stats.HitRate)
_, _ = fmt.Fprintf(w, "avg_response_time_ms %.2f\n", float64(stats.AvgResponseTime.Nanoseconds())/1e6)
_, _ = fmt.Fprintf(w, "total_bytes_served %d\n", stats.TotalBytesServed)
_, _ = fmt.Fprintf(w, "total_bytes_saved %d\n", stats.TotalBytesSaved)
_, _ = fmt.Fprintf(w, "memory_cache_size %d\n", stats.MemoryCacheSize)
_, _ = fmt.Fprintf(w, "disk_cache_size %d\n", stats.DiskCacheSize)
_, _ = fmt.Fprintf(w, "disk_tier_ready %d\n", stats.DiskTierReady)
_, _ = fmt.Fprintf(w, "uptime_seconds %.2f\n", stats.Uptime.Seconds())
writeFloat(w, "hit_rate", "Cache hits divided by total requests.", "gauge", "%.4f", stats.HitRate)
writeFloat(w, "avg_response_time_ms", "Average response time in milliseconds.", "gauge", "%.2f", float64(stats.AvgResponseTime.Nanoseconds())/1e6)
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())
}
func writeHelpType(w http.ResponseWriter, name, help, typ string) {
_, _ = fmt.Fprintf(w, "# HELP %s %s\n# TYPE %s %s\n", name, help, name, typ)
}
func writeInt(w http.ResponseWriter, name, help, typ string, v int64) {
writeHelpType(w, name, help, typ)
_, _ = fmt.Fprintf(w, "%s %d\n", name, v)
}
func writeFloat(w http.ResponseWriter, name, help, typ, valFmt string, v float64) {
writeHelpType(w, name, help, typ)
_, _ = fmt.Fprintf(w, "%s "+valFmt+"\n", name, v)
}
+177
View File
@@ -4,7 +4,9 @@ import (
"bytes"
"errors"
"net/http/httptest"
"strings"
"testing"
"time"
)
func TestCapacityPressureEventsWriteTextAndReset(t *testing.T) {
@@ -47,6 +49,15 @@ func TestCapacityPressureEventsWriteTextAndReset(t *testing.T) {
if !bytes.Contains(body, []byte("evictions 2")) {
t.Errorf("WriteText missing evictions 2: %q", rec.Body.String())
}
if !bytes.Contains(body, []byte("# HELP capacity_pressure_events")) {
t.Errorf("WriteText missing # HELP capacity_pressure_events: %q", rec.Body.String())
}
if !bytes.Contains(body, []byte("# TYPE capacity_pressure_events counter")) {
t.Errorf("WriteText missing # TYPE capacity_pressure_events counter: %q", rec.Body.String())
}
if !bytes.Contains(body, []byte("# TYPE evictions counter")) {
t.Errorf("WriteText missing # TYPE evictions counter: %q", rec.Body.String())
}
m.Reset()
st = m.GetStats()
@@ -61,3 +72,169 @@ func TestNoteSoftEvictionNilMetrics(t *testing.T) {
NoteSoftEviction(nil, "memory", 10)
NoteNoSpace(nil, errors.New("ENOSPC"))
}
func TestWriteTextPrometheusExposition(t *testing.T) {
t.Parallel()
st := &Stats{
TotalRequests: 10,
CacheHits: 4,
CacheMisses: 6,
NegativeCacheHits: 1,
CacheCoalesced: 2,
RangeCache: 3,
RangeUpstream: 5,
Errors: 1,
RateLimited: 1,
UpstreamErrors: 1,
CacheWriteFailures: 1,
MemoryCacheHits: 2,
DiskCacheHits: 2,
Promotions: 1,
Evictions: 1,
CapacityPressureEvents: 1,
ServiceErrors: map[string]int64{"steam": 2},
ServiceRequests: map[string]int64{"steam": 7},
HitRate: 0.4,
AvgResponseTime: 2 * time.Millisecond,
TotalBytesServed: 100,
TotalBytesSaved: 50,
MemoryCacheSize: 8,
DiskCacheSize: 16,
MemoryCacheCapacity: 8,
DiskCacheCapacity: 32,
DiskCacheFullRatio: 0.5,
DiskTierReady: 1,
Uptime: 3 * time.Second,
}
rec := httptest.NewRecorder()
WriteText(rec, st)
body := rec.Body.String()
if strings.Contains(body, "# SteamCache2 Metrics") {
t.Error("non-standard # SteamCache2 Metrics banner must not be present")
}
counters := []string{
"total_requests", "cache_hits", "cache_misses", "negative_cache_hits",
"cache_coalesced", "range_cache", "range_upstream", "errors", "rate_limited",
"upstream_errors", "cache_write_failures", "memory_cache_hits", "disk_cache_hits",
"promotions", "evictions", "capacity_pressure_events", "service_errors",
"service_requests", "total_bytes_served", "total_bytes_saved",
}
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 {
assertHelpType(t, body, name, "counter")
}
for _, name := range gauges {
assertHelpType(t, body, name, "gauge")
}
if !strings.Contains(body, `service_errors{service="steam"} 2`) {
t.Errorf("missing labeled service_errors sample: %q", body)
}
if !strings.Contains(body, `service_requests{service="steam"} 7`) {
t.Errorf("missing labeled service_requests sample: %q", body)
}
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)
}
if !strings.Contains(body, "negative_cache_hits 1\n") {
t.Errorf("missing negative_cache_hits 1 sample: %q", body)
}
}
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 + " "
typeLine := "# TYPE " + name + " " + typ
iHelp := strings.Index(body, help)
if iHelp < 0 {
t.Errorf("missing %q", help)
return
}
iType := strings.Index(body[iHelp:], typeLine)
if iType < 0 {
t.Errorf("missing %q after HELP for %s", typeLine, name)
return
}
afterType := body[iHelp+iType+len(typeLine):]
if !strings.HasPrefix(afterType, "\n") {
t.Errorf("# TYPE %s not followed by newline", name)
return
}
sample := afterType[1:]
if !strings.HasPrefix(sample, name+" ") && !strings.HasPrefix(sample, name+"{") {
t.Errorf("sample for %s does not follow TYPE; next line starts %q", name, firstLine(sample))
}
}
func firstLine(s string) string {
if i := strings.IndexByte(s, '\n'); i >= 0 {
return s[:i]
}
return s
}
+12
View File
@@ -61,6 +61,9 @@ func TestNegativeCache404(t *testing.T) {
if !bytes.Contains(mrec.Body.Bytes(), []byte("negative_cache_hits")) {
t.Errorf("/metrics missing negative_cache_hits:\n%s", mrec.Body.String())
}
if ct := mrec.Header().Get("Content-Type"); ct != "text/plain; version=0.0.4; charset=utf-8" {
t.Errorf("/metrics Content-Type=%q, want Prometheus text 0.0.4", ct)
}
}
func TestNegativeCache410(t *testing.T) {
@@ -207,4 +210,13 @@ func TestWriteTextNegativeCacheHits(t *testing.T) {
if !bytes.Contains(out, []byte("negative_cache_hits 0\n")) {
t.Errorf("/metrics missing negative_cache_hits 0:\n%s", out)
}
if ct := resp.Header.Get("Content-Type"); ct != "text/plain; version=0.0.4; charset=utf-8" {
t.Errorf("/metrics Content-Type=%q, want Prometheus text 0.0.4", ct)
}
if !bytes.Contains(out, []byte("# HELP negative_cache_hits ")) {
t.Errorf("/metrics missing # HELP negative_cache_hits:\n%s", out)
}
if !bytes.Contains(out, []byte("# TYPE negative_cache_hits counter")) {
t.Errorf("/metrics missing # TYPE negative_cache_hits counter:\n%s", out)
}
}
+12
View File
@@ -313,6 +313,18 @@ func TestRangeMetricsWriteText(t *testing.T) {
if !strings.Contains(text, "range_upstream 1\n") {
t.Errorf("/metrics missing 'range_upstream 1' line:\n%s", text)
}
if ct := resp.Header.Get("Content-Type"); ct != "text/plain; version=0.0.4; charset=utf-8" {
t.Errorf("/metrics Content-Type=%q, want Prometheus text 0.0.4", ct)
}
if !strings.Contains(text, "# HELP range_cache ") {
t.Errorf("/metrics missing # HELP range_cache:\n%s", text)
}
if !strings.Contains(text, "# TYPE range_cache counter") {
t.Errorf("/metrics missing # TYPE range_cache counter:\n%s", text)
}
if !strings.Contains(text, "# TYPE range_upstream counter") {
t.Errorf("/metrics missing # TYPE range_upstream counter:\n%s", text)
}
}
// TestStreamCachedResponseRange206 is a focused unit test for streamCachedResponse:
+5
View File
@@ -352,6 +352,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 {
+140
View File
@@ -562,6 +562,12 @@ func TestMetrics(t *testing.T) {
if !bytes.Contains(rec.Body.Bytes(), []byte("negative_cache_hits")) {
t.Error("WriteText output missing negative_cache_hits")
}
if !bytes.Contains(rec.Body.Bytes(), []byte("# HELP total_requests")) {
t.Error("WriteText output missing # HELP total_requests")
}
if !bytes.Contains(rec.Body.Bytes(), []byte("# TYPE total_requests counter")) {
t.Error("WriteText output missing # TYPE total_requests counter")
}
}
// Removed old TestKeyGeneration - replaced with TestURLHashing that uses SHA256
@@ -1128,6 +1134,12 @@ func TestDiskOnlyDelayedAttach(t *testing.T) {
if !bytes.Contains(rec.Body.Bytes(), []byte("capacity_pressure_events")) {
t.Errorf("WriteText output missing capacity_pressure_events: %q", rec.Body.String())
}
if !bytes.Contains(rec.Body.Bytes(), []byte("# TYPE disk_tier_ready gauge")) {
t.Errorf("WriteText output missing # TYPE disk_tier_ready gauge: %q", rec.Body.String())
}
if !bytes.Contains(rec.Body.Bytes(), []byte("# TYPE capacity_pressure_events counter")) {
t.Errorf("WriteText output missing # TYPE capacity_pressure_events counter: %q", rec.Body.String())
}
}
// TestDiskTierSignalMemoryOnly covers memory-only mode: DiskTierReady=1 (N/A, not
@@ -1485,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, "")
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, "")
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)
}
})
}